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