fix(runner): harden warm authority recovery

This commit is contained in:
Dotta 2026-09-05 12:26:40 -05:00
parent ee11cd65d9
commit 7a6f5b6daf
9 changed files with 310 additions and 33 deletions

View File

@ -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<Vec<PolledEvent>, DurableRunnerError> {
self.restore()?;
if self

View File

@ -226,6 +226,11 @@ fn apply_authority_rotation(
pub trait CommandExecutor {
fn execute(&mut self, command: &Command) -> Result<CommandExecution, DurableRunnerError>;
/// 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<Vec<PolledEvent>, DurableRunnerError> {
Ok(Vec::new())
}
@ -449,6 +454,7 @@ pub fn run_durable_runner<E: CommandExecutor>(
}
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<E: CommandExecutor>(
&mut endpoint,
next,
)?;
executor.rotate_authority(&config);
disconnected_since = Some(Instant::now());
break;
}

View File

@ -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<Vec<PolledEvent>, DurableRunnerError> {
self.poll_provider()?;
Ok(self

View File

@ -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() {

View File

@ -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<Vec<PolledEvent>, DurableRunnerError> {
self.poll_provider()?;
Ok(self

View File

@ -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<string, DurablePrpControlPlane>();
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<void>((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<void>((resolveClose) =>
server.close(() => resolveClose()),
);
}
await rm(stateDirectory, { recursive: true, force: true });
}
}, 30_000);
it.each([
{
binding: "runner instance",

View File

@ -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: {

View File

@ -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: <TResult1 = unknown, TResult2 = never>(
onfulfilled?:
((value: unknown[]) => TResult1 | PromiseLike<TResult1>) | null,
onrejected?:
((reason: unknown) => TResult2 | PromiseLike<TResult2>) | null,
) => Promise.resolve(values).then(onfulfilled, onrejected),
};
return query;
};
let persistedLeaseMetadata: Record<string, unknown> | null = null;
const db = {
select: () =>
rows([{ runnerProfileJson: { nativeWorkspaceSync: reference } }]),
transaction: async (callback: (tx: unknown) => Promise<unknown>) =>
callback({
select: () => rows([{ metadata: { retained: true } }]),
update: () => ({
set: (value: { metadata: Record<string, unknown> }) => ({
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,
},
});
});
});

View File

@ -652,6 +652,20 @@ function leaseStamp(input: {
return stamp;
}
function finalizedWorkspaceStamp(input: {
descriptor: NativeWorkspaceSyncDescriptor;
hostSha256: string;
}): Record<string, unknown> {
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<AdapterExecutionTarget, { transport: "sandbox" }>;
@ -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,