← Files Codex ReplayARCHIVED FILE

mcp/codex-worker.mjs

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

↓ Download file

#!/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