diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index 6d8ec4af30..3e3cd076ba 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -484,6 +484,10 @@ An explicit recovery action is a typed liveness repair path for a source issue. A terminal native failure can retain a result accepted before checkpoint or cleanup failed. That result is historical evidence, not a live controller. New user input may start a fresh turn after the controller and execution environment have stopped and ordinary admission checks pass. Preserve the failed run, its result, and its recovery budget. Do not commit the old result, infer action outcomes, or replay the failed turn. A message saved while cleanup is pending must be reconsidered after verified cleanup and delivered once. Admission must consume its deferred receipt in the same transaction that creates the new run, including in agent chat, where later messages keep their separate turns. Completing that run must not promote the consumed receipt again. Workspace-export repair retains its separate saved-result recovery path. +When a legacy turn admitted by an explicit user message fails or times out, a bounded transient retry needs its own durable continuation authorization. Initial admission binds the message's exact body digest and revision. Scheduling requires that original binding and revalidates the human author, source task, assignee, and execution ownership, then records the successor's receipt in the same transaction as its run. Historical receipts without a message binding cannot authorize automatic retries and are not backfilled. The parent receipt remains intact. Dispatch checks the exact successor, retry parent, and unchanged message binding again. Deleted or edited input, changed actors or ownership, superseded work, and duplicate successors cannot renew authority. One-run Retry and Interrupt intents retain their separate admission contracts; they are not renewable user-message authority. Other retry policies cannot borrow this receipt: their scheduling fails closed. Successful explicit turns do not create a prose-only missing-comment follow-up with spent authority. + +Explicit retry admission checks the local parent adapter controller and verified environment cleanup. The existing queue-first release policy may retain only its own terminal execution claim while that cleanup is pending; finalization and the periodic stale-lock sweep reconsider the same policy after cleanup. Generic terminal-lock cleanup and stale checkout adoption leave these claims for that settlement. New user input, reassignment, and a successor's claim take precedence. Exhausted retries and revoked or missing authorization release the terminal claim without requesting another retry. Pending or failed environment cleanup cannot mint a successor receipt. + A new user message can continue a terminal native run whose process fields were cleared before local stop receipts existed. Admission must verify the exact run, runner, workspace, and provider session in the retained suspended state, with no active provider turn, pending tool call, or undelivered output. Missing or mismatched state keeps the hold. A later recorded process launch also keeps the hold until its stop is verified. Normal assignment, decision, controller, environment cleanup, and active-run gates still apply. The message starts one fresh conversation turn; it does not replay the failed run, reset its recovery budget, or certify unknown action outcomes. The task thread exposes the guarded Retry action for failed or timed-out conversation runs and native preparation cancelled before provider startup when the never-started proof is verified. The server projects this eligibility on the recovery notice and rechecks it on Retry. Cleanup quarantine retains its inspection path, and non-conversation reconciliation gates remain enforced. Pending decisions, active execution, pause, budget, dependency, and ownership gates remain in force. A refused Retry reports its reason inline and remains available for another attempt. diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index d98b6d4e93..80d2011761 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -1,10 +1,12 @@ import * as executionContinuation from "../services/execution-continuation.js"; +import * as environmentOrchestrator from "../services/environment-run-orchestrator.js"; +import { remoteTerminationReceipt } from "../services/remote-execution-termination.js"; import { legacyDispositionFingerprint, LEGACY_DISPOSITION_REPAIR_INSTRUCTION } from "../services/recovery/legacy-continuation.js"; import * as controllerLeases from "../services/legacy-controller-lease.js"; import * as instructionWorkingCopies from "../services/agent-instruction-working-copies.js"; import * as runEvents from "../services/heartbeat-run-events.js"; import { instanceSettingsService } from "../services/instance-settings.js"; -import { randomUUID } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import { terminalizeLegacyExecution } from "../services/legacy-execution-recovery.js"; import { issueService } from "../services/issues.js"; import { getExecutionBlocker } from "../services/execution-blocker.js"; @@ -1489,6 +1491,207 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { return { companyId, agentId, runId, wakeupRequestId, issueId }; } + it.each(["timeout", "upstream", "cleanup_pending", "cleanup_pending_edited", "cleanup_pending_exhausted", "edited", "exhausted", "new_message", "no_claim", "reassigned", "superseded", "stopped"])( + "settles explicit retry admission through real executor cleanup: %s", async scenario => { + const { companyId, agentId, runId, issueId } = await seedQueuedIssueRunFixture(); + const sourceRunId = randomUUID(), commentId = randomUUID(), nextCommentId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: sourceRunId, companyId, agentId, status: "failed", + createdAt: new Date("2020-01-01T00:00:00Z"), finishedAt: new Date("2020-01-01T00:00:01Z"), + contextSnapshot: { issueId } }); + const [comment] = await db.insert(issueComments).values({ id: commentId, companyId, issueId, + authorType: "user", authorUserId: "responsible-user", body: "Continue the task." }).returning(); + await db.insert(issueRecoveryActions).values({ companyId, sourceIssueId: issueId, + kind: "active_run_watchdog", cause: "legacy_execution_requires_reconciliation", fingerprint: sourceRunId, + status: "resolved", outcome: "cancelled", nextAction: "User continued.", evidence: { explicitUserContinuation: { + runId, previousRunId: sourceRunId, commentId, actorId: "responsible-user", recordedAt: new Date().toISOString(), + commentUpdatedAt: comment!.updatedAt.toISOString(), commentBodyHash: createHash("sha256").update(comment!.body).digest("hex"), + } } }); + await db.update(heartbeatRuns).set({ createdAt: new Date(), contextSnapshot: { + issueId, taskId: issueId, wakeReason: "issue_commented", wakeCommentId: commentId, + previousRunId: sourceRunId, explicitUserContinuation: { previousRunId: sourceRunId, commentId }, + } }).where(eq(heartbeatRuns.id, runId)); + const identity = { id: randomUUID(), companyId, heartbeatRunId: runId, provider: "daytona", providerLeaseId: "fixture-explicit-sandbox" }; + const supersedingRunId = randomUUID(); + let heartbeat: ReturnType; + let cancellation: ReturnType | undefined; + let cleanupObserved = false; + let cleanupPending = scenario.startsWith("cleanup_pending"); + const originalFactory = environmentOrchestrator.environmentRunOrchestrator; + const release = vi.fn(async () => { + const [lease] = await db.select().from(environmentLeases).where(eq(environmentLeases.id, identity.id)); + if (!lease) return { released: [], errors: [] }; + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + if (!cleanupObserved) { + cleanupObserved = true; + const denied = await db.select().from(heartbeatRunEvents).where(and(eq(heartbeatRunEvents.runId, runId), + eq(heartbeatRunEvents.message, "The automatic retry could not revalidate the original user continuation."))); + // Upstream failure first tries the direct scheduler, then the release + // tail. Both must defer admission until the executor releases ownership. + if (scenario === "upstream") expect(denied.length).toBeGreaterThanOrEqual(2); + if (scenario === "timeout") expect(denied.length).toBeGreaterThanOrEqual(1); + if (scenario === "new_message") { + await db.insert(issueComments).values({ id: nextCommentId, companyId, issueId, authorType: "user", + authorUserId: "responsible-user", body: "Use these new instructions." }); + await heartbeat.wakeup(agentId, { source: "automation", reason: "issue_commented", + requestedByActorType: "user", requestedByActorId: "responsible-user", + payload: { issueId, commentId: nextCommentId }, + contextSnapshot: { issueId, commentId: nextCommentId, wakeReason: "issue_commented" } }); + } + if (scenario === "reassigned") { + const [other] = await db.insert(agents).values({ companyId, name: "New assignee", role: "engineer", + adapterType: "codex_local" }).returning(); + await db.update(issues).set({ assigneeAgentId: other!.id }).where(eq(issues.id, issueId)); + } + if (scenario === "superseded") { + await db.insert(heartbeatRuns).values({ id: supersedingRunId, companyId, agentId, status: "running", + contextSnapshot: { issueId } }); + await db.update(issues).set({ executionRunId: supersedingRunId }).where(eq(issues.id, issueId)); + } + } + await db.update(environmentLeases).set({ status: "released", cleanupStatus: "success", releasedAt: new Date() }) + .where(and(eq(environmentLeases.companyId, companyId), eq(environmentLeases.heartbeatRunId, runId), + eq(environmentLeases.provider, "local"))); + await db.update(environmentLeases).set(cleanupPending + ? { status: "pending_cleanup", cleanupStatus: "failed", releasedAt: new Date() } + : { status: "released", cleanupStatus: "success", releasedAt: new Date(), metadata: { + remoteExecutionTermination: remoteTerminationReceipt(identity, { providerLeaseId: identity.providerLeaseId, state: "destroyed" }), + } }).where(eq(environmentLeases.id, identity.id)); + return { released: [], errors: [] }; + }); + const factory = vi.spyOn(environmentOrchestrator, "environmentRunOrchestrator").mockImplementation((...args) => ({ + ...originalFactory(...args), releaseForRun: release, + })); + heartbeat = heartbeatService(db); + mockAdapterExecute.mockImplementationOnce(async (input?: unknown) => { + const invocation = input as { onCancellationReady: () => Promise; onDispatch: () => void }; + await invocation.onCancellationReady(); + invocation.onDispatch(); + expect(adapterExecutionControls.has(runId)).toBe(true); + await db.insert(environmentLeases).values({ ...identity, status: "active", leasePolicy: "ephemeral" }); + if (scenario === "edited") await db.update(issueComments).set({ body: "Changed after dispatch" }).where(eq(issueComments.id, commentId)); + if (scenario === "no_claim") await db.update(issues).set({ executionRunId: null }).where(eq(issues.id, issueId)); + if (scenario === "exhausted") await db.update(heartbeatRuns).set({ scheduledRetryAttempt: 2, + scheduledRetryReason: "transient_failure" }).where(eq(heartbeatRuns.id, runId)); + if (scenario === "stopped") { + cancellation = heartbeat.cancelRun(runId, "Operator Stop"); + void cancellation.catch(() => undefined); + await waitForValue(async () => adapterExecutionControls.get(runId)?.controller.signal.aborted); + } + return { exitCode: 1, signal: scenario === "stopped" ? "SIGTERM" : null, + timedOut: scenario !== "upstream" && scenario !== "stopped", + errorMessage: scenario === "upstream" ? "Upstream temporarily unavailable" : "Execution timed out", + errorCode: scenario === "upstream" ? "upstream_unavailable" : "timeout", + resultJson: scenario === "upstream" ? { errorFamily: "transient_upstream" } + : scenario === "stopped" ? { executionCancellation: { state: "acknowledged" } } : {}, + summary: "Stopped turn", provider: "test", model: "test-model" }; + }); + mockAdapterExecute.mockImplementation(async () => { + await db.update(issues).set({ status: "done" }).where(eq(issues.id, issueId)); + return { exitCode: 0, signal: null, timedOut: false, errorMessage: null, + summary: "New message handled.", provider: "test", model: "test-model" }; + }); + try { + await heartbeat.resumeQueuedRuns(); + await heartbeat.drainActiveRunExecutions(); + if (cancellation) expect((await cancellation)?.status).toBe("cancelled"); + expect(release).toHaveBeenCalled(); + expect(adapterExecutionControls.has(runId)).toBe(false); + const children = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId)); + const receipts = await db.select().from(issueRecoveryActions).where(and(eq(issueRecoveryActions.sourceIssueId, issueId), + eq(issueRecoveryActions.cause, "explicit_user_continuation_retry"))); + const [task] = await db.select().from(issues).where(eq(issues.id, issueId)); + if (scenario === "timeout" || scenario === "upstream") { + expect(children).toHaveLength(1); + expect(receipts).toHaveLength(1); + expect(task!.executionRunId).toBe(children[0]!.id); + await expect(executionContinuation.buildExecutionContinuation({ db, companyId, issueId, agentId, + runId: children[0]!.id, context: children[0]!.contextSnapshot!, summary: null, exposeLowTrustRaw: false })) + .resolves.toMatchObject({ interruptedRunId: runId }); + expect(await heartbeat.scheduleBoundedRetry(runId)).toMatchObject({ outcome: "scheduled", run: { id: children[0]!.id } }); + } else { + expect(children).toHaveLength(0); + expect(receipts).toHaveLength(0); + expect(task!.executionRunId).toBe(cleanupPending ? runId : scenario === "superseded" ? supersedingRunId : null); + } + if (scenario === "new_message") expect(await db.select().from(heartbeatRuns).where(and( + eq(heartbeatRuns.companyId, companyId), sql`${heartbeatRuns.contextSnapshot}->>'wakeCommentId' = ${nextCommentId}`))) + .toHaveLength(1); + if (cleanupPending) { + if (scenario === "cleanup_pending") { + const svc = issueService(db); + const [other] = await db.insert(agents).values({ companyId, name: "Other worker", role: "engineer", + adapterType: "codex_local" }).returning(); + const actorRunId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: actorRunId, companyId, agentId, status: "running", contextSnapshot: {} }); + // Neither a rejected checkout nor stale checkout adoption may erase + // the exact pending retry claim, including mixed checkout owners. + for (const checkoutRunId of [null, runId, sourceRunId]) { + await db.update(issues).set({ checkoutRunId }).where(eq(issues.id, issueId)); + expect(await svc.clearExecutionRunIfTerminal(issueId)).toBe(false); + expect(await svc.clearCheckoutRunIfTerminal(issueId)).toBe(false); + await expect(svc.checkout(issueId, other!.id, ["in_progress"], actorRunId)).rejects.toMatchObject({ status: 409 }); + await expect(svc.assertCheckoutOwner(issueId, other!.id, actorRunId)).rejects.toMatchObject({ status: 409 }); + await expect(svc.checkout(issueId, agentId, ["in_progress"], actorRunId)).rejects.toMatchObject({ status: 409 }); + await expect(svc.assertCheckoutOwner(issueId, agentId, actorRunId)).rejects.toMatchObject({ status: 409 }); + expect((await db.select().from(issues).where(eq(issues.id, issueId)))[0]!.executionRunId).toBe(runId); + } + await db.update(issues).set({ checkoutRunId: null }).where(eq(issues.id, issueId)); + await db.update(heartbeatRuns).set({ status: "failed", finishedAt: new Date() }).where(eq(heartbeatRuns.id, actorRunId)); + } + await heartbeat.sweepStaleIssueLocks(); + expect((await db.select().from(issues).where(eq(issues.id, issueId)))[0]!.executionRunId).toBe(runId); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + expect(await db.select().from(issueRecoveryActions).where(and(eq(issueRecoveryActions.sourceIssueId, issueId), + eq(issueRecoveryActions.cause, "explicit_user_continuation_retry")))).toHaveLength(0); + if (scenario === "cleanup_pending_edited") await db.update(issueComments).set({ body: "Changed while cleanup waited" }) + .where(eq(issueComments.id, commentId)); + if (scenario === "cleanup_pending_exhausted") await db.update(heartbeatRuns).set({ scheduledRetryAttempt: 2, + scheduledRetryReason: "transient_failure" }).where(eq(heartbeatRuns.id, runId)); + cleanupPending = false; + await heartbeat.releaseEnvironmentLeasesForRun({ runId, companyId, agentId, status: "timed_out" }); + await heartbeat.sweepStaleIssueLocks(); + await heartbeat.sweepStaleIssueLocks(); + const afterCleanup = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId)); + expect(afterCleanup).toHaveLength(scenario === "cleanup_pending" ? 1 : 0); + expect(await db.select().from(issueRecoveryActions).where(and(eq(issueRecoveryActions.sourceIssueId, issueId), + eq(issueRecoveryActions.cause, "explicit_user_continuation_retry")))).toHaveLength(afterCleanup.length); + expect((await db.select().from(issues).where(eq(issues.id, issueId)))[0]!.executionRunId).toBe(afterCleanup[0]?.id ?? null); + } + } finally { factory.mockRestore(); } + }, + ); + + it("does not reuse a completed explicit turn's receipt for a missing-comment follow-up", async () => { + const { companyId, agentId, runId, issueId } = await seedQueuedIssueRunFixture(); + const sourceRunId = randomUUID(), commentId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: sourceRunId, companyId, agentId, status: "failed", + finishedAt: new Date(), contextSnapshot: { issueId } }); + await db.insert(issueComments).values({ id: commentId, companyId, issueId, + authorType: "user", authorUserId: "responsible-user", body: "Continue the task." }); + await db.insert(issueRecoveryActions).values({ companyId, sourceIssueId: issueId, + kind: "active_run_watchdog", cause: "legacy_execution_requires_reconciliation", fingerprint: sourceRunId, + status: "resolved", outcome: "cancelled", nextAction: "User continued.", + evidence: { explicitUserContinuation: { runId, previousRunId: sourceRunId, commentId, + actorId: "responsible-user", recordedAt: new Date().toISOString() } } }); + await db.update(heartbeatRuns).set({ contextSnapshot: { + issueId, taskId: issueId, wakeReason: "issue_assigned", + explicitUserContinuation: { previousRunId: sourceRunId, commentId }, + } }).where(eq(heartbeatRuns.id, runId)); + mockAdapterExecute.mockImplementationOnce(async () => { + await db.update(issues).set({ status: "done" }).where(eq(issues.id, issueId)); + return { exitCode: 0, signal: null, timedOut: false, errorMessage: null, + summary: "Task completed.", provider: "test", model: "test-model" }; + }); + const heartbeat = heartbeatService(db); + await heartbeat.resumeQueuedRuns(); + await waitForRunToSettle(heartbeat, runId); + await heartbeat.waitForRunExecutionDrain(runId); + expect(await heartbeat.getRun(runId)).toMatchObject({ status: "succeeded", issueCommentStatus: "not_applicable" }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + expect(await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.companyId, companyId), + eq(agentWakeupRequests.reason, "missing_issue_comment")))).toHaveLength(0); + }); + it("persists the normalized failure without permanently blocking the conversation", async () => { mockAdapterExecute.mockResolvedValueOnce({ exitCode: 1, diff --git a/server/src/__tests__/recovery-stale-issue-lock-sweep.test.ts b/server/src/__tests__/recovery-stale-issue-lock-sweep.test.ts index 1fd20c88cd..b2068ca567 100644 --- a/server/src/__tests__/recovery-stale-issue-lock-sweep.test.ts +++ b/server/src/__tests__/recovery-stale-issue-lock-sweep.test.ts @@ -110,6 +110,31 @@ describeEmbeddedPostgres("recovery sweepStaleIssueLocks", () => { return { companyId, agentId, failedRunId, runningRunId }; } + it.each([true, false])("preserves explicit retry ownership when finalization wins the sweep write (settler: %s)", async withSettler => { + const { companyId, agentId, runningRunId, failedRunId } = await seed(); + const issueId = randomUUID(); + await db.insert(issues).values({ id: issueId, companyId, title: "Retry finalization race", + status: "in_progress", assigneeAgentId: agentId, executionRunId: runningRunId }); + await db.update(heartbeatRuns).set({ processPid: 2_000_000_000, + contextSnapshot: { issueId, explicitUserContinuation: { previousRunId: failedRunId, commentId: randomUUID() } }, + }).where(eq(heartbeatRuns.id, runningRunId)); + const beforeOrphanedRunTerminalWrite = vi.fn(async () => { + await db.update(heartbeatRuns).set({ status: "timed_out", finishedAt: new Date() }) + .where(eq(heartbeatRuns.id, runningRunId)); + }); + const settleExplicitContinuationRetry = vi.fn(async () => {}); + const result = await recoveryService(db, { enqueueWakeup: vi.fn(), beforeOrphanedRunTerminalWrite, + ...(withSettler ? { settleExplicitContinuationRetry } : {}), + }).sweepStaleIssueLocks(); + expect(beforeOrphanedRunTerminalWrite).toHaveBeenCalledOnce(); + expect(result).toEqual({ cleared: 0, issueIds: [], terminalizedRunIds: [] }); + expect(settleExplicitContinuationRetry).toHaveBeenCalledTimes(withSettler ? 1 : 0); + if (withSettler) expect(settleExplicitContinuationRetry).toHaveBeenCalledWith(expect.objectContaining({ + id: runningRunId, status: "timed_out", + })); + expect((await db.select().from(issues).where(eq(issues.id, issueId)))[0]!.executionRunId).toBe(runningRunId); + }); + it("clears lock columns when checkoutRunId points at a terminal heartbeat run", async () => { const { companyId, agentId, failedRunId } = await seed(); const issueId = randomUUID(); diff --git a/server/src/modules/wake-queue/adapters/postgres.ts b/server/src/modules/wake-queue/adapters/postgres.ts index 57d256e70b..d2624b2512 100644 --- a/server/src/modules/wake-queue/adapters/postgres.ts +++ b/server/src/modules/wake-queue/adapters/postgres.ts @@ -1196,6 +1196,18 @@ export function createPostgresWakeQueueAdapter(db: Db, deps: WakeQueuePostgresAd const locked: LockedIssueExecution = { primaryIssue: toIssueSnapshot(issueRow), run: runSnapshot, recoveryOnly }; const result = await fn(locked, { host: buildHost(tx, deps), transaction: buildTransaction(tx, deps, db, run) }); + // Explicit retry admission requires the exact source's task claim. + // Preserve it only when the existing queue-first recovery policy + // requests that retry. The clear/restore stays inside this task lock. + if (run.runtimeMode === "legacy" && ["failed", "timed_out"].includes(run.status) && + run.contextSnapshot?.explicitUserContinuation && issueRow.executionRunId === run.id && + result.postCommitEffects.some(effect => effect.kind === "conversation_retry_requested" && effect.runId === run.id)) { + await tx.update(issues).set({ executionRunId: run.id, + executionAgentNameKey: issueRow.executionAgentNameKey, + executionLockedAt: issueRow.executionLockedAt, + }).where(and(eq(issues.companyId, input.companyId), eq(issues.id, issueRow.id), + sql`${issues.executionRunId} is null`)); + } return { ...result, run: runSnapshot }; }); }, diff --git a/server/src/services/execution-continuation.ts b/server/src/services/execution-continuation.ts index 1afbd5033b..acda541cd6 100644 --- a/server/src/services/execution-continuation.ts +++ b/server/src/services/execution-continuation.ts @@ -397,6 +397,15 @@ export async function buildExecutionContinuation(input: { priorRuns.some(run => run.id === wake.runId && run.retryOfRunId === failedRunId)) : rows.some(comment => comment.id === value.commentId && comment.authorType === "user" && + (!("commentUpdatedAt" in value || "commentBodyHash" in value || value.automaticRetry) || ( + value.commentUpdatedAt === comment.updatedAt.toISOString() && + value.commentBodyHash === createHash("sha256").update(comment.body).digest("hex") + )) && + (!value.automaticRetry || ( + object(value.automaticRetry).sourceRunId === explicitUserSource && + comment.body.trim().length > 0 && issue.executionRunId === input.runId && + priorRuns.some(run => run.id === value.runId && run.retryOfRunId === explicitUserSource) + )) && (value.queuedCommentInterruptId ? interruptQueues.some(queue => queue.id === value.queuedCommentInterruptId && queue.runId === value.runId && diff --git a/server/src/services/explicit-continuation-retry-claim.ts b/server/src/services/explicit-continuation-retry-claim.ts new file mode 100644 index 0000000000..9c87d38839 --- /dev/null +++ b/server/src/services/explicit-continuation-retry-claim.ts @@ -0,0 +1,14 @@ +import type { heartbeatRuns, issues } from "@paperclipai/db"; + +/** + * This claim must be settled by queue-first recovery, not generic terminal-lock + * cleanup. It is pending intent only; admission still validates its authority. + */ +export function isExplicitContinuationRetryClaim( + issue: Pick, + run: Pick | null | undefined, +) { + return Boolean(run && issue.executionRunId === run.id && run.companyId === issue.companyId && + run.runtimeMode === "legacy" && ["failed", "timed_out"].includes(run.status) && + run.contextSnapshot?.issueId === issue.id && run.contextSnapshot.explicitUserContinuation); +} diff --git a/server/src/services/explicit-native-continuation.test.ts b/server/src/services/explicit-native-continuation.test.ts index 3ecd4596b4..013cab6b90 100644 --- a/server/src/services/explicit-native-continuation.test.ts +++ b/server/src/services/explicit-native-continuation.test.ts @@ -10,7 +10,7 @@ import { createHash, randomUUID } from "node:crypto"; import { mkdtemp, mkdir, writeFile, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { and, eq } from "drizzle-orm"; +import { and, eq, sql } from "drizzle-orm"; import { beforeAll, afterAll, describe, it, expect, vi } from "vitest"; import { approvals, issueApprovals, issueThreadInteractions, chatConversations, chatEndpoints, toolApplications, toolConnections, @@ -19,6 +19,8 @@ import { } from "@paperclipai/db"; import { startEmbeddedPostgresTestDatabase, getEmbeddedPostgresTestSupport } from "../__tests__/helpers/embedded-postgres.js"; import { admitExplicitNativeContinuation } from "./explicit-native-continuation.js"; +import { adapterExecutionControls, createAdapterExecutionControl } from "./adapter-execution-control.js"; +import { CONVERSATION_CONTINUATION_POLICY } from "./conversation-continuation.js"; import { buildExecutionContinuation } from "./execution-continuation.js"; import { heartbeatService, persistHeartbeatRunProcessMetadata, type HeartbeatEnvironmentRuntime } from "./heartbeat.js"; import { getExecutionBlocker } from "./execution-blocker.js"; @@ -386,6 +388,197 @@ const support = await getEmbeddedPostgresTestSupport(); return result; }); + async function seedTimedOutExplicitTurn() { + const f = await seed(); + const explicitUserContinuation = await admit(f); + expect(explicitUserContinuation).not.toBeNull(); + const [receipt] = await db.select().from(issueRecoveryActions) + .where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + await db.update(agents).set({ adapterType: "claude_local" }).where(eq(agents.id, f.agentId)); + const now = new Date(); + const [parent] = await db.update(heartbeatRuns).set({ runtimeMode: "legacy", status: "timed_out", + errorCode: "timeout", startedAt: now, finishedAt: now, + resultJson: { conversationContinuation: CONVERSATION_CONTINUATION_POLICY }, + contextSnapshot: { issueId: f.issueId, wakeReason: "issue_commented", wakeCommentId: f.commentId, + previousRunId: f.sourceRunId, explicitUserContinuation }, + }).where(eq(heartbeatRuns.id, f.successorRunId)).returning(); + await db.update(issues).set({ status: "in_progress", executionRunId: parent.id }) + .where(eq(issues.id, f.issueId)); + return { ...f, parent, receipt, now }; + } + const dispatchExplicitRetry = (f: Fixture, run: typeof heartbeatRuns.$inferSelect) => buildExecutionContinuation({ + db, companyId: f.companyId, issueId: f.issueId, agentId: f.agentId, runId: run.id, + context: run.contextSnapshot!, summary: null, exposeLowTrustRaw: false, + }); + + it("re-admits an explicit user turn's timeout retry with its own durable authorization", async () => { + const f = await seedTimedOutExplicitTurn(); + await expect(dispatchExplicitRetry(f, f.parent)).resolves.toMatchObject({ interruptedRunId: f.sourceRunId }); + const heartbeat = heartbeatService(db); + const [scheduled, competing] = await Promise.all([ + heartbeat.scheduleBoundedRetry(f.parent.id, { now: f.now, delayMs: 1000 }), + heartbeat.scheduleBoundedRetry(f.parent.id, { now: f.now, delayMs: 1000 }), + ]); + expect(scheduled.outcome).toBe("scheduled"); + if (scheduled.outcome !== "scheduled") return; + expect(competing).toMatchObject({ outcome: "scheduled", run: { id: scheduled.run.id } }); + expect(scheduled.run.contextSnapshot?.explicitUserContinuation).toEqual({ + previousRunId: f.parent.id, commentId: f.commentId, + }); + // Dispatch validates the new run, not a borrowed parent receipt. + await expect(dispatchExplicitRetry(f, scheduled.run)).resolves.toMatchObject({ interruptedRunId: f.parent.id }); + const receipts = await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + expect(receipts).toHaveLength(2); + expect(receipts.find(row => row.id === f.receipt.id)?.evidence).toEqual(f.receipt.evidence); + expect(receipts.find(row => row.id !== f.receipt.id)?.evidence.explicitUserContinuation).toMatchObject({ + runId: scheduled.run.id, previousRunId: f.parent.id, actorId: "board", commentId: f.commentId, + automaticRetry: { sourceRunId: f.parent.id, sourceAuthorizationId: f.receipt.id }, + }); + expect(await heartbeat.scheduleBoundedRetry(f.parent.id, { now: f.now })).toMatchObject({ + outcome: "scheduled", run: { id: scheduled.run.id }, + }); + expect(await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, f.issueId))).toHaveLength(2); + // Removing the new receipt reproduces the original dispatch failure. + await db.delete(issueRecoveryActions).where(eq(issueRecoveryActions.id, receipts.find(row => row.id !== f.receipt.id)!.id)); + await expect(dispatchExplicitRetry(f, scheduled.run)).rejects.toThrow("continuation_user_authorization_missing"); + }); + + it("waits for the parent adapter controller even when no environment lease remains", async () => { + const f = await seedTimedOutExplicitTurn(); + const heartbeat = heartbeatService(db); + const control = createAdapterExecutionControl(); + adapterExecutionControls.set(f.parent.id, control); + try { + expect(await heartbeat.scheduleBoundedRetry(f.parent.id)).toMatchObject({ outcome: "not_scheduled" }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, f.parent.id))).toHaveLength(0); + expect(await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, f.issueId))).toHaveLength(1); + } finally { + adapterExecutionControls.delete(f.parent.id); + control.finish(); + } + expect(await heartbeat.scheduleBoundedRetry(f.parent.id)).toMatchObject({ outcome: "scheduled" }); + }); + + it("binds each bounded explicit retry to its immediate parent without replacing prior receipts", async () => { + const f = await seedTimedOutExplicitTurn(); + const heartbeat = heartbeatService(db); + const first = await heartbeat.scheduleBoundedRetry(f.parent.id, { now: f.now }); + expect(first.outcome).toBe("scheduled"); + if (first.outcome !== "scheduled") return; + const nextNow = new Date(f.now.getTime() + 2000); + await db.update(heartbeatRuns).set({ status: "timed_out", errorCode: "timeout", startedAt: f.now, + finishedAt: nextNow }).where(eq(heartbeatRuns.id, first.run.id)); + const second = await heartbeat.scheduleBoundedRetry(first.run.id, { now: nextNow }); + expect(second.outcome).toBe("scheduled"); + if (second.outcome !== "scheduled") return; + expect(second.run).toMatchObject({ retryOfRunId: first.run.id, scheduledRetryAttempt: 2 }); + await expect(dispatchExplicitRetry(f, second.run)).resolves.toMatchObject({ interruptedRunId: first.run.id }); + const receipts = await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + expect(receipts).toHaveLength(3); + expect(receipts.find(row => row.id === f.receipt.id)?.evidence).toEqual(f.receipt.evidence); + await db.update(heartbeatRuns).set({ status: "timed_out", errorCode: "timeout", startedAt: nextNow, + finishedAt: nextNow }).where(eq(heartbeatRuns.id, second.run.id)); + expect(await heartbeat.scheduleBoundedRetry(second.run.id, { now: nextNow })).toMatchObject({ outcome: "retry_exhausted" }); + expect(await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, f.issueId))).toHaveLength(3); + }); + + it("keeps sub-millisecond source ordering when authorizing a timeout retry", async () => { + const f = await seedTimedOutExplicitTurn(); + await db.execute(sql`update heartbeat_runs set created_at = ${f.now.toISOString()}::timestamptz + + case when id = ${f.sourceRunId} then interval '100 microseconds' else interval '500 microseconds' end + where company_id = ${f.companyId} and id in (${f.sourceRunId}, ${f.parent.id})`); + const result = await heartbeatService(db).scheduleBoundedRetry(f.parent.id, { now: f.now }); + expect(result.outcome).toBe("scheduled"); + }); + + it.each(["max_turns_continuation", "interaction_continuation_infra_retry", "workspace_busy"])( + "does not lend explicit user authority to the separate %s retry contract", async retryReason => { + const f = await seedTimedOutExplicitTurn(); + const result = await heartbeatService(db).scheduleBoundedRetry(f.parent.id, { now: f.now, retryReason }); + expect(result.outcome).not.toBe("scheduled"); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, f.parent.id))).toHaveLength(0); + }, + ); + + it.each(["deleted", "edited", "same_millisecond_edit", "empty", "actor", "run_authored", "legacy_receipt", "receipt_run", "receipt_company", "source", "task", "reassigned", "done", "superseded", "native_superseding", "interrupted_intent"])( + "refuses to re-admit an automatic explicit retry after %s changes", async kind => { + const f = await seedTimedOutExplicitTurn(); + if (["deleted", "edited", "same_millisecond_edit", "empty", "actor", "run_authored"].includes(kind)) await db.update(issueComments).set({ + ...(kind === "deleted" ? { deletedAt: new Date() } : {}), + ...(kind === "edited" ? { body: "Different instructions", updatedAt: new Date(f.now.getTime() + 1000) } : {}), + ...(kind === "same_millisecond_edit" ? { body: "Different instructions within the recorded millisecond" } : {}), + ...(kind === "empty" ? { body: " " } : {}), + ...(kind === "actor" ? { authorUserId: "another-user" } : {}), + ...(kind === "run_authored" ? { createdByRunId: f.parent.id } : {}), + }).where(eq(issueComments.id, f.commentId)); + if (["receipt_run", "interrupted_intent"].includes(kind)) await db.update(issueRecoveryActions).set({ evidence: { + ...f.receipt.evidence, explicitUserContinuation: { + ...(f.receipt.evidence.explicitUserContinuation as Record), + ...(kind === "receipt_run" ? { runId: randomUUID() } : { queuedCommentInterruptId: randomUUID() }), + }, + } }).where(eq(issueRecoveryActions.id, f.receipt.id)); + if (kind === "receipt_company") await db.update(issueRecoveryActions).set({ companyId: (await seed()).companyId }) + .where(eq(issueRecoveryActions.id, f.receipt.id)); + if (kind === "legacy_receipt") { + const { commentBodyHash: _bodyHash, commentUpdatedAt: _updatedAt, ...legacy } = + f.receipt.evidence.explicitUserContinuation as Record; + await db.update(issueRecoveryActions).set({ evidence: { ...f.receipt.evidence, + explicitUserContinuation: legacy } }).where(eq(issueRecoveryActions.id, f.receipt.id)); + } + if (kind === "source" || kind === "task") await db.update(heartbeatRuns).set({ contextSnapshot: { + ...f.parent.contextSnapshot, + ...(kind === "task" ? { issueId: randomUUID() } : { explicitUserContinuation: { previousRunId: randomUUID(), commentId: f.commentId } }), + } }).where(eq(heartbeatRuns.id, f.parent.id)); + if (kind === "reassigned" || kind === "done") await db.update(issues).set(kind === "done" + ? { status: "done" } : { assigneeAgentId: null }).where(eq(issues.id, f.issueId)); + if (kind === "superseded") await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, + status: "succeeded", contextSnapshot: { issueId: f.issueId }, createdAt: new Date(f.now.getTime() + 1000) }); + if (kind === "native_superseding") await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, + status: "running", runtimeMode: "native", nativeIssueId: f.issueId, contextSnapshot: {} }); + const result = await heartbeatService(db).scheduleBoundedRetry(f.parent.id, { now: f.now }); + expect(result.outcome).not.toBe("scheduled"); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, f.parent.id))).toHaveLength(0); + expect(await db.select().from(issueRecoveryActions).where(and(eq(issueRecoveryActions.sourceIssueId, f.issueId), + eq(issueRecoveryActions.cause, "explicit_user_continuation_retry")))).toHaveLength(0); + }, + ); + + it("binds the initial user-message admission before any automatic retry", async () => { + const f = await seedTimedOutExplicitTurn(); + await expect(dispatchExplicitRetry(f, f.parent)).resolves.toMatchObject({ interruptedRunId: f.sourceRunId }); + await db.update(issueComments).set({ body: "Changed after the original admission" }) + .where(eq(issueComments.id, f.commentId)); + await expect(dispatchExplicitRetry(f, f.parent)).rejects.toThrow("continuation_user_authorization_missing"); + }); + + it.each(["deleted", "edited", "same_millisecond_edit", "actor", "receipt_run", "receipt_source", "retry_parent", "execution_owner"])( + "revalidates the successor's exact authorization at dispatch: %s", async kind => { + const f = await seedTimedOutExplicitTurn(); + const scheduled = await heartbeatService(db).scheduleBoundedRetry(f.parent.id, { now: f.now }); + expect(scheduled.outcome).toBe("scheduled"); + if (scheduled.outcome !== "scheduled") return; + if (["deleted", "edited", "same_millisecond_edit", "actor"].includes(kind)) await db.update(issueComments).set({ + ...(kind === "deleted" ? { deletedAt: new Date() } : {}), + ...(kind === "edited" ? { body: "Changed after scheduling", updatedAt: new Date(f.now.getTime() + 1000) } : {}), + ...(kind === "same_millisecond_edit" ? { body: "Changed within the recorded millisecond" } : {}), + ...(kind === "actor" ? { authorUserId: "another-user" } : {}), + }).where(eq(issueComments.id, f.commentId)); + if (kind === "receipt_run" || kind === "receipt_source") { + const [receipt] = await db.select().from(issueRecoveryActions).where(and( + eq(issueRecoveryActions.sourceIssueId, f.issueId), eq(issueRecoveryActions.cause, "explicit_user_continuation_retry"))); + await db.update(issueRecoveryActions).set({ evidence: { ...receipt.evidence, explicitUserContinuation: { + ...(receipt.evidence.explicitUserContinuation as Record), + ...(kind === "receipt_run" ? { runId: f.parent.id } : { previousRunId: f.sourceRunId }), + } } }).where(eq(issueRecoveryActions.id, receipt.id)); + } + if (kind === "retry_parent") await db.update(heartbeatRuns).set({ retryOfRunId: f.sourceRunId }) + .where(eq(heartbeatRuns.id, scheduled.run.id)); + if (kind === "execution_owner") await db.update(issues).set({ executionRunId: f.parent.id }) + .where(eq(issues.id, f.issueId)); + await expect(dispatchExplicitRetry(f, scheduled.run)).rejects.toThrow("continuation_user_authorization_missing"); + }, + ); + it.each(["active", "removed", "unavailable", "paused", "disabled"])("shows usable recovery guidance for external chat cancellations: %s", async kind => { const f = await seed(); await db.delete(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); diff --git a/server/src/services/explicit-native-continuation.ts b/server/src/services/explicit-native-continuation.ts index bf1f66530f..bbfd4bb2c9 100644 --- a/server/src/services/explicit-native-continuation.ts +++ b/server/src/services/explicit-native-continuation.ts @@ -1,3 +1,4 @@ +import { createHash } from "node:crypto"; import { appendHeartbeatRunEvent } from "./heartbeat-run-events.js"; import { canContinueCancelledRun } from "./run-cancellation.js"; import { readQueuedInteractionResponse } from "./queued-interaction-response.js"; @@ -292,6 +293,8 @@ export async function admitExplicitNativeContinuation(input: { }); } const authorization = { actorId, commentId, ...(response ? { interactionId: response.source.interactionId } : {}), ...(retry ? { failedRunId: input.failedRunId } : {}), + ...(comment ? { commentUpdatedAt: comment.updatedAt.toISOString(), + commentBodyHash: createHash("sha256").update(comment.body).digest("hex") } : {}), ...(queuedInterrupt ? { queuedCommentInterruptId: input.queuedCommentInterruptId } : {}), ...(queuedRequest ? { queuedCommentRequestId: input.queuedCommentRequestId } : {}), runId: input.successorRunId, previousRunId: previous.id, recordedAt: new Date().toISOString() }; @@ -329,3 +332,91 @@ export async function admitExplicitNativeContinuation(input: { }); return { previousRunId: previous.id, commentId, ...(retry ? { failedRunId: input.failedRunId! } : {}) }; } + +/** Re-admit a bounded automatic retry without lending it its parent's receipt. + * The caller holds the task lock and creates the successor in this transaction. + * Only unchanged user messages are renewable; Retry/Interrupt intents keep their + * separate one-run contracts. The scheduler still owns retry and cleanup policy. + */ +export async function admitExplicitContinuationRetry(input: { + db: Db; companyId: string; issueId: string; agentId: string; + parentRunId: string; successorRunId: string; now: Date; +}): Promise<{ previousRunId: string; commentId: string } | null> { + const { db, companyId, issueId, agentId } = input; + const [task] = await db.select().from(issues).where(and( + eq(issues.companyId, companyId), eq(issues.id, issueId), + )).for("update"); + const [parent] = await db.select().from(heartbeatRuns).where(and( + eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.id, input.parentRunId), + eq(heartbeatRuns.agentId, agentId), + )).for("update"); + if (!task || task.assigneeAgentId !== agentId || task.executionRunId !== input.parentRunId || + ["done", "cancelled"].includes(task.status) || !parent || + !["failed", "timed_out"].includes(parent.status) || !parent.finishedAt || + parent.runtimeMode !== "legacy" || parent.contextSnapshot?.issueId !== issueId || + (parent.nativeIssueId !== null && parent.nativeIssueId !== issueId) || + adapterExecutionControls.has(parent.id) || + await getExecutionBlocker(db, companyId, issueId)) return null; + const context = parent.contextSnapshot; + const explicit = context.explicitUserContinuation as Record | undefined; + if (!explicit || !z.string().guid().safeParse(explicit.previousRunId).success || + !z.string().guid().safeParse(explicit.commentId).success || explicit.failedRunId) return null; + const commentId = explicit.commentId as string; + const [comment] = await db.select().from(issueComments).where(and( + eq(issueComments.companyId, companyId), eq(issueComments.issueId, issueId), + eq(issueComments.id, commentId), eq(issueComments.authorType, "user"), + isNull(issueComments.createdByRunId), isNull(issueComments.deletedAt), + )); + if (!comment?.authorUserId || !comment.body.trim()) return null; + const receipts = await db.select().from(issueRecoveryActions).where(and( + eq(issueRecoveryActions.companyId, companyId), eq(issueRecoveryActions.sourceIssueId, issueId), + eq(issueRecoveryActions.status, "resolved"), + sql`${issueRecoveryActions.evidence}->'explicitUserContinuation'->>'runId' = ${parent.id}`, + )); + const receipt = receipts.find(row => { + const auth = row.evidence.explicitUserContinuation as Record | undefined; + return auth && auth.previousRunId === explicit.previousRunId && auth.commentId === commentId && + auth.actorId === comment.authorUserId && !auth.failedRunId && !auth.queuedCommentInterruptId && + auth.commentUpdatedAt === comment.updatedAt.toISOString() && + auth.commentBodyHash === createHash("sha256").update(comment.body).digest("hex"); + }); + if (!receipt) return null; + const [superseding] = await db.select({ id: heartbeatRuns.id }).from(heartbeatRuns).where(and( + eq(heartbeatRuns.companyId, companyId), + or(eq(heartbeatRuns.nativeIssueId, issueId), sql`${heartbeatRuns.contextSnapshot}->>'issueId' = ${issueId}`), + ne(heartbeatRuns.id, parent.id), ne(heartbeatRuns.id, input.successorRunId), + or(eq(heartbeatRuns.retryOfRunId, parent.id), + // Keep PostgreSQL's timestamp precision: Date would round the parent + // down and mistake an older run in the same millisecond for a successor. + sql`${heartbeatRuns.createdAt} > (select source.created_at from heartbeat_runs source + where source.company_id = ${companyId} and source.id = ${parent.id})`, + inArray(heartbeatRuns.status, ["queued", "running", "scheduled_retry"])), + )).limit(1); + if (superseding) return null; + // Reuse dispatch's complete source/actor/receipt validation on the parent. + // It proves only the parent; the new receipt below authorizes the successor. + try { + await buildExecutionContinuation({ db, companyId, issueId, agentId, runId: parent.id, + context, summary: null, exposeLowTrustRaw: false }); + } catch (error) { + if (error instanceof Error && ["continuation_user_authorization_missing", + "continuation_source_context_missing", "continuation_task_ownership_changed"].includes(error.message)) return null; + throw error; + } + const continuation = { previousRunId: parent.id, commentId }; + await db.insert(issueRecoveryActions).values({ + companyId, sourceIssueId: issueId, kind: "active_run_watchdog", + cause: "explicit_user_continuation_retry", fingerprint: input.successorRunId, + status: "resolved", outcome: "retry_authorized", resolvedAt: input.now, + nextAction: "The bounded retry has its own authorization for the unchanged user message.", + evidence: { explicitUserContinuation: { + ...continuation, actorId: comment.authorUserId, runId: input.successorRunId, + commentUpdatedAt: comment.updatedAt.toISOString(), + commentBodyHash: createHash("sha256").update(comment.body).digest("hex"), + recordedAt: input.now.toISOString(), automaticRetry: { + sourceRunId: parent.id, sourceAuthorizationId: receipt.id, + }, + } }, + }); + return continuation; +} diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 4c7f3a281c..1c5cf0518f 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -32,7 +32,7 @@ import { legacyControllerBootId, legacyControllerClaim, renewLegacyControllerLea import { completeTerminatedRemoteNativeSessionCleanup } from "../vendor/paperclip-runner/index.js"; import { hasRemoteTerminationReceipt, remoteExecutionHasStopped, remoteTerminationReceipt, stoppedRemoteCleanupScopes } from "./remote-execution-termination.js"; import { applyConnectorSkills, prepareConnectorSkillDelivery, resolveConnectorAssignments } from "./connector-runtime.js"; -import { admitExplicitNativeContinuation, undeliveredLegacyUserCommentIds } from "./explicit-native-continuation.js"; +import { admitExplicitNativeContinuation, admitExplicitContinuationRetry, undeliveredLegacyUserCommentIds } from "./explicit-native-continuation.js"; import { canRetryStoppedRun, isCancelledNativeStartup } from "./cancelled-native-startup.js"; import { connectionIntentService } from "./connection-intents.js"; import { managedAiSessionFingerprintConfig, prepareManagedAiRuntime, assertManagedAiProjectAuth, stripAiAuthBindings, isAiConnectionBusy, AI_AUTH_ENV_KEYS } from "./ai-connection-runtime.js"; @@ -9649,6 +9649,7 @@ export function heartbeatService( const recovery = recoveryService(db, { enqueueWakeup, liveRunExecutions, + settleExplicitContinuationRetry: releaseIssueExecutionAndPromote, scheduleRecoveryRetry: async (runId) => { const [run] = await db .select() @@ -9866,10 +9867,19 @@ export function heartbeatService( )); const agent = source ? await getAgent(source.agentId) : null; if (source && agent && agent.companyId === source.companyId) { - await scheduleBoundedRetryForRun(source, agent, effect.reviewParticipant ? { + const retry = await scheduleBoundedRetryForRun(source, agent, effect.reviewParticipant ? { retryReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, wakeReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_WAKE_REASON, } : undefined); + const issueId = readNonEmptyString(source.contextSnapshot?.issueId); + if (retry.outcome !== "scheduled" && source.contextSnapshot?.explicitUserContinuation && issueId && + !adapterExecutionControls.has(source.id) && !(await getExecutionBlocker(db, source.companyId, issueId))) { + // Cleanup has settled, so exhaustion or revoked/missing authority + // must not retain a terminal claim. Do not request another retry. + const settled = await wakeQueue.releaseIssueExecution({ companyId: source.companyId, + runId: source.id, now: new Date(), suppressImmediateRecovery: true }); + await applyWakeQueuePostCommitEffects(settled.postCommitEffects); + } } } else if (effect.kind === "run_queued") { publishLiveEvent({ @@ -14347,7 +14357,9 @@ export function heartbeatService( presentationDecision?: RunPresentationDecision | null, ) { const contextSnapshot = parseObject(run.contextSnapshot); - if (readNonEmptyString(contextSnapshot.goalControlRequestId)) { + // The explicit receipt admitted one turn. A prose-only follow-up cannot + // reuse it or renew it under the separate transient-failure retry policy. + if (readNonEmptyString(contextSnapshot.goalControlRequestId) || contextSnapshot.explicitUserContinuation) { if (run.issueCommentStatus !== "not_applicable") { await patchRunIssueCommentStatus(run.id, { issueCommentStatus: "not_applicable", @@ -15702,6 +15714,7 @@ export function heartbeatService( | "issue_cancelled" | "issue_terminal_status" | "issue_not_in_progress" + | "continuation_user_authorization_missing" | "issue_execution_lock_changed"; issueId: string | null; details: Record; @@ -15979,6 +15992,22 @@ export function heartbeatService( } } + const scheduledRunId = randomUUID(); + if (contextSnapshot.explicitUserContinuation) { + const continuation = issueId && retryReason === "transient_failure" ? await admitExplicitContinuationRetry({ + db: tx as unknown as Db, companyId: run.companyId, issueId, agentId: run.agentId, + parentRunId: run.id, successorRunId: scheduledRunId, now, + }) : null; + if (!continuation) return { + outcome: "not_scheduled", issueId, + errorCode: "continuation_user_authorization_missing", + reason: "The automatic retry could not revalidate the original user continuation.", + details: {}, + }; + retryContextSnapshot.explicitUserContinuation = continuation; + retryContextSnapshot.previousRunId = continuation.previousRunId; + } + const wakeupRequest = await tx .insert(agentWakeupRequests) .values({ @@ -16029,6 +16058,7 @@ export function heartbeatService( const scheduledRun = await tx .insert(heartbeatRuns) .values({ + id: scheduledRunId, companyId: run.companyId, agentId: run.agentId, invocationSource: "automation", @@ -26767,6 +26797,15 @@ export function heartbeatService( } // Terminalization precedes lease and adapter cleanup. Only now is the // owner gone; retry pending input for ordinary completions as well as Stop. + if (latestRun?.runtimeMode === "legacy" && ["failed", "timed_out"].includes(latestRun.status) && + latestRun.contextSnapshot?.explicitUserContinuation) { + // Re-run the same queue-first recovery decision after cleanup. Its + // earlier retry request could not authorize work while this executor + // still held its controller or environment lease. + await releaseIssueExecutionAndPromote(latestRun).catch(err => { + logger.error({ err, runId: run.id }, "failed to settle explicit continuation after cleanup"); + }); + } if (latestRun?.runtimeMode === "legacy" && isHeartbeatRunTerminalStatus(latestRun.status)) { const [pending] = await db.select({ id: agentWakeupRequests.id, payload: agentWakeupRequests.payload }).from(agentWakeupRequests).where(and( eq(agentWakeupRequests.companyId, run.companyId), eq(agentWakeupRequests.agentId, run.agentId), diff --git a/server/src/services/issues.ts b/server/src/services/issues.ts index fbf9868285..839b2945b4 100644 --- a/server/src/services/issues.ts +++ b/server/src/services/issues.ts @@ -9,6 +9,7 @@ import { executionProjectionsForRuns } from "./execution-projection.js"; import type { ExecutionProjection } from "@paperclipai/shared"; import { Buffer } from "node:buffer"; import { createHash, randomUUID } from "node:crypto"; +import { isExplicitContinuationRetryClaim } from "./explicit-continuation-retry-claim.js"; import { markdownToPlainText, parseMarkdown } from "chat"; import { and, @@ -7497,6 +7498,7 @@ export function issueService(db: Db) { const lockedIssue = await tx .select({ id: issues.id, + companyId: issues.companyId, status: issues.status, assigneeAgentId: issues.assigneeAgentId, checkoutRunId: issues.checkoutRunId, @@ -7528,7 +7530,7 @@ export function issueService(db: Db) { ]); const [existingRun, actorRun] = await Promise.all([ tx - .select({ status: heartbeatRuns.status }) + .select() .from(heartbeatRuns) .where(eq(heartbeatRuns.id, input.expectedCheckoutRunId)) .then((rows) => rows[0] ?? null), @@ -7540,6 +7542,17 @@ export function issueService(db: Db) { ]); const stale = !existingRun || TERMINAL_HEARTBEAT_RUN_STATUSES.has(existingRun.status); + if (isExplicitContinuationRetryClaim(lockedIssue, existingRun)) { + return { adopted: null, latest: lockedIssue }; + } + if (lockedIssue.executionRunId && lockedIssue.executionRunId !== input.expectedCheckoutRunId) { + const executionRun = await tx.select().from(heartbeatRuns) + .where(eq(heartbeatRuns.id, lockedIssue.executionRunId)).for("update") + .then(rows => rows[0] ?? null); + if (isExplicitContinuationRetryClaim(lockedIssue, executionRun)) { + return { adopted: null, latest: lockedIssue }; + } + } const actorLive = actorRun && !TERMINAL_HEARTBEAT_RUN_STATUSES.has(actorRun.status); if (!stale || !actorLive) { @@ -7649,7 +7662,7 @@ export function issueService(db: Db) { sql`select ${issues.id} from ${issues} where ${issues.id} = ${issueId} for update`, ); const issue = await tx - .select({ executionRunId: issues.executionRunId }) + .select({ id: issues.id, companyId: issues.companyId, executionRunId: issues.executionRunId }) .from(issues) .where(eq(issues.id, issueId)) .then((rows) => rows[0] ?? null); @@ -7659,11 +7672,12 @@ export function issueService(db: Db) { sql`select ${heartbeatRuns.id} from ${heartbeatRuns} where ${heartbeatRuns.id} = ${issue.executionRunId} for update`, ); const run = await tx - .select({ status: heartbeatRuns.status }) + .select() .from(heartbeatRuns) .where(eq(heartbeatRuns.id, issue.executionRunId)) .then((rows) => rows[0] ?? null); if (run && !TERMINAL_HEARTBEAT_RUN_STATUSES.has(run.status)) return false; + if (isExplicitContinuationRetryClaim(issue, run)) return false; const updated = await tx .update(issues) @@ -7689,8 +7703,7 @@ export function issueService(db: Db) { // Symmetric to clearExecutionRunIfTerminal. Clears checkoutRunId (and the // bundled execution lock cols) when the row's checkoutRunId points at a // heartbeat run that is terminal or no longer exists. No assignee/status - // precondition: a terminal run holds no real claim regardless of who is - // assigned or what status the issue is currently in. + // precondition. Explicit retry claims remain owned by queue-first settlement. async function clearCheckoutRunIfTerminal(issueId: string): Promise { return db.transaction(async (tx) => { await tx.execute( @@ -7698,6 +7711,8 @@ export function issueService(db: Db) { ); const issue = await tx .select({ + id: issues.id, + companyId: issues.companyId, checkoutRunId: issues.checkoutRunId, executionRunId: issues.executionRunId, }) @@ -7710,11 +7725,12 @@ export function issueService(db: Db) { sql`select ${heartbeatRuns.id} from ${heartbeatRuns} where ${heartbeatRuns.id} = ${issue.checkoutRunId} for update`, ); const run = await tx - .select({ status: heartbeatRuns.status }) + .select() .from(heartbeatRuns) .where(eq(heartbeatRuns.id, issue.checkoutRunId)) .then((rows) => rows[0] ?? null); if (run && !TERMINAL_HEARTBEAT_RUN_STATUSES.has(run.status)) return false; + if (isExplicitContinuationRetryClaim(issue, run)) return false; if ( issue.executionRunId && @@ -7724,7 +7740,7 @@ export function issueService(db: Db) { sql`select ${heartbeatRuns.id} from ${heartbeatRuns} where ${heartbeatRuns.id} = ${issue.executionRunId} for update`, ); const executionRun = await tx - .select({ status: heartbeatRuns.status }) + .select() .from(heartbeatRuns) .where(eq(heartbeatRuns.id, issue.executionRunId)) .then((rows) => rows[0] ?? null); @@ -7733,6 +7749,7 @@ export function issueService(db: Db) { !TERMINAL_HEARTBEAT_RUN_STATUSES.has(executionRun.status) ) return false; + if (isExplicitContinuationRetryClaim(issue, executionRun)) return false; } const updated = await tx @@ -11600,10 +11617,10 @@ export function issueService(db: Db) { current.executionRunId !== checkoutRunId && (current.assigneeAgentId === agentId || current.assigneeAgentId == null) ) { - const stale = await isTerminalOrMissingHeartbeatRun( - current.executionRunId, - ); - if (stale) { + const executionRun = await db.select().from(heartbeatRuns) + .where(eq(heartbeatRuns.id, current.executionRunId)).then(rows => rows[0] ?? null); + const stale = !executionRun || TERMINAL_HEARTBEAT_RUN_STATUSES.has(executionRun.status); + if (stale && !isExplicitContinuationRetryClaim({ ...current, companyId: issueCompany.companyId }, executionRun)) { const now = new Date(); const adoptionSet: Record = { assigneeAgentId: agentId, diff --git a/server/src/services/recovery/service.ts b/server/src/services/recovery/service.ts index 99ed604dee..cfda20e5bc 100644 --- a/server/src/services/recovery/service.ts +++ b/server/src/services/recovery/service.ts @@ -3,6 +3,7 @@ import { isNativeWorkspaceExportRepairCause } from "@paperclipai/shared"; import { settleSlackConversation } from "../slack-conversation-lifecycle.js"; import { externalConversationStateSql } from "../slack-conversation-state.js"; import { executionRetryAccounting } from "../execution-recovery-attempt.js"; +import { isExplicitContinuationRetryClaim } from "../explicit-continuation-retry-claim.js"; import { decideLegacyContinuation, legacyDispositionEpisode, legacyDispositionFingerprint, LEGACY_DISPOSITION_REPAIR_INSTRUCTION, type LegacyDispositionEpisode, @@ -912,6 +913,8 @@ export function recoveryService( scheduleRecoveryRetry?: ( runId: string, ) => Promise; + /** Settle retained explicit retry claims through the queue-first release policy. */ + settleExplicitContinuationRetry?: (run: typeof heartbeatRuns.$inferSelect) => Promise; /** * Whether a failed or interrupted run has consumed every bounded * transient retry, so `scheduleRecoveryRetry` can no longer produce a @@ -6098,6 +6101,7 @@ export function recoveryService( : []; const runStatusById = new Map(); for (const row of runRows) runStatusById.set(row.id, row.status); + const runById = new Map(runRows.map(row => [row.id, row])); // Collect the runs that a non-terminal issue still references. Such a run is // the live run of an active issue. A different, terminal issue can also hold @@ -6159,6 +6163,18 @@ export function recoveryService( }; for (const issue of candidates) { + const originalOwner = issue.executionRunId ? runById.get(issue.executionRunId) : undefined; + // The pre-pass can lose its terminal write to the executor and observe a + // newer terminal status. Do not test claim ownership using the old status. + const owner = originalOwner ? { ...originalOwner, + status: runStatusById.get(originalOwner.id) ?? originalOwner.status } : undefined; + if (owner && isExplicitContinuationRetryClaim(issue, owner)) { + // A terminal row can still own cleanup and a pending bounded retry. + // Re-enter the same policy rather than erasing its exact-owner proof. + // That policy handles pending cleanup, newer input, and final denial. + await deps.settleExplicitContinuationRetry?.(owner); + continue; + } if ( !isCleanable(issue.checkoutRunId) || !isCleanable(issue.executionRunId)