Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,12 @@ namespace CodeSpace.Core.Services.Agents.Mcp;
/// connection — sees the COMMITTED <c>Expired</c> terminal and replays it (not the pre-commit <c>AwaitingApproval</c>,
/// 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.</para>
///
/// <para>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.</para>
/// </summary>
public interface IToolApprovalExpiryService
{
Expand Down Expand Up @@ -59,17 +65,45 @@ public async Task<int> 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);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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.</para>
/// <para>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 <c>ParentRunId</c> is also a rerun's lineage, finished by definition.</para>
/// </summary>
private async Task<int> 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)
Expand Down Expand Up @@ -352,6 +345,32 @@ private async Task<int> RedispatchStrandedSuspendedAsync(CancellationToken cance
return redispatched;
}

/// <summary>
/// The stranded Suspended runs to revive this tick, oldest first, at most <see cref="BatchSize"/> — 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.
/// </summary>
private async Task<List<Guid>> 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;
}

/// <summary>
/// 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
Expand Down
Loading
Loading