mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-02 02:07:25 +08:00
fix(recovery): escalate an issue whose interrupted run has spent its retry budget (#14046)
## Thinking Path
> - Paperclip is the open source app people use to manage AI agents for
work.
> - The stranded-issue sweeper is what puts an assigned issue back on a
live path when its run dies.
> - A failed or interrupted run gets a bounded transient retry; when
that budget is spent, the retry scheduler queues nothing and reports the
exhaustion.
> - The sweeper treated that "nothing queued" like every other "nothing
queued" and skipped the issue, on every tick, forever: `in_progress`, no
run, no path, no notice.
> - Three server restarts in a row (a deploy storm) are enough to spend
the budget, so this is reachable in ordinary operation.
> - This pull request makes the sweeper escalate a spent budget to a
board-owned recovery action, the same visible `blocked` state its other
dead ends use.
> - The benefit is that no issue can sit assigned and silent after its
retries run out.
## Linked Issues or Issue Description
No existing issue found. Related work: #14028 (this branch's
predecessor: queue drain-time wakes, retry interrupted corrective runs)
fixed two neighbouring gaps but not this one.
**What happened?**
A routine-created issue's run was interrupted by a graceful server
shutdown, and both bounded transient retries were interrupted by the
next two shutdowns. The retry scheduler logged "Bounded retry exhausted
after 2 scheduled attempts; no further automatic retry will be queued"
and stopped. On every tick after that, `reconcileStrandedAssignedIssues`
reached the generic continuation lane, `enqueueStrandedIssueRecovery`
delegated to `scheduleRecoveryRetry`, which returned null (exhausted),
and the sweeper counted the issue as `skipped`. The issue stayed
`in_progress` with no run, no recovery action, and no comment for hours
until a person woke the agent by hand. The same hole exists in the
assigned-`todo` dispatch lane.
**Expected behavior**
When the transient retry budget for an interrupted or failed run is
spent, the sweeper should treat it as the dead end it is: escalate the
issue to `blocked` with a board-owned recovery action and a notice,
exactly as it does when a continuation retry chain or an assignment
retry chain is exhausted. Leaving the budget spent across restarts is
correct and unchanged; leaving the issue silent is not.
**Steps to reproduce**
1. Assign an agent using a conversation adapter (its interrupted runs
carry `conversationContinuation`, so legacy reconciliation does not
terminalize them) an issue and let a run start.
2. Interrupt the run with a graceful shutdown, let the transient retry
start, interrupt it, and repeat once more so `scheduledRetryAttempt`
reaches 2.
3. Run the stranded-issue sweep. Before this change: `skipped` every
tick, issue `in_progress`, no run, no recovery action. After:
`escalated`, issue `blocked`, one active board-owned
`issue_recovery_actions` row, a "No live execution path" notice.
**Paperclip version or commit**
master at bd6caf51bb (2026-09-25).
**Deployment mode**
Managed cloud instance restarted by fleet deploys; the code path is the
same for any operator whose server restarts more often than the retry
budget allows.
## What Changed
- `recoveryService` gains an optional `transientRetryBudgetSpent(run)`
dependency; `heartbeatService` wires it as
`executionFailureRetryCount(run) >=
BOUNDED_TRANSIENT_HEARTBEAT_RETRY_MAX_ATTEMPTS`, the same check
`scheduleBoundedRetryForRun` applies.
- `enqueueStrandedIssueRecovery` takes an optional `outcome`
out-parameter and sets `retryExhausted` when the failed predecessor's
retry returned nothing **because** the budget is spent. A null return
without it still means another authority owns the run (native runtime,
legacy reconciliation) and the caller leaves it alone; the deliberate
"failure recovery cannot fall through into the continuation queue" rule
is unchanged.
- The generic `in_progress` continuation lane and the assigned-`todo`
dispatch lane escalate on `retryExhausted` via
`escalateStrandedAssignedIssue` with a "No live execution path" notice
(danger tone), which creates the board-owned source-scoped recovery
action and moves the issue to `blocked`. Every other
`enqueueStrandedIssueRecovery` caller is unchanged.
- Tests (`heartbeat-process-recovery.test.ts`): an `in_progress` issue
with a spent budget escalates (blocked, one active board-owned action,
notice, no successor run, idempotent on the next sweep); an interrupted
run with budget remaining still gets its transient retry; an assigned
`todo` issue with a spent dispatch budget escalates. The existing guard
"does not reset an exhausted incident budget on server restart" keeps
its no-successor-run assertion and now expects the board escalation
instead of nothing, with a comment on why.
## Verification
```
pnpm -r --filter './packages/**' build
cd server
npx vitest run src/__tests__/heartbeat-process-recovery.test.ts \
src/__tests__/issue-recovery-actions.test.ts \
src/services/recovery/successful-run-handoff.test.ts \
src/__tests__/heartbeat-task-drain-admission-release.test.ts \
src/__tests__/attention-service.test.ts \
src/__tests__/heartbeat-comment-wake-batching.test.ts
npx tsc --noEmit -p tsconfig.json
```
Live reproduction: the exact stuck state (issue `in_progress`, latest
run `interrupted` with `scheduledRetryAttempt` 2, the "Bounded retry
exhausted" lifecycle event, sweeper `skipped` every tick) was observed
on a managed instance running current master before this change was
written.
## Risks
- Behavior change is limited to runs whose transient budget is already
spent, which previously produced no action at all. Nothing new is
retried; the change only adds the escalation, so no retry loop can be
introduced.
- The board-owned action spawns no run. Resolving it (restore to the
owner) re-dispatches through the existing recovery-action routes, the
same flow as every other stranded escalation.
- Native-runtime and legacy-reconciliation predecessors are untouched:
they return before the exhaustion check.
## Model Used
Claude (Anthropic) — `claude-fable-5-1`, extended thinking, tool use
(Claude Code CLI). The incident diagnosis and the choice to escalate
rather than re-dispatch were steered by the maintainer.
## Checklist
- [x] I have included a thinking path that traces from project context
to this change
- [x] I have specified the model used (with version and capability
details)
- [x] I have checked ROADMAP.md and confirmed this PR does not duplicate
planned core work
- [x] I have searched GitHub for duplicate or related PRs and linked
them above
- [x] I have either (a) linked existing issues with `Fixes: #` / `Closes
#` / `Refs #` OR (b) described the issue in-PR following the relevant
issue template
- [x] I have not referenced internal/instance-local Paperclip issues or
links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip`
URLs)
- [x] My branch name describes the change (e.g. `docs/...`, `fix/...`)
and contains no internal Paperclip ticket id or instance-derived details
- [x] I have run tests locally and they pass
- [x] I have added or updated tests where applicable
- [x] I have updated relevant documentation to reflect my changes
- [x] I have considered and documented any risks above
- [ ] All Paperclip CI gates are green
- [ ] Greptile is 5/5 with no open P2s, recommendations, or follow-ups
- [x] I will address all Greptile and reviewer comments before
requesting merge
---------
Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
@@ -3745,21 +3745,40 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => {
|
||||
})
|
||||
.where(eq(heartbeatRuns.id, runId));
|
||||
await heartbeatService(db).drainRunningRunsForShutdown("SIGTERM");
|
||||
await heartbeatService(db).reconcileStrandedAssignedIssues();
|
||||
const reconciled = await heartbeatService(db).reconcileStrandedAssignedIssues();
|
||||
// The budget stays spent across the restart: no successor run is
|
||||
// queued, by the process-loss retry or by the sweeper.
|
||||
expect(
|
||||
await db
|
||||
.select()
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.agentId, agentId)),
|
||||
).toHaveLength(1);
|
||||
// A spent budget is no longer left to sit: the sweeper escalates the
|
||||
// issue to a board-owned recovery action instead of skipping it on
|
||||
// every tick with no live path. The action spawns no run of its own.
|
||||
expect(reconciled.escalated).toBe(1);
|
||||
const recoveryActions = await db
|
||||
.select()
|
||||
.from(issueRecoveryActions)
|
||||
.where(eq(issueRecoveryActions.sourceIssueId, issueId));
|
||||
expect(recoveryActions).toHaveLength(1);
|
||||
expect(recoveryActions[0]).toMatchObject({ status: "active", ownerType: "board", ownerAgentId: null });
|
||||
const escalatedIssue = await db
|
||||
.select({ status: issues.status })
|
||||
.from(issues)
|
||||
.where(eq(issues.id, issueId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(escalatedIssue?.status).toBe("blocked");
|
||||
// And the escalation is idempotent across further sweeps.
|
||||
const again = await heartbeatService(db).reconcileStrandedAssignedIssues();
|
||||
expect(again.escalated).toBe(0);
|
||||
expect(
|
||||
await db
|
||||
.select()
|
||||
.from(issueRecoveryActions)
|
||||
.where(eq(issueRecoveryActions.sourceIssueId, issueId)),
|
||||
).toEqual([]);
|
||||
const run = await heartbeatService(db).getRun(runId);
|
||||
expect(await getExecutionBlocker(db, run!.companyId, issueId)).toBeNull();
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.agentId, agentId)),
|
||||
).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("releases active environment leases when an orphaned run is reaped", async () => {
|
||||
@@ -6006,6 +6025,176 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => {
|
||||
expect(sourceIssue?.status).toBe("blocked");
|
||||
});
|
||||
|
||||
it("escalates an in_progress issue whose interrupted run has spent its transient retry budget", async () => {
|
||||
// Three deploy restarts in a row interrupted a run and both of its
|
||||
// bounded transient retries. The retry scheduler then reports the
|
||||
// budget exhausted and queues nothing, and the sweeper used to skip
|
||||
// the issue on every tick: in_progress, no run, no path, no notice.
|
||||
const { companyId, agentId, runId, issueId } =
|
||||
await seedStrandedIssueFixture({
|
||||
status: "in_progress",
|
||||
runStatus: "failed",
|
||||
runErrorCode: "server_shutdown_interrupted",
|
||||
runError: "Interrupted by graceful server shutdown (SIGTERM)",
|
||||
// A conversation adapter's interrupted run keeps its session for
|
||||
// continuation, so legacy reconciliation does not terminalize it;
|
||||
// the sweeper reaches the retry lane with a spent budget.
|
||||
resultJson: {
|
||||
stopReason: "interrupted",
|
||||
conversationContinuation: "continue_conversation_v1",
|
||||
},
|
||||
});
|
||||
const originalRunId = randomUUID();
|
||||
await db
|
||||
.update(heartbeatRuns)
|
||||
.set({
|
||||
status: "interrupted",
|
||||
scheduledRetryAttempt: 2,
|
||||
scheduledRetryReason: "transient_failure",
|
||||
contextSnapshot: {
|
||||
issueId,
|
||||
taskId: issueId,
|
||||
wakeReason: "transient_failure_retry",
|
||||
retryOfRunId: originalRunId,
|
||||
},
|
||||
})
|
||||
.where(eq(heartbeatRuns.id, runId));
|
||||
|
||||
const heartbeat = heartbeatService(db);
|
||||
const result = await heartbeat.reconcileStrandedAssignedIssues();
|
||||
expect(result.escalated).toBe(1);
|
||||
expect(result.continuationRequeued).toBe(0);
|
||||
expect(result.skipped).toBe(0);
|
||||
expect(result.issueIds).toEqual([issueId]);
|
||||
const successors = await db
|
||||
.select({ id: heartbeatRuns.id })
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.retryOfRunId, runId));
|
||||
expect(successors).toHaveLength(0);
|
||||
|
||||
const sourceIssue = await db
|
||||
.select()
|
||||
.from(issues)
|
||||
.where(eq(issues.id, issueId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(sourceIssue?.status).toBe("blocked");
|
||||
expect(sourceIssue?.assigneeAgentId).toBe(agentId);
|
||||
const recoveryActions = await db
|
||||
.select()
|
||||
.from(issueRecoveryActions)
|
||||
.where(eq(issueRecoveryActions.sourceIssueId, issueId));
|
||||
expect(recoveryActions).toHaveLength(1);
|
||||
expect(recoveryActions[0]).toMatchObject({
|
||||
companyId,
|
||||
status: "active",
|
||||
previousOwnerAgentId: agentId,
|
||||
});
|
||||
const comments = await db
|
||||
.select()
|
||||
.from(issueComments)
|
||||
.where(eq(issueComments.issueId, issueId));
|
||||
expect(comments).toHaveLength(1);
|
||||
expect(comments[0]?.body).toContain("bounded retry budget");
|
||||
expect(comments[0]?.body).toContain("no live execution path");
|
||||
expect(comments[0]?.authorType).toBe("system");
|
||||
expect(comments[0]?.presentation).toMatchObject({
|
||||
kind: "system_notice",
|
||||
tone: "danger",
|
||||
});
|
||||
|
||||
// The escalated issue is board-owned now: a second sweep leaves it alone.
|
||||
const again = await heartbeat.reconcileStrandedAssignedIssues();
|
||||
expect(again.escalated).toBe(0);
|
||||
expect(again.continuationRequeued).toBe(0);
|
||||
});
|
||||
|
||||
it("keeps retrying an interrupted run while its transient retry budget remains", async () => {
|
||||
const { agentId, runId, issueId } = await seedStrandedIssueFixture({
|
||||
status: "in_progress",
|
||||
runStatus: "failed",
|
||||
runErrorCode: "server_shutdown_interrupted",
|
||||
runError: "Interrupted by graceful server shutdown (SIGTERM)",
|
||||
});
|
||||
await db
|
||||
.update(heartbeatRuns)
|
||||
.set({
|
||||
status: "interrupted",
|
||||
scheduledRetryAttempt: 1,
|
||||
scheduledRetryReason: "transient_failure",
|
||||
contextSnapshot: {
|
||||
issueId,
|
||||
taskId: issueId,
|
||||
wakeReason: "transient_failure_retry",
|
||||
},
|
||||
})
|
||||
.where(eq(heartbeatRuns.id, runId));
|
||||
|
||||
const heartbeat = heartbeatService(db);
|
||||
const result = await heartbeat.reconcileStrandedAssignedIssues();
|
||||
expect(result.continuationRequeued).toBe(1);
|
||||
expect(result.escalated).toBe(0);
|
||||
const successor = await db
|
||||
.select()
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.retryOfRunId, runId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(successor).toMatchObject({
|
||||
agentId,
|
||||
status: "scheduled_retry",
|
||||
scheduledRetryAttempt: 2,
|
||||
scheduledRetryReason: "transient_failure",
|
||||
});
|
||||
const sourceIssue = await db
|
||||
.select({ status: issues.status })
|
||||
.from(issues)
|
||||
.where(eq(issues.id, issueId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(sourceIssue?.status).toBe("in_progress");
|
||||
});
|
||||
|
||||
it("escalates an assigned todo issue whose lost dispatch has spent its transient retry budget", async () => {
|
||||
const { agentId, runId, issueId } = await seedStrandedIssueFixture({
|
||||
status: "todo",
|
||||
runStatus: "failed",
|
||||
runErrorCode: "server_shutdown_interrupted",
|
||||
runError: "Interrupted by graceful server shutdown (SIGTERM)",
|
||||
resultJson: {
|
||||
stopReason: "interrupted",
|
||||
conversationContinuation: "continue_conversation_v1",
|
||||
},
|
||||
});
|
||||
await db
|
||||
.update(heartbeatRuns)
|
||||
.set({
|
||||
status: "interrupted",
|
||||
scheduledRetryAttempt: 2,
|
||||
scheduledRetryReason: "transient_failure",
|
||||
contextSnapshot: {
|
||||
issueId,
|
||||
taskId: issueId,
|
||||
wakeReason: "transient_failure_retry",
|
||||
},
|
||||
})
|
||||
.where(eq(heartbeatRuns.id, runId));
|
||||
|
||||
const heartbeat = heartbeatService(db);
|
||||
const result = await heartbeat.reconcileStrandedAssignedIssues();
|
||||
expect(result.escalated).toBe(1);
|
||||
expect(result.dispatchRequeued).toBe(0);
|
||||
const sourceIssue = await db
|
||||
.select()
|
||||
.from(issues)
|
||||
.where(eq(issues.id, issueId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(sourceIssue?.status).toBe("blocked");
|
||||
expect(sourceIssue?.assigneeAgentId).toBe(agentId);
|
||||
const comments = await db
|
||||
.select()
|
||||
.from(issueComments)
|
||||
.where(eq(issueComments.issueId, issueId));
|
||||
expect(comments[0]?.body).toContain("bounded retry budget");
|
||||
});
|
||||
|
||||
it("escalates an exhausted successful handoff run that still leaves no disposition", async () => {
|
||||
const { companyId, agentId, runId, issueId } =
|
||||
await seedStrandedIssueFixture({
|
||||
@@ -8809,6 +8998,148 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => {
|
||||
).not.toHaveProperty("modelProfile");
|
||||
});
|
||||
|
||||
it("escalates an execution-review participant whose interrupted run has spent its transient retry budget", async () => {
|
||||
const { companyId, agentId, issueId, runId, wakeupRequestId } =
|
||||
await seedInReviewParticipantRunFixture();
|
||||
const finishedAt = new Date("2026-03-19T00:05:00.000Z");
|
||||
await db
|
||||
.update(heartbeatRuns)
|
||||
.set({
|
||||
status: "interrupted",
|
||||
errorCode: "server_shutdown_interrupted",
|
||||
error: "Interrupted by graceful server shutdown (SIGTERM)",
|
||||
scheduledRetryAttempt: 2,
|
||||
scheduledRetryReason: "transient_failure",
|
||||
resultJson: {
|
||||
stopReason: "interrupted",
|
||||
conversationContinuation: "continue_conversation_v1",
|
||||
},
|
||||
contextSnapshot: {
|
||||
issueId,
|
||||
taskId: issueId,
|
||||
wakeReason: "transient_failure_retry",
|
||||
},
|
||||
startedAt: new Date("2026-03-19T00:00:00.000Z"),
|
||||
finishedAt,
|
||||
updatedAt: finishedAt,
|
||||
})
|
||||
.where(eq(heartbeatRuns.id, runId));
|
||||
await db
|
||||
.update(agentWakeupRequests)
|
||||
.set({ status: "cancelled", finishedAt, updatedAt: finishedAt })
|
||||
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
||||
|
||||
const heartbeat = heartbeatService(db);
|
||||
const result = await heartbeat.reconcileStrandedAssignedIssues();
|
||||
expect(result.escalated).toBe(1);
|
||||
expect(result.reviewParticipantRequeued).toBe(0);
|
||||
expect(result.issueIds).toEqual([issueId]);
|
||||
const successors = await db
|
||||
.select({ id: heartbeatRuns.id })
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.retryOfRunId, runId));
|
||||
expect(successors).toHaveLength(0);
|
||||
const sourceIssue = await db
|
||||
.select()
|
||||
.from(issues)
|
||||
.where(eq(issues.id, issueId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(sourceIssue?.status).toBe("blocked");
|
||||
const recoveryActions = await db
|
||||
.select()
|
||||
.from(issueRecoveryActions)
|
||||
.where(eq(issueRecoveryActions.sourceIssueId, issueId));
|
||||
expect(recoveryActions).toHaveLength(1);
|
||||
expect(recoveryActions[0]).toMatchObject({
|
||||
companyId,
|
||||
status: "active",
|
||||
cause: "execution_review_participant_recovery",
|
||||
previousOwnerAgentId: agentId,
|
||||
});
|
||||
const again = await heartbeat.reconcileStrandedAssignedIssues();
|
||||
expect(again.escalated).toBe(0);
|
||||
});
|
||||
|
||||
it("does not block an issue whose review participant spent its retry budget while another agent is still working it", async () => {
|
||||
// `blocked` is issue-wide. The exhausted participant is a dead end for
|
||||
// the review stage, but another agent with a live run on the same issue
|
||||
// must not be blocked under it: the escalation re-reads the live path
|
||||
// for ANY agent and stands down.
|
||||
const { companyId, issueId, runId, wakeupRequestId } =
|
||||
await seedInReviewParticipantRunFixture();
|
||||
const finishedAt = new Date("2026-03-19T00:05:00.000Z");
|
||||
await db
|
||||
.update(heartbeatRuns)
|
||||
.set({
|
||||
status: "interrupted",
|
||||
errorCode: "server_shutdown_interrupted",
|
||||
error: "Interrupted by graceful server shutdown (SIGTERM)",
|
||||
scheduledRetryAttempt: 2,
|
||||
scheduledRetryReason: "transient_failure",
|
||||
resultJson: {
|
||||
stopReason: "interrupted",
|
||||
conversationContinuation: "continue_conversation_v1",
|
||||
},
|
||||
contextSnapshot: {
|
||||
issueId,
|
||||
taskId: issueId,
|
||||
wakeReason: "transient_failure_retry",
|
||||
},
|
||||
startedAt: new Date("2026-03-19T00:00:00.000Z"),
|
||||
finishedAt,
|
||||
updatedAt: finishedAt,
|
||||
})
|
||||
.where(eq(heartbeatRuns.id, runId));
|
||||
await db
|
||||
.update(agentWakeupRequests)
|
||||
.set({ status: "cancelled", finishedAt, updatedAt: finishedAt })
|
||||
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
||||
const otherAgentId = randomUUID();
|
||||
await db.insert(agents).values({
|
||||
id: otherAgentId,
|
||||
companyId,
|
||||
name: "CodexImplementor",
|
||||
role: "engineer",
|
||||
status: "running",
|
||||
adapterType: "codex_local",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
});
|
||||
await db.insert(heartbeatRuns).values({
|
||||
id: randomUUID(),
|
||||
companyId,
|
||||
agentId: otherAgentId,
|
||||
invocationSource: "automation",
|
||||
triggerDetail: "system",
|
||||
status: "running",
|
||||
contextSnapshot: {
|
||||
issueId,
|
||||
taskId: issueId,
|
||||
wakeReason: "issue_commented",
|
||||
},
|
||||
startedAt: new Date("2026-03-19T00:10:00.000Z"),
|
||||
createdAt: new Date(Date.now() + 1_000),
|
||||
updatedAt: new Date("2026-03-19T00:10:00.000Z"),
|
||||
});
|
||||
|
||||
const result = await heartbeatService(db).reconcileStrandedAssignedIssues();
|
||||
expect(result.escalated).toBe(0);
|
||||
expect(result.reviewParticipantRequeued).toBe(0);
|
||||
const sourceIssue = await db
|
||||
.select({ status: issues.status })
|
||||
.from(issues)
|
||||
.where(eq(issues.id, issueId))
|
||||
.then((rows) => rows[0] ?? null);
|
||||
expect(sourceIssue?.status).toBe("in_review");
|
||||
expect(
|
||||
await db
|
||||
.select()
|
||||
.from(issueRecoveryActions)
|
||||
.where(eq(issueRecoveryActions.sourceIssueId, issueId)),
|
||||
).toEqual([]);
|
||||
});
|
||||
|
||||
it("re-enqueues a stranded execution-review participant when another agent has the latest issue run", async () => {
|
||||
const { companyId, agentId, issueId, runId, wakeupRequestId, stageId } =
|
||||
await seedInReviewParticipantRunFixture();
|
||||
|
||||
@@ -9499,6 +9499,12 @@ export function heartbeatService(
|
||||
const result = await scheduleBoundedRetryForRun(run, agent);
|
||||
return result.outcome === "scheduled" ? result.run : null;
|
||||
},
|
||||
// Mirrors scheduleBoundedRetryForRun's transient budget check: a failed
|
||||
// or interrupted run that has already consumed every bounded transient
|
||||
// attempt cannot be retried again through this lane.
|
||||
transientRetryBudgetSpent: (run) =>
|
||||
executionFailureRetryCount(run) >=
|
||||
BOUNDED_TRANSIENT_HEARTBEAT_RETRY_MAX_ATTEMPTS,
|
||||
});
|
||||
const runDispatch = createRunDispatch(db);
|
||||
|
||||
|
||||
@@ -901,6 +901,16 @@ export function recoveryService(
|
||||
scheduleRecoveryRetry?: (
|
||||
runId: string,
|
||||
) => Promise<typeof heartbeatRuns.$inferSelect | null>;
|
||||
/**
|
||||
* Whether a failed or interrupted run has consumed every bounded
|
||||
* transient retry, so `scheduleRecoveryRetry` can no longer produce a
|
||||
* successor for it. Lets the sweeper tell "no retry because the budget
|
||||
* is spent" (escalate) from "no retry because something else owns the
|
||||
* run" (leave alone).
|
||||
*/
|
||||
transientRetryBudgetSpent?: (
|
||||
run: typeof heartbeatRuns.$inferSelect,
|
||||
) => boolean;
|
||||
liveRunExecutions?: Readonly<{ has(id: string): boolean }>;
|
||||
beforeOrphanedRunTerminalWrite?: (runId: string) => Promise<void>;
|
||||
},
|
||||
@@ -1875,6 +1885,13 @@ export function recoveryService(
|
||||
source: string;
|
||||
retryOfRunId?: string | null;
|
||||
extraContext?: Record<string, unknown>;
|
||||
/**
|
||||
* Out-parameter: set `retryExhausted` when no successor was queued
|
||||
* because the failed predecessor has spent its bounded transient retry
|
||||
* budget. A null return without it means another authority owns the
|
||||
* run (native runtime, reconciliation) and the caller must leave it.
|
||||
*/
|
||||
outcome?: { retryExhausted?: boolean };
|
||||
}) {
|
||||
if (input.retryOfRunId) {
|
||||
const [predecessor] = await db
|
||||
@@ -1903,8 +1920,19 @@ export function recoveryService(
|
||||
});
|
||||
return null;
|
||||
}
|
||||
if (deps.scheduleRecoveryRetry)
|
||||
return deps.scheduleRecoveryRetry(predecessor.id);
|
||||
if (deps.scheduleRecoveryRetry) {
|
||||
const retry = await deps.scheduleRecoveryRetry(predecessor.id);
|
||||
if (retry) return retry;
|
||||
// A spent budget is the one "no retry" the sweeper must not wait
|
||||
// out: nothing else will ever queue a successor for this run, so
|
||||
// it reports the exhaustion instead of leaving the issue with no
|
||||
// live path (2026-09-25: three deploy restarts in a row spent the
|
||||
// budget and the issue sat in_progress with no run, unescalated).
|
||||
if (input.outcome && deps.transientRetryBudgetSpent?.(predecessor)) {
|
||||
input.outcome.retryExhausted = true;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
@@ -4855,6 +4883,7 @@ export function recoveryService(
|
||||
continue;
|
||||
}
|
||||
|
||||
const reviewOutcome: { retryExhausted?: boolean } = {};
|
||||
const queued = await enqueueStrandedIssueRecovery({
|
||||
issueId: issue.id,
|
||||
agentId: participantAgentId,
|
||||
@@ -4868,10 +4897,35 @@ export function recoveryService(
|
||||
reviewRecoveryInstruction:
|
||||
"The previous reviewer run ended while this execution-review stage was still pending. Submit the review decision now, or mark the issue blocked with the exact unblock action.",
|
||||
},
|
||||
outcome: reviewOutcome,
|
||||
});
|
||||
if (queued) {
|
||||
result.reviewParticipantRequeued += 1;
|
||||
result.issueIds.push(issue.id);
|
||||
} else if (
|
||||
reviewOutcome.retryExhausted &&
|
||||
!(await hasActiveExecutionPath(issue.companyId, issue.id, null))
|
||||
) {
|
||||
// Same exhaustion as the other lanes: the reviewer run's bounded
|
||||
// retries are spent, so escalate as the review-recovery failure it
|
||||
// is instead of skipping on every sweep with no live path. The
|
||||
// live-path re-read guards the write against a run or wake that
|
||||
// started after the loop's check — for ANY agent, not only the
|
||||
// participant: `blocked` is issue-wide, so another agent still
|
||||
// working the issue must never be blocked under.
|
||||
const updated = await escalateStrandedAssignedIssue({
|
||||
issue,
|
||||
previousStatus: "in_review",
|
||||
latestRun: participantLatestRun,
|
||||
notice: buildExecutionReviewParticipantRecoveryNoticeSeed(),
|
||||
recoveryCause: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_REASON,
|
||||
});
|
||||
if (updated) {
|
||||
result.escalated += 1;
|
||||
result.issueIds.push(issue.id);
|
||||
} else {
|
||||
result.skipped += 1;
|
||||
}
|
||||
} else {
|
||||
result.skipped += 1;
|
||||
}
|
||||
@@ -4950,6 +5004,7 @@ export function recoveryService(
|
||||
continue;
|
||||
}
|
||||
|
||||
const dispatchOutcome: { retryExhausted?: boolean } = {};
|
||||
const queued = await enqueueStrandedIssueRecovery({
|
||||
issueId: issue.id,
|
||||
agentId,
|
||||
@@ -4957,10 +5012,39 @@ export function recoveryService(
|
||||
retryReason: "assignment_recovery",
|
||||
source: "issue.assignment_recovery",
|
||||
retryOfRunId: latestRun.id,
|
||||
outcome: dispatchOutcome,
|
||||
});
|
||||
if (queued) {
|
||||
result.dispatchRequeued += 1;
|
||||
result.issueIds.push(issue.id);
|
||||
} else if (
|
||||
dispatchOutcome.retryExhausted &&
|
||||
!(await hasActiveExecutionPath(issue.companyId, issue.id, null))
|
||||
) {
|
||||
// Same exhaustion as the in_progress lane: the lost dispatch's
|
||||
// bounded retries are spent, so escalate instead of skipping on
|
||||
// every sweep with no live path. The live-path re-read guards the
|
||||
// `blocked` write against a run or wake that started after the
|
||||
// loop's check.
|
||||
const updated = await escalateStrandedAssignedIssue({
|
||||
issue,
|
||||
previousStatus: "todo",
|
||||
latestRun,
|
||||
notice: {
|
||||
body:
|
||||
"Paperclip automatically retried dispatch for this assigned `todo` issue after a lost wake/run, " +
|
||||
"but the bounded retry budget is spent and it still has no live execution path. " +
|
||||
"Moving it to `blocked` so it is visible for intervention.",
|
||||
title: "No live execution path",
|
||||
tone: "danger",
|
||||
},
|
||||
});
|
||||
if (updated) {
|
||||
result.escalated += 1;
|
||||
result.issueIds.push(issue.id);
|
||||
} else {
|
||||
result.skipped += 1;
|
||||
}
|
||||
} else {
|
||||
result.skipped += 1;
|
||||
}
|
||||
@@ -5217,6 +5301,7 @@ export function recoveryService(
|
||||
continue;
|
||||
}
|
||||
|
||||
const recoveryOutcome: { retryExhausted?: boolean } = {};
|
||||
const queued = await enqueueStrandedIssueRecovery({
|
||||
issueId: issue.id,
|
||||
agentId,
|
||||
@@ -5224,10 +5309,40 @@ export function recoveryService(
|
||||
retryReason: "issue_continuation_needed",
|
||||
source: "issue.continuation_recovery",
|
||||
retryOfRunId: latestRun?.id ?? issue.checkoutRunId ?? null,
|
||||
outcome: recoveryOutcome,
|
||||
});
|
||||
if (queued) {
|
||||
result.continuationRequeued += 1;
|
||||
result.issueIds.push(issue.id);
|
||||
} else if (
|
||||
recoveryOutcome.retryExhausted &&
|
||||
!(await hasActiveExecutionPath(issue.companyId, issue.id, null))
|
||||
) {
|
||||
// The failed run's bounded transient retries are all spent, so no
|
||||
// successor will ever be queued for it. Escalate rather than skip
|
||||
// on every sweep forever: the issue would otherwise stay
|
||||
// `in_progress` with no run and no path until a person noticed.
|
||||
// The live-path re-read guards the `blocked` write against a run or
|
||||
// wake that started after the loop's check.
|
||||
const updated = await escalateStrandedAssignedIssue({
|
||||
issue,
|
||||
previousStatus: "in_progress",
|
||||
latestRun,
|
||||
notice: {
|
||||
body:
|
||||
"Paperclip retried this issue's run after it ended without finishing, but the bounded retry budget " +
|
||||
"is spent and it still has no live execution path. " +
|
||||
"Moving it to `blocked` so it is visible for intervention.",
|
||||
title: "No live execution path",
|
||||
tone: "danger",
|
||||
},
|
||||
});
|
||||
if (updated) {
|
||||
result.escalated += 1;
|
||||
result.issueIds.push(issue.id);
|
||||
} else {
|
||||
result.skipped += 1;
|
||||
}
|
||||
} else {
|
||||
result.skipped += 1;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user