fix(runner): acknowledge terminal transitions durably

This commit is contained in:
Dotta 2026-09-02 23:57:10 -05:00
parent a02087a5c1
commit d2c9977065
3 changed files with 324 additions and 73 deletions

View File

@ -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<E: CommandExecutor>(
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<E: CommandExecutor>(
&config,
&mut executor,
&mut transport,
&connection,
&welcome.pending_commands,
);
}
@ -310,7 +314,20 @@ pub fn run_durable_runner<E: CommandExecutor>(
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<E: CommandExecutor>(
.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<E: CommandExecutor>(
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<E: CommandExecutor>(
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<E: CommandExecutor>(
}
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<E: CommandExecutor>(
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<E: CommandExecutor>(
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<E: CommandExecutor>(
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<E: CommandExecutor>(
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<E: CommandExecutor>(
let _ = executor.shutdown();
return Err(error);
}
executor.shutdown()
finish_terminal_transition_after_ack(state, store, executor)
}
fn stop_after_terminal_result_delivery_failure<E: CommandExecutor>(
@ -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!(

View File

@ -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 });
}
});
});

View File

@ -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<string, unknown>,
@ -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<RunnerProcessResult>((resolveCompletion, rejectCompletion) => {
child.once("error", rejectCompletion);
child.once("exit", (code, signal) => resolveCompletion({ code, signal, stdout, stderr }));
});
const completion = new Promise<RunnerProcessResult>(
(resolveCompletion, rejectCompletion) => {
child.once("error", rejectCompletion);
child.once("exit", (code, signal) =>
resolveCompletion({ code, signal, stdout, stderr }),
);
},
);
return withRestart({ child, completion });
}