import { PassThrough } from 'node:stream'; import { describe, expect, it, vi } from 'vitest'; import { JsonRpcPeer, JsonRpcRequestError } from './jsonRpcPeer'; /** 移植自 Noobi.ai src/main/jsonRpcPeer.test.ts:纯内存流,不依赖 codex 二进制与网络 */ function makePeer() { const input = new PassThrough(); const output = new PassThrough(); output.setEncoding('utf8'); const peer = new JsonRpcPeer(input, output); peer.start(); return { input, output, peer }; } describe('JsonRpcPeer', () => { it('把响应匹配回对应的 pending 请求', async () => { const { input, output, peer } = makePeer(); const written = new Promise((resolve) => output.once('data', resolve)); const pending = peer.request<{ models: unknown[] }>('model/list'); const request = JSON.parse(await written) as { id: number; method: string }; expect(request.method).toBe('model/list'); input.write(`${JSON.stringify({ id: request.id, result: { models: [] } })}\n`); await expect(pending).resolves.toEqual({ models: [] }); peer.close(); }); it('区分通知与服务端请求,且能交错处理', async () => { const { input, peer } = makePeer(); const notification = vi.fn(); const serverRequest = vi.fn(); peer.on('notification', notification); peer.on('serverRequest', serverRequest); input.write('{"method":"turn/started","params":{"turnId":"turn-1"}}\n'); input.write('{"id":"approval-1","method":"item/commandExecution/requestApproval","params":{}}\n'); await new Promise((resolve) => setImmediate(resolve)); expect(notification).toHaveBeenCalledWith({ method: 'turn/started', params: { turnId: 'turn-1' }, }); expect(serverRequest).toHaveBeenCalledWith({ id: 'approval-1', method: 'item/commandExecution/requestApproval', params: {}, }); peer.close(); }); it('把 App Server 返回的协议错误抛给调用方', async () => { const { input, output, peer } = makePeer(); const written = new Promise((resolve) => output.once('data', resolve)); const pending = peer.request('thread/start'); const request = JSON.parse(await written) as { id: number }; input.write(`${JSON.stringify({ id: request.id, error: { code: -32602, message: 'invalid params' }, })}\n`); await expect(pending).rejects.toBeInstanceOf(JsonRpcRequestError); peer.close(); }); it('遇到非法 JSON 只报错、不中断流', async () => { const { input, peer } = makePeer(); const protocolError = vi.fn(); const notification = vi.fn(); peer.on('protocolError', protocolError); peer.on('notification', notification); input.write('{not-json}\n'); input.write('{"method":"warning","params":{"message":"still alive"}}\n'); await new Promise((resolve) => setImmediate(resolve)); expect(protocolError).toHaveBeenCalledOnce(); expect(notification).toHaveBeenCalledOnce(); peer.close(); }); it('关闭后拒绝全部 pending 请求并广播 closed', async () => { const { peer } = makePeer(); const closed = vi.fn(); peer.on('closed', closed); const pending = peer.request('model/list'); peer.close(new Error('子进程退出')); await expect(pending).rejects.toThrow('子进程退出'); expect(closed).toHaveBeenCalledOnce(); await expect(peer.request('model/list')).rejects.toThrow('已关闭'); }); });