Files
metona-ai-desktop/electron/harness/adapters/shared/sse-stream.ts
T
thzxx e4d81d8247 feat: 升级至 v0.3.1 — 全量代码审计修复 + 安全增强
本次升级基于完整代码审查,修复 Critical/High/Medium/Low 四级共 96 项问题,
并通过返工审计修复 10 项遗留问题,tsc 双端类型检查零错误。

Critical (10/10 完成):
- C-4: command.ts 接入 shell-quote 进行 token-level 注入检测,替代原有正则匹配
  可防御 r"m" -rf /、$'rm'、$(echo rm) 等字符串拼接绕过

High (11/11 完成):
- 竞态保护、Promise.allSettled、AbortController 资源泄漏、IPC 参数校验等

Medium (55/55 完成):
- 事务保护、敏感数据脱敏、枚举校验、MUI v9 Stack prop 迁移、
  React 组件 cancelled 标志、类型收窄等

Low (20/20 完成):
- 辅助方法提取(flushToolCallBuffer/scoreAndPushMemory/tryAddColumn 等)
- nanoid 统一替代 Date.now()+Math.random()
- confirm() 替换为 MUI Dialog、useMemo 缓存、魔法数字命名化等

返工审计修复 (10/10 完成):
- L-11: LogsSettings 残留的原生 confirm()/alert() 全部替换为 MUI Dialog/Alert
- M-53: MemoryViewer handleSearch 独立 ref,修复 searching 状态卡死
- M-42: 脱敏短值(length <= 4)泄露修复
- M-47: tasks:update 补全 title/description 类型校验
- L-9: ollama.adapter 非流式路径 nanoid 统一
- M-45: audit:query limit 策略与 memory:listAll 一致化
- SettingsModal handleConfirmRemove 补全 try/catch + loadServers cleanup
- L-15: CommandPalette useMemo 补全 sessions 响应式依赖
- useAgentStream 事件类型补全 seq/timestamp 字段

新增依赖: shell-quote + @types/shell-quote
版本号: 0.3.0 -> 0.3.1
2026-07-13 22:36:58 +08:00

270 lines
8.4 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 共享此工具。
*
* SSE 格式:data: {json}\n\n
* 结束标记:data: [DONE]
*/
import { nanoid } from 'nanoid';
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 {
// JSON 解析失败,跳过该工具调用
}
}
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,
cacheHitTokens: chunk.usage.prompt_cache_hit_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 {
// 跳过解析失败的行
}
}
}
}
/**
* 解析 OpenAI 兼容的非流式 JSON 响应 → MetonaResponse
*/
export function parseOpenAICompatibleResponse(
data: Record<string, unknown>,
requestId: string,
provider: string,
defaultModel: string,
): {
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,
cacheHitTokens: usage?.prompt_cache_hit_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';
default: return 'stop';
}
}