| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192 |
- 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<JsonRpcId, PendingRequest>();
- #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<T>(method: string, params?: unknown, timeoutMs = 30_000): Promise<T> {
- if (this.#closed) throw new Error('Codex 协议通道已关闭');
- const id = this.#nextId++;
- return new Promise<T>((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<string, unknown>): 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<string, unknown>;
- try {
- const parsed = JSON.parse(line) as unknown;
- if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) {
- throw new Error('协议消息不是对象');
- }
- message = parsed as Record<string, unknown>;
- } 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<string, unknown>;
- return typeof candidate.code === 'number' && typeof candidate.message === 'string';
- }
- function asError(value: unknown): Error {
- return value instanceof Error ? value : new Error(String(value));
- }
|