From 7a6f5b6daf1b9a673eb301237e69e2aa93b3028b Mon Sep 17 00:00:00 2001 From: Dotta Date: Sat, 5 Sep 2026 12:26:40 -0500 Subject: [PATCH] fix(runner): harden warm authority recovery --- .../runner-core/src/acpx_provider_backend.rs | 7 + .../crates/runner-core/src/durable/runner.rs | 7 + .../src/managed_provider_backend.rs | 4 + .../src/native_provider_backend.rs | 15 ++ .../runner-core/src/provider_backend.rs | 4 + .../src/live/runnerd-codex-transport.test.ts | 93 +++++++++++++ .../src/live/runnerd-codex-transport.ts | 33 +++-- .../__tests__/native-workspace-sync.test.ts | 130 ++++++++++++++++-- .../native-runtime/native-workspace-sync.ts | 50 +++++-- 9 files changed, 310 insertions(+), 33 deletions(-) diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/acpx_provider_backend.rs b/packages/paperclip-runner/runner/crates/runner-core/src/acpx_provider_backend.rs index 402b71a094..92e7323b49 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/acpx_provider_backend.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/acpx_provider_backend.rs @@ -1345,6 +1345,13 @@ impl CommandExecutor for AcpxCommandExecutor { } } + fn rotate_authority(&mut self, config: &DurableRunnerConfig) { + self.context.run_id = config.run_id.clone(); + self.context.normalized_session_id = config.normalized_session_id.clone(); + self.context.turn_id = config.turn_id.clone(); + self.context.item_id = config.item_id.clone(); + } + fn poll_events(&mut self) -> Result, DurableRunnerError> { self.restore()?; if self diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs b/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs index 9fc8d04fef..f88dc0bbc2 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs @@ -226,6 +226,11 @@ fn apply_authority_rotation( pub trait CommandExecutor { fn execute(&mut self, command: &Command) -> Result; + /// Advances provider-side event correlation after a durable `run.attach` + /// has moved runnerd to the next run-bound authority. The runner validates + /// and persists the new authority before invoking this infallible hook. + fn rotate_authority(&mut self, _config: &DurableRunnerConfig) {} + fn poll_events(&mut self) -> Result, DurableRunnerError> { Ok(Vec::new()) } @@ -449,6 +454,7 @@ pub fn run_durable_runner( } if let Some(next) = authority_rotation { apply_authority_rotation(&mut state, &store, &mut config, &mut endpoint, next)?; + executor.rotate_authority(&config); disconnected_since = Some(Instant::now()); continue; } @@ -595,6 +601,7 @@ pub fn run_durable_runner( &mut endpoint, next, )?; + executor.rotate_authority(&config); disconnected_since = Some(Instant::now()); break; } diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/managed_provider_backend.rs b/packages/paperclip-runner/runner/crates/runner-core/src/managed_provider_backend.rs index 5456fc9824..dcdfb5ed0e 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/managed_provider_backend.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/managed_provider_backend.rs @@ -1670,6 +1670,10 @@ impl CommandExecutor for ManagedProviderCommandExecutor { } } + fn rotate_authority(&mut self, config: &DurableRunnerConfig) { + self.config = config.clone(); + } + fn poll_events(&mut self) -> Result, DurableRunnerError> { self.poll_provider()?; Ok(self diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/native_provider_backend.rs b/packages/paperclip-runner/runner/crates/runner-core/src/native_provider_backend.rs index 959b06c8d1..4d54b48a20 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/native_provider_backend.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/native_provider_backend.rs @@ -35,6 +35,14 @@ impl CommandExecutor for SelectedExecutor { } } + fn rotate_authority(&mut self, config: &DurableRunnerConfig) { + match self { + Self::LocalFacade(executor) => executor.rotate_authority(config), + Self::Acpx(executor) => executor.rotate_authority(config), + Self::Managed(executor) => executor.rotate_authority(config), + } + } + fn acknowledge_events(&mut self, count: usize) -> Result<(), DurableRunnerError> { match self { Self::LocalFacade(executor) => executor.acknowledge_events(count), @@ -163,6 +171,13 @@ impl CommandExecutor for NativeProviderCommandExecutor { .map_or_else(|| Ok(Vec::new()), CommandExecutor::poll_events) } + fn rotate_authority(&mut self, config: &DurableRunnerConfig) { + self.config = config.clone(); + if let Some(executor) = self.selected.as_mut() { + executor.rotate_authority(config); + } + } + fn acknowledge_events(&mut self, count: usize) -> Result<(), DurableRunnerError> { self.select_recovery()?; if let Some(executor) = self.selected.as_mut() { diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs b/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs index e1b7059132..c335667184 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs @@ -3030,6 +3030,10 @@ impl CommandExecutor for CodexCommandExecutor { } } + fn rotate_authority(&mut self, config: &DurableRunnerConfig) { + self.event_identity = Some(ProviderEventIdentity::from_config(config)); + } + fn poll_events(&mut self) -> Result, DurableRunnerError> { self.poll_provider()?; Ok(self 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 e352ee3193..b8ffc2d367 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.test.ts @@ -2080,6 +2080,99 @@ it("rotates PRP authority in place for a warm cross-run attachment", async () => } }, 30_000); +it("releases both PRP authorities when warm rotation activation fails", async () => { + const stateDirectory = await mkdtemp( + join(tmpdir(), "runnerd-warm-attach-activation-failure-"), + ); + const server = createServer(); + const authorities = new Map(); + const released: string[] = []; + server.on("upgrade", (request, socket, head) => { + const route = request.url ?? ""; + const authority = authorities.get(route); + if (!authority) { + socket.destroy(); + return; + } + authority.handleUpgrade(request, socket, route, head); + }); + await new Promise((resolveListen) => + server.listen(0, "127.0.0.1", resolveListen), + ); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("Expected warm activation failure test listener"); + } + let registrationCount = 0; + const bundle = createCapabilityRunnerdCodexTransport({ + runnerBinary: defaultCapabilityRunnerdBinary(), + codexCommand: fakeCodex, + codexArgs: fakeCodexArgs(stateDirectory), + stateDirectory, + lifecyclePolicy: { mode: "warm", idleTimeoutMs: 60_000 }, + controlPlaneRegistration: async (authority) => { + registrationCount += 1; + const route = `/runner-${registrationCount}`; + authorities.set(route, authority); + return { + connectUrl: `ws://127.0.0.1:${address.port}${route}`, + ...(registrationCount === 1 + ? {} + : { + activate: () => { + throw new Error("rotation activation failed"); + }, + }), + release: () => { + released.push(route); + if (authorities.get(route) === authority) authorities.delete(route); + }, + }; + }, + }); + bundle.transport.setServerRequestHandler(async () => ({ + success: true, + contentItems: [], + })); + let runnerPid: number | null = null; + try { + await bundle.transport.request("thread/start", { + cwd: tmpdir(), + dynamicTools: codexSemanticToolSpecs(), + }); + runnerPid = bundle.evidence().runnerPid; + + await expect( + bundle.transport.attachRun!({ + runId: "run-warm-activation-failure", + turnId: "turn-warm-activation-failure", + itemId: "item-warm-activation-failure", + }), + ).rejects.toThrow("rotation activation failed"); + expect(new Set(released)).toEqual(new Set(["/runner-1", "/runner-2"])); + expect(authorities.size).toBe(0); + await expect(bundle.transport.request("thread/read", {})).rejects.toThrow( + "rotation activation failed", + ); + } finally { + await bundle.transport.close().catch(() => undefined); + if (runnerPid) { + try { + process.kill(-runnerPid, "SIGKILL"); + } catch { + // A successful durable close already stopped the runner process group. + } + } + server.closeAllConnections(); + if (server.listening) { + await new Promise((resolveClose) => + server.close(() => resolveClose()), + ); + } + await rm(stateDirectory, { recursive: true, force: true }); + } +}, 30_000); + it.each([ { binding: "runner instance", diff --git a/packages/paperclip-runner/src/live/runnerd-codex-transport.ts b/packages/paperclip-runner/src/live/runnerd-codex-transport.ts index f7a3987b9f..fe384d7336 100644 --- a/packages/paperclip-runner/src/live/runnerd-codex-transport.ts +++ b/packages/paperclip-runner/src/live/runnerd-codex-transport.ts @@ -2296,16 +2296,31 @@ class DurablePrpCodexTransport implements CodexAppServerTransport { this.#eventIndex = 0; this.#durableTurnId = desired.turnId; this.#controlPlaneRelease = registration?.release ?? null; - await registration?.activate?.(); - if (registration?.failure) { - void registration.failure.catch((error: unknown) => { - this.#failTransport( - error instanceof Error ? error : new Error(String(error)), - ); - }); + let previousReleased = false; + try { + await registration?.activate?.(); + if (registration?.failure) { + void registration.failure.catch((error: unknown) => { + this.#failTransport( + error instanceof Error ? error : new Error(String(error)), + ); + }); + } + await previousRelease?.(); + previousReleased = true; + await this.#awaitRegistrationReady(registration?.ready); + } catch (error) { + const failure = error instanceof Error ? error : new Error(String(error)); + this.#controlPlaneRelease = null; + await Promise.allSettled([ + Promise.resolve().then(() => registration?.release()), + ...(previousReleased + ? [] + : [Promise.resolve().then(() => previousRelease?.())]), + ]); + this.#failTransport(failure); + throw failure; } - await previousRelease?.(); - await this.#awaitRegistrationReady(registration?.ready); } async resolveRuntimeRequest(input: { diff --git a/server/src/__tests__/native-workspace-sync.test.ts b/server/src/__tests__/native-workspace-sync.test.ts index 78fa717764..a55ee16512 100644 --- a/server/src/__tests__/native-workspace-sync.test.ts +++ b/server/src/__tests__/native-workspace-sync.test.ts @@ -1,7 +1,7 @@ import { mkdtemp, readdir, rm } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; -import { afterEach, describe, expect, it } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; import { directorySnapshotSha256, serializeDirectorySnapshot, @@ -11,6 +11,7 @@ import { classifyNativeWorkspaceInbound, nativeWorkspaceSyncInternals, readNativeWorkspaceSyncReference, + resumeNativeWorkspaceSync, } from "../services/native-runtime/native-workspace-sync.js"; const digest = "a".repeat(64); @@ -27,9 +28,9 @@ describe("native workspace sync durable metadata", () => { delete process.env.PAPERCLIP_INSTANCE_ID; else process.env.PAPERCLIP_INSTANCE_ID = originalPaperclipInstanceId; await Promise.all( - cleanupDirs.splice(0).map((directory) => - rm(directory, { recursive: true, force: true }), - ), + cleanupDirs + .splice(0) + .map((directory) => rm(directory, { recursive: true, force: true })), ); }); @@ -135,7 +136,10 @@ describe("native workspace sync durable metadata", () => { const baseline = { exclude: [".paperclip-runtime"], entries: new Map([ - ["continuity.txt", { kind: "file" as const, mode: 0o644, hash: digest }], + [ + "continuity.txt", + { kind: "file" as const, mode: 0o644, hash: digest }, + ], ]), }; const descriptor = { @@ -160,12 +164,10 @@ describe("native workspace sync durable metadata", () => { resourceDisposition: "keep_running" as const, }; - const first = await nativeWorkspaceSyncInternals.writeDescriptor( - descriptor, - ); - const second = await nativeWorkspaceSyncInternals.writeDescriptor( - descriptor, - ); + const first = + await nativeWorkspaceSyncInternals.writeDescriptor(descriptor); + const second = + await nativeWorkspaceSyncInternals.writeDescriptor(descriptor); expect(second).toEqual(first); const files = await readdir( @@ -186,4 +188,110 @@ describe("native workspace sync durable metadata", () => { }), ).resolves.toMatchObject({ descriptor }); }); + + it("repairs finalized remote and lease stamps after an interrupted commit", async () => { + const paperclipHome = await mkdtemp( + path.join(os.tmpdir(), "paperclip-native-workspace-sync-repair-"), + ); + cleanupDirs.push(paperclipHome); + process.env.PAPERCLIP_HOME = paperclipHome; + process.env.PAPERCLIP_INSTANCE_ID = "descriptor-repair-test"; + const baseline = { + exclude: [".paperclip-runtime"], + entries: new Map([ + [ + "continuity.txt", + { kind: "file" as const, mode: 0o644, hash: digest }, + ], + ]), + }; + const finalHostSha256 = "b".repeat(64); + const descriptor = { + schema: "paperclip.native-workspace-sync/v1" as const, + binding: { + runId: "run-finalized-repair", + companyId: "company-1", + workspaceId: "workspace-1", + leaseId: "lease-1", + providerLeaseId: "sandbox-1", + localCwd: path.join(paperclipHome, "workspace"), + remoteCwd: "/workspace", + }, + state: "finalized" as const, + baselineSha256: directorySnapshotSha256(baseline), + baseline: serializeDirectorySnapshot(baseline), + gitSnapshot: null, + seed: null, + createdAt: "2026-01-01T00:00:00.000Z", + finalizedAt: "2026-01-01T00:01:00.000Z", + finalHostSha256, + resourceDisposition: "keep_running" as const, + }; + const reference = + await nativeWorkspaceSyncInternals.writeDescriptor(descriptor); + const rows = (values: unknown[]) => { + const query = { + from: () => query, + where: () => query, + for: () => query, + limit: () => query, + then: ( + onfulfilled?: + ((value: unknown[]) => TResult1 | PromiseLike) | null, + onrejected?: + ((reason: unknown) => TResult2 | PromiseLike) | null, + ) => Promise.resolve(values).then(onfulfilled, onrejected), + }; + return query; + }; + let persistedLeaseMetadata: Record | null = null; + const db = { + select: () => + rows([{ runnerProfileJson: { nativeWorkspaceSync: reference } }]), + transaction: async (callback: (tx: unknown) => Promise) => + callback({ + select: () => rows([{ metadata: { retained: true } }]), + update: () => ({ + set: (value: { metadata: Record }) => ({ + where: async () => { + persistedLeaseMetadata = value.metadata; + }, + }), + }), + }), + }; + const execute = vi + .fn() + .mockResolvedValueOnce({ timedOut: false, exitCode: 1, stdout: "" }) + .mockResolvedValueOnce({ timedOut: false, exitCode: 0, stdout: "" }); + + await expect( + resumeNativeWorkspaceSync({ + db: db as never, + runId: descriptor.binding.runId, + target: { + kind: "remote", + transport: "sandbox", + remoteCwd: descriptor.binding.remoteCwd, + sandboxLeaseAcquisition: { + providerLeaseId: descriptor.binding.providerLeaseId, + }, + runner: { execute }, + } as never, + }), + ).resolves.toBe(true); + + expect(execute).toHaveBeenCalledTimes(2); + expect(persistedLeaseMetadata).toMatchObject({ + retained: true, + nativeWorkspaceSync: { + schema: "paperclip.native-workspace-stamp/v1", + workspaceId: descriptor.binding.workspaceId, + providerLeaseId: descriptor.binding.providerLeaseId, + remoteCwd: descriptor.binding.remoteCwd, + hostSha256: finalHostSha256, + finalizedRunId: descriptor.binding.runId, + }, + }); + }); }); diff --git a/server/src/services/native-runtime/native-workspace-sync.ts b/server/src/services/native-runtime/native-workspace-sync.ts index 0850c73ec8..e8f9f95884 100644 --- a/server/src/services/native-runtime/native-workspace-sync.ts +++ b/server/src/services/native-runtime/native-workspace-sync.ts @@ -652,6 +652,20 @@ function leaseStamp(input: { return stamp; } +function finalizedWorkspaceStamp(input: { + descriptor: NativeWorkspaceSyncDescriptor; + hostSha256: string; +}): Record { + return { + schema: STAMP_SCHEMA, + workspaceId: input.descriptor.binding.workspaceId, + providerLeaseId: input.descriptor.binding.providerLeaseId, + remoteCwd: input.descriptor.binding.remoteCwd, + hostSha256: input.hostSha256, + finalizedRunId: input.descriptor.binding.runId, + }; +} + async function prepareRuntime(input: { runId: string; target: Extract; @@ -690,14 +704,10 @@ async function finalizePreparedRuntime(input: { }), ); const finalHostSha256 = directorySnapshotSha256(finalSnapshot); - const stamp = { - schema: STAMP_SCHEMA, - workspaceId: input.descriptor.binding.workspaceId, - providerLeaseId: input.descriptor.binding.providerLeaseId, - remoteCwd: input.descriptor.binding.remoteCwd, + const stamp = finalizedWorkspaceStamp({ + descriptor: input.descriptor, hostSha256: finalHostSha256, - finalizedRunId: input.runId, - }; + }); await writeRemoteStamp({ target: input.target, stamp }); const finalizedDescriptor: NativeWorkspaceSyncDescriptor = { ...input.descriptor, @@ -932,12 +942,6 @@ export async function resumeNativeWorkspaceSync(input: { ); if (!reference) return false; const existing = await readDescriptor({ runId: input.runId, reference }); - if ( - existing.descriptor.state === "finalized" && - existing.descriptor.finalHostSha256 - ) { - return true; - } const providerLeaseId = input.target.sandboxLeaseAcquisition?.providerLeaseId ?? reference.providerLeaseId; @@ -947,6 +951,26 @@ export async function resumeNativeWorkspaceSync(input: { ) { throw new Error("workspace_sync_out_unrecoverable"); } + if ( + existing.descriptor.state === "finalized" && + existing.descriptor.finalHostSha256 + ) { + const stamp = finalizedWorkspaceStamp({ + descriptor: existing.descriptor, + hostSha256: existing.descriptor.finalHostSha256, + }); + if ( + !(await remoteStampMatches({ target: input.target, expected: stamp })) + ) { + await writeRemoteStamp({ target: input.target, stamp }); + } + await persistLeaseStamp({ + db: input.db, + leaseId: existing.descriptor.binding.leaseId, + stamp, + }); + return true; + } const runtime = await prepareRuntime({ runId: input.runId, target: input.target,