fix(runner): replay durable attachment outcome

This commit is contained in:
Dotta 2026-09-03 04:23:40 -05:00
parent 620bb16f3b
commit 43558e7783
2 changed files with 94 additions and 9 deletions

View File

@ -47,6 +47,7 @@ import {
rehydrateRunnerdUsageNotification,
rehydrateRunnerdWorkspaceChangeNotification,
runnerdLaunchProfileInternals,
runnerdRecoveryInternals,
resolveRunnerdAcpxPermissionMode,
resolveRunnerdSessionIdentity,
resolveSourceCodexHome,
@ -56,6 +57,40 @@ import {
withCodexCollaborationRuntimeInstructions,
} from "./runnerd-codex-transport.js";
it("replays the durable run attachment outcome and latest provider identity", () => {
expect(
runnerdRecoveryInternals.recoveredRunAttachment({
commands: [
{ commandId: "prepare", type: "run.prepare", status: "completed" },
{ commandId: "attach", type: "run.attach", status: "failed" },
],
committedEvents: [{ eventType: "session.started" }],
}),
).toEqual({
commandId: "attach",
status: "failed",
providerIdentityEventIndex: -1,
});
expect(
runnerdRecoveryInternals.recoveredRunAttachment({
commands: [
{ commandId: "attach", type: "run.attach", status: "completed" },
],
committedEvents: [
{ eventType: "runner.reconciled" },
{ eventType: "session.started" },
{ eventType: "runner.diagnostic" },
{ eventType: "session.resumed" },
],
}),
).toEqual({
commandId: "attach",
status: "completed",
providerIdentityEventIndex: 3,
});
});
it("keeps ACPX terminal tools under the reserved runner-owned catalog", () => {
const tools = [
{

View File

@ -63,9 +63,7 @@ import {
// URL directory conversion preserves a trailing separator while path-derived
// build artifacts do not. Normalize once so a source build cannot be
// misclassified as an external provider pack by a string-only comparison.
const packageRoot = resolve(
fileURLToPath(new URL("../..", import.meta.url)),
);
const packageRoot = resolve(fileURLToPath(new URL("../..", import.meta.url)));
const executableSuffix = process.platform === "win32" ? ".exe" : "";
const MAX_NOTIFICATION_COUNT = 2_048;
const MAX_NOTIFICATION_BYTES = 4 * 1024 * 1024;
@ -293,6 +291,43 @@ function rotatedRunAttachPayload(
return payload;
}
function recoveredRunAttachment(state: {
commands: readonly {
commandId: string;
type: string;
status: string;
}[];
committedEvents: readonly { eventType: string }[];
}): {
commandId: string;
status: string;
providerIdentityEventIndex: number;
} | null {
const command = [...state.commands]
.reverse()
.find((candidate) => candidate.type === "run.attach");
if (!command) return null;
let providerIdentityEventIndex = -1;
if (command.status === "completed") {
for (let index = state.committedEvents.length - 1; index >= 0; index -= 1) {
const eventType = state.committedEvents[index]?.eventType;
if (
eventType === "harness.ready" ||
eventType === "session.started" ||
eventType === "session.resumed"
) {
providerIdentityEventIndex = index;
break;
}
}
}
return {
commandId: command.commandId,
status: command.status,
providerIdentityEventIndex,
};
}
function bridgedCodexQuestionParams(
request: Record<string, unknown>,
method: string,
@ -1356,8 +1391,7 @@ export function createCapabilityRunnerdProviderEnvironment(input: {
input.acpxSidecarPath ??
input.options.acpxSidecarPath ??
resolve(packageRoot, "dist", "cli", "acpx-runtime-sidecar.cjs");
const providerPackageAuthority =
acpxProviderPackageAuthority(sidecarPath);
const providerPackageAuthority = acpxProviderPackageAuthority(sidecarPath);
return {
...createSanitizedAcpxSpawnInput(
input.options.environment,
@ -1367,8 +1401,7 @@ export function createCapabilityRunnerdProviderEnvironment(input: {
// The verified sidecar bundle cannot use import.meta.url while Node
// executes it through /proc/self/fd. Anchor its closed provider package
// lookups at the package that owns the already-authenticated bundle.
PAPERCLIP_ACPX_PROVIDER_PACKAGE_ROOT:
providerPackageAuthority.root,
PAPERCLIP_ACPX_PROVIDER_PACKAGE_ROOT: providerPackageAuthority.root,
PAPERCLIP_ACPX_PROVIDER_PACKAGE_MANIFEST:
providerPackageAuthority.manifest,
...(input.options.providerRecoveryPolicy ===
@ -2678,7 +2711,6 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
connectionLeaseTtlMs: 60 * 60 * 1_000,
});
this.#core = core;
this.#eventIndex = core.store.state.committedEvents.length;
if (rotatedAuthority) {
core.queueCommand(
"run.attach",
@ -2690,6 +2722,18 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
),
);
}
const committedEvents = core.store.state.committedEvents;
const runAttachment = recoveredRunAttachment(core.store.state);
// A controller retry can open the exact authority after run.attach has
// already reached a durable outcome. Re-observe that command instead of
// silently waiting for an identity that a failed command can never emit.
// If attachment completed, replay only its latest identity event into the
// transport's in-memory evidence; session events are consumed internally
// and are not duplicated onto the provider notification stream.
this.#eventIndex =
runAttachment !== null && runAttachment.providerIdentityEventIndex >= 0
? runAttachment.providerIdentityEventIndex
: committedEvents.length;
const registration = this.options.controlPlaneRegistration
? await this.options.controlPlaneRegistration(core)
: null;
@ -2748,7 +2792,9 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
this.#evidence.runnerProcessGroupId = null;
this.#publish();
this.#pump = setInterval(() => this.#pumpEventsSafely(), 5);
if (rotatedAuthority) await this.#waitCommand("run.attach");
if (runAttachment) {
await this.#waitCommand("run.attach", runAttachment.commandId);
}
await this.#waitForProviderIdentity();
this.#startupComplete = true;
this.#diagnostic(
@ -3453,3 +3499,7 @@ export const runnerdLaunchProfileInternals = Object.freeze({
acpxRunnerLaunchProfile,
resolveBuildOwnedCliArtifact,
});
export const runnerdRecoveryInternals = Object.freeze({
recoveredRunAttachment,
});