Files
OmniRoute/tests/unit/codex-app-server.test.ts
Armin Anton” ∴ 8f390efffd feat(codex): self-contained codex app-server transport (executor + provider + sign-in) (#11205)
Merged after conflict resolution: the 5 conflicting test files were the base-red drains that #11201 already landed on the tip — kept the tip versions; the feature content is untouched. Validated on the combined batch board + this branch: codex-app-server + codex-gpt56-catalog 25/25, typecheck:core clean, docs-counts green (351 providers), provider-consistency 268/351/0. The opt-in codex-app-server transport (JSON-RPC-over-WS, turn/completed-awaited close, Responses SSE bridge) leaves the default codex path untouched. Thank you @arminanton — a 3.4k-line transport with the docs wave and tests to match!
2026-08-23 10:20:06 -03:00

703 lines
29 KiB
TypeScript

/**
* Unit tests for the Codex app-server WS transport (CodexAppServerExecutor).
*
* Everything is exercised against a MOCK ws transport (no live connection):
* - gating: isCodexAppServerRequired selects the app-server path only when
* codexTransport==="app-server" (+ config + flag on)
* - lifecycle: the turn emits initialize → thread/start → turn/start in order
* - stall-guard: an inbound server approval request is auto-approved
* - mapping: notifications map to the correct AdapterEvents
* - bridge: streaming output is a valid SSE Response
*/
import test from "node:test";
import assert from "node:assert/strict";
import { isCodexAppServerRequired } from "../../open-sse/executors/codex.ts";
import { CodexAppServerExecutor } from "../../open-sse/executors/codex-app-server.ts";
import {
CodexAppServerClient,
type CodexWreqWebSocket,
} from "../../open-sse/executors/codex/appServerClient.ts";
import {
translateNotification,
translateToolCall,
dynamicToolWireName,
mapUsage,
} from "../../open-sse/executors/codex/appServerEvents.ts";
import { resolveAppServerConfig } from "../../open-sse/executors/codex/appServerConfig.ts";
import { probeCodexAppServerAuth } from "../../open-sse/executors/codex/appServerAuthProbe.ts";
import type { AdapterEvent } from "../../open-sse/vendor/codex-chatgpt-web/types.ts";
import type { ExecuteInput } from "../../open-sse/executors/base.ts";
// ── A scriptable fake wreq WebSocket ────────────────────────────────────────
// Records every frame the client sends, and lets the test drive server frames in.
interface FakeSocketController {
socket: CodexWreqWebSocket;
sent: Array<Record<string, unknown>>;
emit: (frame: Record<string, unknown>) => void;
emitError: (message: string) => void;
emitClose: () => void;
closed: boolean;
}
function makeFakeSocket(): FakeSocketController {
const sent: Array<Record<string, unknown>> = [];
const ctrl: FakeSocketController = {
sent,
closed: false,
socket: null as unknown as CodexWreqWebSocket,
emit: () => {},
emitError: () => {},
emitClose: () => {},
};
const socket: CodexWreqWebSocket = {
send: (data: string) => {
sent.push(JSON.parse(data));
},
close: () => {
ctrl.closed = true;
},
onmessage: null,
onerror: null,
onclose: null,
};
ctrl.socket = socket;
ctrl.emit = (frame) => socket.onmessage?.({ data: JSON.stringify(frame) });
ctrl.emitError = (message) => socket.onerror?.({ message });
ctrl.emitClose = () => socket.onclose?.();
return ctrl;
}
/** A websocketFn that hands out a pre-made fake socket and records the connect opts. */
function fakeTransport(ctrl: FakeSocketController) {
const calls: Array<{ url: string; opts?: Record<string, unknown> }> = [];
const fn = async (url: string, opts?: Record<string, unknown>) => {
calls.push({ url, opts });
return ctrl.socket;
};
return { fn, calls };
}
const APP_SERVER_PSD = {
codexTransport: "app-server",
codexAppServerUrl: "ws://ts-egress:1456",
codexAppServerToken: "deadbeef",
codexAppServerCwd: "/tmp",
};
function makeExecuteInput(overrides: Partial<ExecuteInput> = {}): ExecuteInput {
return {
model: "gpt-5.5",
body: { input: "hello there" },
stream: true,
credentials: { providerSpecificData: { ...APP_SERVER_PSD } },
...overrides,
} as ExecuteInput;
}
// ── Gating ──────────────────────────────────────────────────────────────────
test("isCodexAppServerRequired: true only when codexTransport==='app-server' + configured", () => {
assert.equal(
isCodexAppServerRequired({ providerSpecificData: { ...APP_SERVER_PSD } }),
true
);
// wrong transport
assert.equal(
isCodexAppServerRequired({
providerSpecificData: { ...APP_SERVER_PSD, codexTransport: "websocket" },
}),
false
);
// no providerSpecificData
assert.equal(isCodexAppServerRequired({}), false);
// transport set but not configured (no url/token, no env)
const prevUrl = process.env.OMNIROUTE_CODEX_APPSERVER_WS;
const prevTok = process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN;
const prevTokFile = process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN_FILE;
delete process.env.OMNIROUTE_CODEX_APPSERVER_WS;
delete process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN;
delete process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN_FILE;
try {
assert.equal(
isCodexAppServerRequired({ providerSpecificData: { codexTransport: "app-server" } }),
false
);
} finally {
if (prevUrl !== undefined) process.env.OMNIROUTE_CODEX_APPSERVER_WS = prevUrl;
if (prevTok !== undefined) process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN = prevTok;
if (prevTokFile !== undefined)
process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN_FILE = prevTokFile;
}
});
test("isCodexAppServerRequired: false when OMNIROUTE_CODEX_APP_SERVER_ENABLED=false", () => {
const prev = process.env.OMNIROUTE_CODEX_APP_SERVER_ENABLED;
process.env.OMNIROUTE_CODEX_APP_SERVER_ENABLED = "false";
try {
assert.equal(
isCodexAppServerRequired({ providerSpecificData: { ...APP_SERVER_PSD } }),
false
);
} finally {
if (prev === undefined) delete process.env.OMNIROUTE_CODEX_APP_SERVER_ENABLED;
else process.env.OMNIROUTE_CODEX_APP_SERVER_ENABLED = prev;
}
});
test("resolveAppServerConfig: env fallback + token-file, ws-scheme validation", () => {
assert.equal(resolveAppServerConfig({ codexAppServerUrl: "http://x", codexAppServerToken: "t" }), null);
const cfg = resolveAppServerConfig({ ...APP_SERVER_PSD });
assert.deepEqual(cfg, { url: "ws://ts-egress:1456", token: "deadbeef", cwd: "/tmp" });
});
// ── Notification → AdapterEvent mapping ─────────────────────────────────────
test("translateNotification: maps deltas, done and error to AdapterEvents", () => {
const events: AdapterEvent[] = [];
const push = (e: AdapterEvent) => events.push(e);
assert.equal(
translateNotification("item/agentMessage/delta", { delta: "Hel" }, push),
false
);
assert.equal(
translateNotification("item/reasoning/textDelta", { delta: "think" }, push),
false
);
// terminal → returns true
assert.equal(
translateNotification(
"turn/completed",
{ turn: { usage: { input_tokens: 10, output_tokens: 5 } } },
push
),
true
);
assert.deepEqual(events[0], { type: "text_delta", text: "Hel" });
assert.deepEqual(events[1], { type: "thinking_delta", thinking: "think" });
assert.equal(events[2].type, "done");
const done = events[2] as Extract<AdapterEvent, { type: "done" }>;
assert.equal(done.endTurn, true);
assert.equal(done.usage?.inputTokens, 10);
assert.equal(done.usage?.outputTokens, 5);
});
test("translateNotification: error notification maps to error event (terminal)", () => {
const events: AdapterEvent[] = [];
const isTerminal = translateNotification(
"error",
{ error: { message: "boom" } },
(e) => events.push(e)
);
assert.equal(isTerminal, true);
assert.equal(events[0].type, "error");
const err = events[0] as Extract<AdapterEvent, { type: "error" }>;
assert.equal(err.message, "boom");
assert.equal(err.status, 502);
});
test("mapUsage: converts snake_case token counts", () => {
const usage = mapUsage({
input_tokens: 100,
cached_input_tokens: 20,
output_tokens: 40,
reasoning_output_tokens: 8,
});
assert.equal(usage?.inputTokens, 100);
assert.equal(usage?.cachedInputTokens, 20);
assert.equal(usage?.cacheReadInputTokens, 20);
assert.equal(usage?.outputTokens, 40);
assert.equal(usage?.reasoningOutputTokens, 8);
assert.equal(mapUsage(undefined), undefined);
});
// ── Client: stall-guard auto-approval ───────────────────────────────────────
test("CodexAppServerClient: server approval request is auto-approved", async () => {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const client = new CodexAppServerClient({ websocketFn: fn });
await client.connect("ws://x", "tok");
// Auth header attached on connect
// (the fake records opts on connect via fakeTransport calls; verified in lifecycle test)
// Server sends an exec approval request with id=99.
ctrl.emit({
jsonrpc: "2.0",
id: 99,
method: "execCommandApproval",
params: { command: ["ls", "-la"], cwd: "/tmp" },
});
const reply = ctrl.sent.find((f) => f.id === 99);
assert.ok(reply, "client must reply to the server approval request");
// OmniRoute is a router: approvals are auto-APPROVED so the model's agentic
// tool calls proceed; the harness downstream is the real execution gate.
assert.equal(
(reply!.result as Record<string, unknown>).decision,
"approved"
);
});
test("CodexAppServerClient: non-approval server request gets a JSON-RPC error", async () => {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const client = new CodexAppServerClient({ websocketFn: fn });
await client.connect("ws://x", "tok");
ctrl.emit({ jsonrpc: "2.0", id: 7, method: "item/tool/call", params: {} });
const reply = ctrl.sent.find((f) => f.id === 7);
assert.ok(reply);
assert.equal((reply!.error as { code: number }).code, -32601);
});
test("CodexAppServerClient: notifications reach the handler; responses settle requests", async () => {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const client = new CodexAppServerClient({ websocketFn: fn });
await client.connect("ws://x", "tok");
const seen: string[] = [];
client.onNotification((method) => seen.push(method));
// Fire a request; the fake echoes an id-matched response.
const reqPromise = client.request("initialize", { clientInfo: {} });
const sentInit = ctrl.sent.find((f) => f.method === "initialize");
assert.ok(sentInit);
ctrl.emit({ jsonrpc: "2.0", id: sentInit!.id, result: { ok: true } });
const result = (await reqPromise) as { ok: boolean };
assert.equal(result.ok, true);
// A method-only frame is a notification.
ctrl.emit({ jsonrpc: "2.0", method: "item/agentMessage/delta", params: { delta: "x" } });
assert.ok(seen.includes("item/agentMessage/delta"));
});
// ── Executor: lifecycle order + streaming SSE Response ───────────────────────
/** Drive a full streaming turn against a fake transport and return the SSE text. */
async function runStreamingTurn(): Promise<{
sent: Array<Record<string, unknown>>;
sseText: string;
}> {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const executor = new CodexAppServerExecutor({ websocketFn: fn });
// Auto-responder: as soon as the client sends a request, emit its response and,
// for turn/start, stream a couple of notifications + turn/completed.
const originalSend = ctrl.socket.send;
ctrl.socket.send = (data: string) => {
originalSend(data);
const frame = JSON.parse(data) as Record<string, unknown>;
if (frame.id == null || !frame.method) return;
queueMicrotask(() => {
if (frame.method === "thread/start") {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: { threadId: "thr_1" } });
} else if (frame.method === "turn/start") {
ctrl.emit({
jsonrpc: "2.0",
method: "item/agentMessage/delta",
params: { delta: "Hello" },
});
ctrl.emit({
jsonrpc: "2.0",
method: "turn/completed",
params: { turn: { usage: { input_tokens: 3, output_tokens: 2 } } },
});
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
} else {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
}
});
};
const result = await executor.execute(makeExecuteInput());
const response = "response" in result ? result.response : result;
assert.equal(response.status, 200);
assert.match(response.headers.get("Content-Type") ?? "", /text\/event-stream/);
const sseText = await response.text();
return { sent: ctrl.sent, sseText };
}
test("CodexAppServerExecutor: streaming turn emits initialize → thread/start → turn/start in order", async () => {
const { sent } = await runStreamingTurn();
const methods = sent.filter((f) => typeof f.method === "string" && f.id != null).map((f) => f.method);
const lifecycle = methods.filter(
(m) => m === "initialize" || m === "thread/start" || m === "turn/start"
);
assert.deepEqual(lifecycle, ["initialize", "thread/start", "turn/start"]);
// thread/start carried the router defaults: approvalPolicy:"never" (codex
// never blocks on its own approval) + sandbox:"danger-full-access" (codex's
// own sandbox does not gate the model; the harness is the real execution gate).
const threadStart = sent.find((f) => f.method === "thread/start");
assert.equal((threadStart!.params as Record<string, unknown>).approvalPolicy, "never");
assert.equal((threadStart!.params as Record<string, unknown>).sandbox, "danger-full-access");
// turn/start carried the text input with text_elements:[]
const turnStart = sent.find((f) => f.method === "turn/start");
const turnParams = turnStart!.params as Record<string, unknown>;
assert.equal(turnParams.threadId, "thr_1");
assert.deepEqual(turnParams.input, [{ type: "text", text: "hello there", text_elements: [] }]);
});
test("CodexAppServerExecutor: streaming output is a valid Responses SSE stream", async () => {
const { sseText } = await runStreamingTurn();
assert.match(sseText, /event: response\.created/);
assert.match(sseText, /response\.output_text\.delta/);
assert.ok(sseText.includes("Hello"));
assert.match(sseText, /event: response\.completed/);
assert.ok(sseText.includes("[DONE]"));
});
test("CodexAppServerExecutor: unconfigured connection returns an in-band error Response", async () => {
const executor = new CodexAppServerExecutor({ websocketFn: async () => makeFakeSocket().socket });
const prevUrl = process.env.OMNIROUTE_CODEX_APPSERVER_WS;
const prevTok = process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN;
const prevTokFile = process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN_FILE;
delete process.env.OMNIROUTE_CODEX_APPSERVER_WS;
delete process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN;
delete process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN_FILE;
try {
const result = await executor.execute(
makeExecuteInput({ credentials: { providerSpecificData: { codexTransport: "app-server" } } })
);
const response = "response" in result ? result.response : result;
assert.equal(response.status, 503);
} finally {
if (prevUrl !== undefined) process.env.OMNIROUTE_CODEX_APPSERVER_WS = prevUrl;
if (prevTok !== undefined) process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN = prevTok;
if (prevTokFile !== undefined)
process.env.OMNIROUTE_CODEX_APPSERVER_WS_TOKEN_FILE = prevTokFile;
}
});
test("CodexAppServerExecutor: non-streaming turn returns a JSON Response", async () => {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const executor = new CodexAppServerExecutor({ websocketFn: fn });
const originalSend = ctrl.socket.send;
ctrl.socket.send = (data: string) => {
originalSend(data);
const frame = JSON.parse(data) as Record<string, unknown>;
if (frame.id == null || !frame.method) return;
queueMicrotask(() => {
if (frame.method === "thread/start") {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: { threadId: "thr_1" } });
} else if (frame.method === "turn/start") {
ctrl.emit({
jsonrpc: "2.0",
method: "item/agentMessage/delta",
params: { delta: "Hi" },
});
ctrl.emit({ jsonrpc: "2.0", method: "turn/completed", params: { turn: {} } });
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
} else {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
}
});
};
const result = await executor.execute(makeExecuteInput({ stream: false }));
const response = "response" in result ? result.response : result;
assert.equal(response.status, 200);
assert.match(response.headers.get("Content-Type") ?? "", /application\/json/);
const body = (await response.json()) as Record<string, unknown>;
assert.ok(Array.isArray(body.output));
});
// ── Tool path: INBOUND advertise + OUTBOUND passthrough ──────────────────────
test("dynamicToolWireName: flattens namespaced tools, passes plain ones through", () => {
assert.equal(dynamicToolWireName("mcp__ctx7", "get_docs"), "mcp__ctx7__get_docs");
assert.equal(dynamicToolWireName(null, "read_file"), "read_file");
assert.equal(dynamicToolWireName(undefined, "read_file"), "read_file");
});
test("translateToolCall: emits tool_call_start/delta/end with callId, wire name, JSON args", () => {
const events: AdapterEvent[] = [];
translateToolCall(
{ callId: "call_42", namespace: null, tool: "get_weather", arguments: { city: "SF" } },
(e) => events.push(e)
);
assert.equal(events.length, 3);
assert.deepEqual(events[0], { type: "tool_call_start", id: "call_42", name: "get_weather" });
assert.deepEqual(events[1], { type: "tool_call_delta", arguments: '{"city":"SF"}' });
assert.deepEqual(events[2], { type: "tool_call_end" });
});
test("translateToolCall: restores MCP namespace into the wire name for the round-trip", () => {
const events: AdapterEvent[] = [];
translateToolCall(
{ callId: "call_9", namespace: "mcp__ctx7", tool: "get_docs", arguments: "{}" },
(e) => events.push(e)
);
const start = events[0] as Extract<AdapterEvent, { type: "tool_call_start" }>;
assert.equal(start.name, "mcp__ctx7__get_docs");
});
/**
* Drive a streaming turn where the harness advertises a function tool and codex
* responds by invoking it via the `item/tool/call` ServerRequest. Assert (a) the
* tool is advertised on thread/start via `dynamicTools`, (b) the app-server request
* is settled, and (c) the SSE stream carries a Responses function_call for the tool.
*/
async function runToolTurn(): Promise<{
sent: Array<Record<string, unknown>>;
sseText: string;
}> {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const executor = new CodexAppServerExecutor({ websocketFn: fn });
const originalSend = ctrl.socket.send;
ctrl.socket.send = (data: string) => {
originalSend(data);
const frame = JSON.parse(data) as Record<string, unknown>;
if (frame.id == null || !frame.method) return;
queueMicrotask(() => {
if (frame.method === "thread/start") {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: { threadId: "thr_1" } });
} else if (frame.method === "turn/start") {
// codex invokes the harness tool via a server → client ServerRequest.
ctrl.emit({
jsonrpc: "2.0",
id: 5000,
method: "item/tool/call",
params: {
threadId: "thr_1",
turnId: "turn_1",
callId: "call_abc",
namespace: null,
tool: "get_weather",
arguments: { city: "SF" },
},
});
// Settle turn/start too (codex would eventually complete; the passthrough
// already ended the turn on our side).
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
} else {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
}
});
};
const input = makeExecuteInput({
body: {
input: "what's the weather?",
tools: [
{
type: "function",
name: "get_weather",
description: "Get the weather for a city",
parameters: {
type: "object",
properties: { city: { type: "string" } },
required: ["city"],
},
},
],
},
});
const result = await executor.execute(input);
const response = "response" in result ? result.response : result;
const sseText = await response.text();
return { sent: ctrl.sent, sseText };
}
test("CodexAppServerExecutor: advertises harness tools to codex via thread/start dynamicTools", async () => {
const { sent } = await runToolTurn();
const threadStart = sent.find((f) => f.method === "thread/start");
assert.ok(threadStart, "thread/start must be sent");
const params = threadStart!.params as Record<string, unknown>;
const dynamicTools = params.dynamicTools as Array<Record<string, unknown>> | undefined;
assert.ok(Array.isArray(dynamicTools), "dynamicTools must be advertised");
assert.equal(dynamicTools!.length, 1);
assert.equal(dynamicTools![0].type, "function");
assert.equal(dynamicTools![0].name, "get_weather");
assert.ok(dynamicTools![0].inputSchema, "spec carries the inputSchema");
// experimentalApi capability opted in on initialize (dynamicTools is experimental)
const init = sent.find((f) => f.method === "initialize");
const caps = (init!.params as Record<string, unknown>).capabilities as Record<string, unknown>;
assert.equal(caps.experimentalApi, true);
});
test("CodexAppServerExecutor: item/tool/call is settled and surfaced as a Responses function_call", async () => {
const { sent, sseText } = await runToolTurn();
// The app-server request (id 5000) must be settled so the socket never stalls.
const toolReply = sent.find((f) => f.id === 5000);
assert.ok(toolReply, "the item/tool/call request id must be settled");
const replyResult = toolReply!.result as Record<string, unknown>;
assert.ok(replyResult, "settled with a DynamicToolCallResponse result");
assert.equal(replyResult.success, false);
assert.ok(Array.isArray(replyResult.contentItems));
// The SSE stream carries the harness function_call for get_weather with its args.
assert.match(sseText, /function_call/);
assert.ok(sseText.includes("get_weather"));
assert.ok(sseText.includes("call_abc"), "the codex callId is relayed as the call_id");
assert.ok(sseText.includes("SF"), "the tool arguments are relayed");
assert.match(sseText, /event: response\.completed/);
assert.ok(sseText.includes("[DONE]"));
});
// REGRESSION (live BUG#3, 2026-08-22): the real codex app-server ACCEPTS a turn
// on turn/start (returns status:"inProgress") and delivers the model output +
// terminal turn/completed LATER as async notifications. The original run() closed
// the WS in its finally-block as soon as `await turn/start` resolved, tearing the
// socket down BEFORE those notifications arrived, so the event queue never closed
// and the request hung until the caller's timeout. The pre-existing mocks hid this
// because they emitted turn/completed in the SAME microtask as the turn/start
// response (completion raced ahead of request-resolution). This test reproduces
// the real ordering: turn/start resolves FIRST, then agentMessage/delta +
// turn/completed fire on a later macrotask. It must still complete (not hang).
test("CodexAppServerExecutor: async post-turn/start completion does not close the socket early (BUG#3)", async () => {
const ctrl = makeFakeSocket();
const { fn } = fakeTransport(ctrl);
const executor = new CodexAppServerExecutor({ websocketFn: fn });
// Model a REAL socket: once closed, it delivers no more frames. The shared
// makeFakeSocket keeps emitting after close (fine for the other tests), but
// this regression turns specifically on the fact that a prematurely-closed
// socket DROPS the later turn/completed — so guard emits on ctrl.closed here.
const emitLive = (frame: Record<string, unknown>) => {
if (ctrl.closed) return; // socket torn down → frame never arrives (real behavior)
ctrl.emit(frame);
};
const originalSend = ctrl.socket.send;
ctrl.socket.send = (data: string) => {
originalSend(data);
const frame = JSON.parse(data) as Record<string, unknown>;
if (frame.id == null || !frame.method) return;
if (frame.method === "thread/start") {
queueMicrotask(() =>
emitLive({ jsonrpc: "2.0", id: frame.id, result: { thread: { id: "thr_async" } } })
);
} else if (frame.method === "turn/start") {
// Resolve turn/start FIRST (status inProgress) …
queueMicrotask(() =>
emitLive({
jsonrpc: "2.0",
id: frame.id,
result: { turn: { id: "t1", status: "inProgress" } },
})
);
// … then, on a LATER macrotask, stream the output + terminal completion.
// Under the OLD code the finally-block closes the socket right after
// turn/start resolves, so ctrl.closed is true here and these frames are
// DROPPED → the queue never closes → execute() hangs (test times out).
setTimeout(() => {
emitLive({
jsonrpc: "2.0",
method: "item/agentMessage/delta",
params: { delta: "ASYNC-OK" },
});
emitLive({
jsonrpc: "2.0",
method: "turn/completed",
params: { turn: { usage: { input_tokens: 1, output_tokens: 1 } } },
});
}, 15);
} else {
queueMicrotask(() => emitLive({ jsonrpc: "2.0", id: frame.id, result: {} }));
}
};
// Non-streaming: execute() awaits events.collect(), which only returns once the
// queue closes on the terminal notification. Under the old (buggy) code the
// socket closed early, the terminal frame was dropped, and this promise never
// resolved. Guard with a timeout so a regression fails loudly, not by hanging.
const result = await Promise.race([
executor.execute(makeExecuteInput({ stream: false })),
new Promise<never>((_, reject) =>
setTimeout(() => reject(new Error("execute() hung: socket closed before async completion (BUG#3 regressed)")), 5000)
),
]);
const response = "response" in result ? result.response : (result as Response);
assert.equal(response.status, 200);
const body = JSON.parse(await response.text()) as {
status?: string;
output?: Array<{ content?: Array<{ text?: string }> }>;
};
assert.equal(body.status, "completed", "the turn completed after the async terminal notification");
const text = body.output?.[0]?.content?.[0]?.text ?? "";
assert.equal(text, "ASYNC-OK", "the model output that arrived AFTER turn/start is present");
});
// ── Layer-2 auth-status probe (probeCodexAppServerAuth) ─────────────────────
// /readyz proves the server PROCESS is up but NOT that its Codex CLI is signed
// in. probeCodexAppServerAuth opens the JSON-RPC WS and reads account/read:
// authenticated → { account: {...} } ; logged out → no account (or auth error).
// Verified against codex 0.149.0: account/read returns
// { account: { type, email, planType }, requiresOpenaiAuth }.
/** A fake websocketFn that answers initialize + account/read with a scripted result. */
function fakeAuthTransport(accountReadResponse: {
result?: Record<string, unknown>;
error?: { code: number; message: string };
}) {
const ctrl = makeFakeSocket();
const originalSend = ctrl.socket.send;
ctrl.socket.send = (data: string) => {
originalSend(data);
const frame = JSON.parse(data) as Record<string, unknown>;
if (frame.id == null || !frame.method) return;
queueMicrotask(() => {
if (frame.method === "initialize") {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: { ok: true } });
} else if (frame.method === "account/read") {
if (accountReadResponse.error) {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, error: accountReadResponse.error });
} else {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: accountReadResponse.result ?? {} });
}
} else {
ctrl.emit({ jsonrpc: "2.0", id: frame.id, result: {} });
}
});
};
const fn = async () => ctrl.socket;
return fn;
}
const AUTH_CONFIG = { url: "ws://ts-egress:1456", token: "deadbeef", cwd: "/tmp" };
test("probeCodexAppServerAuth: account with email → authenticated", async () => {
const fn = fakeAuthTransport({
result: { account: { type: "chatgpt", email: "user@example.com", planType: "pro" }, requiresOpenaiAuth: true },
});
const status = await probeCodexAppServerAuth(AUTH_CONFIG, fn, 3000);
assert.equal(status.state, "authenticated");
if (status.state === "authenticated") {
assert.equal(status.account.email, "user@example.com");
assert.equal(status.account.planType, "pro");
}
});
test("probeCodexAppServerAuth: no account → logged_out", async () => {
const fn = fakeAuthTransport({ result: { requiresOpenaiAuth: true } }); // no `account`
const status = await probeCodexAppServerAuth(AUTH_CONFIG, fn, 3000);
assert.equal(status.state, "logged_out");
});
test("probeCodexAppServerAuth: auth-error on account/read → logged_out", async () => {
const fn = fakeAuthTransport({ error: { code: -32000, message: "AuthRequiredError: please login" } });
const status = await probeCodexAppServerAuth(AUTH_CONFIG, fn, 3000);
assert.equal(status.state, "logged_out");
});
test("probeCodexAppServerAuth: no transport → unknown (does not throw)", async () => {
const status = await probeCodexAppServerAuth(AUTH_CONFIG, null, 3000);
assert.equal(status.state, "unknown");
});