jsonRpcPeer.ts 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192
  1. import { EventEmitter } from 'node:events';
  2. import { createInterface, type Interface } from 'node:readline';
  3. import type { Readable, Writable } from 'node:stream';
  4. /**
  5. * 移植自 Noobi.ai src/main/jsonRpcPeer.ts(Codex app-server 的换行分隔 JSONL 通道),逻辑保持不变。
  6. */
  7. export type JsonRpcId = string | number;
  8. export interface JsonRpcErrorShape {
  9. code: number;
  10. message: string;
  11. data?: unknown;
  12. }
  13. export interface JsonRpcServerRequest {
  14. id: JsonRpcId;
  15. method: string;
  16. params?: unknown;
  17. }
  18. interface PendingRequest {
  19. method: string;
  20. resolve(value: unknown): void;
  21. reject(reason: Error): void;
  22. timer: NodeJS.Timeout;
  23. }
  24. const MAX_OUTBOUND_PROTOCOL_LINE_BYTES = 16 * 1024 * 1024;
  25. // 大 diff / 长工具输出会整行走 JSONL,留足余量但不允许无上限分配
  26. const MAX_INBOUND_PROTOCOL_LINE_BYTES = 48 * 1024 * 1024;
  27. export class JsonRpcRequestError extends Error {
  28. readonly code: number;
  29. readonly data: unknown;
  30. constructor(method: string, error: JsonRpcErrorShape) {
  31. super(`${method}: ${error.message}`);
  32. this.name = 'JsonRpcRequestError';
  33. this.code = error.code;
  34. this.data = error.data;
  35. }
  36. }
  37. export class JsonRpcPeer extends EventEmitter {
  38. readonly #input: Readable;
  39. readonly #output: Writable;
  40. readonly #pending = new Map<JsonRpcId, PendingRequest>();
  41. #reader: Interface | null = null;
  42. #nextId = 1;
  43. #closed = false;
  44. constructor(input: Readable, output: Writable) {
  45. super();
  46. this.#input = input;
  47. this.#output = output;
  48. }
  49. start(): void {
  50. if (this.#reader) return;
  51. this.#reader = createInterface({ input: this.#input, crlfDelay: Infinity });
  52. this.#reader.on('line', (line) => this.#handleLine(line));
  53. this.#reader.on('close', () => this.close(new Error('Codex 协议流已关闭')));
  54. }
  55. async request<T>(method: string, params?: unknown, timeoutMs = 30_000): Promise<T> {
  56. if (this.#closed) throw new Error('Codex 协议通道已关闭');
  57. const id = this.#nextId++;
  58. return new Promise<T>((resolve, reject) => {
  59. const timer = setTimeout(() => {
  60. this.#pending.delete(id);
  61. reject(new Error(`${method} 在 ${timeoutMs}ms 内超时`));
  62. }, timeoutMs);
  63. timer.unref();
  64. this.#pending.set(id, {
  65. method,
  66. resolve: (value) => resolve(value as T),
  67. reject,
  68. timer,
  69. });
  70. try {
  71. this.#write({ id, method, ...(params === undefined ? {} : { params }) });
  72. } catch (error) {
  73. clearTimeout(timer);
  74. this.#pending.delete(id);
  75. reject(asError(error));
  76. }
  77. });
  78. }
  79. notify(method: string, params?: unknown): void {
  80. this.#write({ method, ...(params === undefined ? {} : { params }) });
  81. }
  82. respond(id: JsonRpcId, result: unknown): void {
  83. this.#write({ id, result });
  84. }
  85. respondError(id: JsonRpcId, error: JsonRpcErrorShape): void {
  86. this.#write({ id, error });
  87. }
  88. endOutput(): void {
  89. if (!this.#closed) this.#output.end();
  90. }
  91. close(reason = new Error('Codex 协议通道已关闭')): void {
  92. if (this.#closed) return;
  93. this.#closed = true;
  94. this.#reader?.close();
  95. this.#reader = null;
  96. for (const pending of this.#pending.values()) {
  97. clearTimeout(pending.timer);
  98. pending.reject(reason);
  99. }
  100. this.#pending.clear();
  101. this.emit('closed', reason);
  102. }
  103. #write(message: Record<string, unknown>): void {
  104. if (this.#closed) throw new Error('Codex 协议通道已关闭');
  105. const line = `${JSON.stringify(message)}\n`;
  106. if (Buffer.byteLength(line, 'utf8') > MAX_OUTBOUND_PROTOCOL_LINE_BYTES) {
  107. throw new Error('Codex 协议消息超过 16 MiB 上限');
  108. }
  109. this.#output.write(line, 'utf8');
  110. }
  111. #handleLine(line: string): void {
  112. if (!line.trim()) return;
  113. if (Buffer.byteLength(line, 'utf8') > MAX_INBOUND_PROTOCOL_LINE_BYTES) {
  114. const error = new Error('Codex 返回的单行协议内容超过 48 MiB 上限');
  115. this.emit('protocolError', error);
  116. this.close(error);
  117. return;
  118. }
  119. let message: Record<string, unknown>;
  120. try {
  121. const parsed = JSON.parse(line) as unknown;
  122. if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) {
  123. throw new Error('协议消息不是对象');
  124. }
  125. message = parsed as Record<string, unknown>;
  126. } catch (error) {
  127. this.emit('protocolError', new Error(`非法的 Codex JSONL: ${asError(error).message}`));
  128. return;
  129. }
  130. const method = typeof message.method === 'string' ? message.method : null;
  131. const id = isJsonRpcId(message.id) ? message.id : null;
  132. if (method) {
  133. if (id !== null) {
  134. this.emit('serverRequest', { id, method, params: message.params } satisfies JsonRpcServerRequest);
  135. } else {
  136. this.emit('notification', { method, params: message.params });
  137. }
  138. return;
  139. }
  140. if (id === null) {
  141. this.emit('protocolError', new Error('Codex 响应缺少 id'));
  142. return;
  143. }
  144. const pending = this.#pending.get(id);
  145. if (!pending) {
  146. this.emit('protocolError', new Error(`Codex 返回了未知的响应 id: ${String(id)}`));
  147. return;
  148. }
  149. clearTimeout(pending.timer);
  150. this.#pending.delete(id);
  151. if (isJsonRpcError(message.error)) {
  152. pending.reject(new JsonRpcRequestError(pending.method, message.error));
  153. } else {
  154. pending.resolve(message.result);
  155. }
  156. }
  157. }
  158. function isJsonRpcId(value: unknown): value is JsonRpcId {
  159. return typeof value === 'string' || (typeof value === 'number' && Number.isFinite(value));
  160. }
  161. function isJsonRpcError(value: unknown): value is JsonRpcErrorShape {
  162. if (!value || typeof value !== 'object' || Array.isArray(value)) return false;
  163. const candidate = value as Record<string, unknown>;
  164. return typeof candidate.code === 'number' && typeof candidate.message === 'string';
  165. }
  166. function asError(value: unknown): Error {
  167. return value instanceof Error ? value : new Error(String(value));
  168. }