| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932 |
- import { EventEmitter } from 'node:events';
- import { spawn, type ChildProcessWithoutNullStreams } from 'node:child_process';
- import { locateCodexBinary, readCodexVersion } from './codexLocator';
- import { JsonRpcPeer, type JsonRpcServerRequest } from './jsonRpcPeer';
- /**
- * 移植自 Noobi.ai src/main/codexAppServer.ts,主要改动:
- * 1. 删掉全部账号相关能力(account/read、account/login/start、account/logout)——本模块不做任何登录;
- * 2. 不再传 --strict-config(用户与程序共写 config.toml,未知字段只应告警不该退出);
- * 3. 新增 provider 注入:用 -c 覆盖 model_provider / model / model_providers.*,api_key 只走子进程环境变量;
- * 4. clientInfo 改为本项目标识。
- */
- export type JsonValue = null | boolean | number | string | JsonValue[] | { [key: string]: JsonValue };
- export interface ModelOption {
- id: string;
- model: string;
- displayName: string;
- description: string;
- isDefault: boolean;
- defaultEffort: string;
- efforts: string[];
- }
- export interface RuntimeCapabilities {
- namespaceTools: boolean;
- imageGeneration: boolean;
- webSearch: boolean;
- }
- export interface RuntimeStatus {
- state: 'stopped' | 'starting' | 'ready' | 'error';
- binaryPath: string | null;
- version: string | null;
- codexHome: string | null;
- models: ModelOption[];
- capabilities: RuntimeCapabilities;
- providerId: string | null;
- defaultModel: string | null;
- /** model/list 拿不到可用模型时为 true */
- degraded: boolean;
- error: string | null;
- }
- export interface CodexMcpServerStatus {
- name: string;
- authStatus: string;
- connected: boolean;
- toolCount: number;
- tools: Array<{ name: string; description: string }>;
- }
- export interface CodexSkillSummary {
- name: string;
- description: string;
- path: string;
- scope: 'user' | 'repo' | 'system' | 'admin' | string;
- enabled: boolean;
- cwd: string;
- }
- /** 一次 spawn 用的模型来源:一律是自定义 provider,Codex 的内置 id 是保留字、不可覆盖 */
- export interface CodexProviderSpec {
- id: string;
- name?: string | null;
- baseUrl: string;
- model: string;
- /** apiKey 只注入子进程环境变量,不落任何盘 */
- apiKey?: string | null;
- envKey?: string;
- httpHeaders?: Record<string, string> | null;
- envHttpHeaders?: Record<string, string> | null;
- /** 仅放行 PROVIDER_EXTRA_ALLOWLIST 里的键,其余在 providerService 就被丢弃 */
- extra?: Record<string, JsonValue> | null;
- }
- export interface CodexRuntimeOptions {
- codexHome: string;
- provider?: CodexProviderSpec | null;
- clientInfo?: { name: string; title: string; version: string };
- }
- export interface StartThreadOptions {
- cwd: string;
- model?: string | null;
- sandbox?: 'read-only' | 'workspace-write' | 'danger-full-access';
- approvalPolicy?: 'untrusted' | 'on-request' | 'never';
- developerInstructions?: string;
- ephemeral?: boolean;
- }
- export interface StartTurnOptions {
- threadId: string;
- prompt: string;
- cwd?: string;
- model?: string | null;
- effort?: string | null;
- approvalPolicy?: 'untrusted' | 'on-request' | 'never';
- skills?: Array<{ name: string; path: string }>;
- /** 仅宿主侧的完成期限,不会发给 App Server */
- timeoutMs?: number;
- }
- export interface TurnResult {
- turnId: string;
- status: string;
- text: string;
- raw: unknown;
- }
- interface SkillsListResponse {
- data: Array<{
- cwd: string;
- skills: Array<{
- name: string;
- description: string;
- path: string;
- scope: string;
- enabled: boolean;
- }>;
- }>;
- }
- interface ListMcpServerStatusResponse {
- data: Array<{
- name: string;
- serverInfo: unknown | null;
- tools: Record<string, unknown>;
- authStatus: string;
- }>;
- nextCursor: string | null;
- }
- interface InitializeResponse {
- userAgent: string;
- codexHome: string;
- platformFamily: string;
- platformOs: string;
- }
- interface ModelListResponse {
- data: Array<{
- id: string;
- model: string;
- displayName: string;
- description: string;
- isDefault: boolean;
- defaultReasoningEffort: string;
- supportedReasoningEfforts: Array<{ reasoningEffort: string }>;
- }>;
- nextCursor: string | null;
- }
- interface ModelProviderCapabilitiesResponse {
- namespaceTools: boolean;
- imageGeneration: boolean;
- webSearch: boolean;
- }
- interface ThreadResponse {
- thread: { id: string };
- model: string;
- }
- interface TurnStartResponse {
- turn: { id: string; status: string };
- }
- interface EarlyTurnState {
- text: string;
- completed: TurnResult | null;
- }
- interface TurnWaiter {
- text: string;
- resolve(result: TurnResult): void;
- reject(error: Error): void;
- timer: NodeJS.Timeout;
- }
- const TURN_TIMEOUT_MS = 20 * 60 * 1_000;
- /** 单条 RPC 超过这个耗时才值得进诊断(jsonRpcPeer 的预算是 30 秒) */
- const SLOW_RPC_MS = 3_000;
- /** 退出事件可能比 stderr 落地慢半拍,等一下再拼错误信息 */
- const EXIT_GRACE_MS = 300;
- const STDERR_TAIL_LINES = 40;
- export const DEFAULT_API_KEY_ENV = 'ZSJZ_CODEX_API_KEY';
- /**
- * 模型管理 config 里允许透传给 Codex 的 provider 字段(对应 rust-v0.155.1
- * model-provider-info/src/lib.rs 的 ModelProviderInfo)。providerService 用它做白名单,
- * 这里做 -c 落参,两处必须同源,否则页面提示"已丢弃"而实际写进去、或反之。
- */
- export const PROVIDER_EXTRA_ALLOWLIST: readonly string[] = Object.freeze([
- 'request_max_retries',
- 'stream_max_retries',
- 'stream_idle_timeout_ms',
- 'websocket_connect_timeout_ms',
- 'supports_websockets',
- 'query_params',
- ]);
- const PROVIDER_EXTRA_ALLOWED = new Set<string>(PROVIDER_EXTRA_ALLOWLIST);
- export class CodexRuntime extends EventEmitter {
- readonly #codexHome: string;
- readonly #clientInfo: { name: string; title: string; version: string };
- #provider: CodexProviderSpec | null;
- #child: ChildProcessWithoutNullStreams | null = null;
- #peer: JsonRpcPeer | null = null;
- #startPromise: Promise<RuntimeStatus> | null = null;
- #runtime: RuntimeStatus = emptyRuntimeStatus();
- #turnWaiters = new Map<string, TurnWaiter>();
- #earlyTurnStates = new Map<string, EarlyTurnState>();
- #runTurnStartsInFlight = 0;
- #generation = 0;
- constructor(options: CodexRuntimeOptions) {
- super();
- this.#codexHome = options.codexHome;
- this.#provider = options.provider ?? null;
- this.#clientInfo = options.clientInfo ?? {
- name: 'zsjz_ai',
- title: '清鉴线索调查工具',
- version: '0.1.0',
- };
- }
- get status(): RuntimeStatus {
- return structuredClone(this.#runtime);
- }
- get provider(): CodexProviderSpec | null {
- return this.#provider ? { ...this.#provider, apiKey: null } : null;
- }
- /** 换 provider 必须重启子进程:api_key 走的是子进程环境变量,无法热更新 */
- async applyProvider(provider: CodexProviderSpec | null): Promise<RuntimeStatus> {
- const previous = this.#provider;
- this.#provider = provider;
- if (this.#child) await this.stop();
- try {
- return await this.start();
- } catch (error) {
- // 失败要回滚:坏 spec 留在内存里,之后每次 start 都会重放同一组参数,
- // 一条选错的模型能把运行时锁死到重启应用为止。
- this.#provider = previous;
- throw error;
- }
- }
- async start(): Promise<RuntimeStatus> {
- if (this.#runtime.state === 'ready' && this.#peer) return this.status;
- if (this.#startPromise) return this.#startPromise;
- const generation = ++this.#generation;
- this.#startPromise = this.#start(generation);
- try {
- return await this.#startPromise;
- } finally {
- this.#startPromise = null;
- }
- }
- async refresh(): Promise<RuntimeStatus> {
- await this.start();
- let models: ModelOption[] = [];
- let degraded = false;
- try {
- models = await this.listModels();
- } catch (error) {
- // 第三方/本地端点常常不实现 model/list,这里降级而不是失败
- degraded = true;
- this.emit('diagnostic', `model/list 不可用:${asError(error).message}`);
- }
- const capabilities = await this.readModelProviderCapabilities().catch(() => emptyCapabilities());
- this.#runtime = {
- ...this.#runtime,
- models,
- capabilities,
- degraded: degraded || models.length === 0,
- // model/list 返回的是 Codex 内置模型目录(与自定义 provider 无关),
- // 实际会用的是 provider 里配置的模型,所以它优先
- defaultModel: this.#provider?.model ?? models.find((item) => item.isDefault)?.model ?? null,
- state: 'ready',
- error: null,
- };
- this.emit('status', this.status);
- return this.status;
- }
- async readModelProviderCapabilities(): Promise<RuntimeCapabilities> {
- const result = await this.#request<ModelProviderCapabilitiesResponse>(
- 'modelProvider/capabilities/read',
- {},
- );
- return {
- namespaceTools: result.namespaceTools === true,
- imageGeneration: result.imageGeneration === true,
- webSearch: result.webSearch === true,
- };
- }
- async listModels(): Promise<ModelOption[]> {
- const models: ModelOption[] = [];
- let cursor: string | null = null;
- for (let page = 0; page < 10; page += 1) {
- const result: ModelListResponse = await this.#request<ModelListResponse>('model/list', {
- cursor,
- limit: 100,
- includeHidden: false,
- });
- models.push(
- ...result.data.map((model) => ({
- id: model.id,
- model: model.model,
- displayName: model.displayName,
- description: model.description,
- isDefault: model.isDefault,
- defaultEffort: model.defaultReasoningEffort,
- efforts: model.supportedReasoningEfforts.map((item) => item.reasoningEffort),
- })),
- );
- cursor = result.nextCursor;
- if (!cursor) break;
- }
- return models;
- }
- async listSkills(options: { cwds?: string[]; forceReload?: boolean } = {}): Promise<CodexSkillSummary[]> {
- await this.start();
- const result = await this.#request<SkillsListResponse>('skills/list', {
- ...(options.cwds?.length ? { cwds: options.cwds } : {}),
- forceReload: options.forceReload ?? false,
- });
- return result.data.flatMap((entry) =>
- entry.skills.map((skill) => ({
- name: skill.name,
- description: skill.description,
- path: skill.path,
- scope: skill.scope,
- enabled: skill.enabled,
- cwd: entry.cwd,
- })),
- );
- }
- async setSkillEnabled(selector: { path?: string; name?: string }, enabled: boolean): Promise<boolean> {
- await this.start();
- const result = await this.#request<{ effectiveEnabled: boolean }>('skills/config/write', {
- ...(selector.path ? { path: selector.path } : {}),
- ...(!selector.path && selector.name ? { name: selector.name } : {}),
- enabled,
- });
- return result.effectiveEnabled;
- }
- async listMcpServerStatuses(): Promise<CodexMcpServerStatus[]> {
- await this.start();
- const statuses: CodexMcpServerStatus[] = [];
- let cursor: string | null = null;
- for (let page = 0; page < 20; page += 1) {
- const result: ListMcpServerStatusResponse = await this.#request<ListMcpServerStatusResponse>(
- 'mcpServerStatus/list',
- {
- cursor,
- limit: 100,
- detail: 'toolsAndAuthOnly',
- },
- );
- statuses.push(
- ...result.data.map((server) => ({
- name: server.name,
- authStatus: server.authStatus,
- connected: server.serverInfo !== null,
- toolCount: Object.keys(server.tools ?? {}).length,
- tools: toToolSummaries(server.tools),
- })),
- );
- cursor = result.nextCursor;
- if (!cursor) break;
- }
- return statuses;
- }
- async readConfig(): Promise<Record<string, unknown>> {
- await this.start();
- const result = await this.#request<{ config: Record<string, unknown> }>('config/read', {
- includeLayers: false,
- });
- return structuredClone(result.config);
- }
- /** 局部写:Codex 自己也持有 config.toml 写权,禁止整文件覆盖 */
- async writeConfigValue(keyPath: string, value: JsonValue): Promise<void> {
- await this.start();
- await this.#request<unknown>('config/value/write', { keyPath, value, mergeStrategy: 'replace' });
- }
- async reloadMcpServers(): Promise<void> {
- await this.start();
- await this.#request<unknown>('config/mcpServer/reload');
- }
- async startThread(options: StartThreadOptions): Promise<string> {
- await this.start();
- const result = await this.#request<ThreadResponse>('thread/start', {
- cwd: options.cwd,
- ...(options.model ? { model: options.model } : {}),
- sandbox: options.sandbox ?? 'workspace-write',
- approvalPolicy: options.approvalPolicy ?? 'never',
- ...(options.developerInstructions ? { developerInstructions: options.developerInstructions } : {}),
- ...(options.ephemeral === undefined ? {} : { ephemeral: options.ephemeral }),
- experimentalRawEvents: false,
- });
- return result.thread.id;
- }
- async resumeThread(threadId: string, options: StartThreadOptions): Promise<string> {
- await this.start();
- const result = await this.#request<ThreadResponse>('thread/resume', {
- threadId,
- cwd: options.cwd,
- ...(options.model ? { model: options.model } : {}),
- sandbox: options.sandbox ?? 'workspace-write',
- approvalPolicy: options.approvalPolicy ?? 'never',
- ...(options.developerInstructions ? { developerInstructions: options.developerInstructions } : {}),
- // excludeTurns 不影响上游收到的历史(实测 true/false 两轮上下文都完整带回)
- excludeTurns: true,
- });
- return result.thread.id;
- }
- async startTurn(options: StartTurnOptions): Promise<string> {
- await this.start();
- const result = await this.#request<TurnStartResponse>('turn/start', {
- threadId: options.threadId,
- input: [
- { type: 'text', text: options.prompt, text_elements: [] },
- ...(options.skills ?? []).map((skill) => ({
- type: 'skill' as const,
- name: skill.name,
- path: skill.path,
- })),
- ],
- ...(options.cwd ? { cwd: options.cwd } : {}),
- ...(options.model ? { model: options.model } : {}),
- ...(options.effort ? { effort: options.effort } : {}),
- ...(options.approvalPolicy ? { approvalPolicy: options.approvalPolicy } : {}),
- });
- return result.turn.id;
- }
- async runTurn(options: StartTurnOptions): Promise<TurnResult> {
- this.#runTurnStartsInFlight += 1;
- let turnId: string;
- try {
- turnId = await this.startTurn(options);
- } catch (error) {
- this.#releaseRunTurnStart();
- throw error;
- }
- // turn/completed 可能早于 turn/start 的响应回来,先存起来
- const startedAt = Date.now();
- const noteTurn = (result: TurnResult): TurnResult => {
- this.emit(
- 'diagnostic',
- `layer=turn turn=${result.turnId} status=${result.status} 用时 ${((Date.now() - startedAt) / 1000).toFixed(1)}s`,
- );
- return result;
- };
- const early = this.#earlyTurnStates.get(turnId);
- this.#earlyTurnStates.delete(turnId);
- if (early?.completed) {
- this.#releaseRunTurnStart();
- return noteTurn(early.completed);
- }
- return new Promise<TurnResult>((resolve, reject) => {
- const timeoutMs =
- options.timeoutMs === undefined
- ? TURN_TIMEOUT_MS
- : Math.min(TURN_TIMEOUT_MS, Math.max(1_000, Math.trunc(options.timeoutMs)));
- const timer = setTimeout(() => {
- this.#turnWaiters.delete(turnId);
- void this.interruptTurn(options.threadId, turnId).catch(() => undefined);
- this.emit(
- 'diagnostic',
- `layer=turn turn=${turnId} 宿主侧期限已到(${Math.round(timeoutMs / 1000)} 秒),已发 turn/interrupt`,
- );
- reject(new Error(`Codex turn ${turnId} 超时`));
- }, timeoutMs);
- timer.unref();
- this.#turnWaiters.set(turnId, {
- text: early?.text ?? '',
- resolve: (result) => resolve(noteTurn(result)),
- reject,
- timer,
- });
- this.#releaseRunTurnStart();
- });
- }
- async interruptTurn(threadId: string, turnId: string): Promise<void> {
- await this.#request<Record<string, never>>('turn/interrupt', { threadId, turnId });
- }
- async unsubscribeThread(threadId: string): Promise<void> {
- await this.#request<{ status: string }>('thread/unsubscribe', { threadId });
- }
- respondToServerRequest(requestId: string | number, result: unknown): void {
- if (!this.#peer) throw new Error('Codex App Server 未运行');
- this.#peer.respond(requestId, result);
- }
- rejectServerRequest(requestId: string | number, code: number, message: string): void {
- if (!this.#peer) throw new Error('Codex App Server 未运行');
- this.#peer.respondError(requestId, { code, message });
- }
- async stop(): Promise<void> {
- this.#generation += 1;
- const child = this.#child;
- const peer = this.#peer;
- this.#child = null;
- this.#peer = null;
- for (const waiter of this.#turnWaiters.values()) {
- clearTimeout(waiter.timer);
- waiter.reject(new Error('Codex App Server 在本轮对话进行中被停止'));
- }
- this.#turnWaiters.clear();
- this.#earlyTurnStates.clear();
- this.#runTurnStartsInFlight = 0;
- if (child && child.exitCode === null) {
- peer?.endOutput();
- await Promise.race([
- new Promise<void>((resolve) => child.once('exit', () => resolve())),
- new Promise<void>((resolve) => setTimeout(resolve, 1_500)),
- ]);
- if (child.exitCode === null) {
- child.kill('SIGTERM');
- await Promise.race([
- new Promise<void>((resolve) => child.once('exit', () => resolve())),
- new Promise<void>((resolve) => setTimeout(resolve, 1_000)),
- ]);
- }
- if (child.exitCode === null) child.kill('SIGKILL');
- }
- peer?.close(new Error('Codex App Server 已停止'));
- this.#runtime = { ...emptyRuntimeStatus(), providerId: this.#provider?.id ?? null };
- this.emit('status', this.status);
- }
- async dispose(): Promise<void> {
- await this.stop();
- this.removeAllListeners();
- }
- async #start(generation: number): Promise<RuntimeStatus> {
- this.#runtime = { ...emptyRuntimeStatus(), state: 'starting', providerId: this.#provider?.id ?? null };
- this.emit('status', this.status);
- let child: ChildProcessWithoutNullStreams | null = null;
- let peer: JsonRpcPeer | null = null;
- let watcher: ExitWatcher | null = null;
- try {
- const binaryPath = await locateCodexBinary();
- this.#assertGeneration(generation);
- const version = readCodexVersion(binaryPath);
- this.#assertGeneration(generation);
- const provider = this.#provider;
- child = spawn(binaryPath, buildArgs(provider), {
- cwd: process.cwd(),
- env: buildEnv(this.#codexHome, provider),
- stdio: ['pipe', 'pipe', 'pipe'],
- windowsHide: true,
- });
- this.#child = child;
- peer = new JsonRpcPeer(child.stdout, child.stdin);
- this.#peer = peer;
- const spawnedChild = child;
- const spawnedPeer = peer;
- peer.on('notification', (notification) => this.#handleNotification(notification));
- peer.on('serverRequest', (request: JsonRpcServerRequest) => this.emit('serverRequest', request));
- peer.on('protocolError', (error) => this.emit('diagnostic', asError(error).message));
- watcher = createExitWatcher(
- spawnedChild,
- (line) => this.emit('diagnostic', sanitizeDiagnostic(line)),
- (error) => {
- if (this.#peer !== spawnedPeer || this.#child !== spawnedChild) return;
- this.#handleExit(spawnedChild, error);
- },
- );
- child.stderr.setEncoding('utf8');
- child.stderr.on('data', (chunk: string) => {
- for (const line of chunk.split(/\r?\n/u).filter(Boolean)) watcher?.noteStderr(line);
- });
- child.once('error', (error) => watcher?.shutdown(error.message));
- child.once('exit', () => watcher?.shutdown('Codex App Server 已退出'));
- peer.once('closed', (error) => watcher?.shutdown(asError(error).message));
- peer.start();
- const initialized = await peer.request<InitializeResponse>('initialize', {
- clientInfo: this.#clientInfo,
- capabilities: { experimentalApi: true },
- });
- this.#assertGeneration(generation);
- peer.notify('initialized');
- this.#runtime = {
- ...emptyRuntimeStatus(),
- state: 'ready',
- binaryPath,
- version,
- codexHome: initialized.codexHome,
- providerId: provider?.id ?? null,
- };
- return await this.refresh();
- } catch (error) {
- // 子进程启动即退出的话真话只在 stderr 里:等它落定,别把「协议流已关闭」丢给用户猜
- const fallback = asError(error).message;
- const message = watcher ? (await watcher.describe(fallback)).message : fallback;
- peer?.close(new Error(message));
- if (this.#peer === peer) this.#peer = null;
- if (this.#child === child) this.#child = null;
- if (child && child.exitCode === null) child.kill('SIGTERM');
- if (generation === this.#generation) {
- this.#runtime = { ...this.#runtime, state: 'error', error: sanitizeDiagnostic(message) };
- this.emit('status', this.status);
- }
- // 换掉原文而不是换掉错误对象:调用方拿到的类型与 stack 不变,只是话变清楚了
- if (error instanceof Error) error.message = message;
- throw error;
- }
- }
- #request<T>(method: string, params?: unknown): Promise<T> {
- if (!this.#peer) throw new Error('Codex App Server 未运行');
- const startedAt = Date.now();
- // 只报慢的和失败的:正常一轮有几十条 RPC,全打等于没打
- return this.#peer.request<T>(method, params).then(
- (value) => {
- const ms = Date.now() - startedAt;
- if (ms > SLOW_RPC_MS) {
- this.emit('diagnostic', `layer=rpc ${method} 用了 ${(ms / 1000).toFixed(1)} 秒`);
- }
- return value;
- },
- (error: unknown) => {
- const ms = Date.now() - startedAt;
- this.emit(
- 'diagnostic',
- `layer=rpc ${method} 失败(${(ms / 1000).toFixed(1)} 秒):${sanitizeDiagnostic(error instanceof Error ? error.message : String(error))}`,
- );
- throw error;
- },
- );
- }
- #handleNotification(notification: { method: string; params?: unknown }): void {
- const params = asRecord(notification.params);
- const turnId = readString(params?.turnId) ?? readString(asRecord(params?.turn)?.id);
- if (turnId) {
- const waiter = this.#turnWaiters.get(turnId);
- if (waiter && notification.method === 'item/agentMessage/delta') {
- waiter.text += readString(params?.delta) ?? '';
- }
- if (waiter && notification.method === 'item/completed') {
- const item = asRecord(params?.item);
- if (item?.type === 'agentMessage' && typeof item.text === 'string') waiter.text = item.text;
- }
- if (waiter && notification.method === 'turn/completed') {
- clearTimeout(waiter.timer);
- this.#turnWaiters.delete(turnId);
- const turn = asRecord(params?.turn);
- waiter.resolve({
- turnId,
- status: readString(turn?.status) ?? 'completed',
- text: waiter.text,
- raw: params,
- });
- }
- if (!waiter && this.#runTurnStartsInFlight > 0) {
- const early = this.#earlyTurnStates.get(turnId) ?? { text: '', completed: null };
- if (notification.method === 'item/agentMessage/delta') {
- early.text += readString(params?.delta) ?? '';
- }
- if (notification.method === 'item/completed') {
- const item = asRecord(params?.item);
- if (item?.type === 'agentMessage' && typeof item.text === 'string') early.text = item.text;
- }
- if (notification.method === 'turn/completed') {
- const turn = asRecord(params?.turn);
- early.completed = {
- turnId,
- status: readString(turn?.status) ?? 'completed',
- text: early.text,
- raw: params,
- };
- }
- this.#earlyTurnStates.set(turnId, early);
- }
- }
- this.emit('notification', notification);
- }
- #handleExit(child: ChildProcessWithoutNullStreams, error: Error): void {
- if (this.#child !== child) return;
- this.#peer?.close(error);
- this.#peer = null;
- this.#child = null;
- this.#runtime = {
- ...this.#runtime,
- state: 'error',
- degraded: true,
- error: sanitizeDiagnostic(error.message),
- };
- for (const waiter of this.#turnWaiters.values()) {
- clearTimeout(waiter.timer);
- waiter.reject(error);
- }
- this.#turnWaiters.clear();
- this.#earlyTurnStates.clear();
- this.#runTurnStartsInFlight = 0;
- this.emit('status', this.status);
- }
- #assertGeneration(generation: number): void {
- if (generation !== this.#generation) throw new Error('Codex App Server 的启动已被取消');
- }
- #releaseRunTurnStart(): void {
- this.#runTurnStartsInFlight = Math.max(0, this.#runTurnStartsInFlight - 1);
- if (this.#runTurnStartsInFlight === 0) this.#earlyTurnStates.clear();
- }
- }
- /** 每个配置项单独一个 -c,避免手写嵌套 inline table */
- export function buildArgs(provider: CodexProviderSpec | null): string[] {
- const args = ['app-server', '--listen', 'stdio://'];
- if (!provider) return args;
- args.push('-c', `model_provider=${tomlString(provider.id)}`);
- args.push('-c', `model=${tomlString(provider.model)}`);
- const prefix = `model_providers.${provider.id}`;
- if (provider.name) args.push('-c', `${prefix}.name=${tomlString(provider.name)}`);
- if (provider.baseUrl) args.push('-c', `${prefix}.base_url=${tomlString(provider.baseUrl)}`);
- // wire_api 只支持 responses:0.155.1 已下线 chat
- args.push('-c', `${prefix}.wire_api="responses"`);
- args.push('-c', `${prefix}.requires_openai_auth=false`);
- if (provider.apiKey) {
- args.push('-c', `${prefix}.env_key=${tomlString(provider.envKey ?? DEFAULT_API_KEY_ENV)}`);
- }
- if (provider.httpHeaders && Object.keys(provider.httpHeaders).length) {
- args.push('-c', `${prefix}.http_headers=${tomlInlineTable(provider.httpHeaders)}`);
- }
- if (provider.envHttpHeaders && Object.keys(provider.envHttpHeaders).length) {
- args.push('-c', `${prefix}.env_http_headers=${tomlInlineTable(provider.envHttpHeaders)}`);
- }
- for (const [key, value] of Object.entries(provider.extra ?? {})) {
- if (!PROVIDER_EXTRA_ALLOWED.has(key) || value === null || value === undefined) continue;
- args.push('-c', `${prefix}.${key}=${tomlValue(value)}`);
- }
- return args;
- }
- function buildEnv(codexHome: string, provider: CodexProviderSpec | null): NodeJS.ProcessEnv {
- const env: NodeJS.ProcessEnv = {
- ...process.env,
- CODEX_HOME: codexHome,
- RUST_LOG: process.env.RUST_LOG ?? 'warn',
- LOG_FORMAT: 'json',
- };
- // 密钥只存在于子进程环境变量,config.toml 里只写 env_key 名字
- if (provider?.apiKey) {
- env[provider.envKey ?? DEFAULT_API_KEY_ENV] = provider.apiKey;
- }
- return env;
- }
- function tomlString(value: string): string {
- return JSON.stringify(value);
- }
- function tomlInlineTable(value: Record<string, string>): string {
- const entries = Object.entries(value).map(([key, item]) => `${tomlKey(key)}=${tomlString(item)}`);
- return `{${entries.join(',')}}`;
- }
- function tomlKey(key: string): string {
- return /^[A-Za-z0-9_-]+$/.test(key) ? key : tomlString(key);
- }
- function tomlValue(value: JsonValue): string {
- if (typeof value === 'string') return tomlString(value);
- if (typeof value === 'number' || typeof value === 'boolean') return String(value);
- if (Array.isArray(value)) {
- return `[${value.filter((item) => item !== null).map(tomlValue).join(',')}]`;
- }
- if (value && typeof value === 'object') {
- // TOML inline table 用 `=` 而不是 JSON 的 `:`,不能直接 JSON.stringify
- return tomlInlineTable(
- Object.fromEntries(
- Object.entries(value as Record<string, JsonValue>)
- .filter(([, item]) => typeof item === 'string')
- .map(([key, item]) => [key, String(item)]),
- ),
- );
- }
- return '""';
- }
- function emptyRuntimeStatus(): RuntimeStatus {
- return {
- state: 'stopped',
- binaryPath: null,
- version: null,
- codexHome: null,
- models: [],
- capabilities: emptyCapabilities(),
- providerId: null,
- defaultModel: null,
- degraded: false,
- error: null,
- };
- }
- function emptyCapabilities(): RuntimeCapabilities {
- return { namespaceTools: false, imageGeneration: false, webSearch: false };
- }
- function asRecord(value: unknown): Record<string, unknown> | null {
- return value && typeof value === 'object' && !Array.isArray(value)
- ? (value as Record<string, unknown>)
- : null;
- }
- function readString(value: unknown): string | null {
- return typeof value === 'string' ? value : null;
- }
- /**
- * mcpServerStatus/list 的 tools 是「工具名 → 详情」映射,详情结构由 Codex 版本决定。
- * 这里只认 name / description,认不出的形态退化成只留 key 名,避免协议一变列表就空掉。
- */
- export function toToolSummaries(
- tools: Record<string, unknown> | null | undefined,
- ): Array<{ name: string; description: string }> {
- return Object.entries(tools ?? {}).map(([key, value]) => {
- const detail = asRecord(value);
- return {
- name: readString(detail?.name) || key,
- description: readString(detail?.description) || '',
- };
- });
- }
- function asError(value: unknown): Error {
- return value instanceof Error ? value : new Error(String(value));
- }
- export function sanitizeDiagnostic(value: string): string {
- return value
- .replace(/sk-[A-Za-z0-9_-]{12,}/gu, 'sk-[redacted]')
- .replace(/Bearer\s+[A-Za-z0-9._~-]+/giu, 'Bearer [redacted]')
- .slice(0, 4_000);
- }
- function delay(ms: number): Promise<void> {
- return new Promise((resolve) => {
- const timer = setTimeout(resolve, ms);
- timer.unref();
- });
- }
- /**
- * Codex 的致命错误由 anyhow 打成纯文本(`Error: ...`),RUST_LOG 的常规日志是 JSON 行;
- * 纯文本里还混着 `WARNING: ...` 这类噪音,所以优先取 Error 行,取不到才退化成第一条纯文本行。
- */
- function pickFatalLine(tail: string[]): string | null {
- const plain = tail
- .map((line) => line.trim())
- .filter((line) => line && !line.startsWith('{'));
- if (!plain.length) return null;
- const fatal = plain.find((line) => /^(error|panic)\b/iu.test(line) || line.includes('panicked at'));
- return sanitizeDiagnostic((fatal ?? plain[0]).slice(0, 500));
- }
- interface ExitWatcher {
- /** 收 stderr:既实时推诊断,也留尾部若干行用于拼退出原因 */
- noteStderr(line: string): void;
- /** 统一退出出口(流关闭 / 进程退出 / spawn 失败),只生效一次 */
- shutdown(fallback: string): void;
- /** 等 stderr 收完,把没有信息量的 fallback 换成带 exit code 与原因原文的错误 */
- describe(fallback: string): Promise<Error>;
- }
- function createExitWatcher(
- child: ChildProcessWithoutNullStreams,
- onDiagnostic: (line: string) => void,
- onExit: (error: Error) => void,
- ): ExitWatcher {
- const tail: string[] = [];
- const closed = new Promise<void>((resolve) => child.once('close', () => resolve()));
- let shuttingDown = false;
- const describe = async (fallback: string): Promise<Error> => {
- // 等不到 close 也得往下走:一次卡住的退出不能把启动的 Promise 挂在那里
- await Promise.race([closed, delay(EXIT_GRACE_MS)]);
- const code = child.exitCode ?? child.signalCode;
- const reason = code === null || code === undefined ? fallback : `Codex App Server 退出(${code})`;
- const detail = pickFatalLine(tail);
- return new Error(detail ? `${reason}:${detail}` : reason);
- };
- return {
- noteStderr(line: string): void {
- tail.push(line);
- if (tail.length > STDERR_TAIL_LINES) tail.shift();
- onDiagnostic(line);
- },
- describe,
- shutdown(fallback: string): void {
- if (shuttingDown) return;
- shuttingDown = true;
- void describe(fallback).then((error) => {
- // 流都关了进程还活着(卡死的子进程)只能自己收
- if (child.exitCode === null && child.signalCode === null) child.kill('SIGTERM');
- onExit(error);
- });
- },
- };
- }
|