/** * 流式上游错误帧 + 全线截断自愈测试(v0.6.4) * * 背景(v0.6.3 审计遗留): * 1. 错误帧黑洞 —— OpenAI 兼容网关中途发送的 `{"error":{...}}` 数据帧被解析器 * 整帧吞掉(零日志),任何上游错误都伪装成"干净的空回复 + 正常 DONE", * 且以普通事件而非异常出现,绕过引擎的重试/故障转移通道。 * 2. 截断自愈只修了 OpenAI 共享层 —— v0.6.3 的 _truncatedArguments 修复未覆盖: * - Anthropic:content_block_stop 解析失败静默 args={};断流时未完成块整体蒸发 * - Ollama:NDJSON 坏参抛错落入外层 catch,工具调用丢弃且同 chunk USAGE/DONE 被跳过 * - 非流式 parseOpenAICompatibleResponse:坏参仍静默 {} * - 引擎兜底缓冲 finalizeToolCallsFromBuffer:坏参静默 {} * * 本文件锁定以下契约: * A. 上游错误帧 → 抛出携带归一化 status 的 SseUpstreamError(可驱动重试判定) * B. finish_reason=content_filter → ContentFilterError(终态、不重试) * C. data:{无空格} 变体正常解析 * D. 非流式/Ollama/Anthropic 截断参数统一转 _truncatedArguments 自愈载荷 * E. Ollama 坏参不再吞掉同 chunk 的 done/USAGE 处理 * F. Anthropic 断流时未完成 tool_use 块 flush 为自愈调用 + DONE */ import { describe, it, expect, vi } from 'vitest'; vi.mock('electron-log', () => ({ default: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, })); import { parseSSEStream, parseOpenAICompatibleResponse, SseUpstreamError } from '../shared/sse-stream'; import { ContentFilterError } from '../base-adapter'; import { OllamaAdapter } from '../ollama.adapter'; import { AnthropicAdapter } from '../anthropic.adapter'; import { MetonaErrorCode, MetonaStreamEventType } from '../../types'; import type { MetonaRequest } from '../../types'; const encoder = new TextEncoder(); function makeStream(lines: string[]): ReadableStream { const payload = lines.join('\n') + '\n'; return new ReadableStream({ start(controller) { controller.enqueue(encoder.encode(payload)); controller.close(); }, }); } async function collectExpectingThrow(stream: ReadableStream): Promise { try { for await (const _ev of parseSSEStream(stream, 'r_test', 's_test', 1)) { void _ev; } } catch (err) { return err; } throw new Error('expected parseSSEStream to throw but it completed normally'); } function sseData(json: unknown): string { return `data: ${JSON.stringify(json)}`; } // ===== A. 上游错误帧 → 抛出结构化异常 ===== describe('parseSSEStream — 上游错误帧(v0.6.4 错误帧黑洞根治)', () => { it('顶层 error 帧(含数值 status)→ 抛出携带该 status 的 SseUpstreamError', async () => { const err = await collectExpectingThrow( makeStream([sseData({ error: { message: 'Gateway timeout', status: 504 } })]), ); expect(err).toBeInstanceOf(SseUpstreamError); expect((err as SseUpstreamError).status).toBe(504); expect((err as Error).message).toContain('Gateway timeout'); }); it('choices[0].error 变体包装也能检出', async () => { const err = await collectExpectingThrow( makeStream([ sseData({ choices: [{ error: { message: 'bad gateway', code: 'upstream_failure', status: 502 } }], }), ]), ); expect(err).toBeInstanceOf(SseUpstreamError); expect((err as SseUpstreamError).status).toBe(502); }); it('字符串型顶层 error 也能检出', async () => { const err = await collectExpectingThrow(makeStream(['data: {"error":"service unavailable"}'])); expect(err).toBeInstanceOf(SseUpstreamError); expect((err as Error).message).toContain('service unavailable'); }); it('providerCode 归一化:rate_limit_exceeded 无数值 status → 映射 429(可重试)', async () => { const err = await collectExpectingThrow( makeStream([ sseData({ error: { code: 'rate_limit_exceeded', message: 'too many requests' } }), ]), ); expect(err).toBeInstanceOf(SseUpstreamError); expect((err as SseUpstreamError).status).toBe(429); expect((err as SseUpstreamError).providerCode).toBe('rate_limit_exceeded'); }); it('insufficient_quota → 402;invalid_api_key → 401(不可重试区间)', async () => { const e1 = await collectExpectingThrow( makeStream([sseData({ error: { code: 'insufficient_quota', message: 'quota exceeded' } })]), ); expect((e1 as SseUpstreamError).status).toBe(402); const e2 = await collectExpectingThrow( makeStream([ sseData({ error: { code: 'invalid_api_key', message: 'Incorrect API key provided' } }), ]), ); expect((e2 as SseUpstreamError).status).toBe(401); }); it('含 content_filter 码的错误帧 → ContentFilterError(复用专用类型)', async () => { const err = await collectExpectingThrow( makeStream([ sseData({ error: { code: 'content_filter', message: 'rejected by safety policy' } }), ]), ); expect(err).toBeInstanceOf(ContentFilterError); }); it('isRetryable 契约对齐:错误带 status 时,engine.isRetryableError 的 429/5xx 判定可直接命中', async () => { // 用与引擎 isRetryableError 相同的判定逻辑验证字段形态 const isRetryableShape = (err: unknown): boolean => { const e = err as { status?: number; message?: string }; if (e.status === 429) return true; if (e.status && e.status >= 500 && e.status < 600) return true; return false; }; const rateLimited = await collectExpectingThrow( makeStream([sseData({ error: { code: 'rate_limit_exceeded', message: 'rl' } })]), ); expect(isRetryableShape(rateLimited)).toBe(true); const authFail = await collectExpectingThrow( makeStream([sseData({ error: { code: 'invalid_api_key', message: 'auth' } })]), ); expect(isRetryableShape(authFail)).toBe(false); }); it('正常数据帧不含 error 字段时不受影响(回归)', async () => { // choices[0] 中存在 delta 但无 error → 正常产出文本增量并 DONE 收尾 const events: string[] = []; const stream = makeStream([ sseData({ choices: [{ delta: { content: 'hello' } }] }), 'data: [DONE]', ]); for await (const ev of parseSSEStream(stream, 'r', 's', 1)) { events.push(ev.type); } expect(events).toContain(MetonaStreamEventType.TEXT_DELTA); expect(events[events.length - 1]).toBe(MetonaStreamEventType.DONE); }); }); // ===== B/C. content_filter 终止映射 + 无空格 data 变体 ===== describe('parseSSEStream — content_filter 与行格式兼容', () => { it('finish_reason=content_filter → 抛出 ContentFilterError(不再当普通结束)', async () => { const err = await collectExpectingThrow( makeStream([sseData({ choices: [{ delta: {}, finish_reason: 'content_filter' }] })]), ); expect(err).toBeInstanceOf(ContentFilterError); }); it('data:{}(无空格)变体被正常解析(此前整帧跳过)', async () => { const events: Array<{ type: string }> = []; const stream = new ReadableStream({ start(controller) { controller.enqueue(encoder.encode('data:{"choices":[{"delta":{"content":"hi"}}]}\n')); controller.enqueue(encoder.encode('data:[DONE]\n')); controller.close(); }, }); for await (const ev of parseSSEStream(stream, 'r', 's', 1)) { events.push({ type: ev.type }); } expect(events.some((e) => e.type === MetonaStreamEventType.TEXT_DELTA)).toBe(true); expect(events[events.length - 1].type).toBe(MetonaStreamEventType.DONE); }); }); // ===== D. 非流式截断自愈同步 ===== describe('parseOpenAICompatibleResponse — 非流式截断自愈(v0.6.4 同步)', () => { it('坏 JSON arguments 不再静默 {},转为 _truncatedArguments 载荷', () => { const result = parseOpenAICompatibleResponse({ choices: [ { message: { role: 'assistant', content: null, tool_calls: [ { id: 'call_1', function: { name: 'write_file', arguments: '{"file_path": "a.html", "con' }, }, ], }, finish_reason: 'tool_calls', }, ], usage: { prompt_tokens: 10, completion_tokens: 5, total_tokens: 15 }, }); expect(result.toolCalls).toHaveLength(1); const args = result.toolCalls![0].args as Record; expect(args._truncatedArguments).toBe(true); expect(String(args._truncatedReason)).toContain('truncated'); }); it('合法对象型 arguments 保持原样(回归)', () => { const result = parseOpenAICompatibleResponse({ choices: [ { message: { role: 'assistant', tool_calls: [{ id: 'c1', function: { name: 'think', arguments: '{"a":1}' } }], }, finish_reason: 'tool_calls', }, ], }); expect(result.toolCalls![0].args).toEqual({ a: 1 }); }); }); // ===== E/F. Ollama NDJSON 与 Anthropic 事件机 ===== /** 构造全局 fetch mock:返回给定行的 NDJSON/SSE 流 */ function mockFetchWithLines(lines: string[]): ReturnType { const payload = encoder.encode(lines.join('\n') + '\n'); const body = new ReadableStream({ start(controller) { controller.enqueue(payload); controller.close(); }, }); const fetchMock = vi.fn().mockResolvedValue( new Response(body, { status: 200, headers: { 'Content-Type': 'application/x-ndjson' }, }), ); vi.stubGlobal('fetch', fetchMock); return fetchMock; } const baseRequest: MetonaRequest = { meta: { sessionId: 's1', iteration: 1, requestId: 'r1', timestamp: Date.now(), agentVersion: 'test' }, systemPrompt: { roleDefinition: 'rd', outputConstraints: '', safetyGuidelines: '' }, messages: [{ role: 'user', content: 'hi', timestamp: Date.now() }], params: { maxTokens: 4096, temperature: 0, stream: true }, }; describe('OllamaAdapter.sendStream — NDJSON 截断自愈(v0.6.4)', () => { it('坏 JSON arguments → _truncatedArguments 工具调用,且后续 done chunk 的 USAGE/DONE 不再被吞掉', async () => { const adapter = new OllamaAdapter({ provider: 'ollama', baseURL: 'http://localhost:11434', defaultModel: 'qwen3', }); mockFetchWithLines([ JSON.stringify({ model: 'qwen3', message: { role: 'assistant', content: '', tool_calls: [{ function: { name: 'write_file', arguments: '{"path": "a.txt", "cont' } }], }, }), // 关键:同一响应流中随后仍有收尾 chunk(原实现外层 catch 会跳过这些处理) JSON.stringify({ model: 'qwen3', message: { role: 'assistant', content: '' }, done: true, done_reason: 'stop', prompt_eval_count: 11, eval_count: 7, }), ]); const events: string[] = []; let usageInputTokens = -1; for await (const ev of adapter.sendStream(baseRequest)) { events.push(ev.type); if (ev.type === MetonaStreamEventType.USAGE) usageInputTokens = ev.usage!.inputTokens ?? 0; } // 流不再被坏参打断:usage 与 done 都到达 expect(usageInputTokens).toBe(11); expect(events[events.length - 1]).toBe(MetonaStreamEventType.DONE); }); it('自愈载荷内容正确(_truncatedArguments=true + reason 含 truncated)', async () => { const adapter = new OllamaAdapter({ provider: 'ollama', baseURL: 'http://localhost:11434', defaultModel: 'qwen3', }); mockFetchWithLines([ JSON.stringify({ model: 'm', message: { role: 'assistant', tool_calls: [{ function: { name: 'read_file', arguments: '{"file_path": "b.t' } }], }, done: false, }), JSON.stringify({ model: 'm', message: { role: 'assistant', content: '' }, done: true }), ]); let completeArgs: Record | undefined; for await (const ev of adapter.sendStream(baseRequest)) { if (ev.type === MetonaStreamEventType.TOOL_CALL_COMPLETE) { completeArgs = ev.toolCall!.args as Record; } } expect(completeArgs).toBeDefined(); expect(completeArgs!._truncatedArguments).toBe(true); expect(String(completeArgs!._truncatedReason)).toContain('truncated'); }); }); describe('AnthropicAdapter.sendStream — 事件机截断自愈 + 断流 flush(v0.6.4)', () => { it('缺口 A:content_block_stop 时坏 JSON → _truncatedArguments(不再静默 {})', async () => { const adapter = new AnthropicAdapter({ provider: 'anthropic', baseURL: 'http://anthropic.test', apiKey: 'sk-test', defaultModel: 'claude-sonnet-4-5', }); mockFetchWithLines([ 'event: content_block_start', sseData({ type: 'content_block_start', index: 0, content_block: { type: 'tool_use', id: 'toolu_1', name: 'write_file' }, }), 'event: content_block_delta', sseData({ type: 'content_block_delta', index: 0, delta: { type: 'input_json_delta', partial_json: '{"file_path": "a.html", "con' }, }), 'event: content_block_stop', sseData({ type: 'content_block_stop', index: 0 }), 'event: message_stop', sseData({ type: 'message_stop' }), ]); let completeArgs: Record | undefined; let completeId: string | undefined; for await (const ev of adapter.sendStream(baseRequest)) { if (ev.type === MetonaStreamEventType.TOOL_CALL_COMPLETE) { completeArgs = ev.toolCall!.args as Record; completeId = ev.toolCall!.id; } } expect(completeArgs).toBeDefined(); expect(completeArgs!._truncatedArguments).toBe(true); // 保留上游原始 block id(非 nanoid 重造) expect(completeId).toBe('toolu_1'); }); it('缺口 B:断流未完成 tool_use 块 → flush 为自愈调用 + 补发 DONE(不再整体蒸发)', async () => { const adapter = new AnthropicAdapter({ provider: 'anthropic', baseURL: 'http://anthropic.test', apiKey: 'sk-test', defaultModel: 'claude-sonnet-4-5', }); // 有 content_block_start,但流在 content_block_stop/message_stop 之前断开 const payload = 'event: message_start\ndata: {"type":"message_start","message":{"role":"assistant","usage":{"input_tokens":42}}}\n\n' + 'event: content_block_start\ndata: {"type":"content_block_start","index":0,"content_block":{"type":"tool_use","id":"toolu_X","name":"edit_file"}}\n\n' + 'event: content_block_delta\ndata: {"type":"content_block_delta","index":0,"delta":{"type":"input_json_delta","partial_json":"{\\"pa"}}\n\n'; const body = new ReadableStream({ start(controller) { controller.enqueue(encoder.encode(payload)); controller.close(); }, }); vi.stubGlobal( 'fetch', vi.fn().mockResolvedValue(new Response(body, { status: 200 })), ); const events: Array<{ type: string; toolCallId?: string; toolCallName?: string }> = []; for await (const ev of adapter.sendStream(baseRequest)) { events.push({ type: ev.type, toolCallId: ev.toolCall?.id, toolCallName: ev.toolCall?.name, }); } const complete = events.find((e) => e.type === MetonaStreamEventType.TOOL_CALL_COMPLETE); // 核心契约:断流前缓冲中的 block 必须以 TOOL_CALL_COMPLETE 产出(引擎才不会误判空回复完成) expect(complete).toBeDefined(); expect(complete!.toolCallId).toBe('toolu_X'); expect(complete!.toolCallName).toBe('edit_file'); expect(events[events.length - 1].type).toBe(MetonaStreamEventType.DONE); }); it('error 事件 → 抛出携带归一化 status 的异常(overloaded → 529 可重试语义)', async () => { const adapter = new AnthropicAdapter({ provider: 'anthropic', baseURL: 'http://anthropic.test', apiKey: 'sk-test', defaultModel: 'claude-sonnet-4-5', }); mockFetchWithLines([ 'event: error', sseData({ type: 'error', error: { type: 'overloaded_error', message: 'Overloaded' } }), ]); let caught: unknown; try { for await (const _ev of adapter.sendStream(baseRequest)) { void _ev; } } catch (err) { caught = err; } expect(caught).toBeDefined(); expect((caught as Error & { status?: number }).status).toBe(529); expect((caught as Error).message).toContain('Overloaded'); }); }); // ===== G. 引擎侧 ERROR 事件保留结构化码 ===== describe('MetonaErrorCode — CONTENT_FILTERED 枚举契约(finish 映射依赖)', () => { it('code 值稳定为 content_filtered', () => { expect(MetonaErrorCode.CONTENT_FILTERED).toBe('content_filtered'); }); }); // ===== H. 引擎集成:错误帧异常进入重试/故障转移通道;ERROR 码映射 CONTENT_FILTERED ===== import { AgentLoopEngine } from '../../agent-loop/engine'; import { TerminationReason } from '../../agent-loop/types'; import type { IMetonaProviderAdapter, MetonaResponse, MetonaStreamEvent } from '../../types'; function textDone(text: string): MetonaStreamEvent[] { return [ { type: MetonaStreamEventType.TEXT_DELTA, requestId: 'r1', sessionId: 's1', iteration: 1, seq: 0, timestamp: Date.now(), delta: text }, { type: MetonaStreamEventType.DONE, requestId: 'r1', sessionId: 's1', iteration: 1, seq: 1, timestamp: Date.now() }, ]; } describe('AgentLoopEngine 集成 — v0.6.4 错误通道单轨化', () => { const userMessage = { role: 'user' as const, content: 'hi', timestamp: Date.now() }; const systemPrompt = { roleDefinition: '', outputConstraints: '', safetyGuidelines: '' }; function scriptedAdapter( behaviors: Array<{ throws?: Error; events?: MetonaStreamEvent[] }>, ): { adapter: IMetonaProviderAdapter; calls: () => number } { let call = 0; const base: IMetonaProviderAdapter = { providerId: 'mock', supportedModels: ['m'], supportsToolCalling: true, supportsThinking: false, getContextWindow: () => 1_000_000, send: async (): Promise => ({ meta: { requestId: 'r', provider: 'mock', model: 'm', latencyMs: 0, timestamp: Date.now() }, content: '', usage: { inputTokens: 1, outputTokens: 1, totalTokens: 2 }, finishReason: 'stop' as never, }), sendStream: async function* (): AsyncIterable { const b = behaviors[Math.min(call, behaviors.length - 1)]; call++; if (b.throws) throw b.throws; for (const ev of b.events ?? []) yield ev; }, setAbortSignal: vi.fn(), healthCheck: async () => true, }; return { adapter: base, calls: () => call }; } it('SseUpstreamError(429) 首次失败 → 引擎指数退避重试后成功(不再落入 UNKNOWN 终态)', async () => { vi.useFakeTimers(); try { const { adapter } = scriptedAdapter([ { throws: new SseUpstreamError('rate limited', { status: 429 }) }, { events: textDone('recovered answer') }, ]); const engine = new AgentLoopEngine({ retryCount: 3 }, adapter); // 推进退避定时器(1s/2s/4s + jitter 上限) const runPromise = engine.runStream(userMessage, 's1', [], systemPrompt); await vi.advanceTimersByTimeAsync(10_000); const output = await runPromise; expect(output.terminationReason).toBe(TerminationReason.COMPLETED); expect(output.finalAnswer).toBe('recovered answer'); } finally { vi.useRealTimers(); } }); it('引擎收到的流内 ERROR 带 code=content_filtered → finish 发出 CONTENT_FILTERED 错误事件', async () => { const { adapter } = scriptedAdapter([ { events: [ { type: MetonaStreamEventType.ERROR, requestId: 'r1', sessionId: 's1', iteration: 1, seq: 0, timestamp: Date.now(), error: { code: MetonaErrorCode.CONTENT_FILTERED, message: '内容被安全审核拦截', retryable: false, }, }, { type: MetonaStreamEventType.DONE, requestId: 'r1', sessionId: 's1', iteration: 1, seq: 1, timestamp: Date.now() }, ], }, ]); const engine = new AgentLoopEngine({ retryCount: 0 }, adapter); const errorEvents: Array<{ code?: string; message?: string }> = []; engine.on('streamEvent', (ev: MetonaStreamEvent) => { if (ev.type === MetonaStreamEventType.ERROR) { errorEvents.push({ code: ev.error?.code, message: ev.error?.message }); } }); const output = await engine.runStream(userMessage, 's1', [], systemPrompt); expect(output.terminationReason).toBe(TerminationReason.ERROR); expect(errorEvents[errorEvents.length - 1]?.code).toBe(MetonaErrorCode.CONTENT_FILTERED); }); });