mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
test: add agents-everywhere live harness fixtures
Co-authored-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Tyler Longwell <tlongwell@block.xyz> Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
co-authored by
Tyler Longwell
parent
8908bd6b71
commit
5da4132865
+84
@@ -0,0 +1,84 @@
|
||||
#!/usr/bin/env node
|
||||
/** Deterministic ACP fixture for agents-everywhere live tests. */
|
||||
import { createInterface } from "node:readline";
|
||||
|
||||
const wakeDelayMs = Number.parseInt(
|
||||
process.env.BUZZ_E2E_FAKE_ACP_WAKE_MS ?? "0",
|
||||
10,
|
||||
);
|
||||
if (!Number.isFinite(wakeDelayMs) || wakeDelayMs < 0) {
|
||||
throw new Error("BUZZ_E2E_FAKE_ACP_WAKE_MS must be a non-negative integer");
|
||||
}
|
||||
|
||||
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
|
||||
const write = (message) => process.stdout.write(`${JSON.stringify(message)}\n`);
|
||||
const textFromPrompt = (params) =>
|
||||
(params?.prompt ?? [])
|
||||
.filter((part) => part?.type === "text" && typeof part.text === "string")
|
||||
.map((part) => part.text)
|
||||
.join("\n");
|
||||
|
||||
let sessionCounter = 0;
|
||||
const input = createInterface({
|
||||
input: process.stdin,
|
||||
crlfDelay: Number.POSITIVE_INFINITY,
|
||||
});
|
||||
for await (const line of input) {
|
||||
if (!line.trim()) continue;
|
||||
const request = JSON.parse(line);
|
||||
if (request.id === undefined || typeof request.method !== "string") continue;
|
||||
|
||||
switch (request.method) {
|
||||
case "initialize":
|
||||
if (wakeDelayMs > 0) await sleep(wakeDelayMs);
|
||||
write({
|
||||
jsonrpc: "2.0",
|
||||
id: request.id,
|
||||
result: { protocolVersion: 2, agentCapabilities: {} },
|
||||
});
|
||||
break;
|
||||
case "session/new":
|
||||
sessionCounter += 1;
|
||||
write({
|
||||
jsonrpc: "2.0",
|
||||
id: request.id,
|
||||
result: { sessionId: `fake-session-${sessionCounter}` },
|
||||
});
|
||||
break;
|
||||
case "session/prompt": {
|
||||
const prompt = textFromPrompt(request.params);
|
||||
const ids = [...prompt.matchAll(/\bAE-ID:([A-Za-z0-9._:-]+)\b/g)].map(
|
||||
(match) => match[1],
|
||||
);
|
||||
write({
|
||||
jsonrpc: "2.0",
|
||||
method: "session/update",
|
||||
params: {
|
||||
sessionId: request.params?.sessionId,
|
||||
update: {
|
||||
sessionUpdate: "agent_message_chunk",
|
||||
content: { type: "text", text: `AE-ACK:${ids.join(",")}` },
|
||||
},
|
||||
},
|
||||
});
|
||||
write({
|
||||
jsonrpc: "2.0",
|
||||
id: request.id,
|
||||
result: { stopReason: "end_turn" },
|
||||
});
|
||||
break;
|
||||
}
|
||||
case "session/cancel":
|
||||
write({ jsonrpc: "2.0", id: request.id, result: {} });
|
||||
break;
|
||||
default:
|
||||
write({
|
||||
jsonrpc: "2.0",
|
||||
id: request.id,
|
||||
error: {
|
||||
code: -32601,
|
||||
message: `Unsupported fixture method: ${request.method}`,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,159 @@
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import { createWriteStream } from "node:fs";
|
||||
import { mkdtemp, readFile, rm } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join, resolve } from "node:path";
|
||||
import { setTimeout as delay } from "node:timers/promises";
|
||||
|
||||
export type RelayPorts = { main: number; health: number; metrics: number };
|
||||
export type RelaySpec = {
|
||||
name: string;
|
||||
ports: RelayPorts;
|
||||
databaseUrl: string;
|
||||
redisUrl: string;
|
||||
};
|
||||
|
||||
type OwnedProcess = { name: string; child: ChildProcess; logPath: string };
|
||||
|
||||
export class TwoRelayHarness {
|
||||
readonly root: string;
|
||||
readonly relays: readonly [RelaySpec, RelaySpec];
|
||||
private readonly processes: OwnedProcess[] = [];
|
||||
|
||||
private constructor(root: string, relays: readonly [RelaySpec, RelaySpec]) {
|
||||
this.root = root;
|
||||
this.relays = relays;
|
||||
}
|
||||
|
||||
static async create(relays: readonly [RelaySpec, RelaySpec]) {
|
||||
return new TwoRelayHarness(
|
||||
await mkdtemp(join(tmpdir(), "buzz-ae-e2e-")),
|
||||
relays,
|
||||
);
|
||||
}
|
||||
|
||||
async startRelays(binary = process.env.BUZZ_E2E_RELAY_BIN) {
|
||||
if (!binary)
|
||||
throw new Error("BUZZ_E2E_RELAY_BIN is required for the live gate");
|
||||
await Promise.all(
|
||||
this.relays.map((relay) => this.startRelay(binary, relay)),
|
||||
);
|
||||
}
|
||||
|
||||
async startAcp(
|
||||
name: string,
|
||||
relayWsUrl: string,
|
||||
privateKey: string,
|
||||
extraEnv: NodeJS.ProcessEnv = {},
|
||||
) {
|
||||
const binary = process.env.BUZZ_E2E_ACP_BIN;
|
||||
if (!binary)
|
||||
throw new Error("BUZZ_E2E_ACP_BIN is required for the live gate");
|
||||
return this.spawnOwned(name, binary, [], {
|
||||
BUZZ_RELAY_URL: relayWsUrl,
|
||||
BUZZ_PRIVATE_KEY: privateKey,
|
||||
BUZZ_ACP_LAZY_POOL: "1",
|
||||
BUZZ_ACP_AGENT_COMMAND: process.execPath,
|
||||
BUZZ_ACP_AGENT_ARGS: resolve("tests/e2e/fixtures/fake-acp-agent.mjs"),
|
||||
...extraEnv,
|
||||
});
|
||||
}
|
||||
|
||||
async logs(): Promise<string> {
|
||||
const chunks = await Promise.all(
|
||||
this.processes.map(async ({ name, logPath }) => {
|
||||
const body = await readFile(logPath, "utf8").catch(() => "<no log>");
|
||||
return `===== ${name} =====\n${body}`;
|
||||
}),
|
||||
);
|
||||
return chunks.join("\n");
|
||||
}
|
||||
|
||||
private signal(child: ChildProcess, signal: NodeJS.Signals) {
|
||||
if (child.exitCode !== null || child.signalCode !== null || !child.pid)
|
||||
return;
|
||||
if (process.platform === "win32") {
|
||||
child.kill(signal);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
process.kill(-child.pid, signal);
|
||||
} catch (error) {
|
||||
const code = (error as NodeJS.ErrnoException).code;
|
||||
if (code !== "ESRCH") throw error;
|
||||
}
|
||||
}
|
||||
|
||||
private async stopChild(child: ChildProcess) {
|
||||
if (child.exitCode !== null || child.signalCode !== null) return;
|
||||
const exited = new Promise<void>((resolveExit) =>
|
||||
child.once("exit", () => resolveExit()),
|
||||
);
|
||||
this.signal(child, "SIGTERM");
|
||||
if (await Promise.race([exited.then(() => true), delay(5_000, false)]))
|
||||
return;
|
||||
this.signal(child, "SIGKILL");
|
||||
if (!(await Promise.race([exited.then(() => true), delay(2_000, false)]))) {
|
||||
throw new Error(`child process ${child.pid ?? "unknown"} did not exit`);
|
||||
}
|
||||
}
|
||||
|
||||
async stop() {
|
||||
await Promise.all(
|
||||
[...this.processes].reverse().map(({ child }) => this.stopChild(child)),
|
||||
);
|
||||
await rm(this.root, { recursive: true, force: true });
|
||||
}
|
||||
|
||||
private async startRelay(binary: string, relay: RelaySpec) {
|
||||
const child = this.spawnOwned(relay.name, binary, [], {
|
||||
DATABASE_URL: relay.databaseUrl,
|
||||
REDIS_URL: relay.redisUrl,
|
||||
RELAY_URL: `ws://127.0.0.1:${relay.ports.main}`,
|
||||
BUZZ_BIND_ADDR: `127.0.0.1:${relay.ports.main}`,
|
||||
BUZZ_HEALTH_PORT: String(relay.ports.health),
|
||||
BUZZ_METRICS_PORT: String(relay.ports.metrics),
|
||||
BUZZ_REQUIRE_AUTH_TOKEN: "false",
|
||||
BUZZ_RECONCILE_CHANNELS: "true",
|
||||
});
|
||||
await this.waitForHealth(relay, child);
|
||||
}
|
||||
|
||||
private spawnOwned(
|
||||
name: string,
|
||||
command: string,
|
||||
args: string[],
|
||||
env: NodeJS.ProcessEnv,
|
||||
) {
|
||||
const logPath = join(this.root, `${name}.log`);
|
||||
const child = spawn(command, args, {
|
||||
cwd: resolve(".."),
|
||||
env: { ...process.env, ...env, RUST_LOG: process.env.RUST_LOG ?? "info" },
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
detached: process.platform !== "win32",
|
||||
});
|
||||
const log = createWriteStream(logPath, { flags: "a" });
|
||||
child.stdout?.pipe(log, { end: false });
|
||||
child.stderr?.pipe(log, { end: false });
|
||||
child.on("exit", () => log.end());
|
||||
this.processes.push({ name, child, logPath });
|
||||
return child;
|
||||
}
|
||||
|
||||
private async waitForHealth(relay: RelaySpec, child: ChildProcess) {
|
||||
const deadline = Date.now() + 30_000;
|
||||
while (Date.now() < deadline) {
|
||||
if (child.exitCode !== null || child.signalCode !== null) {
|
||||
throw new Error(`${relay.name} exited before readiness`);
|
||||
}
|
||||
try {
|
||||
const response = await fetch(
|
||||
`http://127.0.0.1:${relay.ports.health}/_readiness`,
|
||||
);
|
||||
if (response.ok) return;
|
||||
} catch {}
|
||||
await delay(100);
|
||||
}
|
||||
throw new Error(`${relay.name} was not ready within 30s`);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user