/** * SSE 流式解析工具 * * 解析 OpenAI 兼容的 Server-Sent Events (SSE) 流式响应, * 产出 MetonaStreamEvent。DeepSeek、Agnes AI 和 MiMo 共享此工具。 * * SSE 格式:data: {json}\n\n * 结束标记:data: [DONE] */ import { nanoid } from 'nanoid'; import log from 'electron-log'; import type { MetonaStreamEvent, MetonaTokenUsage } from '../../types'; import { MetonaStreamEventType } from '../../types'; /** * L-4 修复: 提取 flushToolCallBuffer 辅助函数,消除 [DONE] 分支和 finish_reason='tool_calls' 分支的重复代码 * * 遍历工具调用缓冲区,对每个缓冲的工具调用: * 1. JSON.parse argsBuffer * 2. yield 一个 TOOL_CALL_COMPLETE 事件 * 3. 清空缓冲区 * * v0.6.3 会话停止修复: argsBuffer 解析失败(流截断致 JSON 半截 — 典型场景: * 模型写大文件时输出 token 达上限 finish_reason=length)时,不再静默丢弃该 * 工具调用。丢弃会让引擎看到"零工具调用 + 零文本"→ 误判为模型已完成 → * COMPLETED + 空回复 → 会话无声停止(main.log 20:31/20:32 两次实锤)。 * 现转为 yield 一个携带截断错误说明的 tool call:工具执行将因参数缺失失败, * 错误结果回传模型 → 模型感知截断后重试/分块写入(ReAct 自愈路径)。 * 无限循环由引擎死循环检测器兜底。 * * @param toolCallsBuffer - 工具调用缓冲区(index → { name, argsBuffer }) * @param requestId - 请求 ID * @param sessionId - 会话 ID * @param iteration - 当前迭代轮次 * @param seqRef - seq 计数器引用(递增) * @yields MetonaStreamEvent */ function* flushToolCallBuffer( toolCallsBuffer: Map, requestId: string, sessionId: string, iteration: number, seqRef: { seq: number }, ): Generator { for (const [, buf] of toolCallsBuffer) { if (buf.argsBuffer) { try { yield { type: MetonaStreamEventType.TOOL_CALL_COMPLETE, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), toolCall: { id: `tc_${nanoid(8)}`, name: buf.name, args: JSON.parse(buf.argsBuffer), iteration, timestamp: Date.now(), }, }; } catch (err) { // v0.6.3: 截断的工具调用转显式错误参数(不丢弃)— 工具执行失败后 // 错误结果回传模型,触发重试/分块写入,替代"静默丢弃→空回复终止会话" const rawTail = buf.argsBuffer.slice(-120); log.warn( `[SSE] Tool call args truncated (unparseable JSON, ${(err as Error).message}). ` + `Forwarding as error to model for self-healing. Tail: ...${rawTail}`, ); yield { type: MetonaStreamEventType.TOOL_CALL_COMPLETE, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), toolCall: { id: `tc_${nanoid(8)}`, name: buf.name, args: { _truncatedArguments: true, _truncatedReason: 'The streamed arguments JSON was truncated before completion ' + '(likely max_tokens output limit reached while generating this tool call). ' + 'The original arguments are lost. Please retry with smaller output ' + '(e.g. write the file in smaller chunks) — do NOT reuse the previous oversized arguments.', }, iteration, timestamp: Date.now(), }, }; } } else { // 空 argsBuffer:模型发了空 arguments(合法 — 无参工具) yield { type: MetonaStreamEventType.TOOL_CALL_COMPLETE, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), toolCall: { id: `tc_${nanoid(8)}`, name: buf.name, args: {}, iteration, timestamp: Date.now(), }, }; } } toolCallsBuffer.clear(); } /** * 解析 OpenAI 兼容 SSE 流式响应 * * @param responseBody - fetch Response.body (ReadableStream) * @param requestId - 对应的请求 ID * @param sessionId - 会话 ID * @param iteration - 当前迭代轮次 * @yields MetonaStreamEvent */ export async function* parseSSEStream( responseBody: ReadableStream, requestId: string, sessionId: string, iteration: number, ): AsyncGenerator { const reader = responseBody.getReader(); const decoder = new TextDecoder(); const seqRef = { seq: 0 }; let buffer = ''; // v0.6.3: 是否收到过 [DONE](流断开兜底用) let sawDone = false; // 工具调用缓冲区:index → { name, argsBuffer } const toolCallsBuffer = new Map(); while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); buffer = lines.pop() ?? ''; for (const line of lines) { const trimmed = line.trim(); if (!trimmed || !trimmed.startsWith('data: ')) continue; const data = trimmed.slice(6); // 流结束 if (data === '[DONE]') { sawDone = true; // L-4 修复: 使用 flushToolCallBuffer 替代重复的遍历代码 yield* flushToolCallBuffer(toolCallsBuffer, requestId, sessionId, iteration, seqRef); yield { type: MetonaStreamEventType.DONE, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), }; return; } try { const chunk = JSON.parse(data); const delta = chunk.choices?.[0]?.delta; // 文本内容增量 if (delta?.content) { yield { type: MetonaStreamEventType.TEXT_DELTA, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), delta: delta.content, }; } // 推理内容增量(Thinking 模式) if (delta?.reasoning_content) { yield { type: MetonaStreamEventType.REASONING_DELTA, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), delta: delta.reasoning_content, }; } // 工具调用增量 — 缓冲拼接 if (delta?.tool_calls) { for (const tc of delta.tool_calls) { const idx = tc.index ?? 0; if (!toolCallsBuffer.has(idx)) { toolCallsBuffer.set(idx, { name: tc.function?.name ?? '', argsBuffer: '' }); } const buf = toolCallsBuffer.get(idx)!; if (tc.function?.name) buf.name = tc.function.name; if (tc.function?.arguments) buf.argsBuffer += tc.function.arguments; yield { type: MetonaStreamEventType.TOOL_CALL_DELTA, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), toolCallDelta: { index: idx, name: tc.function?.name, argsDelta: tc.function?.arguments, }, }; } } // Token 使用统计 / finish_reason if (chunk.usage) { const usage: MetonaTokenUsage = { inputTokens: chunk.usage.prompt_tokens ?? 0, outputTokens: chunk.usage.completion_tokens ?? 0, totalTokens: chunk.usage.total_tokens ?? 0, reasoningTokens: chunk.usage.completion_tokens_details?.reasoning_tokens, // DeepSeek: prompt_cache_hit_tokens / prompt_cache_miss_tokens // MiMo: prompt_tokens_details.cached_tokens cacheHitTokens: chunk.usage.prompt_cache_hit_tokens ?? chunk.usage.prompt_tokens_details?.cached_tokens, cacheMissTokens: chunk.usage.prompt_cache_miss_tokens, }; yield { type: MetonaStreamEventType.USAGE, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), usage, }; } // 非 [DONE] 但 finish_reason 为 tool_calls 时提前 flush 缓冲区 const finishReason = chunk.choices?.[0]?.finish_reason as string | undefined; if (finishReason === 'tool_calls') { // L-4 修复: 使用 flushToolCallBuffer 替代重复的遍历代码 yield* flushToolCallBuffer(toolCallsBuffer, requestId, sessionId, iteration, seqRef); } // v0.6.3 归因: 输出 token 上限截断(长工具参数/长文本的常见根因)显式落日志 if (finishReason === 'length') { log.warn( `[SSE] finish_reason=length — output truncated by max_tokens limit ` + `(accumulated argsBuffer: ${[...toolCallsBuffer.values()].reduce((n, b) => n + b.argsBuffer.length, 0)} chars, ` + `model may retry with smaller output)`, ); } } catch (parseErr) { // P2-8 修复: 不再静默跳过,记录 warning 便于排查 SSE 数据损坏 log.warn( `[SSE] Failed to parse stream line: ${(parseErr as Error).message}`, line.slice(0, 200), ); } } } // v0.6.3 流断开兜底: read() done 但从未收到 [DONE](连接中断/服务端异常收尾)。 // 原实现直接结束 generator —— 工具缓冲不 flush、DONE 事件缺失(引擎侧等待 // 流收尾的路径行为未定义,且缓冲的工具调用整体丢失)。补 flush + DONE, // 截断的参数由 flushToolCallBuffer 转为错误结果回传模型自愈。 if (!sawDone) { log.warn( '[SSE] Stream ended without [DONE] marker — flushing buffers (connection likely dropped)', ); yield* flushToolCallBuffer(toolCallsBuffer, requestId, sessionId, iteration, seqRef); yield { type: MetonaStreamEventType.DONE, requestId, sessionId, iteration, seq: seqRef.seq++, timestamp: Date.now(), }; } } /** * 解析 OpenAI 兼容的非流式 JSON 响应 → MetonaResponse */ export function parseOpenAICompatibleResponse(data: Record): { content: string; reasoningContent?: string; toolCalls?: Array<{ id: string; name: string; args: Record; iteration: number; timestamp: number; }>; finishReason: string; usage: MetonaTokenUsage; } { const choice = (data.choices as Array>)?.[0]; const message = choice?.message as Record | undefined; const usage = data.usage as Record | undefined; const rawToolCalls = message?.tool_calls as Array> | undefined; return { content: (message?.content as string) ?? '', reasoningContent: message?.reasoning_content as string | undefined, toolCalls: rawToolCalls?.map((tc) => { const fn = tc.function as Record; let args: Record = {}; const rawArgs = fn?.arguments; if (typeof rawArgs === 'string') { try { args = JSON.parse(rawArgs); } catch { args = {}; } } else if (rawArgs && typeof rawArgs === 'object') { args = rawArgs as Record; } return { id: tc.id as string, name: fn.name as string, args, iteration: 0, timestamp: Date.now(), }; }), finishReason: mapOpenAIFinishReason(choice?.finish_reason as string), usage: { inputTokens: (usage?.prompt_tokens as number) ?? 0, outputTokens: (usage?.completion_tokens as number) ?? 0, totalTokens: (usage?.total_tokens as number) ?? 0, reasoningTokens: (usage?.completion_tokens_details as Record) ?.reasoning_tokens as number | undefined, // DeepSeek: prompt_cache_hit_tokens / MiMo: prompt_tokens_details.cached_tokens cacheHitTokens: (usage?.prompt_cache_hit_tokens as number | undefined) ?? ((usage?.prompt_tokens_details as Record | undefined)?.cached_tokens as | number | undefined), cacheMissTokens: usage?.prompt_cache_miss_tokens as number | undefined, }, }; } function mapOpenAIFinishReason(reason: string): string { switch (reason) { case 'stop': return 'stop'; case 'length': return 'length'; case 'tool_calls': return 'tool_calls'; case 'content_filter': return 'content_filter'; // MiMo 特有:检测到复读截断 case 'repetition_truncation': return 'stop'; default: return 'stop'; } }