codexRuntime.ts 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942
  1. import { EventEmitter } from 'node:events';
  2. import { spawn, type ChildProcessWithoutNullStreams } from 'node:child_process';
  3. import { locateCodexBinary, readCodexVersion } from './codexLocator';
  4. import { JsonRpcPeer, type JsonRpcServerRequest } from './jsonRpcPeer';
  5. /**
  6. * 移植自 Noobi.ai src/main/codexAppServer.ts,主要改动:
  7. * 1. 删掉全部账号相关能力(account/read、account/login/start、account/logout)——本模块不做任何登录;
  8. * 2. 不再传 --strict-config(用户与程序共写 config.toml,未知字段只应告警不该退出);
  9. * 3. 新增 provider 注入:用 -c 覆盖 model_provider / model / model_providers.*,api_key 只走子进程环境变量;
  10. * 4. clientInfo 改为本项目标识。
  11. */
  12. export type JsonValue = null | boolean | number | string | JsonValue[] | { [key: string]: JsonValue };
  13. export interface ModelOption {
  14. id: string;
  15. model: string;
  16. displayName: string;
  17. description: string;
  18. isDefault: boolean;
  19. defaultEffort: string;
  20. efforts: string[];
  21. }
  22. export interface RuntimeCapabilities {
  23. namespaceTools: boolean;
  24. imageGeneration: boolean;
  25. webSearch: boolean;
  26. }
  27. export interface RuntimeStatus {
  28. state: 'stopped' | 'starting' | 'ready' | 'error';
  29. binaryPath: string | null;
  30. version: string | null;
  31. codexHome: string | null;
  32. models: ModelOption[];
  33. capabilities: RuntimeCapabilities;
  34. providerId: string | null;
  35. defaultModel: string | null;
  36. /** model/list 拿不到可用模型时为 true */
  37. degraded: boolean;
  38. error: string | null;
  39. }
  40. export interface CodexMcpServerStatus {
  41. name: string;
  42. authStatus: string;
  43. connected: boolean;
  44. toolCount: number;
  45. tools: Array<{ name: string; description: string }>;
  46. }
  47. export interface CodexSkillSummary {
  48. name: string;
  49. description: string;
  50. path: string;
  51. scope: 'user' | 'repo' | 'system' | 'admin' | string;
  52. enabled: boolean;
  53. cwd: string;
  54. }
  55. /** 一次 spawn 用的模型来源:一律是自定义 provider,Codex 的内置 id 是保留字、不可覆盖 */
  56. export interface CodexProviderSpec {
  57. id: string;
  58. name?: string | null;
  59. baseUrl: string;
  60. model: string;
  61. /** apiKey 只注入子进程环境变量,不落任何盘 */
  62. apiKey?: string | null;
  63. envKey?: string;
  64. httpHeaders?: Record<string, string> | null;
  65. envHttpHeaders?: Record<string, string> | null;
  66. /** 仅放行 PROVIDER_EXTRA_ALLOWLIST 里的键,其余在 providerService 就被丢弃 */
  67. extra?: Record<string, JsonValue> | null;
  68. /**
  69. * model_catalog_json 的绝对路径(providerService 落盘后填入)。
  70. * 自定义模型不在 codex 内置目录里,没有目录就会用兜底元数据并打
  71. * "Model metadata not found" 警告;目录整体替换内置目录,-c 随进程启动生效。
  72. */
  73. modelCatalogPath?: string | null;
  74. }
  75. export interface CodexRuntimeOptions {
  76. codexHome: string;
  77. provider?: CodexProviderSpec | null;
  78. clientInfo?: { name: string; title: string; version: string };
  79. }
  80. export interface StartThreadOptions {
  81. cwd: string;
  82. model?: string | null;
  83. sandbox?: 'read-only' | 'workspace-write' | 'danger-full-access';
  84. approvalPolicy?: 'untrusted' | 'on-request' | 'never';
  85. developerInstructions?: string;
  86. ephemeral?: boolean;
  87. }
  88. export interface StartTurnOptions {
  89. threadId: string;
  90. prompt: string;
  91. cwd?: string;
  92. model?: string | null;
  93. effort?: string | null;
  94. approvalPolicy?: 'untrusted' | 'on-request' | 'never';
  95. skills?: Array<{ name: string; path: string }>;
  96. /** 仅宿主侧的完成期限,不会发给 App Server */
  97. timeoutMs?: number;
  98. }
  99. export interface TurnResult {
  100. turnId: string;
  101. status: string;
  102. text: string;
  103. raw: unknown;
  104. }
  105. interface SkillsListResponse {
  106. data: Array<{
  107. cwd: string;
  108. skills: Array<{
  109. name: string;
  110. description: string;
  111. path: string;
  112. scope: string;
  113. enabled: boolean;
  114. }>;
  115. }>;
  116. }
  117. interface ListMcpServerStatusResponse {
  118. data: Array<{
  119. name: string;
  120. serverInfo: unknown | null;
  121. tools: Record<string, unknown>;
  122. authStatus: string;
  123. }>;
  124. nextCursor: string | null;
  125. }
  126. interface InitializeResponse {
  127. userAgent: string;
  128. codexHome: string;
  129. platformFamily: string;
  130. platformOs: string;
  131. }
  132. interface ModelListResponse {
  133. data: Array<{
  134. id: string;
  135. model: string;
  136. displayName: string;
  137. description: string;
  138. isDefault: boolean;
  139. defaultReasoningEffort: string;
  140. supportedReasoningEfforts: Array<{ reasoningEffort: string }>;
  141. }>;
  142. nextCursor: string | null;
  143. }
  144. interface ModelProviderCapabilitiesResponse {
  145. namespaceTools: boolean;
  146. imageGeneration: boolean;
  147. webSearch: boolean;
  148. }
  149. interface ThreadResponse {
  150. thread: { id: string };
  151. model: string;
  152. }
  153. interface TurnStartResponse {
  154. turn: { id: string; status: string };
  155. }
  156. interface EarlyTurnState {
  157. text: string;
  158. completed: TurnResult | null;
  159. }
  160. interface TurnWaiter {
  161. text: string;
  162. resolve(result: TurnResult): void;
  163. reject(error: Error): void;
  164. timer: NodeJS.Timeout;
  165. }
  166. const TURN_TIMEOUT_MS = 20 * 60 * 1_000;
  167. /** 单条 RPC 超过这个耗时才值得进诊断(jsonRpcPeer 的预算是 30 秒) */
  168. const SLOW_RPC_MS = 3_000;
  169. /** 退出事件可能比 stderr 落地慢半拍,等一下再拼错误信息 */
  170. const EXIT_GRACE_MS = 300;
  171. const STDERR_TAIL_LINES = 40;
  172. export const DEFAULT_API_KEY_ENV = 'ZSJZ_CODEX_API_KEY';
  173. /**
  174. * 模型管理 config 里允许透传给 Codex 的 provider 字段(对应 rust-v0.155.1
  175. * model-provider-info/src/lib.rs 的 ModelProviderInfo)。providerService 用它做白名单,
  176. * 这里做 -c 落参,两处必须同源,否则页面提示"已丢弃"而实际写进去、或反之。
  177. */
  178. export const PROVIDER_EXTRA_ALLOWLIST: readonly string[] = Object.freeze([
  179. 'request_max_retries',
  180. 'stream_max_retries',
  181. 'stream_idle_timeout_ms',
  182. 'websocket_connect_timeout_ms',
  183. 'supports_websockets',
  184. 'query_params',
  185. ]);
  186. const PROVIDER_EXTRA_ALLOWED = new Set<string>(PROVIDER_EXTRA_ALLOWLIST);
  187. export class CodexRuntime extends EventEmitter {
  188. readonly #codexHome: string;
  189. readonly #clientInfo: { name: string; title: string; version: string };
  190. #provider: CodexProviderSpec | null;
  191. #child: ChildProcessWithoutNullStreams | null = null;
  192. #peer: JsonRpcPeer | null = null;
  193. #startPromise: Promise<RuntimeStatus> | null = null;
  194. #runtime: RuntimeStatus = emptyRuntimeStatus();
  195. #turnWaiters = new Map<string, TurnWaiter>();
  196. #earlyTurnStates = new Map<string, EarlyTurnState>();
  197. #runTurnStartsInFlight = 0;
  198. #generation = 0;
  199. constructor(options: CodexRuntimeOptions) {
  200. super();
  201. this.#codexHome = options.codexHome;
  202. this.#provider = options.provider ?? null;
  203. this.#clientInfo = options.clientInfo ?? {
  204. name: 'zsjz_ai',
  205. title: '清鉴线索调查工具',
  206. version: '0.1.0',
  207. };
  208. }
  209. get status(): RuntimeStatus {
  210. return structuredClone(this.#runtime);
  211. }
  212. get provider(): CodexProviderSpec | null {
  213. return this.#provider ? { ...this.#provider, apiKey: null } : null;
  214. }
  215. /** 换 provider 必须重启子进程:api_key 走的是子进程环境变量,无法热更新 */
  216. async applyProvider(provider: CodexProviderSpec | null): Promise<RuntimeStatus> {
  217. const previous = this.#provider;
  218. this.#provider = provider;
  219. if (this.#child) await this.stop();
  220. try {
  221. return await this.start();
  222. } catch (error) {
  223. // 失败要回滚:坏 spec 留在内存里,之后每次 start 都会重放同一组参数,
  224. // 一条选错的模型能把运行时锁死到重启应用为止。
  225. this.#provider = previous;
  226. throw error;
  227. }
  228. }
  229. async start(): Promise<RuntimeStatus> {
  230. if (this.#runtime.state === 'ready' && this.#peer) return this.status;
  231. if (this.#startPromise) return this.#startPromise;
  232. const generation = ++this.#generation;
  233. this.#startPromise = this.#start(generation);
  234. try {
  235. return await this.#startPromise;
  236. } finally {
  237. this.#startPromise = null;
  238. }
  239. }
  240. async refresh(): Promise<RuntimeStatus> {
  241. await this.start();
  242. let models: ModelOption[] = [];
  243. let degraded = false;
  244. try {
  245. models = await this.listModels();
  246. } catch (error) {
  247. // 第三方/本地端点常常不实现 model/list,这里降级而不是失败
  248. degraded = true;
  249. this.emit('diagnostic', `model/list 不可用:${asError(error).message}`);
  250. }
  251. const capabilities = await this.readModelProviderCapabilities().catch(() => emptyCapabilities());
  252. this.#runtime = {
  253. ...this.#runtime,
  254. models,
  255. capabilities,
  256. degraded: degraded || models.length === 0,
  257. // defaultModel 只取 provider 里配置的模型。绝不回退到 model/list 的内置目录
  258. // 默认值(gpt-6-astra 之类)—— 那是 Codex 自家目录的模型,与「模型管理」无关,
  259. // 一旦漏到 UI 上就会让用户以为应用在用 GPT 模型。
  260. defaultModel: this.#provider?.model ?? null,
  261. state: 'ready',
  262. error: null,
  263. };
  264. this.emit('status', this.status);
  265. return this.status;
  266. }
  267. async readModelProviderCapabilities(): Promise<RuntimeCapabilities> {
  268. const result = await this.#request<ModelProviderCapabilitiesResponse>(
  269. 'modelProvider/capabilities/read',
  270. {},
  271. );
  272. return {
  273. namespaceTools: result.namespaceTools === true,
  274. imageGeneration: result.imageGeneration === true,
  275. webSearch: result.webSearch === true,
  276. };
  277. }
  278. async listModels(): Promise<ModelOption[]> {
  279. const models: ModelOption[] = [];
  280. let cursor: string | null = null;
  281. for (let page = 0; page < 10; page += 1) {
  282. const result: ModelListResponse = await this.#request<ModelListResponse>('model/list', {
  283. cursor,
  284. limit: 100,
  285. includeHidden: false,
  286. });
  287. models.push(
  288. ...result.data.map((model) => ({
  289. id: model.id,
  290. model: model.model,
  291. displayName: model.displayName,
  292. description: model.description,
  293. isDefault: model.isDefault,
  294. defaultEffort: model.defaultReasoningEffort,
  295. efforts: model.supportedReasoningEfforts.map((item) => item.reasoningEffort),
  296. })),
  297. );
  298. cursor = result.nextCursor;
  299. if (!cursor) break;
  300. }
  301. return models;
  302. }
  303. async listSkills(options: { cwds?: string[]; forceReload?: boolean } = {}): Promise<CodexSkillSummary[]> {
  304. await this.start();
  305. const result = await this.#request<SkillsListResponse>('skills/list', {
  306. ...(options.cwds?.length ? { cwds: options.cwds } : {}),
  307. forceReload: options.forceReload ?? false,
  308. });
  309. return result.data.flatMap((entry) =>
  310. entry.skills.map((skill) => ({
  311. name: skill.name,
  312. description: skill.description,
  313. path: skill.path,
  314. scope: skill.scope,
  315. enabled: skill.enabled,
  316. cwd: entry.cwd,
  317. })),
  318. );
  319. }
  320. async setSkillEnabled(selector: { path?: string; name?: string }, enabled: boolean): Promise<boolean> {
  321. await this.start();
  322. const result = await this.#request<{ effectiveEnabled: boolean }>('skills/config/write', {
  323. ...(selector.path ? { path: selector.path } : {}),
  324. ...(!selector.path && selector.name ? { name: selector.name } : {}),
  325. enabled,
  326. });
  327. return result.effectiveEnabled;
  328. }
  329. async listMcpServerStatuses(): Promise<CodexMcpServerStatus[]> {
  330. await this.start();
  331. const statuses: CodexMcpServerStatus[] = [];
  332. let cursor: string | null = null;
  333. for (let page = 0; page < 20; page += 1) {
  334. const result: ListMcpServerStatusResponse = await this.#request<ListMcpServerStatusResponse>(
  335. 'mcpServerStatus/list',
  336. {
  337. cursor,
  338. limit: 100,
  339. detail: 'toolsAndAuthOnly',
  340. },
  341. );
  342. statuses.push(
  343. ...result.data.map((server) => ({
  344. name: server.name,
  345. authStatus: server.authStatus,
  346. connected: server.serverInfo !== null,
  347. toolCount: Object.keys(server.tools ?? {}).length,
  348. tools: toToolSummaries(server.tools),
  349. })),
  350. );
  351. cursor = result.nextCursor;
  352. if (!cursor) break;
  353. }
  354. return statuses;
  355. }
  356. async readConfig(): Promise<Record<string, unknown>> {
  357. await this.start();
  358. const result = await this.#request<{ config: Record<string, unknown> }>('config/read', {
  359. includeLayers: false,
  360. });
  361. return structuredClone(result.config);
  362. }
  363. /** 局部写:Codex 自己也持有 config.toml 写权,禁止整文件覆盖 */
  364. async writeConfigValue(keyPath: string, value: JsonValue): Promise<void> {
  365. await this.start();
  366. await this.#request<unknown>('config/value/write', { keyPath, value, mergeStrategy: 'replace' });
  367. }
  368. async reloadMcpServers(): Promise<void> {
  369. await this.start();
  370. await this.#request<unknown>('config/mcpServer/reload');
  371. }
  372. async startThread(options: StartThreadOptions): Promise<string> {
  373. await this.start();
  374. const result = await this.#request<ThreadResponse>('thread/start', {
  375. cwd: options.cwd,
  376. ...(options.model ? { model: options.model } : {}),
  377. sandbox: options.sandbox ?? 'workspace-write',
  378. approvalPolicy: options.approvalPolicy ?? 'never',
  379. ...(options.developerInstructions ? { developerInstructions: options.developerInstructions } : {}),
  380. ...(options.ephemeral === undefined ? {} : { ephemeral: options.ephemeral }),
  381. experimentalRawEvents: false,
  382. });
  383. return result.thread.id;
  384. }
  385. async resumeThread(threadId: string, options: StartThreadOptions): Promise<string> {
  386. await this.start();
  387. const result = await this.#request<ThreadResponse>('thread/resume', {
  388. threadId,
  389. cwd: options.cwd,
  390. ...(options.model ? { model: options.model } : {}),
  391. sandbox: options.sandbox ?? 'workspace-write',
  392. approvalPolicy: options.approvalPolicy ?? 'never',
  393. ...(options.developerInstructions ? { developerInstructions: options.developerInstructions } : {}),
  394. // excludeTurns 不影响上游收到的历史(实测 true/false 两轮上下文都完整带回)
  395. excludeTurns: true,
  396. });
  397. return result.thread.id;
  398. }
  399. async startTurn(options: StartTurnOptions): Promise<string> {
  400. await this.start();
  401. const result = await this.#request<TurnStartResponse>('turn/start', {
  402. threadId: options.threadId,
  403. input: [
  404. { type: 'text', text: options.prompt, text_elements: [] },
  405. ...(options.skills ?? []).map((skill) => ({
  406. type: 'skill' as const,
  407. name: skill.name,
  408. path: skill.path,
  409. })),
  410. ],
  411. ...(options.cwd ? { cwd: options.cwd } : {}),
  412. ...(options.model ? { model: options.model } : {}),
  413. ...(options.effort ? { effort: options.effort } : {}),
  414. ...(options.approvalPolicy ? { approvalPolicy: options.approvalPolicy } : {}),
  415. });
  416. return result.turn.id;
  417. }
  418. async runTurn(options: StartTurnOptions): Promise<TurnResult> {
  419. this.#runTurnStartsInFlight += 1;
  420. let turnId: string;
  421. try {
  422. turnId = await this.startTurn(options);
  423. } catch (error) {
  424. this.#releaseRunTurnStart();
  425. throw error;
  426. }
  427. // turn/completed 可能早于 turn/start 的响应回来,先存起来
  428. const startedAt = Date.now();
  429. const noteTurn = (result: TurnResult): TurnResult => {
  430. this.emit(
  431. 'diagnostic',
  432. `layer=turn turn=${result.turnId} status=${result.status} 用时 ${((Date.now() - startedAt) / 1000).toFixed(1)}s`,
  433. );
  434. return result;
  435. };
  436. const early = this.#earlyTurnStates.get(turnId);
  437. this.#earlyTurnStates.delete(turnId);
  438. if (early?.completed) {
  439. this.#releaseRunTurnStart();
  440. return noteTurn(early.completed);
  441. }
  442. return new Promise<TurnResult>((resolve, reject) => {
  443. const timeoutMs =
  444. options.timeoutMs === undefined
  445. ? TURN_TIMEOUT_MS
  446. : Math.min(TURN_TIMEOUT_MS, Math.max(1_000, Math.trunc(options.timeoutMs)));
  447. const timer = setTimeout(() => {
  448. this.#turnWaiters.delete(turnId);
  449. void this.interruptTurn(options.threadId, turnId).catch(() => undefined);
  450. this.emit(
  451. 'diagnostic',
  452. `layer=turn turn=${turnId} 宿主侧期限已到(${Math.round(timeoutMs / 1000)} 秒),已发 turn/interrupt`,
  453. );
  454. reject(new Error(`Codex turn ${turnId} 超时`));
  455. }, timeoutMs);
  456. timer.unref();
  457. this.#turnWaiters.set(turnId, {
  458. text: early?.text ?? '',
  459. resolve: (result) => resolve(noteTurn(result)),
  460. reject,
  461. timer,
  462. });
  463. this.#releaseRunTurnStart();
  464. });
  465. }
  466. async interruptTurn(threadId: string, turnId: string): Promise<void> {
  467. await this.#request<Record<string, never>>('turn/interrupt', { threadId, turnId });
  468. }
  469. async unsubscribeThread(threadId: string): Promise<void> {
  470. await this.#request<{ status: string }>('thread/unsubscribe', { threadId });
  471. }
  472. respondToServerRequest(requestId: string | number, result: unknown): void {
  473. if (!this.#peer) throw new Error('Codex App Server 未运行');
  474. this.#peer.respond(requestId, result);
  475. }
  476. rejectServerRequest(requestId: string | number, code: number, message: string): void {
  477. if (!this.#peer) throw new Error('Codex App Server 未运行');
  478. this.#peer.respondError(requestId, { code, message });
  479. }
  480. async stop(): Promise<void> {
  481. this.#generation += 1;
  482. const child = this.#child;
  483. const peer = this.#peer;
  484. this.#child = null;
  485. this.#peer = null;
  486. for (const waiter of this.#turnWaiters.values()) {
  487. clearTimeout(waiter.timer);
  488. waiter.reject(new Error('Codex App Server 在本轮对话进行中被停止'));
  489. }
  490. this.#turnWaiters.clear();
  491. this.#earlyTurnStates.clear();
  492. this.#runTurnStartsInFlight = 0;
  493. if (child && child.exitCode === null) {
  494. peer?.endOutput();
  495. await Promise.race([
  496. new Promise<void>((resolve) => child.once('exit', () => resolve())),
  497. new Promise<void>((resolve) => setTimeout(resolve, 1_500)),
  498. ]);
  499. if (child.exitCode === null) {
  500. child.kill('SIGTERM');
  501. await Promise.race([
  502. new Promise<void>((resolve) => child.once('exit', () => resolve())),
  503. new Promise<void>((resolve) => setTimeout(resolve, 1_000)),
  504. ]);
  505. }
  506. if (child.exitCode === null) child.kill('SIGKILL');
  507. }
  508. peer?.close(new Error('Codex App Server 已停止'));
  509. this.#runtime = { ...emptyRuntimeStatus(), providerId: this.#provider?.id ?? null };
  510. this.emit('status', this.status);
  511. }
  512. async dispose(): Promise<void> {
  513. await this.stop();
  514. this.removeAllListeners();
  515. }
  516. async #start(generation: number): Promise<RuntimeStatus> {
  517. this.#runtime = { ...emptyRuntimeStatus(), state: 'starting', providerId: this.#provider?.id ?? null };
  518. this.emit('status', this.status);
  519. let child: ChildProcessWithoutNullStreams | null = null;
  520. let peer: JsonRpcPeer | null = null;
  521. let watcher: ExitWatcher | null = null;
  522. try {
  523. const binaryPath = await locateCodexBinary();
  524. this.#assertGeneration(generation);
  525. const version = readCodexVersion(binaryPath);
  526. this.#assertGeneration(generation);
  527. const provider = this.#provider;
  528. child = spawn(binaryPath, buildArgs(provider), {
  529. cwd: process.cwd(),
  530. env: buildEnv(this.#codexHome, provider),
  531. stdio: ['pipe', 'pipe', 'pipe'],
  532. windowsHide: true,
  533. });
  534. this.#child = child;
  535. peer = new JsonRpcPeer(child.stdout, child.stdin);
  536. this.#peer = peer;
  537. const spawnedChild = child;
  538. const spawnedPeer = peer;
  539. peer.on('notification', (notification) => this.#handleNotification(notification));
  540. peer.on('serverRequest', (request: JsonRpcServerRequest) => this.emit('serverRequest', request));
  541. peer.on('protocolError', (error) => this.emit('diagnostic', asError(error).message));
  542. watcher = createExitWatcher(
  543. spawnedChild,
  544. (line) => this.emit('diagnostic', sanitizeDiagnostic(line)),
  545. (error) => {
  546. if (this.#peer !== spawnedPeer || this.#child !== spawnedChild) return;
  547. this.#handleExit(spawnedChild, error);
  548. },
  549. );
  550. child.stderr.setEncoding('utf8');
  551. child.stderr.on('data', (chunk: string) => {
  552. for (const line of chunk.split(/\r?\n/u).filter(Boolean)) watcher?.noteStderr(line);
  553. });
  554. child.once('error', (error) => watcher?.shutdown(error.message));
  555. child.once('exit', () => watcher?.shutdown('Codex App Server 已退出'));
  556. peer.once('closed', (error) => watcher?.shutdown(asError(error).message));
  557. peer.start();
  558. const initialized = await peer.request<InitializeResponse>('initialize', {
  559. clientInfo: this.#clientInfo,
  560. capabilities: { experimentalApi: true },
  561. });
  562. this.#assertGeneration(generation);
  563. peer.notify('initialized');
  564. this.#runtime = {
  565. ...emptyRuntimeStatus(),
  566. state: 'ready',
  567. binaryPath,
  568. version,
  569. codexHome: initialized.codexHome,
  570. providerId: provider?.id ?? null,
  571. };
  572. return await this.refresh();
  573. } catch (error) {
  574. // 子进程启动即退出的话真话只在 stderr 里:等它落定,别把「协议流已关闭」丢给用户猜
  575. const fallback = asError(error).message;
  576. const message = watcher ? (await watcher.describe(fallback)).message : fallback;
  577. peer?.close(new Error(message));
  578. if (this.#peer === peer) this.#peer = null;
  579. if (this.#child === child) this.#child = null;
  580. if (child && child.exitCode === null) child.kill('SIGTERM');
  581. if (generation === this.#generation) {
  582. this.#runtime = { ...this.#runtime, state: 'error', error: sanitizeDiagnostic(message) };
  583. this.emit('status', this.status);
  584. }
  585. // 换掉原文而不是换掉错误对象:调用方拿到的类型与 stack 不变,只是话变清楚了
  586. if (error instanceof Error) error.message = message;
  587. throw error;
  588. }
  589. }
  590. #request<T>(method: string, params?: unknown): Promise<T> {
  591. if (!this.#peer) throw new Error('Codex App Server 未运行');
  592. const startedAt = Date.now();
  593. // 只报慢的和失败的:正常一轮有几十条 RPC,全打等于没打
  594. return this.#peer.request<T>(method, params).then(
  595. (value) => {
  596. const ms = Date.now() - startedAt;
  597. if (ms > SLOW_RPC_MS) {
  598. this.emit('diagnostic', `layer=rpc ${method} 用了 ${(ms / 1000).toFixed(1)} 秒`);
  599. }
  600. return value;
  601. },
  602. (error: unknown) => {
  603. const ms = Date.now() - startedAt;
  604. this.emit(
  605. 'diagnostic',
  606. `layer=rpc ${method} 失败(${(ms / 1000).toFixed(1)} 秒):${sanitizeDiagnostic(error instanceof Error ? error.message : String(error))}`,
  607. );
  608. throw error;
  609. },
  610. );
  611. }
  612. #handleNotification(notification: { method: string; params?: unknown }): void {
  613. const params = asRecord(notification.params);
  614. const turnId = readString(params?.turnId) ?? readString(asRecord(params?.turn)?.id);
  615. if (turnId) {
  616. const waiter = this.#turnWaiters.get(turnId);
  617. if (waiter && notification.method === 'item/agentMessage/delta') {
  618. waiter.text += readString(params?.delta) ?? '';
  619. }
  620. if (waiter && notification.method === 'item/completed') {
  621. const item = asRecord(params?.item);
  622. if (item?.type === 'agentMessage' && typeof item.text === 'string') waiter.text = item.text;
  623. }
  624. if (waiter && notification.method === 'turn/completed') {
  625. clearTimeout(waiter.timer);
  626. this.#turnWaiters.delete(turnId);
  627. const turn = asRecord(params?.turn);
  628. waiter.resolve({
  629. turnId,
  630. status: readString(turn?.status) ?? 'completed',
  631. text: waiter.text,
  632. raw: params,
  633. });
  634. }
  635. if (!waiter && this.#runTurnStartsInFlight > 0) {
  636. const early = this.#earlyTurnStates.get(turnId) ?? { text: '', completed: null };
  637. if (notification.method === 'item/agentMessage/delta') {
  638. early.text += readString(params?.delta) ?? '';
  639. }
  640. if (notification.method === 'item/completed') {
  641. const item = asRecord(params?.item);
  642. if (item?.type === 'agentMessage' && typeof item.text === 'string') early.text = item.text;
  643. }
  644. if (notification.method === 'turn/completed') {
  645. const turn = asRecord(params?.turn);
  646. early.completed = {
  647. turnId,
  648. status: readString(turn?.status) ?? 'completed',
  649. text: early.text,
  650. raw: params,
  651. };
  652. }
  653. this.#earlyTurnStates.set(turnId, early);
  654. }
  655. }
  656. this.emit('notification', notification);
  657. }
  658. #handleExit(child: ChildProcessWithoutNullStreams, error: Error): void {
  659. if (this.#child !== child) return;
  660. this.#peer?.close(error);
  661. this.#peer = null;
  662. this.#child = null;
  663. this.#runtime = {
  664. ...this.#runtime,
  665. state: 'error',
  666. degraded: true,
  667. error: sanitizeDiagnostic(error.message),
  668. };
  669. for (const waiter of this.#turnWaiters.values()) {
  670. clearTimeout(waiter.timer);
  671. waiter.reject(error);
  672. }
  673. this.#turnWaiters.clear();
  674. this.#earlyTurnStates.clear();
  675. this.#runTurnStartsInFlight = 0;
  676. this.emit('status', this.status);
  677. }
  678. #assertGeneration(generation: number): void {
  679. if (generation !== this.#generation) throw new Error('Codex App Server 的启动已被取消');
  680. }
  681. #releaseRunTurnStart(): void {
  682. this.#runTurnStartsInFlight = Math.max(0, this.#runTurnStartsInFlight - 1);
  683. if (this.#runTurnStartsInFlight === 0) this.#earlyTurnStates.clear();
  684. }
  685. }
  686. /** 每个配置项单独一个 -c,避免手写嵌套 inline table */
  687. export function buildArgs(provider: CodexProviderSpec | null): string[] {
  688. const args = ['app-server', '--listen', 'stdio://'];
  689. if (!provider) return args;
  690. args.push('-c', `model_provider=${tomlString(provider.id)}`);
  691. args.push('-c', `model=${tomlString(provider.model)}`);
  692. const prefix = `model_providers.${provider.id}`;
  693. if (provider.name) args.push('-c', `${prefix}.name=${tomlString(provider.name)}`);
  694. if (provider.baseUrl) args.push('-c', `${prefix}.base_url=${tomlString(provider.baseUrl)}`);
  695. // wire_api 只支持 responses:0.155.1 已下线 chat
  696. args.push('-c', `${prefix}.wire_api="responses"`);
  697. args.push('-c', `${prefix}.requires_openai_auth=false`);
  698. if (provider.apiKey) {
  699. args.push('-c', `${prefix}.env_key=${tomlString(provider.envKey ?? DEFAULT_API_KEY_ENV)}`);
  700. }
  701. if (provider.httpHeaders && Object.keys(provider.httpHeaders).length) {
  702. args.push('-c', `${prefix}.http_headers=${tomlInlineTable(provider.httpHeaders)}`);
  703. }
  704. if (provider.envHttpHeaders && Object.keys(provider.envHttpHeaders).length) {
  705. args.push('-c', `${prefix}.env_http_headers=${tomlInlineTable(provider.envHttpHeaders)}`);
  706. }
  707. for (const [key, value] of Object.entries(provider.extra ?? {})) {
  708. if (!PROVIDER_EXTRA_ALLOWED.has(key) || value === null || value === undefined) continue;
  709. args.push('-c', `${prefix}.${key}=${tomlValue(value)}`);
  710. }
  711. if (provider.modelCatalogPath) {
  712. args.push('-c', `model_catalog_json=${tomlString(provider.modelCatalogPath)}`);
  713. }
  714. return args;
  715. }
  716. function buildEnv(codexHome: string, provider: CodexProviderSpec | null): NodeJS.ProcessEnv {
  717. const env: NodeJS.ProcessEnv = {
  718. ...process.env,
  719. CODEX_HOME: codexHome,
  720. RUST_LOG: process.env.RUST_LOG ?? 'warn',
  721. LOG_FORMAT: 'json',
  722. };
  723. // 密钥只存在于子进程环境变量,config.toml 里只写 env_key 名字
  724. if (provider?.apiKey) {
  725. env[provider.envKey ?? DEFAULT_API_KEY_ENV] = provider.apiKey;
  726. }
  727. return env;
  728. }
  729. function tomlString(value: string): string {
  730. return JSON.stringify(value);
  731. }
  732. function tomlInlineTable(value: Record<string, string>): string {
  733. const entries = Object.entries(value).map(([key, item]) => `${tomlKey(key)}=${tomlString(item)}`);
  734. return `{${entries.join(',')}}`;
  735. }
  736. function tomlKey(key: string): string {
  737. return /^[A-Za-z0-9_-]+$/.test(key) ? key : tomlString(key);
  738. }
  739. function tomlValue(value: JsonValue): string {
  740. if (typeof value === 'string') return tomlString(value);
  741. if (typeof value === 'number' || typeof value === 'boolean') return String(value);
  742. if (Array.isArray(value)) {
  743. return `[${value.filter((item) => item !== null).map(tomlValue).join(',')}]`;
  744. }
  745. if (value && typeof value === 'object') {
  746. // TOML inline table 用 `=` 而不是 JSON 的 `:`,不能直接 JSON.stringify
  747. return tomlInlineTable(
  748. Object.fromEntries(
  749. Object.entries(value as Record<string, JsonValue>)
  750. .filter(([, item]) => typeof item === 'string')
  751. .map(([key, item]) => [key, String(item)]),
  752. ),
  753. );
  754. }
  755. return '""';
  756. }
  757. function emptyRuntimeStatus(): RuntimeStatus {
  758. return {
  759. state: 'stopped',
  760. binaryPath: null,
  761. version: null,
  762. codexHome: null,
  763. models: [],
  764. capabilities: emptyCapabilities(),
  765. providerId: null,
  766. defaultModel: null,
  767. degraded: false,
  768. error: null,
  769. };
  770. }
  771. function emptyCapabilities(): RuntimeCapabilities {
  772. return { namespaceTools: false, imageGeneration: false, webSearch: false };
  773. }
  774. function asRecord(value: unknown): Record<string, unknown> | null {
  775. return value && typeof value === 'object' && !Array.isArray(value)
  776. ? (value as Record<string, unknown>)
  777. : null;
  778. }
  779. function readString(value: unknown): string | null {
  780. return typeof value === 'string' ? value : null;
  781. }
  782. /**
  783. * mcpServerStatus/list 的 tools 是「工具名 → 详情」映射,详情结构由 Codex 版本决定。
  784. * 这里只认 name / description,认不出的形态退化成只留 key 名,避免协议一变列表就空掉。
  785. */
  786. export function toToolSummaries(
  787. tools: Record<string, unknown> | null | undefined,
  788. ): Array<{ name: string; description: string }> {
  789. return Object.entries(tools ?? {}).map(([key, value]) => {
  790. const detail = asRecord(value);
  791. return {
  792. name: readString(detail?.name) || key,
  793. description: readString(detail?.description) || '',
  794. };
  795. });
  796. }
  797. function asError(value: unknown): Error {
  798. return value instanceof Error ? value : new Error(String(value));
  799. }
  800. export function sanitizeDiagnostic(value: string): string {
  801. return value
  802. .replace(/sk-[A-Za-z0-9_-]{12,}/gu, 'sk-[redacted]')
  803. .replace(/Bearer\s+[A-Za-z0-9._~-]+/giu, 'Bearer [redacted]')
  804. .slice(0, 4_000);
  805. }
  806. function delay(ms: number): Promise<void> {
  807. return new Promise((resolve) => {
  808. const timer = setTimeout(resolve, ms);
  809. timer.unref();
  810. });
  811. }
  812. /**
  813. * Codex 的致命错误由 anyhow 打成纯文本(`Error: ...`),RUST_LOG 的常规日志是 JSON 行;
  814. * 纯文本里还混着 `WARNING: ...` 这类噪音,所以优先取 Error 行,取不到才退化成第一条纯文本行。
  815. */
  816. function pickFatalLine(tail: string[]): string | null {
  817. const plain = tail
  818. .map((line) => line.trim())
  819. .filter((line) => line && !line.startsWith('{'));
  820. if (!plain.length) return null;
  821. const fatal = plain.find((line) => /^(error|panic)\b/iu.test(line) || line.includes('panicked at'));
  822. return sanitizeDiagnostic((fatal ?? plain[0]).slice(0, 500));
  823. }
  824. interface ExitWatcher {
  825. /** 收 stderr:既实时推诊断,也留尾部若干行用于拼退出原因 */
  826. noteStderr(line: string): void;
  827. /** 统一退出出口(流关闭 / 进程退出 / spawn 失败),只生效一次 */
  828. shutdown(fallback: string): void;
  829. /** 等 stderr 收完,把没有信息量的 fallback 换成带 exit code 与原因原文的错误 */
  830. describe(fallback: string): Promise<Error>;
  831. }
  832. function createExitWatcher(
  833. child: ChildProcessWithoutNullStreams,
  834. onDiagnostic: (line: string) => void,
  835. onExit: (error: Error) => void,
  836. ): ExitWatcher {
  837. const tail: string[] = [];
  838. const closed = new Promise<void>((resolve) => child.once('close', () => resolve()));
  839. let shuttingDown = false;
  840. const describe = async (fallback: string): Promise<Error> => {
  841. // 等不到 close 也得往下走:一次卡住的退出不能把启动的 Promise 挂在那里
  842. await Promise.race([closed, delay(EXIT_GRACE_MS)]);
  843. const code = child.exitCode ?? child.signalCode;
  844. const reason = code === null || code === undefined ? fallback : `Codex App Server 退出(${code})`;
  845. const detail = pickFatalLine(tail);
  846. return new Error(detail ? `${reason}:${detail}` : reason);
  847. };
  848. return {
  849. noteStderr(line: string): void {
  850. tail.push(line);
  851. if (tail.length > STDERR_TAIL_LINES) tail.shift();
  852. onDiagnostic(line);
  853. },
  854. describe,
  855. shutdown(fallback: string): void {
  856. if (shuttingDown) return;
  857. shuttingDown = true;
  858. void describe(fallback).then((error) => {
  859. // 流都关了进程还活着(卡死的子进程)只能自己收
  860. if (child.exitCode === null && child.signalCode === null) child.kill('SIGTERM');
  861. onExit(error);
  862. });
  863. },
  864. };
  865. }