fix(runner): restore stopped ACPX session recovery
This commit is contained in:
parent
d13b9ede00
commit
449b8d5e29
|
|
@ -1159,6 +1159,23 @@ impl AcpxCommandExecutor {
|
|||
state.active_turn_id = None;
|
||||
self.session = None;
|
||||
self.save_state()?;
|
||||
} else if self.state.as_ref().is_some_and(|state| {
|
||||
state.lifecycle == "prepared"
|
||||
&& state.identity.is_some()
|
||||
&& !state.provider_exit_unconfirmed
|
||||
&& state.active_turn_id.is_none()
|
||||
}) {
|
||||
// turn.stop deliberately leaves an already-reaped provider in a
|
||||
// non-recoverable `prepared` state while runner.drain crosses the
|
||||
// durable event barrier. Once the following runner.suspend reaches
|
||||
// this boundary, publish the exact stopped checkpoint as
|
||||
// recoverable instead of reporting a no-op success that can never
|
||||
// emit session.resumed in the replacement runner.
|
||||
self.state
|
||||
.as_mut()
|
||||
.expect("ACPX stopped provider state remains available")
|
||||
.lifecycle = "suspended".to_owned();
|
||||
self.save_state()?;
|
||||
}
|
||||
Ok(CommandExecution::result(json!({"status": "completed"})))
|
||||
}
|
||||
|
|
@ -1956,7 +1973,7 @@ mod tests {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn unconfirmed_suspension_state_is_recoverable_but_not_attachable() {
|
||||
fn unconfirmed_suspension_state_becomes_recoverable_only_after_cleanup_and_suspend() {
|
||||
let directory = temporary_directory("suspension-fence-pending");
|
||||
let (provider_lifetime_fence_candidates, original_lifetime_fence) =
|
||||
reserve_provider_lifetime_fence();
|
||||
|
|
@ -2018,6 +2035,14 @@ mod tests {
|
|||
let recovered = executor.state.as_ref().unwrap();
|
||||
assert_eq!(recovered.lifecycle, "prepared");
|
||||
assert!(!recovered.provider_exit_unconfirmed);
|
||||
|
||||
executor.suspend().unwrap();
|
||||
let suspended: AcpxDurableState = serde_json::from_slice(
|
||||
&fs::read(executor.state_path()).expect("read suspended ACPX state"),
|
||||
)
|
||||
.expect("parse suspended ACPX state");
|
||||
assert_eq!(suspended.lifecycle, "suspended");
|
||||
assert!(!suspended.provider_exit_unconfirmed);
|
||||
fs::remove_dir_all(directory).unwrap();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -54,8 +54,10 @@ import type {
|
|||
} from "./codex-driver-types.js";
|
||||
import {
|
||||
boundedText,
|
||||
canonicalJson,
|
||||
codexSemanticToolSpecs,
|
||||
differingJsonPaths,
|
||||
parseProviderIdentity,
|
||||
record,
|
||||
text,
|
||||
} from "./codex-driver-values.js";
|
||||
|
|
@ -372,6 +374,17 @@ export class CodexAppServerDriver implements HarnessDriver {
|
|||
reason: "provider resumed a different provider session",
|
||||
};
|
||||
}
|
||||
if (
|
||||
snapshot.providerIdentity !== undefined &&
|
||||
canonicalJson(opened.providerIdentity) !==
|
||||
canonicalJson(snapshot.providerIdentity)
|
||||
) {
|
||||
await cancellation.wait(cancellation.close());
|
||||
return {
|
||||
recovered: false,
|
||||
reason: "provider resumed with a different tagged session identity",
|
||||
};
|
||||
}
|
||||
const checkpointedActiveTurnId = snapshot.activeTurnId ?? null;
|
||||
// A terminal fingerprint is the durable provider fact. A crash can
|
||||
// persist it before the following active-turn clear reaches the same
|
||||
|
|
@ -670,9 +683,11 @@ export class CodexAppServerDriver implements HarnessDriver {
|
|||
"Codex thread response changed the assigned working directory",
|
||||
);
|
||||
}
|
||||
const providerIdentity = parseProviderIdentity(thread.providerIdentity);
|
||||
return {
|
||||
threadId,
|
||||
providerSessionId,
|
||||
...(providerIdentity === undefined ? {} : { providerIdentity }),
|
||||
collaborationMode,
|
||||
context: {
|
||||
protocolVersion: CODEX_CODEX_PROTOCOL_VERSION,
|
||||
|
|
|
|||
|
|
@ -44,6 +44,45 @@ import {
|
|||
} from "./codex-app-server-driver.test-support.js";
|
||||
|
||||
describe("Codex app-server Codex driver", () => {
|
||||
it("persists and verifies the tagged runnerd provider identity on recovery", async () => {
|
||||
const providerIdentity = {
|
||||
kind: "acpx",
|
||||
normalizedSessionId: "normalized-tagged-recovery",
|
||||
acpxRecordId: "acpx-record-1",
|
||||
backendSessionId: "backend-session-1",
|
||||
agentSessionId: "agent-session-1",
|
||||
profileDigest: `sha256:${"a".repeat(64)}`,
|
||||
workspaceDigest: `sha256:${"b".repeat(64)}`,
|
||||
requestedModel: "gpt-5.6-sol",
|
||||
effectiveModel: "gpt-5.6-sol",
|
||||
permissionMode: "approve-all",
|
||||
providerLifetimeFenceCandidates: [60_001, 60_002, 60_003],
|
||||
};
|
||||
const first = new FakeCodexTransport(
|
||||
"thread-1",
|
||||
"provider-session-1",
|
||||
providerIdentity,
|
||||
);
|
||||
const second = new FakeCodexTransport("thread-1", "provider-session-1", {
|
||||
...providerIdentity,
|
||||
backendSessionId: "backend-session-2",
|
||||
});
|
||||
const driver = makeDriver([first, second]);
|
||||
const original = await driver.openSession({
|
||||
runId: "run-tagged-recovery",
|
||||
normalizedSessionId: "normalized-tagged-recovery",
|
||||
workingDirectory: WORKSPACE,
|
||||
});
|
||||
const snapshot = await original.snapshot();
|
||||
expect(snapshot.providerIdentity).toEqual(providerIdentity);
|
||||
await original.close({ reason: "transport lost" });
|
||||
|
||||
await expect(driver.recoverSession?.(snapshot)).resolves.toEqual({
|
||||
recovered: false,
|
||||
reason: "provider resumed with a different tagged session identity",
|
||||
});
|
||||
});
|
||||
|
||||
it("resumes and reconciles the exact provider thread after transport loss", async () => {
|
||||
const first = new FakeCodexTransport();
|
||||
const second = new FakeCodexTransport();
|
||||
|
|
|
|||
|
|
@ -110,6 +110,7 @@ export class FakeCodexTransport implements CodexAppServerTransport {
|
|||
constructor(
|
||||
readonly threadId = "thread-1",
|
||||
readonly providerSessionId = "provider-session-1",
|
||||
readonly providerIdentity?: Record<string, unknown>,
|
||||
) {}
|
||||
|
||||
async request(
|
||||
|
|
@ -148,6 +149,9 @@ export class FakeCodexTransport implements CodexAppServerTransport {
|
|||
thread: {
|
||||
id: this.threadId,
|
||||
sessionId: this.providerSessionId,
|
||||
...(this.providerIdentity === undefined
|
||||
? {}
|
||||
: { providerIdentity: structuredClone(this.providerIdentity) }),
|
||||
modelProvider: "openai",
|
||||
cwd: WORKSPACE,
|
||||
turns: [],
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ import type {
|
|||
HarnessRuntimeRequest,
|
||||
HarnessRuntimeRequestResolution,
|
||||
HarnessThreadLineageEntry,
|
||||
PersistedHarnessProviderIdentity,
|
||||
PersistedHarnessSession,
|
||||
} from "../../contracts/harness-driver.js";
|
||||
import type {
|
||||
|
|
@ -87,6 +88,7 @@ export interface TerminalReplayConflict {
|
|||
export interface OpenedCodexThread {
|
||||
threadId: string;
|
||||
providerSessionId: string | null;
|
||||
providerIdentity?: PersistedHarnessProviderIdentity;
|
||||
collaborationMode: Record<string, unknown> | null;
|
||||
context: CodexModelContextSnapshot;
|
||||
lineage: HarnessThreadLineageEntry;
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import type { PersistedHarnessProviderIdentity } from "../../contracts/harness-driver.js";
|
||||
import type { NativeUserMessage } from "../../contracts/types.js";
|
||||
import {
|
||||
CODEX_BLOCK_RESULT_PROVIDER_INPUT_SCHEMA,
|
||||
|
|
@ -21,6 +22,75 @@ export function text(value: unknown, fallback = ""): string {
|
|||
return typeof value === "string" ? value : fallback;
|
||||
}
|
||||
|
||||
export function parseProviderIdentity(
|
||||
value: unknown,
|
||||
): PersistedHarnessProviderIdentity | undefined {
|
||||
const identity = record(value);
|
||||
if (identity.kind !== "acpx") return undefined;
|
||||
const requiredStrings = [
|
||||
"normalizedSessionId",
|
||||
"acpxRecordId",
|
||||
"backendSessionId",
|
||||
"agentSessionId",
|
||||
"profileDigest",
|
||||
"workspaceDigest",
|
||||
"requestedModel",
|
||||
"effectiveModel",
|
||||
] as const;
|
||||
if (
|
||||
requiredStrings.some(
|
||||
(key) =>
|
||||
typeof identity[key] !== "string" ||
|
||||
identity[key].length === 0 ||
|
||||
identity[key].length > 240,
|
||||
)
|
||||
) {
|
||||
throw new Error("ACPX provider identity is incomplete");
|
||||
}
|
||||
const permissionMode = identity.permissionMode;
|
||||
if (
|
||||
permissionMode !== undefined &&
|
||||
permissionMode !== "approve-all" &&
|
||||
permissionMode !== "approve-reads" &&
|
||||
permissionMode !== "deny-all"
|
||||
) {
|
||||
throw new Error(
|
||||
"ACPX provider identity contains an invalid permission mode",
|
||||
);
|
||||
}
|
||||
const fenceCandidates = identity.providerLifetimeFenceCandidates;
|
||||
if (
|
||||
!Array.isArray(fenceCandidates) ||
|
||||
fenceCandidates.length !== 3 ||
|
||||
fenceCandidates.some(
|
||||
(candidate) =>
|
||||
!Number.isInteger(candidate) ||
|
||||
candidate < 49_152 ||
|
||||
candidate > 65_535,
|
||||
) ||
|
||||
new Set(fenceCandidates).size !== 3
|
||||
) {
|
||||
throw new Error("ACPX provider identity contains invalid lifetime fences");
|
||||
}
|
||||
return {
|
||||
kind: "acpx",
|
||||
normalizedSessionId: identity.normalizedSessionId as string,
|
||||
acpxRecordId: identity.acpxRecordId as string,
|
||||
backendSessionId: identity.backendSessionId as string,
|
||||
agentSessionId: identity.agentSessionId as string,
|
||||
profileDigest: identity.profileDigest as string,
|
||||
workspaceDigest: identity.workspaceDigest as string,
|
||||
requestedModel: identity.requestedModel as string,
|
||||
effectiveModel: identity.effectiveModel as string,
|
||||
...(permissionMode === undefined ? {} : { permissionMode }),
|
||||
providerLifetimeFenceCandidates: fenceCandidates as [
|
||||
number,
|
||||
number,
|
||||
number,
|
||||
],
|
||||
};
|
||||
}
|
||||
|
||||
export function boundedText(
|
||||
value: unknown,
|
||||
fallback = "unknown",
|
||||
|
|
|
|||
|
|
@ -634,6 +634,9 @@ export class CodexHarnessSession extends CodexSessionState implements HarnessSes
|
|||
driverKind: this.driverKind,
|
||||
driverSessionId: this.opened.threadId,
|
||||
providerSessionId: this.opened.providerSessionId,
|
||||
...(this.opened.providerIdentity === undefined
|
||||
? {}
|
||||
: { providerIdentity: structuredClone(this.opened.providerIdentity) }),
|
||||
runId: this.runId,
|
||||
normalizedSessionId: this.normalizedSessionId,
|
||||
activeTurnId: this.activeTurnId,
|
||||
|
|
|
|||
Loading…
Reference in New Issue