From a3cfb37b13d23230beaab6441114c71220dfe009 Mon Sep 17 00:00:00 2001 From: "Mars.P" Date: Wed, 30 Sep 2026 22:13:27 +0800 Subject: [PATCH] Isolate approval expiry follow-ups; skip orphaned suspended children The approval sweep wakes each expired row's blocked call and mirrors its card, neither step guarded. Under the job's transactional command the post-commit drain already kept a failure to its own row, but logged it without the ledger id, and a wake that threw skipped the same row's mirror. With no transaction open, one card's failed save ended the follow-ups of every later row: their calls stayed unwoken and their cards Open, and no later tick selects an Expired row again. Each step now catches and logs its own failure against the ledger id, as the decision sweep's do, and a shutdown's cancellation still stops the sweep. The command stays transactional. The stranded-Suspended sweep revived a sub-workflow child whose parent had finished (its cancel failed in the parent's teardown, then its waits were closed), so it ran for a parent no one was waiting on. It now leaves such children alone, keyed on source type because ParentRunId is also a rerun's lineage, and logs how many it left. The exclusion is in the query so a child left Suspended cannot take a place in the batch every tick and crowd out the runs that are wanted. --- .../Agents/Mcp/IToolApprovalExpiryService.cs | 54 ++++- .../StuckRunReconcilerService.cs | 43 ++-- .../Agents/ToolApprovalExpiryServiceTests.cs | 162 ++++++++++++++ .../Workflows/StuckRunReconcilerFlowTests.cs | 122 ++++++++++- .../Agents/ToolApprovalExpiryServiceTests.cs | 197 ++++++++++++++++++ 5 files changed, 552 insertions(+), 26 deletions(-) create mode 100644 backend/tests/CodeSpace.UnitTests/Agents/ToolApprovalExpiryServiceTests.cs diff --git a/backend/src/CodeSpace.Core/Services/Agents/Mcp/IToolApprovalExpiryService.cs b/backend/src/CodeSpace.Core/Services/Agents/Mcp/IToolApprovalExpiryService.cs index f2f4e7f79..b0375636e 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/Mcp/IToolApprovalExpiryService.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/Mcp/IToolApprovalExpiryService.cs @@ -20,6 +20,12 @@ namespace CodeSpace.Core.Services.Agents.Mcp; /// connection — sees the COMMITTED Expired terminal and replays it (not the pre-commit AwaitingApproval, /// which would lose the same-pod fast-path the design promises). Called outside a transaction (ad-hoc / tests), the /// deferred action runs inline since the CAS already auto-committed. +/// +/// Each row's wake and card mirror are best-effort, the two on their own. The drain swallows an action that throws, +/// one row at a time and without naming the row; a caller with no transaction runs them inline with nothing to swallow, +/// and there one card's failed mirror would cost every later row its wake and its mirror. So each step catches its own +/// failure and logs the ledger row: a failed mirror still leaves its call woken, and the next row woken and mirrored. The +/// row stands Expired either way; a card left Open is refused a click, the row being the authority. /// public interface IToolApprovalExpiryService { @@ -59,17 +65,45 @@ public async Task ExpireDueAsync(DateTimeOffset now, CancellationToken canc private async Task ResolveAsync(ExpiredToolApproval row, CancellationToken cancellationToken) { - // Best-effort SAME-POD fast-path: wake a handler blocked on THIS pod immediately. It re-reads the now-committed - // Expired terminal on its own connection and replays it. Cross-pod it harmlessly returns false — and that's - // fine: the durable Expired row + the blocked call's bounded-elapse → pending-ticket → a re-call that replays - // the Expired terminal IS the cross-pod guarantee. No wake, no decision lost. - _waiters.TrySignal(row.LedgerId, ToolApprovalOutcome.Expired); + WakeQuietly(row); - // Best-effort + idempotent: mirror the approval card to timed-out (no-ops if a human already resolved it, or if - // no card was ever posted). The ledger row is the authority; this is its display mirror. - if (row.ApprovalMessageId is { } messageId) - await _interactions.MarkTimedOutAsync(messageId, "expired", cancellationToken).ConfigureAwait(false); + await MirrorQuietlyAsync(row, cancellationToken).ConfigureAwait(false); + + _logger.LogInformation("Tool call approval expired. LedgerId={LedgerId} TeamId={TeamId}", row.LedgerId, row.TeamId); + } + + // Best-effort SAME-POD fast-path: wake a handler blocked on THIS pod immediately. It re-reads the now-committed + // Expired terminal on its own connection and replays it. Cross-pod it harmlessly returns false — and that's + // fine: the durable Expired row + the blocked call's bounded-elapse → pending-ticket → a re-call that replays + // the Expired terminal IS the cross-pod guarantee. No wake, no decision lost. A wake that throws costs its card + // nothing: the mirror is its own step. + private void WakeQuietly(ExpiredToolApproval row) + { + try + { + _waiters.TrySignal(row.LedgerId, ToolApprovalOutcome.Expired); + } + catch (Exception exception) + { + _logger.LogWarning(exception, "Tool call approval {LedgerId} expired, but waking its call failed; the call reads the Expired terminal when it next looks", row.LedgerId); + } + } - _logger.LogInformation("Tool call approval expired and mirrored. LedgerId={LedgerId} TeamId={TeamId}", row.LedgerId, row.TeamId); + // Best-effort + idempotent: mirror the approval card to timed-out (no-ops if a human already resolved it, or if no + // card was ever posted). The ledger row is the authority; this is its display mirror. A failed mirror leaves nothing + // tracked on the context (the interaction service forgets the card whose save failed), so the next row's mirror is + // not the one that writes it. + private async Task MirrorQuietlyAsync(ExpiredToolApproval row, CancellationToken cancellationToken) + { + if (row.ApprovalMessageId is not { } messageId) return; + + try + { + await _interactions.MarkTimedOutAsync(messageId, "expired", cancellationToken).ConfigureAwait(false); + } + catch (Exception exception) when (exception is not OperationCanceledException || !cancellationToken.IsCancellationRequested) + { + _logger.LogWarning(exception, "Tool call approval {LedgerId} expired, but mirroring its card failed; the ledger row stands as the verdict", row.LedgerId); + } } } diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Reconciliation/StuckRunReconcilerService.cs b/backend/src/CodeSpace.Core/Services/Workflows/Reconciliation/StuckRunReconcilerService.cs index 447ddf715..6f4a5efa2 100644 --- a/backend/src/CodeSpace.Core/Services/Workflows/Reconciliation/StuckRunReconcilerService.cs +++ b/backend/src/CodeSpace.Core/Services/Workflows/Reconciliation/StuckRunReconcilerService.cs @@ -308,21 +308,14 @@ select run.Id /// concurrent resume that already drove the run (its flip won) leaves us with 0 rows → we skip, /// no double-dispatch. The resolved waits rehydrate as the suspended nodes' ResumePayloads on the /// re-walk, so the run continues from where it stranded. + /// Never a sub-workflow child whose parent has finished — one whose cancel failed during its parent's stop and whose + /// waits were then closed, say. Revived, it would run under the finished parent toward a wait that parent no longer holds + /// open; no one is waiting for it. It is left as it is and counted in the log, so it does not linger unseen. The guard is + /// by source type, as the Pending sweep's is, because ParentRunId is also a rerun's lineage, finished by definition. /// private async Task RedispatchStrandedSuspendedAsync(CancellationToken cancellationToken) { - var threshold = DateTimeOffset.UtcNow - SuspendedStrandedAfter; - - var strandedIds = await _db.WorkflowRun.AsNoTracking() - .Where(r => r.Status == WorkflowRunStatus.Suspended - && r.LastModifiedDate < threshold - && r.CompletionParkedAt == null - && !_db.WorkflowRunWait.Any(w => w.RunId == r.Id && w.Status == WorkflowWaitStatuses.Pending)) - .OrderBy(r => r.LastModifiedDate) - .Take(BatchSize) - .Select(r => r.Id) - .ToListAsync(cancellationToken) - .ConfigureAwait(false); + var strandedIds = await FindStrandedSuspendedIdsAsync(cancellationToken).ConfigureAwait(false); var redispatched = 0; foreach (var runId in strandedIds) @@ -352,6 +345,32 @@ private async Task RedispatchStrandedSuspendedAsync(CancellationToken cance return redispatched; } + /// + /// The stranded Suspended runs to revive this tick, oldest first, at most — leaving out each + /// sub-workflow child whose parent has finished, and saying in the log how many it left. The exclusion is in the query, + /// not the loop: a child left alone stays Suspended, so one that took a place in the batch would take it every tick and + /// crowd out the runs that are wanted. + /// + private async Task> FindStrandedSuspendedIdsAsync(CancellationToken cancellationToken) + { + var threshold = DateTimeOffset.UtcNow - SuspendedStrandedAfter; + + var stranded = _db.WorkflowRun.AsNoTracking() + .Where(r => r.Status == WorkflowRunStatus.Suspended + && r.LastModifiedDate < threshold + && r.CompletionParkedAt == null + && !_db.WorkflowRunWait.Any(w => w.RunId == r.Id && w.Status == WorkflowWaitStatuses.Pending)) + .Select(r => new { r.Id, r.LastModifiedDate, ParentFinished = r.SourceType == WorkflowRunSourceTypes.ChildWorkflow && _db.WorkflowRun.Any(p => p.Id == r.ParentRunId && TerminalRunStatuses.Contains(p.Status)) }); + + var strandedIds = await stranded.Where(r => !r.ParentFinished).OrderBy(r => r.LastModifiedDate).Take(BatchSize).Select(r => r.Id).ToListAsync(cancellationToken).ConfigureAwait(false); + + var leftAlone = await stranded.CountAsync(r => r.ParentFinished, cancellationToken).ConfigureAwait(false); + + if (leftAlone > 0) _logger.LogInformation("StuckRunReconciler: left {Count} stranded Suspended sub-workflow(s) alone because their parent has finished", leftAlone); + + return strandedIds; + } + /// /// Free active-rerun leases whose fork reached a TERMINAL state — the complete backstop for the engine's inline /// release (which the cancel paths skip, and which is best-effort). A crashed fork is first flipped to Failure by diff --git a/backend/tests/CodeSpace.IntegrationTests/Agents/ToolApprovalExpiryServiceTests.cs b/backend/tests/CodeSpace.IntegrationTests/Agents/ToolApprovalExpiryServiceTests.cs index f08642e96..ed7a1487d 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Agents/ToolApprovalExpiryServiceTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Agents/ToolApprovalExpiryServiceTests.cs @@ -1,3 +1,5 @@ +using System.Collections.Concurrent; +using System.Data.Common; using System.Text.Json; using Autofac; using CodeSpace.Core.Persistence.Db; @@ -13,6 +15,8 @@ using CodeSpace.Messages.Enums; using MediatR; using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Diagnostics; +using Microsoft.Extensions.Logging; using Npgsql; using Shouldly; @@ -117,6 +121,58 @@ public async Task Mediator_dispatch_commits_the_expired_CAS_BEFORE_it_wakes_a_sa customMessage: "the durable CAS must be COMMITTED before the wake — a fresh-connection read at signal time saw the row as Expired, not the pre-commit AwaitingApproval. If this fails, the signal is firing inside the command transaction (see ExpireDueAsync deferring via IPostCommitActions)."); } + [Theory] + [InlineData(false)] // the service called with no transaction open: each row's follow-ups run inline, as the sweep reaches it + [InlineData(true)] // the recurring job's command: one transaction, the follow-ups drained once it commits + public async Task A_card_mirror_that_fails_to_save_leaves_the_other_expired_approvals_woken_and_mirrored(bool throughTheCommand) + { + // The sweep expires its whole batch, then wakes each row's blocked call and mirrors its card. With no drain to + // swallow it, one card's failed save used to end those follow-ups right there: every later row unwoken, its card + // Open, and no tick to come back for it, since the sweep selects only rows still awaiting approval. The drain + // contains that per row but cannot name the row. And a card whose save failed must not stay tracked: the next + // card's save would write it as if it had landed. + var (teamId, channelId) = await SeedTeamChannelAsync(); + var first = await ParkOverdueApprovalAsync(teamId, channelId, LongAgo); + var faulted = await ParkOverdueApprovalAsync(teamId, channelId, LongAgo.AddSeconds(1)); + var last = await ParkOverdueApprovalAsync(teamId, channelId, LongAgo.AddSeconds(2)); + + var log = new RecordingLogger(); + using var scope = BeginScopeFailingCardSaveOf(faulted.MessageId, log); + var waiters = scope.Resolve(); + var firstCall = waiters.Register(first.LedgerId); + var faultedCall = waiters.Register(faulted.LedgerId); + var lastCall = waiters.Register(last.LedgerId); + + try + { + var expired = 0; + (await Record.ExceptionAsync(async () => expired = await SweepAsync(scope, throughTheCommand))).ShouldBeNull("one card's failed mirror is that card's, not the sweep's"); + + // >= not == : the tally is deployment-wide; the rows this test owns are the proof. + expired.ShouldBeGreaterThanOrEqualTo(3, "all three approvals were durably expired, whatever their follow-ups did"); + + foreach (var approval in new[] { first, faulted, last }) + (await ReadRowAsync(approval.LedgerId)).Status.ShouldBe(ToolCallLedgerStatus.Expired, "the ledger row is the authority, and it stands"); + + (await WokenAsync(firstCall, "the first approval's blocked call")).ShouldBe(ToolApprovalOutcome.Expired); + (await WokenAsync(faultedCall, "the blocked call of the approval whose card failed to mirror")).ShouldBe(ToolApprovalOutcome.Expired, customMessage: "the wake is its own step: a failed mirror costs the call nothing"); + (await WokenAsync(lastCall, "the last approval's blocked call, behind the failed card")).ShouldBe(ToolApprovalOutcome.Expired, customMessage: "check ToolApprovalExpiryService reached the rows after the failed mirror"); + + (await ReadInteractionStateAsync(first.MessageId)).ShouldBe(InteractionState.Resolved, "the card before the failed one is mirrored"); + (await ReadInteractionStateAsync(last.MessageId)).ShouldBe(InteractionState.Resolved, "and the card after it is mirrored too"); + (await ReadInteractionStateAsync(faulted.MessageId)).ShouldBe(InteractionState.Open, "only the failed card's display mirror was lost — a late click on it is refused, its row being Expired — and no later save wrote it"); + + log.Entries.Where(e => e.Level == LogLevel.Warning).Select(e => e.Message).ShouldContain(m => m.Contains(faulted.LedgerId.ToString()), "the failed mirror is logged against its own ledger row"); + scope.Resolve().ChangeTracker.HasChanges().ShouldBeFalse("a card whose save failed is forgotten, not left tracked for the next save"); + } + finally + { + waiters.Remove(first.LedgerId); + waiters.Remove(faulted.LedgerId); + waiters.Remove(last.LedgerId); + } + } + // ─── Park a real approval card through the real handler (times out fast → row stays AwaitingApproval, then back-date the deadline) ─── private async Task<(Guid LedgerId, Guid MessageId)> ParkApprovalAsync(Guid teamId, Guid runId, Guid channelId) @@ -157,6 +213,112 @@ private McpRequestHandler Handler(ILifetimeScope scope, Guid teamId, Guid runId, scope.Resolve(), 0, governanceEnabled: true, approvalConversationId: channelId, scope.Resolve(), scope.Resolve(), scope.Resolve()); + /// Overdue ahead of anything else overdue in this shared database, so the sweep reaches these rows first. + private static readonly DateTimeOffset LongAgo = new(2000, 1, 1, 0, 0, 0, TimeSpan.Zero); + + /// An undecided approval exactly as a parked call leaves it — claimed, parked with its token and deadline, its card posted and recorded on the row — without the handler's bounded wait. + private async Task<(Guid LedgerId, Guid MessageId)> ParkOverdueApprovalAsync(Guid teamId, Guid channelId, DateTimeOffset deadlineAt) + { + using var scope = _fixture.BeginScope(); + var ledger = scope.Resolve(); + var token = Guid.NewGuid().ToString("N"); + + var claim = await ledger.TryClaimAsync(Guid.NewGuid(), teamId, "git.open_pr", Guid.NewGuid().ToString("N"), "input-hash", 0, CancellationToken.None); + (await ledger.TryBeginApprovalAsync(claim.LedgerId, teamId, token, deadlineAt, CancellationToken.None)).ShouldBeTrue("fixture check: the claimed row parks for approval"); + + var card = new MessageInteraction + { + Component = new ActionButtonsComponent + { + Buttons = new List + { + new() { Key = "approve", Label = "Approve", Style = InteractionButtonStyle.Primary }, // McpRequestHandler.ApprovalButtonsConfig + new() { Key = "reject", Label = "Reject", Style = InteractionButtonStyle.Danger, RequiresComment = true }, + }, + }, + Target = new ToolCallApprovalTarget { Token = token }, + AllowedResponderUserIds = null, + Resolve = new ResolvePolicy(), + }; + + var posted = await scope.Resolve().PostAsBotAsync(channelId, "Agent run requests approval to run **git.open_pr**. Approve to let it proceed, or reject to refuse it.", card, CancellationToken.None); + await ledger.SetApprovalMessageAsync(claim.LedgerId, teamId, posted.Id, CancellationToken.None); + + return (claim.LedgerId, posted.Id); + } + + /// A scope whose database commands pass the interceptor that fails one card's mirror, and whose expiry service logs to . + private ILifetimeScope BeginScopeFailingCardSaveOf(Guid messageId, ILogger log) + { + DbContextOptions production; + using (var probe = _fixture.BeginScope()) + production = probe.Resolve>(); + + var options = new DbContextOptionsBuilder(production).AddInterceptors(new FailCardSaveOf(messageId)).Options; + + return _fixture.BeginScope(b => + { + b.RegisterInstance(options).As>().SingleInstance(); + b.RegisterInstance>(log); + }); + } + + /// One sweep: the service called straight, with no transaction open, or the command the recurring job sends, through the mediator's own transaction and post-commit drain. + private static async Task SweepAsync(ILifetimeScope scope, bool throughTheCommand) => + throughTheCommand + ? (await scope.Resolve().Send(new ExpireStaleToolApprovalsCommand(), CancellationToken.None)).Expired + : await scope.Resolve().ExpireDueAsync(DateTimeOffset.UtcNow, CancellationToken.None); + + /// The outcome was woken with; fails by its name when it is not woken within five seconds. + private static async Task WokenAsync(IToolApprovalWaiter call, string signal) + { + var first = await Task.WhenAny(call.Completion, Task.Delay(TimeSpan.FromSeconds(5))); + + (first == call.Completion).ShouldBeTrue($"{signal} was not woken within 5s — check ToolApprovalExpiryService.ResolveAsync reached it"); + + return await call.Completion; + } + + /// Fails the save of one card's timed-out mirror, once. + private sealed class FailCardSaveOf : DbCommandInterceptor + { + private readonly Guid _messageId; + private int _fired; + + public FailCardSaveOf(Guid messageId) { _messageId = messageId; } + + public override ValueTask> NonQueryExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) + { + Fail(command); + return ValueTask.FromResult(result); + } + + public override ValueTask> ReaderExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) + { + Fail(command); + return ValueTask.FromResult(result); + } + + private void Fail(DbCommand command) + { + if (command.CommandText.Contains("UPDATE message", StringComparison.Ordinal) && Carries(command, _messageId) && Interlocked.Exchange(ref _fired, 1) == 0) + throw new TimeoutException($"saving card {_messageId} timed out"); + } + + private static bool Carries(DbCommand command, Guid value) => command.Parameters.Cast().Any(parameter => parameter.Value is Guid id && id == value); + } + + private sealed class RecordingLogger : ILogger + { + public ConcurrentQueue<(LogLevel Level, string Message, Exception? Exception)> Entries { get; } = new(); + + public IDisposable? BeginScope(TState state) where TState : notnull => null; + + public bool IsEnabled(LogLevel logLevel) => true; + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) => Entries.Enqueue((logLevel, formatter(state, exception), exception)); + } + // ─── Reads ─────────────────────────────────────────────────────────────────── private async Task ReadRowAsync(Guid ledgerId) diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/StuckRunReconcilerFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/StuckRunReconcilerFlowTests.cs index b3ea7b326..1a412e4bd 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/StuckRunReconcilerFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/StuckRunReconcilerFlowTests.cs @@ -11,6 +11,7 @@ using CodeSpace.Messages.Enums; using MediatR; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; using Shouldly; namespace CodeSpace.IntegrationTests.Workflows; @@ -625,6 +626,82 @@ public async Task Suspended_with_zero_pending_waits_but_within_grace_window_is_N "window must protect it so we don't race the concurrent Suspended→Pending flip"); } + [Theory] + [InlineData(WorkflowRunStatus.Cancelled)] + [InlineData(WorkflowRunStatus.Failure)] + [InlineData(WorkflowRunStatus.Success)] + public async Task A_stranded_suspended_sub_workflow_under_a_finished_parent_is_never_redispatched(WorkflowRunStatus finished) + { + // A child whose cancel failed during its parent's stop, and whose waits were then closed, sits Suspended with no + // pending wait — exactly what this sweep revives. Revived, it ran under a finished parent toward a wait that parent + // no longer holds open, spending on work no one is waiting for. The guard the Pending sweep has: a finished parent's + // child is not work anyone wants done. + var (teamId, userId) = await WorkflowsTestSeed.SeedTeamAsync(_fixture); + var workflowId = await CreateWorkflowAsync(teamId, userId); + var parentId = await StageStuckRunAsync(workflowId, teamId, finished, createdAgo: TimeSpan.FromHours(1)); + var childId = await StageStrandedSuspendedUnderAsync(workflowId, teamId, parentId, WorkflowRunSourceTypes.ChildWorkflow); + + var versionBefore = await ReadRowVersionAsync(childId); + + await ReconcileAsync(); + + (await ReadRowVersionAsync(childId)).ShouldBe(versionBefore, $"a stranded sub-workflow under a {finished} parent must not be revived — nor touched at all"); + (await ReadStatusAsync(childId)).ShouldBe(WorkflowRunStatus.Suspended); + } + + [Theory] + [InlineData(WorkflowRunStatus.Running)] // the parent walking on while its child is stranded + [InlineData(WorkflowRunStatus.Suspended)] // the parent parked on its child's Subworkflow wait — how production holds a live parent + public async Task A_stranded_suspended_sub_workflow_under_a_live_parent_is_still_redispatched(WorkflowRunStatus live) + { + // The guard is on the parent having finished. A child stranded under a parent that still lives is what this sweep + // exists to revive — that parent is waiting on it. + var (teamId, userId) = await WorkflowsTestSeed.SeedTeamAsync(_fixture); + var workflowId = await CreateWorkflowAsync(teamId, userId); + var parentId = await StageStuckRunAsync(workflowId, teamId, live, createdAgo: TimeSpan.FromMinutes(1), startedAtAgo: TimeSpan.FromMinutes(1)); + var childId = await StageStrandedSuspendedUnderAsync(workflowId, teamId, parentId, WorkflowRunSourceTypes.ChildWorkflow); + + if (live == WorkflowRunStatus.Suspended) await SeedSubworkflowWaitAsync(parentId, childId); + + await ReconcileAsync(); + + (await ReadStatusAsync(childId)).ShouldBe(WorkflowRunStatus.Enqueued, $"a stranded sub-workflow under a {live} parent is revived like any stranded run"); + } + + [Fact] + public async Task A_stranded_suspended_rerun_of_a_finished_run_is_still_redispatched() + { + // The other meaning of ParentRunId: a rerun's parent is its lineage, which has finished by definition. The + // finished-parent guard is for sub-workflow children only — a rerun it swallowed would never be revived. + var (teamId, userId) = await WorkflowsTestSeed.SeedTeamAsync(_fixture); + var workflowId = await CreateWorkflowAsync(teamId, userId); + var originalId = await StageStuckRunAsync(workflowId, teamId, WorkflowRunStatus.Failure, createdAgo: TimeSpan.FromHours(1)); + var rerunId = await StageStrandedSuspendedUnderAsync(workflowId, teamId, originalId, WorkflowRunSourceTypes.Rerun); + + await ReconcileAsync(); + + (await ReadStatusAsync(rerunId)).ShouldBe(WorkflowRunStatus.Enqueued, "a stranded rerun of a finished run is revived like any stranded run"); + } + + [Fact] + public async Task The_sweep_says_how_many_stranded_sub_workflows_it_left_alone() + { + // The children it leaves Suspended stay Suspended, so without a line they would linger unseen. + var (teamId, userId) = await WorkflowsTestSeed.SeedTeamAsync(_fixture); + var workflowId = await CreateWorkflowAsync(teamId, userId); + var parentId = await StageStuckRunAsync(workflowId, teamId, WorkflowRunStatus.Cancelled, createdAgo: TimeSpan.FromHours(1)); + await StageStrandedSuspendedUnderAsync(workflowId, teamId, parentId, WorkflowRunSourceTypes.ChildWorkflow); + + var log = new RecordingLogger(); + + await ReconcileAsync(log); + + var entry = log.Entries.Where(e => e.Level == LogLevel.Information && e.Message.Contains("parent has finished")).ShouldHaveSingleItem("one line per sweep says so"); + + // >= not == : the count is deployment-wide (see the class note); the child this test staged is one of them. + ((int)entry.Properties["Count"]!).ShouldBeGreaterThanOrEqualTo(1, "the line carries the count as a structured property, not interpolated text"); + } + // ─── Helpers ────────────────────────────────────────────────────────────────── private async Task RunEngineAsync(Guid runId) @@ -895,9 +972,9 @@ private async Task SeedWaitAsync(Guid runId, string nodeId, string status) WaitKind = WorkflowWaitKinds.Approval, Token = Guid.NewGuid().ToString("N"), Status = status, - PayloadJson = status == WorkflowWaitStatuses.Resolved ? "{}" : null, + PayloadJson = status != WorkflowWaitStatuses.Pending ? "{}" : null, CreatedAt = DateTimeOffset.UtcNow, - ResolvedAt = status == WorkflowWaitStatuses.Resolved ? DateTimeOffset.UtcNow : null, + ResolvedAt = status != WorkflowWaitStatuses.Pending ? DateTimeOffset.UtcNow : null, }); await db.SaveChangesAsync(); @@ -1130,12 +1207,29 @@ private async Task StageStuckPendingUnderAsync(Guid workflowId, Guid teamI { var runId = await StageStuckRunAsync(workflowId, teamId, WorkflowRunStatus.Pending, createdAgo: StuckRunReconcilerService.PendingStuckAfter + TimeSpan.FromMinutes(1)); - using var scope = _fixture.BeginScope(); - await scope.Resolve().WorkflowRun.Where(r => r.Id == runId).ExecuteUpdateAsync(s => s.SetProperty(r => r.ParentRunId, parentRunId).SetProperty(r => r.SourceType, sourceType)); + await LinkUnderAsync(runId, parentRunId, sourceType); + + return runId; + } + + /// A Suspended run stranded past the grace window — no pending wait, the one wait it had closed Discarded, as a stop's teardown closes it — started under as : a sub-workflow child, or a rerun whose parent is only its lineage. + private async Task StageStrandedSuspendedUnderAsync(Guid workflowId, Guid teamId, Guid parentRunId, string sourceType) + { + var runId = await StageStuckRunAsync(workflowId, teamId, WorkflowRunStatus.Suspended, createdAgo: StuckRunReconcilerService.SuspendedStrandedAfter + TimeSpan.FromMinutes(5), backdateLastModified: true); + await SeedWaitAsync(runId, "start", WorkflowWaitStatuses.Discarded); + + await LinkUnderAsync(runId, parentRunId, sourceType); return runId; } + /// Point the run at its parent and say what kind of child it is. Written with ExecuteUpdate, which leaves the backdated timestamps the staging set alone. + private async Task LinkUnderAsync(Guid runId, Guid parentRunId, string sourceType) + { + using var scope = _fixture.BeginScope(); + await scope.Resolve().WorkflowRun.Where(r => r.Id == runId).ExecuteUpdateAsync(s => s.SetProperty(r => r.ParentRunId, parentRunId).SetProperty(r => r.SourceType, sourceType)); + } + private async Task SeedLedgerRecordAsync(Guid runId, string recordType, DateTimeOffset occurredAt) { using var scope = _fixture.BeginScope(); @@ -1164,6 +1258,13 @@ private async Task ReconcileAsync() return await mediator.Send(new ReconcileStuckRunsCommand()); } + /// As , in a scope whose reconciler logs to . + private async Task ReconcileAsync(ILogger log) + { + using var scope = _fixture.BeginScope(b => b.RegisterInstance(log).As>()); + return await scope.Resolve().Send(new ReconcileStuckRunsCommand()); + } + /// The Postgres row version. Changes on any write, so it answers "was this row touched" rather than "what does it say now". private async Task ReadRowVersionAsync(Guid runId) { @@ -1198,4 +1299,17 @@ private async Task ContinueAsync(Guid runId, Guid teamId) using var scope = _fixture.BeginScope(); return await scope.Resolve().ContinueRunAsync(runId, teamId, CancellationToken.None); } + + /// Keeps each entry's level, formatted message and structured properties — the named template values an interpolated message would not carry. + private sealed class RecordingLogger : ILogger + { + public List<(LogLevel Level, string Message, IReadOnlyDictionary Properties)> Entries { get; } = new(); + + public IDisposable? BeginScope(TState state) where TState : notnull => null; + + public bool IsEnabled(LogLevel logLevel) => true; + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) => + Entries.Add((logLevel, formatter(state, exception), (state as IEnumerable> ?? []).ToDictionary(p => p.Key, p => p.Value))); + } } diff --git a/backend/tests/CodeSpace.UnitTests/Agents/ToolApprovalExpiryServiceTests.cs b/backend/tests/CodeSpace.UnitTests/Agents/ToolApprovalExpiryServiceTests.cs new file mode 100644 index 000000000..7e1fcca34 --- /dev/null +++ b/backend/tests/CodeSpace.UnitTests/Agents/ToolApprovalExpiryServiceTests.cs @@ -0,0 +1,197 @@ +using CodeSpace.Core.Middlewares.Transactional; +using CodeSpace.Core.Persistence.Entities; +using CodeSpace.Core.Services.Agents.Mcp; +using CodeSpace.Core.Services.Chat; +using CodeSpace.Messages.Agents; +using CodeSpace.Messages.Decisions; +using Microsoft.Extensions.Logging; +using Shouldly; + +namespace CodeSpace.UnitTests.Agents; + +/// +/// 🟢 Unit: the approval sweep's follow-ups — the in-process wake and the card mirror — stay best-effort per row and per +/// step. A caller with no transaction runs them inline as the sweep reaches each row, and nothing there swallows a +/// failure: one card's failed mirror must still leave that row woken and every later row woken and mirrored, a failed wake +/// must still leave its own card mirrored, and the count of approvals durably expired stands. Each swallowed failure names +/// its ledger row in the log, since the drain that would otherwise catch it cannot say which row it was. +/// +[Trait("Category", "Unit")] +public sealed class ToolApprovalExpiryServiceTests +{ + [Fact] + public async Task A_failed_card_mirror_leaves_every_other_expired_approval_woken_and_mirrored() + { + var first = Expired(); + var second = Expired(); + var third = Expired(); + var waiters = new RecordingWaiters(); + var cards = new CardsFailingFor(second.ApprovalMessageId!.Value); + var log = new CapturingLogger(); + var service = new ToolApprovalExpiryService(new ExpiredLedger(first, second, third), waiters, cards, new InlinePostCommitActions(), log); + + (await service.ExpireDueAsync(DateTimeOffset.UtcNow, CancellationToken.None)).ShouldBe(3, "all three approvals were durably expired, whatever their follow-ups did"); + + waiters.Signalled.ShouldBe(new[] { (first.LedgerId, ToolApprovalOutcome.Expired), (second.LedgerId, ToolApprovalOutcome.Expired), (third.LedgerId, ToolApprovalOutcome.Expired) }, ignoreOrder: false, customMessage: "each approval's blocked call is woken with Expired, the row whose mirror failed included"); + cards.Mirrored.ShouldBe(new[] { first.ApprovalMessageId!.Value, third.ApprovalMessageId!.Value }, ignoreOrder: false, customMessage: "the second card's mirror failed, and the third's still ran"); + + var warning = log.Entries.ShouldHaveSingleItem("the one failed step is logged once"); + warning.Level.ShouldBe(LogLevel.Warning); + warning.Exception.ShouldBeOfType("the failure itself travels with the entry"); + warning.Properties["LedgerId"].ShouldBe(second.LedgerId, "the entry names the row whose card was not mirrored"); + } + + [Fact] + public async Task A_failed_wake_still_mirrors_its_card() + { + var approval = Expired(); + var cards = new CardsFailingFor(Guid.Empty); + var log = new CapturingLogger(); + var service = new ToolApprovalExpiryService(new ExpiredLedger(approval), new WaitersFailingFor(approval.LedgerId), cards, new InlinePostCommitActions(), log); + + (await service.ExpireDueAsync(DateTimeOffset.UtcNow, CancellationToken.None)).ShouldBe(1, "the approval was durably expired, whatever its follow-ups did"); + + cards.Mirrored.ShouldBe(new[] { approval.ApprovalMessageId!.Value }, ignoreOrder: false, customMessage: "the wake and the mirror are separate steps: a wake that threw costs the card nothing"); + log.Entries.ShouldHaveSingleItem("the one failed step is logged once").Properties["LedgerId"].ShouldBe(approval.LedgerId, "the entry names the row whose call was not woken"); + } + + [Fact] + public async Task An_approval_that_never_posted_a_card_is_woken_and_has_nothing_to_mirror() + { + var approval = Expired(withCard: false); + var waiters = new RecordingWaiters(); + var cards = new CardsFailingFor(Guid.Empty); + var service = new ToolApprovalExpiryService(new ExpiredLedger(approval), waiters, cards, new InlinePostCommitActions(), new CapturingLogger()); + + (await service.ExpireDueAsync(DateTimeOffset.UtcNow, CancellationToken.None)).ShouldBe(1); + + waiters.Signalled.ShouldBe(new[] { (approval.LedgerId, ToolApprovalOutcome.Expired) }, ignoreOrder: false, customMessage: "the blocked call is woken whether or not a card was ever posted"); + cards.Mirrored.ShouldBeEmpty("no card was recorded on the row, so there is nothing to mirror"); + } + + [Fact] + public async Task A_shutdown_that_cancels_a_mirror_stops_the_sweep_instead_of_being_logged_away() + { + using var shutdown = new CancellationTokenSource(); + var first = Expired(); + var second = Expired(); + var waiters = new RecordingWaiters(); + var log = new CapturingLogger(); + var service = new ToolApprovalExpiryService(new ExpiredLedger(first, second), waiters, new CardsCancelledBy(shutdown), new InlinePostCommitActions(), log); + + await Should.ThrowAsync(() => service.ExpireDueAsync(DateTimeOffset.UtcNow, shutdown.Token)); + + waiters.Signalled.ShouldBe(new[] { (first.LedgerId, ToolApprovalOutcome.Expired) }, ignoreOrder: false, customMessage: "the sweep stopped at the row the shutdown reached; the rows it never got to are not swept on a token that is cancelled"); + log.Entries.ShouldBeEmpty("a stop is not a failed mirror, so nothing is logged as one"); + } + + private static ExpiredToolApproval Expired(bool withCard = true) => new() { LedgerId = Guid.NewGuid(), TeamId = Guid.NewGuid(), ApprovalMessageId = withCard ? Guid.NewGuid() : null }; + + /// No transaction is open, as when the service is called outside its command: an action runs the moment it is handed over. + private sealed class InlinePostCommitActions : IPostCommitActions + { + public Task RunAfterCommitAsync(Func action, CancellationToken cancellationToken) => action(cancellationToken); + public Task RunAllAsync(CancellationToken cancellationToken) => Task.CompletedTask; + public int CreateCheckpoint() => 0; + public void RollbackTo(int checkpoint) { } + } + + private sealed class RecordingWaiters : IToolApprovalWaiterRegistry + { + public List<(Guid LedgerId, ToolApprovalOutcome Outcome)> Signalled { get; } = new(); + + public IToolApprovalWaiter Register(Guid ledgerId) => throw new NotSupportedException(); + public bool TrySignal(Guid ledgerId, ToolApprovalOutcome outcome) + { + Signalled.Add((ledgerId, outcome)); + return true; + } + public void Remove(Guid ledgerId) => throw new NotSupportedException(); + } + + private sealed class WaitersFailingFor : IToolApprovalWaiterRegistry + { + private readonly Guid _failing; + + public WaitersFailingFor(Guid failing) { _failing = failing; } + + public IToolApprovalWaiter Register(Guid ledgerId) => throw new NotSupportedException(); + public bool TrySignal(Guid ledgerId, ToolApprovalOutcome outcome) => ledgerId == _failing ? throw new InvalidOperationException("the waiter could not be signalled") : true; + public void Remove(Guid ledgerId) => throw new NotSupportedException(); + } + + private sealed class CardsFailingFor : IMessageInteractionService + { + private readonly Guid _failing; + + public CardsFailingFor(Guid failing) { _failing = failing; } + + public List Mirrored { get; } = new(); + + public Task RespondAsync(Guid teamId, Guid messageId, string responseKey, Guid actorUserId, string? comment, IReadOnlyDictionary? values, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task MarkTimedOutAsync(Guid messageId, string responseKey, CancellationToken cancellationToken) + { + if (messageId == _failing) throw new InvalidOperationException("the card could not be mirrored"); + + Mirrored.Add(messageId); + return Task.CompletedTask; + } + } + + /// A host shutting down under the sweep: the first card mirror it reaches cancels the token the sweep runs on, and is itself cancelled by it. + private sealed class CardsCancelledBy : IMessageInteractionService + { + private readonly CancellationTokenSource _shutdown; + + public CardsCancelledBy(CancellationTokenSource shutdown) { _shutdown = shutdown; } + + public Task RespondAsync(Guid teamId, Guid messageId, string responseKey, Guid actorUserId, string? comment, IReadOnlyDictionary? values, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task MarkTimedOutAsync(Guid messageId, string responseKey, CancellationToken cancellationToken) + { + _shutdown.Cancel(); + + throw new OperationCanceledException(_shutdown.Token); + } + } + + /// Captures each warning's level, exception and structured properties — the named template values an interpolated message would not carry. The per-row Information line is not what these tests are about, so nothing below Warning is kept. + private sealed class CapturingLogger : ILogger + { + public List<(LogLevel Level, Exception? Exception, IReadOnlyDictionary Properties)> Entries { get; } = new(); + + public IDisposable? BeginScope(TState state) where TState : notnull => null; + public bool IsEnabled(LogLevel logLevel) => logLevel >= LogLevel.Warning; + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) + { + if (logLevel < LogLevel.Warning) return; + + Entries.Add((logLevel, exception, (state as IEnumerable> ?? []).ToDictionary(p => p.Key, p => p.Value))); + } + } + + /// A ledger whose sweep expired exactly the approvals handed in; every other ledger method is unreachable here. + private sealed class ExpiredLedger : IToolCallLedgerService + { + private readonly IReadOnlyList _expired; + + public ExpiredLedger(params ExpiredToolApproval[] expired) { _expired = expired; } + + public Task> ExpireStaleApprovalsAsync(DateTimeOffset now, CancellationToken cancellationToken) => Task.FromResult(_expired); + + public IAsyncEnumerable ExpireStaleDecisionsAsync(DateTimeOffset now, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task ExpireStaleToolCallsAsync(DateTimeOffset now, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task TryClaimAsync(Guid agentRunId, Guid teamId, string toolKind, string idempotencyKey, string inputHash, long fenceEpoch, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task RecordTerminalAsync(Guid ledgerId, Guid teamId, ToolCallLedgerStatus status, string? resultJson, string? error, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task TryBeginApprovalAsync(Guid ledgerId, Guid teamId, string approvalToken, DateTimeOffset deadlineAt, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task SetApprovalMessageAsync(Guid ledgerId, Guid teamId, Guid messageId, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task TryBeginExecutionAsync(Guid ledgerId, Guid teamId, long fenceEpoch, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task ReadApprovalStateAsync(Guid ledgerId, Guid teamId, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task ReadTerminalForReplayAsync(Guid ledgerId, Guid agentRunId, Guid teamId, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task TryAnswerDecisionAsync(Guid ledgerId, Guid teamId, string answerJson, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task SetDecisionEnvelopeAsync(Guid ledgerId, Guid teamId, string envelopeJson, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task CountPendingDecisionsAsync(Guid agentRunId, Guid teamId, string excludeIdempotencyKey, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task FindBlockingDecisionIdAsync(Guid agentRunId, CancellationToken cancellationToken) => throw new NotSupportedException(); + public Task> GetForRunAsync(Guid agentRunId, Guid teamId, CancellationToken cancellationToken) => throw new NotSupportedException(); + } +}