From 83838d74eaf3752cbbcdd247fb9041c7b913754b Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sat, 3 Oct 2026 19:01:07 -0600 Subject: [PATCH] test: commit ACP replay status before the agent's answers leave The ACP replay agent sent each inbound frame before recording it in its status file, and rewrote the status in place. A client that reacted to the final answer could stop the agent or read the status first, seeing an earlier cursor ("did not consume all frames") or a truncated file ("Failed to decode ACP replay status"). The agent now consumes a batch of inbound frames (and a trailing runtime exit), commits the status atomically via rename, then sends the batch. Mismatches are recorded before the error answer is sent. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/scripts/acp-replay-agent.test.ts | 83 +++++++++++++++++++ apps/server/scripts/acp-replay-agent.ts | 61 ++++++++------ .../Adapters/AcpRegistryAdapterV2.test.ts | 4 +- 3 files changed, 120 insertions(+), 28 deletions(-) create mode 100644 apps/server/scripts/acp-replay-agent.test.ts diff --git a/apps/server/scripts/acp-replay-agent.test.ts b/apps/server/scripts/acp-replay-agent.test.ts new file mode 100644 index 000000000..16bcb118b --- /dev/null +++ b/apps/server/scripts/acp-replay-agent.test.ts @@ -0,0 +1,83 @@ +// @effect-diagnostics nodeBuiltinImport:off +import * as NodeChildProcess from "node:child_process"; +import * as NodeFS from "node:fs"; +import * as NodeOS from "node:os"; +import * as NodePath from "node:path"; +import * as NodeReadline from "node:readline"; +import * as NodeURL from "node:url"; + +import { assert, it } from "@effect/vitest"; + +const scriptPath = NodeURL.fileURLToPath(new URL("./acp-replay-agent.ts", import.meta.url)); + +const entries = [ + { type: "expect_outbound", frame: { kind: "request", method: "initialize", params: "" } }, + { type: "emit_inbound", frame: { kind: "response", method: "initialize", result: {} } }, + { type: "expect_outbound", frame: { kind: "request", method: "session/new", params: "" } }, + { + type: "emit_inbound", + frame: { kind: "response", method: "session/new", result: { sessionId: "s" } }, + }, + { type: "runtime_exit", status: "success" }, +]; + +interface Status { + readonly cursor: number; + readonly total: number; + readonly failure?: unknown; +} + +const readStatus = (statusPath: string): Status => + JSON.parse(NodeFS.readFileSync(statusPath, "utf8")) as Status; + +// Harnesses read the status the moment the client has seen a frame, and may +// stop the agent at that moment, so the status must already count the frame. +it("commits the replay status before the client can observe the frames it counts", async () => { + const dir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-acp-replay-agent-")); + const statusPath = NodePath.join(dir, "status.json"); + const transcriptPath = NodePath.join(dir, "transcript.json"); + NodeFS.writeFileSync(transcriptPath, JSON.stringify({ scenario: "status-order", entries })); + const child = NodeChildProcess.spawn( + process.execPath, + ["--experimental-strip-types", scriptPath], + { + env: { + ...process.env, + T3_ACP_REPLAY_TRANSCRIPT_PATH: transcriptPath, + T3_ACP_REPLAY_STATUS_PATH: statusPath, + }, + stdio: ["pipe", "pipe", "inherit"], + }, + ); + const exited = new Promise((resolve) => child.once("exit", () => resolve())); + const observed: Array = []; + const lines = NodeReadline.createInterface({ input: child.stdout }); + lines.on("line", (line) => { + const { id, error } = JSON.parse(line) as { readonly id: number; readonly error?: unknown }; + observed.push(readStatus(statusPath)); + if (id === 1 && error === undefined) { + child.stdin.write( + `${JSON.stringify({ jsonrpc: "2.0", id: 2, method: "session/new", params: {} })}\n`, + ); + } else { + // Stop the agent the instant its last answer arrives, as session teardown does. + child.kill("SIGKILL"); + } + }); + child.stdin.write( + `${JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize", params: {} })}\n`, + ); + await exited; + try { + assert.deepEqual( + observed.map(({ cursor, total, failure }) => ({ cursor, total, failure })), + [ + { cursor: 2, total: 5, failure: undefined }, + { cursor: 5, total: 5, failure: undefined }, + ], + ); + assert.deepInclude(readStatus(statusPath), { cursor: 5, total: 5 }); + } finally { + NodeFS.rmSync(dir, { recursive: true, force: true }); + } +}); diff --git a/apps/server/scripts/acp-replay-agent.ts b/apps/server/scripts/acp-replay-agent.ts index 4cff15e35..940a5d69c 100644 --- a/apps/server/scripts/acp-replay-agent.ts +++ b/apps/server/scripts/acp-replay-agent.ts @@ -56,9 +56,12 @@ let nextAgentRequestId = 1; const pendingClientRequestIds = new Map(); const pendingAgentRequestMethods = new Map(); +// The harness may read the status while this agent is mid-write or after it was +// killed mid-write, so a status is committed whole by renaming a sibling file. function writeStatus(failure?: unknown): void { + const pendingStatusPath = `${replayStatusPath}.${process.pid}.pending`; NodeFS.writeFileSync( - replayStatusPath, + pendingStatusPath, JSON.stringify({ scenario: transcript.scenario, cursor, @@ -67,6 +70,7 @@ function writeStatus(failure?: unknown): void { }), "utf8", ); + NodeFS.renameSync(pendingStatusPath, replayStatusPath); } function stableStringify(value: unknown): string { @@ -138,11 +142,6 @@ function stopWithFailure(detail: string, actual?: unknown): void { process.stdin.pause(); } -function advance(): void { - cursor += 1; - writeStatus(); -} - function send(message: JsonRpcMessage): void { process.stdout.write(`${JSON.stringify(message)}\n`); } @@ -184,56 +183,64 @@ function materializeInbound(value: unknown): unknown { ); } -function emitInbound(recorded: LogicalFrame): void { +function inboundMessage(recorded: LogicalFrame): JsonRpcMessage | undefined { const frame = materializeInbound(recorded) as LogicalFrame; switch (frame.kind) { case "notification": - send({ + return { jsonrpc: "2.0", method: frame.method, ...(frame.params === undefined ? {} : { params: frame.params }), - }); - return; + }; case "request": { const id = nextAgentRequestId; nextAgentRequestId += 1; pendingAgentRequestMethods.set(String(id), frame.method); - send({ + return { jsonrpc: "2.0", id, method: frame.method, ...(frame.params === undefined ? {} : { params: frame.params }), headers: [], - }); - return; + }; } case "response": { const id = pendingClientRequestId(frame.method); if (id === undefined) { stopWithFailure(`No pending client request for ${frame.method}`, frame); - return; + return undefined; } pendingClientRequestIds.delete(frame.method); - send({ + return { jsonrpc: "2.0", id, ...(frame.result === undefined ? {} : { result: frame.result }), ...(frame.error === undefined ? {} : { error: frame.error }), - }); + }; } } } +// Commits the status before any frame of the batch leaves. The client reacts to +// a frame (and may stop this agent) as soon as it arrives, so every frame it +// has seen, and the trailing runtime exit, must already count as consumed. function flushInbound(): void { + const batch = consumeInbound(); + if (!stopped) writeStatus(); + for (const message of batch) send(message); +} + +function consumeInbound(): ReadonlyArray { + const batch: Array = []; while (!stopped) { const entry = transcript.entries[cursor]; - if (entry === undefined || entry.type === "expect_outbound") return; + if (entry === undefined || entry.type === "expect_outbound") return batch; if (entry.type === "runtime_exit") { if (entry.status !== "success" && entry.status !== "cancelled") { stopWithFailure(`Recorded runtime exit was ${entry.status ?? "unknown"}`, entry.error); - return; + return batch; } - advance(); + cursor += 1; continue; } const frame = entry.frame as LogicalFrame; @@ -244,12 +251,14 @@ function flushInbound(): void { typeof frame.method !== "string" ) { stopWithFailure("Invalid emit_inbound logical ACP frame", entry.frame); - return; + return batch; } - emitInbound(frame); - if (stopped) return; - advance(); + const message = inboundMessage(frame); + if (message === undefined) return batch; + batch.push(message); + cursor += 1; } + return batch; } function handleMessage(message: JsonRpcMessage): void { @@ -261,6 +270,8 @@ function handleMessage(message: JsonRpcMessage): void { } const entry = transcript.entries[cursor]; if (entry?.type !== "expect_outbound" || !matchesExpected(entry.frame, actual)) { + // Record the mismatch before answering, for the same reason as flushInbound. + stopWithFailure("Unexpected outbound ACP frame", actual); if (actual.kind === "request" && message.id !== undefined && message.id !== null) { send({ jsonrpc: "2.0", @@ -268,7 +279,6 @@ function handleMessage(message: JsonRpcMessage): void { error: { code: -32603, message: "ACP replay frame mismatch" }, }); } - stopWithFailure("Unexpected outbound ACP frame", actual); return; } if (actual.kind === "request" && message.id !== undefined && message.id !== null) { @@ -276,11 +286,10 @@ function handleMessage(message: JsonRpcMessage): void { } else if (actual.kind === "response" && message.id !== undefined && message.id !== null) { pendingAgentRequestMethods.delete(String(message.id)); } - advance(); + cursor += 1; flushInbound(); } -writeStatus(); flushInbound(); const input = NodeReadline.createInterface({ input: process.stdin, crlfDelay: Infinity }); diff --git a/apps/server/src/orchestration-v2/Adapters/AcpRegistryAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/AcpRegistryAdapterV2.test.ts index e822ea4ef..f4095b43e 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpRegistryAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpRegistryAdapterV2.test.ts @@ -225,8 +225,8 @@ describe("AcpRegistryAdapterV2", () => { }), }) .pipe(Effect.scoped); - // Closing the session stops the agent, which writes its replay status in - // the same tick as its last answer. The script must be consumed exactly. + // Closing the session stops the agent, which commits its replay status + // before sending its last answer. The script must be consumed exactly. yield* makeAcpReplayCompletenessAssertion(fileSystem, statusPath, transcript); });