diff --git a/packages/paperclip-runner/src/live/live-session.test.ts b/packages/paperclip-runner/src/live/live-session.test.ts index f0fa761ee3..5a1478f770 100644 --- a/packages/paperclip-runner/src/live/live-session.test.ts +++ b/packages/paperclip-runner/src/live/live-session.test.ts @@ -66,6 +66,8 @@ interface FakeProviderState { usageRunDelta: Record | null; onTurnStart?: () => Promise; onUsage?: (queue: AsyncNotifications, turnId: string) => void | Promise; + /** Delays the fake `turn/interrupt` reply, to model a slow transport round trip. */ + interruptDelayMs: number; } class FakeCapabilityCodexTransport implements CodexAppServerTransport { @@ -132,6 +134,9 @@ class FakeCapabilityCodexTransport implements CodexAppServerTransport { return { turn: { id: turnId, status: "inProgress" } }; } if (method === "turn/interrupt") { + if (this.state.interruptDelayMs > 0) { + await new Promise((resolve) => setTimeout(resolve, this.state.interruptDelayMs)); + } const turnId = String(params.turnId); this.state.turns.set(turnId, "interrupted"); this.notificationsQueue.push({ @@ -318,6 +323,7 @@ function providerState(): FakeProviderState { holdAfterTool: false, closeError: null, usageRunDelta: null, + interruptDelayMs: 0, }; } @@ -1312,6 +1318,47 @@ describe("Capability live runnerd and Codex session", () => { }); }); + it("does not raise an unhandled rejection when interrupt() outlasts the turn timeout during reconcileActiveTurn", async () => { + const state = providerState(); + const store = new InMemoryCapabilityLiveSessionStore(); + const firstService = new CapabilityLiveSessionService({ + store, + transportFactory: fakeTransportFactory(state), + }); + const first = await firstService.create({ + runId: "run-slow-interrupt-reconcile", + sessionId: "session-slow-interrupt-reconcile", + attemptId: "attempt-slow-interrupt-killed", + turnTimeoutMs: 20, + }); + state.holdAfterTool = true; + const killedTurn = captureTurnRejection(first.sendMessage("Apply idempotent progress once.")); + await vi.waitFor(async () => { + expect((await store.load(first.id))?.mockState).toContain("progress-governed-once"); + }); + await state.transports[0]!.close(); + await expect(killedTurn).resolves.toMatchObject({ message: expect.stringContaining("timed out") }); + + const resumedService = new CapabilityLiveSessionService({ + store, + transportFactory: fakeTransportFactory(state), + }); + const resumed = await resumedService.resume({ + sessionId: first.id, + attemptId: "attempt-slow-interrupt-resumed", + resumeOf: "attempt-slow-interrupt-killed", + }); + + // Make the fake transport's turn/interrupt reply outlast the 20 ms turn + // timeout. reconcileActiveTurn() awaits interrupt() first, so the turn + // waiter's timer can reject before interrupt() resolves. The rejection + // handler must already be in place at that moment; otherwise Node + // reports an unhandled rejection and Vitest fails the whole file, even + // though the assertion below is correct. + state.interruptDelayMs = 200; + await expect(resumed.reconcileActiveTurn()).rejects.toThrow(/timed out/); + }); + it("persists a resumed turnTimeoutMs override so a later resume that omits it keeps the value", async () => { const state = providerState(); const store = new InMemoryCapabilityLiveSessionStore(); diff --git a/packages/paperclip-runner/src/live/live-session.ts b/packages/paperclip-runner/src/live/live-session.ts index 8cc9c15c2a..821352e71f 100644 --- a/packages/paperclip-runner/src/live/live-session.ts +++ b/packages/paperclip-runner/src/live/live-session.ts @@ -1550,7 +1550,6 @@ export class CapabilityLiveSession { // Arm the provider timeout only after bounded preflight succeeds. The // admission token excludes concurrent sends before this point. const terminal = this.#armTurnWaiter(); - void terminal.catch(() => undefined); try { response = await admission.transport.request("turn/start", { threadId: this.#providerThreadId, @@ -1996,7 +1995,7 @@ export class CapabilityLiveSession { #armTurnWaiter(): Promise> { if (this.#turnWaiter !== null) throw new Error("Capability live session already has a turn waiter"); - return new Promise>((resolve, reject) => { + const terminal = new Promise>((resolve, reject) => { const timer = setTimeout(() => { const waiter = this.#turnWaiter; this.#turnWaiter = null; @@ -2011,6 +2010,13 @@ export class CapabilityLiveSession { }, this.#config.turnTimeoutMs); this.#turnWaiter = { resolve, reject, timer, assistantText: "", draftId: null }; }); + // A caller may await this promise only after another `await` of its own + // (see `reconcileActiveTurn`). The timer above can reject before that + // point, so attach a no-op handler here, at creation, on every call + // site. `.catch()` returns a new promise; the original stays rejected + // and a later `await terminal` still observes it. + terminal.catch(() => undefined); + return terminal; } /** Interrupts and durably reconciles a checkpointed active turn after restart. */ 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 9fbbdff3b7..8199bbbd75 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts @@ -7341,6 +7341,47 @@ it("probes an exact-authority resume and confirms its live provider identity", a } }, 30_000); +it("still fails closed when a real close grace period cannot fit a durable suspension round trip", async () => { + const stateDirectory = await mkdtemp( + join(tmpdir(), "runnerd-close-grace-too-small-"), + ); + const identity = { + runnerInstanceId: "runner-close-grace-too-small", + environmentLeaseId: "lease-close-grace-too-small", + runId: "run-close-grace-too-small", + normalizedSessionId: "session-close-grace-too-small", + turnId: "turn-close-grace-too-small", + itemId: "item-close-grace-too-small", + }; + const bundle = createCapabilityRunnerdCodexTransport({ + runnerBinary: defaultCapabilityRunnerdBinary(), + codexCommand: fakeCodex, + codexArgs: fakeCodexArgs(stateDirectory), + stateDirectory, + // No real durable command round trip can complete this fast. A wider + // budget for the provider-drain proof must not turn this barrier into + // one that always passes; it still needs the actual proof to arrive. + closeGraceMs: 1, + lifecyclePolicy: { mode: "per_turn", idleTimeoutMs: null }, + prpIdentity: identity, + }); + bundle.transport.setServerRequestHandler(async () => ({ + success: true, + contentItems: [], + })); + try { + await bundle.transport.request("thread/start", { + cwd: tmpdir(), + dynamicTools: [], + }); + await expect(bundle.transport.close()).rejects.toThrow( + "runner did not durably suspend before checkpoint", + ); + } finally { + await rm(stateDirectory, { recursive: true, force: true }); + } +}, 30_000); + it("cold-restores a suspended provider session under its durable run binding", async () => { const stateDirectory = await mkdtemp(join(tmpdir(), "runnerd-cold-attach-")); const tracePath = join(stateDirectory, "provider-trace.ndjson"); diff --git a/packages/paperclip-runner/src/live/runnerd-codex-transport.ts b/packages/paperclip-runner/src/live/runnerd-codex-transport.ts index 7fc5354161..f754359c53 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.ts @@ -3828,9 +3828,7 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { } } - async #stopActiveProviderTurnBeforeSuspend( - deadline: number, - ): Promise { + async #stopActiveProviderTurnBeforeSuspend(deadline: number): Promise { const state = this.#providerDrainState(); const core = this.#core; const inferredActiveProviderTurnId = @@ -3848,7 +3846,7 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { activeProviderTurnId === null || core === null ) { - return false; + return; } const commandId = `command_close_stop_${randomUUID().replaceAll("-", "")}`; core.queueCommand( @@ -3866,19 +3864,18 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { this.#diagnostic( `stopped active provider turn ${activeProviderTurnId} before runner suspension`, ); - return true; + return; } if (command !== undefined && command.status !== "pending") { this.#diagnostic( `provider turn stop ${command.status} before runner suspension`, ); - return false; + return; } - if (await this.#runnerHasExited()) return false; + if (await this.#runnerHasExited()) return; await new Promise((resolveWait) => setTimeout(resolveWait, 5)); } this.#diagnostic("provider turn stop timed out before runner suspension"); - return false; } async #drainSettledProviderEventsBeforeSuspend( @@ -3982,13 +3979,15 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { this.#pumpEventsSafely(); await new Promise((resolveWait) => setTimeout(resolveWait, 5)); } - const stoppedActiveTurn = - await this.#stopActiveProviderTurnBeforeSuspend(preparationDeadline); + // The drain always needs one real command round trip to the runner + // process, whether or not a turn was active: stopping an active + // turn only changes how much trailing event traffic that round + // trip may need to carry. Give both cases the same budget so a + // slow-but-idle runner is not held to a tighter deadline than a + // runner that just stopped a turn. + await this.#stopActiveProviderTurnBeforeSuspend(preparationDeadline); providerDrained = await this.#drainSettledProviderEventsBeforeSuspend( - Math.min( - stoppedActiveTurn ? 5_000 : 1_000, - Math.max(0, preparationDeadline - Date.now()), - ), + Math.min(5_000, Math.max(0, preparationDeadline - Date.now())), ); } // Local durable roots are reused too. Process exit alone cannot prove