Files
metona-ai-desktop/electron/harness/adapters/shared/sse-stream.ts
T
thzxx 2230bcec3f feat: v0.4.0 四阶段迭代 — 安全加固 + 工程基线 + 架构重构 + 双 Provider 扩展
P0 安全修复:
- API Key 加密存储(safeStorage 密钥链,版本化前缀,历史明文平滑兼容)
- 间接提示注入防护(SecurityScanHook 工具结果深扫描,网络工具脱敏/本地工具警示分级)
- error:report IPC 断链修复(渲染进程错误上报落 electron-log + 审计)
- abort 信号贯通工具层(run_command/dev-tools 子进程随会话中断终止)
- run_command 沙箱加固(cd 系统目录/敏感文件读取拦截 + chcp 前缀剥离防解析退化)
- .env 真实生效(dotenv 回退加载,应用内配置优先)

P1 工程基础:
- ESLint 9 flat config + 全部 34 条存量 warnings 清零(零容忍基线)
- 测试基线 118 用例 11 文件(token/文件防护/权限/沙箱/注入/命令/引擎/注册表/审计链/摘要分层)
- test:electron 双模式(ELECTRON_RUN_AS_NODE 跑 Electron ABI,SQLite 套件全执行)
- SessionRecorder 多会话隔离 + 9 种 TRACE 事件补全(含最终轮 iteration_end)
- Provider 故障转移(重试耗尽/不可重试一次性切换 fallback + 前端通知)
- MCP 真就绪(等待全部连接完成再广播 tools:ready)
- SLO/HealthChecker 真实接入(60s 巡检 + 托盘状态)
- CONFIG_DEFAULTS 单一来源(消除 SEED 双源漂移)

P2 架构升级:
- handlers.ts 1940 行拆分为 13 个 IPC 域模块(防重入注册 + 多窗口广播)
- AgentEngineManager 每会话独立引擎(LRU 30 + adapter 工厂隔离 abort 信号)
- TaskOrchestrator EngineProvider 改造 + abortByParent 联动中断 SubAgent
- 会话摘要分层上下文(session_summaries 滚动摘要 + 截断游标清理防因果污染)
- 消息编辑重发/重新生成(truncateAfter IPC + store 动作 + UI)
- Markdown 导出 / WebSearch 并行抓取(并发 3)/ 记忆 TF 缓存 / 版本构建期注入

P3 能力扩展:
- OpenAI Adapter(o 系列推理模型 reasoning_effort/max_completion_tokens)
- Anthropic Adapter(原生 Messages API:tool_use 块/角色合并/thinking budget/图片 base64/SSE 事件机)
- 设置页/Onboarding 六 Provider 全链路接入
2026-08-20 23:17:02 +08:00

282 lines
9.6 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
* 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. 清空缓冲区
*
* @param toolCallsBuffer - 工具调用缓冲区(index → { name, argsBuffer }
* @param requestId - 请求 ID
* @param sessionId - 会话 ID
* @param iteration - 当前迭代轮次
* @param seqRef - seq 计数器引用(递增)
* @yields MetonaStreamEvent
*/
function* flushToolCallBuffer(
toolCallsBuffer: Map<number, { name: string; argsBuffer: string }>,
requestId: string,
sessionId: string,
iteration: number,
seqRef: { seq: number },
): Generator<MetonaStreamEvent> {
for (const [, buf] of toolCallsBuffer) {
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: buf.argsBuffer ? JSON.parse(buf.argsBuffer) : {},
iteration,
timestamp: Date.now(),
},
};
} catch (err) {
// #25 修复: 不再静默丢弃 JSON 解析失败的工具调用
// 审查修复: 不再 yield ERROR 事件,因为 Engine 收到 ERROR 会 throw 中断整个请求,
// 导致一个好的工具调用 JSON 解析失败就丢弃所有工具调用。
// 改为 log.warn 记录后 continue 跳过这条坏的工具调用,继续处理 buffer 中剩余的。
const rawPreview = buf.argsBuffer?.slice(0, 200) ?? '';
log.warn(`[SSE] Tool call JSON parse failed: ${(err as Error).message}`, rawPreview);
continue;
}
}
toolCallsBuffer.clear();
}
/**
* 解析 OpenAI 兼容 SSE 流式响应
*
* @param responseBody - fetch Response.body (ReadableStream<Uint8Array>)
* @param requestId - 对应的请求 ID
* @param sessionId - 会话 ID
* @param iteration - 当前迭代轮次
* @yields MetonaStreamEvent
*/
export async function* parseSSEStream(
responseBody: ReadableStream<Uint8Array>,
requestId: string,
sessionId: string,
iteration: number,
): AsyncGenerator<MetonaStreamEvent> {
const reader = responseBody.getReader();
const decoder = new TextDecoder();
const seqRef = { seq: 0 };
let buffer = '';
// 工具调用缓冲区:index → { name, argsBuffer }
const toolCallsBuffer = new Map<number, { name: string; argsBuffer: string }>();
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]') {
// 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);
}
} catch (parseErr) {
// P2-8 修复: 不再静默跳过,记录 warning 便于排查 SSE 数据损坏
log.warn(`[SSE] Failed to parse stream line: ${(parseErr as Error).message}`, line.slice(0, 200));
}
}
}
}
/**
* 解析 OpenAI 兼容的非流式 JSON 响应 → MetonaResponse
*/
export function parseOpenAICompatibleResponse(
data: Record<string, unknown>,
): {
content: string;
reasoningContent?: string;
toolCalls?: Array<{ id: string; name: string; args: Record<string, unknown>; iteration: number; timestamp: number }>;
finishReason: string;
usage: MetonaTokenUsage;
} {
const choice = (data.choices as Array<Record<string, unknown>>)?.[0];
const message = choice?.message as Record<string, unknown> | undefined;
const usage = data.usage as Record<string, unknown> | undefined;
const rawToolCalls = message?.tool_calls as Array<Record<string, unknown>> | undefined;
return {
content: (message?.content as string) ?? '',
reasoningContent: message?.reasoning_content as string | undefined,
toolCalls: rawToolCalls?.map((tc) => {
const fn = tc.function as Record<string, unknown>;
let args: Record<string, unknown> = {};
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<string, unknown>;
}
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<string, unknown>)?.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<string, unknown> | 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';
}
}