Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
83 changes: 83 additions & 0 deletions apps/server/scripts/acp-replay-agent.test.ts
Original file line number Diff line number Diff line change
@@ -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: "<any>" } },
{ type: "emit_inbound", frame: { kind: "response", method: "initialize", result: {} } },
{ type: "expect_outbound", frame: { kind: "request", method: "session/new", params: "<any>" } },
{
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<void>((resolve) => child.once("exit", () => resolve()));
const observed: Array<Status> = [];
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 });
}
});
61 changes: 35 additions & 26 deletions apps/server/scripts/acp-replay-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,9 +56,12 @@ let nextAgentRequestId = 1;
const pendingClientRequestIds = new Map<string, string | number>();
const pendingAgentRequestMethods = new Map<string, string>();

// 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,
Expand All @@ -67,6 +70,7 @@ function writeStatus(failure?: unknown): void {
}),
"utf8",
);
NodeFS.renameSync(pendingStatusPath, replayStatusPath);
}

function stableStringify(value: unknown): string {
Expand Down Expand Up @@ -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`);
}
Expand Down Expand Up @@ -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<JsonRpcMessage> {
const batch: Array<JsonRpcMessage> = [];
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;
Expand All @@ -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 {
Expand All @@ -261,26 +270,26 @@ 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",
id: message.id,
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) {
pendingClientRequestIds.set(actual.method, message.id);
} 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 });
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});

Expand Down
Loading