import { EventEmitter } from 'node:events'; import { createInterface, type Interface } from 'node:readline'; import type { Readable, Writable } from 'node:stream'; /** * 移植自 Noobi.ai src/main/jsonRpcPeer.ts(Codex app-server 的换行分隔 JSONL 通道),逻辑保持不变。 */ export type JsonRpcId = string | number; export interface JsonRpcErrorShape { code: number; message: string; data?: unknown; } export interface JsonRpcServerRequest { id: JsonRpcId; method: string; params?: unknown; } interface PendingRequest { method: string; resolve(value: unknown): void; reject(reason: Error): void; timer: NodeJS.Timeout; } const MAX_OUTBOUND_PROTOCOL_LINE_BYTES = 16 * 1024 * 1024; // 大 diff / 长工具输出会整行走 JSONL,留足余量但不允许无上限分配 const MAX_INBOUND_PROTOCOL_LINE_BYTES = 48 * 1024 * 1024; export class JsonRpcRequestError extends Error { readonly code: number; readonly data: unknown; constructor(method: string, error: JsonRpcErrorShape) { super(`${method}: ${error.message}`); this.name = 'JsonRpcRequestError'; this.code = error.code; this.data = error.data; } } export class JsonRpcPeer extends EventEmitter { readonly #input: Readable; readonly #output: Writable; readonly #pending = new Map(); #reader: Interface | null = null; #nextId = 1; #closed = false; constructor(input: Readable, output: Writable) { super(); this.#input = input; this.#output = output; } start(): void { if (this.#reader) return; this.#reader = createInterface({ input: this.#input, crlfDelay: Infinity }); this.#reader.on('line', (line) => this.#handleLine(line)); this.#reader.on('close', () => this.close(new Error('Codex 协议流已关闭'))); } async request(method: string, params?: unknown, timeoutMs = 30_000): Promise { if (this.#closed) throw new Error('Codex 协议通道已关闭'); const id = this.#nextId++; return new Promise((resolve, reject) => { const timer = setTimeout(() => { this.#pending.delete(id); reject(new Error(`${method} 在 ${timeoutMs}ms 内超时`)); }, timeoutMs); timer.unref(); this.#pending.set(id, { method, resolve: (value) => resolve(value as T), reject, timer, }); try { this.#write({ id, method, ...(params === undefined ? {} : { params }) }); } catch (error) { clearTimeout(timer); this.#pending.delete(id); reject(asError(error)); } }); } notify(method: string, params?: unknown): void { this.#write({ method, ...(params === undefined ? {} : { params }) }); } respond(id: JsonRpcId, result: unknown): void { this.#write({ id, result }); } respondError(id: JsonRpcId, error: JsonRpcErrorShape): void { this.#write({ id, error }); } endOutput(): void { if (!this.#closed) this.#output.end(); } close(reason = new Error('Codex 协议通道已关闭')): void { if (this.#closed) return; this.#closed = true; this.#reader?.close(); this.#reader = null; for (const pending of this.#pending.values()) { clearTimeout(pending.timer); pending.reject(reason); } this.#pending.clear(); this.emit('closed', reason); } #write(message: Record): void { if (this.#closed) throw new Error('Codex 协议通道已关闭'); const line = `${JSON.stringify(message)}\n`; if (Buffer.byteLength(line, 'utf8') > MAX_OUTBOUND_PROTOCOL_LINE_BYTES) { throw new Error('Codex 协议消息超过 16 MiB 上限'); } this.#output.write(line, 'utf8'); } #handleLine(line: string): void { if (!line.trim()) return; if (Buffer.byteLength(line, 'utf8') > MAX_INBOUND_PROTOCOL_LINE_BYTES) { const error = new Error('Codex 返回的单行协议内容超过 48 MiB 上限'); this.emit('protocolError', error); this.close(error); return; } let message: Record; try { const parsed = JSON.parse(line) as unknown; if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { throw new Error('协议消息不是对象'); } message = parsed as Record; } catch (error) { this.emit('protocolError', new Error(`非法的 Codex JSONL: ${asError(error).message}`)); return; } const method = typeof message.method === 'string' ? message.method : null; const id = isJsonRpcId(message.id) ? message.id : null; if (method) { if (id !== null) { this.emit('serverRequest', { id, method, params: message.params } satisfies JsonRpcServerRequest); } else { this.emit('notification', { method, params: message.params }); } return; } if (id === null) { this.emit('protocolError', new Error('Codex 响应缺少 id')); return; } const pending = this.#pending.get(id); if (!pending) { this.emit('protocolError', new Error(`Codex 返回了未知的响应 id: ${String(id)}`)); return; } clearTimeout(pending.timer); this.#pending.delete(id); if (isJsonRpcError(message.error)) { pending.reject(new JsonRpcRequestError(pending.method, message.error)); } else { pending.resolve(message.result); } } } function isJsonRpcId(value: unknown): value is JsonRpcId { return typeof value === 'string' || (typeof value === 'number' && Number.isFinite(value)); } function isJsonRpcError(value: unknown): value is JsonRpcErrorShape { if (!value || typeof value !== 'object' || Array.isArray(value)) return false; const candidate = value as Record; return typeof candidate.code === 'number' && typeof candidate.message === 'string'; } function asError(value: unknown): Error { return value instanceof Error ? value : new Error(String(value)); }