From fdec4ccf329cb23459f3eeecbb15a8369b6e31b3 Mon Sep 17 00:00:00 2001 From: "Mars.P" Date: Sun, 20 Sep 2026 07:03:47 +0800 Subject: [PATCH 1/3] Honour the worker's landing budget in the log capture drain A worker tear-down and its log capture spend one deadline. Two things in the capture bridge ignored that. The live capture loop's only cancellation point was a poll wait, and Task.WhenAny hands back a cancelled delay instead of raising it. The loop's other await is the source read, which answers "nothing new" out of a length comparison whenever the spool has grown by less than one minimum segment -- no I/O and no token. So a capture whose observer was cancelled spun at full speed forever and the tear-down awaiting it never returned. Only a source that ignores the token the way production's does can reach this, which is why the unit fake now ignores it too. The final drain then waited out a destination that was refusing writes for as long as its caller's token allowed. Under a tear-down that token IS the lease-landing budget, so a deployment whose log destination was misconfigured spent the whole of it retrying and landed no verdict at all. The drain now learns that the host is going away -- the same IHostApplicationLifetime signal the tear-down arm itself acts on -- and parks on the first refusal: locally, leaving the row Open at its fence with the stall marker it already wrote, which is the shape the recovery sweep finishes. The live stall cadence is untouched. The drain's ceilings move onto the injected TimeProvider so a test advances them instead of sleeping through them; the poll cadence stays on the wall clock, because it decides nothing. A retry test that raced a 300ms wall budget against a 50ms wall backoff was flaking on loaded runners; it now spends virtual time, which load cannot. The deployed cadence is pinned by a literal test. --- .../Services/Agents/AgentRunExecutor.cs | 5 + .../AgentRunLogCaptureBridge.cs | 88 ++++++++-- .../IAgentRunLogCaptureBridge.cs | 12 ++ .../AgentRunExecutorCredentialBrokerTests.cs | 143 ++++++++++++++++ .../AgentRunLogCaptureBackpressureTests.cs | 100 +++++++++++- .../Agents/AgentRunLogCaptureBridgeTests.cs | 152 ++++++++++++++++-- 6 files changed, 465 insertions(+), 35 deletions(-) diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs index b5c50ce85..b64a80a75 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs @@ -5087,6 +5087,11 @@ private async Task OpenLogCaptureAsync(LogCaptureCon { TeamId = context.TeamId, AgentRunId = context.RunId, ActorId = context.ActorId, WorkerFenceEpoch = context.WorkerFenceEpoch, Handle = handle, Source = source, Redactor = context.Redactor, + // The capture drain and this run's terminal write spend ONE budget on a worker tear-down + // (ShutdownLeaseLandingBudget), so the drain has to know a tear-down is happening: a destination that + // is refusing writes would otherwise wait out every second the landing needed. Same predicate the + // tear-down arm itself acts on — the HOST's 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..f489d4f3b 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,6 +221,16 @@ 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 @@ -218,6 +240,7 @@ private async Task CaptureLoopAsync(AgentRunLogCaptureOpenRequest request, Guid foreach (var stream in streams.Where(value => !value.Terminal)) 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)) { @@ -296,9 +319,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 +342,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 +355,27 @@ 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. + /// + /// The drain and the run's terminal write spend ONE deadline (the worker's lease-landing budget). Waiting out + /// a destination that is refusing writes spends all of it, and the run then lands nothing: a misconfigured log + /// destination would cost the operator the VERDICT as well as the log tail, which is the wrong trade by a wide + /// margin. 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 recovery sweep + /// finishes. Nothing is terminalized, so nothing a later observer could have completed is foreclosed. + /// + private bool ParkedForHostShutdown(AgentRunLogCaptureOpenRequest request, CaptureStream stream, RemoteStall stall) + { + if (!request.HostShutdown.IsCancellationRequested) return false; + + stream.Terminal = 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 recovery sweep so the run's own landing keeps the rest of the shared budget", request.AgentRunId, stream.Metadata.StreamId, stall.Attempts, stall.Code); + + 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) { @@ -620,7 +667,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,8 +701,7 @@ public async Task ObserveAsync(Func !value.Terminal)) _owner._logger.LogWarning("Agent run {RunId} source final drain exceeded its shadow budget and remains Open for reconciliation", _request.AgentRunId); return result; @@ -674,6 +721,15 @@ 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 diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs index 0e24269f9..3a9fb4f36 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs @@ -24,6 +24,18 @@ 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), which is what makes this + /// capture's drain budget SHARED rather than its own: a worker that is going away lands the run's verdict out of + /// the same seconds this drain is spending, and a destination that is refusing writes would otherwise spend all of + /// them. Raised ⇒ the final drain stops waiting out a refusal and parks — the stream stays Open at its own fence + /// with its stall marker, for the recovery sweep to finish. + /// + /// 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..f05d70ecf 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs @@ -16,6 +16,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; @@ -379,6 +380,148 @@ 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`"); + } + + /// + /// 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 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..03f547540 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,68 @@ 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 spent NOTHING of the budget it shares with the run's landing — a single backoff here is a second the terminal write does not get"); + 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_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 +352,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" }; @@ -303,6 +372,31 @@ private static byte[] Payload(string? secret, int length) } /// Explicit timeout with the watched signal named, so a failure says what never happened rather than only that time ran out (Rule 12.10). + /// 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); + } + } + private static async Task WaitAsync(Func condition, string signal) { var watch = Stopwatch.StartNew(); diff --git a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs index 84700f980..8e896c854 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 * 2, "ONE total finalization budget bounds a provider that never completes — N segments must not multiply it"); 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) From 2a686755c0d22561553e94b55ed148643de6a972 Mon Sep 17 00:00:00 2001 From: "Mars.P" Date: Sun, 20 Sep 2026 08:14:27 +0800 Subject: [PATCH 2/3] Park the capture drain's other wait, and say what waiting costs Review of the first pass found three things. The drain has two ways to wait forever, and only one was parked. A source that never seals terminalizes nothing, so the final loop polls it to the finalization budget -- thirty seconds of a worker that is already leaving, spent on bytes a later owner can read just as well. It parks there too now. A park is not a conclusion, so it stops borrowing the flag for one. CaptureStream.Parked is its own field: the row is still Open at its own fence, nothing terminal was written, and the drain's "exceeded its shadow budget" warning no longer claims a stream that stopped early on purpose and already said so. "One shared budget" was too broad. ShutdownLeaseLandingBudget is minted inside the tear-down arm, which runs from the cancellation the live drain is sitting on -- so a live drain that will not end is a landing that never STARTS, and only the session the tear-down opens to re-fold the dead agent actually spends the landing's own seconds. Both are stated where the claim is made. Test-side: the backpressure suite's own fake source stopped observing a token production does not observe either -- the same fixture defect the first pass fixed one copy of. The existing healthy-destination drain arm was vacuous about capture, because with no route the bridge fails every stream at open and hands back a passthrough; it now seeds a writable destination, which makes it the integration witness for the live loop's cancellation. A new arm covers the park end to end for a run that FINISHES while the host is stopping, where the observer returns normally and the final drain is the first thing to see it. --- .../Services/Agents/AgentRunExecutor.cs | 12 ++- .../AgentRunLogCaptureBridge.cs | 73 ++++++++++++---- .../IAgentRunLogCaptureBridge.cs | 14 +-- .../AgentRunExecutorCredentialBrokerTests.cs | 87 +++++++++++++++++++ .../AgentRunLogCaptureBackpressureTests.cs | 42 ++++++++- .../Agents/AgentRunLogCaptureBridgeTests.cs | 2 +- 6 files changed, 200 insertions(+), 30 deletions(-) diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs index b64a80a75..add042785 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs @@ -5087,10 +5087,14 @@ private async Task OpenLogCaptureAsync(LogCaptureCon { TeamId = context.TeamId, AgentRunId = context.RunId, ActorId = context.ActorId, WorkerFenceEpoch = context.WorkerFenceEpoch, Handle = handle, Source = source, Redactor = context.Redactor, - // The capture drain and this run's terminal write spend ONE budget on a worker tear-down - // (ShutdownLeaseLandingBudget), so the drain has to know a tear-down is happening: a destination that - // is refusing writes would otherwise wait out every second the landing needed. Same predicate the - // tear-down arm itself acts on — the HOST's lifetime, never a cancelled job token. + // 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); } diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs index f489d4f3b..3f671edac 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs @@ -237,16 +237,20 @@ private async Task CaptureLoopAsync(AgentRunLogCaptureOpenRequest request, Guid { 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)) + while (streams.Any(value => value.Draining)) { - 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 (ParkedIncompleteSourcesForHostShutdown(request, streams)) break; + + await Task.Delay(PollInterval, cancellationToken).ConfigureAwait(false); } } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { } @@ -260,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; @@ -359,19 +363,43 @@ private async Task FlushBacklogAsync(AgentRunLogCaptureOpenRequest /// Stop draining this stream because the HOST is going away — the one thing that ends a final drain before its own /// budget does. /// - /// The drain and the run's terminal write spend ONE deadline (the worker's lease-landing budget). Waiting out - /// a destination that is refusing writes spends all of it, and the run then lands nothing: a misconfigured log - /// destination would cost the operator the VERDICT as well as the log tail, which is the wrong trade by a wide - /// margin. 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 recovery sweep - /// finishes. Nothing is terminalized, so nothing a later observer could have completed is foreclosed. + /// 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. /// private bool ParkedForHostShutdown(AgentRunLogCaptureOpenRequest request, CaptureStream stream, RemoteStall stall) { if (!request.HostShutdown.IsCancellationRequested) return false; - stream.Terminal = 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 recovery sweep so the run's own landing keeps the rest of the shared budget", request.AgentRunId, stream.Metadata.StreamId, stall.Attempts, stall.Code); + 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; + } + + /// + /// The final drain's OTHER unbounded wait, and the same answer. 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. That is 30 seconds of a worker that is already going away, spent on a source whose + /// bytes a later owner can read just as well. + /// + 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; } @@ -597,7 +625,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); } @@ -702,7 +730,7 @@ public async Task ObserveAsync(Func !value.Terminal)) + 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; } @@ -716,7 +744,7 @@ public async Task ObserveAsync(FuncThe 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 3a9fb4f36..39ddc195b 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/IAgentRunLogCaptureBridge.cs @@ -26,11 +26,15 @@ public sealed record AgentRunLogCaptureOpenRequest public required SecretRedactor Redactor { get; init; } /// - /// The HOST's own tear-down signal (IHostApplicationLifetime.ApplicationStopping), which is what makes this - /// capture's drain budget SHARED rather than its own: a worker that is going away lands the run's verdict out of - /// the same seconds this drain is spending, and a destination that is refusing writes would otherwise spend all of - /// them. Raised ⇒ the final drain stops waiting out a refusal and parks — the stream stays Open at its own fence - /// with its stall marker, for the recovery sweep to finish. + /// 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. diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunExecutorCredentialBrokerTests.cs index f05d70ecf..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; @@ -282,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 @@ -366,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"); @@ -450,6 +459,69 @@ await AwaitWithinAsync(execution, UnwritableDestinationDrainCeiling, 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 @@ -458,6 +530,21 @@ await AwaitWithinAsync(execution, UnwritableDestinationDrainCeiling, /// 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 { diff --git a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs index 03f547540..b342ffa90 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs @@ -229,7 +229,7 @@ public async Task A_drain_under_a_host_that_is_going_away_parks_on_the_first_ref 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 spent NOTHING of the budget it shares with the run's landing — a single backoff here is a second the terminal write does not get"); + "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"); @@ -238,6 +238,37 @@ public async Task A_drain_under_a_host_that_is_going_away_parks_on_the_first_ref 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_drain_on_a_healthy_host_still_waits_out_the_same_refusal() { @@ -371,7 +402,6 @@ private static byte[] Payload(string? secret, int length) return bytes; } - /// Explicit timeout with the watched signal named, so a failure says what never happened rather than only that time ran out (Rule 12.10). /// 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) { @@ -397,6 +427,7 @@ private static async Task AdvanceUntilAsync(FakeTimeProvider clock, Task work, s } } + /// 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) { var watch = Stopwatch.StartNew(); @@ -465,12 +496,15 @@ 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; + + /// 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) 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 8e896c854..796e1f230 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBridgeTests.cs @@ -266,7 +266,7 @@ public async Task Blocking_capture_backend_is_cancelled_by_one_total_shadow_budg var observed = await observing; spent.ShouldBeGreaterThanOrEqualTo(TimeSpan.FromMilliseconds(150), "the drain gave the provider its whole budget before giving up"); - spent.ShouldBeLessThan(TimeSpan.FromMilliseconds(150) + AdvanceStep * 2, "ONE total finalization budget bounds a provider that never completes — N segments must not multiply it"); + 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)); From 4132a4ff4b10aa9eda748a3716698cbadc9ce515 Mon Sep 17 00:00:00 2001 From: "Mars.P" Date: Mon, 21 Sep 2026 15:06:00 +0800 Subject: [PATCH 3/3] Let a late seal arrive before parking the drain for a shutdown The two park sites do not face the same odds, and giving them the same impatience was wrong. An append park happens because the destination REFUSED, and the next attempt is overwhelmingly likely to refuse again. The incomplete-source park waits on a local copier finishing -- very often the one draining the FIFO of an agent this same drain just killed, which is exactly the kind of thing that completes on the next poll. Parking that on the first pass reports a capture whose every byte landed as incomplete. Without an EndOfSource there is no final-drain receipt, so CompleteRunAsync skips the stream and the recovery sweep's terminal grace later stamps it CaptureFailed -- it terminalizes, it does not re-read the spool. A quarter second of scheduling would have cost the operator the whole log. So the incomplete-source park is consulted only from the second pass: one retry, one PollInterval, 250ms of a ten-second landing. The 30s the fix removed stays removed. The asymmetry in what the two leave behind is now stated where each one is. The append park lands on top of the remote_stall marker its own refusal already wrote, so the Room can say why. This one writes nothing durable, because nothing refused anything and a stall marker would send an operator hunting a storage incident that never happened -- at the cost, named rather than hidden, of the stream reading as "Finalizing" with no reason until the sweep's terminal grace elapses. --- .../AgentRunLogCaptureBridge.cs | 37 ++++++++++++++--- .../AgentRunLogCaptureBackpressureTests.cs | 41 ++++++++++++++++++- 2 files changed, 71 insertions(+), 7 deletions(-) diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs index 3f671edac..65c02ab9a 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureBridge.cs @@ -242,13 +242,13 @@ private async Task CaptureLoopAsync(AgentRunLogCaptureOpenRequest request, Guid await Task.WhenAny(Task.Delay(PollInterval, cancellationToken), finish).ConfigureAwait(false); cancellationToken.ThrowIfCancellationRequested(); } - while (streams.Any(value => value.Draining)) + for (var pass = 1; streams.Any(value => value.Draining); pass++) { foreach (var stream in streams.Where(value => value.Draining)) await PumpAsync(request, captureSessionId, stream, final: true, cancellationToken).ConfigureAwait(false); if (!streams.Any(value => value.Draining)) break; - if (ParkedIncompleteSourcesForHostShutdown(request, streams)) break; + if (pass > SourceSealGracePasses && ParkedIncompleteSourcesForHostShutdown(request, streams)) break; await Task.Delay(PollInterval, cancellationToken).ConfigureAwait(false); } @@ -374,6 +374,10 @@ private async Task FlushBacklogAsync(AgentRunLogCaptureOpenRequest /// 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) { @@ -386,10 +390,31 @@ private bool ParkedForHostShutdown(AgentRunLogCaptureOpenRequest request, Captur } /// - /// The final drain's OTHER unbounded wait, and the same answer. 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. That is 30 seconds of a worker that is already going away, spent on a source whose - /// bytes a later owner can read just as well. + /// 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) { diff --git a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs index b342ffa90..d9d2681e1 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/AgentRunLogCaptureBackpressureTests.cs @@ -269,6 +269,39 @@ public async Task A_drain_whose_source_never_seals_also_parks_for_a_host_that_is 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() { @@ -499,12 +532,18 @@ public IReadOnlyList DescribeLogs(SandboxHandle han /// 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) { var bytes = _sources[request.SourceKey]; var available = bytes.LongLength - request.OffsetBytes; - if (available == 0 && request.FinalDrain && EmitEndOfSource) 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)));