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 b5412fc392..6d0ae23dfc 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 @@ -45,6 +45,8 @@ enum CommandLifecycle { Shutdown, } +const TERMINAL_RESULT_ACK_TIMEOUT: Duration = Duration::from_secs(2); + fn sleep_for_reconnect(base: Duration, max_delay: Duration, attempt: &mut u32) { let multiplier = 1_u128 << (*attempt).min(5); let uncapped = base.as_millis().saturating_mul(multiplier); @@ -267,6 +269,7 @@ pub fn run_durable_runner( if let Some(acked_source_seq) = welcome.acked_source_seq { state.apply_ack(acked_source_seq)?; } + let connection = welcome.connection; if state.pending_terminal_delivery.is_some() { return reconcile_pending_terminal_delivery( &mut state, @@ -274,6 +277,7 @@ pub fn run_durable_runner( &config, &mut executor, &mut transport, + &connection, &welcome.pending_commands, ); } @@ -310,7 +314,20 @@ pub fn run_durable_runner( break; } if lifecycle.durable_state().is_some() { - complete_terminal_delivery_after_send(&mut state, &store, &result)?; + if let Err(error) = wait_for_terminal_result_ack( + &mut transport, + &mut state, + &store, + &connection, + &result, + ) { + return stop_after_terminal_result_delivery_failure( + &mut state, + &store, + &mut executor, + error, + ); + } // A terminal lifecycle command is the final command this // process may accept. Flush its already-durable outbox below, // then release the executor without observing later commands. @@ -338,8 +355,7 @@ pub fn run_durable_runner( .filter(|_| !disconnected) { debug_assert_eq!(state.lifecycle, durable_lifecycle); - executor.shutdown()?; - return Ok(()); + return finish_terminal_transition_after_ack(&mut state, &store, &mut executor); } if disconnected { disconnected_since.get_or_insert_with(Instant::now); @@ -350,8 +366,6 @@ pub fn run_durable_runner( sleep_before_deadline(config.reconnect_delay, reconnect_deadline); continue; } - - let connection = welcome.connection; loop { if started.elapsed() >= config.max_runtime { break; @@ -439,7 +453,20 @@ pub fn run_durable_runner( break; } if lifecycle.durable_state().is_some() { - complete_terminal_delivery_after_send(&mut state, &store, &result)?; + if let Err(error) = wait_for_terminal_result_ack( + &mut transport, + &mut state, + &store, + &connection, + &result, + ) { + return stop_after_terminal_result_delivery_failure( + &mut state, + &store, + &mut executor, + error, + ); + } } if let Err(error) = send_outbox(&mut transport, &state, &mut sent_source_seq) { state.record_diagnostic( @@ -459,8 +486,11 @@ pub fn run_durable_runner( } if let Some(durable_lifecycle) = lifecycle.durable_state() { debug_assert_eq!(state.lifecycle, durable_lifecycle); - executor.shutdown()?; - return Ok(()); + return finish_terminal_transition_after_ack( + &mut state, + &store, + &mut executor, + ); } } Some("revoke") => { @@ -545,33 +575,108 @@ fn persist_lifecycle_before_command_delivery( store.save(state) } -fn complete_terminal_delivery_after_send( +fn complete_terminal_delivery_after_cleanup( state: &mut DurableState, store: &DurableStateStore, - result: &StoredCommandResult, ) -> Result<(), DurableRunnerError> { let pending = state.pending_terminal_delivery.as_ref().ok_or_else(|| { - DurableRunnerError::invalid("terminal result delivery has no durable recovery fence") + DurableRunnerError::invalid("terminal cleanup has no durable recovery fence") })?; - if pending.command_id != result.command_id - || pending.controller_seq != result.controller_seq - || pending.command_type != result.command_type - || state.lifecycle != pending.lifecycle - { + if state.lifecycle != pending.lifecycle { return Err(DurableRunnerError::invalid( - "terminal result delivery does not match its durable recovery fence", + "terminal cleanup does not match its durable recovery fence", )); } state.pending_terminal_delivery = None; store.save(state) } +fn finish_terminal_transition_after_ack( + state: &mut DurableState, + store: &DurableStateStore, + executor: &mut E, +) -> Result<(), DurableRunnerError> { + // Keep the durable fence through provider cleanup. If cleanup fails, a + // replacement may authenticate only to retry terminal reconciliation and + // cannot restore the suspended runner to ready. + executor.shutdown()?; + complete_terminal_delivery_after_cleanup(state, store) +} + +fn wait_for_terminal_result_ack( + transport: &mut AuthenticatedTransport, + state: &mut DurableState, + store: &DurableStateStore, + connection: &ConnectionMetadata, + result: &StoredCommandResult, +) -> Result<(), DurableRunnerError> { + let deadline = Instant::now() + TERMINAL_RESULT_ACK_TIMEOUT; + while Instant::now() < deadline { + let Some(message) = transport.receive_json()? else { + continue; + }; + validate_control_identity(&message, state, Some(connection))?; + match message.get("kind").and_then(Value::as_str) { + Some("command_result_ack") => { + let payload = message + .get("payload") + .and_then(Value::as_object) + .ok_or_else(|| { + DurableRunnerError::invalid( + "terminal command result acknowledgement payload is required", + ) + })?; + if payload.get("commandId").and_then(Value::as_str) + != Some(result.command_id.as_str()) + || payload.get("commandType").and_then(Value::as_str) + != Some(result.command_type.as_str()) + || payload.get("controllerSeq").and_then(Value::as_u64) + != Some(result.controller_seq) + || payload.get("status").and_then(Value::as_str) != Some(result.status.as_str()) + { + return Err(DurableRunnerError::invalid( + "terminal command result acknowledgement changed its durable identity", + )); + } + return Ok(()); + } + Some("ack") => { + let acked = message + .pointer("/payload/ackedSourceSeq") + .and_then(Value::as_u64) + .ok_or_else(|| DurableRunnerError::invalid("ACK cursor is required"))?; + state.apply_ack(acked)?; + store.save(state)?; + } + Some("ping") => transport.send_json(&control_envelope( + state, + connection, + "pong", + json!({ + "lifecycle": state.lifecycle, + "ackedSourceSeq": state.acked_source_seq, + "outboxBytes": state.outbox_bytes(), + }), + ))?, + _ => { + return Err(DurableRunnerError::invalid( + "controller sent a non-acknowledgement after a terminal command result", + )); + } + } + } + Err(DurableRunnerError::invalid( + "terminal command result acknowledgement timed out", + )) +} + fn reconcile_pending_terminal_delivery( state: &mut DurableState, store: &DurableStateStore, config: &DurableRunnerConfig, executor: &mut E, transport: &mut AuthenticatedTransport, + connection: &ConnectionMetadata, pending_commands: &[Command], ) -> Result<(), DurableRunnerError> { let pending = state.pending_terminal_delivery.clone().ok_or_else(|| { @@ -597,7 +702,11 @@ fn reconcile_pending_terminal_delivery( if let Err(error) = transport.send_json(&command_result_envelope(state, &result)) { return stop_after_terminal_result_delivery_failure(state, store, executor, error); } - complete_terminal_delivery_after_send(state, store, &result)?; + if let Err(error) = + wait_for_terminal_result_ack(transport, state, store, connection, &result) + { + return stop_after_terminal_result_delivery_failure(state, store, executor, error); + } } else { // An authenticated welcome is the controller's authoritative pending // set. Absence means the prior write reached the controller even if @@ -605,7 +714,6 @@ fn reconcile_pending_terminal_delivery( state.record_diagnostic( "controller confirmed the pending terminal result was already delivered", ); - state.pending_terminal_delivery = None; store.save(state)?; } @@ -616,7 +724,7 @@ fn reconcile_pending_terminal_delivery( let _ = executor.shutdown(); return Err(error); } - executor.shutdown() + finish_terminal_transition_after_ack(state, store, executor) } fn stop_after_terminal_result_delivery_failure( @@ -982,7 +1090,7 @@ mod tests { } #[test] - fn successful_terminal_result_delivery_clears_the_recovery_fence() { + fn successful_terminal_cleanup_clears_the_recovery_fence() { let directory = std::env::temp_dir().join(format!( "paperclip-runner-terminal-result-delivered-{}", std::process::id() @@ -1004,7 +1112,7 @@ mod tests { ) .unwrap(); assert!(state.pending_terminal_delivery.is_some()); - complete_terminal_delivery_after_send(&mut state, &store, &result).unwrap(); + complete_terminal_delivery_after_cleanup(&mut state, &store).unwrap(); let (recovered, existed) = store.load_or_create(&config).unwrap(); assert!(existed); @@ -1013,6 +1121,39 @@ mod tests { fs::remove_dir_all(directory).unwrap(); } + #[test] + fn failed_terminal_cleanup_keeps_the_recovery_fence() { + let directory = std::env::temp_dir().join(format!( + "paperclip-runner-terminal-cleanup-failed-{}", + std::process::id() + )); + let _ = fs::remove_dir_all(&directory); + let config = config(directory.clone()); + let store = DurableStateStore::new(&directory).unwrap(); + let (mut state, _) = store.load_or_create(&config).unwrap(); + let mut executor = ShutdownFailingExecutor; + let command = command("runner.suspend"); + let (result, lifecycle) = + process_command(&mut state, &store, &config, &mut executor, &command).unwrap(); + persist_lifecycle_before_command_delivery( + &mut state, + &store, + lifecycle.durable_state().unwrap(), + &result, + ) + .unwrap(); + + let error = finish_terminal_transition_after_ack(&mut state, &store, &mut executor) + .expect_err("cleanup failure remains fenced"); + let (recovered, existed) = store.load_or_create(&config).unwrap(); + + assert!(error.to_string().contains("terminal cleanup failure")); + assert!(existed); + assert_eq!(recovered.lifecycle, "suspended"); + assert!(recovered.pending_terminal_delivery.is_some()); + fs::remove_dir_all(directory).unwrap(); + } + #[test] fn event_batch_keeps_accepted_prefix_and_unacknowledged_suffix() { let directory = std::env::temp_dir().join(format!( diff --git a/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.test.ts b/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.test.ts index c6d8025adb..fcf8b7309e 100644 --- a/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.test.ts +++ b/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.test.ts @@ -123,20 +123,22 @@ it("pins the OpenCode launch profile in runner startup arguments and restarts", handle.restart("replacement-ticket"); expect(launches).toHaveLength(2); for (const launch of launches) { - expect(launch.args).toEqual(expect.arrayContaining([ - "--opencode-proxy-command", - profile.command, - "--opencode-proxy-command-sha256", - profile.commandSha256, - "--opencode-proxy-script", - profile.proxyScript, - "--opencode-proxy-script-sha256", - profile.proxyScriptSha256, - "--opencode-executable", - profile.executable, - "--opencode-executable-sha256", - profile.executableSha256, - ])); + expect(launch.args).toEqual( + expect.arrayContaining([ + "--opencode-proxy-command", + profile.command, + "--opencode-proxy-command-sha256", + profile.commandSha256, + "--opencode-proxy-script", + profile.proxyScript, + "--opencode-proxy-script-sha256", + profile.proxyScriptSha256, + "--opencode-executable", + profile.executable, + "--opencode-executable-sha256", + profile.executableSha256, + ]), + ); } }); @@ -228,9 +230,9 @@ it("preserves the controller-selected ACPX provider package root", () => { }); expect(launches).toHaveLength(1); - expect( - launches[0]!.environment.PAPERCLIP_ACPX_PROVIDER_PACKAGE_ROOT, - ).toBe("/verified/provider-pack"); + expect(launches[0]!.environment.PAPERCLIP_ACPX_PROVIDER_PACKAGE_ROOT).toBe( + "/verified/provider-pack", + ); expect(launches[0]!.environment.NODE_PATH).toBeUndefined(); }); @@ -967,4 +969,63 @@ describe.sequential("DurablePrpControlPlane", () => { rmSync(root, { recursive: true, force: true }); } }); + + it("acknowledges terminal command results after persisting them", async () => { + const root = mkdtempSync(resolve(tmpdir(), "paperclip-prp-terminal-ack-")); + const controlPlane = new DurablePrpControlPlane({ + stateDirectory: root, + identity, + expectedRunnerVersion, + expectedRunnerDigest, + }); + try { + await controlPlane.start(); + const command = controlPlane.queueCommand( + "runner.suspend", + {}, + "command-suspend-1", + ); + const client = await authenticate( + controlPlane, + controlPlane.issueBootstrapTicket(), + ); + const terminalResult = { + protocol: "paperclip.runner", + version: 1, + kind: "command_result", + payload: { + commandId: command.commandId, + commandType: command.type, + controllerSeq: command.controllerSeq, + status: "completed", + result: { suspended: true }, + }, + }; + + sendSecure(client!, terminalResult); + await expect(receiveSecure(client!)).resolves.toMatchObject({ + kind: "command_result_ack", + payload: { + commandId: "command-suspend-1", + commandType: "runner.suspend", + controllerSeq: command.controllerSeq, + status: "completed", + }, + }); + expect(controlPlane.store.state.commands).toMatchObject([ + { commandId: "command-suspend-1", status: "completed" }, + ]); + + sendSecure(client!, terminalResult); + await expect(receiveSecure(client!)).resolves.toMatchObject({ + kind: "command_result_ack", + payload: { commandId: "command-suspend-1" }, + }); + expect(controlPlane.store.state.duplicateCommandResults).toBe(1); + client?.socket.destroy(); + } finally { + await controlPlane.stop(); + rmSync(root, { recursive: true, force: true }); + } + }); }); diff --git a/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.ts b/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.ts index 5a7e8ee74b..dc2360d9fa 100644 --- a/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.ts +++ b/packages/paperclip-runner/src/control-plane/durable-prp-control-plane.ts @@ -295,21 +295,33 @@ function canonicalJson( depth = 0, ): string { state.nodes += 1; - if (depth > MAX_CANONICAL_JSON_DEPTH || state.nodes > MAX_CANONICAL_JSON_NODES) { + if ( + depth > MAX_CANONICAL_JSON_DEPTH || + state.nodes > MAX_CANONICAL_JSON_NODES + ) { throw new Error("durable_prp_canonical_json_too_large"); } - if (value === null || typeof value === "boolean" || typeof value === "string") { + if ( + value === null || + typeof value === "boolean" || + typeof value === "string" + ) { return JSON.stringify(value) ?? "null"; } if (typeof value === "number") { - if (!Number.isFinite(value)) throw new Error("durable_prp_canonical_json_invalid"); + if (!Number.isFinite(value)) + throw new Error("durable_prp_canonical_json_invalid"); return JSON.stringify(value) ?? "null"; } if (typeof value !== "object" || ancestors.has(value)) { throw new Error("durable_prp_canonical_json_invalid"); } const prototype = Object.getPrototypeOf(value); - if (!Array.isArray(value) && prototype !== Object.prototype && prototype !== null) { + if ( + !Array.isArray(value) && + prototype !== Object.prototype && + prototype !== null + ) { throw new Error("durable_prp_canonical_json_invalid"); } ancestors.add(value); @@ -1511,10 +1523,7 @@ export class DurablePrpControlPlane { this.#welcome(connection, leaseToken); } - #welcome( - connection: AuthorityConnection, - leaseToken: string | null, - ): void { + #welcome(connection: AuthorityConnection, leaseToken: string | null): void { const lease = connection.lease; if (lease === null || connection.connectionId === null) { connection.close(); @@ -1666,15 +1675,42 @@ export class DurablePrpControlPlane { } this.#store.state.duplicateCommandResults += 1; this.#store.save(); + this.#ackTerminalCommandResult(connection, command); this.#sendNextCommand(connection); return; } command.status = status; command.result = structuredClone(result); this.#store.save(); + this.#ackTerminalCommandResult(connection, command); this.#sendNextCommand(connection); } + #ackTerminalCommandResult( + connection: AuthorityConnection, + command: DurableRecoveryCoreCommand, + ): void { + if ( + command.type !== "runner.suspend" && + command.type !== "runner.shutdown" + ) { + return; + } + connection.sendJson( + this.#controlEnvelope( + connection, + `command_result_ack_${command.controllerSeq}`, + "command_result_ack", + { + commandId: command.commandId, + commandType: command.type, + controllerSeq: command.controllerSeq, + status: command.status, + }, + ), + ); + } + async #event( connection: AuthorityConnection, envelope: Record, @@ -1957,26 +1993,30 @@ export function spawnRunner(options: { environment?: NodeJS.ProcessEnv; processLauncher?: (spec: RunnerProcessLaunchSpec) => RunnerProcessHandle; }): RunnerProcessHandle { - const connection = options.connection ?? (options.connectUrl - ? { mode: "connect" as const, connectUrl: options.connectUrl } - : null); - if (connection === null) throw new Error("runner process connection is required"); - const connectionArgs = connection.mode === "connect" - ? [ - "--connect-url", - connection.connectUrl, - ...(connection.caBundlePath === undefined - ? [] - : ["--ca-bundle-path", connection.caBundlePath]), - ] - : [ - "--listen-address", - connection.listenAddress, - "--listen-port", - String(connection.listenPort), - "--listen-path", - connection.listenPath, - ]; + const connection = + options.connection ?? + (options.connectUrl + ? { mode: "connect" as const, connectUrl: options.connectUrl } + : null); + if (connection === null) + throw new Error("runner process connection is required"); + const connectionArgs = + connection.mode === "connect" + ? [ + "--connect-url", + connection.connectUrl, + ...(connection.caBundlePath === undefined + ? [] + : ["--ca-bundle-path", connection.caBundlePath]), + ] + : [ + "--listen-address", + connection.listenAddress, + "--listen-port", + String(connection.listenPort), + "--listen-path", + connection.listenPath, + ]; const args = [ ...connectionArgs, "--state-dir", @@ -2049,7 +2089,10 @@ export function spawnRunner(options: { if (options.lifecyclePolicy !== undefined) { args.push("--lifecycle-mode", options.lifecyclePolicy.mode); if (options.lifecyclePolicy.mode === "warm") { - args.push("--idle-timeout-ms", String(options.lifecyclePolicy.idleTimeoutMs)); + args.push( + "--idle-timeout-ms", + String(options.lifecyclePolicy.idleTimeoutMs), + ); } } @@ -2060,7 +2103,9 @@ export function spawnRunner(options: { restart: (ticket) => spawnRunner({ ...options, ticket }), }); if (options.processLauncher !== undefined) { - return withRestart(options.processLauncher({ command, args, cwd: packageRoot, environment })); + return withRestart( + options.processLauncher({ command, args, cwd: packageRoot, environment }), + ); } const child = spawn(command, args, { @@ -2076,10 +2121,14 @@ export function spawnRunner(options: { child.stderr.setEncoding("utf8").on("data", (chunk: string) => { stderr = `${stderr}${chunk}`.slice(-16_384); }); - const completion = new Promise((resolveCompletion, rejectCompletion) => { - child.once("error", rejectCompletion); - child.once("exit", (code, signal) => resolveCompletion({ code, signal, stdout, stderr })); - }); + const completion = new Promise( + (resolveCompletion, rejectCompletion) => { + child.once("error", rejectCompletion); + child.once("exit", (code, signal) => + resolveCompletion({ code, signal, stdout, stderr }), + ); + }, + ); return withRestart({ child, completion }); }