diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRetryContinuity.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRetryContinuity.cs index 72b73c1e8..55c3fda66 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRetryContinuity.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRetryContinuity.cs @@ -11,7 +11,7 @@ namespace CodeSpace.Core.Services.Agents; public static class AgentRetryContinuity { /// The honest-redo line: fires ONLY when a resumed conversation exists but the workspace was NOT pinned to a prior pushed branch — never on a genuine cold-start retry (no prior attempt at all), which stays byte-identical. - public const string HonestNoContinuityHint = "Note: your prior attempt's conversation is restored, but its git changes were NOT preserved in this workspace (your prior attempt pushed no branch of its own) — you must redo any relevant file changes from scratch."; + public const string HonestNoContinuityHint = "Note: your prior attempt's conversation is restored, but its git changes were NOT preserved in this workspace (your prior attempt pushed no branch of its own for this repository) — you must redo any relevant file changes from scratch."; /// Append to a resumed task's goal. One composition, so the two lanes cannot drift on the separator either. public static string WithHonestNoContinuityHint(string goal) => $"{goal}\n\n{HonestNoContinuityHint}"; diff --git a/backend/src/CodeSpace.Core/Services/Agents/Credentials/Broker/LoopbackModelCredentialBroker.cs b/backend/src/CodeSpace.Core/Services/Agents/Credentials/Broker/LoopbackModelCredentialBroker.cs index fd21c221c..22a00f800 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/Credentials/Broker/LoopbackModelCredentialBroker.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/Credentials/Broker/LoopbackModelCredentialBroker.cs @@ -499,6 +499,16 @@ private async Task AcceptAsync(Lease lease) // matters more here than it used to — a revoke now closes a listener per FINISHED RUN, where before a // listener only ever closed at process teardown — and the only correct response to any of them is to stop // serving an address that no longer exists, and to stop CLAIMING it. + // + // A close can also land between the request above and this re-registration. The managed listener checks + // its state and only then queues the wait, and Close() completes-and-clears that queue in between, so the + // wait joins a closed listener's queue and this loop never resumes. That strands nothing: every close path + // takes the lease out of _byRun (supersede, revoke, sweep and drop before closing, dispose right after), + // and what is left is a cycle no root reaches — the closed listener's queue, the wait, this state machine, + // the lease, the listener — because Close() has already unhooked the listener from the endpoint manager's + // statics. It is collected, decrypted key and all, exactly as a loop that woke and returned would be + // (checked on .NET 10 by stranding a loop this way: its lease's weak reference cleared, while a lease a + // table still held stayed alive). try { context = await lease.Listener.GetContextAsync().ConfigureAwait(false); } catch (Exception exception) { DropIfStillServing(lease, exception); return; } diff --git a/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.Rehydrate.cs b/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.Rehydrate.cs index 4977ef974..1e1bf704c 100644 --- a/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.Rehydrate.cs +++ b/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.Rehydrate.cs @@ -11,6 +11,7 @@ using CodeSpace.Messages.Dtos.Decisions; using CodeSpace.Messages.Review; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace CodeSpace.Core.Services.Supervisor; @@ -1567,20 +1568,34 @@ private async Task GradeStopTargetsWithHeartbeatAsync(Guid super } } - /// The heartbeat loop itself: sleeps, logs, repeats — until fires (grading finished). A cancellation mid-sleep is the expected exit, never propagated as a fault. Internal + clock-parameterized so a unit test drives the sleep on a fake instead of racing the wall clock; REQUIRED rather than defaulting to the system clock for the reason gives — a default is how a call site keeps the wall clock without saying so. - internal async Task RunGradingHeartbeatLoopAsync(Guid supervisorRunId, string nodeId, TimeSpan interval, CancellationToken cancellationToken, TimeProvider timeProvider, string message = "Supervisor stop acceptance grading is still in progress.") + /// + /// The heartbeat loop itself: sleeps, pulses, repeats — until fires (grading + /// finished). It is , so it NEVER completes faulted: a cancellation is the + /// expected exit, and a pulse that fails is logged as a warning and retried on the next interval. That is + /// load-bearing rather than tidy — every call site awaits this task in the finally around the grade it + /// protects, so a fault here would REPLACE the grade with the error of a missed log line. + /// + /// Internal + clock-parameterized so a unit test drives the sleep on a fake + /// instead of racing the wall clock; REQUIRED rather than defaulting to the system clock for the reason + /// gives — a default is how a call site keeps the wall clock without saying so. + /// + internal Task RunGradingHeartbeatLoopAsync(Guid supervisorRunId, string nodeId, TimeSpan interval, CancellationToken cancellationToken, TimeProvider timeProvider, string message = "Supervisor stop acceptance grading is still in progress.") => + HeartbeatLoop.RunAsync(ct => PulseGradingHeartbeatAsync(supervisorRunId, nodeId, message, ct), interval, exception => _logger.LogWarning(exception, "Supervisor run {SupervisorRunId} node {NodeId}: a grading heartbeat pulse failed; grading continues and the next pulse retries", supervisorRunId, nodeId), cancellationToken, timeProvider); + + /// + /// One pulse, written through a record logger from a DI scope of its OWN. The pulse runs concurrently with the + /// grade, and this service's scope holds the one DbContext the grade is using — its manifest stamps and recorded + /// judge calls go through it — so a pulse on this scope's logger that fell due mid-query died on EF's "a second + /// operation was started on this context" guard. A scope per pulse rather than one per loop, so an insert that + /// failed never waits in a change tracker for the next pulse's save to retry it. + /// + private async Task PulseGradingHeartbeatAsync(Guid supervisorRunId, string nodeId, string message, CancellationToken cancellationToken) { - try - { - while (true) - { - await Task.Delay(interval, timeProvider, cancellationToken).ConfigureAwait(false); + using var scope = _scopeFactory.CreateScope(); - await _recordLogger.LogAsync(supervisorRunId, nodeId, Workflows.Lifecycle.LogLevel.Info, - message, cancellationToken).ConfigureAwait(false); - } - } - catch (OperationCanceledException) { } + var records = scope.ServiceProvider.GetRequiredService(); + + await records.LogAsync(supervisorRunId, nodeId, Workflows.Lifecycle.LogLevel.Info, message, cancellationToken).ConfigureAwait(false); } /// diff --git a/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.cs b/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.cs index 5881125dd..95c8dd925 100644 --- a/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.cs +++ b/backend/src/CodeSpace.Core/Services/Supervisor/SupervisorTurnService.cs @@ -10,6 +10,7 @@ using CodeSpace.Messages.Budget; using CodeSpace.Messages.Dtos.Agents; using CodeSpace.Messages.Plans; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace CodeSpace.Core.Services.Supervisor; @@ -41,6 +42,9 @@ public sealed partial class SupervisorTurnService : ISupervisorTurnService, ISco private readonly ISupervisorPublishedBranchResolver _publishedBranches; private readonly Completion.ICompletionAssessmentComposer _completion; + /// Opens the scope each grading-heartbeat pulse writes through — never this service's own, whose DbContext the grade the pulse runs beside is using (see ). + private readonly IServiceScopeFactory _scopeFactory; + /// C1 — the rubric judge a BRANCHLESS LlmJudge stop reads its summary with. OPTIONAL so the many hand-built test doubles keep compiling; DI always supplies it, and a null one fails the gate CLOSED rather than passing it silently. private readonly Review.IRubricJudge? _rubricJudge; @@ -49,7 +53,7 @@ public sealed partial class SupervisorTurnService : ISupervisorTurnService, ISco private readonly ILogger _logger; - public SupervisorTurnService(ISupervisorDecisionLog ledger, ISupervisorDecider decider, ISupervisorActionExecutor executor, CodeSpaceDbContext db, ISupervisorAcceptanceGrader acceptanceGrader, IDecisionQueueService decisionQueue, IDecisionArbiter arbiter, IDecisionAnswerService decisionAnswer, Plans.IWorkPlanService workPlans, Workflows.Lifecycle.IRunRecordLogger recordLogger, Workflows.Artifacts.IArtifactOffloader offloader, IPublishManifestStore manifests, ISupervisorPublishedBranchResolver publishedBranches, Completion.ICompletionAssessmentComposer completion, Workflows.Budget.IBudgetLedger budget, Learning.ILessonReader lessons, ILogger logger, Review.IRubricJudge? rubricJudge = null, Completion.IModeProfileRegistry? modes = null) + public SupervisorTurnService(ISupervisorDecisionLog ledger, ISupervisorDecider decider, ISupervisorActionExecutor executor, CodeSpaceDbContext db, ISupervisorAcceptanceGrader acceptanceGrader, IDecisionQueueService decisionQueue, IDecisionArbiter arbiter, IDecisionAnswerService decisionAnswer, Plans.IWorkPlanService workPlans, Workflows.Lifecycle.IRunRecordLogger recordLogger, Workflows.Artifacts.IArtifactOffloader offloader, IPublishManifestStore manifests, ISupervisorPublishedBranchResolver publishedBranches, Completion.ICompletionAssessmentComposer completion, Workflows.Budget.IBudgetLedger budget, Learning.ILessonReader lessons, IServiceScopeFactory scopeFactory, ILogger logger, Review.IRubricJudge? rubricJudge = null, Completion.IModeProfileRegistry? modes = null) { _ledger = ledger; _decider = decider; @@ -67,6 +71,7 @@ public SupervisorTurnService(ISupervisorDecisionLog ledger, ISupervisorDecider d _manifests = manifests; _publishedBranches = publishedBranches; _completion = completion; + _scopeFactory = scopeFactory; _rubricJudge = rubricJudge; _modes = modes; _logger = logger; diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Artifacts/Retention/ArtifactRetentionPolicy.cs b/backend/src/CodeSpace.Core/Services/Workflows/Artifacts/Retention/ArtifactRetentionPolicy.cs index 5f6df4193..14324803e 100644 --- a/backend/src/CodeSpace.Core/Services/Workflows/Artifacts/Retention/ArtifactRetentionPolicy.cs +++ b/backend/src/CodeSpace.Core/Services/Workflows/Artifacts/Retention/ArtifactRetentionPolicy.cs @@ -62,8 +62,7 @@ public static class ArtifactRetentionPolicy [SessionTranscriptCheckpoint.Class] = SessionTranscriptCheckpoint, }; - /// The rule for , or null when the running policy does not register it — which the reaper reads as "cannot tell" and keeps. - /// The rule for a class NAME, or null when this build registers none — including a name a rolled-back build wrote that this one has never heard of. Null settles as keep. + /// The rule for the class NAME , or null when this build registers none — including a name a rolled-back build wrote that this one has never heard of. The reaper reads null as "cannot tell" and keeps. public static ArtifactRetentionRule? For(string value) => Enum.TryParse(value, ignoreCase: false, out var parsed) && Rules.TryGetValue(parsed, out var rule) ? rule : null; /// diff --git a/backend/tests/CodeSpace.IntegrationTests/Agents/ModelPricingUnderCapFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Agents/ModelPricingUnderCapFlowTests.cs index b40b4d853..73041834b 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Agents/ModelPricingUnderCapFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Agents/ModelPricingUnderCapFlowTests.cs @@ -764,7 +764,7 @@ private static SupervisorTurnService BuildService(ILifetimeScope scope, ISupervi scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); private sealed class AlwaysSpawnDecider : ISupervisorDecider { diff --git a/backend/tests/CodeSpace.IntegrationTests/Learning/LessonArmSupervisorFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Learning/LessonArmSupervisorFlowTests.cs index 11ae63669..737579c9a 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Learning/LessonArmSupervisorFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Learning/LessonArmSupervisorFlowTests.cs @@ -219,7 +219,7 @@ private async Task SeedLessonAsync(Guid teamId) scope.Resolve(), scope.Resolve(), lessons ?? scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); private static Lesson Lesson(Guid teamId, string howToApply) => new() { diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/Supervisor/RealModelChecksBeforeCriticE2ETests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/Supervisor/RealModelChecksBeforeCriticE2ETests.cs index 9e7e37da0..30fb5462e 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/Supervisor/RealModelChecksBeforeCriticE2ETests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/Supervisor/RealModelChecksBeforeCriticE2ETests.cs @@ -105,7 +105,7 @@ private async Task RunPlanTurnAsync(Guid runId, Guid teamId, Guid reviewerRowId, scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); var goalConfig = new SupervisorGoalConfig { Goal = Goal, DecisionReviewMode = ReviewMode.Gate, ReviewerModelId = reviewerRowId }; diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorAcceptanceFoldFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorAcceptanceFoldFlowTests.cs index 452c48d63..37fed8ff4 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorAcceptanceFoldFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorAcceptanceFoldFlowTests.cs @@ -729,7 +729,7 @@ public async Task The_real_grade_drives_the_terminal_stop_status_over_a_real_rep scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); result = await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, GoalConfig(repoId, acceptanceChecks: null), CancellationToken.None); } @@ -773,7 +773,7 @@ public async Task The_real_operator_floor_gates_a_clean_runs_terminal_stop_over_ scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); result = await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, GoalConfig(repoId, acceptanceChecks: new[] { "sh", "check.sh" }), CancellationToken.None); } @@ -821,7 +821,7 @@ public async Task The_real_operator_floor_gates_a_multi_repo_stop_all_or_nothing scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); result = await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, GoalConfig(repoA, acceptanceChecks: new[] { "sh", "check.sh" }), CancellationToken.None); } @@ -878,7 +878,7 @@ public async Task The_real_operator_floor_voids_a_heads_rewrite_of_the_check_scr scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); result = await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, GoalConfig(repoId, acceptanceChecks: new[] { "sh", "check.sh" }), CancellationToken.None); } @@ -1217,7 +1217,7 @@ private async Task RehydrateAsync(Guid runId, Guid teamId scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); return await service.RehydrateFromDecisionLogAsync(runId, teamId, NodeId, Goal, goalConfig, CancellationToken.None); } @@ -1244,7 +1244,7 @@ private async Task RunStopTurnWithGraderAsync(Guid runId, scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), - scope.Resolve(), scope.Resolve>()); + scope.Resolve(), scope.Resolve(), scope.Resolve>()); return await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, goalConfig, CancellationToken.None); } diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorArbiterDrainFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorArbiterDrainFlowTests.cs index 08b6e0371..b087a5c9d 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorArbiterDrainFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorArbiterDrainFlowTests.cs @@ -219,7 +219,7 @@ public async Task The_frozen_in_flight_replay_bypasses_the_arbiter_drain() scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); private async Task RunTurnAsync(Guid runId, Guid teamId, IDecisionArbiter arbiter) { diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorChecksBeforeCriticFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorChecksBeforeCriticFlowTests.cs index 41d4489c5..d13b42223 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorChecksBeforeCriticFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorChecksBeforeCriticFlowTests.cs @@ -96,7 +96,7 @@ private async Task RunPlanTurnAsync(Guid runId, Guid teamId, RecordingCritic cri scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); var goalConfig = new SupervisorGoalConfig { Goal = Goal, DecisionReviewMode = ReviewMode.Gate }; diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDeliveryGateFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDeliveryGateFlowTests.cs index c6cb153c8..b9a10f4e0 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDeliveryGateFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDeliveryGateFlowTests.cs @@ -572,7 +572,7 @@ private async Task RunTurnAsync(Guid runId, Guid teamId scope.Resolve(), scope.Resolve(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); [Fact] public async Task A_model_minted_gate_card_cannot_drive_the_adjudication_release() diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDependencyOrderingFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDependencyOrderingFlowTests.cs index bdb8f8b87..3df2a7fa1 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDependencyOrderingFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorDependencyOrderingFlowTests.cs @@ -149,7 +149,7 @@ private async Task RunSpawnTurnAsync(Guid runId, Guid teamId, SupervisorRational scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, goalConfig: null, CancellationToken.None); } diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorGradingHeartbeatIsolationFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorGradingHeartbeatIsolationFlowTests.cs new file mode 100644 index 000000000..d2dece17e --- /dev/null +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorGradingHeartbeatIsolationFlowTests.cs @@ -0,0 +1,142 @@ +using Autofac; +using CodeSpace.Core.Persistence.Db; +using CodeSpace.Core.Services.Supervisor; +using CodeSpace.Core.Services.Workflows.Lifecycle; +using CodeSpace.IntegrationTests.Infrastructure; +using CodeSpace.IntegrationTests.Workflows.Infrastructure; +using CodeSpace.Messages.Agents.Benchmark; +using CodeSpace.Messages.Commands.Workflows; +using CodeSpace.Messages.Constants; +using MediatR; +using Microsoft.EntityFrameworkCore; +using Npgsql; +using Shouldly; + +namespace CodeSpace.IntegrationTests.Workflows; + +/// +/// 🟢 High fidelity — the production container's own scope wiring, a real and real +/// Postgres. The grading heartbeat () runs +/// CONCURRENTLY with the grade it protects, and one DI scope holds ONE : the turn +/// service's own reads and writes (the unit grade's manifest stamps, the judge's recorded model calls) and every +/// resolved from that scope write through the same instance. A pulse that fell due +/// while a grade query was in flight on it died on EF's "a second operation was started on this context" guard — +/// and, awaited in the grade's finally, replaced the grade with that error. +/// +/// The test holds a real query open on the grade scope's context (blocked on an advisory lock only the test +/// releases), lets a pulse fall due, and releases only once that pulse has settled — so the verdict is decided by +/// which context the pulse wrote through, never by how the runner scheduled it. The wall clock decides only WHEN +/// the first pulse fires; every wait on it is bounded and names what it watched. +/// +[Collection(PostgresCollection.Name)] +[Trait("Category", "Integration")] +public class SupervisorGradingHeartbeatIsolationFlowTests +{ + private const string NodeId = "sup"; + private const string Graded = "graded while a pulse fell due"; + + /// The advisory-lock class this test's key lives under — a two-key lock never contends with the single-key pg_advisory_xact_lock the run-record admission trigger takes per run. + private const int LockClass = 1_976_1994; + + private static readonly TimeSpan PulseInterval = TimeSpan.FromMilliseconds(100); + private static readonly TimeSpan SettleBound = TimeSpan.FromSeconds(30); + + private readonly PostgresFixture _fixture; + + public SupervisorGradingHeartbeatIsolationFlowTests(PostgresFixture fixture) { _fixture = fixture; } + + [Fact] + public async Task A_pulse_that_falls_due_mid_query_lands_on_its_own_context_and_the_grade_survives() + { + var runId = await SeedRunAsync(); + var lockKey = Random.Shared.Next(); + + using var gradeScope = _fixture.BeginScope(); + var service = gradeScope.Resolve(); + + await using var holder = await HoldLockAsync(lockKey); + + // The grade's own DB work, in flight on the scope's context for as long as the holder keeps the lock. + var gradeQuery = gradeScope.Resolve().Database.ExecuteSqlRawAsync("SELECT pg_advisory_xact_lock({0}, {1})", LockClass, lockKey); + + // The premise, at the DI level: a record logger from the grade's own scope writes through the grade's own + // context, so a write on it collides with the query above. Without this collision the test proves nothing. + var collision = await Should.ThrowAsync(() => gradeScope.Resolve().LogAsync(runId, NodeId, LogLevel.Info, "premise probe", CancellationToken.None)); + collision.Message.ShouldContain("second operation was started on this context", customMessage: "premise: the scope's record logger shares the grade's DbContext — the collision the heartbeat must be kept out of"); + + using var heartbeatCts = new CancellationTokenSource(); + var heartbeat = service.RunGradingHeartbeatLoopAsync(runId, NodeId, PulseInterval, heartbeatCts.Token, TimeProvider.System); + + BenchmarkGrade grade; + try + { + await PulseSettledAsync(runId, heartbeat); + + await holder.DisposeAsync(); // release: the grade's query takes the lock and completes + await gradeQuery; + + grade = new BenchmarkGrade { Passed = true, Detail = Graded }; + } + finally + { + // The production call sites' own finally shape: whatever the loop ends with is what this await surfaces. + heartbeatCts.Cancel(); + + try { await heartbeat; } + catch (OperationCanceledException) { } + } + + grade.Detail.ShouldBe(Graded, "a heartbeat pulse can never replace the grade it protects"); + + (await PulseCountAsync(runId)).ShouldBeGreaterThan(0, "a pulse landed while the grade's query still held the scope's context — only a pulse on its OWN context can"); + } + + /// + /// Waits until a pulse row for has landed, or the loop itself has ended — a loop that + /// died on its first pulse has settled too, and the call-site finally then surfaces what killed it. + /// + private async Task PulseSettledAsync(Guid runId, Task heartbeat) + { + var deadline = DateTime.UtcNow + SettleBound; + + while (DateTime.UtcNow < deadline) + { + if (heartbeat.IsCompleted || await PulseCountAsync(runId) > 0) return; + + await Task.Delay(50); + } + + throw new TimeoutException($"no grading-heartbeat pulse landed for run {runId} within {SettleBound.TotalSeconds}s while the grade's query held its scope's context, and the loop is still running — every pulse is failing. Look for the SupervisorTurnService pulse warnings: a pulse that writes through the grade's own context dies on EF's second-operation guard."); + } + + private async Task PulseCountAsync(Guid runId) + { + using var scope = _fixture.BeginScope(); + + return await scope.Resolve().WorkflowRunRecord.AsNoTracking().CountAsync(r => r.RunId == runId && r.NodeId == NodeId && r.RecordType == WorkflowRunRecordTypes.Log); + } + + /// Takes the lock on an unpooled connection of its own, so disposing it ends the session and releases the lock on every path — including a failed assertion. + private async Task HoldLockAsync(int lockKey) + { + var connection = new NpgsqlConnection(new NpgsqlConnectionStringBuilder(_fixture.ConnectionString) { Pooling = false }.ConnectionString); + await connection.OpenAsync(); + + await using var take = new NpgsqlCommand("SELECT pg_advisory_lock(@class, @key)", connection); + take.Parameters.AddWithValue("class", LockClass); + take.Parameters.AddWithValue("key", lockKey); + await take.ExecuteNonQueryAsync(); + + return connection; + } + + private async Task SeedRunAsync() + { + var (teamId, userId) = await WorkflowsTestSeed.SeedTeamAsync(_fixture, inProcessPool: false); + + using var scope = _fixture.BeginScopeAs(userId, teamId); + var workflowId = await scope.Resolve().Send(new CreateWorkflowCommand { Name = $"grading-heartbeat-{Guid.NewGuid():N}"[..24], Definition = WorkflowsTestSeed.MinimalDefinition(), Activations = Array.Empty(), Enabled = true }); + + return await WorkflowsTestSeed.SeedManualRunAsync(_fixture, workflowId, teamId); + } +} diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorLedgerDirectTerminalOutputFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorLedgerDirectTerminalOutputFlowTests.cs index c1f3e475e..8af583463 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorLedgerDirectTerminalOutputFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorLedgerDirectTerminalOutputFlowTests.cs @@ -191,7 +191,7 @@ private async Task RunTurnAsync(Guid runId, Guid teamId, I scope.Resolve(), scope.Resolve(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); private sealed class AlwaysStopDecider : ISupervisorDecider { diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorMergeWithholdFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorMergeWithholdFlowTests.cs index 72f2c350a..431e07059 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorMergeWithholdFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorMergeWithholdFlowTests.cs @@ -331,7 +331,7 @@ private static SupervisorAgentResult Unit(Guid agentRunId, string producedBranch scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, GoalConfig(), CancellationToken.None); } diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPayloadReaskFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPayloadReaskFlowTests.cs index 422c2b4fd..da2dc11bf 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPayloadReaskFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPayloadReaskFlowTests.cs @@ -220,7 +220,7 @@ private async Task CurrentWorkPlanItemsAsync(Guid runId, Guid teamId) scope.Resolve(), scope.Resolve(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); private static LlmSupervisorDecider NewDecider(ILifetimeScope scope, IStructuredLLMClient client) => new( new LLMClientRegistry(new ILLMClient[] { (ILLMClient)client }), diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanDeliveryFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanDeliveryFlowTests.cs index 66fe89c11..2032e72db 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanDeliveryFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanDeliveryFlowTests.cs @@ -195,7 +195,7 @@ private async Task LatestPlanPayloadAsync(Guid runId, Guid teamId) scope.Resolve(), scope.Resolve(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); /// A decider that always authors a plan with one subtask, proposing the given delivery contract (or none). private sealed class AlwaysPlanDecider : ISupervisorDecider diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanValidatorFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanValidatorFlowTests.cs index 4ee213228..32d29e14e 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanValidatorFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPlanValidatorFlowTests.cs @@ -81,7 +81,7 @@ private async Task RunPlanTurnAsync(Guid runId, Guid teamId, params (string Id, scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); await service.RunTurnAsync(runId, teamId, NodeId, Goal, conversationId: null, goalConfig: null, CancellationToken.None); } diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPublishGateFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPublishGateFlowTests.cs index a923269a7..4d4afde9e 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPublishGateFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorPublishGateFlowTests.cs @@ -300,7 +300,7 @@ private async Task RunStopTurnAsync(Guid runId scope.Resolve(), scope.Resolve(), scope.Resolve(), - scope.Resolve>()); + scope.Resolve(), scope.Resolve>()); private sealed record SupervisorDecisionRecordSnapshot(string Kind, string PayloadJson, string? OutcomeJson); diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorUnitAcceptanceFoldFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorUnitAcceptanceFoldFlowTests.cs index aab88f7b4..17790c08d 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorUnitAcceptanceFoldFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/SupervisorUnitAcceptanceFoldFlowTests.cs @@ -1524,7 +1524,7 @@ private async Task RehydrateAsync(Guid runId, Guid teamId scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), scope.Resolve(), new AdmitAllBudgetLedger(), scope.Resolve(), - scope.Resolve>(), + scope.Resolve(), scope.Resolve>(), rubricJudge: null, modes: scope.Resolve()); diff --git a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorArbiterDrainTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorArbiterDrainTests.cs index ccc2a6df5..ae3215a07 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorArbiterDrainTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorArbiterDrainTests.cs @@ -164,7 +164,7 @@ public async Task A_no_spawn_rehydrate_skips_the_queue_read_entirely() var queue = new FakeDecisionQueue(); var ledger = new FakeSupervisorDecisionLog(); ledger.SeedTerminal(runId, TeamId, SupervisorDecisionKinds.Plan, """{"subtasks":["a"]}""", """{"planned":["a"]}"""); - var service = new SupervisorTurnService(ledger, new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), queue, new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReaderStub(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), queue, new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReaderStub(), null!, NullLogger.Instance); var context = await service.RehydrateFromDecisionLogAsync(runId, TeamId, "sup", "goal", goalConfig: null, CancellationToken.None); @@ -187,7 +187,7 @@ public async Task The_arbiter_call_is_metered_against_the_runs_own_budget_ledger }); var ledger = new AdmitAllBudgetLedger(); var runId = Guid.NewGuid(); - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), arbiter, new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), ledger, new NoLessonsReaderStub(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), arbiter, new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), ledger, new NoLessonsReaderStub(), null!, NullLogger.Instance); var context = new SupervisorTurnContext { SupervisorRunId = runId, TeamId = TeamId, NodeId = "sup", Goal = "ship it", SupervisorModelId = BrainModelId, MaxCostUsd = 7.5m, PendingChildDecisions = new[] { Pending() } }; @@ -230,7 +230,7 @@ public async Task The_drain_runs_before_the_delivery_decider_and_falls_through_t private static SupervisorTurnService Drain(FakeDecisionArbiter arbiter, FakeDecisionAnswerService answer) => new(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), arbiter, answer, new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), - new NoLessonsReaderStub(), NullLogger.Instance); + new NoLessonsReaderStub(), null!, NullLogger.Instance); private static SupervisorTurnContext Context(params PendingDecision[] pending) => new() { diff --git a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBoundsServiceTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBoundsServiceTests.cs index a180c672c..c0e17254e 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBoundsServiceTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBoundsServiceTests.cs @@ -132,7 +132,7 @@ public async Task A_spawns_policy_parks_the_spawn_for_a_human_instead_of_creatin ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Plan, """{"subtasks":[{"id":"a","title":"A","instruction":"do"}]}""", "{}"); var executor = new CountingExecutor(); - var service = new SupervisorTurnService(ledger, new AlwaysSpawnDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReaderStub(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysSpawnDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReaderStub(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "g", null, Config(approvalPolicy: "spawns"), CancellationToken.None); @@ -151,7 +151,7 @@ public async Task A_none_policy_spawns_without_a_gate() ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Plan, """{"subtasks":[{"id":"a","title":"A","instruction":"do"}]}""", "{}"); var executor = new CountingExecutor(); - var service = new SupervisorTurnService(ledger, new AlwaysSpawnDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReaderStub(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysSpawnDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReaderStub(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "g", null, Config(approvalPolicy: "none"), CancellationToken.None); @@ -163,7 +163,7 @@ public async Task A_none_policy_spawns_without_a_gate() private SupervisorTurnService Service(FakeSupervisorDecisionLog ledger, ISupervisorDecider decider) => new(ledger, decider, new CountingExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), - new NoLessonsReaderStub(), NullLogger.Instance); + new NoLessonsReaderStub(), null!, NullLogger.Instance); private static SupervisorGoalConfig Config(int? maxTotalSpawns = null, int? maxNoProgress = null, string? approvalPolicy = null) => new() { MaxTotalSpawns = maxTotalSpawns, MaxNoProgressDecisions = maxNoProgress, ApprovalPolicy = approvalPolicy }; diff --git a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBranchlessStopGradeTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBranchlessStopGradeTests.cs index ce62d77aa..daddc68c1 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBranchlessStopGradeTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorBranchlessStopGradeTests.cs @@ -232,7 +232,7 @@ private static async Task GradeAsync(CapturingGrader grader, string stop // scope carries. Every other seam is untouched by ApplyStopAcceptanceGradeAsync. var service = new SupervisorTurnService(null!, null!, null!, db: Infrastructure.EmptyTestDb.New(), grader, null!, null!, null!, null!, null!, null!, new NoManifests(), new FakeSupervisorPublishedBranchResolver(), null!, new AdmitAllBudgetLedger(), - null!, NullLogger.Instance, rubricJudge); + null!, null!, NullLogger.Instance, rubricJudge); var context = new SupervisorTurnContext { diff --git a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorDependencyStagingTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorDependencyStagingTests.cs index 518f78b7d..456afc9ec 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorDependencyStagingTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorDependencyStagingTests.cs @@ -446,8 +446,8 @@ public void A_unit_resumed_from_a_host_loss_checkpoint_carries_its_provenance_an { // A checkpoint is not a captured transcript. The attempt that wrote it never finished: its machine is gone, // so the conversation may describe turns the checkpoint never saw and edits the new sandbox does not - // contain — which is true even when a branch WAS pushed, because the unpublished remainder died with the - // host. The ordinary honest-redo line only covers "no branch to continue from", a smaller claim. The whole + // contain — which is true even when a branch WAS pushed, because the unpublished remainder was lost with that + // attempt. The ordinary honest-redo line only covers "no branch to continue from", a smaller claim. The whole // goal is asserted, so each arm pins exactly WHICH tree sentence follows the preamble. // MUTATION: drop the CheckpointAt branch from ApplyResumeRecord → the task carries no provenance, is not // marked a checkpoint (so an unreadable ref would FAIL the attempt instead of degrading), and the goal says diff --git a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorTurnServiceTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorTurnServiceTests.cs index d24bb4bc6..61e0ad26c 100644 --- a/backend/tests/CodeSpace.UnitTests/Agents/SupervisorTurnServiceTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Agents/SupervisorTurnServiceTests.cs @@ -133,7 +133,7 @@ public async Task The_no_progress_guard_forces_a_clean_terminal_stop() ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Plan, $$"""{"turn":{{i}}}""", "{}"); // A decider that would NEVER stop on its own — proving the bound, not the decider, terminates. - var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: null, goalConfig: null, CancellationToken.None); @@ -149,7 +149,7 @@ public async Task A_budget_ledger_refusal_forces_the_cost_cap_stop_not_an_infra_ // land the SAME cost-cap terminal the realized-spend bound reaches, one call earlier; never a park-and- // retry loop (the ledger will refuse forever) and never an exception-shaped run failure. var ledger = new FakeSupervisorDecisionLog(); - var service = new SupervisorTurnService(ledger, new BudgetRefusedDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new BudgetRefusedDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: null, goalConfig: null, CancellationToken.None); @@ -189,7 +189,7 @@ public async Task A_merge_executor_runs_under_its_own_recorded_and_budgeted_synt { var ledger = new FakeSupervisorDecisionLog(); var executor = new ScopeObservingExecutor(); - var service = new SupervisorTurnService(ledger, new MergeDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new MergeDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: null, goalConfig: null, CancellationToken.None); @@ -216,7 +216,7 @@ public async Task A_no_progress_forced_stop_under_a_required_delivery_contract_i for (var i = 0; i < SupervisorLane.DefaultMaxNoProgressDecisions; i++) ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Plan, $$"""{"turn":{{i}}}""", "{}"); - var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: Guid.NewGuid(), goalConfig: DeliveryGoalConfig(), CancellationToken.None); @@ -234,7 +234,7 @@ public async Task A_forced_stop_over_an_unsatisfied_publish_parks_on_the_deliver // The latest publish ran and found NOTHING to open a PR from — unsatisfied, and never satisfied by absence (H1). ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Publish, "{}", """{"pullRequests":[]}"""); - var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: Guid.NewGuid(), goalConfig: DeliveryGoalConfig(), CancellationToken.None); @@ -253,7 +253,7 @@ public async Task A_forced_stop_over_an_unsatisfied_publish_with_no_conversation ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Publish, "{}", """{"pullRequests":[]}"""); - var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: null, goalConfig: DeliveryGoalConfig(), CancellationToken.None); @@ -277,7 +277,7 @@ public async Task An_unanswered_gate_card_on_the_tape_fuses_the_forced_stop_back ledger.SeedTerminal(_runId, _teamId, SupervisorDecisionKinds.Publish, "{}", """{"pullRequests":[]}"""); var goalConfig = DeliveryGoalConfig(); - var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var context = await service.RehydrateFromDecisionLogAsync(_runId, _teamId, "sup", "goal", goalConfig, CancellationToken.None); // The gate's own card, already degraded once: terminal, question pinned to the gate prefix, NO answer. @@ -296,7 +296,7 @@ public async Task A_depth_capped_run_is_never_gated_into_delivering() // parent. Gating it would let a run that should never have taken a single decision open PRs, and (with no // conversation) erase DepthCapExceeded behind DeliveryAdjudicationUnavailable in the scorecard. var ledger = new FakeSupervisorDecisionLog(); - var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var goalConfig = DeliveryGoalConfig(); var context = await service.RehydrateFromDecisionLogAsync(_runId, _teamId, "sup", "goal", goalConfig, CancellationToken.None); @@ -395,7 +395,7 @@ public void A_governance_denied_side_effecting_decision_force_stops_and_stages_n // — the same forward-compat exposure a future irreversible/merge-PR policy would open. Asserts the gate // turns the denied side effect into a force-STOP carrying the GovernanceDenied reason and stages NO agent. var executor = new CountingExecutor(); - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var context = new SupervisorTurnContext { Goal = "goal", TurnNumber = 0, ApprovalPolicy = (SupervisorApprovalPolicy)999 }; var spawn = new SupervisorDecision { Kind = kind, PayloadJson = """{"subtaskIds":["a","b"]}""" }; @@ -417,7 +417,7 @@ public void A_resolve_parks_for_approval_under_the_autonomous_policy_but_a_spawn // The safety floor wired end-to-end through the REAL gate: under None (autonomous) a spawn runs unchanged, // but a resolve — which dispatches an agent to autonomously RE-MERGE code — escalates to a human approval // card (it parks), because GateSideEffectingDecision passes irreversible=IsIrreversible(kind) for resolve. - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new CountingExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new CountingExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var context = new SupervisorTurnContext { Goal = "goal", SupervisorRunId = _runId, TeamId = _teamId, NodeId = "sup", TurnNumber = 1, ApprovalPolicy = SupervisorApprovalPolicy.None, ConversationId = Guid.NewGuid() }; @@ -774,7 +774,7 @@ public async Task Replaying_a_settled_turn_does_not_double_execute() { var ledger = new FakeSupervisorDecisionLog(); var executor = new CountingExecutor(); - var service = new SupervisorTurnService(ledger, new StubSupervisorDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, new StubSupervisorDecider(), executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); // First pass: turn 0 (plan) executes once + records terminal. await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: null, goalConfig: null, CancellationToken.None); @@ -813,7 +813,7 @@ public async Task Replaying_an_in_flight_turn_re_executes_the_frozen_decision_ev var decider = new NonDeterministicDecider(SupervisorDecisionKinds.Plan, plannedA, SupervisorDecisionKinds.Stop, """{"reason":"divergent-B"}"""); var executor = new CountingExecutor(); - var service = new SupervisorTurnService(ledger, decider, executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(ledger, decider, executor, db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var result = await service.RunTurnAsync(_runId, _teamId, "sup", "goal", conversationId: null, goalConfig: null, CancellationToken.None); @@ -831,12 +831,12 @@ public async Task Replaying_an_in_flight_turn_re_executes_the_frozen_decision_ev } private SupervisorTurnService Service(FakeSupervisorDecisionLog ledger) => - new(ledger, new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + new(ledger, new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); // ── L4 P1 stop-acceptance test helpers ────────────────────────────────────────── private SupervisorTurnService ServiceWith(FakeSupervisorDecisionLog ledger, ISupervisorDecider decider, FakeAcceptanceGrader grader) => - new(ledger, decider, new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), grader, new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + new(ledger, decider, new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), grader, new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); private static SupervisorDecision StopDecision() => new() { Kind = SupervisorDecisionKinds.Stop, PayloadJson = "{}" }; @@ -942,7 +942,7 @@ public async Task An_unconfirmed_plan_preempts_the_decider_with_the_confirmation var row = PlanRow(CodeSpace.Messages.Plans.WorkPlanStatuses.Authored); var store = new FakeWorkPlanStore(row); // AlwaysPlanDecider would PLAN — the ask_human coming back proves the gate preempted the brain entirely. - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var decision = await service.ChooseDecisionAsync(GateContext(TerminalPlan()), SupervisorGoalPlan.From(null), depth: 0, CancellationToken.None); @@ -962,7 +962,7 @@ public async Task A_just_answered_confirmation_releases_to_the_decider_and_settl { var row = PlanRow(CodeSpace.Messages.Plans.WorkPlanStatuses.AwaitingConfirmation); var store = new FakeWorkPlanStore(row); - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var decision = await service.ChooseDecisionAsync(GateContext(TerminalPlan(), AnsweredConfirmation(answer)), SupervisorGoalPlan.From(null), depth: 0, CancellationToken.None); @@ -976,7 +976,7 @@ public async Task No_conversation_surface_force_stops_rather_than_degrading_the_ // A task launch wires no conversation — the injected card would DEGRADE to a no-surface self-advance and // agents would spawn unconfirmed (the review's blocker). The gate must stop the run instead. var store = new FakeWorkPlanStore(PlanRow(CodeSpace.Messages.Plans.WorkPlanStatuses.Authored)); - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var decision = await service.ChooseDecisionAsync(GateContext(conversationId: null, TerminalPlan()), SupervisorGoalPlan.From(null), depth: 0, CancellationToken.None); @@ -991,7 +991,7 @@ public async Task No_conversation_surface_force_stops_rather_than_degrading_the_ public void A_spawn_while_the_latest_plan_stands_rejected_is_refused(string kind) { // Prompt-following is not a guarantee — the structural floor refuses to execute a REJECTED plan version. - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var context = GateContext(TerminalPlan(), AnsweredConfirmation("revise: do not touch the DB")); var gated = service.ApplyPostDecisionGate(context, SupervisorGoalPlan.From(null), new SupervisorDecision { Kind = kind, PayloadJson = """{"subtaskIds":["sa","sb"]}""" }); @@ -1005,7 +1005,7 @@ public void A_stop_owing_an_approved_amendment_is_rewritten_into_its_retry() { // B5 (MAJOR-5): the fold's already-graded guard means an amendment only affects a FUTURE attempt — a stop // before that attempt terminalizes on the dead oracle's verdict, silently dropping what a human co-signed. - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var card = SupervisorAmendAcceptance.IntoAskHuman(new SupervisorAmendAcceptancePayload { @@ -1028,7 +1028,7 @@ public void A_stop_owing_an_approved_amendment_is_rewritten_into_its_retry() [Fact] public void A_stop_after_the_owed_retry_ran_passes_the_obligation_gate() { - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var card = SupervisorAmendAcceptance.IntoAskHuman(new SupervisorAmendAcceptancePayload { @@ -1053,7 +1053,7 @@ public void A_forced_stop_names_the_obligation_it_drops() // B5 (A2 ruling): a bound may stop a run before it consumed an approved amendment — the server never spends // into a tripped cap — but the drop is NAMED on the stop, never silent. Driven through the invalid-plan // forced stop (a pure ApplyPostDecisionGate path), the same GateForcedStop every bound routes through. - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var card = SupervisorAmendAcceptance.IntoAskHuman(new SupervisorAmendAcceptancePayload { @@ -1084,7 +1084,7 @@ private static SupervisorPriorDecision StagingPrior(long seq, string subtaskId) [Fact] public void A_revised_plan_clears_the_rejected_floor() { - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var context = GateContext(TerminalPlan(), AnsweredConfirmation("revise: merge"), TerminalPlan()); var spawn = new SupervisorDecision { Kind = SupervisorDecisionKinds.Spawn, PayloadJson = """{"subtaskIds":["sa"]}""" }; @@ -1100,7 +1100,7 @@ public async Task The_release_flip_lands_even_when_a_pre_bound_stop_takes_the_tu // strand AwaitingConfirmation — the answered confirmation must settle FIRST. var row = PlanRow(CodeSpace.Messages.Plans.WorkPlanStatuses.AwaitingConfirmation); var store = new FakeWorkPlanStore(row); - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); // NoProgressDecisions at the cap → the pre-decision no-progress guard force-stops this turn. var context = GateContext(TerminalPlan(), AnsweredConfirmation("approve")) with { NoProgressDecisions = SupervisorLane.DefaultMaxNoProgressDecisions }; @@ -1115,7 +1115,7 @@ public async Task The_release_flip_lands_even_when_a_pre_bound_stop_takes_the_tu public async Task A_missing_plan_row_degrades_open_rather_than_parking_on_an_unreviewable_card() { var store = new FakeWorkPlanStore(); // no row — unreachable in production (persist precedes the terminal plan) - var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), NullLogger.Instance); + var service = new SupervisorTurnService(new FakeSupervisorDecisionLog(), new AlwaysPlanDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), store, null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), new NoLessonsReader(), null!, NullLogger.Instance); var decision = await service.ChooseDecisionAsync(GateContext(TerminalPlan()), SupervisorGoalPlan.From(null), depth: 0, CancellationToken.None); diff --git a/backend/tests/CodeSpace.UnitTests/Infrastructure/HeartbeatClock.cs b/backend/tests/CodeSpace.UnitTests/Infrastructure/HeartbeatClock.cs new file mode 100644 index 000000000..37efe2f68 --- /dev/null +++ b/backend/tests/CodeSpace.UnitTests/Infrastructure/HeartbeatClock.cs @@ -0,0 +1,33 @@ +using Microsoft.Extensions.Time.Testing; +using Shouldly; + +namespace CodeSpace.UnitTests.Infrastructure; + +/// +/// A fake clock for a heartbeat loop that says when the loop has ARMED its next sleep. A loop registers its next timer +/// only after the previous beat's work returns, so an advance issued before that registration moves a clock nothing is +/// waiting on yet: the beat is then missed, or lands late. Waiting for the timer before every advance lets the fake +/// clock decide not only HOW MANY beats land but exactly WHEN each one does — which is what lets a test pin the interval +/// itself instead of only the count. The same rendezvous AgentProgressLeaseTests uses for its renewal loop. +/// +public sealed class HeartbeatClock : TimeProvider +{ + private readonly FakeTimeProvider _time = new(); + private readonly SemaphoreSlim _armed = new(0); + + public override DateTimeOffset GetUtcNow() => _time.GetUtcNow(); + public override long GetTimestamp() => _time.GetTimestamp(); + public override long TimestampFrequency => _time.TimestampFrequency; + + public void Advance(TimeSpan amount) => _time.Advance(amount); + + public override ITimer CreateTimer(TimerCallback callback, object? state, TimeSpan dueTime, TimeSpan period) + { + var timer = _time.CreateTimer(callback, state, dueTime, period); + _armed.Release(); + return timer; + } + + /// Returns once the loop is asleep on a timer armed at the clock's current instant. The bound only detects a loop that never re-arms; it never lets time pass on the fake clock. + public async Task NextTimerAsync() => (await _armed.WaitAsync(TimeSpan.FromSeconds(10))).ShouldBeTrue("the heartbeat loop must arm its next sleep — it stopped repeating, or it is not sleeping on the clock it was handed"); +} diff --git a/backend/tests/CodeSpace.UnitTests/Learning/LessonArmFreezeTests.cs b/backend/tests/CodeSpace.UnitTests/Learning/LessonArmFreezeTests.cs index cde93d336..d2de15a87 100644 --- a/backend/tests/CodeSpace.UnitTests/Learning/LessonArmFreezeTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Learning/LessonArmFreezeTests.cs @@ -112,5 +112,5 @@ public Task> ListCurrentAsync(LessonReadRequest request, C } private static SupervisorTurnService Service(FakeSupervisorDecisionLog ledger, ILessonReader lessons) => - new(ledger, new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), lessons, NullLogger.Instance); + new(ledger, new StubSupervisorDecider(), new StubSupervisorActionExecutor(), db: Infrastructure.EmptyTestDb.New(), new FakeAcceptanceGrader(), new FakeDecisionQueue(), new FakeDecisionArbiter(), new FakeDecisionAnswerService(), new FakeWorkPlanStore(), null!, null!, new FakePublishManifestStore(), new FakeSupervisorPublishedBranchResolver(), new NullCompletionComposer(), new AdmitAllBudgetLedger(), lessons, null!, NullLogger.Instance); } diff --git a/backend/tests/CodeSpace.UnitTests/Supervisor/SupervisorGradingHeartbeatTests.cs b/backend/tests/CodeSpace.UnitTests/Supervisor/SupervisorGradingHeartbeatTests.cs index 01833779c..b78dc06a6 100644 --- a/backend/tests/CodeSpace.UnitTests/Supervisor/SupervisorGradingHeartbeatTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Supervisor/SupervisorGradingHeartbeatTests.cs @@ -1,7 +1,9 @@ using CodeSpace.Core.Services.Supervisor; using CodeSpace.Core.Services.Workflows.Lifecycle; +using CodeSpace.Core.Services.Workflows.Reconciliation; +using CodeSpace.Messages.Agents.Benchmark; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; -using Microsoft.Extensions.Time.Testing; using Shouldly; using System.Text.Json; @@ -9,12 +11,12 @@ namespace CodeSpace.UnitTests.Supervisor; /// /// 🟢 Unit: the P1.3 grading-heartbeat loop () — -/// pins the cancellation contract that keeps a long acceptance grade from looking abandoned to the reconciler -/// WITHOUT waiting out the real 90s production interval: a decides when each interval -/// has elapsed, so the production value itself costs nothing. The DB-observable effect (a real ledger row landing, -/// and the reconciler reading it as liveness) is proved at the integration tier — this pins the pure loop mechanics: -/// it logs once per elapsed interval while un-cancelled, stops the instant it's cancelled (mid-sleep or between -/// ticks), and never lets OperationCanceledException escape. +/// pins the contract that keeps a long acceptance grade from looking abandoned to the reconciler WITHOUT waiting out +/// the real 90s production interval: a fake clock decides when each interval has elapsed, so the production value +/// itself costs nothing. The DB-observable effects (a real ledger row landing on its own context while the grade holds +/// the scope's, and the reconciler reading such a row as liveness) are proved at the integration tier — this pins the +/// loop mechanics: one pulse per elapsed interval and not a moment sooner, each on a DI scope of its own, a failed pulse +/// reported and survived rather than surfaced into the grade, and a quiet stop the instant it is cancelled. /// /// These once slept on the wall clock — "three 15ms ticks within 10s", "cancel after 30ms" — which proved the /// loop repeats only as reliably as the runner happened to schedule it. On the fake clock the same properties are @@ -25,30 +27,40 @@ public class SupervisorGradingHeartbeatTests { private static readonly Guid RunId = Guid.NewGuid(); private const string NodeId = "sup"; + private const string Graded = "graded while every pulse failed"; + private static readonly TimeSpan Interval = SupervisorLane.AcceptanceGradeHeartbeatInterval; - // RunGradingHeartbeatLoopAsync touches ONLY _recordLogger — every other dependency is stored by the ctor - // (plain field assignment, no eager calls) and never read on this path, so null! is safe here exactly as the - // existing SupervisorTurnServiceTests already pass null! for db/offloader on paths that don't touch them. - private static SupervisorTurnService Service(IRunRecordLogger logger) => - new(null!, null!, null!, db: Infrastructure.EmptyTestDb.New(), null!, null!, null!, null!, null!, logger, null!, null!, null!, new NullCompletionComposer(), null!, null!, NullLogger.Instance); + /// How far short of a full interval the "not yet" check stands: a loop sleeping even this much less than it was handed pulses inside that window. + private static readonly TimeSpan Step = TimeSpan.FromSeconds(1); + + // The loop touches only the scope factory (each pulse's record logger) and the ILogger (a failed pulse's warning) — + // every other dependency is stored by the ctor and never read on this path, so null! is safe here exactly as the + // existing SupervisorTurnServiceTests pass null! for seams a path doesn't touch. The ctor's OWN record logger is a + // recording fake rather than null!, so a pulse that wrote through it — the grade's shared context — is caught. + private static SupervisorTurnService Service(PulseScopes pulses, RecordingLogger? injected = null, CapturingLogger? warnings = null) => + new(null!, null!, null!, db: Infrastructure.EmptyTestDb.New(), null!, null!, null!, null!, null!, injected ?? new RecordingLogger(), null!, null!, null!, new NullCompletionComposer(), null!, null!, pulses, warnings ?? (Microsoft.Extensions.Logging.ILogger)NullLogger.Instance); [Fact] public async Task The_loop_logs_once_per_elapsed_interval_while_uncancelled() { - var time = new FakeTimeProvider(); - var interval = SupervisorLane.AcceptanceGradeHeartbeatInterval; + var time = new Infrastructure.HeartbeatClock(); var logger = new RecordingLogger(); using var cts = new CancellationTokenSource(); - var loop = Service(logger).RunGradingHeartbeatLoopAsync(RunId, NodeId, interval, cts.Token, time); - - logger.Calls.ShouldBeEmpty("the first heartbeat is owed only once a full interval of grading has passed"); + var loop = Service(new PulseScopes(logger)).RunGradingHeartbeatLoopAsync(RunId, NodeId, Interval, cts.Token, time); - // Three OBSERVED ticks prove it is a REPEATING loop, not a one-shot — and on the fake clock the count after - // i intervals is exactly i, where the wall clock could only ever promise "at least". + // Three OBSERVED ticks prove it is a REPEATING loop, not a one-shot; the two advances around each one pin WHEN it + // lands, not only that it did — nothing a step short of the interval, exactly one the moment it completes. A loop + // that sleeps anything but the interval it was handed reds on one side or the other. for (var i = 1; i <= 3; i++) { - await AdvanceUntilLoggedAsync(time, logger, interval, i); + await time.NextTimerAsync(); + + time.Advance(Interval - Step); + (await logger.Logged.WaitAsync(TimeSpan.FromMilliseconds(100))).ShouldBeFalse($"heartbeat {i} is owed only once a FULL interval of grading has passed — one {Step.TotalSeconds}s early means the loop sleeps less than it was handed"); + + time.Advance(Step); + (await logger.Logged.WaitAsync(TimeSpan.FromSeconds(10))).ShouldBeTrue($"heartbeat {i} is owed the instant its interval completes"); logger.Calls.Count.ShouldBe(i, $"exactly one heartbeat per elapsed interval — after {i} interval(s) there must be {i}"); } @@ -62,15 +74,17 @@ public async Task The_loop_logs_once_per_elapsed_interval_while_uncancelled() [Fact] public async Task Cancelling_stops_the_loop_without_throwing() { - var time = new FakeTimeProvider(); - var interval = SupervisorLane.AcceptanceGradeHeartbeatInterval; + var time = new Infrastructure.HeartbeatClock(); var logger = new RecordingLogger(); using var cts = new CancellationTokenSource(); - var loop = Service(logger).RunGradingHeartbeatLoopAsync(RunId, NodeId, interval, cts.Token, time); + var loop = Service(new PulseScopes(logger)).RunGradingHeartbeatLoopAsync(RunId, NodeId, Interval, cts.Token, time); // One tick first, so the cancel below lands MID-SLEEP on the next interval rather than before the loop ran. - await AdvanceUntilLoggedAsync(time, logger, interval, 1); + await time.NextTimerAsync(); + time.Advance(Interval); + (await logger.Logged.WaitAsync(TimeSpan.FromSeconds(10))).ShouldBeTrue("the first heartbeat is owed once its interval has passed"); + await time.NextTimerAsync(); cts.Cancel(); @@ -82,7 +96,7 @@ public async Task Cancelling_stops_the_loop_without_throwing() escaped.ShouldBeNull("a cancel mid-sleep must end the loop at once and quietly — a TimeoutException means it outlived the grade it protects, a TaskCanceledException that the cancel escaped to the caller"); - time.Advance(interval * 3); + time.Advance(Interval * 3); (await logger.Logged.WaitAsync(TimeSpan.FromMilliseconds(200))).ShouldBeFalse("a cancelled loop logs no more heartbeats, however much time passes"); logger.Calls.Count.ShouldBe(1); } @@ -91,37 +105,150 @@ public async Task Cancelling_stops_the_loop_without_throwing() public async Task An_already_cancelled_token_produces_zero_heartbeats() { var logger = new RecordingLogger(); + var pulses = new PulseScopes(logger); using var cts = new CancellationTokenSource(); cts.Cancel(); - await Service(logger).RunGradingHeartbeatLoopAsync(RunId, NodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, cts.Token, new FakeTimeProvider()).WaitAsync(TimeSpan.FromSeconds(10)); + await Service(pulses).RunGradingHeartbeatLoopAsync(RunId, NodeId, Interval, cts.Token, new Infrastructure.HeartbeatClock()).WaitAsync(TimeSpan.FromSeconds(10)); logger.Calls.ShouldBeEmpty("a grade that finishes before the FIRST tick never needs a heartbeat"); + pulses.Created.ShouldBe(0, "no pulse, no scope"); + } + + /// + /// Every call site awaits this loop in the finally around the grade it protects, so a loop that ends faulted + /// REPLACES the grade with its error. The live shape was a pulse writing through the grade's own DbContext mid-query + /// — EF's second-operation guard, an — but any failed write is the same shape. + /// A failed pulse is a warning naming the run and node, the next pulse still fires, and the grade comes back graded. + /// + [Fact] + public async Task A_failed_pulse_is_reported_and_never_replaces_the_grade() + { + var time = new Infrastructure.HeartbeatClock(); + var failing = new RecordingLogger(new InvalidOperationException("A second operation was started on this context instance before a previous operation completed.")); + var warnings = new CapturingLogger(); + using var heartbeatCts = new CancellationTokenSource(); + + var heartbeat = Service(new PulseScopes(failing), warnings: warnings).RunGradingHeartbeatLoopAsync(RunId, NodeId, Interval, heartbeatCts.Token, time); + + BenchmarkGrade grade; + try + { + for (var pulse = 1; pulse <= 3; pulse++) + { + await time.NextTimerAsync(); + time.Advance(Interval); + + (await failing.Logged.WaitAsync(TimeSpan.FromSeconds(10))).ShouldBeTrue($"pulse {pulse} must still fire after {pulse - 1} failed"); + } + + grade = new BenchmarkGrade { Passed = true, Detail = Graded }; + } + finally + { + // The call sites' own shape: whatever this await surfaces is what the grade becomes. + heartbeatCts.Cancel(); + + try { await heartbeat; } + catch (OperationCanceledException) { } + } + + grade.Detail.ShouldBe(Graded, "a failed pulse is a missed log line — it must never replace the grade it was keeping alive"); + heartbeat.IsCompletedSuccessfully.ShouldBeTrue("the loop ends only by cancellation, and quietly"); + + warnings.Entries.Count.ShouldBe(3, "one warning per failed pulse — none swallowed silently, none doubled"); + warnings.Entries.ShouldAllBe(w => w.Level == Microsoft.Extensions.Logging.LogLevel.Warning && w.Exception is InvalidOperationException); + warnings.Entries.ShouldAllBe(w => Equals(w.Properties["SupervisorRunId"], RunId) && Equals(w.Properties["NodeId"], NodeId), "a structured template naming the run and node, not an interpolated string"); + } + + /// + /// The pulse runs concurrently with the grade, and the service's own scope holds the DbContext the grade is using, so + /// each pulse writes through a scope of its OWN — one per pulse, disposed once the pulse returns — and never through the + /// constructor-injected record logger, which shares that context. The collision itself is proved against real + /// Postgres by SupervisorGradingHeartbeatIsolationFlowTests; this pins the wiring that avoids it. + /// + [Fact] + public async Task Each_pulse_writes_through_a_record_logger_from_a_scope_of_its_own() + { + var time = new Infrastructure.HeartbeatClock(); + var scoped = new RecordingLogger(); + var injected = new RecordingLogger(); + var pulses = new PulseScopes(scoped); + using var cts = new CancellationTokenSource(); + + var loop = Service(pulses, injected).RunGradingHeartbeatLoopAsync(RunId, NodeId, Interval, cts.Token, time); + + await time.NextTimerAsync(); + + for (var pulse = 1; pulse <= 3; pulse++) + { + time.Advance(Interval); + + // Re-armed ⇒ the pulse's write has returned and its scope's using block has closed. + await time.NextTimerAsync(); + + scoped.Calls.Count.ShouldBe(pulse, "every pulse writes through the record logger its own scope serves"); + pulses.Created.ShouldBe(pulse, "one fresh scope per pulse — not one per loop"); + pulses.Disposed.ShouldBe(pulse, "each pulse's scope is disposed as soon as its write returns"); + } + + cts.Cancel(); + await loop.WaitAsync(TimeSpan.FromSeconds(10)); + + injected.Calls.ShouldBeEmpty("a pulse never writes through the service's own record logger — that one shares the grade's DbContext"); } /// - /// Advances the fake clock until the loop logs heartbeat , rather than advancing one - /// whole interval and assuming the loop was already listening. - /// - /// The loop arms its next delay only AFTER the previous heartbeat's write returns, so a single Advance can - /// land before that registration and be missed — the clock then never moves again. Nudging in tenths of an - /// interval cannot fire a delay early, so the exact count asserted at the call site still means one heartbeat - /// per interval. Same shape as HeartbeatLoopTests. + /// The cadence production runs this heartbeat at, as a NUMBER. The fake clock above proves the loop honours whatever + /// interval it is handed and says nothing about which one production hands it — the split AgentRunLivenessTests + /// makes for the agent heartbeat. 's own doc calls the + /// value pinned; nothing pinned it. /// - private static async Task AdvanceUntilLoggedAsync(FakeTimeProvider time, RecordingLogger logger, TimeSpan interval, int ordinal) + [Fact] + public void The_production_cadence_is_pinned_well_inside_the_liveness_window() { - for (var nudge = 0; nudge < 200; nudge++) + SupervisorLane.AcceptanceGradeHeartbeatInterval.ShouldBe(TimeSpan.FromSeconds(90)); + + (StuckRunReconcilerService.LedgerLivenessWindow / SupervisorLane.AcceptanceGradeHeartbeatInterval).ShouldBeGreaterThanOrEqualTo(3, "a failed pulse is only a warning, so the liveness window must outlast a missed pulse or two — three pulses per window absorbs two in a row"); + } + + /// The scope factory a pulse opens its scope through: every scope serves as its , and the factory counts the scopes opened and disposed. + private sealed class PulseScopes(IRunRecordLogger logger) : IServiceScopeFactory + { + private int _created; + private int _disposed; + + public int Created => Volatile.Read(ref _created); + public int Disposed => Volatile.Read(ref _disposed); + + public IServiceScope CreateScope() { - if (await logger.Logged.WaitAsync(TimeSpan.FromMilliseconds(10))) return; + Interlocked.Increment(ref _created); + return new Scope(this, logger); + } - time.Advance(interval / 10); + private sealed class Scope(PulseScopes owner, IRunRecordLogger logger) : IServiceScope, IServiceProvider + { + public IServiceProvider ServiceProvider => this; + public object? GetService(Type serviceType) => serviceType == typeof(IRunRecordLogger) ? logger : null; + public void Dispose() => Interlocked.Increment(ref owner._disposed); } + } + + /// Captures each entry's level, exception and structured properties — the named template values an interpolated message would not carry. + private sealed class CapturingLogger : Microsoft.Extensions.Logging.ILogger + { + public List<(Microsoft.Extensions.Logging.LogLevel Level, Exception? Exception, IReadOnlyDictionary Properties)> Entries { get; } = new(); + + public IDisposable? BeginScope(TState state) where TState : notnull => null; + public bool IsEnabled(Microsoft.Extensions.Logging.LogLevel logLevel) => true; - throw new TimeoutException($"heartbeat {ordinal} never landed after advancing the fake clock well past its interval — the loop stopped repeating, or RunGradingHeartbeatLoopAsync is not sleeping on the TimeProvider it is handed"); + public void Log(Microsoft.Extensions.Logging.LogLevel logLevel, Microsoft.Extensions.Logging.EventId eventId, TState state, Exception? exception, Func formatter) => + Entries.Add((logLevel, exception, (state as IEnumerable> ?? []).ToDictionary(p => p.Key, p => p.Value))); } - /// Minimal fake — every method a harmless no-op except , which records each call and releases once it has. The grading-heartbeat path touches ONLY LogAsync; every other member exists solely to satisfy the interface. - private sealed class RecordingLogger : IRunRecordLogger + /// Minimal fake — every method a harmless no-op except , which records each call, releases once it has, and then fails with when one is given (a pulse whose write failed). The grading-heartbeat path touches ONLY LogAsync; every other member exists solely to satisfy the interface. + private sealed class RecordingLogger(Exception? fault = null) : IRunRecordLogger { public List<(Guid RunId, string? NodeId, LogLevel Level, string Message)> Calls { get; } = new(); @@ -132,7 +259,7 @@ public Task LogAsync(Guid runId, string? nodeId, LogLevel level, string message, { Calls.Add((runId, nodeId, level, message)); Logged.Release(); - return Task.CompletedTask; + return fault is null ? Task.CompletedTask : Task.FromException(fault); } public Task RunQueuedAsync(Guid runId, string sourceType, Guid? actorId, CancellationToken cancellationToken) => Task.CompletedTask; diff --git a/backend/tests/CodeSpace.UnitTests/Workflows/AgentCodeNodeTests.cs b/backend/tests/CodeSpace.UnitTests/Workflows/AgentCodeNodeTests.cs index 639409571..eeb97aba6 100644 --- a/backend/tests/CodeSpace.UnitTests/Workflows/AgentCodeNodeTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Workflows/AgentCodeNodeTests.cs @@ -1244,7 +1244,7 @@ public void The_honest_no_continuity_hint_is_pinned_verbatim() // agent the same thing about a tree that does not carry its prior work. The literal is pinned because the // supervisor's own behaviour test asserts this exact wording — a reword must be a visible decision. AgentRetryContinuity.HonestNoContinuityHint.ShouldBe( - "Note: your prior attempt's conversation is restored, but its git changes were NOT preserved in this workspace (your prior attempt pushed no branch of its own) — you must redo any relevant file changes from scratch."); + "Note: your prior attempt's conversation is restored, but its git changes were NOT preserved in this workspace (your prior attempt pushed no branch of its own for this repository) — you must redo any relevant file changes from scratch."); } [Fact] diff --git a/backend/tests/CodeSpace.UnitTests/Workflows/HeartbeatLoopTests.cs b/backend/tests/CodeSpace.UnitTests/Workflows/HeartbeatLoopTests.cs index c13c95374..fb02e182a 100644 --- a/backend/tests/CodeSpace.UnitTests/Workflows/HeartbeatLoopTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Workflows/HeartbeatLoopTests.cs @@ -1,4 +1,5 @@ using CodeSpace.Core.Services.Agents; +using CodeSpace.UnitTests.Infrastructure; using Microsoft.Extensions.Time.Testing; using Shouldly; @@ -20,10 +21,13 @@ namespace CodeSpace.UnitTests.Workflows; [Trait("Category", "Unit")] public class HeartbeatLoopTests { + /// How far short of a full interval the "not yet" check stands. + private static readonly TimeSpan Step = TimeSpan.FromSeconds(1); + [Fact] public async Task Pings_once_per_interval_until_cancelled() { - var time = new FakeTimeProvider(); + var time = new HeartbeatClock(); var interval = TimeSpan.FromSeconds(30); var pinged = new SemaphoreSlim(0); var count = 0; @@ -40,7 +44,7 @@ public async Task Pings_once_per_interval_until_cancelled() for (var i = 1; i <= 3; i++) { - await AdvanceUntilPingedAsync(time, pinged, interval, i); + await AdvanceOneIntervalAsync(time, pinged, interval, i); Volatile.Read(ref count).ShouldBe(i, $"exactly one ping per elapsed interval — after {i} interval(s) there must be {i}, not 'at least' {i}"); } @@ -53,31 +57,30 @@ public async Task Pings_once_per_interval_until_cancelled() } /// - /// Advances the fake clock until the loop signals , rather than advancing once and - /// assuming it was listening. + /// Advances the fake clock across ONE interval in two moves and pins what each must do: a short of + /// the interval the loop must still be asleep, and the step that completes it must produce exactly the beat. A loop + /// that sleeps anything but the interval it was handed reds on one side or the other — where nudging in tenths until + /// SOMETHING arrived let any sleep between a tenth of the interval and twenty of them pass. /// - /// The loop arms its next timer inside Task.Delay AFTER the previous ping returns, so a - /// single Advance can land in the window before that registration and be missed entirely — the - /// clock then never moves again and the wait burns its full timeout. That is the race this test - /// kept losing. Nudging in fractions of an interval cannot fire a timer early, and the count - /// assertion at the call site is what still proves one ping per interval. + /// The loop arms its next timer inside Task.Delay AFTER the previous ping returns, so an Advance issued before + /// that registration moves a clock nothing is waiting on yet — the race this test once kept losing. Waiting for the + /// timer first closes it without letting any time pass. /// - private static async Task AdvanceUntilPingedAsync(FakeTimeProvider time, SemaphoreSlim ticked, TimeSpan interval, int ordinal) + private static async Task AdvanceOneIntervalAsync(HeartbeatClock time, SemaphoreSlim ticked, TimeSpan interval, int ordinal) { - for (var nudge = 0; nudge < 200; nudge++) - { - if (await ticked.WaitAsync(TimeSpan.FromMilliseconds(10))) return; + await time.NextTimerAsync(); - time.Advance(interval / 10); - } + time.Advance(interval - Step); + (await ticked.WaitAsync(TimeSpan.FromMilliseconds(100))).ShouldBeFalse($"beat {ordinal} arrived {Step.TotalSeconds}s before its interval completed — the loop sleeps less than it was handed"); - throw new TimeoutException($"ping {ordinal} never arrived after advancing the fake clock well past its interval"); + time.Advance(Step); + (await ticked.WaitAsync(TimeSpan.FromSeconds(10))).ShouldBeTrue($"beat {ordinal} never arrived once its interval completed — the loop sleeps longer than it was handed, or stopped repeating"); } [Fact] public async Task A_failing_ping_is_reported_but_does_not_kill_the_loop() { - var time = new FakeTimeProvider(); + var time = new HeartbeatClock(); var interval = TimeSpan.FromSeconds(30); var reported = new SemaphoreSlim(0); var pings = 0; @@ -95,7 +98,7 @@ public async Task A_failing_ping_is_reported_but_does_not_kill_the_loop() for (var i = 1; i <= 3; i++) { - await AdvanceUntilPingedAsync(time, reported, interval, i); + await AdvanceOneIntervalAsync(time, reported, interval, i); Volatile.Read(ref pings).ShouldBe(i, "a throwing ping must not stop, skip, or double the cadence"); Volatile.Read(ref errors).ShouldBe(i, "every failed ping is reported exactly once — none aborted the loop");