From 5da4132865faaea8688aabf5784fe5c624ad42e8 Mon Sep 17 00:00:00 2001 From: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co> Date: Sun, 19 Jul 2026 10:43:37 -0400 Subject: [PATCH] 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 Signed-off-by: Tyler Longwell --- desktop/tests/e2e/fixtures/fake-acp-agent.mjs | 84 +++++++++ desktop/tests/e2e/helpers/twoRelayHarness.ts | 159 ++++++++++++++++++ 2 files changed, 243 insertions(+) create mode 100755 desktop/tests/e2e/fixtures/fake-acp-agent.mjs create mode 100644 desktop/tests/e2e/helpers/twoRelayHarness.ts diff --git a/desktop/tests/e2e/fixtures/fake-acp-agent.mjs b/desktop/tests/e2e/fixtures/fake-acp-agent.mjs new file mode 100755 index 000000000..c06999b19 --- /dev/null +++ b/desktop/tests/e2e/fixtures/fake-acp-agent.mjs @@ -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}`, + }, + }); + } +} diff --git a/desktop/tests/e2e/helpers/twoRelayHarness.ts b/desktop/tests/e2e/helpers/twoRelayHarness.ts new file mode 100644 index 000000000..335d20247 --- /dev/null +++ b/desktop/tests/e2e/helpers/twoRelayHarness.ts @@ -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 { + const chunks = await Promise.all( + this.processes.map(async ({ name, logPath }) => { + const body = await readFile(logPath, "utf8").catch(() => ""); + 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((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`); + } +}