diff --git a/doc/DEVELOPING.md b/doc/DEVELOPING.md index e145ac21fa..8927e66d41 100644 --- a/doc/DEVELOPING.md +++ b/doc/DEVELOPING.md @@ -230,7 +230,7 @@ at least one identity source. Supported-platform process probes fail explicitly instead of silently treating a live PID as either the original owner or a recycled process when identity cannot be established. -Use `--drain-required` only when the deploy intentionally requires the old terminate-and-retry behavior. Without that flag, the old server verifies that the marker targets its own PID, stops new scheduler work, waits for any queue-claim callback already in flight, snapshots currently running heartbeat run IDs and child PIDs, and skips the shutdown drain so eligible detached local-agent processes can keep running. ACP-backed local runs use server-owned stdio and cannot survive their parent server, so the old server instead persists their complete snapshot, changes the marker to `drainRequired` with `drainReason: "active_acp_run"`, and drains only those runs to queued retries. Detached CLI runs remain eligible for adoption during the same mixed restart. If an ACP process terminates but its terminal run update does not persist, startup classifies it as lost with reason `selective_drain_not_finalized` rather than treating the drain as successful. On startup the new server writes `$PAPERCLIP_HOME/instances/${PAPERCLIP_INSTANCE_ID:-default}/hot-restart-report.json` with `previousServerPid`, `newServerPid`, `previousServerVersion`, `newServerVersion`, `drainReason`, `adoptedRunIds`, `finalizedWhileDownRunIds`, `lostRunIds`, and per-run classifications before the normal orphan reaper runs. +Use `--drain-required` only when the deploy intentionally requires the old terminate-and-retry behavior. Without that flag, the old server verifies that the marker targets its own PID, stops new scheduler work, waits for any queue-claim callback already in flight, snapshots currently running heartbeat run IDs and child PIDs, and skips the shutdown drain so eligible detached local-agent processes can keep running. ACP-backed local runs use server-owned stdio and cannot survive their parent server, so the old server instead persists their complete snapshot, changes the marker to `drainRequired` with `drainReason: "active_acp_run"`, and drains only those runs to bounded conversation retries. These retries resume the prior session when compatible, otherwise carry the full task conversation into a fresh session. They do not automatically replay tool calls or require receipts for every prior action. Detached CLI runs remain eligible for adoption during the same mixed restart. If an ACP process terminates but its terminal run update does not persist, startup classifies it as lost with reason `selective_drain_not_finalized` rather than treating the drain as successful. On startup the new server writes `$PAPERCLIP_HOME/instances/${PAPERCLIP_INSTANCE_ID:-default}/hot-restart-report.json` with `previousServerPid`, `newServerPid`, `previousServerVersion`, `newServerVersion`, `drainReason`, `adoptedRunIds`, `finalizedWhileDownRunIds`, `lostRunIds`, and per-run classifications before the normal orphan reaper runs. When Paperclip manages embedded PostgreSQL, it suppresses that dependency's eager `SIGINT`/`SIGTERM` cleanup hooks. Paperclip owns signal ordering so the heartbeat diff --git a/doc/SPEC-implementation.md b/doc/SPEC-implementation.md index a2e39acd95..c404003f89 100644 --- a/doc/SPEC-implementation.md +++ b/doc/SPEC-implementation.md @@ -41,7 +41,7 @@ These decisions close open questions from `SPEC.md` for V1. | Communication | Tasks + comments only (no separate chat system) | | Task ownership | Single assignee; atomic checkout required for `in_progress` transition | | Task watchdogs | A task watchdog is an explicitly configured, issue-subtree-scoped verification and recovery capacity. It may restore live task paths inside the watched subtree; for issue-thread interaction resolution it is an ordinary agent subject to the same audience and containment checks, not board authority, active-run output monitoring, or general liveness recovery. | -| Recovery | Liveness/watchdog recovery preserves explicit ownership: retry lost execution continuity where safe, otherwise open visible source-scoped recovery actions by default, use issue-backed recovery only for independent repair work, or require human escalation (see `doc/execution-semantics.md`) | +| Recovery | Liveness/watchdog recovery preserves explicit ownership: continue interrupted local conversations with bounded fresh turns and preserved history, never replay tool calls automatically; retain native ownership and real execution gates; otherwise open visible source-scoped recovery actions by default, use issue-backed recovery only for independent repair work, or require human escalation (see `doc/execution-semantics.md`) | | Agent adapters | Built-in `process`, `http`, local CLI/session adapters, and OpenClaw gateway support; external adapters can also be loaded through the adapter plugin flow | | Plugin framework | Local/self-hosted early plugin runtime is in scope; cloud marketplace and packaged public distribution remain out of scope | | Auth | Mode-dependent human auth (`local_trusted` implicit board in current code; authenticated mode uses sessions), API keys for agents | diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index 5b14ed6600..47be6583c5 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -481,7 +481,7 @@ Agent-assigned `in_review` with no typed participant is only healthy when one of An `in_review` issue is stalled when it has no typed participant, no pending interaction or approval, no user owner, no active monitor, no active run, no queued wake, and no explicit recovery action. Paperclip should surface that state as recovery work rather than silently completing the issue or leaving blocker chains parked indefinitely. -When an execution-policy review stage has a pending agent participant, the participant's run is part of the review path only while it is live or queued. If that participant run reaches a terminal state while `executionState.status` remains `pending`, no decision has been recorded. After a successful run with no review decision, Paperclip should queue one bounded normal-model recovery wake for the same participant when the agent is invokable and no other review path exists. A failed participant instead follows the provider-continuity rules below: positive bootstrap evidence or a validated native resume/replacement can permit bounded recovery; uncertain effects use the automatic no-replay disposition while preserving the original assignee. If that recovery run also finishes while the stage remains pending, or the participant cannot be invoked, Paperclip must move the source issue to an explicit blocked/recovery path instead of leaving `in_review` to drift silently. +When an execution-policy review stage has a pending agent participant, the participant's run is part of the review path only while it is live or queued. If that participant run reaches a terminal state while `executionState.status` remains `pending`, no decision has been recorded. After a successful run with no review decision, Paperclip should queue one bounded normal-model recovery wake for the same participant when the agent is invokable and no other review path exists. A failed participant instead follows the provider-continuity rules below: local conversational adapters can start a bounded continuation turn, while native sessions use validated resume/replacement. Other adapters retain their action-recovery gates. The original assignee stays unchanged. If that recovery run also finishes while the stage remains pending, or the participant cannot be invoked, Paperclip must move the source issue to an explicit blocked/recovery path instead of leaving `in_review` to drift silently. ### Issue monitors @@ -791,8 +791,8 @@ Auto-recovery is allowed when ownership is clear and the control plane only lost Examples: -- requeue one dispatch wake for an assigned `todo` issue whose latest run failed, timed out, or was cancelled only when the provider-continuity rules establish safe recovery -- requeue one continuation wake for an assigned `in_progress` issue whose live execution path disappeared only with the required continuity and action-outcome evidence +- requeue one dispatch wake for an assigned `todo` issue whose latest run failed, timed out, or was cancelled under the bounded conversation or provider-continuity rules below +- requeue one continuation wake for an assigned `in_progress` issue whose live execution path disappeared under the bounded conversation or provider-continuity rules below - assign an orphan blocker back to its creator when that blocker is already preventing other work Auto-recovery preserves the existing owner. It does not choose a replacement agent. @@ -815,7 +815,7 @@ to recover recorded output. ### Provider continuity and bounded finalization -A permanently unusable established provider session may be replaced only with evidence that its predecessor is stopped and fenced, completed results and workspace state are preserved, required task history is available, and pending effects have been reconciled. A provider-native shell command or external write without a reliable outcome receipt is unknown. Unknown effects, integrity failures, and unverified process ownership never authorize speculative replay. Once automatic recovery is ruled out, Paperclip selects a conservative default: preserve recorded work, stop the affected task, and retain a durable no-replay hold. Unknown action outcomes remain unknown. No reconciliation form or user diagnosis is required. +A permanently unusable native runner session may be replaced only with evidence that its predecessor is stopped and fenced, completed results and workspace state are preserved, required task history is available, and pending effects have been reconciled. A provider-native shell command or external write without a reliable outcome receipt is unknown. Unknown effects, integrity failures, and unverified process ownership never authorize speculative replay. Once automatic recovery is ruled out, Paperclip selects a conservative default: preserve recorded work, stop the affected task, and retain a durable no-replay hold. Unknown action outcomes remain unknown. No reconciliation form or user diagnosis is required. Bootstrap retries, exact-checkpoint resumes, and fresh replacement sessions share three total provider attempts, including the original attempt. Linked run IDs, controller restarts, and duplicate wakes do not reset this budget. Automatic attempts retain the 30-second delay. Replacement scheduling and predecessor lineage commit together, with one successor per predecessor and admission through the normal task locks, authorization, pause, approval, and budget gates. @@ -823,13 +823,19 @@ Provider execution and control-plane finalization have different clocks. A healt Every continuation carries the triggering request, ordered user direction, interaction outcomes, completed work, and explicit history coverage. A delivered message remains part of the task's request after its connection or approval resolves. The original title is background; a completed Notion read does not satisfy a later Gmail request. Author and source-trust boundaries survive rendering into both native and legacy prompts. Missing required history must be fetched before dispatch rather than described as complete. -Legacy adapters without a verified resume capability use the same automatic no-replay disposition after provider failure. An availability error family (including quota or upstream overload) is not proof that earlier actions did not happen. The compatible adapter result field `executionRecovery: { kind: "bootstrap", providerWorkStarted: false }` can establish a pre-provider retry; the server records the same evidence for failures before adapter dispatch. Bootstrap retries and process-loss bootstrap retries use the same durable counter and delay. Productive max-turn continuation remains a separate execution boundary rather than a failed provider incident. A pre-dispatch wait for a confirmed live workspace holder is also a resource wait, not a provider failure: explicit `workspace_wait` evidence preserves that wait path without consuming the failure incident budget. +### Interrupted conversation continuation -The server projection remains available for execution diagnostics. Normal working, finishing, and interaction waits add no badges or cards to task lists or feeds. Active transcript headers keep saying Working during automatic retry and execution confirmation; attempts, causes, and recovery decisions belong in the run log. There is no reconciliation dialog. Safe recovery remains automatic. If it cannot continue safely, the source-scoped recovery record resolves with a blocked no-replay disposition and the ordinary task status becomes blocked, preserving its owner. Resolving this record does not grant replay authority: dispatch continues enforcing the durable hold. Replacement history remains inspectable and the composer stays usable. +An interrupted conversation does not permanently block its task. For local conversational adapters, Paperclip starts a new bounded turn with the existing session when compatible, or the full task conversation when the session is unavailable. The prompt says: “Your previous run was interrupted. Continue from where you left off.” The agent decides what remains from the history and latest user request. Paperclip never automatically replays recorded tool calls. Unknown past action outcomes are not a task-wide execution gate, and no action-reconciliation questionnaire is required. -An operator Stop reaches embedded ACP execution through its run-owned cancellation signal. The response waits for adapter settlement; acknowledgment requires the local provider to have exited. A deadline or failed cleanup never grants continuation permission. A persistent local ACP session can record an interrupted checkpoint only after acknowledged cancellation, complete tool reporting with settled reads (or no tools), and successful cleanup. Writes, shell commands, incomplete client-operation receipts, forced cancellation, and lost transports retain the ordinary no-replay hold. Continuation must restore the same compatible session; an unavailable checkpoint cannot fall back to a new session. A restored provider receives the current run identity, API credential, and scratch environment. Run-owned scratch paths rotate without changing session identity, while user configuration changes still invalidate compatibility. +Shutdown, process loss, and provider failure use the existing durable failure retry counter and delay. Ordinary failure recovery permits at most two automatic retries in a failure chain. Accepted-interaction infrastructure recovery retains its existing bounded policy. Repeated scheduler visits reuse the same successor; restarting the server does not reset the counter. After exhaustion, automatic attempts stop. A new explicit user message can start a fresh run and failure budget. Productive max-turn continuation and confirmed workspace waits keep their separate existing semantics. -Stop alone does not promote deferred messages. A subsequent explicit wake adopts pending comment IDs atomically in order through the existing queue. The task's ordered continuation history remains authoritative. A subtree pause still requires Resume; the text “go” has no special bypass. Task detail exposes the effective execution blocker, including a recovery record resolved with replay blocked, using the same predicate as dispatch and Resume. A cancelled run that never started says “Couldn't start” instead of claiming successful completion without an answer. Historical ambiguous executions remain held. A queued message or healthy child task cannot clear an execution reconciliation hold during a generic recovery sweep. +Real gates still apply: company and task ownership, active provider ownership, budget limits, agent availability, dependencies, pending approval/review paths, and explicit pause holds. Native runner reattachment and finalization retain their existing ownership protocol. Process, HTTP, and gateway adapters retain their recovery rules because invoking those adapters can itself repeat an external action rather than start a conversation turn. + +An operator Stop still waits for local provider termination. Stop alone never promotes deferred comments or starts an automatic continuation. Once stopped, the next explicit wake adopts pending comment IDs in order through the existing queue. A compatible saved ACP session can resume, and an unavailable or incompatible session can start fresh with the full task context. Run credentials and scratch paths remain scoped to the new run. A subtree pause requires Resume; a message does not bypass it. + +Historical legacy interruption holds for conversational adapters no longer block new messages or Resume. Classification uses the run’s saved adapter invocation or continuation policy, never the agent’s current adapter settings. Missing historical adapter evidence retains the hold. A terminal row with a live predecessor process or unreleased environment lease still blocks actual admission and Resume. Retry scheduling can happen before cleanup, but grants no execution authority. Recovery folds their obsolete no-replay bookkeeping without changing task ownership, status, or automatically waking old work. The audit trail remains readable. Native integrity and ownership holds, and non-conversational adapter holds, remain enforced. + +The server projection remains available for diagnostics. Normal working, finishing, and interaction waits add no badges or cards to task lists or feeds. Active transcript headers keep saying Working during automatic retry and execution confirmation; attempts, causes, and recovery decisions belong in the run log. Recovery uses the existing transcript and run log rather than adding a reconciliation form. A cancelled run that never started says “Couldn't start” instead of implying that the agent answered. ### Codex startup and provider state diff --git a/packages/adapter-utils/src/acpx-engine/execute.ts b/packages/adapter-utils/src/acpx-engine/execute.ts index e7168640dd..ae10093fcb 100644 --- a/packages/adapter-utils/src/acpx-engine/execute.ts +++ b/packages/adapter-utils/src/acpx-engine/execute.ts @@ -3910,6 +3910,7 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { let forcedStop = false; let runtimeStopConfirmed = false; let safeInterruptedSession = false; + let preserveInterruptedSession = false; const interruptionTools = new Map(); let incompleteToolInventory = false; try { @@ -4144,9 +4145,6 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { const previousParams = parseObject(ctx.runtime.sessionParams); const canResume = isCompatibleSession(previousParams, prepared); - if (previousParams.interruptedCheckpoint === true && !canResume) { - throw new Error("The interrupted session is no longer compatible. Its action history must be checked before starting a new session."); - } const resumeSessionId = canResume ? asString(previousParams.acpSessionId, "") || undefined : undefined; // Borrow the warm entry without removing it, so an overlapping run of the // same session still sees it. The borrow clears the entry's idle timer, so @@ -4347,7 +4345,7 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { }), }); } catch (err) { - if (!resumeSessionId || !isResumeFailure(err) || previousParams.interruptedCheckpoint === true) throw err; + if (!resumeSessionId || !isResumeFailure(err)) throw err; clearSession = true; resumedSession = false; await ctx.onLog( @@ -4394,10 +4392,11 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { parentContext: prepared.stepMetrics.parentContext, }); } - // A compatible warm handle reuses the already-running ACP agent and does - if (previousParams.interruptedCheckpoint === true && handle?.backendSessionId !== resumeSessionId) { - throw new Error("The provider did not restore the interrupted session; refusing a fresh-session fallback."); + if (resumeSessionId && handle?.backendSessionId !== resumeSessionId) { + resumedSession = false; + clearSession = true; } + // A compatible warm handle reuses the already-running ACP agent and does // not emit another spawn event. Persist its known identity on this run // before the next prompt starts so every running heartbeat is adoptable. if (handle && cached && processIdentitySink.latest && ctx.onSpawn) { @@ -4883,13 +4882,18 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { eventCostUsd, }); const failedTurn = terminal.status === "failed" || terminal.status === "cancelled" || timedOut; - // A provider-native command/write has no reliable external outcome - // receipt. Only settled reads (or a turn with no tools) can establish - // automatic interrupted-session continuity here. - safeInterruptedSession = ctx.signal?.aborted === true && !forcedStop && !timedOut && !channelLost + // ACPX can defer session/load until runTurn. Forget an unavailable + // session so the next bounded turn receives the full task conversation. + const sessionUnavailable = terminal.status === "failed" && + terminal.error.detailCode === "SESSION_RESUME_REQUIRED"; + if (sessionUnavailable) clearSession = true; + // Saving a conversation is independent from certifying tool outcomes. + // Its next turn receives history, not a replay of pending tool calls. + preserveInterruptedSession = ctx.signal?.aborted === true && !forcedStop && !timedOut && !channelLost && (terminal.status === "cancelled" || terminal.status === "completed") && prepared.mode === "persistent" && !prepared.processSessionBridge - && Boolean(sessionHandle.backendSessionId) + && Boolean(sessionHandle.backendSessionId); + safeInterruptedSession = preserveInterruptedSession && !incompleteToolInventory && [...interruptionTools.values()].every((tool) => tool.kind === "read" && tool.status === "completed"); // Record how the settlement `endSession` step closes the runtime for this @@ -4911,7 +4915,7 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { : failedTurn ? `paperclip turn ${terminal.status}` : "paperclip completed turn cleanup", - discardPersistentState: (terminal.status === "cancelled" && !safeInterruptedSession) || timedOut || channelLost, + discardPersistentState: sessionUnavailable || (terminal.status === "cancelled" && !preserveInterruptedSession) || timedOut || channelLost, dropWarmEntry: false, recordCloseError: false, cancelTurnReason: null, diff --git a/packages/adapter-utils/src/acpx-engine/operator-stop.test.ts b/packages/adapter-utils/src/acpx-engine/operator-stop.test.ts index ee5b70beb7..304ed26bb7 100644 --- a/packages/adapter-utils/src/acpx-engine/operator-stop.test.ts +++ b/packages/adapter-utils/src/acpx-engine/operator-stop.test.ts @@ -39,7 +39,7 @@ it('stops an actual ACP process and resumes its established session with the new expect(params?.interruptedCheckpoint).toBe(true); const next = await execute({ ...ctx, runId: 'follow-up', authToken: 'follow-up-test-token', signal: undefined, context: { prompt: 'List recent Drive files' }, runtime: { ...ctx.runtime, sessionParams: params } }); - expect(next.exitCode).toBe(0); + expect(next.exitCode, JSON.stringify(next)).toBe(0); const prompts = (await fs.readFile(path.join(root, 'prompts'), 'utf8')).trim().split('\n').map(line => JSON.parse(line)); expect(prompts).toHaveLength(2); expect(prompts[1].sessionId).toBe(prompts[0].sessionId); @@ -47,7 +47,7 @@ it('stops an actual ACP process and resumes its established session with the new expect(launches.map(launch => launch.runId)).toEqual(['stop-test', 'follow-up']); expect(launches[1].tokenHash).toBe(createHash('sha256').update('follow-up-test-token').digest('hex')); }); -it('stops writes but does not authorize replay when the interrupted tool has no outcome', async () => { +it('continues after an interrupted write without replaying that write', async () => { const { root, ctx, abort, started, execute } = await setup('write'); const running = execute(ctx); await started; @@ -57,8 +57,11 @@ it('stops writes but does not authorize replay when the interrupted tool has no expect(result.resultJson?.executionCancellation).toMatchObject({ state: 'acknowledged' }); expect(result.executionRecovery).toBeUndefined(); const before = await fs.readFile(path.join(root, 'writes'), 'utf8'); - await new Promise(resolve => setTimeout(resolve, 5000)); + const next = await execute({ ...ctx, runId: 'follow-up', signal: undefined, context: { prompt: 'What happened?' }, + runtime: { ...ctx.runtime, sessionParams: sessionCodec.serialize(result.sessionParams ?? null) } }); + expect(next.exitCode, JSON.stringify(next)).toBe(0); expect(await fs.readFile(path.join(root, 'writes'), 'utf8')).toBe(before); + expect((await fs.readFile(path.join(root, 'completed'), 'utf8')).trim()).toBe('follow-up'); }, 15000); it('does not dispatch a provider when Stop precedes startup', async () => { const { root, ctx, abort, execute } = await setup(); @@ -67,7 +70,7 @@ it('does not dispatch a provider when Stop precedes startup', async () => { await expect(fs.access(path.join(root, 'prompts'))).rejects.toThrow(); }); -it.each(['missing session', 'changed configuration'])('refuses fresh-session fallback after Stop: %s', async (change) => { +it.each(['missing session', 'changed configuration'])('starts a new turn when the interrupted session cannot resume: %s', async (change) => { const { root, ctx, abort, started, execute } = await setup(); const running = execute(ctx); await started; @@ -75,12 +78,17 @@ it.each(['missing session', 'changed configuration'])('refuses fresh-session fal const result = await running; expect(result.executionRecovery?.kind).toBe('interrupted'); if (change === 'missing session') await fs.rm(path.join(root, 'session')); - const next = await execute({ ...ctx, signal: undefined, + let next = await execute({ ...ctx, signal: undefined, config: change === 'changed configuration' ? { ...ctx.config, env: { ...(ctx.config.env as object), SETTING: 'changed' } } : ctx.config, runtime: { ...ctx.runtime, sessionParams: sessionCodec.serialize(result.sessionParams ?? null) }, }); - expect(next.exitCode).not.toBe(0); - expect((await fs.readFile(path.join(root, 'prompts'), 'utf8')).trim().split('\n')).toHaveLength(1); + if (change === 'missing session') { + expect(next.clearSession, JSON.stringify(next)).toBe(true); + next = await execute({ ...ctx, runId: 'fresh-follow-up', signal: undefined, + runtime: { ...ctx.runtime, sessionParams: null } }); + } + expect(next.exitCode, JSON.stringify(next)).toBe(0); + expect((await fs.readFile(path.join(root, 'prompts'), 'utf8')).trim().split('\n')).toHaveLength(2); }); it('keeps the Stop deadline active after cancellation returns until provider exit', async () => { diff --git a/packages/adapter-utils/src/server-utils.ts b/packages/adapter-utils/src/server-utils.ts index e7abb155ac..93bab9ae63 100644 --- a/packages/adapter-utils/src/server-utils.ts +++ b/packages/adapter-utils/src/server-utils.ts @@ -2404,6 +2404,9 @@ export function renderPaperclipWakePrompt( ]; if (normalized.executionContinuation) { + if (normalized.executionContinuation.interruptedRunId) { + lines.push("", "Your previous run was interrupted. Continue from where you left off using the conversation history and the latest user request. Prior tool calls are history, not commands to replay. Decide what remains and take the next appropriate step."); + } const { resumeDelta, ...snapshot } = normalized.executionContinuation; const continuation = resumedSession && resumeDelta ? { ...snapshot, messages: resumeDelta.messages, coverage: { ...snapshot.coverage, kind: "task_history_delta", baseRunId: resumeDelta.baseRunId }, diff --git a/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts b/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts index a6a00d4286..34d080ace4 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts @@ -448,8 +448,14 @@ it.each([ expect(runnerPid).toBeGreaterThan(0); expect(providerPid).toBeGreaterThan(0); await bundle.detachControllerForRestart(); - process.kill(-runnerPid, "SIGKILL"); - process.kill(-providerPid, "SIGKILL"); + // The provider can exit when its runner dies. An already-gone process + // group satisfies teardown; still fail on other signal errors and join below. + for (const pid of [runnerPid, providerPid]) { + try { process.kill(-pid, "SIGKILL"); } + catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ESRCH") throw error; + } + } await vi.waitFor(() => { expect(dead(runnerPid)).toBe(true); expect(dead(providerPid)).toBe(true); diff --git a/packages/paperclip-runner/test/fixtures/fake-final-burst-codex-app-server.mjs b/packages/paperclip-runner/test/fixtures/fake-final-burst-codex-app-server.mjs index 106f346ee5..fbe492a852 100644 --- a/packages/paperclip-runner/test/fixtures/fake-final-burst-codex-app-server.mjs +++ b/packages/paperclip-runner/test/fixtures/fake-final-burst-codex-app-server.mjs @@ -87,6 +87,7 @@ createInterface({ input: process.stdin }).on("line", (line) => { modelProvider: "fixture-no-provider", thread: { id: state.threadId, + status: { type: Object.values(state.turns).some(turn => turn.status === "inProgress") ? "active" : "idle" }, sessionId: "final-burst-fixture", turns: Object.entries(state.turns).map(([turnId, turn]) => ({ id: turnId, diff --git a/packages/shared/src/types/execution-continuation.ts b/packages/shared/src/types/execution-continuation.ts index 6505773947..0d980b9cd5 100644 --- a/packages/shared/src/types/execution-continuation.ts +++ b/packages/shared/src/types/execution-continuation.ts @@ -35,6 +35,8 @@ export interface ExecutionContinuationEnvelope { }; recoveryOutcomes?: Array<{ recoveryActionId: string; decision: unknown }>; completedWork: string | null; + /** Start a new turn from history; never replay prior tool calls automatically. */ + interruptedRunId?: string; /** Completed mutations are context, never instructions to replay them. */ completedActions?: Array<{ runId: string; diff --git a/packages/shared/src/types/execution-projection.ts b/packages/shared/src/types/execution-projection.ts index 4d75380438..2ba1361341 100644 --- a/packages/shared/src/types/execution-projection.ts +++ b/packages/shared/src/types/execution-projection.ts @@ -1,5 +1,6 @@ export interface ExecutionBlocker { - recoveryActionId: string; + /** Null when live execution authority itself blocks continuation. */ + recoveryActionId: string | null; runId: string | null; agentId: string | null; cause: string; diff --git a/scripts/mcp-fixtures/servers/acp-stop-agent.mjs b/scripts/mcp-fixtures/servers/acp-stop-agent.mjs index c31cdfcdea..5a8e92a957 100644 --- a/scripts/mcp-fixtures/servers/acp-stop-agent.mjs +++ b/scripts/mcp-fixtures/servers/acp-stop-agent.mjs @@ -22,7 +22,7 @@ async function request(message) { return { sessionId }; } case 'session/load': - if (fs.readFileSync(`${root}/session`, 'utf8') !== message.params.sessionId) throw Error('Unknown session'); + if (!fs.existsSync(`${root}/session`) || fs.readFileSync(`${root}/session`, 'utf8') !== message.params.sessionId) throw Error('Unknown session'); return {}; case 'session/prompt': { fs.appendFileSync(`${root}/prompts`, `${JSON.stringify(message.params)}\n`); diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index 495d489670..e12e62ddb2 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -650,7 +650,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { async function seedRunFixture(input?: { adapterType?: string; agentStatus?: "paused" | "idle" | "running"; - runStatus?: "running" | "queued" | "failed"; + runStatus?: "running" | "queued" | "failed" | "interrupted"; processPid?: number | null; processGroupId?: number | null; processLossRetryCount?: number; @@ -719,10 +719,14 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { ...(input?.runtimeMode ? { runtimeMode: input.runtimeMode } : {}), errorCode: input?.runErrorCode ?? null, error: input?.runError ?? null, + nextEventSeq: 2, startedAt: now, updatedAt: new Date("2026-03-19T00:00:00.000Z"), }); + await db.insert(heartbeatRunEvents).values({ companyId, agentId, runId, + seq: 1, eventType: "adapter.invoke", payload: { adapterType: input?.adapterType ?? "codex_local" } }); + if (input?.includeIssue !== false) { await db.insert(issues).values({ id: issueId, @@ -1459,7 +1463,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { return { companyId, agentId, runId, wakeupRequestId, issueId }; } - it("persists the normalized failure and exposes an operator recovery action", async () => { + it("persists the normalized failure without permanently blocking the conversation", async () => { mockAdapterExecute.mockResolvedValueOnce({ exitCode: 1, signal: null, @@ -1483,12 +1487,6 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { .from(agentRuntimeState) .where(eq(agentRuntimeState.agentId, agentId)) .then((rows) => rows[0] ?? null); - const agent = await db - .select({ status: agents.status, errorReason: agents.errorReason }) - .from(agents) - .where(eq(agents.id, agentId)) - .then((rows) => rows[0] ?? null); - const recoveryRun = await db .select() .from(heartbeatRuns) @@ -1497,16 +1495,8 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(run).toMatchObject({ status: "failed", error: "Adapter failed" }); expect(runtime?.lastError).toBe("Adapter failed"); - expect(recoveryRun).toBeNull(); - const [action] = await db - .select() - .from(issueRecoveryActions) - .where(eq(issueRecoveryActions.sourceIssueId, issueId)); - expect(action).toMatchObject({ - cause: "legacy_execution_requires_reconciliation", - ownerType: "board", - returnOwnerAgentId: agentId, - }); + expect(recoveryRun).toMatchObject({ retryOfRunId: runId, scheduledRetryAttempt: 1 }); + expect(await getExecutionBlocker(db, companyId, issueId)).toBeNull(); const missingCommentWakeups = await db .select({ id: agentWakeupRequests.id }) .from(agentWakeupRequests) @@ -1517,7 +1507,6 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { ), ); expect(missingCommentWakeups).toHaveLength(0); - expect(agent?.status).not.toBe("running"); }); it("does not immediately continue a low-trust preflight setup failure", async () => { @@ -1624,11 +1613,71 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(recoveryRuns).toHaveLength(0); }); + it("does not relabel a lost process run when its agent changes to a conversation adapter", async () => { + const { companyId, agentId, runId } = await seedRunFixture({ adapterType: "process", processPid: 99999999 }); + await db.update(agents).set({ adapterType: "codex_local" }).where(eq(agents.id, agentId)); + await heartbeatService(db).reapOrphanedRuns(); + const [run] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)); + expect(run.status).toBe("failed"); + expect(run.resultJson?.conversationContinuation).toBeUndefined(); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + }); + + it.each([false, true])("waits for a live terminal predecessor before a new turn (historical hold: %s)", async withHold => { + const child = spawnAliveProcess(); + childProcesses.add(child); + const { runId, issueId, companyId } = await seedRunFixture({ + agentStatus: "idle", runStatus: "interrupted", runErrorCode: "server_shutdown_interrupted", processPid: child.pid!, + }); + await db.update(heartbeatRuns).set({ resultJson: { conversationContinuation: "continue_conversation_v1" } }).where(eq(heartbeatRuns.id, runId)); + if (withHold) await db.insert(issueRecoveryActions).values({ + companyId, sourceIssueId: issueId, kind: "active_run_watchdog", status: "active", ownerType: "board", + cause: "legacy_execution_requires_reconciliation", fingerprint: runId, + evidence: { runId, automaticRecovery: { replay: "blocked" } }, nextAction: "Wait for the previous execution to stop.", + }); + expect(await getExecutionBlocker(db, companyId, issueId)).toMatchObject({ runId, cause: "execution_owner_active" }); + // Scheduling itself has no execution authority. Actual queued admission + // must still block while the predecessor owns a process or lease. + const queuedId = randomUUID(); + const previous = (await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)))[0]!; + await db.insert(heartbeatRuns).values({ id: queuedId, companyId, agentId: previous.agentId, + status: "queued", contextSnapshot: { issueId, wakeReason: "issue_commented" } }); + const { createPostgresRunDispatchAdapter } = await import("../modules/run-dispatch/adapters/postgres.js"); + expect(await createPostgresRunDispatchAdapter(db).cancelStaleQueuedRun({ companyId, runId: queuedId, + expectedStatus: "queued", now: new Date() })).toMatchObject({ outcome: "cancelled", errorCode: "execution_reconciliation_required" }); + const { settleUnrecoverableExecutions } = await import("../services/execution-recovery-resolution.js"); + await settleUnrecoverableExecutions(db); + if (withHold) expect((await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.sourceIssueId, issueId)))[0].status).toBe("active"); + expect(isPidAlive(child.pid!)).toBe(true); + child.kill("SIGTERM"); + await waitForPidExit(child.pid!); + expect(await getExecutionBlocker(db, companyId, issueId)).toBeNull(); + }); + + it("waits for a terminal predecessor's environment lease to be released", async () => { + const { companyId, issueId, runId } = await seedRunFixture({ agentStatus: "idle", runStatus: "interrupted" }); + await db.update(heartbeatRuns).set({ resultJson: { conversationContinuation: "continue_conversation_v1" } }).where(eq(heartbeatRuns.id, runId)); + const [lease] = await db.insert(environmentLeases).values({ companyId, issueId, heartbeatRunId: runId, status: "active" }).returning(); + expect(await getExecutionBlocker(db, companyId, issueId)).toMatchObject({ runId, cause: "execution_owner_active" }); + // Scheduling itself has no execution authority. Actual queued admission + // must still block while the predecessor owns a process or lease. + const queuedId = randomUUID(); + const previous = (await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)))[0]!; + await db.insert(heartbeatRuns).values({ id: queuedId, companyId, agentId: previous.agentId, + status: "queued", contextSnapshot: { issueId, wakeReason: "issue_commented" } }); + const { createPostgresRunDispatchAdapter } = await import("../modules/run-dispatch/adapters/postgres.js"); + expect(await createPostgresRunDispatchAdapter(db).cancelStaleQueuedRun({ companyId, runId: queuedId, + expectedStatus: "queued", now: new Date() })).toMatchObject({ outcome: "cancelled", errorCode: "execution_reconciliation_required" }); + await db.update(environmentLeases).set({ releasedAt: new Date(), status: "released" }).where(eq(environmentLeases.id, lease!.id)); + expect(await getExecutionBlocker(db, companyId, issueId)).toBeNull(); + }); + it("keeps an unsafe Stop blocked when recovery sees a deferred human comment", async () => { const { companyId, agentId, issueId, runId } = await seedStrandedIssueFixture({ status: "in_progress", runStatus: "cancelled", resultJson: { executionCancellation: { state: "acknowledged" } }, }); + await db.update(agents).set({ adapterType: "process" }).where(eq(agents.id, agentId)); const [run] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)); await terminalizeLegacyExecution({ db, run, status: "cancelled" }); const wakeId = randomUUID(); @@ -2296,7 +2345,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(runs).toHaveLength(0); }); - it("holds an unknown dead-provider outcome for reconciliation", async () => { + it("schedules one conversation continuation after losing the provider", async () => { const { agentId, runId, issueId } = await seedRunFixture({ agentStatus: "idle", processPid: 999_999_999, @@ -2317,11 +2366,11 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { .select() .from(heartbeatRuns) .where(eq(heartbeatRuns.agentId, agentId)); - expect(runs).toHaveLength(1); + expect(runs).toHaveLength(2); const failedRun = runs.find((row) => row.id === runId); const retryRuns = runs.filter((row) => row.retryOfRunId === runId); - expect(retryRuns).toHaveLength(0); + expect(retryRuns).toHaveLength(1); const retryRun = retryRuns[0]; expect(failedRun?.status).toBe("failed"); expect(failedRun?.errorCode).toBe("process_lost"); @@ -2336,18 +2385,15 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { .select() .from(issueRecoveryActions) .where(eq(issueRecoveryActions.sourceIssueId, issueId)); - expect(action).toMatchObject({ - cause: "legacy_execution_requires_reconciliation", - returnOwnerAgentId: agentId, - ownerType: "board", - }); + expect(action).toBeUndefined(); + expect(retryRun).toMatchObject({ status: "scheduled_retry", scheduledRetryAttempt: 1 }); await heartbeat.reapOrphanedRuns(); expect( await db .select() .from(heartbeatRuns) .where(eq(heartbeatRuns.agentId, agentId)), - ).toHaveLength(1); + ).toHaveLength(2); const issue = await waitForValue(async () => db @@ -2769,7 +2815,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { [runId], ); expect(drain.interruptedRunIds).toEqual([runId]); - expect(drain.retryRunIds).toHaveLength(0); + expect(drain.retryRunIds).toHaveLength(1); await waitForPidExit(child.pid!); const reconciliation = await heartbeat.reconcileHotRestartAdoption( @@ -2791,7 +2837,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { status: "interrupted", errorCode: "server_shutdown_interrupted", }); - expect(runs).toHaveLength(1); + expect(runs).toHaveLength(2); const report = JSON.parse( await fs.readFile(resolveHotRestartReportPath(home), "utf8"), @@ -3390,6 +3436,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { it("terminalizes an unsupported legacy session on shutdown without speculative replay", async () => { const { agentId, runId, issueId, wakeupRequestId } = await seedRunFixture({ + adapterType: "process", agentStatus: "running", }); const result = await heartbeatService(db).drainRunningRunsForShutdown( @@ -3541,6 +3588,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { it("does not duplicate a legacy reconciliation action across repeated shutdowns", async () => { const { agentId, runId, issueId } = await seedRunFixture({ + adapterType: "process", agentStatus: "running", }); const heartbeat = heartbeatService(db); @@ -3587,12 +3635,9 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { .select() .from(issueRecoveryActions) .where(eq(issueRecoveryActions.sourceIssueId, issueId)), - ).toEqual([ - expect.objectContaining({ - cause: "legacy_execution_requires_reconciliation", - evidence: expect.objectContaining({ attempt: 3 }), - }), - ]); + ).toEqual([]); + const run = await heartbeatService(db).getRun(runId); + expect(await getExecutionBlocker(db, run!.companyId, issueId)).toBeNull(); }); it("releases active environment leases when an orphaned run is reaped", async () => { @@ -3658,6 +3703,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { it("does not bypass unknown process outcomes through immediate continuation recovery", async () => { const { agentId, runId, issueId } = await seedRunFixture({ + adapterType: "process", agentStatus: "idle", processPid: 999_999_999, processLossRetryCount: 1, @@ -3688,6 +3734,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { it("blocks failed recovery work in place during immediate terminal-run cleanup", async () => { const sourceIssueId = randomUUID(); const { companyId, agentId, runId, issueId } = await seedRunFixture({ + adapterType: "process", agentStatus: "idle", processPid: 999_999_999, processLossRetryCount: 1, @@ -3824,7 +3871,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(comments).toHaveLength(0); }); - it("does not treat a transient remote-compaction failure as evidence of safe replay", async () => { + it("allows conversation continuation after a transient remote-compaction failure", async () => { mockAdapterExecute.mockResolvedValueOnce({ exitCode: 1, signal: null, @@ -3845,33 +3892,22 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { await heartbeat.resumeQueuedRuns(); await waitForRunToSettle(heartbeat, runId); + await heartbeat.waitForRunExecutionDrain(runId); - expect( - await db - .select() - .from(heartbeatRuns) - .where(eq(heartbeatRuns.agentId, agentId)), - ).toEqual([ - expect.objectContaining({ - id: runId, - status: "failed", - errorCode: "adapter_failed", - resultJson: expect.objectContaining({ - errorFamily: "transient_upstream", - }), - }), - ]); + const runs = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, agentId)); + expect(runs).toHaveLength(2); + expect(runs.find(run => run.id === runId)).toMatchObject({ + status: "failed", errorCode: "adapter_failed", resultJson: { errorFamily: "transient_upstream" }, + }); + expect(runs.find(run => run.retryOfRunId === runId)).toMatchObject({ + status: "scheduled_retry", scheduledRetryAttempt: 1, + }); expect( await db .select() .from(issueRecoveryActions) .where(eq(issueRecoveryActions.sourceIssueId, issueId)), - ).toEqual([ - expect.objectContaining({ - cause: "legacy_execution_requires_reconciliation", - ownerType: "board", - }), - ]); + ).toEqual([]); expect(mockAdapterExecute).toHaveBeenCalledTimes(1); }); @@ -7611,7 +7647,7 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { const events = await db .select() .from(heartbeatRunEvents) - .where(eq(heartbeatRunEvents.runId, runId)); + .where(and(eq(heartbeatRunEvents.runId, runId), eq(heartbeatRunEvents.eventType, "lifecycle"))); expect(events).toHaveLength(1); expect(events[0]).toMatchObject({ eventType: "lifecycle", diff --git a/server/src/__tests__/heartbeat-retry-scheduling.test.ts b/server/src/__tests__/heartbeat-retry-scheduling.test.ts index 90c89dec06..490e01e150 100644 --- a/server/src/__tests__/heartbeat-retry-scheduling.test.ts +++ b/server/src/__tests__/heartbeat-retry-scheduling.test.ts @@ -3,6 +3,9 @@ import { and, eq, inArray, sql } from "drizzle-orm"; import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest"; import { agents, + approvals, + issueApprovals, + issueThreadInteractions, agentRuntimeState, agentWakeupRequests, activityLog, @@ -164,6 +167,7 @@ describeEmbeddedPostgres("heartbeat bounded retry scheduling", () => { await db.delete(environmentLeases); await db.delete(issueRelations); await db.delete(issues); + await db.delete(approvals); await db.delete(executionWorkspaces); await db.delete(projects); await cleanupHeartbeatRunDependents(); @@ -445,6 +449,87 @@ describeEmbeddedPostgres("heartbeat bounded retry scheduling", () => { return { companyId, agentId, issueId, runId, now }; } + it("bounds interrupted conversations across restarts and concurrent scheduling", async () => { + const { companyId, issueId, runId, now } = await seedMaxTurnFixture(); + const resultJson = { conversationContinuation: "continue_conversation_v1" }; + await db.update(heartbeatRuns).set({ status: "interrupted", errorCode: "server_shutdown_interrupted", resultJson }) + .where(eq(heartbeatRuns.id, runId)); + let predecessor = runId; + for (const attempt of [1, 2]) { + const restarted = heartbeatService(db); + const outcomes = await Promise.all([ + restarted.scheduleBoundedRetry(predecessor, { now, random: () => 0 }), + restarted.scheduleBoundedRetry(predecessor, { now, random: () => 0 }), + ]); + expect(outcomes.every(outcome => outcome.outcome === "scheduled")).toBe(true); + const children = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, predecessor)); + expect(children).toHaveLength(1); + expect(children[0]).toMatchObject({ scheduledRetryAttempt: attempt }); + predecessor = children[0]!.id; + await db.update(heartbeatRuns).set({ status: "interrupted", finishedAt: now, resultJson }) + .where(eq(heartbeatRuns.id, predecessor)); + } + expect(await heartbeatService(db).scheduleBoundedRetry(predecessor, { now })) + .toMatchObject({ outcome: "retry_exhausted" }); + await heartbeatService(db).reconcileStrandedAssignedIssues(); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.companyId, companyId))).toHaveLength(3); + // Exhaustion leaves the task available to a new explicit request. + const { getExecutionBlocker } = await import("../services/execution-blocker.js"); + expect(await getExecutionBlocker(db, companyId, issueId)).toBeNull(); + }); + + it.each(["dependency", "disabled", "reassigned"])("respects the %s gate for interrupted conversations", async gate => { + const { companyId, agentId, issueId, runId, now } = await seedMaxTurnFixture(); + await db.update(heartbeatRuns).set({ status: "interrupted", errorCode: "process_lost", + resultJson: { conversationContinuation: "continue_conversation_v1" } }).where(eq(heartbeatRuns.id, runId)); + if (gate === "dependency") { + const blockerId = randomUUID(); + await db.insert(issues).values({ id: blockerId, companyId, title: "Required work", status: "todo" }); + await db.insert(issueRelations).values({ companyId, issueId: blockerId, relatedIssueId: issueId, type: "blocks" }); + } else if (gate === "disabled") { + await db.update(agents).set({ runtimeConfig: { heartbeat: { wakeOnDemand: false } } }).where(eq(agents.id, agentId)); + } else { + await db.update(issues).set({ assigneeAgentId: null }).where(eq(issues.id, issueId)); + } + expect(await heartbeat.scheduleBoundedRetry(runId, { now })).toMatchObject({ outcome: "not_scheduled" }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(0); + }); + + it.each([ + ["interaction", false], ["approval", false], ["interaction", true], ["approval", true], + ] as const)("waits for a pending %s before continuing (already scheduled: %s)", async (kind, alreadyScheduled) => { + const { companyId, issueId, runId, now } = await seedMaxTurnFixture(); + await db.update(heartbeatRuns).set({ status: "interrupted", errorCode: "process_lost", + resultJson: { conversationContinuation: "continue_conversation_v1" } }).where(eq(heartbeatRuns.id, runId)); + let retryRunId = runId; + if (alreadyScheduled) { + const scheduled = await heartbeat.scheduleBoundedRetry(runId, { now, random: () => 0 }); + expect(scheduled.outcome).toBe("scheduled"); + if (scheduled.outcome !== "scheduled") throw new Error("Expected a retry"); + retryRunId = scheduled.run.id; + } + if (kind === "interaction") { + await db.insert(issueThreadInteractions).values({ companyId, issueId, kind: "ask_user_questions", + status: "pending", payload: { version: 1, questions: [] } }); + } else { + const approvalId = randomUUID(); + await db.insert(approvals).values({ id: approvalId, companyId, type: "hire_agent", status: "pending", payload: {} }); + await db.insert(issueApprovals).values({ companyId, issueId, approvalId }); + } + if (alreadyScheduled) { + const adapter = createPostgresRunDispatchAdapter(db); + expect(await adapter.promoteOrCancelDueRetry({ companyId, runId: retryRunId, now: new Date(now.getTime() + 60_000) })) + .toMatchObject({ outcome: "gate_suppressed", errorCode: "issue_waiting_for_response" }); + const stopped = await heartbeat.getRun(retryRunId); + expect(stopped?.status).toBe("cancelled"); + const { legacyExecutionNeedsReconciliation } = await import("../services/legacy-execution-recovery.js"); + expect(legacyExecutionNeedsReconciliation(stopped!)).toBe(false); + } else { + expect(await heartbeat.scheduleBoundedRetry(runId, { now })) + .toMatchObject({ outcome: "not_scheduled", errorCode: "issue_waiting_for_response" }); + } + }); + it("schedules a retry with durable metadata and only promotes it when due", async () => { const companyId = randomUUID(); const agentId = randomUUID(); diff --git a/server/src/__tests__/issue-tree-control-routes.test.ts b/server/src/__tests__/issue-tree-control-routes.test.ts index bea4486453..3b22b9d985 100644 --- a/server/src/__tests__/issue-tree-control-routes.test.ts +++ b/server/src/__tests__/issue-tree-control-routes.test.ts @@ -19,6 +19,8 @@ const mockTreeControlService = vi.hoisted(() => ({ const mockReplayBlocks = vi.hoisted(() => vi.fn()); const mockReplayWhere = vi.hoisted(() => vi.fn()); +const mockExecutionBlocker = vi.hoisted(() => vi.fn()); +vi.mock("../services/execution-blocker.js", () => ({ getExecutionBlocker: mockExecutionBlocker })); const mockLogActivity = vi.hoisted(() => vi.fn()); const mockHeartbeatService = vi.hoisted(() => ({ @@ -47,7 +49,7 @@ async function createApp(actor: Record) { const query = { from: () => query, innerJoin: () => query, - where: (predicate: unknown) => { mockReplayWhere(predicate); return query; }, + where: (predicate: unknown) => { mockReplayWhere(predicate); return mockReplayBlocks(); }, limit: mockReplayBlocks, }; app.use("/api", issueTreeControlRoutes({ select: () => query } as any)); @@ -59,6 +61,7 @@ describe("issue tree control routes", () => { beforeEach(() => { vi.clearAllMocks(); mockReplayBlocks.mockResolvedValue([]); + mockExecutionBlocker.mockResolvedValue(null); mockTreeControlService.getHold.mockResolvedValue(null); mockIssueService.getById.mockResolvedValue({ id: "11111111-1111-4111-8111-111111111111", @@ -229,7 +232,9 @@ describe("issue tree control routes", () => { mode: "pause", members: [{ issueId: rootId }], }); - mockReplayBlocks.mockResolvedValue([{ identifier: "TEST-1" }]); + mockReplayBlocks.mockResolvedValue([{ id: rootId, identifier: "TEST-1" }]); + mockExecutionBlocker.mockResolvedValue({ nextAction: "Resume without waking agents until the previous execution stops.", + recoveryActionId: null, cause: "execution_owner_active", runId: "previous-run", agentId: "agent-1" }); const app = await createApp({ type: "board", userId: "user-1", diff --git a/server/src/modules/run-dispatch/adapters/postgres.test.ts b/server/src/modules/run-dispatch/adapters/postgres.test.ts index 5bb8742372..8a4e9a764f 100644 --- a/server/src/modules/run-dispatch/adapters/postgres.test.ts +++ b/server/src/modules/run-dispatch/adapters/postgres.test.ts @@ -2,6 +2,7 @@ import { randomUUID } from "node:crypto"; import { eq, sql } from "drizzle-orm"; import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest"; import { + activityLog, agents, companies, createDb, @@ -21,6 +22,7 @@ import { startEmbeddedPostgresTestDatabase, } from "../../../__tests__/helpers/embedded-postgres.js"; import { createPostgresRunDispatchAdapter } from "./postgres.js"; +import { settleUnrecoverableExecutions } from "../../../services/execution-recovery-resolution.js"; import { getExecutionBlocker } from "../../../services/execution-blocker.js"; // Proves the DB-to-facts mapping this adapter owns for each state the two @@ -49,6 +51,7 @@ describeEmbeddedPostgres("run-dispatch postgres adapter", () => { }, 20_000); afterEach(async () => { + await db.delete(activityLog); await db.delete(issueDocuments); await db.delete(documentRevisions); await db.delete(documents); @@ -770,6 +773,59 @@ describeEmbeddedPostgres("run-dispatch postgres adapter", () => { await expect(adapter.cancelStaleQueuedRun({ companyId, runId, expectedStatus: "queued", now: new Date() })).resolves.toMatchObject({ outcome: "cancelled", errorCode: "execution_reconciliation_required" }); }); + it.each(["active", "resolved"])("allows a new message through a historical %s interruption hold", async status => { + const { companyId, agentId } = await seedCompanyAndAgent(); + await db.update(agents).set({ adapterType: "process" }).where(eq(agents.id, agentId)); + const issueId = randomUUID(), previousRunId = randomUUID(), runId = randomUUID(); + await seedIssue({ companyId, issueId, status: "blocked", assigneeAgentId: agentId }); + await db.insert(heartbeatRuns).values([ + { id: previousRunId, companyId, agentId, status: "interrupted", errorCode: "server_shutdown_interrupted", contextSnapshot: { issueId } }, + { id: runId, companyId, agentId, status: "queued", contextSnapshot: { issueId, wakeReason: "issue_commented" } }, + ]); + await db.insert(heartbeatRunEvents).values({ companyId, agentId, runId: previousRunId, + seq: 1, eventType: "adapter.invoke", payload: { adapterType: "codex_local" } }); + const [action] = await db.insert(issueRecoveryActions).values({ companyId, sourceIssueId: issueId, + kind: "active_run_watchdog", ownerType: "board", cause: "legacy_execution_requires_reconciliation", status, + evidence: { runId: previousRunId, automaticRecovery: { replay: "blocked", actionOutcome: "unknown" } }, + fingerprint: previousRunId, nextAction: "Automatic recovery stopped.", + }).returning(); + expect(await getExecutionBlocker(db, companyId, issueId)).toBeNull(); + const adapter = createPostgresRunDispatchAdapter(db); + expect(await adapter.cancelStaleQueuedRun({ companyId, runId, expectedStatus: "queued", now: new Date() })).toMatchObject({ outcome: "not_stale" }); + await settleUnrecoverableExecutions(db); + const [resolved] = await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.id, action!.id)); + expect(resolved).toMatchObject({ status: "resolved", outcome: "cancelled", evidence: { runId: previousRunId } }); + expect(resolved.evidence.automaticRecovery).toMatchObject({ replay: "conversation_continuation", actionOutcome: "unknown" }); + await settleUnrecoverableExecutions(db); + const audit = await db.select().from(activityLog).where(eq(activityLog.entityId, issueId)); + expect(audit).toHaveLength(1); + expect(audit[0]).toMatchObject({ companyId, action: "issue.execution_recovery_settled" }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, previousRunId))).toHaveLength(0); + // The upgrade does not silently resume historical blocked work. + expect((await db.select().from(issues).where(eq(issues.id, issueId)))[0].status).toBe("blocked"); + }); + + it.each(["process", "http", null])("keeps a historical %s hold after switching to a conversation adapter", async historicalAdapter => { + const { companyId, agentId } = await seedCompanyAndAgent(); + await db.update(agents).set({ adapterType: "codex_local" }).where(eq(agents.id, agentId)); + const issueId = randomUUID(), previousRunId = randomUUID(); + await seedIssue({ companyId, issueId, status: "blocked", assigneeAgentId: agentId }); + await db.insert(heartbeatRuns).values({ id: previousRunId, companyId, agentId, + status: "interrupted", errorCode: "server_shutdown_interrupted", contextSnapshot: { issueId } }); + if (historicalAdapter) await db.insert(heartbeatRunEvents).values({ companyId, agentId, runId: previousRunId, + seq: 1, eventType: "adapter.invoke", payload: { adapterType: historicalAdapter } }); + const [action] = await db.insert(issueRecoveryActions).values({ companyId, sourceIssueId: issueId, + kind: "active_run_watchdog", ownerType: "board", cause: "legacy_execution_requires_reconciliation", status: "active", + evidence: { runId: previousRunId, automaticRecovery: { replay: "blocked", actionOutcome: "unknown" } }, + fingerprint: previousRunId, nextAction: "Inspect previous execution.", + }).returning(); + expect(await getExecutionBlocker(db, companyId, issueId)).toMatchObject({ recoveryActionId: action!.id }); + await settleUnrecoverableExecutions(db); + expect(await getExecutionBlocker(db, companyId, issueId)).toMatchObject({ recoveryActionId: action!.id }); + const [retained] = await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.id, action!.id)); + expect(retained.evidence.automaticRecovery).toMatchObject({ replay: "blocked", actionOutcome: "unknown" }); + }); + it("links the stopped run's agent instead of its return owner, within the same company", async () => { const { companyId, agentId: ownerId } = await seedCompanyAndAgent(); const reviewerId = randomUUID(), issueId = randomUUID(), runId = randomUUID(); diff --git a/server/src/modules/run-dispatch/adapters/postgres.ts b/server/src/modules/run-dispatch/adapters/postgres.ts index 24b8eee2e2..55fa7b2735 100644 --- a/server/src/modules/run-dispatch/adapters/postgres.ts +++ b/server/src/modules/run-dispatch/adapters/postgres.ts @@ -1,9 +1,13 @@ +import { hasConversationContinuationPolicy } from "../../../services/conversation-continuation.js"; import { getExecutionBlocker } from "../../../services/execution-blocker.js"; import { and, asc, eq, gte, inArray, lte, or, sql } from "drizzle-orm"; import type { Db } from "@paperclipai/db"; import { agentWakeupRequests, agents, + approvals, + issueApprovals, + issueThreadInteractions, heartbeatRuns, issueRecoveryActions, issues, @@ -61,6 +65,7 @@ import { RunDispatchApplicationError } from "../application/types.js"; type HeartbeatRun = typeof heartbeatRuns.$inferSelect; type LoadGateFactsInput = { + conversationContinuation: boolean; runId: string; companyId: string; agentId: string; @@ -338,6 +343,21 @@ export function createPostgresRunDispatchAdapter( facts.issueAssigneeAgentId = issue.assigneeAgentId; facts.issueExecutionRunId = issue.executionRunId; facts.issueCheckoutRunId = issue.checkoutRunId; + if (input.conversationContinuation) { + const [interactions, linkedApprovals] = await Promise.all([ + dbOrTx.select({ id: issueThreadInteractions.id }).from(issueThreadInteractions).where(and( + eq(issueThreadInteractions.companyId, input.companyId), + eq(issueThreadInteractions.issueId, issueId), eq(issueThreadInteractions.status, "pending"), + )).limit(1), + dbOrTx.select({ id: approvals.id }).from(issueApprovals).innerJoin(approvals, and( + eq(approvals.id, issueApprovals.approvalId), eq(approvals.companyId, issueApprovals.companyId), + )).where(and( + eq(issueApprovals.companyId, input.companyId), eq(issueApprovals.issueId, issueId), + inArray(approvals.status, ["pending", "revision_requested"]), + )).limit(1), + ]); + facts.pendingResponse = interactions.length > 0 ? "interaction" : linkedApprovals.length > 0 ? "approval" : null; + } facts.reviewParticipant = buildReviewParticipantFacts({ isInReview: issue.status === "in_review", executionState: parseIssueExecutionState(issue.executionState), @@ -402,6 +422,7 @@ export function createPostgresRunDispatchAdapter( runId: run.id, companyId: run.companyId, agentId: run.agentId, + conversationContinuation: run.runtimeMode === "legacy" && hasConversationContinuationPolicy(run.resultJson), contextSnapshot: parseObject(run.contextSnapshot), scheduledRetryReason: run.scheduledRetryReason, retryReasonOverride: input.retryReasonOverride, @@ -685,6 +706,7 @@ export function createPostgresRunDispatchAdapter( runId: run.id, companyId: run.companyId, agentId: run.agentId, + conversationContinuation: run.runtimeMode === "legacy" && hasConversationContinuationPolicy(run.resultJson), contextSnapshot: parseObject(run.contextSnapshot), scheduledRetryReason: run.scheduledRetryReason, retryReasonOverride: run.scheduledRetryReason, @@ -723,6 +745,7 @@ export function createPostgresRunDispatchAdapter( // issue suppresses a max-turn continuation, but every other retry // reason proceeds to promotion anyway. const isLegacyMissingIssueException = + !hasConversationContinuationPolicy(run.resultJson) && !gate.allowed && gate.errorCode === "issue_not_found" && factsResult.facts.retryReasonKind !== "max_turn_continuation" && diff --git a/server/src/modules/run-dispatch/domain/policy.ts b/server/src/modules/run-dispatch/domain/policy.ts index 4b99335ee6..dcdf1c48fd 100644 --- a/server/src/modules/run-dispatch/domain/policy.ts +++ b/server/src/modules/run-dispatch/domain/policy.ts @@ -60,7 +60,8 @@ export type ScheduledRetryGateErrorCode = | "issue_review_participant_changed" | "issue_paused" | "issue_dependencies_blocked" - | "issue_disposition_repair_superseded"; + | "issue_disposition_repair_superseded" + | "issue_waiting_for_response"; export type GateDecision = | { allowed: true } @@ -100,6 +101,8 @@ export type ScheduledRetryFacts = { /** Present only when retryReasonKind is "disposition_repair". */ dispositionRepair: DispositionRepairFacts | null; + /** A conversation retry must wait for unresolved questions and approvals. */ + pendingResponse?: "interaction" | "approval" | null; }; export type QueuedRunStalenessErrorCode = @@ -456,6 +459,16 @@ export function decideScheduledRetryGate( }; } + if (facts.pendingResponse) { + return { + allowed: false, + reason: "Conversation retry is waiting for a response to a pending question or approval", + errorCode: "issue_waiting_for_response", + issueId: facts.issueId, + details: { waitingFor: facts.pendingResponse }, + }; + } + return { allowed: true }; } diff --git a/server/src/modules/wake-queue/adapters/postgres.test.ts b/server/src/modules/wake-queue/adapters/postgres.test.ts index fd8b7e1d4e..f6979079d2 100644 --- a/server/src/modules/wake-queue/adapters/postgres.test.ts +++ b/server/src/modules/wake-queue/adapters/postgres.test.ts @@ -22,6 +22,7 @@ import { createWakeAdmissionWriter, } from "./postgres.js"; import type { WakeQueuePostgresAdapterDeps } from "./postgres.js"; +import { createReleaseIssueExecution } from "../application/use-cases.js"; import type { TransactionScope } from "../application/ports.js"; // Proves the atomicity and company-scope properties the security review @@ -163,6 +164,38 @@ describeEmbeddedPostgres("wake-queue postgres adapter", () => { return id; } + for (const hasDeferredMessage of [false, true]) { + it(`plans conversation recovery during owner cleanup without draining messages (queued=${hasDeferredMessage})`, async () => { + const companyId = await seedCompany(); + const agentId = await seedAgent({ companyId }); + const issueId = await seedIssue({ companyId, assigneeAgentId: agentId }); + const runId = await seedRun({ companyId, agentId, status: "failed", contextSnapshot: { issueId } }); + await db.update(heartbeatRuns).set({ + processPid: process.pid, + resultJson: { conversationContinuation: "continue_conversation_v1" }, + }).where(eq(heartbeatRuns.id, runId)); + const wakeId = hasDeferredMessage ? await seedDeferredWake({ companyId, agentId, issueId }) : null; + const release = createReleaseIssueExecution({ + issueLock: createPostgresWakeQueueAdapter(db, stubDeps), + recovery: { + escalateStrandedAssignedIssue: async () => { throw new Error("unexpected escalation"); }, + escalateStrandedRecoveryIssueInPlace: async () => { throw new Error("unexpected escalation"); }, + }, + }); + const result = await release({ companyId, runId, now: new Date() }); + expect(result.outcome.kind).toBe("released"); + expect(result.postCommitEffects.every((effect) => effect.kind === "conversation_retry_requested")).toBe(true); + if (!wakeId) expect(result.postCommitEffects).toEqual([ + { kind: "conversation_retry_requested", companyId, runId, reviewParticipant: false }, + ]); + if (wakeId) { + const [wake] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, wakeId)); + expect(wake.status).toBe("deferred_issue_execution"); + } + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.companyId, companyId))).toHaveLength(1); + }); + } + it("leaves deferred work untouched until the effective execution hold clears", async () => { const companyId = await seedCompany(); const agentId = await seedAgent({ companyId }); diff --git a/server/src/modules/wake-queue/adapters/postgres.ts b/server/src/modules/wake-queue/adapters/postgres.ts index 4550810ec7..abeccf1b12 100644 --- a/server/src/modules/wake-queue/adapters/postgres.ts +++ b/server/src/modules/wake-queue/adapters/postgres.ts @@ -12,6 +12,7 @@ import { issues, nativeRunFinalizations, } from "@paperclipai/db"; +import { hasConversationContinuationPolicy } from "../../../services/conversation-continuation.js"; import { legacyExecutionNeedsReconciliation } from "../../../services/legacy-execution-recovery.js"; import { authorizeFailedChatRunRetryWake, @@ -85,6 +86,7 @@ function toRunSnapshot(row: HeartbeatRunRow): RunSnapshot { agentId: row.agentId, status: row.status, runtimeMode: row.runtimeMode, + conversationContinuation: row.runtimeMode === "legacy" && hasConversationContinuationPolicy(row.resultJson), errorCode: row.errorCode, responsibleUserId: row.responsibleUserId, contextSnapshot: parseObject(row.contextSnapshot), @@ -1061,12 +1063,18 @@ export function createPostgresWakeQueueAdapter(db: Db, deps: WakeQueuePostgresAd return { outcome: { kind: "released" }, postCommitEffects: [], run: runSnapshot }; } - // A release must leave deferred messages intact while replay is held. - if (await getExecutionBlocker(tx, issueRow.companyId, issueRow.id)) { + // A release must leave deferred messages intact while execution is held. + // The finishing conversation may still own its lease until finally cleanup. + // Allow only bounded retry planning in that case; admission stays gated. + const executionBlocker = await getExecutionBlocker(tx, issueRow.companyId, issueRow.id); + const recoveryOnly = Boolean(executionBlocker && + executionBlocker.cause === "execution_owner_active" && executionBlocker.runId === run.id && + runSnapshot.conversationContinuation && ["failed", "timed_out", "interrupted"].includes(run.status)); + if (executionBlocker && !recoveryOnly) { return { outcome: { kind: "released" }, postCommitEffects: [], run: runSnapshot }; } - const locked: LockedIssueExecution = { primaryIssue: toIssueSnapshot(issueRow), run: runSnapshot }; + const locked: LockedIssueExecution = { primaryIssue: toIssueSnapshot(issueRow), run: runSnapshot, recoveryOnly }; const result = await fn(locked, { host: buildHost(tx, deps), transaction: buildTransaction(tx, deps, db, run) }); return { ...result, run: runSnapshot }; }); diff --git a/server/src/modules/wake-queue/application/ports.ts b/server/src/modules/wake-queue/application/ports.ts index 85a085a94b..0dfc4e8f25 100644 --- a/server/src/modules/wake-queue/application/ports.ts +++ b/server/src/modules/wake-queue/application/ports.ts @@ -12,6 +12,8 @@ export type { InvokableAgentSnapshot, IssueSnapshot, ReleaseRecoveryBlockedNotic /** The primary issue a locked release resolves to, plus the finishing run the lock step already loaded. */ export type LockedIssueExecution = { + /** Plan bounded recovery without draining messages while the finishing owner cleans up. */ + recoveryOnly?: boolean; primaryIssue: IssueSnapshot; run: RunSnapshot; }; diff --git a/server/src/modules/wake-queue/application/types.ts b/server/src/modules/wake-queue/application/types.ts index 90b9df83e4..e2c6cdb4d1 100644 --- a/server/src/modules/wake-queue/application/types.ts +++ b/server/src/modules/wake-queue/application/types.ts @@ -21,6 +21,8 @@ export type RunSnapshot = { contextSnapshot: Record; /** `resultJson.configurationIncomplete`, already parsed; non-null only on a configuration-incomplete failed run. */ configurationIncompletePayload: Record | null; + /** The host schedules failed conversation turns with its durable retry budget. */ + conversationContinuation?: boolean; }; export type IssueSnapshot = { @@ -67,7 +69,12 @@ export type IssueReopenedEffect = { }; /** Explicit post-commit work a caller applies only after the release transaction commits. */ -export type PostCommitEffect = RunQueuedEffect | IssueReopenedEffect; +export type PostCommitEffect = RunQueuedEffect | IssueReopenedEffect | { + kind: "conversation_retry_requested"; + companyId: string; + runId: string; + reviewParticipant: boolean; +}; export type ReleaseOutcome = | { kind: "released" } diff --git a/server/src/modules/wake-queue/application/use-cases.ts b/server/src/modules/wake-queue/application/use-cases.ts index f4a2a9a666..487510a59d 100644 --- a/server/src/modules/wake-queue/application/use-cases.ts +++ b/server/src/modules/wake-queue/application/use-cases.ts @@ -139,6 +139,10 @@ async function runReleaseDrain( const issue = locked.primaryIssue; const postCommitEffects: PostCommitEffect[] = []; + if (locked.recoveryOnly) { + return runReleaseRecoveryTail(issue, run, ports.host, ports.transaction, input, postCommitEffects); + } + // Each `continue` path below leaves the wake row off the // `deferred_issue_execution` status, so the next queue read cannot // return that same row again. That invariant is what ends this loop. @@ -554,6 +558,19 @@ async function runReleaseRecoveryTail( "wake-queue: queued a recovery run with no invokable recovery agent", ); + if (run.conversationContinuation && ["failed", "timed_out", "interrupted"].includes(run.status)) { + // Do not create an uncounted immediate successor inside the issue lock. + // The host's idempotent scheduler claims it after commit with the same + // retry counter used by restart and process-loss recovery. + postCommitEffects.push({ + kind: "conversation_retry_requested", + companyId: run.companyId, + runId: run.id, + reviewParticipant: decision.kind === "queue_review_participant_recovery", + }); + return { outcome: { kind: "released" }, postCommitEffects }; + } + const sessionBefore = await host.resolveSessionBeforeForWakeup({ companyId: issue.companyId, agentId: recoveryAgent.id, diff --git a/server/src/routes/issue-tree-control.ts b/server/src/routes/issue-tree-control.ts index 65939df94d..5a2acae573 100644 --- a/server/src/routes/issue-tree-control.ts +++ b/server/src/routes/issue-tree-control.ts @@ -1,12 +1,11 @@ import { Router } from "express"; import type { Request } from "express"; import { - issueRecoveryActions, issues as issueRows, type Db, } from "@paperclipai/db"; import { and, eq, inArray, isNotNull } from "drizzle-orm"; -import { executionBlockerPredicate } from "../services/execution-blocker.js"; +import { getExecutionBlocker } from "../services/execution-blocker.js"; import { conflict } from "../errors.js"; import { createIssueTreeHoldSchema, @@ -387,30 +386,15 @@ export function issueTreeControlRoutes(db: Db) { .map((member) => member.issueId) : []; if (issueIds.length > 0) { - const [blocked] = await db - .select({ identifier: issueRows.identifier }) - .from(issueRecoveryActions) - .innerJoin( - issueRows, - and( - eq(issueRows.id, issueRecoveryActions.sourceIssueId), - eq(issueRows.companyId, root.companyId), - ), - ) - .where( - and( - eq(issueRecoveryActions.companyId, root.companyId), - inArray(issueRecoveryActions.sourceIssueId, issueIds), - inArray(issueRows.status, RESUME_EXECUTABLE_STATUSES), - isNotNull(issueRows.assigneeAgentId), - executionBlockerPredicate(), - ), - ) - .limit(1); - if (blocked) - throw conflict( - `Cannot wake ${blocked.identifier ?? "this task"} until its stopped execution is reconciled. Resume without waking agents, or review the stopped run first.`, - ); + const candidates = await db.select({ id: issueRows.id, identifier: issueRows.identifier }) + .from(issueRows).where(and( + eq(issueRows.companyId, root.companyId), inArray(issueRows.id, issueIds), + inArray(issueRows.status, RESUME_EXECUTABLE_STATUSES), isNotNull(issueRows.assigneeAgentId), + )); + for (const task of candidates) { + const blocked = await getExecutionBlocker(db, root.companyId, task.id); + if (blocked) throw conflict(`Cannot wake ${task.identifier ?? "this task"}: ${blocked.nextAction}`); + } } } const actor = getActorInfo(req); diff --git a/server/src/services/conversation-continuation.ts b/server/src/services/conversation-continuation.ts new file mode 100644 index 0000000000..5b89201b27 --- /dev/null +++ b/server/src/services/conversation-continuation.ts @@ -0,0 +1,118 @@ +import { and, desc, eq, inArray, isNotNull, or, sql } from "drizzle-orm"; +import { environmentLeases, heartbeatRunEvents, heartbeatRuns, issueRecoveryActions, type Db } from "@paperclipai/db"; +import { readProcessStartedAt } from "./hot-restart.js"; + +// These adapters accept a conversation turn. Retrying a process or webhook can +// replay the action itself, so those adapters retain their recovery contract. +export const CONVERSATION_ADAPTER_TYPES = [ + "claude_local", "codex_local", "cursor", "gemini_local", "opencode_local", + "pi_local", "grok_local", "kimi_local", "hermes_local", +] as const; + +export function isConversationAdapter(adapterType: string): boolean { + return (CONVERSATION_ADAPTER_TYPES as readonly string[]).includes(adapterType); +} + +export const CONVERSATION_CONTINUATION_POLICY = "continue_conversation_v1"; + +export function hasConversationContinuationPolicy(result: Record | null | undefined): boolean { + return result?.conversationContinuation === CONVERSATION_CONTINUATION_POLICY; +} + +function conversationRunPredicate() { + return or( + sql`${heartbeatRuns.resultJson}->>'conversationContinuation' = ${CONVERSATION_CONTINUATION_POLICY}`, + sql`exists ( + select 1 from ${heartbeatRunEvents} + where ${heartbeatRunEvents.companyId} = ${heartbeatRuns.companyId} + and ${heartbeatRunEvents.runId} = ${heartbeatRuns.id} + and ${heartbeatRunEvents.eventType} = 'adapter.invoke' + and ${inArray(sql`${heartbeatRunEvents.payload}->>'adapterType'`, [...CONVERSATION_ADAPTER_TYPES])} + )`, + ); +} + +/** Recovery must not infer the old adapter from the agent's mutable settings. */ +export async function runUsedConversationAdapter(db: Db, run: typeof heartbeatRuns.$inferSelect): Promise { + if (hasConversationContinuationPolicy(run.resultJson)) return true; + const [invocation] = await db.select({ payload: heartbeatRunEvents.payload }).from(heartbeatRunEvents) + .where(and(eq(heartbeatRunEvents.companyId, run.companyId), eq(heartbeatRunEvents.runId, run.id), + eq(heartbeatRunEvents.eventType, "adapter.invoke"))) + .orderBy(desc(heartbeatRunEvents.seq)).limit(1); + const adapterType = invocation?.payload?.adapterType; + return typeof adapterType === "string" && isConversationAdapter(adapterType); +} + +/** Only immutable run evidence can retire a historical conversation hold. + * An agent's current adapter can differ from the one that executed this run. + * Missing evidence retains the hold; the current agent is never a fallback. + */ +export function conversationRecoveryActionPredicate() { + return and( + eq(issueRecoveryActions.cause, "legacy_execution_requires_reconciliation"), + sql`exists ( + select 1 from ${heartbeatRuns} + where ${heartbeatRuns.companyId} = ${issueRecoveryActions.companyId} + and ${heartbeatRuns.id}::text = ${issueRecoveryActions.evidence}->>'runId' + and coalesce(${heartbeatRuns.nativeIssueId}::text, ${heartbeatRuns.contextSnapshot}->>'issueId') = ${issueRecoveryActions.sourceIssueId}::text + and ${heartbeatRuns.runtimeMode} = 'legacy' + and ${inArray(heartbeatRuns.status, ['failed', 'timed_out', 'interrupted', 'cancelled'])} + and ${conversationRunPredicate()} + and ${or( + sql`${heartbeatRuns.resultJson}->>'conversationContinuation' = ${CONVERSATION_CONTINUATION_POLICY}`, + eq(heartbeatRuns.status, "interrupted"), + inArray(heartbeatRuns.errorCode, ["process_lost", "server_shutdown_interrupted", "execution_reconciliation_required"]), + and(eq(heartbeatRuns.status, "cancelled"), sql`${heartbeatRuns.resultJson}->'executionCancellation'->>'state' = 'acknowledged'`), + )} + )`, + ); +} + +/** OS liveness probes do not signal or stop the process. Unknown ownership holds. */ +function processMayBeAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (error) { + return (error as NodeJS.ErrnoException).code !== "ESRCH"; + } +} + +/** A terminal conversation row does not prove that its execution authority ended. + * Other adapters keep their existing bootstrap and ownership protocols. + */ +export async function getConversationOwnershipBlocker(db: Db, companyId: string, issueId: string) { + const activeLease = sql`exists (select 1 from ${environmentLeases} + where ${environmentLeases.companyId} = "heartbeat_runs"."company_id" + and ${environmentLeases.heartbeatRunId} = "heartbeat_runs"."id" + and ${environmentLeases.releasedAt} is null)`; + const candidates = await db.select({ run: heartbeatRuns, activeLease }).from(heartbeatRuns) + .where(and( + eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.runtimeMode, "legacy"), + conversationRunPredicate(), + sql`coalesce(${heartbeatRuns.nativeIssueId}::text, ${heartbeatRuns.contextSnapshot}->>'issueId') = ${issueId}`, + inArray(heartbeatRuns.status, ["failed", "timed_out", "interrupted", "cancelled"]), + or(isNotNull(heartbeatRuns.processPid), isNotNull(heartbeatRuns.processGroupId), activeLease), + )).orderBy(desc(heartbeatRuns.createdAt), desc(heartbeatRuns.id)); + for (const { run, activeLease: leaseHeld } of candidates) { + let pidAlive = run.processPid !== null && processMayBeAlive(run.processPid); + if (pidAlive && run.processStartedAt) { + // A recycled PID cannot keep an old task blocked. An unreadable identity + // stays conservative; the original process may still own execution. + const observed = await readProcessStartedAt(run.processPid!).catch(() => null); + if (observed && new Date(observed).getTime() !== run.processStartedAt.getTime()) pidAlive = false; + } + const groupAlive = run.processGroupId !== null && processMayBeAlive(-run.processGroupId); + if (pidAlive || groupAlive || leaseHeld) { + return { + runId: run.id, + agentId: run.agentId, + cause: "execution_owner_active", + nextAction: pidAlive || groupAlive + ? "The previous provider process is still running. Stop it before continuing this task." + : "The previous execution has not released its environment lease. Wait for cleanup before continuing this task.", + }; + } + } + return null; +} diff --git a/server/src/services/execution-blocker.ts b/server/src/services/execution-blocker.ts index b12ae7f7e8..ad9f46072f 100644 --- a/server/src/services/execution-blocker.ts +++ b/server/src/services/execution-blocker.ts @@ -1,4 +1,5 @@ -import { and, desc, eq, inArray, or, sql } from "drizzle-orm"; +import { and, desc, eq, inArray, not, or, sql } from "drizzle-orm"; +import { conversationRecoveryActionPredicate, getConversationOwnershipBlocker } from "./conversation-continuation.js"; import { z } from "zod"; import { heartbeatRuns, issueRecoveryActions, type Db } from "@paperclipai/db"; import { EXECUTION_RECONCILIATION_CAUSES, type ExecutionBlocker } from "@paperclipai/shared"; @@ -6,6 +7,7 @@ import { EXECUTION_RECONCILIATION_CAUSES, type ExecutionBlocker } from "@papercl /** Resolved recovery bookkeeping can still carry an effective no-replay hold. */ export function executionBlockerPredicate() { return and( + not(conversationRecoveryActionPredicate()!), inArray(issueRecoveryActions.cause, [...EXECUTION_RECONCILIATION_CAUSES]), or(inArray(issueRecoveryActions.status, ["active", "escalated"]), sql`${issueRecoveryActions.evidence}->'automaticRecovery'->>'replay' = 'blocked'`), @@ -13,6 +15,8 @@ export function executionBlockerPredicate() { } export async function getExecutionBlocker(db: Db, companyId: string, issueId: string): Promise { + const ownership = await getConversationOwnershipBlocker(db, companyId, issueId); + if (ownership) return { ...ownership, recoveryActionId: null }; const [action] = await db.select().from(issueRecoveryActions).where(and( eq(issueRecoveryActions.companyId, companyId), eq(issueRecoveryActions.sourceIssueId, issueId), diff --git a/server/src/services/execution-continuation.test.ts b/server/src/services/execution-continuation.test.ts index 12e4fb62d4..ca3f03b8e8 100644 --- a/server/src/services/execution-continuation.test.ts +++ b/server/src/services/execution-continuation.test.ts @@ -133,6 +133,41 @@ const support = await getEmbeddedPostgresTestSupport(); summary: "Notion read completed.", exposeLowTrustRaw: false, }); + it("cancelled admission must not hide the interrupted execution", async () => { + const rejectedId = randomUUID(); + await db.update(heartbeatRuns).set({ status: "interrupted", errorCode: "server_shutdown_interrupted", createdAt: new Date("2026-09-08T10:00:00Z") }).where(eq(heartbeatRuns.id, runId)); + await db.insert(heartbeatRuns).values({ id: rejectedId, companyId, agentId, + status: "cancelled", errorCode: "execution_reconciliation_required", + contextSnapshot: { issueId }, createdAt: new Date("2026-09-08T11:00:00Z") }); + try { + const envelope = await build(); + expect(envelope.interruptedRunId).toBe(runId); + } finally { + await db.delete(heartbeatRuns).where(eq(heartbeatRuns.id, rejectedId)); + await db.update(heartbeatRuns).set({ status: "failed", errorCode: null }).where(eq(heartbeatRuns.id, runId)); + } + }); + + it("preserves the latest user request and adds an interruption notice to fresh and resumed turns", async () => { + await db.update(heartbeatRuns).set({ status: "interrupted", errorCode: "server_shutdown_interrupted" }).where(eq(heartbeatRuns.id, runId)); + try { + const envelope = await buildExecutionContinuation({ db, companyId, issueId, agentId, + context: { retryOfRunId: runId, wakeReason: "retry_failed_run" }, + summary: "Deployment completed. Verification remains.", exposeLowTrustRaw: false }); + expect(envelope.interruptedRunId).toBe(runId); + expect(envelope.objective).toBe("Focus the Gmail summary on launch decisions."); + expect(envelope.messages.map(message => message.id)).toContain(gmailId); + for (const resumedSession of [true, false]) { + const prompt = renderPaperclipWakePrompt({ executionContinuation: envelope }, { resumedSession }); + expect(prompt).toContain("Your previous run was interrupted. Continue from where you left off"); + expect(prompt).toContain("Prior tool calls are history, not commands to replay"); + expect(prompt).toContain("Deployment completed. Verification remains."); + } + } finally { + await db.update(heartbeatRuns).set({ status: "failed", errorCode: null }).where(eq(heartbeatRuns.id, runId)); + } + }); + it("keeps Local CLI run-authored comments as history without promoting them to human direction", async () => { const id = randomUUID(); await db.insert(issueComments).values({ id, companyId, issueId, authorType: "user", diff --git a/server/src/services/execution-continuation.ts b/server/src/services/execution-continuation.ts index 307d925630..524005af8a 100644 --- a/server/src/services/execution-continuation.ts +++ b/server/src/services/execution-continuation.ts @@ -9,6 +9,7 @@ import { } from "@paperclipai/db"; import type { ExecutionContinuationEnvelope } from "@paperclipai/shared"; import { sanitizeQuarantinedCommentForHigherTrust } from "./source-trust.js"; +import { hasConversationContinuationPolicy } from "./conversation-continuation.js"; const object = (v: unknown): Record => v && typeof v === "object" && !Array.isArray(v) @@ -208,7 +209,7 @@ export async function buildExecutionContinuation(input: { row.authorType === "user" && !row.createdByRunId && !row.deleted && row.body.trim().length > 0, ); const priorRuns = await db - .select({ id: heartbeatRuns.id, result: heartbeatRuns.resultJson }) + .select({ id: heartbeatRuns.id, result: heartbeatRuns.resultJson, status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode }) .from(heartbeatRuns) .where( and( @@ -249,7 +250,16 @@ export async function buildExecutionContinuation(input: { eq(issueRecoveryActions.status, "resolved"), ), ); + const lastTerminal = priorRuns.findLast((run) => + ["succeeded", "failed", "timed_out", "interrupted", "cancelled"].includes(run.status) && + !(run.status === "cancelled" && run.errorCode === "execution_reconciliation_required"), + ); + const interruptedRunId = lastTerminal && lastTerminal.status !== "succeeded" && + (hasConversationContinuationPolicy(lastTerminal.result) || + lastTerminal.status === "interrupted" || lastTerminal.errorCode === "process_lost") + ? lastTerminal.id : undefined; return { + ...(interruptedRunId ? { interruptedRunId } : {}), ...(resumeDelta ? { resumeDelta } : {}), recoveryOutcomes: reconciliations .filter((row) => row.evidence.executionReconciliation) diff --git a/server/src/services/execution-recovery-resolution.ts b/server/src/services/execution-recovery-resolution.ts index 0d569e0d56..deabf16880 100644 --- a/server/src/services/execution-recovery-resolution.ts +++ b/server/src/services/execution-recovery-resolution.ts @@ -1,8 +1,9 @@ import { randomUUID } from "node:crypto"; +import { conversationRecoveryActionPredicate, getConversationOwnershipBlocker } from "./conversation-continuation.js"; import { persistActivity } from "./activity-log.js"; import { appendHeartbeatRunEvent } from "./heartbeat-run-events.js"; import { logger } from "../middleware/logger.js"; -import { and, eq, inArray, isNull, or, sql } from "drizzle-orm"; +import { and, eq, inArray, isNull, not, or, sql } from "drizzle-orm"; import { chatActions, environmentLeases, @@ -305,8 +306,47 @@ export async function settleUnrecoverableExecutions( now = new Date(), options: { failpoint?: (phase: "persisted") => void } = {}, ) { + // Fold obsolete conversation holds without waking historical work on upgrade. + // Keep their evidence and record the policy change in the task's activity log. + const obsoleteConversationHold = and( + conversationRecoveryActionPredicate(), + or( + inArray(issueRecoveryActions.status, ["active", "escalated"]), + sql`${issueRecoveryActions.evidence}->'automaticRecovery'->>'replay' = 'blocked'`, + ), + ); + await db.transaction(async tx => { + const foldable = await tx.select().from(issueRecoveryActions).where(obsoleteConversationHold) + .limit(25).for("update", { skipLocked: true }); + for (const candidate of foldable) { + if (await getConversationOwnershipBlocker(tx as unknown as Db, candidate.companyId, candidate.sourceIssueId)) continue; + const [action] = await tx.update(issueRecoveryActions).set({ + status: "resolved", + outcome: "cancelled", + resolvedAt: now, + updatedAt: now, + nextAction: "Automatic attempts stopped. Send a new message to continue the conversation.", + resolutionNote: "Conversation continuation does not replay prior tool calls.", + wakePolicy: null, + monitorPolicy: null, + evidence: sql`case when ${issueRecoveryActions.evidence} ? 'automaticRecovery' + then jsonb_set(${issueRecoveryActions.evidence}, '{automaticRecovery,replay}', '"conversation_continuation"'::jsonb) + else ${issueRecoveryActions.evidence} end`, + }).where(and(obsoleteConversationHold, eq(issueRecoveryActions.id, candidate.id))).returning(); + if (!action) continue; + await persistActivity(tx as unknown as Db, { + companyId: action.companyId, + actorType: "system", + actorId: "execution-recovery", + action: "issue.execution_recovery_settled", + entityType: "issue", + entityId: action.sourceIssueId, + details: { recoveryActionId: action.id, outcome: "cancelled", continuation: "conversation" }, + }); + } + }); // Filter eligibility before applying the batch limit. A queue of sessions - // still awaiting safe replacement must not starve settled incidents behind it. + // awaiting replacement must not starve settled incidents behind it. const candidates = await db .select({ action: issueRecoveryActions }) .from(issueRecoveryActions) @@ -327,6 +367,7 @@ export async function settleUnrecoverableExecutions( ) .where( and( + not(conversationRecoveryActionPredicate()!), inArray(issueRecoveryActions.status, ["active", "escalated"]), eq(issueRecoveryActions.kind, "active_run_watchdog"), inArray(issueRecoveryActions.cause, [ diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index fe6da253eb..e6bbae1c3b 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -1,4 +1,5 @@ import { getExecutionBlocker } from "./execution-blocker.js"; +import { CONVERSATION_CONTINUATION_POLICY, runUsedConversationAdapter, hasConversationContinuationPolicy, isConversationAdapter } from "./conversation-continuation.js"; import { recordExecutionWait } from "./execution-wait.js"; import { legacyExecutionNeedsReconciliation, @@ -9420,7 +9421,18 @@ export function heartbeatService( // issue's activity entry. async function applyWakeQueuePostCommitEffects(effects: WakeQueuePostCommitEffect[]) { for (const effect of effects) { - if (effect.kind === "run_queued") { + if (effect.kind === "conversation_retry_requested") { + const [source] = await db.select().from(heartbeatRuns).where(and( + eq(heartbeatRuns.companyId, effect.companyId), eq(heartbeatRuns.id, effect.runId), + )); + const agent = source ? await getAgent(source.agentId) : null; + if (source && agent && agent.companyId === source.companyId) { + await scheduleBoundedRetryForRun(source, agent, effect.reviewParticipant ? { + retryReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, + wakeReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_WAKE_REASON, + } : undefined); + } + } else if (effect.kind === "run_queued") { publishLiveEvent({ companyId: effect.run.companyId, type: "heartbeat.run.queued", @@ -14472,7 +14484,7 @@ export function heartbeatService( restartSuspendedRunIds.push(run.id); continue; } - const message = `Interrupted by graceful server shutdown (${signal}); recovery requires verified provider continuity`; + const message = `Interrupted by graceful server shutdown (${signal})`; const running = runningProcesses.get(run.id); try { if (run.runtimeMode === "native") { @@ -14510,6 +14522,7 @@ export function heartbeatService( errorCode: "server_shutdown_interrupted", signal, resultJson: mergeRunStopMetadataForAgent(agent, "interrupted", { + conversationContinuationEligible: await runUsedConversationAdapter(db, run), resultJson: persistedCancellationResult, errorCode: "server_shutdown_interrupted", errorMessage: message, @@ -14732,6 +14745,7 @@ export function heartbeatService( : baseSchedule; const requiresIssueGate = + hasConversationContinuationPolicy(run.resultJson) || retryReason === MAX_TURN_CONTINUATION_RETRY_REASON || retryReason === INTERACTION_CONTINUATION_INFRA_RETRY_REASON; if (requiresIssueGate) { @@ -15179,6 +15193,8 @@ export function heartbeatService( status: "scheduled_retry", wakeupRequestId: wakeupRequest.id, contextSnapshot: retryContextSnapshot, + ...(hasConversationContinuationPolicy(run.resultJson) + ? { resultJson: { conversationContinuation: CONVERSATION_CONTINUATION_POLICY } } : {}), responsibleUserId, sessionIdBefore: sessionBefore, retryOfRunId: run.id, @@ -15302,6 +15318,7 @@ export function heartbeatService( .update(issues) .set({ executionRunId: scheduledRun.id, + checkoutRunId: sql`case when ${issues.checkoutRunId} = ${run.id} then null else ${issues.checkoutRunId} end`, executionAgentNameKey: normalizeAgentNameKey(agent.name), executionLockedAt: now, ...(detachWorkspaceFromIssue @@ -17213,6 +17230,7 @@ export function heartbeatService( outcome: "succeeded" | "interrupted" | "failed" | "cancelled" | "timed_out", options?: { resultJson?: Record | null; + conversationContinuationEligible?: boolean; errorCode?: string | null; errorMessage?: string | null; }, @@ -17224,10 +17242,17 @@ export function heartbeatService( errorCode: options?.errorCode ?? null, errorMessage: options?.errorMessage ?? null, }); - return mergeHeartbeatRunStopMetadata( + const result = mergeHeartbeatRunStopMetadata( options?.resultJson ?? null, stopMetadata, ); + const cancellationAcknowledged = + parseObject(result?.executionCancellation).state === "acknowledged"; + return options?.conversationContinuationEligible !== false && outcome !== "succeeded" && + (outcome !== "cancelled" || cancellationAcknowledged) && + isConversationAdapter(agent.adapterType) + ? { ...result, conversationContinuation: CONVERSATION_CONTINUATION_POLICY } + : result; } function countValue(value: unknown) { @@ -18246,6 +18271,7 @@ export function heartbeatService( (!!run.processPid || !!run.processGroupId)) || monitorDispatchLostWithoutFutureWake); const baseMessage = buildProcessLossMessage(run); + const conversationContinuationEligible = await runUsedConversationAdapter(db, run); const failureWrite = await setRunStatusFromLive( run.id, @@ -18260,6 +18286,7 @@ export function heartbeatService( { adapterType, adapterConfig }, "failed", { + conversationContinuationEligible, resultJson: parseObject(run.resultJson), errorCode: "process_lost", errorMessage: shouldRetry diff --git a/server/src/services/legacy-execution-recovery.test.ts b/server/src/services/legacy-execution-recovery.test.ts index d1419c317d..3567091ffb 100644 --- a/server/src/services/legacy-execution-recovery.test.ts +++ b/server/src/services/legacy-execution-recovery.test.ts @@ -29,3 +29,12 @@ it("retains the hold until the provider actually acknowledges cancellation", () ...stopped.resultJson, executionCancellation: { state: "requested" }, } })).toBe(true); }); + + it("continues a conversation without requiring receipts, even after automatic attempts are exhausted", () => { + for (const status of ["failed", "timed_out", "interrupted", "cancelled"]) { + expect(legacyExecutionNeedsReconciliation({ + runtimeMode: "legacy", status, errorCode: "process_lost", scheduledRetryAttempt: 2, + resultJson: { conversationContinuation: "continue_conversation_v1" }, + })).toBe(false); + } +}); diff --git a/server/src/services/legacy-execution-recovery.ts b/server/src/services/legacy-execution-recovery.ts index 1785b12d60..3b9eb80f74 100644 --- a/server/src/services/legacy-execution-recovery.ts +++ b/server/src/services/legacy-execution-recovery.ts @@ -1,4 +1,5 @@ import { normalizeMaxTurnStopReason } from "./heartbeat-stop-metadata.js"; +import { hasConversationContinuationPolicy } from "./conversation-continuation.js"; import { randomUUID } from "node:crypto"; import { and, eq, inArray, sql } from "drizzle-orm"; import { heartbeatRuns, issueRecoveryActions, issues, type Db } from "@paperclipai/db"; @@ -18,6 +19,9 @@ export function legacyExecutionNeedsReconciliation( !["failed", "timed_out", "interrupted", "cancelled"].includes(run.status) ) return false; + // A fresh conversation turn lets the agent decide what remains. The retry + // scheduler, not an action-outcome hold, owns the automatic attempt limit. + if (hasConversationContinuationPolicy(run.resultJson)) return false; // Productive turn-budget continuation is not a failed provider session. if (normalizeMaxTurnStopReason(run.resultJson?.stopReason) ?? normalizeMaxTurnStopReason(run.errorCode)) return false; const evidence = run.resultJson?.executionRecovery as diff --git a/tests/e2e/acp-stop-continuation.spec.ts b/tests/e2e/acp-stop-continuation.spec.ts index 6772749792..c23ee8e3c1 100644 --- a/tests/e2e/acp-stop-continuation.spec.ts +++ b/tests/e2e/acp-stop-continuation.spec.ts @@ -10,7 +10,7 @@ async function json(response: APIResponse) { } for (const { unfinishedWrite, pause } of [{ unfinishedWrite: false, pause: false }, { unfinishedWrite: true, pause: false }, { unfinishedWrite: false, pause: true }]) { - test(`embedded ACP Stop: ${unfinishedWrite ? "unknown action stays visibly blocked" : pause ? "composer pause requires Resume before continuation" : "go continues the same session with queued input"}`, async ({ page, request }) => { + test(`embedded ACP Stop: ${unfinishedWrite ? "unknown action continues without replaying the write" : pause ? "composer pause requires Resume before continuation" : "go continues the same session with queued input"}`, async ({ page, request }) => { test.setTimeout(120_000); const root = await mkdtemp(path.join(os.tmpdir(), "paperclip-stop-browser-")); const company = await json(await request.post("/api/companies", { data: { name: `ACP Stop ${Date.now()}` } })); @@ -53,7 +53,6 @@ for (const { unfinishedWrite, pause } of [{ unfinishedWrite: false, pause: false expect(stopped.resultJson.executionCancellation.state).toBe("acknowledged"); const writesAtStop = unfinishedWrite ? await readFile(path.join(root, "writes"), "utf8") : null; await page.reload(); - if (unfinishedWrite) await expect(page.getByText("Work cannot start.", { exact: false })).toBeVisible(); if (pause) { await expect(page.getByTestId("paused-composer-takeover")).toBeVisible(); await expect(editor).toHaveCount(0); @@ -70,39 +69,20 @@ for (const { unfinishedWrite, pause } of [{ unfinishedWrite: false, pause: false await editor.fill("go"); await page.getByRole("button", { name: "Send", exact: true }).click(); } - if (unfinishedWrite) { - await expect(page.getByText("Work cannot start.", { exact: false })).toBeVisible(); - await expect(editor).toHaveText(""); - // Repeated user wakes preserve input without manufacturing attempts. - for (const message of ["go again", "still waiting"]) { - await editor.fill(message); - const response = page.waitForResponse((candidate) => - candidate.request().method() === "POST" && candidate.url().endsWith(`/api/issues/${issue.identifier}/comments`)); - await page.getByRole("button", { name: "Send", exact: true }).click(); - expect((await response).ok()).toBe(true); - await expect(editor).toHaveText(""); - } - expect((await json(await request.get(`/api/issues/${issue.id}`))).executionBlocker).toBeTruthy(); - await page.waitForTimeout(1000); - await expect(page.getByText("Couldn't start", { exact: false })).toHaveCount(0); - expect(await json(await request.get(`/api/issues/${issue.id}/runs`))).toHaveLength(1); - expect(await readFile(path.join(root, "writes"), "utf8")).toBe(writesAtStop); - expect((await readFile(path.join(root, "prompts"), "utf8")).trim().split("\n")).toHaveLength(1); - } else { - await expect(page.getByText("Answered the pending follow-up once.", { exact: false })).toBeVisible({ timeout: 30_000 }); - await expect.poll(async () => (await json(await request.get(`/api/issues/${issue.id}/live-runs`))).length).toBe(0); - const prompts = (await readFile(path.join(root, "prompts"), "utf8")).trim().split("\n").map(line => JSON.parse(line)); - expect(prompts).toHaveLength(2); - expect(new Set(prompts.map(prompt => prompt.sessionId)).size).toBe(1); - // Resume delivers the queued follow-up in the same provider session. - const continuationPrompts = pause ? prompts.slice(1) : [prompts.at(-1)]; - expect(JSON.stringify(continuationPrompts)).toContain("List my recent Drive files."); - if (!pause) expect(JSON.stringify(continuationPrompts)).toContain("go"); - expect(await readFile(path.join(root, "completed"), "utf8")).toBe("follow-up\n"); - const completedIssue = await json(await request.get(`/api/issues/${issue.id}`)); - expect(completedIssue.executionBlocker).toBeNull(); - expect(completedIssue.status).toBe("done"); - } + await expect(page.getByText("Answered the pending follow-up once.", { exact: false })).toBeVisible({ timeout: 30_000 }); + await expect.poll(async () => (await json(await request.get(`/api/issues/${issue.id}/live-runs`))).length).toBe(0); + const prompts = (await readFile(path.join(root, "prompts"), "utf8")).trim().split("\n").map(line => JSON.parse(line)); + expect(prompts).toHaveLength(2); + expect(new Set(prompts.map(prompt => prompt.sessionId)).size).toBe(1); + // Resume delivers the queued follow-up in the same provider session. + const continuationPrompts = pause ? prompts.slice(1) : [prompts.at(-1)]; + expect(JSON.stringify(continuationPrompts)).toContain("List my recent Drive files."); + if (!pause) expect(JSON.stringify(continuationPrompts)).toContain("go"); + expect(await readFile(path.join(root, "completed"), "utf8")).toBe("follow-up\n"); + const completedIssue = await json(await request.get(`/api/issues/${issue.id}`)); + expect(completedIssue.executionBlocker).toBeNull(); + expect(completedIssue.status).toBe("done"); + if (unfinishedWrite) expect(await readFile(path.join(root, "writes"), "utf8")).toBe(writesAtStop); await expect(page.getByRole("dialog")).toHaveCount(0); } finally { await request.patch(`/api/companies/${company.id}`, { data: { status: "archived" } });