diff --git a/doc/run-log-events.md b/doc/run-log-events.md index df12d118b6..e468d7ff86 100644 --- a/doc/run-log-events.md +++ b/doc/run-log-events.md @@ -159,6 +159,15 @@ Provider identity diagnostics remain in the local run log. They record the notif Recovery lifecycle events retain the original structured failure code, retry attempt, next retry time, and predecessor/successor identifiers. Durable status delivery uses an idempotency marker; delivery grants no provider authority. Failed publication is retried without repeating provider work. These records are not first-party Telemetry. +If execution-continuation setup finds that a task no longer exists, is closed, +or its owner changed, the existing cancellation settlement records +`continuation_task_ownership_changed`. The run and wake request become cancelled +before adapter dispatch, with the run-log message +`stale execution continuation cancelled before dispatch`. Immediate recovery is +suppressed. Missing source context, authorization failures, and other setup +errors retain their failure classification. An untyped error with the same +message is also still a failure; cancellation requires the typed ownership guard. + Bounded retry exhaustion writes one lifecycle receipt per run, retry reason, scheduled attempt, and retry limit. Repeated or concurrent recovery checks reuse that receipt, including receipts from earlier builds, without advancing the event diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index e7325a17a4..4eca1577e6 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -1,3 +1,4 @@ +import * as executionContinuation from "../services/execution-continuation.js"; import { legacyDispositionFingerprint, LEGACY_DISPOSITION_REPAIR_INSTRUCTION } from "../services/recovery/legacy-continuation.js"; import * as controllerLeases from "../services/legacy-controller-lease.js"; import { instanceSettingsService } from "../services/instance-settings.js"; @@ -1528,6 +1529,62 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(missingCommentWakeups).toHaveLength(0); }); + it.each([ + "continuation_task_ownership_changed", + ] as const)("cancels obsolete continuation setup: %s", async (code) => { + const { agentId, runId, wakeupRequestId } = await seedQueuedIssueRunFixture(); + const build = vi.spyOn(executionContinuation, "buildExecutionContinuation") + .mockRejectedValueOnce(new executionContinuation.StaleExecutionContinuationError(code)); + try { + const heartbeat = heartbeatService(db); + await heartbeat.resumeQueuedRuns(); + await waitForRunToSettle(heartbeat, runId); + await heartbeat.waitForRunExecutionDrain(runId); + expect(build).toHaveBeenCalled(); + expect(await heartbeat.getRun(runId)).toMatchObject({ status: "cancelled", errorCode: code }); + const [wakeup] = await db.select().from(agentWakeupRequests) + .where(eq(agentWakeupRequests.id, wakeupRequestId)); + expect(wakeup.status).toBe("cancelled"); + const [agent] = await db.select().from(agents).where(eq(agents.id, agentId)); + expect(agent.status).not.toBe("error"); + expect(mockAdapterExecute.mock.calls.some( + ([input]) => (input as { runId?: string } | undefined)?.runId === runId, + )).toBe(false); + const runs = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, agentId)); + expect(runs).toHaveLength(1); + } finally { + build.mockRestore(); + } + }); + + it.each([ + "continuation_source_context_missing", + "continuation_user_authorization_missing", + "continuation_task_ownership_changed", + ])("retains untyped continuation setup failures: %s", async (message) => { + const { runId, wakeupRequestId } = await seedQueuedIssueRunFixture(); + const build = vi.spyOn(executionContinuation, "buildExecutionContinuation") + .mockRejectedValueOnce(new Error(message)); + try { + const heartbeat = heartbeatService(db); + await heartbeat.resumeQueuedRuns(); + await waitForRunToSettle(heartbeat, runId); + await heartbeat.waitForRunExecutionDrain(runId); + expect(build).toHaveBeenCalled(); + expect(await heartbeat.getRun(runId)).toMatchObject({ + status: "failed", errorCode: "setup_failed", error: message, + }); + const [wakeup] = await db.select().from(agentWakeupRequests) + .where(eq(agentWakeupRequests.id, wakeupRequestId)); + expect(wakeup.status).toBe("failed"); + expect(mockAdapterExecute.mock.calls.some( + ([input]) => (input as { runId?: string } | undefined)?.runId === runId, + )).toBe(false); + } finally { + build.mockRestore(); + } + }); + it("does not immediately continue a low-trust preflight setup failure", async () => { const { agentId, runId, issueId, companyId } = await seedQueuedIssueRunFixture(); diff --git a/server/src/services/execution-continuation.test.ts b/server/src/services/execution-continuation.test.ts index 97336ffa1f..047038e53d 100644 --- a/server/src/services/execution-continuation.test.ts +++ b/server/src/services/execution-continuation.test.ts @@ -15,7 +15,20 @@ import { getEmbeddedPostgresTestSupport, startEmbeddedPostgresTestDatabase, } from "../__tests__/helpers/embedded-postgres.js"; -import { buildExecutionContinuation, currentContinuationOrigins, projectHumanInteractionResponse } from "./execution-continuation.js"; +import { StaleExecutionContinuationError, buildExecutionContinuation, currentContinuationOrigins, projectHumanInteractionResponse } from "./execution-continuation.js"; + +const expectStaleContinuation = async ( + run: () => Promise, + code: StaleExecutionContinuationError["code"], +) => { + await expect(run()).rejects.toThrow(StaleExecutionContinuationError); + await expect(run()).rejects.toMatchObject({ code }); +}; +const expectMissingContinuationContext = async (run: () => Promise) => { + const result = run(); + await expect(result).rejects.toThrow("continuation_source_context_missing"); + await expect(result).rejects.not.toBeInstanceOf(StaleExecutionContinuationError); +}; const support = await getEmbeddedPostgresTestSupport(); (support.supported ? describe : describe.skip)( "authorized continuation context", @@ -174,9 +187,10 @@ const support = await getEmbeddedPostgresTestSupport(); const [source] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)); await db.update(heartbeatRuns).set({ contextSnapshot: { issueId: randomUUID() } }).where(eq(heartbeatRuns.id, runId)); try { - await expect(buildExecutionContinuation({ db, companyId, issueId, agentId, - context: { interruptedRunId: runId }, summary: null, exposeLowTrustRaw: false })) - .rejects.toThrow("continuation_source_context_missing"); + await expectMissingContinuationContext( + () => buildExecutionContinuation({ db, companyId, issueId, agentId, + context: { interruptedRunId: runId }, summary: null, exposeLowTrustRaw: false }), + ); } finally { await db.update(heartbeatRuns).set({ contextSnapshot: source.contextSnapshot }).where(eq(heartbeatRuns.id, runId)); } @@ -336,41 +350,63 @@ const support = await getEmbeddedPostgresTestSupport(); expect(freshPrompt).not.toContain('"resumeDelta"'); }); it("fails closed when required originating context is missing", async () => { - await expect( - buildExecutionContinuation({ - db, - companyId, - issueId, - agentId, - context: { commentId: randomUUID() }, - summary: null, - exposeLowTrustRaw: false, - }), - ).rejects.toThrow("continuation_source_context_missing"); + await expectMissingContinuationContext( + () => + buildExecutionContinuation({ + db, + companyId, + issueId, + agentId, + context: { commentId: randomUUID() }, + summary: null, + exposeLowTrustRaw: false, + }), + ); + }); + it.each(["done", "cancelled"])("rejects continuation after the task becomes %s", async (status) => { + await db.update(issues).set({ status }).where(eq(issues.id, issueId)); + try { + await expectStaleContinuation( + () => buildExecutionContinuation({ + db, companyId, issueId, agentId, + context: { wakeReason: "issue_commented", commentId: gmailId }, + summary: null, exposeLowTrustRaw: false, + }), + "continuation_task_ownership_changed", + ); + const [issue] = await db.select().from(issues).where(eq(issues.id, issueId)); + expect(issue).toMatchObject({ status, assigneeAgentId: agentId }); + } finally { + await db.update(issues).set({ status: "in_progress" }).where(eq(issues.id, issueId)); + } }); it("rejects another company and an invalidated task owner", async () => { - await expect( - buildExecutionContinuation({ - db, - companyId: randomUUID(), - issueId, - agentId, - context: {}, - summary: null, - exposeLowTrustRaw: false, - }), - ).rejects.toThrow("continuation_task_ownership_changed"); - await expect( - buildExecutionContinuation({ - db, - companyId, - issueId, - agentId: randomUUID(), - context: {}, - summary: null, - exposeLowTrustRaw: false, - }), - ).rejects.toThrow("continuation_task_ownership_changed"); + await expectStaleContinuation( + () => + buildExecutionContinuation({ + db, + companyId: randomUUID(), + issueId, + agentId, + context: {}, + summary: null, + exposeLowTrustRaw: false, + }), + "continuation_task_ownership_changed", + ); + await expectStaleContinuation( + () => + buildExecutionContinuation({ + db, + companyId, + issueId, + agentId: randomUUID(), + context: {}, + summary: null, + exposeLowTrustRaw: false, + }), + "continuation_task_ownership_changed", + ); }); }, ); diff --git a/server/src/services/execution-continuation.ts b/server/src/services/execution-continuation.ts index 82aa0753f9..89701777c4 100644 --- a/server/src/services/execution-continuation.ts +++ b/server/src/services/execution-continuation.ts @@ -15,6 +15,13 @@ import { hasConversationContinuationPolicy } from "./conversation-continuation.j import { queuedCommentIdsFromWakePayload } from "./issue-queued-comment-queue.js"; import { childReviewOutcomes } from "./native-runtime/child-review-outcomes.js"; +export class StaleExecutionContinuationError extends Error { + constructor(readonly code: "continuation_task_ownership_changed") { + super(code); + this.name = "StaleExecutionContinuationError"; + } +} + const object = (v: unknown): Record => v && typeof v === "object" && !Array.isArray(v) ? (v as Record) @@ -127,7 +134,7 @@ export async function buildExecutionContinuation(input: { issue.assigneeAgentId !== input.agentId || ["done", "cancelled"].includes(issue.status) ) - throw new Error("continuation_task_ownership_changed"); + throw new StaleExecutionContinuationError("continuation_task_ownership_changed"); const rows = await db .select() .from(issueComments) diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index d25c617b47..8c2979dd13 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -37,7 +37,7 @@ import { import { executionFailureRetryCount, executionRetryAttemptCount, accountingForScheduledRetry } from "./execution-recovery-attempt.js"; import { buildHeartbeatRunStatusLiveEventPayload } from "./heartbeat-run-status-payload.js"; export { buildHeartbeatRunStatusLiveEventPayload } from "./heartbeat-run-status-payload.js"; -import { buildExecutionContinuation } from "./execution-continuation.js"; +import { buildExecutionContinuation, StaleExecutionContinuationError } from "./execution-continuation.js"; import { renderPaperclipWakePrompt } from "@paperclipai/adapter-utils/server-utils"; import { PROJECT_REPOSITORIES_DIR, readGitWorkspaceSnapshot } from "@paperclipai/adapter-utils/git-workspace-sync"; import { isWorkspaceGitScanError, WorkspaceGitScanError, WORKSPACE_GIT_SCAN_ERROR_CODES } from "./workspace-git-operation-scheduler.js"; @@ -25679,6 +25679,15 @@ export function heartbeatService( ? outerErr.reason : "adopted_runner_authentication_timeout", }).catch(() => undefined); + } else if (outerErr instanceof StaleExecutionContinuationError) { + // The queued continuation became obsolete before adapter dispatch. + // Use cancellation settlement so wakeup, issue ownership, agent state, + // and notifications agree; do not retry work for the previous owner. + await cancelRunInternal(run.id, outerErr.code, { + errorCode: outerErr.code, + eventMessage: "stale execution continuation cancelled before dispatch", + suppressImmediateRecovery: true, + }); } else if (isWorkspaceBusyDeferral(outerErr)) { // Expected contention on a shared project workspace, not a // failure: park the run as a bounded scheduled retry and leave the