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(); + } +}