Files
metona-ai-desktop/electron/harness/adapters/shared/sse-stream.ts
T
thzxx 3940716dc2
CI / 类型检查 + Lint + 单元测试 (push) Failing after 5m45s
CI / 全量测试 (Electron ABI) (push) Failing after 5m22s
CI / 产物编译验证 (push) Successful in 10m3s
feat: v0.7.0 四阶段全量迭代 — 修复面收口 · 安全纵深 · 架构还债 · 能力演进
P1 修复面收口: v0.6.3 截断自愈推全量(Anthropic/Ollama/非流式/引擎兜底); SSE 上游错误帧检测进重试通道;
clearMessages 摘要游标根治; truncateResult 内联图片白名单统一; 前端四 bug(确认弹窗锁死/MemoryViewer/
Virtuoso Footer/abort 尾部过滤) + reasoning 缓冲跨迭代污染; 托盘通知过滤与新建会话死链接线

P2 安全纵深: MCP 审批闭环(ConfirmationHook×PolicyEngine 联动+重名拒注册); SSRF 收敛 ssrf-guard 共享模块
(web_fetch 双通道校验+重定向终态复检); Electron 加固(preload CJS 化→sandbox:true/CSP/权限白名单/will-navigate);
run_command cmd.exe 白名单通道元字符守门; diff_viewer 10MB 预检; Anthropic thinking 预算下限; Agnes 思考显式关闭

P3 架构还债: OpenAICompatibleAdapter 中间基类收敛四家样板; 错误分类单轨化(删 mapError/getFetchSignal,
超时显式 ETIMEDOUT); PRAGMA user_version 迁移版本化; 死代码清理专项(cn.ts/SHORTCUTS/ContextMenu 分支/
getWindowState/modifiedArgs/sandbox 空壳); i18next 引入; a11y 第一轮; SearXNG 页批量草稿模型统一

P4 能力演进: Ollama pull 可取消/capabilities 探测/num_ctx 实测缓存; UpdateService feed 比对式自动更新
(app:updateCheck IPC + StatusBar 入口); MiMo providerOptions(web_search 服务端工具/strict JSON);
web_fetch extract_mode=markdown(turndown); network.proxyUrl 全局代理(Chromium sessions+undici dispatcher)

测试: 264 → 507 用例(Electron ABI 全绿零跳过), 覆盖引擎压缩管线/重试竞速/MEMORY.md 闸门/file_editor 五操作/
filesystem 七工具实体夹具/git 真实仓库/SSE 错误帧/全线截断自愈/Provider 请求形态矩阵/SSRF 表测/钩子分级矩阵/
OutputValidator 全量/SLO 指标/MCP 安全纯函数/task_manager 链路/渲染层纯域/i18n 桥契约
2026-08-27 17:06:58 +08:00

562 lines
21 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';
import { ContentFilterError } from '../base-adapter';
/**
* v0.6.4 根治「上游错误帧黑洞」:
*
* OpenAI 兼容网关常在中途发送 `data: {"error":{...}}` 数据帧(网关超时、限流、
* 配额耗尽、鉴权失效等)。旧解析器整条处理链只从 `chunk.choices?.[0]?.delta` 取数,
* 这类帧两个分支都不命中、零日志 —— 结果任何上游错误都伪装成"干净的空回复 +
* 正常 DONE",且异常以普通事件而非异常形式出现,绕过了引擎 chatStreamWithRetry
* 的重试/故障转移通道。
*
* 现契约:上游错误帧在解析层**直接抛出**携带结构化 status/providerCode 的
* SseUpstreamError —— 异常沿 async generator 传播进 chatStreamWithRetry 的 catch
* 使 429/5xx 自动走指数退避重试、401 等不可重试错误走 fallback 故障转移,
* 与 HTTP 状态码路径的行为完全对齐(错误分类单轨化的流式半边)。
*/
export class SseUpstreamError extends Error {
/** 归一化后的 HTTP status(当帧内无数值 status 时按 providerCode 推断) */
readonly status?: number;
/** 上游原始错误码(如 "rate_limit_exceeded" / "insufficient_quota" */
readonly providerCode?: string;
constructor(message: string, options?: { status?: number; providerCode?: string }) {
super(message);
this.name = 'SseUpstreamError';
this.status = options?.status;
this.providerCode = options?.providerCode;
}
}
/** v0.6.4: OpenAI 兼容流帧的最小结构化类型(仅承载本解析器实际消费的字段) */
interface SseStreamFrame {
choices?: Array<{
delta?: {
content?: string;
reasoning_content?: string;
tool_calls?: Array<{
index?: number;
function?: { name?: string; arguments?: string };
}>;
};
finish_reason?: string;
}>;
usage?: {
prompt_tokens?: number;
completion_tokens?: number;
total_tokens?: number;
prompt_cache_hit_tokens?: number;
prompt_cache_miss_tokens?: number;
completion_tokens_details?: { reasoning_tokens?: number };
prompt_tokens_details?: { cached_tokens?: number };
};
}
/**
* 从一条已 JSON.parse 的 SSE 数据帧中提取上游错误信息。
* 兼容三种形态:
* 1. 顶层 `{error:{message,status}}` — OpenAI 兼容网关最常见
* 2. `{choices:[{error:{...}}]}` — 少数代理的变体包装
* 3. `{error:"plain string"}` — 极简实现
* 返回 null 表示该帧不含错误(正常数据帧)。
*/
function extractUpstreamErrorFrame(
chunk: unknown,
): { message: string; status?: number; code?: string } | null {
if (!chunk || typeof chunk !== 'object') return null;
const c = chunk as Record<string, unknown>;
let errObj: unknown = c.error;
if (
(!errObj || typeof errObj !== 'object') &&
Array.isArray(c.choices) &&
c.choices.length > 0
) {
errObj = (c.choices[0] as Record<string, unknown> | undefined)?.error;
}
if (typeof errObj === 'string') {
return errObj.trim() ? { message: errObj } : null;
}
if (!errObj || typeof errObj !== 'object') return null;
const e = errObj as Record<string, unknown>;
const rawStatus =
typeof e.status === 'number'
? e.status
: typeof e.status_code === 'number'
? e.status_code
: undefined;
const rawCode = typeof e.code === 'string' ? e.code : typeof e.type === 'string' ? e.type : '';
const message =
typeof e.message === 'string'
? e.message
: typeof e.msg === 'string'
? e.msg
: JSON.stringify(errObj);
// 无消息且无状态码的无害空对象不算错误(防御性)
if (!message && rawStatus === undefined && !rawCode) return null;
return { message: message || `upstream error (${rawCode || rawStatus})`, status: rawStatus, code: rawCode || undefined };
}
/** 上游字符串错误码 → 归一化 HTTP status(用于帧内缺失数值 status 时仍能驱动重试判定) */
function providerCodeToStatus(code: string): number | undefined {
const c = code.toLowerCase();
if (/rate_limit|too_many/.test(c)) return 429;
if (/quota|insufficient|billing|exceeded_balance/.test(c)) return 402;
if (/invalid_api_key|api_key_invalid|unauthorized|authentication/.test(c)) return 401;
if (/forbidden|permission/.test(c)) return 403;
if (/model_not_found|no_such_model/.test(c)) return 404;
if (/overloaded|capacity|unavailable/.test(c)) return 503;
if (/server_error|internal_error|internal/.test(c)) return 500;
// 请求级 400 家族(无效参数/上下文超限)— 不映射到可重试区间,保持非重试语义
return undefined;
}
/**
* 由提取出的错误信息构造待抛出的 SseUpstreamError /
* ContentFilterErrorcontent_filter 类直接复用既有专用类型)。
*/
function makeUpstreamThrowable(info: { message: string; status?: number; code?: string }): Error {
if (info.code && info.code.toLowerCase().includes('content_filter')) {
return new ContentFilterError(info.message, 'SSE 流中收到上游安全审核错误');
}
const status = info.status ?? (info.code ? providerCodeToStatus(info.code) : undefined);
return new SseUpstreamError(`upstream_error${info.code ? ` (${info.code})` : ''}: ${info.message}`, {
status,
providerCode: info.code,
});
}
/**
* 截断参数自愈载荷的唯一构造点(v0.6.4: 非流式/Ollama NDJSON/Anthropic 全线共用)。
*
* 原 v0.6.3 只覆盖了 OpenAI 共享 SSE 层;此处抽出为共享函数后,所有协议路径的
* 截断工具调用统一转为显式错误参数 —— 工具执行失败 → 错误结果回传模型 → 模型
* 重试/分块写入(ReAct 自愈闭环),彻底消灭"静默丢弃 → 空回复 → 会话无声终止"。
*/
export function truncatedArgumentsPayload(
parseErrorMessage: string,
rawTailSample: string,
): Record<string, unknown> {
return {
_truncatedArguments: true,
_truncatedReason:
'The tool-call arguments JSON was truncated before completion ' +
'(likely max_tokens output limit reached while generating this tool call). ' +
'The original arguments are lost and cannot be recovered. Please retry with a ' +
'smaller output (e.g. write the file in smaller chunks) — do NOT reuse or repeat ' +
'the previous oversized arguments.' +
` [parser: ${parseErrorMessage}; tail sample: ...${rawTailSample}]`,
};
}
/**
* 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<number, { name: string; argsBuffer: string }>,
requestId: string,
sessionId: string,
iteration: number,
seqRef: { seq: number },
): Generator<MetonaStreamEvent> {
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 → v0.6.4: 截断的工具调用转显式错误参数(不丢弃)— 工具执行失败后
// 错误结果回传模型,触发重试/分块写入。载荷构造已收敛到共享的
// truncatedArgumentsPayloadOllama NDJSON / Anthropic 事件机 / 非流式同步复用)。
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: truncatedArgumentsPayload((err as Error).message, rawTail),
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<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 = '';
// v0.6.3: 是否收到过 [DONE](流断开兜底用)
let sawDone = false;
// 工具调用缓冲区: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();
// v0.6.4: 兼容 `data:{...}`(无空格)变体 — 部分代理网关不带空格,
// 原实现的 startsWith('data: ') 会将其整帧跳过
if (!trimmed || !trimmed.startsWith('data:')) continue;
const data = trimmed.slice(5).trim();
if (!data) continue;
// 流结束
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;
}
// v0.6.4: 结构化类型承载帧内容(替代 JSON.parse 的隐式 any,杜绝字段漂移)
let chunk: SseStreamFrame;
try {
chunk = JSON.parse(data) as SseStreamFrame;
} catch (parseErr) {
// P2-8 修复: 不再静默跳过,记录 warning 便于排查 SSE 数据损坏
log.warn(
`[SSE] Failed to parse stream line: ${(parseErr as Error).message}`,
line.slice(0, 200),
);
continue;
}
// v0.6.4: 上游错误帧检测(根治"错误帧黑洞")。解析失败直接 throw
// 异常沿 chatStreamWithRetry 的 catch 走重试/故障转移通道;
// 工具调用缓冲不 flush —— 重试会从零重建整个响应。
const upstreamError = extractUpstreamErrorFrame(chunk);
if (upstreamError) {
log.warn(
`[SSE] Upstream error frame received: code=${upstreamError.code ?? 'n/a'} status=${upstreamError.status ?? 'n/a'} message=${upstreamError.message.slice(0, 300)} — throwing for retry/failover handling`,
);
throw makeUpstreamThrowable(upstreamError);
}
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
const usageRaw = chunk.usage;
if (usageRaw) {
const usage: MetonaTokenUsage = {
inputTokens: usageRaw.prompt_tokens ?? 0,
outputTokens: usageRaw.completion_tokens ?? 0,
totalTokens: usageRaw.total_tokens ?? 0,
reasoningTokens: usageRaw.completion_tokens_details?.reasoning_tokens,
// DeepSeek: prompt_cache_hit_tokens / prompt_cache_miss_tokens
// MiMo: prompt_tokens_details.cached_tokens
cacheHitTokens:
usageRaw.prompt_cache_hit_tokens ?? usageRaw.prompt_tokens_details?.cached_tokens,
cacheMissTokens: usageRaw.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)`,
);
}
// v0.6.4: 流式 content_filter 终止映射 —— 非流式路径早已支持
// throwHttpError → ContentFilterError),流式此前既不映射也不打日志,
// 引擎拿到普通结束、用户看不到拦截原因。抛专用类型使 finish() 映射为
// CONTENT_FILTERED 错误码 + 友好提示,且不会被重试逻辑反复重放。
if (finishReason === 'content_filter') {
log.warn('[SSE] finish_reason=content_filter — provider safety filter terminated the response');
throw new ContentFilterError(
'流式响应被 Provider 安全审核终止',
'SSE stream (finish_reason=content_filter)',
);
}
}
}
// 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<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 (err) {
// v0.6.4: 非流式路径与流式截断自愈对齐 —— 此前坏参静默降级 {} 与流式
// 的显式自愈行为不一致(v0.6.3 只修了流式半边)。空参数会让工具以
// "缺少必要参数"泛化失败,模型无法得知发生了截断;现在统一转为
// _truncatedArguments 错误参数,触发模型分块重试。
const sample = rawArgs.slice(-120);
log.warn(
`[SSE] Non-stream tool call args truncated (unparseable JSON, ${(err as Error).message}). Tail: ...${sample}`,
);
args = truncatedArgumentsPayload((err as Error).message, sample);
}
} 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';
}
}