From 04f74b4732845d1534658eecb4a1c592f713769d Mon Sep 17 00:00:00 2001 From: Dotta Date: Sat, 12 Sep 2026 20:29:32 -0500 Subject: [PATCH] Persist replacement Codex process ownership before recovery succeeds Co-Authored-By: Paperclip --- .../src/live/runnerd-codex-transport.test.ts | 71 +++++++++++++++++++ .../src/live/runnerd-codex-transport.ts | 14 ++-- 2 files changed, 80 insertions(+), 5 deletions(-) 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 e1ed1b2383..69532d1918 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts @@ -2085,6 +2085,77 @@ it("publishes spawned runner ownership before waiting for provider startup", asy } }); +it.each([false, true])("persists each recovered runner identity before declaring recovery complete (save fails: %s)", async (failOwnershipSave) => { + const stateDirectory = await mkdtemp(join(tmpdir(), "runner-replacement-ownership-")); + const handles: durableControlPlane.RunnerProcessHandle[] = []; + const realSpawn = durableControlPlane.spawnRunner; + const wrap = (handle: durableControlPlane.RunnerProcessHandle): durableControlPlane.RunnerProcessHandle => { + const wrapped = { + ...handle, + startedAt: `2026-09-01T10:00:0${handles.length}.000Z`, + restart: handle.restart ? (ticket: string) => wrap(handle.restart!(ticket)) : undefined, + }; + handles.push(wrapped); + return wrapped; + }; + const spawnSpy = vi.spyOn(durableControlPlane, "spawnRunner").mockImplementation((options) => wrap(realSpawn(options))); + let releaseOwnership!: () => void; + const ownershipBarrier = new Promise((resolve) => { releaseOwnership = resolve; }); + let durableOwner: { pid: number; processGroupId: number | null; startedAt: string } | null = null; + const onSpawn = vi.fn(async (owner: NonNullable) => { + if (handles.length === 2) { + await ownershipBarrier; + if (failOwnershipSave) throw new Error("fixture durable ownership unavailable"); + } + durableOwner = structuredClone(owner); + }); + const diagnostics: string[] = []; + const bundle = createCapabilityRunnerdCodexTransport({ + runnerBinary: defaultCapabilityRunnerdBinary(), codexCommand: fakeCodex, + codexArgs: fakeCodexArgs(stateDirectory), stateDirectory, + lifecyclePolicy: { mode: "warm", idleTimeoutMs: 60_000 }, + runnerReconnectGraceMs: 5_000, onSpawn, + onDiagnostic: (message) => diagnostics.push(message), + }); + try { + await bundle.transport.request("thread/start", { cwd: tmpdir(), dynamicTools: codexSemanticToolSpecs() }); + const originalOwner = structuredClone(durableOwner); + handles[0]!.child.kill("SIGKILL"); + await vi.waitFor(() => expect(onSpawn).toHaveBeenCalledTimes(2)); + expect(durableOwner).toEqual(originalOwner); + expect(diagnostics).not.toContain("runner process restored its durable PRP session"); + releaseOwnership(); + if (failOwnershipSave) { + await vi.waitFor(() => expect(diagnostics).toContain("native_runner_process_ownership_failed: fixture durable ownership unavailable")); + expect(durableOwner).toEqual(originalOwner); + expect(diagnostics).not.toContain("runner process restored its durable PRP session"); + await expect(bundle.transport.request("thread/start", { cwd: tmpdir() })).rejects.toThrow("native_runner_process_ownership_failed"); + expect(handles).toHaveLength(2); + return; + } + await vi.waitFor(() => expect(diagnostics).toContain("runner process restored its durable PRP session")); + expect(durableOwner).toEqual({ pid: handles[1]!.child.pid, processGroupId: handles[1]!.processGroupId ?? null, startedAt: handles[1]!.startedAt }); + expect(bundle.transport.processInfo?.()).toMatchObject(durableOwner!); + expect(durableOwner).not.toEqual(originalOwner); + // A second recovery must advance durable ownership again, not resurrect + // either the first process or its original start timestamp. + handles[1]!.child.kill("SIGKILL"); + await vi.waitFor(() => expect(onSpawn).toHaveBeenCalledTimes(3)); + await vi.waitFor(() => expect(diagnostics.filter((message) => message === "runner process restored its durable PRP session")).toHaveLength(2)); + expect(durableOwner).toEqual({ pid: handles[2]!.child.pid, processGroupId: handles[2]!.processGroupId ?? null, startedAt: handles[2]!.startedAt }); + expect(bundle.transport.processInfo?.()).toMatchObject(durableOwner!); + } finally { + releaseOwnership(); + await bundle.transport.close().catch(() => undefined); + for (const handle of handles) { + if (handle.child.exitCode === null) handle.child.kill("SIGKILL"); + await handle.completion.catch(() => undefined); + } + spawnSpy.mockRestore(); + await rm(stateDirectory, { recursive: true, force: true }); + } +}, 15_000); + it("launches runnerd with its production durable outbox limits", () => { expect(runnerdLaunchProfileInternals.maxOutboxBytes).toBe(16 * 1024 * 1024); expect(runnerdLaunchProfileInternals.p0ReserveBytes).toBe(1024 * 1024); diff --git a/packages/paperclip-runner/src/live/runnerd-codex-transport.ts b/packages/paperclip-runner/src/live/runnerd-codex-transport.ts index 03860fdede..b97e566316 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.ts @@ -6472,11 +6472,6 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { this.#evidence.runnerExited = false; this.#evidence.runnerExitCode = null; this.#evidence.runnerSignal = null; - this.#evidence.runnerPid = recoveredHandle.child.pid ?? null; - this.#evidence.runnerProcessGroupId = - recoveredHandle.processGroupId ?? null; - this.#publish(); - let processSettled = false; const completion = recoveredHandle.completion.then( (result) => { @@ -6490,6 +6485,15 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { return false; }, ); + try { + await this.#publishSpawnedProcess(recoveredHandle); + } catch (error) { + this.#failTransport(new Error( + `native_runner_process_ownership_failed: ${error instanceof Error ? error.message : String(error)}`, + )); + return; + } + if (this.#closed || this.#failure !== null) return; const authenticated = (async () => { while ( !processSettled &&