diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs index b5c50ce85..add042785 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs @@ -5087,6 +5087,15 @@ private async Task OpenLogCaptureAsync(LogCaptureCon { TeamId = context.TeamId, AgentRunId = context.RunId, ActorId = context.ActorId, WorkerFenceEpoch = context.WorkerFenceEpoch, Handle = handle, Source = source, Redactor = context.Redactor, + // A drain that waits out a destination which is refusing writes costs this run its VERDICT on a + // tear-down, in one of two ways depending on which capture is draining. For the session opened here on + // the LIVE path the token is the job's, and a drain that will not end is a tear-down that never STARTS + // — EndBrokeredAttemptOnShutdownAsync runs from the OperationCanceledException this drain is sitting + // on. For the session this very method opens again under EndBrokeredAttemptOnShutdownAsync the token + // IS ShutdownLeaseLandingBudget, and every second the drain spends is one the terminal write does not + // get. So the drain is told, with the same predicate the tear-down arm itself acts on — the HOST's own + // lifetime, never a cancelled job token. + HostShutdown = _lifetime?.ApplicationStopping ?? CancellationToken.None, }, cancellationToken).ConfigureAwait(false); } catch (Exception exception) when (exception is not OperationCanceledException || !cancellationToken.IsCancellationRequested) diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs index a6cf95524..65c02ab9a 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs @@ -16,9 +16,21 @@ public sealed class AgentRunLogCaptureBridge : IAgentRunLogCaptureBridge internal const int MaximumSegmentBytes = 1024 * 1024; private const int MaximumReadsPerPoll = 8; - private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(250); - private static readonly TimeSpan DefaultOperationTimeout = TimeSpan.FromSeconds(5); - private static readonly TimeSpan DefaultFinalizationBudget = TimeSpan.FromSeconds(30); + + /// + /// How often the loop asks the source for more. Deliberately the ONE wait in this class left on the wall clock: + /// it is a liveness cadence rather than a ceiling — nothing is decided by how many times it ticked — and keeping + /// it real is what lets a test drive the ceilings on a while the loop keeps turning. + /// Every wait that a durable outcome DOES depend on (the operation timeout, the finalization budget, the + /// transient-append backoff) is on the injected clock. + /// + /// This and the two ceilings below are the DEPLOYED cadence of every capture, and they are internal for the + /// same reason is: a test pins them (InternalsVisibleTo) so that moving a wait + /// onto another clock — which is exactly what happened here — cannot quietly change what a deployment waits. + /// + internal static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(250); + internal static readonly TimeSpan DefaultOperationTimeout = TimeSpan.FromSeconds(5); + internal static readonly TimeSpan DefaultFinalizationBudget = TimeSpan.FromSeconds(30); /// The source proved complete but its own size cap cut it short: everything captured stays readable, and the stream terminalizes Truncated rather than claiming a whole capture it knows it does not have. private static readonly CaptureFailure Truncation = new("source-truncated", "The durable sandbox log source reached its spool size cap; the captured bytes are the head of a longer output.", AgentRunLogStreamState.Truncated); @@ -55,8 +67,8 @@ internal AgentRunLogCaptureBridge(IAgentRunLogService logs, IAgentRunLogStorageR public async Task OpenAsync(AgentRunLogCaptureOpenRequest request, CancellationToken cancellationToken) { - using var operation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - operation.CancelAfter(_operationTimeout); + using var budget = new CancellationTokenSource(_operationTimeout, _clock); + using var operation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, budget.Token); var captureToken = operation.Token; try { @@ -119,8 +131,8 @@ public async Task OpenAsync(AgentRunLogCaptureOpenRe public async Task RecordGapAsync(AgentRunLogCaptureGapRequest request, CancellationToken cancellationToken) { - using var operation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - operation.CancelAfter(_operationTimeout); + using var budget = new CancellationTokenSource(_operationTimeout, _clock); + using var operation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, budget.Token); var captureToken = operation.Token; try { @@ -151,8 +163,8 @@ public async Task RecordGapAsync(AgentRunLogCaptureGapRequest request, Cancellat public async Task CompleteRunAsync(Guid teamId, Guid agentRunId, long workerFenceEpoch, CancellationToken cancellationToken) { if (teamId == Guid.Empty || agentRunId == Guid.Empty || workerFenceEpoch <= 0) return; - using var finalization = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - finalization.CancelAfter(_finalizationBudget); + using var budget = new CancellationTokenSource(_finalizationBudget, _clock); + using var finalization = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, budget.Token); var captureToken = finalization.Token; IReadOnlyList streams; try { streams = await _logs.ListCaptureHeadsAsync(teamId, agentRunId, captureToken).ConfigureAwait(false); } @@ -209,21 +221,36 @@ public async Task CompleteRunAsync(Guid teamId, Guid agentRunId, long workerFenc } } + /// + /// Tail the source into its streams until the observer finishes, then drain what is left. + /// + /// The poll wait below is the live half's ONLY cancellation point, and it has to be asked explicitly: + /// hands back a cancelled delay rather than raising it, and a source that + /// answers "nothing new" without touching the token — which is what the local spool's own read does whenever the + /// file has grown by less than one minimum segment — raises nothing either. So a capture whose observer was + /// cancelled used to spin here at full speed forever, and the tear-down waiting on this task never returned: a + /// worker with a reachable log destination hung on shutdown instead of landing its runs. + /// private async Task CaptureLoopAsync(AgentRunLogCaptureOpenRequest request, Guid captureSessionId, IReadOnlyList streams, Task finish, CancellationToken cancellationToken) { try { while (!finish.IsCompleted) { - foreach (var stream in streams.Where(value => !value.Terminal)) + foreach (var stream in streams.Where(value => value.Draining)) await PumpAsync(request, captureSessionId, stream, final: false, cancellationToken).ConfigureAwait(false); await Task.WhenAny(Task.Delay(PollInterval, cancellationToken), finish).ConfigureAwait(false); + cancellationToken.ThrowIfCancellationRequested(); } - while (streams.Any(value => !value.Terminal)) + for (var pass = 1; streams.Any(value => value.Draining); pass++) { - foreach (var stream in streams.Where(value => !value.Terminal)) + foreach (var stream in streams.Where(value => value.Draining)) await PumpAsync(request, captureSessionId, stream, final: true, cancellationToken).ConfigureAwait(false); - if (streams.Any(value => !value.Terminal)) await Task.Delay(PollInterval, cancellationToken).ConfigureAwait(false); + + if (!streams.Any(value => value.Draining)) break; + if (pass > SourceSealGracePasses && ParkedIncompleteSourcesForHostShutdown(request, streams)) break; + + await Task.Delay(PollInterval, cancellationToken).ConfigureAwait(false); } } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { } @@ -237,7 +264,7 @@ private async Task CaptureLoopAsync(AgentRunLogCaptureOpenRequest request, Guid private async Task PumpAsync(AgentRunLogCaptureOpenRequest request, Guid captureSessionId, CaptureStream stream, bool final, CancellationToken cancellationToken) { var reads = 0; - while (!stream.Terminal && (final || reads < MaximumReadsPerPoll)) + while (stream.Draining && (final || reads < MaximumReadsPerPoll)) { var drain = await FlushBacklogAsync(request, captureSessionId, stream, final, cancellationToken).ConfigureAwait(false); if (drain == DrainOutcome.Stopped) return; @@ -296,9 +323,10 @@ private async Task PumpAsync(AgentRunLogCaptureOpenRequest request, Guid capture /// is the wait; the caller may keep filling the window from the sandbox spool /// while it lasts, up to . /// - /// The FINAL drain keeps its own tighter loop and never parks: it is already bounded by the finalization - /// budget, and a stream the budget cancels stays Open and reconcilable, which is strictly better than a terminal - /// verdict a later observer could not undo. + /// The FINAL drain keeps its own tighter loop and never parks TERMINALLY: it is already bounded by the + /// finalization budget, and a stream the budget cancels stays Open and reconcilable, which is strictly better than + /// a terminal verdict a later observer could not undo. The one thing that stops it early is the host going away + /// () — and that too leaves the stream Open rather than terminalizing it. /// private async Task FlushBacklogAsync(AgentRunLogCaptureOpenRequest request, Guid captureSessionId, CaptureStream stream, bool final, CancellationToken cancellationToken) { @@ -318,7 +346,9 @@ private async Task FlushBacklogAsync(AgentRunLogCaptureOpenRequest await MarkStallAsync(request, captureSessionId, stream, stall, cancellationToken).ConfigureAwait(false); if (final) { - await Task.Delay(AppendRetryDelay(stall.Attempts), cancellationToken).ConfigureAwait(false); + if (ParkedForHostShutdown(request, stream, stall)) return DrainOutcome.Stopped; + + await Task.Delay(AppendRetryDelay(stall.Attempts), _clock, cancellationToken).ConfigureAwait(false); continue; } if (!Exhausted(stream, stall)) return DrainOutcome.Holding; @@ -329,6 +359,76 @@ private async Task FlushBacklogAsync(AgentRunLogCaptureOpenRequest return DrainOutcome.Drained; } + /// + /// Stop draining this stream because the HOST is going away — the one thing that ends a final drain before its own + /// budget does. + /// + /// What waiting costs on a tear-down, stated exactly, because it differs by which capture is draining. A + /// session opened for the LIVE run holds the job's own token, and a drain that will not end there is a + /// AgentRunExecutor tear-down arm that never STARTS — it runs from the cancellation this drain is refusing + /// to observe. The session the tear-down itself opens to re-fold the dead agent holds + /// ShutdownLeaseLandingBudget, and every second this drain spends there is a second the terminal write does + /// not get. Either way the log tail is worth less than the verdict: the tail is still in the spool and still + /// committable by a later owner, while an unlanded run is a degrade nobody can see. + /// + /// So the wait ends here, and only LOCALLY — the durable row keeps its Open state at its own fence, with the + /// stall marker already wrote, which is exactly the shape the capture recovery sweep + /// re-claims. Nothing is terminalized, so nothing a later observer could have completed is foreclosed. + /// + /// It parks on the FIRST refusal, with no grace, because a destination that refused once is overwhelmingly + /// likely to refuse again — unlike , whose wait is on a local + /// copier and therefore gets a pass. + /// + private bool ParkedForHostShutdown(AgentRunLogCaptureOpenRequest request, CaptureStream stream, RemoteStall stall) + { + if (!request.HostShutdown.IsCancellationRequested) return false; + + stream.Parked = true; + _logger.LogWarning("Agent run {RunId} log stream {StreamId} parked its final drain after {Attempts} refusal(s) ({Problem}) because this host is shutting down; it stays Open at its fence for the capture recovery sweep, so the run's own landing is not spent waiting on a destination that is refusing writes", request.AgentRunId, stream.Metadata.StreamId, stall.Attempts, stall.Code); + + return true; + } + + /// + /// How many final-drain passes a source gets to produce its seal before a host tear-down stops waiting for it. + /// ONE, which costs a quarter second of a ten-second landing and is what separates a copier that has not been + /// reaped yet from a source that is genuinely not coming: an agent this drain itself just killed + /// (AttachAsync's timeout / stalled / vanished returns) can easily still have a cat draining its + /// FIFO when the first pass reads, and parking on that first read would call a capture whose every byte landed + /// incomplete over 250ms of scheduling. + /// + private const int SourceSealGracePasses = 1; + + /// + /// The final drain's OTHER unbounded wait. A source that is not sealed yet — or that is answering retryably — + /// terminalizes nothing, so the loop above keeps polling it until the finalization budget cancels the whole + /// capture: 30 seconds of a worker that is already going away, spent on a source whose bytes a later owner can + /// read just as well. + /// + /// It is not the same bet as the append park, and it does not get the same odds. An append park + /// happens because the destination REFUSED — the next attempt is very unlikely to differ, so it parks at once. + /// Here the wait is on a local copier finishing, which is exactly the kind of thing that completes on the next + /// poll; hence , so this is consulted only from the second pass. + /// + /// The durable state the two leave differs too, deliberately. An append park lands on top of the + /// remote_stall marker already wrote, so the Room can say WHY. This one writes + /// nothing: no destination refused anything here, and a stall marker would send an operator hunting a storage + /// incident that never happened. The cost is honest and worth naming — until the capture recovery sweep's terminal + /// grace elapses, such a stream reads as "Finalizing" with no reason attached. + /// + private bool ParkedIncompleteSourcesForHostShutdown(AgentRunLogCaptureOpenRequest request, IReadOnlyList streams) + { + if (!request.HostShutdown.IsCancellationRequested) return false; + + foreach (var stream in streams.Where(value => value.Draining)) + { + stream.Parked = true; + _logger.LogWarning("Agent run {RunId} log stream {StreamId} parked its final drain with the source still incomplete because this host is shutting down; it stays Open at its fence for the capture recovery sweep", request.AgentRunId, stream.Metadata.StreamId); + } + + return true; + } + /// One attempt at the queue's head segment. The head is dequeued ONLY on a committed receipt, so nothing is ever consumed by a failure. private async Task AppendHeadAsync(AgentRunLogCaptureOpenRequest request, Guid captureSessionId, CaptureStream stream, CancellationToken cancellationToken) { @@ -550,7 +650,7 @@ private async Task FailStreamAsync(AgentRunLogCaptureOpenRequest request, Guid c private async Task FailStreamsAsync(AgentRunLogCaptureOpenRequest request, Guid captureSessionId, IReadOnlyList streams, CaptureFailure failure, CancellationToken cancellationToken) { - foreach (var stream in streams.Where(value => !value.Terminal)) + foreach (var stream in streams.Where(value => value.Draining)) await FailStreamAsync(request, captureSessionId, stream, failure, cancellationToken).ConfigureAwait(false); } @@ -620,7 +720,8 @@ private async Task DeclareExpectedStreamsAsync(AgentRunLogCaptureDeclarati private static bool Valid(AgentRunLogCaptureOpenRequest request, IReadOnlyList descriptors) => request.TeamId != Guid.Empty && request.AgentRunId != Guid.Empty && request.ActorId != Guid.Empty && request.WorkerFenceEpoch > 0 && request.Handle.AgentRunLogCaptureSessionId is { } sessionId && sessionId != Guid.Empty && descriptors.Count > 0 && descriptors.Select(value => value.SourceKey).Distinct(StringComparer.Ordinal).Count() == descriptors.Count && descriptors.Select(value => value.StreamKind).Distinct(StringComparer.Ordinal).Count() == descriptors.Count; private static bool Valid(AgentRunLogCaptureGapRequest request, IReadOnlyList descriptors) => request.TeamId != Guid.Empty && request.AgentRunId != Guid.Empty && request.WorkerFenceEpoch > 0 && request.Handle.AgentRunLogCaptureSessionId is { } sessionId && sessionId != Guid.Empty && request.ErrorCode is { Length: > 0 and <= 128 } && request.ErrorMessage is { Length: > 0 and <= 2048 } && descriptors.Count > 0 && descriptors.Select(value => value.SourceKey).Distinct(StringComparer.Ordinal).Count() == descriptors.Count && descriptors.Select(value => value.StreamKind).Distinct(StringComparer.Ordinal).Count() == descriptors.Count; private static bool IsProcessStream(string streamKind) => streamKind is AgentRunLogKinds.StandardOutput or AgentRunLogKinds.StandardError; - private static TimeSpan AppendRetryDelay(int attempt) => TimeSpan.FromMilliseconds(Math.Min(1000, 50 * Math.Pow(2, Math.Min(Math.Max(attempt - 1, 0), 5)))); + /// The wait before a refused segment's next offer on the FINAL drain: 50ms doubled per prior attempt, capped at one second. Internal so a test both PINS it and advances a virtual clock by exactly it, rather than mirroring the number and drifting. + internal static TimeSpan AppendRetryDelay(int attempt) => TimeSpan.FromMilliseconds(Math.Min(1000, 50 * Math.Pow(2, Math.Min(Math.Max(attempt - 1, 0), 5)))); private static string Held(TimeSpan value) => $"{value.TotalMinutes.ToString("F1", CultureInfo.InvariantCulture)}m"; private static string Code(T value) where T : struct, Enum => string.Concat(value.ToString().Select((character, index) => char.IsUpper(character) && index > 0 ? $"-{char.ToLowerInvariant(character)}" : char.ToLowerInvariant(character).ToString())); @@ -653,9 +754,8 @@ public async Task ObserveAsync(Func !value.Terminal)) + await DrainWithinBudgetAsync(capture, captureCts).ConfigureAwait(false); + if (_streams.Any(value => value.Draining)) _owner._logger.LogWarning("Agent run {RunId} source final drain exceeded its shadow budget and remains Open for reconciliation", _request.AgentRunId); return result; } @@ -669,11 +769,20 @@ public async Task ObserveAsync(FuncAwait the final drain under the finalization budget, on the INJECTED clock: a test advances the ceiling instead of sleeping through it, so a loaded runner can no longer spend a budget before the retry it is measuring. Production's clock is the system one, so the ceiling is the same it always was. + private async Task DrainWithinBudgetAsync(Task capture, CancellationTokenSource captureCts) + { + using var budget = new CancellationTokenSource(_owner._finalizationBudget, _owner._clock); + using var expiry = budget.Token.Register(static state => ((CancellationTokenSource)state!).Cancel(), captureCts); + + await capture.ConfigureAwait(false); + } } private sealed class NoopCaptureSession(SandboxHandle handle) : IAgentRunLogCaptureSession @@ -709,8 +818,21 @@ public CaptureStream(SandboxDurableLogDescriptor descriptor, AgentRunLogMetadata /// The live outage, or null when the remote is answering. public RemoteStall? Stall { get; set; } + + /// This stream concluded — a committed final-drain receipt, or a durable terminal verdict already written for it. public bool Terminal { get; set; } + /// + /// This stream stopped draining because the HOST is going away, NOT because anything about it concluded. + /// Distinct from on purpose: the row is still Open at its own fence for the capture + /// recovery sweep, no terminal verdict was written, and the drain's "exceeded its shadow budget" warning must + /// not claim a stream that deliberately stopped early and already said so in its own log line. + /// + public bool Parked { get; set; } + + /// Whether this stream still wants pumping: neither concluded nor parked. + public bool Draining => !Terminal && !Parked; + public void Enqueue(PendingAppend append) { Backlog.Enqueue(append); diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs index 0e24269f9..39ddc195b 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs @@ -24,6 +24,22 @@ public sealed record AgentRunLogCaptureOpenRequest public required SandboxHandle Handle { get; init; } public required ISandboxDurableLogSource Source { get; init; } public required SecretRedactor Redactor { get; init; } + + /// + /// The HOST's own tear-down signal (IHostApplicationLifetime.ApplicationStopping). Raised ⇒ the FINAL drain + /// stops waiting out a destination that is refusing writes and parks instead: the stream stays Open at its own + /// fence with its stall marker, for the capture recovery sweep to re-claim and finish. + /// + /// What the waiting costs depends on which capture is draining, and it is never nothing. A session opened + /// for the live run holds the job's own token, so a drain that will not end is a worker tear-down that never + /// STARTS — it runs from the cancellation that drain is refusing to observe. A session the tear-down itself opens + /// to re-fold the dead agent holds the worker's lease-landing budget, so every second the drain spends there is a + /// second the run's terminal write does not get. + /// + /// Default (None) means no host tear-down is in progress, which is every ordinary run: the live and final + /// retry cadences are exactly what they were. + /// + public CancellationToken HostShutdown { get; init; } } public sealed record AgentRunLogCaptureGapRequest diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs index a57b9b130..67592a95a 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs @@ -1,3 +1,4 @@ +using System.Diagnostics; using System.Net; using System.Net.Http; using System.Text.Json; @@ -16,6 +17,7 @@ using CodeSpace.IntegrationTests.Infrastructure; using CodeSpace.IntegrationTests.Workflows.Infrastructure; using CodeSpace.Messages.Agents; +using CodeSpace.Messages.Constants; using CodeSpace.Messages.Enums; using Microsoft.EntityFrameworkCore; using Shouldly; @@ -281,6 +283,13 @@ public async Task A_worker_shutting_down_ends_every_brokered_run_whose_address_i var credId = await SeedModelCredentialAsync(teamId, BrokeredProvider, "sk-shutdown-drain-fixture"); var runIds = new[] { await CreateRunWithCredentialAsync(teamId, credId), await CreateRunWithCredentialAsync(teamId, credId) }; + // A destination that can actually take the bytes. Without one the bridge resolves Unavailable, fails every + // stream at open (before a single byte is read) and hands back a passthrough session — so the capture loop + // never runs and the terminal-state claim at the end of this test is satisfied by a row written during open. + // With it, the drain below is a real capture being torn down, which is what this arm says it covers. + using var destination = new WritableLogDestination(); + await SeedAgentRunLogRouteAsync(teamId, destination.RootPath); + // A broker that records NO re-bind address — the pre-upgrade generation, and the only generation this arm // still owns. A run whose handle DOES carry one is left running for the next worker instead (see // A_worker_shutting_down_leaves_a_rebindable_brokered_run_for_the_next_worker); what remains here is the run @@ -365,6 +374,7 @@ await AwaitWithinAsync(Task.WhenAll(executions), TimeSpan.FromSeconds(120), .Where(x => x.AgentRunId == runId).Select(x => new { x.State, x.WorkerFenceEpoch }).ToListAsync(); streams.ShouldNotBeEmpty($"run {runId} opened no log stream, so every claim below is about the empty set — check that the capture bridge reached the executor rather than its passthrough branch"); + streams.Count.ShouldBe(2, $"run {runId} captured fewer than its two process streams, so this arm is not exercising the capture it claims to"); streams.ShouldAllBe(x => x.State != AgentRunLogStreamState.Open, $"run {runId} left a capture stream Open after the drain landed it terminal; nothing writes one of those afterwards (the landing leaves the fence unchanged, so even the owner-loss statement cannot match it) and the Room reports it as still finalizing forever"); @@ -379,6 +389,226 @@ await AwaitWithinAsync(Task.WhenAll(executions), TimeSpan.FromSeconds(120), } } + /// + /// The same drain, with the ONE difference an operator can actually misconfigure: the team's Agent Run log route + /// points at a destination that refuses every write. The capture bridge is the deployed one, its storage resolves + /// Ready (the route IS Active), and only the bytes have nowhere to go — so the drain's final flush meets a + /// provider that answers "retryable" forever. + /// + /// A drain that treats that as something to wait out spends the whole landing on it and lands nothing: the + /// verdict is what the operator needs, and the log tail is what they can afford to lose. So the run must still + /// land TYPED inside the budget, and the stream it could not flush stays Open at its fence for the recovery + /// sweep — the one state a later observer can still act on. + /// + [Fact] + public async Task A_worker_shutting_down_lands_its_brokered_run_within_the_budget_even_when_the_log_destination_refuses_every_write() + { + if (OperatingSystem.IsWindows()) return; + + var teamId = await SeedTeamAsync(); + var credId = await SeedModelCredentialAsync(teamId, BrokeredProvider, "sk-unwritable-log-destination-fixture"); + var runId = await CreateRunWithCredentialAsync(teamId, credId); + + using var destination = new UnwritableLogDestination(); + await SeedAgentRunLogRouteAsync(teamId, destination.RootPath); + + using var broker = new AddresslessBroker(new LoopbackModelCredentialBroker()); + using var shutdown = new CancellationTokenSource(); + var lifetime = new FakeHostLifetime(); + using var captureScope = _fixture.BeginScope(); + + var execution = ExecuteUntilShutdownAsync(runId, broker, shutdown.Token, lifetime, productionCapturePlanes: true, logCapture: captureScope.Resolve()); + + await WaitUntilAsync(() => broker.HasLease(runId), TimeSpan.FromSeconds(30), "the run never opened its credential lease"); + await WaitUntilAsync(() => HandleOf(runId) is not null, TimeSpan.FromSeconds(30), "the run never persisted a durable handle, so there was no launched agent for a shutdown to account for"); + await WaitUntilAsync(() => HasLogStream(runId), TimeSpan.FromSeconds(60), + "no agent_run_log_stream row appeared, so the capture bridge never opened one — this arm is about a destination that refuses WRITES, so a route that does not even resolve would make every claim below vacuous"); + + var pid = HandleOf(runId)!.ProcessId; + ProcessIsAlive(pid).ShouldBeTrue("precondition: the agent is alive at the moment the worker is told to go"); + + lifetime.Stop(); + shutdown.Cancel(); + await AwaitWithinAsync(execution, UnwritableDestinationDrainCeiling, + $"the draining executor never returned with an unwritable log destination — the landing is bounded by ShutdownLeaseLandingBudget, so a wait past {UnwritableDestinationDrainCeiling.TotalSeconds}s means the capture drain is retrying a dead destination instead of observing the shutdown token (diagnose by attaching to the test host and looking for AgentRunLogCaptureBridge.FlushBacklogAsync on a thread)"); + + using var scope = _fixture.BeginScope(); + var run = await scope.Resolve().GetAsync(runId, CancellationToken.None); + + run.Status.ShouldBe(AgentRunStatus.Failed, + "the run was left Running by a worker taking its model access with it — a log destination nobody can write to must not cost the operator the VERDICT as well as the log tail"); + var result = JsonSerializer.Deserialize(run.ResultJson!, AgentJson.Options).ShouldNotBeNull(); + result.ExitReason.ShouldBe(CodeSpace.Messages.Failures.FailureCodes.ModelCredentialLeaseLost, + "the landing has to name the real cause; a storage outage during the drain is not what ended this run"); + + var streams = await scope.Resolve().AgentRunLogStream.AsNoTracking() + .Where(x => x.AgentRunId == runId).Select(x => new { x.StreamKind, x.State, x.WorkerFenceEpoch, x.RemoteStallSince, x.RemoteStallCode }).ToListAsync(); + + streams.ShouldNotBeEmpty($"run {runId} opened no log stream, so every claim below is about the empty set"); + streams.ShouldAllBe(x => x.State != AgentRunLogStreamState.CaptureFailed && x.State != AgentRunLogStreamState.Truncated, + $"run {runId} stamped a terminal LOSS on a stream whose bytes the destination only refused transiently — nothing a later observer could still commit may be foreclosed by a drain that ran out of budget"); + + var held = streams.Where(x => x.StreamKind == AgentRunLogKinds.StandardOutput).ToList().ShouldHaveSingleItem($"run {runId} has no stdout log stream, so the claims below are about nothing"); + held.State.ShouldBe(AgentRunLogStreamState.Open, + $"run {runId}'s stdout stream held bytes the destination refused, so it must stay Open — that is the one state the recovery sweep can still finish"); + held.WorkerFenceEpoch.ShouldBe(run.FenceEpoch, + "the parked stream stays at the fence the landing left unchanged, which is what makes it reachable by the recovery sweep"); + held.RemoteStallSince.ShouldNotBeNull($"run {runId} parked a stream without saying why — a stream left Open with no stall marker reads as 'still finalizing' forever, and the Room has nothing to show the operator"); + held.RemoteStallCode.ShouldNotBeNull($"run {runId} recorded a stall with no code, so nothing tells an operator which destination fault held the bytes"); + + await WaitUntilAsync(() => !ProcessIsAlive(pid), TimeSpan.FromSeconds(15), $"the agent (pid {pid}) was still alive after its run was landed lease-lost; diagnose with `ps -p {pid} -o pid,stat,etime,command`"); + } + + /// + /// The park's other reachable shape, and the one the drain arm above cannot show: an agent that FINISHES while + /// the host is stopping. Its observer returns normally, so the capture takes its ordinary final-drain path — with + /// ApplicationStopping already raised and a destination that refuses every write. + /// + /// Nothing here is cancelled: the job token stays live and the run lands its own ordinary verdict. What is + /// under test is the seconds between the agent exiting and that verdict being written, which a drain that waits + /// the destination out spends on its full finalization budget while the process is being torn down around it. + /// + [Fact] + public async Task A_run_that_finishes_while_the_host_is_stopping_parks_its_capture_instead_of_draining_to_the_budget() + { + if (OperatingSystem.IsWindows()) return; + + var teamId = await SeedTeamAsync(); + var credId = await SeedModelCredentialAsync(teamId, BrokeredProvider, "sk-finish-during-stop-fixture"); + var runId = await CreateRunWithCredentialAsync(teamId, credId); + + using var destination = new UnwritableLogDestination(); + await SeedAgentRunLogRouteAsync(teamId, destination.RootPath); + + using var release = new TempDir(); + var releaseFile = Path.Combine(release.Path, "release"); + using var broker = new LoopbackModelCredentialBroker(); + var lifetime = new FakeHostLifetime(); + using var captureScope = _fixture.BeginScope(); + + var harness = new BrokerableScriptedHarness(BrokeredProvider, $"echo '{ShutdownFactLine}'; while [ ! -f '{releaseFile}' ]; do sleep 0.2; done; echo done"); + var execution = ExecuteAsync(runId, harness, logCapture: captureScope.Resolve(), credentialBroker: broker, lifetime: lifetime, productionCapturePlanes: true); + + await WaitUntilAsync(() => HasLogStream(runId), TimeSpan.FromSeconds(60), + "no agent_run_log_stream row appeared, so the capture bridge never opened one and every claim below is about the empty set"); + + // The host announces it is stopping while the agent is still working, and only THEN does the agent finish — + // so the final drain below is the first thing in this run to see ApplicationStopping raised. + lifetime.Stop(); + await File.WriteAllTextAsync(releaseFile, "go"); + + var watch = Stopwatch.StartNew(); + await AwaitWithinAsync(execution, ParkedDrainCeiling, + $"the run did not land within {ParkedDrainCeiling.TotalSeconds}s of its agent exiting — the capture drain is waiting out a destination that refuses every write instead of parking, which on a stopping host costs the whole 30s finalization budget"); + + using var scope = _fixture.BeginScope(); + var run = await scope.Resolve().GetAsync(runId, CancellationToken.None); + + run.Status.ShouldBe(AgentRunStatus.Succeeded, "the agent exited 0; a log destination nobody can write to is not this run's verdict"); + watch.Elapsed.ShouldBeLessThan(ParkedDrainCeiling, "and it landed without waiting the destination out"); + + var streams = await scope.Resolve().AgentRunLogStream.AsNoTracking() + .Where(x => x.AgentRunId == runId).Select(x => new { x.StreamKind, x.State, x.WorkerFenceEpoch, x.RemoteStallSince, x.RemoteStallCode }).ToListAsync(); + + var held = streams.Where(x => x.StreamKind == AgentRunLogKinds.StandardOutput).ToList().ShouldHaveSingleItem($"run {runId} has no stdout log stream"); + held.State.ShouldBe(AgentRunLogStreamState.Open, "the parked stream stays Open — the one state the capture recovery sweep can still finish"); + held.WorkerFenceEpoch.ShouldBe(run.FenceEpoch, "at the fence the run still holds, which is what makes it reachable by that sweep"); + held.RemoteStallSince.ShouldNotBeNull($"run {runId} parked without saying why; an Open stream with no stall marker reads as \"still finalizing\" forever"); + held.RemoteStallCode.ShouldNotBeNull(); + streams.ShouldAllBe(x => x.State != AgentRunLogStreamState.CaptureFailed, + "nothing was permanently refused, so no stream may carry a terminal loss a later owner could have avoided"); + } + + /// How long after its agent exits a run may take to land while the host is stopping. Far below the bridge's own 30s finalization budget, which is exactly what a drain that waits an unwritable destination out would spend. + private static readonly TimeSpan ParkedDrainCeiling = TimeSpan.FromSeconds(20); + + /// + /// How long the whole tear-down may take when the capture destination refuses every write. Well above the + /// ShutdownLeaseLandingBudget the landing itself is bounded by (so a loaded runner does not fail it) and + /// far below the 120s the other arms allow, because the point of this arm is the DIFFERENCE between a bounded + /// drain and one that waits out a provider that is never coming back. + /// + private static readonly TimeSpan UnwritableDestinationDrainCeiling = TimeSpan.FromSeconds(45); + + /// A destination that takes the bytes — an ordinary empty directory the local-rwx provider can write under, cleaned up with the test. + private sealed class WritableLogDestination : IDisposable + { + private readonly string _root = Path.Combine(Path.GetTempPath(), "cs-log-dest-" + Guid.NewGuid().ToString("N")); + + public WritableLogDestination() => Directory.CreateDirectory(_root); + + public string RootPath => _root; + + public void Dispose() + { + try { Directory.Delete(_root, recursive: true); } catch { /* best-effort */ } + } + } + + /// A destination that resolves and activates exactly like a real one, and whose root can never be created: its parent is a regular FILE, so every Directory.CreateDirectory under it is an IOException the provider reports as retryable. + private sealed class UnwritableLogDestination : IDisposable + { + private readonly string _root = Path.Combine(Path.GetTempPath(), "cs-dead-log-dest-" + Guid.NewGuid().ToString("N")); + + public UnwritableLogDestination() + { + Directory.CreateDirectory(_root); + File.WriteAllText(Path.Combine(_root, "blocked"), "not a directory"); + } + + public string RootPath => Path.Combine(_root, "blocked", "store"); + + public void Dispose() + { + try { Directory.Delete(_root, recursive: true); } catch { /* best-effort */ } + } + } + + /// Route this team's Agent Run log data class at — an Active route over an Active profile, so the resolver answers Ready and only the WRITE fails. + private async Task SeedAgentRunLogRouteAsync(Guid teamId, string rootPath) + { + using var scope = _fixture.BeginScope(); + var db = scope.Resolve(); + var now = DateTimeOffset.UtcNow; + var profileId = Guid.NewGuid(); + using var document = JsonDocument.Parse(JsonSerializer.Serialize(new { rootPath })); + var canonicalConfig = CodeSpace.Core.Services.Workflows.Artifacts.Profiles.StorageProfileRules.CanonicalJson(document.RootElement); + using var canonical = JsonDocument.Parse(canonicalConfig); + + var profile = new StorageProfile + { + Id = profileId, TeamId = teamId, StableName = $"dead-log-dest-{profileId:N}", State = StorageProfileState.Active, + CurrentRevision = 1, CreatedDate = now, CreatedBy = SystemUsers.SeederId, LastModifiedDate = now, LastModifiedBy = SystemUsers.SeederId, + }; + profile.Revisions.Add(new StorageProfileRevision + { + Id = Guid.NewGuid(), TeamId = teamId, StorageProfileId = profileId, Revision = 1, + ProviderTypeKey = LocalRwxArtifactStorageDriverFactory.TypeKey, NonSecretConfigJson = canonicalConfig, CredentialRef = null, + NamespaceFingerprint = CodeSpace.Core.Services.Workflows.Artifacts.Profiles.StorageProfileRules.NamespaceFingerprint(LocalRwxArtifactStorageDriverFactory.TypeKey, canonical.RootElement), + CreatedDate = now, CreatedBy = SystemUsers.SeederId, + }); + db.StorageProfile.Add(profile); + + var route = new StorageRoute + { + Id = Guid.NewGuid(), TeamId = teamId, DataClassTypeKey = AgentRunLogStorageResolver.DataClassTypeKey, + CurrentRevision = 1, State = StorageRouteState.Draft, CreatedDate = now, CreatedBy = SystemUsers.SeederId, + LastModifiedDate = now, LastModifiedBy = SystemUsers.SeederId, + }; + route.Revisions.Add(new StorageRouteRevision + { + Id = Guid.NewGuid(), TeamId = teamId, StorageRouteId = route.Id, Revision = 1, + StorageProfileId = profileId, ProfileRevisionMode = StorageProfileRevisionMode.CurrentAtWrite, + CreatedDate = now, CreatedBy = SystemUsers.SeederId, + }); + db.StorageRoute.Add(route); + await db.SaveChangesAsync(); + + route.State = StorageRouteState.Active; + route.LastModifiedDate = DateTimeOffset.UtcNow; + await db.SaveChangesAsync(); + } + [Fact] public async Task A_worker_shutting_down_leaves_a_rebindable_brokered_run_for_the_next_worker() { diff --git a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs index 25994d889..d9d2681e1 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs @@ -186,7 +186,14 @@ public async Task An_outage_that_begins_during_the_final_drain_still_says_so_wit var bridge = Bridge(logs, clock, stalls, gaps, finalizationBudget: TimeSpan.FromSeconds(2)); var capture = await bridge.OpenAsync(Request(source), CancellationToken.None); - await capture.ObserveAsync((_, _) => { logs.RemoteUnavailable = true; return Task.FromResult(Result()); }, CancellationToken.None); + var observing = capture.ObserveAsync((_, _) => { logs.RemoteUnavailable = true; return Task.FromResult(Result()); }, CancellationToken.None); + + // Both of the drain's waits — the transient-append backoff and the finalization budget — are on THIS clock, so + // the drain is retrying the outage and will keep doing so until the test says the budget is spent. That is + // what makes the assertions below about the drain's behaviour rather than about how loaded the runner is. + await WaitAsync(() => stalls.Held.Count > 0, "the bridge never marked the outage that started on the final drain"); + clock.Advance(TimeSpan.FromSeconds(2)); + await observing; stalls.Held.ShouldNotBeEmpty("an outage that starts on the final drain is the same outage; the Room cannot read \"Finalizing\" through it"); stalls.Held[0].StallCode.ShouldBe("capture-backend-unavailable"); @@ -196,6 +203,132 @@ public async Task An_outage_that_begins_during_the_final_drain_still_says_so_wit gaps.Gaps.ShouldBeEmpty("a marker is not a park, and only a park may declare a span lost"); } + [Fact] + public async Task A_drain_under_a_host_that_is_going_away_parks_on_the_first_refusal_instead_of_spending_the_landing_budget() + { + // The drain and the run's own terminal write share ONE deadline on a worker tear-down, so a destination that + // is refusing writes must not be waited out here: the verdict an operator needs costs seconds this retry + // would otherwise spend, and the log tail it buys is recoverable while the verdict is not. Every wait the + // drain could take is on this clock, so "did not wait" is observable rather than inferred from elapsed time. + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + var logs = new FakeLogService { CurrentFence = 1 }; + var stalls = new FakeStallWriter(); + var gaps = new FakeCompletenessWriter(); + var source = new FakeLogSource(); + source.Set("stdout", Payload(null, 4096)); // under one minimum segment: nothing is offered until the FINAL drain + source.Set("stderr", []); + var bridge = Bridge(logs, clock, stalls, gaps, finalizationBudget: TimeSpan.FromSeconds(30)); + using var shutdown = new CancellationTokenSource(); + shutdown.Cancel(); + + var capture = await bridge.OpenAsync(Request(source, hostShutdown: shutdown.Token), CancellationToken.None); + var observing = capture.ObserveAsync((_, _) => { logs.RemoteUnavailable = true; return Task.FromResult(Result()); }, CancellationToken.None); + + // Bounded and loud: with no clock to advance, a drain that decides to WAIT instead of parking waits on virtual + // time nobody is going to spend, and an unbounded await would report that as silence (Rule 12.10). + await AwaitWithinAsync(observing, "the drain never returned, so it is waiting out the refusal on a clock this test never advances — the shutdown park did not fire"); + + clock.GetUtcNow().ShouldBe(DateTimeOffset.UnixEpoch, + "the drain waited for NOTHING — on a tear-down every second it waits either delays the landing's start or is taken straight out of the landing's own budget, depending on which capture is draining"); + logs.AppendAttempts.ShouldBe(1, "the segment was offered once and then parked; a second offer means the tear-down is still waiting the destination out"); + var stdout = logs.Head(AgentRunLogKinds.StandardOutput).Metadata; + stdout.State.ShouldBe(AgentRunLogStreamState.Open, "parking for a shutdown is local: the row stays Open at its own fence, which is the one state the recovery sweep can still finish"); + stdout.ErrorCode.ShouldBeNull("nothing was permanently refused, so no terminal cause may be invented"); + stalls.Held.ShouldNotBeEmpty("a stream left Open with no stall marker reads as \"still finalizing\" forever"); + gaps.Gaps.ShouldBeEmpty("a park for shutdown is not a declared loss: the held span is still in the spool and still committable by whoever finishes the stream"); + } + + [Fact] + public async Task A_drain_whose_source_never_seals_also_parks_for_a_host_that_is_going_away() + { + // The final drain has TWO ways to wait forever, and a refusal is only one of them. A source that never seals + // terminalizes nothing, so the drain loop keeps polling it until the finalization budget cancels the whole + // capture — 30 seconds, in production, of a worker that is already leaving. Nothing is refused here, so there + // is no stall marker to park behind: the park has to come from the loop itself. + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + var logs = new FakeLogService { CurrentFence = 1 }; + var stalls = new FakeStallWriter(); + var gaps = new FakeCompletenessWriter(); + var source = new FakeLogSource { EmitEndOfSource = false }; + source.Set("stdout", Payload(null, 4096)); + source.Set("stderr", []); + var bridge = Bridge(logs, clock, stalls, gaps, finalizationBudget: TimeSpan.FromSeconds(30)); + using var shutdown = new CancellationTokenSource(); + shutdown.Cancel(); + + var capture = await bridge.OpenAsync(Request(source, hostShutdown: shutdown.Token), CancellationToken.None); + var observing = capture.ObserveAsync((_, _) => Task.FromResult(Result()), CancellationToken.None); + + await AwaitWithinAsync(observing, "the drain never returned, so it is still polling a source that will never seal — on this clock only the finalization budget could end it, and this test never advances one"); + + clock.GetUtcNow().ShouldBe(DateTimeOffset.UnixEpoch, "the drain waited for nothing: a worker that is leaving must not spend its finalization budget on a source a later owner can read just as well"); + var stdout = logs.Head(AgentRunLogKinds.StandardOutput).Metadata; + stdout.State.ShouldBe(AgentRunLogStreamState.Open, "the row stays Open at its fence for the capture recovery sweep"); + stdout.ErrorCode.ShouldBeNull("nothing failed — the host left; inventing a terminal cause would foreclose what a later owner can still finish"); + stalls.Held.ShouldBeEmpty("nothing was refused, so nothing may claim the destination is stalled"); + gaps.Gaps.ShouldBeEmpty("a park for shutdown is not a declared loss"); + } + + [Fact] + public async Task A_source_that_seals_one_poll_late_is_finalized_rather_than_parked_when_the_host_is_stopping() + { + // The park above must not be an impatience. The drain's own kill is what usually ends a source, and the copier + // draining that agent's FIFO can easily still be running when the first final pass reads — so a park with no + // grace would call a capture whose every byte landed incomplete, and the recovery sweep would later stamp it + // CaptureFailed rather than re-read a spool it does not re-read. One pass of grace is the difference, and it + // costs a quarter second of a ten-second landing. + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + var logs = new FakeLogService { CurrentFence = 1 }; + var stalls = new FakeStallWriter(); + var gaps = new FakeCompletenessWriter(); + var source = new FakeLogSource { SealAfterFinalReads = 1 }; + source.Set("stdout", Payload(null, 4096)); + source.Set("stderr", []); + var bridge = Bridge(logs, clock, stalls, gaps, finalizationBudget: TimeSpan.FromSeconds(30)); + using var shutdown = new CancellationTokenSource(); + shutdown.Cancel(); + + var capture = await bridge.OpenAsync(Request(source, hostShutdown: shutdown.Token), CancellationToken.None); + var observing = capture.ObserveAsync((_, _) => Task.FromResult(Result()), CancellationToken.None); + + await AwaitWithinAsync(observing, "the drain never returned while waiting one poll for a seal that does arrive"); + await bridge.CompleteRunAsync(TeamId, RunId, 1, CancellationToken.None); + + var stdout = logs.Head(AgentRunLogKinds.StandardOutput); + stdout.CaptureFinalizedAt.ShouldNotBeNull( + "the source sealed on the second pass and the drain was still there to take the receipt — park it on the first and this capture is reported incomplete for 250ms of scheduling"); + stdout.Metadata.State.ShouldBe(AgentRunLogStreamState.Completed, "a capture with a final-drain receipt terminalizes Completed, not left for the sweep to fail"); + logs.Bytes(AgentRunLogKinds.StandardOutput).Length.ShouldBe(4096, "and every byte is there"); + stalls.Held.ShouldBeEmpty("nothing was refused; a stall marker here would send an operator hunting a storage incident that never happened"); + } + + [Fact] + public async Task A_drain_on_a_healthy_host_still_waits_out_the_same_refusal() + { + // The other half of the arm above, and the thing that keeps it from being a blanket "stop retrying": with no + // host tear-down the SAME refusal is waited out exactly as before, which is what makes the transient outage + // survivable for an ordinary run (#1969's cadence is untouched). + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + var logs = new FakeLogService { CurrentFence = 1 }; + var stalls = new FakeStallWriter(); + var source = new FakeLogSource(); + source.Set("stdout", Payload(null, 4096)); + source.Set("stderr", []); + var bridge = Bridge(logs, clock, stalls, finalizationBudget: TimeSpan.FromSeconds(30)); + + var capture = await bridge.OpenAsync(Request(source), CancellationToken.None); + var observing = capture.ObserveAsync((_, _) => { logs.RemoteUnavailable = true; return Task.FromResult(Result()); }, CancellationToken.None); + + await WaitAsync(() => logs.AppendAttempts == 1, "the drain never offered the segment at all"); + logs.RemoteUnavailable = false; + await AdvanceUntilAsync(clock, observing, "the refused segment was never offered again, so an ordinary run's drain is no longer waiting a transient outage out"); + await observing; + + logs.AppendAttempts.ShouldBe(2, "an ordinary run's final drain still offers a transiently refused segment again"); + logs.Bytes(AgentRunLogKinds.StandardOutput).Length.ShouldBe(4096, "and the wait was worth making: every held byte landed"); + logs.Head(AgentRunLogKinds.StandardOutput).CaptureFinalizedAt.ShouldNotBeNull(); + } + [Fact] public async Task A_permanent_refusal_still_terminalizes_immediately_instead_of_waiting_out_a_park_window() { @@ -283,11 +416,11 @@ private static AgentRunLogCaptureBridge Bridge(FakeLogService logs, FakeTimeProv new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(200), finalizationBudget ?? TimeSpan.FromSeconds(5)) { Backpressure = backpressure ?? CaptureBackpressureOptions.Default }, stalls, gaps ?? new FakeCompletenessWriter(), clock); - private static AgentRunLogCaptureOpenRequest Request(ISandboxDurableLogSource source, SecretRedactor? redactor = null) => new() + private static AgentRunLogCaptureOpenRequest Request(ISandboxDurableLogSource source, SecretRedactor? redactor = null, CancellationToken hostShutdown = default) => new() { TeamId = TeamId, AgentRunId = RunId, ActorId = ActorId, WorkerFenceEpoch = 1, Handle = new SandboxHandle { Kind = "fake", ProcessId = 1, SpoolDirectory = "/opaque", Deadline = DateTimeOffset.MaxValue, AgentRunLogCaptureSessionId = Guid.NewGuid() }, - Source = source, Redactor = redactor ?? SecretRedactor.None, + Source = source, Redactor = redactor ?? SecretRedactor.None, HostShutdown = hostShutdown, }; private static SandboxResult Result() => new() { Status = SandboxStatus.Success, ExitCode = 0, Stdout = "legacy", Stderr = "legacy-error" }; @@ -302,6 +435,31 @@ private static byte[] Payload(string? secret, int length) return bytes; } + /// Await work that must settle on its own — no clock to advance — with a deadline and a message naming what did not (Rule 12.10). + private static async Task AwaitWithinAsync(Task work, string signal) + { + if (await Task.WhenAny(work, Task.Delay(Patience)).ConfigureAwait(false) != work) + throw new Xunit.Sdk.XunitException($"{signal} (waited {Patience.TotalSeconds:F0}s)"); + + await work.ConfigureAwait(false); + } + + /// + /// Push the virtual clock forward in small steps until settles. Stepped rather than + /// jumped because a single jump can land in the window between a wait being decided on and its timer being armed, + /// which leaves the timer due AFTER the jump and the test waiting on a moment that never comes. + /// + private static async Task AdvanceUntilAsync(FakeTimeProvider clock, Task work, string signal) + { + var watch = Stopwatch.StartNew(); + while (!work.IsCompleted) + { + if (watch.Elapsed > Patience) throw new Xunit.Sdk.XunitException($"{signal} (advanced the virtual clock to {clock.GetUtcNow():O} over {Patience.TotalSeconds:F0}s of real time)"); + clock.Advance(TimeSpan.FromMilliseconds(10)); + await Task.Delay(5); + } + } + /// Explicit timeout with the watched signal named, so a failure says what never happened rather than only that time ran out (Rule 12.10). private static async Task WaitAsync(Func condition, string signal) { @@ -371,12 +529,21 @@ public IReadOnlyList DescribeLogs(SandboxHandle han new("stderr", AgentRunLogKinds.StandardError, AgentRunLogRepresentations.PlainTextContentType, AgentRunLogRepresentations.Utf8ContentEncoding, "fake-spool/v1"), ]; + /// False models a spool the copier has not sealed yet: the final drain gets "not yet" forever, which is the OTHER way it can run to its budget without anything being refused. + public bool EmitEndOfSource { get; init; } = true; + + /// How many final-drain reads AT THE END of each source answer "not yet" before its seal appears — the copier this drain's own kill has not reaped yet. Counted per source key, because the two streams are drained independently. + public int SealAfterFinalReads { get; init; } + + private readonly ConcurrentDictionary _finalReadsAtEnd = new(StringComparer.Ordinal); + + /// Does NOT observe the token, because production's does not: the local spool's read answers "nothing new" out of a length comparison, with no I/O and no cancellation check. A fake that threw here is what kept a capture-loop cancellation bug out of reach of this suite. public Task ReadAsync(SandboxDurableLogReadRequest request, CancellationToken cancellationToken) { - cancellationToken.ThrowIfCancellationRequested(); var bytes = _sources[request.SourceKey]; var available = bytes.LongLength - request.OffsetBytes; - if (available == 0 && request.FinalDrain) return Task.FromResult(new SandboxDurableLogReadResult.EndOfSource()); + if (available == 0 && request.FinalDrain && EmitEndOfSource && _finalReadsAtEnd.AddOrUpdate(request.SourceKey, 1, (_, seen) => seen + 1) > SealAfterFinalReads) + return Task.FromResult(new SandboxDurableLogReadResult.EndOfSource()); if (available == 0 || (!request.FinalDrain && available < request.MinimumBytes)) return Task.FromResult(new SandboxDurableLogReadResult.NoData()); var length = (int)Math.Min(available, request.MaximumBytes); return Task.FromResult(new SandboxDurableLogReadResult.Available(bytes.AsMemory((int)request.OffsetBytes, length))); diff --git a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs index 84700f980..796e1f230 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs @@ -8,6 +8,7 @@ using CodeSpace.Core.Services.Agents.Sandbox; using CodeSpace.Messages.Agents; using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Time.Testing; using Shouldly; namespace CodeSpace.UnitTests.Agents; @@ -89,7 +90,8 @@ public async Task Cancelled_observer_leaves_open_source_and_higher_fence_reattac var first = await bridge.OpenAsync(Request(source, 1, sessionId, redactor), CancellationToken.None); using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(600))) { - await Should.ThrowAsync(() => first.ObserveAsync(async (_, token) => { await Task.Delay(Timeout.InfiniteTimeSpan, token); return Result(); }, cts.Token)); + await AwaitCancelledWithinAsync(first.ObserveAsync(async (_, token) => { await Task.Delay(Timeout.InfiniteTimeSpan, token); return Result(); }, cts.Token), + "the first capture session never returned after its observer was cancelled, so nothing below about the re-attach can run"); } logs.Heads.Single(head => head.Metadata.StreamKind == AgentRunLogKinds.StandardOutput).CaptureFinalizedAt.ShouldBeNull(); @@ -106,6 +108,38 @@ public async Task Cancelled_observer_leaves_open_source_and_higher_fence_reattac logs.Heads.Single(head => head.Metadata.StreamKind == AgentRunLogKinds.StandardOutput).CaptureFinalizedAt.ShouldNotBeNull(); } + [Fact] + public async Task A_cancelled_observer_stops_the_capture_loop_even_though_the_source_never_looks_at_the_token() + { + // The worker-shutdown shape, reduced to the two facts that produce it: an observer that is cancelled, and a + // source whose "nothing new" answer touches no token (production's local spool read, whenever the file has + // grown by less than one minimum segment). The capture loop's poll wait was its only cancellation point, and + // Task.WhenAny hands back a cancelled delay rather than raising it — so this used to spin at full speed + // forever, and the tear-down that was awaiting the capture never returned. It is bounded and loud here + // because an unbounded await on it reports a wedged drain as silence (Rule 12.10). + var logs = new FakeLogService { CurrentFence = 1 }; + var source = new FakeLogSource(); + source.Set("stdout", Enumerable.Repeat((byte)'q', 4096).ToArray()); // under one minimum segment: every live read answers NoData + source.Set("stderr", []); + var session = await Bridge(logs).OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None); + using var cancelled = new CancellationTokenSource(); + + var observing = session.ObserveAsync(async (_, token) => { cancelled.Cancel(); await Task.Delay(Timeout.InfiniteTimeSpan, token); return Result(); }, cancelled.Token); + + await AwaitCancelledWithinAsync(observing, "the capture session never returned after its observer was cancelled — its loop is consuming the cancellation instead of observing it, which is what wedges a worker on shutdown"); + logs.Heads.ShouldAllBe(head => head.Metadata.State == AgentRunLogStreamState.Open && head.CaptureFinalizedAt == null, + "a cancelled capture leaves its streams Open and reconcilable; it may not invent a terminal verdict for bytes that are still in the spool"); + } + + /// Expect to end in a cancellation, within a deadline, with a message naming what did not happen (Rule 12.10). + private static async Task AwaitCancelledWithinAsync(Task work, string signal) + { + if (await Task.WhenAny(work, Task.Delay(Patience)).ConfigureAwait(false) != work) + throw new Xunit.Sdk.XunitException($"{signal} (waited {Patience.TotalSeconds:F0}s)"); + + await Should.ThrowAsync(() => work); + } + [Fact] public async Task Stale_fence_and_deterministic_corruption_never_change_harness_result_and_failure_is_durable() { @@ -218,18 +252,21 @@ public async Task Capture_health_reloads_the_monotonic_head_after_a_concurrent_r [Fact] public async Task Blocking_capture_backend_is_cancelled_by_one_total_shadow_budget_without_changing_the_sandbox_result() { + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); var logs = new FakeLogService { CurrentFence = 1, BlockAppend = true }; var source = new FakeLogSource(); source.Set("stdout", Enumerable.Repeat((byte)'x', 300 * 1024).ToArray()); source.Set("stderr", []); - var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance, new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromMilliseconds(150))); + var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance, new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromMilliseconds(150)), clock: clock); var expected = Result(); - var watch = Stopwatch.StartNew(); var session = await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None); - var observed = await session.ObserveAsync((_, _) => Task.FromResult(expected), CancellationToken.None); + var observing = session.ObserveAsync((_, _) => Task.FromResult(expected), CancellationToken.None); + var spent = await AdvanceUntilAsync(clock, observing, "the provider never completes, so only the finalization budget can end this drain — it did not"); + var observed = await observing; - watch.Elapsed.ShouldBeLessThan(TimeSpan.FromSeconds(2), "one total finalization budget bounds a provider that never completes, instead of N segments multiplying the provider default"); + spent.ShouldBeGreaterThanOrEqualTo(TimeSpan.FromMilliseconds(150), "the drain gave the provider its whole budget before giving up"); + spent.ShouldBeLessThan(TimeSpan.FromMilliseconds(150) + AdvanceStep * 4, "ONE total finalization budget bounds a provider that never completes — N segments must not multiply it. The slack is a few advance steps, not a second budget: doubling the ceiling would still fail this."); observed.ShouldBeSameAs(expected); logs.Heads.ShouldAllBe(head => head.Metadata.State == AgentRunLogStreamState.Open && head.CaptureFinalizedAt == null, "timeout remains durably reconcilable Open state, never incomplete-but-Completed"); logs.ObservedOperationTimeout.ShouldBe(TimeSpan.FromMilliseconds(40)); @@ -238,23 +275,56 @@ public async Task Blocking_capture_backend_is_cancelled_by_one_total_shadow_budg [Fact] public async Task Final_drain_retries_one_transient_append_and_preserves_the_complete_source() { + // The backoff (50ms) and the budget it has to fit inside (300ms) are BOTH on this clock. Measured on the wall + // clock they were 50ms and 300ms of real time, so a loaded runner could spend the budget before the one retry + // fired and the drain would honestly report a single attempt — the test failing for a reason that had nothing + // to do with the bridge. Virtual time cannot be spent by load; only the advance below spends it. + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); var logs = new FakeLogService { CurrentFence = 1, RetryableAppendFailures = 1 }; var source = new FakeLogSource(); source.Set("stdout", "eventually-durable"u8.ToArray()); source.Set("stderr", []); var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance, - new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromMilliseconds(300))); + new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), FinalizationBudget), clock: clock); var expected = Result(); - var observed = await (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) + var observing = (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) .ObserveAsync((_, _) => Task.FromResult(expected), CancellationToken.None); + var spent = await AdvanceUntilAsync(clock, observing, "the refused segment was never offered again — the drain is not retrying a transient refusal at all"); + var observed = await observing; + observed.ShouldBeSameAs(expected); - logs.AppendAttempts.ShouldBeGreaterThanOrEqualTo(2); + logs.AppendAttempts.ShouldBe(2, "the refused segment is offered again exactly once — the first attempt plus the retry this test advanced the clock to"); + spent.ShouldBeGreaterThanOrEqualTo(FirstAppendBackoff, "the retry waits the backoff out rather than hammering the destination"); + spent.ShouldBeLessThan(FinalizationBudget, "and it lands INSIDE the finalization budget — shrink the budget below the backoff and this drain is cancelled with one attempt, which is what the wall clock used to do at random"); logs.Bytes(AgentRunLogKinds.StandardOutput).ShouldBe("eventually-durable"u8.ToArray()); logs.Heads.Single(value => value.Metadata.StreamKind == AgentRunLogKinds.StandardOutput).CaptureFinalizedAt.ShouldNotBeNull(); } + /// The wait before a refused segment's SECOND offer, asked of production rather than mirrored, so the bound below stays true after the backoff changes. + private static readonly TimeSpan FirstAppendBackoff = AgentRunLogCaptureBridge.AppendRetryDelay(1); + + /// The ceiling that retry has to fit inside. Narrowed from production's 30s only to keep the virtual clock's travel short; what matters is that it is comfortably longer than one backoff. + private static readonly TimeSpan FinalizationBudget = TimeSpan.FromMilliseconds(300); + + /// + /// The cadence every deployed capture actually runs at. Pinned to literals on purpose (Rule 8): moving a wait onto + /// another clock, or "tidying" a timer, is a one-line edit that changes how long every worker in the fleet spends + /// on a drain — and nothing else in the suite would notice, because every other test supplies its own narrowed + /// options. The numbers below are the ones that shipped; changing one means changing this line in the same PR. + /// + [Fact] + public void The_deployed_capture_cadence_is_unchanged() + { + AgentRunLogCaptureBridge.PollInterval.ShouldBe(TimeSpan.FromMilliseconds(250), "how often a live capture asks its source for more"); + AgentRunLogCaptureBridge.DefaultOperationTimeout.ShouldBe(TimeSpan.FromSeconds(5), "one capture metadata/storage call's ceiling"); + AgentRunLogCaptureBridge.DefaultFinalizationBudget.ShouldBe(TimeSpan.FromSeconds(30), "the whole final drain's ceiling — what a worker spends landing a run's log tail"); + AgentRunLogCaptureBridge.AppendRetryDelay(1).ShouldBe(TimeSpan.FromMilliseconds(50), "the first backoff after a transient refusal on the final drain"); + AgentRunLogCaptureBridge.AppendRetryDelay(2).ShouldBe(TimeSpan.FromMilliseconds(100), "and it doubles"); + AgentRunLogCaptureBridge.AppendRetryDelay(99).ShouldBe(TimeSpan.FromSeconds(1), "and stops doubling at one second, so a long outage costs attempts rather than latency"); + } + [Fact] public async Task A_non_retryable_append_rejection_terminalizes_the_stream_instead_of_retrying_to_the_budget() { @@ -266,16 +336,16 @@ public async Task A_non_retryable_append_rejection_terminalizes_the_stream_inste var source = new FakeLogSource(); source.Set("stdout", Enumerable.Repeat((byte)'p', 300 * 1024).ToArray()); source.Set("stderr", []); + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance, - new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromSeconds(4))); + new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromSeconds(4)), clock: clock); var expected = Result(); - var watch = Stopwatch.StartNew(); var observed = await (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) .ObserveAsync((_, _) => Task.FromResult(expected), CancellationToken.None); observed.ShouldBeSameAs(expected); - watch.Elapsed.ShouldBeLessThan(TimeSpan.FromSeconds(3), "a permanent rejection must not be retried until the finalization budget cancels the capture"); + clock.GetUtcNow().ShouldBe(DateTimeOffset.UnixEpoch, "a permanent rejection must not be retried at all: the drain settled without spending one tick of its budget, and nothing but this test can spend one"); var stdout = logs.Heads.Single(value => value.Metadata.StreamKind == AgentRunLogKinds.StandardOutput); stdout.Metadata.State.ShouldBe(AgentRunLogStreamState.CaptureFailed); stdout.Metadata.ErrorCode.ShouldBe("capture-backend-unavailable"); @@ -289,14 +359,18 @@ public async Task Final_drain_budget_exhaustion_on_transient_append_leaves_an_op var source = new FakeLogSource(); source.Set("stdout", "still-in-native-spool"u8.ToArray()); source.Set("stderr", []); + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), recovery, NullLogger.Instance, - new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(20), TimeSpan.FromMilliseconds(120))); + new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(20), TimeSpan.FromMilliseconds(120)), clock: clock); var expected = Result(); - var observed = await (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) + var observing = (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) .ObserveAsync((_, _) => Task.FromResult(expected), CancellationToken.None); + var spent = await AdvanceUntilAsync(clock, observing, "a destination that refuses every offer can only be given up on by the finalization budget — it never was"); + var observed = await observing; observed.ShouldBeSameAs(expected); + spent.ShouldBeGreaterThanOrEqualTo(TimeSpan.FromMilliseconds(120), "the drain kept offering the segment until its budget, rather than giving up on the first refusal"); var stdout = logs.Heads.Single(value => value.Metadata.StreamKind == AgentRunLogKinds.StandardOutput); stdout.Metadata.State.ShouldBe(AgentRunLogStreamState.Open); stdout.CaptureFinalizedAt.ShouldBeNull(); @@ -311,11 +385,15 @@ public async Task Transient_no_data_never_authorizes_finalization_and_returns_on var source = new FakeLogSource { EmitEndOfSource = false }; source.Set("stdout", []); source.Set("stderr", []); - var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance, new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromMilliseconds(100))); + var clock = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + var bridge = new AgentRunLogCaptureBridge(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance, new AgentRunLogCaptureBridgeOptions(TimeSpan.FromMilliseconds(40), TimeSpan.FromMilliseconds(100)), clock: clock); var expected = Result(); - var observed = await (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) + var observing = (await bridge.OpenAsync(Request(source, 1, Guid.NewGuid()), CancellationToken.None)) .ObserveAsync((_, _) => Task.FromResult(expected), CancellationToken.None); + await AdvanceUntilAsync(clock, observing, "a source that keeps answering \"not yet\" can only be given up on by the finalization budget — it never was"); + var observed = await observing; + observed.ShouldBeSameAs(expected); logs.Heads.ShouldAllBe(head => head.Metadata.State == AgentRunLogStreamState.Open && head.CaptureFinalizedAt == null); @@ -400,6 +478,43 @@ public async Task A_truncated_source_terminalizes_its_stream_as_Truncated_with_e private static AgentRunLogCaptureBridge Bridge(FakeLogService logs) => new(logs, new ReadyStorageResolver(), new FakeRecoveryService(), NullLogger.Instance); + /// How long a virtual-time test may spend in REAL seconds before it is a hang rather than a slow runner. Nothing is asserted against it; it only keeps a wedged drain from being reported as five silent minutes (Rule 12.10). + private static readonly TimeSpan Patience = TimeSpan.FromSeconds(20); + + /// Wait for something the drain does on its own, with no clock to advance — the first offer of a segment, say. Real time, because only the ceilings are virtual. + private static async Task WaitAsync(Func condition, string signal) + { + var watch = Stopwatch.StartNew(); + while (watch.Elapsed < Patience) + { + if (condition()) return; + await Task.Delay(5); + } + throw new Xunit.Sdk.XunitException($"{signal} (waited {Patience.TotalSeconds:F0}s)"); + } + + /// One step of virtual time. Small enough that a ceiling is never overshot by more than this, and stepped rather than jumped because a single jump can land in the window between a wait being decided on and its timer being armed — which would leave a timer due AFTER the jump and a test waiting for a moment that never comes. + private static readonly TimeSpan AdvanceStep = TimeSpan.FromMilliseconds(10); + + /// + /// Push the virtual clock forward in steps until settles, and + /// answer how much virtual time that took. The bridge's ceilings are on this clock, so one expires only because a + /// test expired it — never because the runner was busy — and the ORDER in which two expire is fixed by their due + /// times rather than by scheduling luck. + /// + private static async Task AdvanceUntilAsync(FakeTimeProvider clock, Task work, string signal) + { + var started = clock.GetUtcNow(); + var watch = Stopwatch.StartNew(); + while (!work.IsCompleted) + { + if (watch.Elapsed > Patience) throw new Xunit.Sdk.XunitException($"{signal} (advanced the virtual clock by {clock.GetUtcNow() - started} over {Patience.TotalSeconds:F0}s of real time)"); + clock.Advance(AdvanceStep); + await Task.Delay(5); + } + return clock.GetUtcNow() - started; + } + private static AgentRunLogCaptureOpenRequest Request(FakeLogSource source, long fence, Guid sessionId, SecretRedactor? redactor = null) => new() { TeamId = TeamId, AgentRunId = RunId, ActorId = ActorId, WorkerFenceEpoch = fence, @@ -458,9 +573,14 @@ public IReadOnlyList DescribeLogs(SandboxHandle han new("stderr", AgentRunLogKinds.StandardError, AgentRunLogRepresentations.PlainTextContentType, AgentRunLogRepresentations.Utf8ContentEncoding, "fake-spool/v1"), ]; + /// + /// Deliberately does NOT observe the token, because production's does not either: the local spool's read + /// answers "nothing new" out of a length comparison, with no I/O and no cancellation check, whenever the file + /// has grown by less than one minimum segment. A fake that threw here made every cancellation bug in the + /// capture loop unreachable from this suite — which is exactly how one shipped. + /// public Task ReadAsync(SandboxDurableLogReadRequest request, CancellationToken cancellationToken) { - cancellationToken.ThrowIfCancellationRequested(); if (!_sources.TryGetValue(request.SourceKey, out var bytes)) return Task.FromResult(new SandboxDurableLogReadResult.Unavailable(new SandboxDurableLogProblem(SandboxDurableLogProblemCode.SourceMissing))); if (request.OffsetBytes > bytes.LongLength)