后端修复: Ollama adapter 移除死代码/修复超时硬编码/iteration硬编码/pullModel无超时/流异常断开补发DONE; SSE解析器支持非字符串arguments; Sandbox fail-closed安全加固; IPC移除未使用变量 新增功能: TaskOrchestrator子任务编排器; DelegateTaskTool委派工具; Agent Loop并行工具执行; 上下文自动压缩; 记忆检索注入System Prompt 前端修复: useAgentStream text_delta thought累积bug; 工具调用状态正确流转; 修复重复stateChange事件; TraceStep.thought正确填充
265 lines
8.1 KiB
TypeScript
265 lines
8.1 KiB
TypeScript
/**
|
|
* 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';
|
|
|
|
/**
|
|
* 解析 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();
|
|
let 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]') {
|
|
// 将缓冲区中未完成拼接的工具调用发送
|
|
for (const [index, buf] of toolCallsBuffer) {
|
|
try {
|
|
yield {
|
|
type: MetonaStreamEventType.TOOL_CALL_COMPLETE,
|
|
requestId,
|
|
sessionId,
|
|
iteration,
|
|
seq: 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();
|
|
|
|
yield {
|
|
type: MetonaStreamEventType.DONE,
|
|
requestId,
|
|
sessionId,
|
|
iteration,
|
|
seq: 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: seq++,
|
|
timestamp: Date.now(),
|
|
delta: delta.content,
|
|
};
|
|
}
|
|
|
|
// 推理内容增量(Thinking 模式)
|
|
if (delta?.reasoning_content) {
|
|
yield {
|
|
type: MetonaStreamEventType.REASONING_DELTA,
|
|
requestId,
|
|
sessionId,
|
|
iteration,
|
|
seq: 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: 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: 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') {
|
|
for (const [index, buf] of toolCallsBuffer) {
|
|
try {
|
|
yield {
|
|
type: MetonaStreamEventType.TOOL_CALL_COMPLETE,
|
|
requestId,
|
|
sessionId,
|
|
iteration,
|
|
seq: 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();
|
|
}
|
|
} 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';
|
|
}
|
|
}
|