← Files Codex ReplayARCHIVED FILE

worker/tests/codex-worker.test.mjs

27.5 KB · Oct 2, 2026 · 00:19 UTC

↓ Download file

import assert from "node:assert/strict";
import { once } from "node:events";
import { chmod, mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import path from "node:path";
import { spawn } from "node:child_process";
import test from "node:test";
import { fileURLToPath } from "node:url";
import {
  executeRunRequest,
  normalizeRunRequest,
  resolveCodexTarget
} from "../src/codex-worker.mjs";

const workerRoot = path.resolve(path.dirname(fileURLToPath(import.meta.url)), "..");
const pluginRoot = path.resolve(workerRoot, "..");
const builtWorkerPath = path.join(pluginRoot, "mcp", "codex-worker.mjs");
const temporaryRoots = [];

test("captures a real filesystem error with its native code and launch stage", async () => {
  const root = await mkdtemp(path.join(tmpdir(), "codex-replay-error-"));
  temporaryRoots.push(root);
  let original;
  try {
    await readFile(path.join(root, "missing-file"));
  } catch (error) {
    original = error;
  }
  assert.equal(original.code, "ENOENT");
  const request = normalizeRunRequest({
    type: "run", id: "native-error", model: "gpt-5.6-sol", prompt: "synthetic",
    workingDirectory: root, sandboxMode: "read-only", networkAccessEnabled: false
  });
  await assert.rejects(executeRunRequest(request, {
    codexFactory: () => { throw original; }
  }), (error) => {
    assert.equal(error.systemCode, "ENOENT");
    assert.equal(error.diagnostic, undefined);
    assert.equal(error.stage, "launch");
    assert.ok(error.elapsedMs >= 0);
    return true;
  });
});

test.after(async () => {
  await Promise.all(
    temporaryRoots.map(async (root) => {
      await rm(root, { recursive: true, force: true });
    })
  );
});

test("retains bounded redacted SDK exception context and the underlying cause", async () => {
  const diagnostics = [];
  const cause = new Error("Operation not permitted: access_token=private-token-value");
  cause.code = "EPERM";
  const original = new Error(`SDK startup failed: Bearer private-bearer-value\n${"warning ".repeat(2_000)}`, { cause });
  const request = normalizeRunRequest({
    type: "run", id: "sdk-exception", model: "gpt-5.6-sol", prompt: "synthetic",
    workingDirectory: "/private/fixture", sandboxMode: "read-only", networkAccessEnabled: false
  });
  await assert.rejects(executeRunRequest(request, {
    diagnostic: (message) => diagnostics.push(message),
    codexFactory: () => { throw original; }
  }), (error) => {
    assert.equal(error.code, "worker_failed");
    assert.equal(error.systemCode, "EPERM");
    assert.equal(error.stage, "launch");
    assert.equal(error.retryable, false);
    assert.equal(error.message, "Operation not permitted: access_token=[REDACTED]");
    return true;
  });
  assert.equal(diagnostics.length, 1);
  assert.match(diagnostics[0], /SDK startup failed: Bearer \[REDACTED\]/);
  assert.ok(diagnostics[0].length <= 8_000 + "Codex worker exception: ".length);
  assert.doesNotMatch(diagnostics[0], /private-token-value|private-bearer-value/);
});

test("probes Codex candidates in configured, PATH, and app-bundle order", async () => {
  const root = await mkdtemp(path.join(tmpdir(), "codex-replay-discovery-"));
  temporaryRoots.push(root);
  const binDirectory = path.join(root, "bin");
  const pathCodex = path.join(binDirectory, "codex");
  const bundledCodex = path.join(
    root,
    "Codex Preview.app",
    "Contents",
    "Resources",
    "codex"
  );
  await Promise.all([
    mkdir(binDirectory, { recursive: true }),
    mkdir(path.dirname(bundledCodex), { recursive: true })
  ]);
  await Promise.all([
    writeFile(pathCodex, ""),
    writeFile(bundledCodex, "")
  ]);
  await Promise.all([chmod(pathCodex, 0o755), chmod(bundledCodex, 0o755)]);

  const configuredProbes = [];
  assert.equal(
    resolveCodexTarget({
      env: { CODEX_CLI_PATH: "/custom/codex", PATH: binDirectory },
      platform: "darwin",
      applicationRoots: [root],
      probe(candidate) {
        configuredProbes.push(candidate);
        return true;
      }
    }),
    "/custom/codex"
  );
  assert.deepEqual(configuredProbes, ["/custom/codex"]);

  const pathProbes = [];
  assert.equal(
    resolveCodexTarget({
      env: { PATH: binDirectory },
      platform: "darwin",
      applicationRoots: [root],
      probe(candidate) {
        pathProbes.push(candidate);
        return true;
      }
    }),
    pathCodex
  );
  assert.deepEqual(pathProbes, [pathCodex]);

  const bundleProbes = [];
  assert.equal(
    resolveCodexTarget({
      env: { PATH: "" },
      platform: "darwin",
      applicationRoots: [root],
      probe(candidate) {
        bundleProbes.push(candidate);
        return true;
      }
    }),
    bundledCodex
  );
  assert.deepEqual(bundleProbes, [bundledCodex]);
});

test("falls back from a broken configured Codex to a healthy app bundle", async () => {
  const root = await mkdtemp(path.join(tmpdir(), "codex-replay-fallback-"));
  temporaryRoots.push(root);
  const configuredCodex = path.join(root, "nvm", "bin", "codex");
  const bundledCodex = path.join(
    root,
    "Codex.app",
    "Contents",
    "Resources",
    "codex"
  );
  await Promise.all([
    mkdir(path.dirname(configuredCodex), { recursive: true }),
    mkdir(path.dirname(bundledCodex), { recursive: true })
  ]);
  await Promise.all([writeFile(configuredCodex, ""), writeFile(bundledCodex, "")]);
  await Promise.all([chmod(configuredCodex, 0o755), chmod(bundledCodex, 0o755)]);

  const probes = [];
  assert.equal(
    resolveCodexTarget({
      env: { CODEX_CLI_PATH: configuredCodex, PATH: "" },
      platform: "darwin",
      applicationRoots: [root],
      probe(candidate) {
        probes.push(candidate);
        return candidate === bundledCodex;
      }
    }),
    bundledCodex
  );
  assert.deepEqual(probes, [configuredCodex, bundledCodex]);
});

test("reports the candidates when no Codex version probe succeeds", () => {
  assert.throws(
    () =>
      resolveCodexTarget({
        env: { CODEX_CLI_PATH: "/broken/codex", PATH: "" },
        platform: "linux",
        applicationRoots: [],
        probe: () => false
      }),
    (error) =>
      error.code === "codex_unavailable" &&
      error.message.includes("/broken/codex")
  );
});

test("normalizes server aliases and keeps SDK events sanitized", async () => {
  const request = normalizeRunRequest({
    type: "run",
    requestId: "fixture-1",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    sandboxMode: "workspace-write",
    networkAccess: false,
    outputSchema: { type: "object" }
  });
  assert.equal(request.id, "fixture-1");
  assert.equal(request.timeoutMs, undefined);
  assert.equal(request.networkAccessEnabled, false);

  const emitted = [];
  let observedOptions;
  let observedTurnOptions;
  const result = await executeRunRequest(request, {
    emit: (event) => emitted.push(event),
    codexFactory: () => ({
      startThread(options) {
        observedOptions = options;
        return {
          id: "fixture-thread",
          async runStreamed(prompt, turnOptions) {
            assert.equal(prompt, "private prompt");
            observedTurnOptions = turnOptions;
            return {
              events: (async function* () {
                yield { type: "thread.started", thread_id: "fixture-thread" };
                yield { type: "turn.started" };
                yield {
                  type: "item.completed",
                  item: {
                    type: "command_execution",
                    command: "secret command",
                    aggregated_output: "secret output",
                    status: "completed"
                  }
                };
                yield {
                  type: "item.completed",
                  item: { type: "agent_message", text: "finished" }
                };
                yield {
                  type: "turn.completed",
                  usage: {
                    input_tokens: 11,
                    cached_input_tokens: 2,
                    output_tokens: 3,
                    reasoning_output_tokens: 1
                  }
                };
              })()
            };
          }
        };
      }
    })
  });

  assert.equal(observedOptions.model, "gpt-5.6-sol");
  assert.equal(observedOptions.approvalPolicy, "never");
  assert.equal(observedOptions.networkAccessEnabled, false);
  assert.equal(observedOptions.sandboxMode, "workspace-write");
  assert.deepEqual(observedTurnOptions.outputSchema, { type: "object" });
  assert.equal(result.threadId, "fixture-thread");
  assert.equal(result.finalResponse, "finished");
  assert.equal(result.usage.input_tokens, 11);
  const serializedEvents = JSON.stringify(emitted);
  assert.doesNotMatch(serializedEvents, /private prompt|secret command|secret output/);
  assert.match(serializedEvents, /thread_started|item_completed/);
});

test("preserves actionable SDK stream errors while redacting credentials", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "stream-error",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });

  const diagnostics = [];
  await assert.rejects(
    executeRunRequest(request, {
      diagnostic: (message) => diagnostics.push(message),
      codexFactory: () => ({
        startThread() {
          return {
            id: "fixture-thread",
            async runStreamed() {
              return {
                events: (async function* () {
                  yield { type: "thread.started", thread_id: "fixture-thread" };
                  yield {
                    type: "error",
                    message: "connection reset: api_key=sk-proj-abcdefghijklmnopqrst"
                  };
                  yield {
                    type: "turn.completed",
                    usage: {
                      input_tokens: 0,
                      cached_input_tokens: 0,
                      output_tokens: 0,
                      reasoning_output_tokens: 0
                    }
                  };
                })()
              };
            }
          };
        }
      })
    }),
    (error) =>
      error instanceof Error &&
      error.code === "stream_error" &&
      error.retryable === true &&
      error.message === "connection reset: api_key=[REDACTED]" &&
      !error.message.includes("abcdefghijklmnopqrst")
  );
  assert.deepEqual(diagnostics, [
    "Codex stream error: connection reset: api_key=[REDACTED]"
  ]);
});

test("extracts actionable provider errors from JSON stream events", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "outdated-codex",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });
  const providerMessage =
    "The 'gpt-5.6-sol' model requires a newer version of Codex. Please upgrade to the latest app or CLI and try again.";

  await assert.rejects(
    executeRunRequest(request, {
      codexFactory: () => ({
        startThread() {
          return {
            async runStreamed() {
              return {
                events: (async function* () {
                  yield {
                    type: "error",
                    message: JSON.stringify({
                      type: "error",
                      status: 400,
                      error: { type: "invalid_request_error", message: providerMessage }
                    })
                  };
                })()
              };
            }
          };
        }
      })
    }),
    (error) =>
      error.code === "stream_error" &&
      error.retryable === false &&
      error.message === providerMessage
  );
});

test("retries structured transient and ambiguous SDK stream errors", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "retryable-stream-errors",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });
  const messages = [
    JSON.stringify({
      type: "error",
      status: 503,
      error: {
        type: "error",
        message: "invalid request while upstream auth service is unavailable"
      }
    }),
    "The stream stopped before the final event."
  ];

  for (const message of messages) {
    await assert.rejects(
      executeRunRequest(request, {
        codexFactory: () => ({
          startThread() {
            return {
              async runStreamed() {
                return {
                  events: (async function* () {
                    yield { type: "error", message };
                  })()
                };
              }
            };
          }
        })
      }),
      (error) => error.code === "stream_error" && error.retryable === true
    );
  }
});

test("retries interrupted SDK iterators and streams that end before completion", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "interrupted-streams",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });
  const eventStreams = [
    (async function* () {
      yield { type: "turn.started" };
      throw new Error("socket hang up");
    })(),
    (async function* () {
      yield { type: "turn.started" };
    })()
  ];

  for (const events of eventStreams) {
    await assert.rejects(
      executeRunRequest(request, {
        codexFactory: () => ({
          startThread() {
            return {
              async runStreamed() {
                return { events };
              }
            };
          }
        })
      }),
      (error) => error.retryable === true
    );
  }
});

test("does not retry permanent SDK iterator failures", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "permanent-iterator-error",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });

  await assert.rejects(
    executeRunRequest(request, {
      codexFactory: () => ({
        startThread() {
          return {
            async runStreamed() {
              return {
                events: (async function* () {
                  throw new Error("Authentication failed because the API key is invalid.");
                })()
              };
            }
          };
        }
      })
    }),
    (error) => error.code === "worker_failed" && error.retryable === false
  );
});

test("does not retry explicit permanent plain-text stream errors", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "permanent-stream-error",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });

  await assert.rejects(
    executeRunRequest(request, {
      codexFactory: () => ({
        startThread() {
          return {
            async runStreamed() {
              return {
                events: (async function* () {
                  yield {
                    type: "error",
                    message: "Authentication failed because the API key is invalid."
                  };
                })()
              };
            }
          };
        }
      })
    }),
    (error) => error.code === "stream_error" && error.retryable === false
  );
});

test("preserves actionable SDK turn failures while redacting bearer tokens", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "turn-failure",
    model: "gpt-5.6-sol",
    prompt: "private prompt",
    workingDirectory: "/private/fixture",
    timeoutMs: 30_000,
    sandboxMode: "read-only",
    networkAccessEnabled: false
  });
  const diagnostics = [];

  await assert.rejects(
    executeRunRequest(request, {
      diagnostic: (message) => diagnostics.push(message),
      codexFactory: () => ({
        startThread() {
          return {
            id: "fixture-thread",
            async runStreamed() {
              return {
                events: (async function* () {
                  yield { type: "thread.started", thread_id: "fixture-thread" };
                  yield {
                    type: "turn.failed",
                    error: { message: "authentication failed: Bearer sensitive-token-value" }
                  };
                })()
              };
            }
          };
        }
      })
    }),
    (error) =>
      error instanceof Error &&
      error.code === "turn_failed" &&
      error.retryable === false &&
      error.message === "authentication failed: Bearer [REDACTED]" &&
      !error.message.includes("sensitive-token-value")
  );
  assert.deepEqual(diagnostics, [
    "Codex turn failed: authentication failed: Bearer [REDACTED]"
  ]);
});

test("returns actionable SDK stream errors through the worker protocol", async () => {
  const fixture = await fakeCodexFixture();
  const child = spawn(process.execPath, [builtWorkerPath], {
    env: {
      ...process.env,
      CODEX_CLI_PATH: fixture.executablePath
    },
    stdio: ["pipe", "pipe", "pipe"]
  });
  const stdout = collectLines(child.stdout);
  const stderr = collectText(child.stderr);
  child.stdin.end(
    `${JSON.stringify({
      type: "run",
      id: "stdio-stream-error",
      model: "gpt-5.6-sol",
      prompt: "STREAM_ERROR",
      workingDirectory: fixture.root,
      timeoutMs: 10_000,
      sandboxMode: "workspace-write",
      networkAccessEnabled: false
    })}\n`
  );

  const [code, signal] = await once(child, "exit");
  assert.equal(code, 1);
  assert.equal(signal, null);
  const messages = await stdout;
  const failure = messages.at(-1);
  assert.equal(failure.type, "failed");
  assert.equal(failure.retryable, true);
  assert.equal(failure.message, "connection reset: fixture provider detail");
  assert.match(await stderr, /Codex stream error: connection reset: fixture provider detail/);
});

test("reports CLI startup failures through the packaged worker without losing stderr", async () => {
  const fixture = await fakeCodexFixture();
  const child = spawn(process.execPath, [builtWorkerPath], {
    env: { ...process.env, CODEX_CLI_PATH: fixture.executablePath },
    stdio: ["pipe", "pipe", "pipe"]
  });
  const stdout = collectLines(child.stdout);
  const stderr = collectText(child.stderr);
  child.stdin.end(`${JSON.stringify({
    type: "run", id: "cli-startup-error", model: "gpt-5.6-sol", prompt: "CLI_STARTUP_ERROR",
    workingDirectory: fixture.root, sandboxMode: "workspace-write", networkAccessEnabled: false
  })}\n`);

  const [code, signal] = await once(child, "exit");
  assert.equal(code, 1);
  assert.equal(signal, null);
  const messages = await stdout;
  const failure = messages.at(-1);
  assert.equal(failure.type, "failed");
  assert.equal(failure.code, "worker_failed");
  assert.equal(failure.stage, "stream");
  assert.equal(failure.retryable, false);
  assert.equal(failure.message, "Error: failed to initialize in-process app-server client: Operation not permitted (os error 1)");
  assert.equal(messages.some((event) => event.phase === "thread_started"), false);
  const details = await stderr;
  assert.match(details, /Codex Exec exited with code 1/);
  assert.match(details, /startup warning: api_key=\[REDACTED\]/);
  assert.match(details, /failed to initialize in-process app-server client/);
  assert.doesNotMatch(JSON.stringify(messages) + details, /fixture-private-api-key/);
});

test("built JSONL worker executes one SDK run and returns the observed result", async () => {
  const fixture = await fakeCodexFixture();
  const child = spawn(process.execPath, [builtWorkerPath], {
    env: {
      ...process.env,
      CODEX_CLI_PATH: fixture.executablePath,
      FAKE_CODEX_ARGS_FILE: fixture.argsPath
    },
    stdio: ["pipe", "pipe", "pipe"]
  });
  const stdout = collectLines(child.stdout);
  const stderr = collectText(child.stderr);
  child.stdin.end(
    `${JSON.stringify({
      type: "run",
      requestId: "stdio-1",
      model: "gpt-5.6-sol",
      prompt: "fixture prompt",
      workingDirectory: fixture.root,
      timeoutSeconds: 10,
      sandboxMode: "workspace-write",
      networkAccess: false
    })}\n`
  );

  const [code, signal] = await once(child, "exit");
  assert.equal(code, 0, await stderr);
  assert.equal(signal, null);
  const messages = await stdout;
  assert.equal(messages[0].type, "ready");
  assert.equal(messages[1].type, "accepted");
  const completed = messages.find((message) => message.type === "completed");
  assert.ok(completed);
  assert.equal(completed.id, "stdio-1");
  assert.equal(completed.threadId, "fixture-thread-id");
  assert.equal(completed.finalResponse, "fixture final response");
  assert.equal(completed.usage.input_tokens, 7);
  const serialized = JSON.stringify(messages);
  assert.doesNotMatch(serialized, /fixture-secret-command|fixture-secret-output/);
  const cliArgs = JSON.parse(await readFile(fixture.argsPath, "utf8"));
  assert.deepEqual(cliArgs.slice(0, 2), [
    "exec",
    "--ignore-rules"
  ]);
});

test("SIGTERM cancels the active SDK run without retrying", async () => {
  const fixture = await fakeCodexFixture();
  const child = spawn(process.execPath, [builtWorkerPath], {
    env: {
      ...process.env,
      CODEX_CLI_PATH: fixture.executablePath,
      FAKE_TREE_TERMINATED_FILE: fixture.terminationPath
    },
    stdio: ["pipe", "pipe", "pipe"]
  });
  const messages = [];
  let buffered = "";
  child.stdout.setEncoding("utf8");
  child.stdout.on("data", (chunk) => {
    buffered += chunk;
    const lines = buffered.split("\n");
    buffered = lines.pop() ?? "";
    for (const line of lines) {
      if (!line) continue;
      const message = JSON.parse(line);
      messages.push(message);
      if (message.type === "lifecycle" && message.phase === "thread_started") {
        child.kill("SIGTERM");
      }
    }
  });
  child.stdin.write(
    `${JSON.stringify({
      type: "run",
      id: "stdio-cancel",
      model: "gpt-5.6-sol",
      prompt: "HANG",
      workingDirectory: fixture.root,
      timeoutMs: 10_000,
      sandboxMode: "workspace-write",
      networkAccessEnabled: false
    })}\n`
  );

  const [code, signal] = await once(child, "exit");
  assert.equal(code, 143);
  assert.equal(signal, null);
  const canceled = messages.find((message) => message.type === "canceled");
  assert.ok(canceled);
  assert.equal(canceled.reason, "SIGTERM");
  assert.equal(messages.filter((message) => message.type === "accepted").length, 1);
  await waitForFile(fixture.terminationPath);
});

test("legacy timeout fields do not limit Codex execution", async () => {
  const request = normalizeRunRequest({
    type: "run",
    id: "unlimited-run",
    model: "gpt-5.6-sol",
    prompt: "Continue until complete.",
    workingDirectory: "/private/fixture",
    timeoutMs: 1,
    sandboxMode: "workspace-write",
    networkAccessEnabled: false
  });
  assert.equal(request.timeoutMs, undefined);

  const result = await executeRunRequest(request, {
    codexFactory: () => ({
      startThread() {
        return {
          id: "unlimited-thread",
          async runStreamed() {
            return {
              events: (async function* () {
                await new Promise((resolve) => setTimeout(resolve, 20));
                yield { type: "turn.completed", usage: {} };
              })()
            };
          }
        };
      }
    })
  });

  assert.equal(result.threadId, "unlimited-thread");
});

async function fakeCodexFixture() {
  const root = await mkdtemp(path.join(tmpdir(), "codex-replay-worker-"));
  temporaryRoots.push(root);
  const executablePath = path.join(root, "fake-codex.mjs");
  const grandchildPath = path.join(root, "fake-grandchild.mjs");
  const argsPath = path.join(root, "args.json");
  const terminationPath = path.join(root, "tree-terminated");
  await writeFile(
    grandchildPath,
    [
      "#!/usr/bin/env node",
      "import { writeFileSync } from 'node:fs';",
      "process.on('SIGTERM', () => {",
      "  writeFileSync(process.env.FAKE_TREE_TERMINATED_FILE, 'terminated');",
      "  process.exit(0);",
      "});",
      "console.log('ready');",
      "setInterval(() => {}, 1000);",
      ""
    ].join("\n")
  );
  await chmod(grandchildPath, 0o755);
  await writeFile(
    executablePath,
    [
      "#!/usr/bin/env node",
      "import { spawn } from 'node:child_process';",
      "import { writeFileSync } from 'node:fs';",
      "if (process.env.FAKE_CODEX_ARGS_FILE) writeFileSync(process.env.FAKE_CODEX_ARGS_FILE, JSON.stringify(process.argv.slice(2)));",
      "let prompt = '';",
      "for await (const chunk of process.stdin) prompt += chunk;",
      "if (prompt.includes('CLI_STARTUP_ERROR')) { console.error('startup warning: api_key=fixture-private-api-key ' + 'warning '.repeat(200)); console.error('Error: failed to initialize in-process app-server client: Operation not permitted (os error 1)'); process.exit(1); }",
      `if (prompt.includes('HANG')) {`,
      `  const grandchild = spawn(${JSON.stringify(grandchildPath)}, [], { stdio: ['ignore', 'pipe', 'ignore'] });`,
      "  await new Promise((resolve, reject) => { grandchild.stdout.once('data', resolve); grandchild.once('error', reject); });",
      "}",
      "console.log(JSON.stringify({ type: 'thread.started', thread_id: 'fixture-thread-id' }));",
      "console.log(JSON.stringify({ type: 'turn.started' }));",
      "if (prompt.includes('HANG')) { setInterval(() => {}, 1000); await new Promise(() => {}); }",
      "if (prompt.includes('STREAM_ERROR')) { console.log(JSON.stringify({ type: 'error', message: 'connection reset: fixture provider detail' })); process.exit(0); }",
      "console.log(JSON.stringify({ type: 'item.completed', item: { type: 'command_execution', command: 'fixture-secret-command', aggregated_output: 'fixture-secret-output', status: 'completed' } }));",
      "console.log(JSON.stringify({ type: 'item.completed', item: { type: 'agent_message', text: 'fixture final response' } }));",
      "console.log(JSON.stringify({ type: 'turn.completed', usage: { input_tokens: 7, cached_input_tokens: 2, output_tokens: 3, reasoning_output_tokens: 1 } }));",
      ""
    ].join("\n")
  );
  await chmod(executablePath, 0o755);
  return { root, executablePath, argsPath, terminationPath };
}

async function waitForFile(filePath) {
  const deadline = Date.now() + 3_000;
  while (Date.now() < deadline) {
    try {
      return await readFile(filePath, "utf8");
    } catch {
      await new Promise((resolve) => setTimeout(resolve, 25));
    }
  }
  assert.fail(`Timed out waiting for ${filePath}`);
}

function collectLines(stream) {
  return new Promise((resolve, reject) => {
    const messages = [];
    let buffered = "";
    stream.setEncoding("utf8");
    stream.on("data", (chunk) => {
      buffered += chunk;
      const lines = buffered.split("\n");
      buffered = lines.pop() ?? "";
      try {
        for (const line of lines) {
          if (line) messages.push(JSON.parse(line));
        }
      } catch (error) {
        reject(error);
      }
    });
    stream.on("end", () => {
      if (buffered.trim()) {
        try {
          messages.push(JSON.parse(buffered));
        } catch (error) {
          reject(error);
          return;
        }
      }
      resolve(messages);
    });
    stream.on("error", reject);
  });
}

function collectText(stream) {
  return new Promise((resolve, reject) => {
    let text = "";
    stream.setEncoding("utf8");
    stream.on("data", (chunk) => {
      text += chunk;
    });
    stream.on("end", () => resolve(text));
    stream.on("error", reject);
  });
}

SHA-256: aa73b45b5476cc8bbc1f0d6d8fd1e60d0c37dd8386ff6a4d65560056b1b7d820