← Files Codex ReplayARCHIVED FILE
mcp/codex-worker.mjs
42.2 KB · Oct 2, 2026 · 00:19 UTC
#!/usr/bin/env node
// src/codex-worker.mjs
import { createInterface } from "node:readline";
import { spawn as spawn2, spawnSync } from "node:child_process";
import { accessSync, constants, readdirSync, statSync as statSync2 } from "node:fs";
import { homedir } from "node:os";
import path3 from "node:path";
import { fileURLToPath, pathToFileURL } from "node:url";
// node_modules/@openai/codex-sdk/dist/index.js
import { promises as fs } from "fs";
import os from "os";
import path from "path";
import { spawn } from "child_process";
import { statSync } from "fs";
import path2 from "path";
import readline from "readline";
import { createRequire } from "module";
async function createOutputSchemaFile(schema) {
if (schema === void 0) {
return { cleanup: async () => {
} };
}
if (!isJsonObject(schema)) {
throw new Error("outputSchema must be a plain JSON object");
}
const schemaDir = await fs.mkdtemp(path.join(os.tmpdir(), "codex-output-schema-"));
const schemaPath = path.join(schemaDir, "schema.json");
const cleanup = async () => {
try {
await fs.rm(schemaDir, { recursive: true, force: true });
} catch {
}
};
try {
await fs.writeFile(schemaPath, JSON.stringify(schema), "utf8");
return { schemaPath, cleanup };
} catch (error) {
await cleanup();
throw error;
}
}
function isJsonObject(value) {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
var Thread = class {
_exec;
_options;
_id;
_threadOptions;
/** Returns the ID of the thread. Populated after the first turn starts. */
get id() {
return this._id;
}
/* @internal */
constructor(exec, options, threadOptions, id = null) {
this._exec = exec;
this._options = options;
this._id = id;
this._threadOptions = threadOptions;
}
/** Provides the input to the agent and streams events as they are produced during the turn. */
async runStreamed(input, turnOptions = {}) {
return { events: this.runStreamedInternal(input, turnOptions) };
}
async *runStreamedInternal(input, turnOptions = {}) {
const { schemaPath, cleanup } = await createOutputSchemaFile(turnOptions.outputSchema);
const options = this._threadOptions;
const { prompt, images } = normalizeInput(input);
const generator = this._exec.run({
input: prompt,
baseUrl: this._options.baseUrl,
apiKey: this._options.apiKey,
threadId: this._id,
images,
model: options?.model,
sandboxMode: options?.sandboxMode,
workingDirectory: options?.workingDirectory,
skipGitRepoCheck: options?.skipGitRepoCheck,
outputSchemaFile: schemaPath,
modelReasoningEffort: options?.modelReasoningEffort,
signal: turnOptions.signal,
networkAccessEnabled: options?.networkAccessEnabled,
webSearchMode: options?.webSearchMode,
webSearchEnabled: options?.webSearchEnabled,
approvalPolicy: options?.approvalPolicy,
additionalDirectories: options?.additionalDirectories
});
try {
for await (const item of generator) {
let parsed;
try {
parsed = JSON.parse(item);
} catch (error) {
throw new Error(`Failed to parse item: ${item}`, { cause: error });
}
if (parsed.type === "thread.started") {
this._id = parsed.thread_id;
}
yield parsed;
}
} finally {
await cleanup();
}
}
/** Provides the input to the agent and returns the completed turn. */
async run(input, turnOptions = {}) {
const generator = this.runStreamedInternal(input, turnOptions);
const items = [];
let finalResponse = "";
let usage = null;
let turnFailure = null;
for await (const event of generator) {
if (event.type === "item.completed") {
if (event.item.type === "agent_message") {
finalResponse = event.item.text;
}
items.push(event.item);
} else if (event.type === "turn.completed") {
usage = event.usage;
} else if (event.type === "turn.failed") {
turnFailure = event.error;
break;
}
}
if (turnFailure) {
throw new Error(turnFailure.message);
}
return { items, finalResponse, usage };
}
};
function normalizeInput(input) {
if (typeof input === "string") {
return { prompt: input, images: [] };
}
const promptParts = [];
const images = [];
for (const item of input) {
if (item.type === "text") {
promptParts.push(item.text);
} else if (item.type === "local_image") {
images.push(item.path);
}
}
return { prompt: promptParts.join("\n\n"), images };
}
var INTERNAL_ORIGINATOR_ENV = "CODEX_INTERNAL_ORIGINATOR_OVERRIDE";
var TYPESCRIPT_SDK_ORIGINATOR = "codex_sdk_ts";
var CODEX_NPM_NAME = "@openai/codex";
var PLATFORM_PACKAGE_BY_TARGET = {
"x86_64-unknown-linux-musl": "@openai/codex-linux-x64",
"aarch64-unknown-linux-musl": "@openai/codex-linux-arm64",
"x86_64-apple-darwin": "@openai/codex-darwin-x64",
"aarch64-apple-darwin": "@openai/codex-darwin-arm64",
"x86_64-pc-windows-msvc": "@openai/codex-win32-x64",
"aarch64-pc-windows-msvc": "@openai/codex-win32-arm64"
};
var moduleRequire = createRequire(import.meta.url);
var CodexExec = class {
executablePath;
pathDirs;
envOverride;
configOverrides;
constructor(executablePath = null, env, configOverrides) {
if (executablePath) {
this.executablePath = executablePath;
this.pathDirs = [];
} else {
const resolved = findCodexPath();
this.executablePath = resolved.executablePath;
this.pathDirs = resolved.pathDirs;
}
this.envOverride = env;
this.configOverrides = configOverrides;
}
async *run(args) {
const commandArgs = ["exec", "--experimental-json"];
if (this.configOverrides) {
for (const override of serializeConfigOverrides(this.configOverrides)) {
commandArgs.push("--config", override);
}
}
if (args.baseUrl) {
commandArgs.push(
"--config",
`openai_base_url=${toTomlValue(args.baseUrl, "openai_base_url")}`
);
}
if (args.model) {
commandArgs.push("--model", args.model);
}
if (args.sandboxMode) {
commandArgs.push("--sandbox", args.sandboxMode);
}
if (args.workingDirectory) {
commandArgs.push("--cd", args.workingDirectory);
}
if (args.additionalDirectories?.length) {
for (const dir of args.additionalDirectories) {
commandArgs.push("--add-dir", dir);
}
}
if (args.skipGitRepoCheck) {
commandArgs.push("--skip-git-repo-check");
}
if (args.outputSchemaFile) {
commandArgs.push("--output-schema", args.outputSchemaFile);
}
if (args.modelReasoningEffort) {
commandArgs.push("--config", `model_reasoning_effort="${args.modelReasoningEffort}"`);
}
if (args.networkAccessEnabled !== void 0) {
commandArgs.push(
"--config",
`sandbox_workspace_write.network_access=${args.networkAccessEnabled}`
);
}
if (args.webSearchMode) {
commandArgs.push("--config", `web_search="${args.webSearchMode}"`);
} else if (args.webSearchEnabled === true) {
commandArgs.push("--config", `web_search="live"`);
} else if (args.webSearchEnabled === false) {
commandArgs.push("--config", `web_search="disabled"`);
}
if (args.approvalPolicy) {
commandArgs.push("--config", `approval_policy="${args.approvalPolicy}"`);
}
if (args.threadId) {
commandArgs.push("resume", args.threadId);
}
if (args.images?.length) {
for (const image of args.images) {
commandArgs.push("--image", image);
}
}
const env = {};
if (this.envOverride) {
Object.assign(env, this.envOverride);
} else {
for (const [key, value] of Object.entries(process.env)) {
if (value !== void 0) {
env[key] = value;
}
}
}
if (!env[INTERNAL_ORIGINATOR_ENV]) {
env[INTERNAL_ORIGINATOR_ENV] = TYPESCRIPT_SDK_ORIGINATOR;
}
if (args.apiKey) {
env.CODEX_API_KEY = args.apiKey;
}
if (this.pathDirs.length > 0) {
prependPathDirs(env, this.pathDirs);
}
const child = spawn(this.executablePath, commandArgs, {
env,
signal: args.signal
});
let spawnError = null;
child.once("error", (err) => spawnError = err);
if (!child.stdin) {
child.kill();
throw new Error("Child process has no stdin");
}
child.stdin.write(args.input);
child.stdin.end();
if (!child.stdout) {
child.kill();
throw new Error("Child process has no stdout");
}
const stderrChunks = [];
if (child.stderr) {
child.stderr.on("data", (data) => {
stderrChunks.push(data);
});
}
const exitPromise = new Promise(
(resolve) => {
child.once("exit", (code, signal) => {
resolve({ code, signal });
});
}
);
const rl = readline.createInterface({
input: child.stdout,
crlfDelay: Infinity
});
try {
for await (const line of rl) {
yield line;
}
if (spawnError) throw spawnError;
const { code, signal } = await exitPromise;
if (code !== 0 || signal) {
const stderrBuffer = Buffer.concat(stderrChunks);
const detail = signal ? `signal ${signal}` : `code ${code ?? 1}`;
throw new Error(`Codex Exec exited with ${detail}: ${stderrBuffer.toString("utf8")}`);
}
} finally {
rl.close();
child.removeAllListeners();
try {
if (!child.killed) child.kill();
} catch {
}
}
}
};
function serializeConfigOverrides(configOverrides) {
const overrides = [];
flattenConfigOverrides(configOverrides, "", overrides);
return overrides;
}
function flattenConfigOverrides(value, prefix, overrides) {
if (!isPlainObject(value)) {
if (prefix) {
overrides.push(`${prefix}=${toTomlValue(value, prefix)}`);
return;
} else {
throw new Error("Codex config overrides must be a plain object");
}
}
const entries = Object.entries(value);
if (!prefix && entries.length === 0) {
return;
}
if (prefix && entries.length === 0) {
overrides.push(`${prefix}={}`);
return;
}
for (const [key, child] of entries) {
if (!key) {
throw new Error("Codex config override keys must be non-empty strings");
}
if (child === void 0) {
continue;
}
const path32 = prefix ? `${prefix}.${key}` : key;
if (isPlainObject(child)) {
flattenConfigOverrides(child, path32, overrides);
} else {
overrides.push(`${path32}=${toTomlValue(child, path32)}`);
}
}
}
function toTomlValue(value, path32) {
if (typeof value === "string") {
return JSON.stringify(value);
} else if (typeof value === "number") {
if (!Number.isFinite(value)) {
throw new Error(`Codex config override at ${path32} must be a finite number`);
}
return `${value}`;
} else if (typeof value === "boolean") {
return value ? "true" : "false";
} else if (Array.isArray(value)) {
const rendered = value.map((item, index) => toTomlValue(item, `${path32}[${index}]`));
return `[${rendered.join(", ")}]`;
} else if (isPlainObject(value)) {
const parts = [];
for (const [key, child] of Object.entries(value)) {
if (!key) {
throw new Error("Codex config override keys must be non-empty strings");
}
if (child === void 0) {
continue;
}
parts.push(`${formatTomlKey(key)} = ${toTomlValue(child, `${path32}.${key}`)}`);
}
return `{${parts.join(", ")}}`;
} else if (value === null) {
throw new Error(`Codex config override at ${path32} cannot be null`);
} else {
const typeName = typeof value;
throw new Error(`Unsupported Codex config override value at ${path32}: ${typeName}`);
}
}
var TOML_BARE_KEY = /^[A-Za-z0-9_-]+$/;
function formatTomlKey(key) {
return TOML_BARE_KEY.test(key) ? key : JSON.stringify(key);
}
function isPlainObject(value) {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
function findCodexPath() {
const { platform, arch } = process;
let targetTriple = null;
switch (platform) {
case "linux":
case "android":
switch (arch) {
case "x64":
targetTriple = "x86_64-unknown-linux-musl";
break;
case "arm64":
targetTriple = "aarch64-unknown-linux-musl";
break;
default:
break;
}
break;
case "darwin":
switch (arch) {
case "x64":
targetTriple = "x86_64-apple-darwin";
break;
case "arm64":
targetTriple = "aarch64-apple-darwin";
break;
default:
break;
}
break;
case "win32":
switch (arch) {
case "x64":
targetTriple = "x86_64-pc-windows-msvc";
break;
case "arm64":
targetTriple = "aarch64-pc-windows-msvc";
break;
default:
break;
}
break;
default:
break;
}
if (!targetTriple) {
throw new Error(`Unsupported platform: ${platform} (${arch})`);
}
const platformPackage = PLATFORM_PACKAGE_BY_TARGET[targetTriple];
if (!platformPackage) {
throw new Error(`Unsupported target triple: ${targetTriple}`);
}
let vendorRoot;
try {
const codexPackageJsonPath = moduleRequire.resolve(`${CODEX_NPM_NAME}/package.json`);
const codexRequire = createRequire(codexPackageJsonPath);
const platformPackageJsonPath = codexRequire.resolve(`${platformPackage}/package.json`);
vendorRoot = path2.join(path2.dirname(platformPackageJsonPath), "vendor");
} catch {
throw new Error(
`Unable to locate Codex CLI binaries. Ensure ${CODEX_NPM_NAME} is installed with optional dependencies.`
);
}
const codexBinaryName = process.platform === "win32" ? "codex.exe" : "codex";
const nativePackage = resolveNativePackage(vendorRoot, targetTriple, codexBinaryName);
if (!nativePackage) {
throw new Error(
`Unable to locate Codex CLI binaries for ${targetTriple}. Ensure ${CODEX_NPM_NAME} is installed with optional dependencies.`
);
}
return nativePackage;
}
function resolveNativePackage(vendorRoot, targetTriple, codexBinaryName) {
const packageRoot = path2.join(vendorRoot, targetTriple);
const packageBinaryPath = path2.join(packageRoot, "bin", codexBinaryName);
if (isFile(packageBinaryPath) && isFile(path2.join(packageRoot, "codex-package.json"))) {
return {
executablePath: packageBinaryPath,
pathDirs: existingDirs(path2.join(packageRoot, "codex-path"))
};
}
const legacyBinaryPath = path2.join(packageRoot, "codex", codexBinaryName);
if (isFile(legacyBinaryPath)) {
return {
executablePath: legacyBinaryPath,
pathDirs: existingDirs(path2.join(packageRoot, "path"))
};
}
return null;
}
function existingDirs(...dirs) {
return dirs.filter(isDirectory);
}
function prependPathDirs(env, pathDirs, platform = process.platform) {
const pathKey = pathEnvKey(env, platform);
if (platform === "win32") {
for (const key of Object.keys(env)) {
if (key.toLowerCase() === "path" && key !== pathKey) {
delete env[key];
}
}
}
const existingEntries = (env[pathKey] ?? "").split(path2.delimiter).filter((entry) => entry.length > 0 && !pathDirs.includes(entry));
env[pathKey] = [...pathDirs, ...existingEntries].join(path2.delimiter);
}
function pathEnvKey(env, platform) {
if (platform !== "win32") {
return "PATH";
}
const matchingKeys = Object.keys(env).filter((key) => key.toLowerCase() === "path");
return matchingKeys.includes("Path") ? "Path" : matchingKeys.at(-1) ?? "PATH";
}
function isFile(filePath) {
try {
return statSync(filePath).isFile();
} catch {
return false;
}
}
function isDirectory(filePath) {
try {
return statSync(filePath).isDirectory();
} catch {
return false;
}
}
var Codex = class {
exec;
options;
constructor(options = {}) {
const { codexPathOverride, env, config } = options;
this.exec = new CodexExec(codexPathOverride, env, config);
this.options = options;
}
/**
* Starts a new conversation with an agent.
* @returns A new thread instance.
*/
startThread(options = {}) {
return new Thread(this.exec, this.options, options);
}
/**
* Resumes a conversation with an agent based on the thread id.
* Threads are persisted in ~/.codex/sessions.
*
* @param id The id of the thread to resume.
* @returns A new thread instance.
*/
resumeThread(id, options = {}) {
return new Thread(this.exec, this.options, options, id);
}
};
// src/codex-worker.mjs
var PROTOCOL_VERSION = 1;
var MAX_PROMPT_CHARS = 2e6;
var MAX_SCHEMA_CHARS = 2e5;
var MAX_FINAL_RESPONSE_CHARS = 2e5;
var MAX_LIFECYCLE_EVENTS = 128;
var MAX_PROVIDER_ERROR_CHARS = 1e3;
var MAX_WORKER_DIAGNOSTIC_CHARS = 8e3;
var CLI_WRAPPER_MODE_ENV = "CODEX_REPLAY_CODEX_WRAPPER";
var CLI_WRAPPER_TARGET_ENV = "CODEX_REPLAY_CODEX_TARGET";
var CLI_WRAPPER_OWNER_ENV = "CODEX_REPLAY_CODEX_OWNER_PID";
var TREE_KILL_GRACE_MS = 2e3;
var CODEX_VERSION_PROBE_TIMEOUT_MS = 5e3;
var CODEX_VERSION_PROBE_MAX_BUFFER = 64 * 1024;
var SAFE_ID = /^[A-Za-z0-9._:-]{1,128}$/;
var SAFE_MODEL = /^[A-Za-z0-9._:-]{1,128}$/;
var SANDBOX_MODES = /* @__PURE__ */ new Set(["read-only", "workspace-write"]);
var REASONING_EFFORTS = /* @__PURE__ */ new Set([
"minimal",
"low",
"medium",
"high",
"xhigh"
]);
var SAFE_ITEM_TYPES = /* @__PURE__ */ new Set([
"agent_message",
"command_execution",
"error",
"file_change",
"mcp_tool_call",
"reasoning",
"todo_list",
"web_search"
]);
var SYSTEM_CODES = /* @__PURE__ */ new Set([
"ECONNRESET",
"ECONNREFUSED",
"ETIMEDOUT",
"ENOTFOUND",
"EAI_AGAIN",
"ENOENT",
"EACCES",
"EPERM",
"EPIPE"
]);
var PERMANENT_SYSTEM_CODES = /* @__PURE__ */ new Set(["ENOENT", "EACCES", "EPERM"]);
var PERMANENT_STREAM_ERROR_KINDS = /* @__PURE__ */ new Set([
"invalid_request_error",
"authentication_error",
"permission_error",
"not_found_error",
"model_not_found",
"insufficient_quota",
"billing_error",
"account_deactivated",
"access_denied"
]);
var RETRYABLE_STREAM_ERROR_KINDS = /* @__PURE__ */ new Set([
"connection_error",
"timeout_error",
"transport_error",
"rate_limit_error",
"server_error",
"service_unavailable_error",
"internal_server_error",
"overloaded_error"
]);
var SafeWorkerError = class extends Error {
constructor(code, message, options = {}) {
const { retryable = false, ...errorOptions } = options;
super(message, errorOptions);
this.name = "SafeWorkerError";
this.code = code;
this.retryable = retryable;
this.systemCode = systemErrorCode(errorOptions.cause);
}
};
function systemErrorCode(error) {
for (let depth = 0; error instanceof Error && depth < 3; depth++, error = error.cause) {
if (SYSTEM_CODES.has(error.code)) return error.code;
}
return "unknown";
}
function normalizeRunRequest(input) {
if (!isRecord(input) || input.type !== "run") {
throw new SafeWorkerError("invalid_request", 'Expected a JSON object with type "run".');
}
const id = resolveAlias(input, "id", "requestId");
if (typeof id !== "string" || !SAFE_ID.test(id)) {
throw new SafeWorkerError(
"invalid_request",
"id or requestId must contain only letters, numbers, dot, underscore, colon, or dash."
);
}
if (typeof input.model !== "string" || !SAFE_MODEL.test(input.model)) {
throw new SafeWorkerError("invalid_request", "model must be a non-empty model identifier.");
}
if (typeof input.prompt !== "string" || input.prompt.trim().length === 0 || input.prompt.length > MAX_PROMPT_CHARS) {
throw new SafeWorkerError(
"invalid_request",
`prompt must contain between 1 and ${MAX_PROMPT_CHARS} characters.`
);
}
if (typeof input.workingDirectory !== "string" || input.workingDirectory.length === 0 || input.workingDirectory.length > 4096 || !path3.isAbsolute(input.workingDirectory)) {
throw new SafeWorkerError("invalid_request", "workingDirectory must be an absolute path.");
}
if (!SANDBOX_MODES.has(input.sandboxMode)) {
throw new SafeWorkerError(
"invalid_request",
'sandboxMode must be "read-only" or "workspace-write".'
);
}
const networkAccessEnabled = resolveAlias(
input,
"networkAccessEnabled",
"networkAccess"
);
if (typeof networkAccessEnabled !== "boolean") {
throw new SafeWorkerError(
"invalid_request",
"networkAccessEnabled or networkAccess must be a boolean."
);
}
if (input.reasoningEffort !== void 0 && !REASONING_EFFORTS.has(input.reasoningEffort)) {
throw new SafeWorkerError("invalid_request", "reasoningEffort is not supported.");
}
if (input.outputSchema !== void 0) {
if (!isRecord(input.outputSchema)) {
throw new SafeWorkerError("invalid_request", "outputSchema must be a JSON object.");
}
let serializedSchema;
try {
serializedSchema = JSON.stringify(input.outputSchema);
} catch {
throw new SafeWorkerError("invalid_request", "outputSchema must be JSON serializable.");
}
if (serializedSchema.length > MAX_SCHEMA_CHARS) {
throw new SafeWorkerError(
"invalid_request",
`outputSchema must not exceed ${MAX_SCHEMA_CHARS} serialized characters.`
);
}
}
return {
type: "run",
id,
model: input.model,
prompt: input.prompt,
workingDirectory: input.workingDirectory,
sandboxMode: input.sandboxMode,
networkAccessEnabled,
...input.reasoningEffort === void 0 ? {} : { reasoningEffort: input.reasoningEffort },
...input.outputSchema === void 0 ? {} : { outputSchema: input.outputSchema }
};
}
async function executeRunRequest(request, {
abortController = new AbortController(),
codexFactory = defaultCodexFactory,
emit = () => {
},
diagnostic = () => {
}
} = {}) {
const lifecycle = createLifecycleEmitter(request.id, emit);
let threadId = null;
let finalResponse = "";
let finalResponseTruncated = false;
let usage = null;
let turnCompleted = false;
let streamWarnings = 0;
const itemCounts = {};
const startedAt = performance.now();
let stage = "launch";
try {
const codex = codexFactory();
const thread = codex.startThread({
model: request.model,
sandboxMode: request.sandboxMode,
approvalPolicy: "never",
networkAccessEnabled: request.networkAccessEnabled,
webSearchMode: "disabled",
skipGitRepoCheck: true,
workingDirectory: request.workingDirectory,
...request.reasoningEffort === void 0 ? {} : { modelReasoningEffort: request.reasoningEffort }
});
stage = "turn_start";
const { events } = await thread.runStreamed(request.prompt, {
signal: abortController.signal,
...request.outputSchema === void 0 ? {} : { outputSchema: request.outputSchema }
});
if (!events || typeof events[Symbol.asyncIterator] !== "function") {
throw new SafeWorkerError(
"invalid_sdk_response",
"Codex SDK did not return an event stream."
);
}
stage = "stream";
for await (const event of events) {
if (abortController.signal.aborted) {
throw abortController.signal.reason;
}
if (!isRecord(event) || typeof event.type !== "string") continue;
if (event.type === "thread.started") {
threadId = safeThreadId(event.thread_id) ?? safeThreadId(thread.id);
lifecycle({ phase: "thread_started", ...threadId ? { threadId } : {} });
} else if (event.type === "turn.started") {
lifecycle({ phase: "turn_started" });
} else if (event.type === "item.completed" && isRecord(event.item)) {
const itemType = safeItemType(event.item.type);
itemCounts[itemType] = (itemCounts[itemType] ?? 0) + 1;
if (event.item.type === "agent_message" && typeof event.item.text === "string") {
const truncated = truncate(event.item.text, MAX_FINAL_RESPONSE_CHARS);
finalResponse = truncated.value;
finalResponseTruncated = truncated.truncated;
}
lifecycle({
phase: "item_completed",
itemType,
itemCount: itemCounts[itemType]
});
} else if (event.type === "turn.completed") {
usage = normalizeUsage(event.usage);
turnCompleted = true;
} else if (event.type === "turn.failed") {
const message = safeProviderErrorMessage(event.error?.message, "Codex turn failed.");
diagnostic(`Codex turn failed: ${message}`);
throw new SafeWorkerError("turn_failed", message);
} else if (event.type === "error") {
const message = safeProviderErrorMessage(
event.message,
"Codex reported an unrecoverable stream error."
);
diagnostic(`Codex stream error: ${message}`);
throw new SafeWorkerError(
"stream_error",
message,
{ retryable: isRetryableStreamError(event) }
);
}
}
if (!turnCompleted) {
throw new SafeWorkerError(
"incomplete_stream",
"Codex event stream ended before turn completion.",
{ retryable: true }
);
}
return {
threadId: threadId ?? safeThreadId(thread.id),
finalResponse,
finalResponseTruncated,
usage,
itemCounts,
streamWarnings
};
} catch (error) {
if (abortController.signal.aborted) {
throw new SafeWorkerError("canceled", "Codex run was canceled.", { cause: error });
}
const message = error instanceof SafeWorkerError ? error.message : describeWorkerError(error, diagnostic);
const failure = error instanceof SafeWorkerError ? error : looksLikeMissingCodex(error) ? new SafeWorkerError(
"codex_unavailable",
"The local Codex executable could not be started.",
{ cause: error }
) : new SafeWorkerError("worker_failed", message, {
cause: error,
retryable: stage === "stream" && isRetryableStreamFailure(error)
});
failure.stage = stage;
failure.elapsedMs = performance.now() - startedAt;
throw failure;
}
}
function startStdioWorker({
input = process.stdin,
output = process.stdout,
errorOutput = process.stderr,
setExitCode = (code) => {
process.exitCode = code;
},
codexFactory = defaultCodexFactory
} = {}) {
const readline2 = createInterface({ input, crlfDelay: Infinity });
let active = null;
let finished = false;
let runStarted = false;
const emit = (message) => {
output.write(`${JSON.stringify(message)}
`);
};
const finish = (exitCode) => {
if (finished) return;
finished = true;
setExitCode(exitCode);
readline2.close();
input.pause?.();
};
const failProtocol = (error, id = null) => {
const safeError = error instanceof SafeWorkerError ? error : new SafeWorkerError("invalid_request", "Invalid worker request.");
emit({
type: "failed",
id,
code: safeError.code,
message: safeError.message,
retryable: safeError.retryable,
systemCode: safeError.systemCode,
stage: "validation"
});
finish(2);
};
const startRun = (inputRequest) => {
let request;
try {
request = normalizeRunRequest(inputRequest);
} catch (error) {
failProtocol(error, safeRequestId(inputRequest));
return;
}
runStarted = true;
const abortController = new AbortController();
active = { id: request.id, abortController };
emit({ type: "accepted", id: request.id });
void executeRunRequest(request, {
abortController,
codexFactory,
emit,
diagnostic: (message) => errorOutput.write(`[codex-worker] ${message}
`)
}).then((result) => {
emit({
type: "completed",
id: request.id,
threadId: result.threadId,
finalResponse: result.finalResponse,
finalResponseTruncated: result.finalResponseTruncated,
usage: result.usage,
itemCounts: result.itemCounts,
streamWarnings: result.streamWarnings
});
finish(0);
}).catch((error) => {
const safeError = error instanceof SafeWorkerError ? error : new SafeWorkerError("worker_failed", "Codex worker failed.");
if (safeError.code === "canceled") {
const reason = active?.abortController.signal.reason;
emit({
type: "canceled",
id: request.id,
reason: isRecord(reason) && typeof reason.signal === "string" ? reason.signal : "requested"
});
finish(
isRecord(reason) && reason.signal === "SIGTERM" ? 143 : isRecord(reason) && reason.signal === "SIGINT" ? 130 : 0
);
return;
}
emit({
type: "failed",
id: request.id,
code: safeError.code,
message: safeError.message,
retryable: safeError.retryable,
systemCode: safeError.systemCode,
stage: safeError.stage ?? "validation",
elapsedMs: safeError.elapsedMs ?? null
});
finish(1);
});
};
readline2.on("line", (line) => {
if (finished || line.trim().length === 0) return;
let message;
try {
message = JSON.parse(line);
} catch {
if (!runStarted) {
failProtocol(new SafeWorkerError("invalid_json", "Input must be valid JSON."));
} else {
emit({
type: "protocol_error",
id: active?.id ?? null,
code: "invalid_json",
message: "Input must be valid JSON."
});
}
return;
}
if (isRecord(message) && message.type === "cancel") {
if (!active || message.id !== active.id) {
emit({
type: "protocol_error",
id: safeRequestId(message),
code: "unknown_run",
message: "No matching run is active."
});
return;
}
active.abortController.abort({ kind: "cancel" });
return;
}
if (runStarted) {
emit({
type: "protocol_error",
id: safeRequestId(message),
code: "run_already_started",
message: "This worker accepts exactly one run request."
});
return;
}
startRun(message);
});
emit({ type: "ready", protocolVersion: PROTOCOL_VERSION });
return {
cancel(signal = "SIGTERM") {
if (!active || active.abortController.signal.aborted) return false;
active.abortController.abort({ kind: "signal", signal });
return true;
},
close() {
if (!finished) finish(0);
}
};
}
function defaultCodexFactory() {
const target = resolveCodexTarget();
const env = Object.fromEntries(
Object.entries(process.env).filter((entry) => entry[1] !== void 0)
);
env[CLI_WRAPPER_MODE_ENV] = "1";
env[CLI_WRAPPER_TARGET_ENV] = target;
env[CLI_WRAPPER_OWNER_ENV] = String(process.pid);
return new Codex({
codexPathOverride: fileURLToPath(import.meta.url),
env
});
}
function resolveCodexTarget({
env = process.env,
platform = process.platform,
applicationRoots = defaultApplicationRoots(env, platform),
probe = probeCodexCandidate
} = {}) {
const candidates = [];
const addCandidate = (candidate) => {
if (candidate && !candidates.includes(candidate)) candidates.push(candidate);
};
const configured = env.CODEX_CLI_PATH?.trim();
addCandidate(configured);
const executableName = platform === "win32" ? "codex.exe" : "codex";
for (const pathMatch of findExecutablesOnPath(executableName, env, platform)) {
addCandidate(pathMatch);
}
if (platform === "darwin") {
for (const applicationRoot of applicationRoots) {
for (const applicationName of codexApplicationNames(applicationRoot)) {
const candidate = path3.join(
applicationRoot,
applicationName,
"Contents",
"Resources",
"codex"
);
if (isExecutableFile(candidate)) addCandidate(candidate);
}
}
}
for (const candidate of candidates) {
try {
if (probe(candidate, env)) return candidate;
} catch {
}
}
const checked = candidates.length > 0 ? ` Checked: ${candidates.join(", ")}.` : "";
throw new SafeWorkerError(
"codex_unavailable",
`No working Codex executable was found.${checked}`
);
}
function defaultApplicationRoots(env, platform) {
if (platform !== "darwin") return [];
const home = env.HOME?.trim() || homedir();
return ["/Applications", path3.join(home, "Applications")];
}
function findExecutablesOnPath(executableName, env, platform) {
const pathValue = platform === "win32" ? Object.entries(env).find(([key]) => key.toLowerCase() === "path")?.[1] : env.PATH;
if (!pathValue) return [];
const matches = [];
for (const directory of pathValue.split(path3.delimiter)) {
if (!directory) continue;
const candidate = path3.join(directory, executableName);
if (isExecutableFile(candidate)) matches.push(candidate);
}
return matches;
}
function probeCodexCandidate(candidate, env) {
const completed = spawnSync(candidate, ["--version"], {
env,
encoding: "utf8",
stdio: ["ignore", "pipe", "pipe"],
timeout: CODEX_VERSION_PROBE_TIMEOUT_MS,
maxBuffer: CODEX_VERSION_PROBE_MAX_BUFFER,
windowsHide: true
});
return completed.status === 0 && completed.error === void 0;
}
function codexApplicationNames(applicationRoot) {
try {
return readdirSync(applicationRoot, { withFileTypes: true }).filter(
(entry) => entry.isDirectory() && /^(?:Codex|ChatGPT).*\.app$/.test(entry.name)
).map((entry) => entry.name).sort((left, right) => applicationPriority(left) - applicationPriority(right));
} catch {
return [];
}
}
function applicationPriority(name) {
if (name === "Codex.app") return 0;
if (name === "ChatGPT.app") return 1;
if (name.startsWith("Codex")) return 2;
return 3;
}
function isExecutableFile(candidate) {
try {
if (!statSync2(candidate).isFile()) return false;
accessSync(candidate, constants.X_OK);
return true;
} catch {
return false;
}
}
function startCodexCliWrapper() {
const target = process.env[CLI_WRAPPER_TARGET_ENV]?.trim();
if (!target) {
process.stderr.write("Missing isolated Codex CLI target.\n");
process.exitCode = 127;
return;
}
const originalArgs = process.argv.slice(2);
if (originalArgs[0] !== "exec") {
process.stderr.write("The isolated Codex wrapper only supports exec.\n");
process.exitCode = 2;
return;
}
const args = [
"exec",
"--ignore-rules",
...originalArgs.slice(1)
];
const detached = process.platform !== "win32";
const child = spawn2(target, args, {
detached,
env: process.env,
stdio: "inherit"
});
let killTimer = null;
const killTree = (signal) => {
if (child.exitCode !== null || child.signalCode !== null) return;
try {
if (detached && Number.isInteger(child.pid)) {
process.kill(-child.pid, signal);
} else {
child.kill(signal);
}
} catch {
}
};
const forwardSignal = (signal) => {
killTree(signal);
if (killTimer === null) {
killTimer = setTimeout(() => killTree("SIGKILL"), TREE_KILL_GRACE_MS);
killTimer.unref?.();
}
};
const signalHandlers = new Map(
["SIGTERM", "SIGINT", "SIGHUP"].map((signal) => {
const handler = () => forwardSignal(signal);
process.on(signal, handler);
return [signal, handler];
})
);
const ownerPid = Number.parseInt(process.env[CLI_WRAPPER_OWNER_ENV] ?? "", 10);
const ownerMonitor = Number.isInteger(ownerPid) && ownerPid > 1 ? setInterval(() => {
try {
process.kill(ownerPid, 0);
} catch {
forwardSignal("SIGTERM");
}
}, 250) : null;
ownerMonitor?.unref?.();
const cleanup = () => {
if (killTimer !== null) clearTimeout(killTimer);
if (ownerMonitor !== null) clearInterval(ownerMonitor);
for (const [signal, handler] of signalHandlers) {
process.off(signal, handler);
}
};
child.once("error", (error) => {
cleanup();
process.stderr.write(`${error.message}
`);
process.exitCode = 127;
});
child.once("exit", (code, signal) => {
cleanup();
process.exitCode = code ?? signalExitCode(signal);
});
}
function signalExitCode(signal) {
return {
SIGHUP: 129,
SIGINT: 130,
SIGTERM: 143,
SIGKILL: 137
}[signal] ?? 1;
}
function createLifecycleEmitter(id, emit) {
let emitted = 0;
return (event) => {
if (emitted >= MAX_LIFECYCLE_EVENTS) return;
emitted += 1;
emit({ type: "lifecycle", id, sequence: emitted, ...event });
};
}
function normalizeUsage(value) {
if (!isRecord(value)) return null;
const usage = {};
for (const key of [
"input_tokens",
"cached_input_tokens",
"output_tokens",
"reasoning_output_tokens"
]) {
if (Number.isSafeInteger(value[key]) && value[key] >= 0) {
usage[key] = value[key];
}
}
return Object.keys(usage).length > 0 ? usage : null;
}
function resolveAlias(input, canonical, alias) {
if (input[canonical] !== void 0 && input[alias] !== void 0 && input[canonical] !== input[alias]) {
throw new SafeWorkerError(
"invalid_request",
`${canonical} and ${alias} must match when both are provided.`
);
}
return input[canonical] ?? input[alias];
}
function safeItemType(value) {
return typeof value === "string" && SAFE_ITEM_TYPES.has(value) ? value : "other";
}
function safeThreadId(value) {
return typeof value === "string" && /^[A-Za-z0-9_-]{1,256}$/.test(value) ? value : null;
}
function safeRequestId(value) {
if (!isRecord(value)) return null;
const id = value.id ?? value.requestId;
return typeof id === "string" && SAFE_ID.test(id) ? id : null;
}
function describeWorkerError(error, diagnostic) {
const messages = [];
for (let depth = 0; error instanceof Error && depth < 3; depth++, error = error.cause) {
if (error.message.trim()) messages.push(error.message.trim());
}
const fallback = "Codex worker failed.";
if (messages.length === 0) return fallback;
diagnostic(`Codex worker exception: ${safeProviderErrorMessage(
messages.join("\nCaused by: "),
fallback,
MAX_WORKER_DIAGNOSTIC_CHARS
)}`);
return safeProviderErrorMessage(messages.at(-1).split(/\r?\n/).at(-1), fallback);
}
function safeProviderErrorMessage(value, fallback, maxChars = MAX_PROVIDER_ERROR_CHARS) {
if (typeof value !== "string" || !value.trim()) return fallback;
let message = value;
const parsed = parseProviderErrorPayload(value);
if (isRecord(parsed?.error) && typeof parsed.error.message === "string") {
message = parsed.error.message;
} else if (isRecord(parsed) && typeof parsed.message === "string") {
message = parsed.message;
}
return redactKnownCredentials(message).replace(/[\u0000-\u001f\u007f]/g, " ").trim().slice(0, maxChars) || fallback;
}
function redactKnownCredentials(message) {
return message.replace(/\bBearer\s+[A-Za-z0-9._~+/-]+=*/gi, "Bearer [REDACTED]").replace(/\bsk-(?:proj-|svcacct-)?[A-Za-z0-9_-]{12,}/g, "[REDACTED]").replace(/\b(api[_-]?key|access[_-]?token|password|secret)\s*[:=]\s*[^\s,;]+/gi, "$1=[REDACTED]");
}
function parseProviderErrorPayload(value) {
if (typeof value !== "string" || !value.trim()) return null;
try {
const parsed = JSON.parse(value);
return isRecord(parsed) ? parsed : null;
} catch {
return null;
}
}
function providerErrorObjects(event) {
if (!isRecord(event)) return [];
const objects = [event];
const parsed = parseProviderErrorPayload(event.message);
if (parsed !== null) objects.push(parsed);
for (const value of [...objects]) {
if (isRecord(value.error)) objects.push(value.error);
}
return objects;
}
function providerErrorStatus(value) {
const status = value.status ?? value.status_code ?? value.statusCode;
if (Number.isInteger(status)) return status;
if (typeof status === "string" && /^[0-9]{3}$/.test(status)) {
return Number(status);
}
return null;
}
function providerErrorKinds(objects) {
const kinds = [];
for (const value of objects) {
for (const field of ["type", "code", "error_code"]) {
if (typeof value[field] === "string") kinds.push(value[field]);
}
}
return kinds;
}
function isRetryableStreamError(event) {
const objects = providerErrorObjects(event);
const kinds = providerErrorKinds(objects).map((kind) => kind.toLowerCase());
if (kinds.some((kind) => PERMANENT_STREAM_ERROR_KINDS.has(kind))) {
return false;
}
if (kinds.some((kind) => RETRYABLE_STREAM_ERROR_KINDS.has(kind))) {
return true;
}
for (const value of objects) {
const status = providerErrorStatus(value);
if (status === 408 || status === 429 || status !== null && status >= 500) {
return true;
}
if (status !== null && status >= 400 && status < 500) {
return false;
}
}
const message = typeof event?.message === "string" ? event.message : "";
if (/auth|unauthoriz|forbidden|api key|quota|billing|usage limit|model.*(?:access|not found)|invalid request/i.test(
message
)) {
return false;
}
return true;
}
function isRetryableStreamFailure(error) {
const code = systemErrorCode(error);
if (PERMANENT_SYSTEM_CODES.has(code)) return false;
if (code !== "unknown") return true;
if (looksLikePermanentLocalFailure(error)) return false;
return isRetryableStreamError(error);
}
function looksLikePermanentLocalFailure(error) {
const message = error instanceof Error ? error.message : String(error);
return /(?:operation not permitted|permission denied|failed to initialize in-process app-server client)/i.test(
message
);
}
function truncate(value, maxChars) {
if (value.length <= maxChars) return { value, truncated: false };
return { value: value.slice(0, maxChars), truncated: true };
}
function looksLikeMissingCodex(error) {
const message = error instanceof Error ? error.message : String(error);
return /(?:ENOENT|Unable to locate Codex CLI binaries|spawn .* not found)/i.test(message);
}
function isRecord(value) {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
function isMainModule() {
const scriptPath = process.argv[1];
return typeof scriptPath === "string" && pathToFileURL(path3.resolve(scriptPath)).href === import.meta.url;
}
if (isMainModule()) {
if (process.env[CLI_WRAPPER_MODE_ENV] === "1") {
startCodexCliWrapper();
} else {
const worker = startStdioWorker();
process.once("SIGTERM", () => {
if (worker.cancel("SIGTERM")) {
setTimeout(() => process.exit(143), 5e3).unref();
} else {
process.exitCode = 143;
worker.close();
}
});
process.once("SIGINT", () => {
if (worker.cancel("SIGINT")) {
setTimeout(() => process.exit(130), 5e3).unref();
} else {
process.exitCode = 130;
worker.close();
}
});
}
}
export {
executeRunRequest,
normalizeRunRequest,
resolveCodexTarget,
startStdioWorker
};
SHA-256: b147ece3ee39bf73cbcbd0f9ed3f6782937e7d90b043734f8016180d2aac9740