diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index a4fa177040..344bfeffd4 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -410,6 +410,8 @@ The handshake failure code is distinct from a session-identity mismatch. A timeo An explicit recovery action is a typed liveness repair path for a source issue. It is the recovery primitive; the action can be rendered directly on the source issue or backed by a separate recovery issue when the repair needs its own work item. +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 existing guarded Retry action for failed or timed-out legacy conversation runs. Where the server supports an explicit new attempt after a stopped legacy conversation, the thread must not hide that action solely because the old run still has a recovery-needed projection. Native and process recovery holds, pending decisions, active execution, and other retry gates remain in force. When a gate hides Retry, the thread says the message is preserved instead of promising an unavailable action. This presentation change does not rewrite historical outcomes or certify prior actions. A valid recovery action must name: diff --git a/server/src/services/explicit-native-continuation.test.ts b/server/src/services/explicit-native-continuation.test.ts index 49b4708eee..e8bad1eb3e 100644 --- a/server/src/services/explicit-native-continuation.test.ts +++ b/server/src/services/explicit-native-continuation.test.ts @@ -1,7 +1,10 @@ import { appendHeartbeatRunEvent } from "./heartbeat-run-events.js"; import { recordNativeLocalProcessStop, hasNativeLocalProcessStop, PROCESS_START_REQUESTED } from "./native-local-process-stop.js"; import { remoteTerminationReceipt } from "./remote-execution-termination.js"; -import { randomUUID } from "node:crypto"; +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 { beforeAll, afterAll, describe, it, expect } from "vitest"; import { @@ -369,6 +372,85 @@ const support = await getEmbeddedPostgresTestSupport(); agentId: f.agentId, status: "queued", contextSnapshot: { issueId: f.issueId, previousRunId: result.previousRunId, forceFreshSession: true } }); return result; }); + + it.each(["suspended", "ready", "wrong_run", "wrong_thread", "active_provider", "pending_tool", "pending_output", "missing_state", "new_launch"])( + "recovers a historical run without process metadata only from exact suspended state (%s)", async kind => { + const f = await seed(); + const stateBase = await mkdtemp(join(tmpdir(), "historical-native-followup-")); + const previous = process.env.PAPERCLIP_RUNNER_STATE_DIR; + process.env.PAPERCLIP_RUNNER_STATE_DIR = stateBase; + try { + const nativeSessionId = randomUUID(), runnerInstanceId = randomUUID(); + const execution = { + schema: "paperclip.native-execution-input.v1", provider: { kind: "codex", model: null }, + binding: { companyId: f.companyId, issueId: f.issueId, agentId: f.agentId, runId: f.sourceRunId, executionWorkspaceId: "workspace" }, + task: { identifier: "TEST", title: "Continue", description: null, prompt: "Continue", workMode: "standard" }, + workspace: { cwd: stateBase, repoUrl: null, repoRef: null, branchName: null }, + session: { normalizedSessionId: nativeSessionId, driverKind: "codex_app_server", protocolVersion: 1, lifecyclePolicy: { mode: "per_turn", idleTimeoutMs: null } }, + completionContract: { id: "contract", sha256: "sha", schemaVersion: "paperclip.completion-contract.v1", + contract: { revision: "1", objective: "Continue", criteria: [{ id: "objective", requirement: "Continue" }] } }, + interactionResponses: [], credentialBindings: [], + }; + await db.update(heartbeatRuns).set({ processPid: null, nativeSessionId, runnerInstanceId, + errorCode: "native_runner_process_exited", runnerProfileJson: { nativeExecutionInput: execution, + sessionCheckpoint: { sessionId: "exact-thread", providerSessionId: "backend-account" } }, + }).where(eq(heartbeatRuns.id, f.sourceRunId)); + const canonical = (value: unknown): string => value && typeof value === "object" && !Array.isArray(value) + ? `{${Object.entries(value).sort(([a], [b]) => a.localeCompare(b)).map(([key, entry]) => `${JSON.stringify(key)}:${canonical(entry)}`).join(",")}}` + : JSON.stringify(value); + const root = join(stateBase, createHash("sha256").update(canonical({ + schema: "paperclip.native-session-scope.v2", companyId: f.companyId, agentId: f.agentId, + workspace: { kind: "managed", executionWorkspaceId: "workspace" }, + provider: { driverKind: "codex_app_server", identity: { kind: "codex" } }, normalizedSessionId: nativeSessionId, + })).digest("hex")); + if (kind !== "missing_state") { + await mkdir(join(root, "control-plane"), { recursive: true }); + await mkdir(join(root, "runner"), { recursive: true }); + const identity = { runId: kind === "wrong_run" ? randomUUID() : f.sourceRunId, runnerInstanceId, + normalizedSessionId: nativeSessionId, environmentLeaseId: "workspace" }; + await writeFile(join(root, "control-plane/control-plane-state.json"), JSON.stringify({ schema: "paperclip.runner.durable.control-plane-state.v1", identity })); + await writeFile(join(root, "runner/runner-state.json"), JSON.stringify({ schema: "paperclip.runner.durable.state.v1", + ...identity, lifecycle: kind === "ready" ? "ready" : "suspended", outbox: kind === "pending_output" ? [{}] : [] })); + await writeFile(join(root, "runner/codex-provider-state.json"), JSON.stringify({ + schema: "paperclip.runner.codex-provider-state.v1", lifecycle: "prepared", + threadId: kind === "wrong_thread" ? "another-thread" : "exact-thread", providerSessionId: "backend-account", + activeProviderTurnId: kind === "active_provider" ? "unfinished-turn" : null, ambiguousTurnStartPending: false, + config: { provider: "codex", driver: "codex_app_server" }, pendingEvents: [], queuedEvents: [], + toolBridge: { pending: kind === "pending_tool" ? { call: {} } : {} }, activeProviderResultFingerprint: null, + })); + } + if (kind === "new_launch") await appendHeartbeatRunEvent(db, { companyId: f.companyId, runId: f.sourceRunId, + agentId: f.agentId, eventType: PROCESS_START_REQUESTED }); + if (kind !== "suspended") { + expect(await admit(f, true)).toBeNull(); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).not.toBeNull(); + return; + } + expect(await admit(f, true)).toMatchObject({ previousRunId: f.sourceRunId }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).not.toBeNull(); + // Exercise the real message admission path while keeping the provider slot occupied. + await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "running" }); + const heartbeat = heartbeatService(db); + for (let n = 0; n < 2; n++) await heartbeat.wakeup(f.agentId, { source: "automation", triggerDetail: "system", + reason: "issue_commented", requestedByActorType: "user", requestedByActorId: "board", + payload: { issueId: f.issueId, commentId: f.commentId }, + contextSnapshot: { issueId: f.issueId, wakeCommentId: f.commentId } }); + const successors = await db.select().from(heartbeatRuns).where(and(eq(heartbeatRuns.companyId, f.companyId), eq(heartbeatRuns.status, "queued"))); + expect(successors).toHaveLength(1); + expect(successors[0].contextSnapshot).toMatchObject({ previousRunId: f.sourceRunId, forceFreshSession: true, wakeCommentId: f.commentId }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toBeNull(); + const [coordinator] = await db.select().from(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); + expect(coordinator).toMatchObject({ phase: "terminal_failure", attempt: 3 }); + const [action] = await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + expect(action.evidence.automaticRecovery).toMatchObject({ actionOutcome: "unknown", replay: "explicit_user_continuation" }); + expect(await hasNativeLocalProcessStop(db, f.companyId, f.sourceRunId)).toBe(false); + } finally { + if (previous === undefined) delete process.env.PAPERCLIP_RUNNER_STATE_DIR; + else process.env.PAPERCLIP_RUNNER_STATE_DIR = previous; + await rm(stateBase, { recursive: true, force: true }); + } + }, + ); async function seedCancelledStartup() { const f = await seed(); await db.update(heartbeatRuns).set({ status: "cancelled", processPid: null, diff --git a/server/src/services/explicit-native-continuation.ts b/server/src/services/explicit-native-continuation.ts index 4be1f26581..fecb224ca0 100644 --- a/server/src/services/explicit-native-continuation.ts +++ b/server/src/services/explicit-native-continuation.ts @@ -1,5 +1,5 @@ import { isCancelledNativeStartup } from "./cancelled-native-startup.js"; -import { hasNativeLocalProcessStop } from "./native-local-process-stop.js"; +import { hasNativeLocalProcessStop, hasHistoricalSuspendedNativeSession } from "./native-local-process-stop.js"; import { completeTerminatedRemoteNativeSessionCleanup } from "../vendor/paperclip-runner/index.js"; import { hasRemoteTerminationReceipt, remoteLeaseCleanupScope } from "./remote-execution-termination.js"; import { z } from "zod"; @@ -212,7 +212,8 @@ export async function admitExplicitNativeContinuation(input: { if (!unusedAdmission && !cancelledStartup) { // A missing process identity is not evidence that a provider exited. if (!run.processPid && !run.processGroupId && - !await hasNativeLocalProcessStop(db, companyId, run.id)) return blocked("process_identity_missing", "The previous run has no verified stop record. Paperclip cannot start this message yet."); + !await hasNativeLocalProcessStop(db, companyId, run.id) && + !await hasHistoricalSuspendedNativeSession(db, run)) return blocked("process_identity_missing", "The previous run has no verified stop record. Paperclip cannot start this message yet."); if (run.processPid && !processStopped(run.processPid)) return blocked("process_running", "Waiting for the previous process to stop. Your message will start automatically."); if (run.processGroupId && !processStopped(-run.processGroupId)) return blocked("process_running", "Waiting for the previous process to stop. Your message will start automatically."); } diff --git a/server/src/services/native-local-process-stop.ts b/server/src/services/native-local-process-stop.ts index 3c82cddeb3..29cbafdf0b 100644 --- a/server/src/services/native-local-process-stop.ts +++ b/server/src/services/native-local-process-stop.ts @@ -59,6 +59,37 @@ export async function hasNativeLocalProcessStop(db: Db, companyId: string, runId return event?.eventType === LOCAL_PROCESS_STOPPED; } +/** Pre-receipt native runs can retain an exact suspended session after their + * mutable process fields were cleared. This is admission evidence for a new + * user turn only, never permission to replay the old run or infer its outcomes. + * The caller must hold the run/controller locks and verify local lease cleanup. + */ +export async function hasHistoricalSuspendedNativeSession(db: Db, run: typeof heartbeatRuns.$inferSelect) { + if (run.runtimeMode !== "native" || run.processPid || run.processGroupId || + !run.nativeSessionId || !run.runnerInstanceId || !run.nativeIssueId) return false; + const [modernProcessEvidence] = await db.select({ id: heartbeatRunEvents.id }).from(heartbeatRunEvents).where(and( + eq(heartbeatRunEvents.companyId, run.companyId), eq(heartbeatRunEvents.runId, run.id), + isNull(heartbeatRunEvents.sourceEventId), + inArray(heartbeatRunEvents.eventType, [PROCESS_START_REQUESTED, PROCESS_IDENTITY_RECORDED, LOCAL_PROCESS_STOPPED]), + )).limit(1); + // A newer launch invalidates an old stop receipt. Never bypass that fence + // with a suspended file that could belong to the earlier process generation. + if (modernProcessEvidence) return false; + const checkpoint = run.runnerProfileJson?.sessionCheckpoint as Record | undefined; + if (checkpoint?.providerSessionId != null && (typeof checkpoint.providerSessionId !== "string" || + !checkpoint.providerSessionId.trim())) return false; + const { nativeFailedRunRetryStateIsSafe } = await import("./native-runtime/native-session-executor.js"); + return nativeFailedRunRetryStateIsSafe({ + execution: run.runnerProfileJson?.nativeExecutionInput, + companyId: run.companyId, issueId: run.nativeIssueId, agentId: run.agentId, runId: run.id, + nativeSessionId: run.nativeSessionId, runnerInstanceId: run.runnerInstanceId, + processPid: null, processGroupId: null, + providerSessionId: typeof checkpoint?.sessionId === "string" ? checkpoint.sessionId : null, + providerBackendSessionId: typeof checkpoint?.providerSessionId === "string" ? checkpoint.providerSessionId : null, + recoveryMode: "exact_checkpoint_resume", allowVerifiedBackup: false, + }); +} + /** Recover the exact stopped identity after the mutable run fields were cleared. */ export async function readNativeLocalProcessStop(db: Db, companyId: string, runId: string) { const [event] = await db.select({ eventType: heartbeatRunEvents.eventType, payload: heartbeatRunEvents.payload })