修复Agent聊天流式错误与等待反馈
完善OpenAI Chat SSE尾包、错误与结束事件处理 事件监听不可用时降级普通回复并继续持久化 补充等待状态、监听拒绝与协议回归测试 将平台LLM测试纳入客户端完整检查
This commit is contained in:
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -21,7 +21,7 @@
|
||||
- 背景:开发用单 Agent 聊天已经能真实调用各 Agent 的 LLM 路由并持久化对话,但 Agent 仍主要表现为同步问答,用户无法明确投递一个任务让某个 Agent 独立运行,也无法同时启动多个 Agent 的工作。
|
||||
- 决策:在现有 `.agent/runtime` 和 `.agent/conversations` 基础上新增单 Agent 后台任务入口。Tauri 命令 `start_game_creator_agent_runtime_task` 立即写入该 Agent 的 runtime state/event/task history,追加用户任务到 `.agent/conversations/agents/<agentId>.jsonl`,随后在 App 进程内启动 tokio task 执行最小 Agent loop:Agent 按轮输出 `thinkingSummary / plan / actions / response`,Runtime 按白名单和项目权限策略执行工具并记录 `action / observation` 事件,再把已有 observation 放回下一轮 prompt,让 Agent 修正计划、继续行动或用空 actions + response 收束;当前后台任务最多执行 3 轮 loop,仍未收束时再按最后计划和全部观察生成最终回复。完成或失败后把 assistant 回复或错误追加回对话,并写入 `.agent/agent.db` 审计记录。工具箱包含只读工具 `memory.read`、`conversation.read`、`asset.list`、`project.index`、`project.diff`、`file.list`、`file.read`、`agent.run_status`,以及受策略保护的写/运行工具 `memory.write`、`file.write`、`command.run_limited`、`blackboard.write`、`agent.message` 和 `agent.delegate`;`memory.write` 可追加或覆盖本 Agent 私有记忆、项目长期/短期记忆或黑板,`file.write` 只能写项目内相对路径,`command.run_limited` 只接受 `game.static_smoke` 并复用本地静态自检安全边界,`blackboard.write` 追加共享黑板,`agent.message` 写目标 Agent 对话,`agent.delegate` 把任务投递到目标 Agent 的独立后台队列;策略拒绝时不执行工具并把 `blocked` observation 回给 Agent;策略要求确认时不执行工具,而是持久化精确待确认动作并暂停该 Agent 队列,待开发者确认或拒绝后在同一 run 续跑。每个 Agent 的任务历史落在 `.agent/runtime/tasks/<agentId>.jsonl`,读 runtime 时按 `runId` 去重返回最近任务,任务视角状态使用 `pending / running / completed / failed`,Runtime state 增加 `nextStep`,UI 在 Runtime 面板和主 Agent 状态卡展示当前任务、动作、下一步与最近任务。不同 Agent 使用独立 `.agent/runtime/locks/<agentId>.lock`,允许并行运行;同一 Agent 已有运行任务时,新任务会先进入该 Agent 的 pending 队列,当前 drain 持锁完成后串行继续下一条 pending。该能力仍不是独立 OS 进程或跨重启离线常驻 worker。
|
||||
- 2026-07-10 补充:后台 Runtime 每次追加 `.agent/runtime/events/<agentId>.jsonl` 后会通过 Tauri `game-creator-agent-runtime-update` 事件广播当前 `AgentRuntimeResult`;开发单 Agent 聊天页、项目内 Agent 对话弹窗和主窗口 Agent 状态列表都只把该事件作为实时 UI 通知并复用前端 runtime 归一化合并,事实源仍是 `.agent/runtime/agents`、`events` 和 `tasks` 文件。
|
||||
- 2026-07-11 补充:开发单 Agent 聊天页保留整页纵向滚动,聊天消息区固定响应式高度并在内部滚动;Runtime 恢复确认区使用独立布局行,避免与 Runtime 详情或聊天内容重叠。Runtime 面板详情可折叠且折叠时不渲染详情 DOM,但状态标题与任务控制按钮继续保留;等待 LLM 时在消息区显示动态状态,连续流式 delta 合并到动画帧更新并跳过重复 Runtime state。OpenAI Chat SSE 的空 `choices` usage 事件不再报错,事件监听不可用或首个文本片段前流式失败时降级普通回复。
|
||||
- 2026-07-11 补充:开发单 Agent 聊天页保留整页纵向滚动,聊天消息区固定响应式高度并在内部滚动;Runtime 恢复确认区使用独立布局行,避免与 Runtime 详情或聊天内容重叠。Runtime 面板详情可折叠且折叠时不渲染详情 DOM,但状态标题与任务控制按钮继续保留;等待 LLM 时在消息区显示动态状态,连续流式 delta 合并到动画帧更新并跳过重复 Runtime state。OpenAI Chat SSE 会收集 usage-only 尾包、保留 finish reason 与上游 error message,收到 `[DONE]` 后立即结束;持久事件订阅失败时显示非致命错误,聊天事件监听不可用或首个文本片段前流式失败时降级普通回复并继续落盘。
|
||||
- 2026-07-11 补充:为缩小单 Agent 与 Codex CLI 在代码任务上的差距,Runtime 工具箱新增 `project.search` 和 `file.patch`,并扩展 `file.read` 的按行分页。`project.search` 在项目内执行有界字面量检索,默认忽略大小写,返回相对路径、行号和匹配行,跳过 `.agent`、敏感配置、依赖和构建目录;权限继承 `file.read`。`file.read` 接受 `startLine / maxLines`,返回带行号的最多 240 行、8,000 字符上下文,允许 Agent 继续分页而不是只看到文件开头约 900 字符。`file.patch` 只做 `oldText -> newText` 精确替换,必须声明预期匹配数,匹配数不符时不写入;它继承 `file.write` 权限,复用项目写锁和 Runtime 动作账本,并追加不含代码正文的 `agent.runtime.file.patch` 审计记录。三者组成“搜索定位 -> 分段读取 -> 局部修改 -> 再次读取验证”的最小代码工作闭环,不开放任意 shell。
|
||||
- 2026-07-11 补充:后台 Agent 的工具规划和最终回复请求对 `LlmError::EmptyResponse` 最多自动重试 3 次(含首次共 4 次请求),与既有 Generator 对上游 HTTP 成功但空 content 的恢复策略一致;连接、协议、鉴权、解析等其他错误不在此处重试,仍按原错误路径失败并落 Runtime 事件。该重试不会重复执行工具动作,只会原样重发尚未得到有效文本的 LLM 请求。
|
||||
- 2026-07-11 调整:后台单 Agent 的 planning loop 上限从 3 轮提升到 6 轮,每轮工具动作上限仍为 3;真实代码任务已证明“搜索定位、分段读取、等待写入确认、写后复读”可能在第 3 轮才进入待确认状态,原上限会让确认后的同 run 没有继续验证余量。`maxLoopIterations` 随新上限写入 Runtime,跨重启待确认动作按已完成轮次继续使用剩余轮次;6 轮后 actions 仍未收束时继续进入 `failed / budget-exhausted`,不会伪装完成。该调整只作用于后台单 Agent Runtime,不改变游戏草案 Generator/Evaluator 的 3 轮上限;本文件和旧实施摘要中“后台 3 轮后整理最终回复”的历史描述由本条取代。
|
||||
|
||||
@@ -52,7 +52,7 @@ Agent Runtime 负责:
|
||||
- 2026-07-10 补充:`recentEvents` 接入前端归一态和 Runtime 状态面板,事件事实源仍是 `.agent/runtime/events/<agentId>.jsonl`;面板按时间展示最近 `thinking_summary / plan / action / observation / response / error` 事件,现在能同时看到 Agent 的计划、最近观察、最近事件、最近工具动作和任务队列。
|
||||
- 2026-07-10 补充:后台 Runtime 每次追加 `.agent/runtime/events/<agentId>.jsonl` 后会通过 Tauri `game-creator-agent-runtime-update` 事件广播当前 `AgentRuntimeResult`,开发单 Agent 聊天页、项目内 Agent 对话弹窗和主窗口 Agent 状态卡用同一套前端归一化逻辑合并状态;该事件只做实时 UI 通知,`.agent/runtime/agents`、`events` 和 `tasks` 仍是重开项目后的事实源。
|
||||
- 2026-07-10 补充:后台 Agent loop 的统一语义事件类型为 `thinking_summary / plan / action / observation / response / error`。普通失败和 loop 预算耗尽都会追加 `error` 事件,并继续保留 `turn.failed / turn.budget_exhausted` 生命周期事件兼容既有读取方;开发窗口、项目内 Agent 对话弹窗和主窗口状态卡通过现有最近事件列表直接展示统一错误事件及其安全详情。状态面板默认保持最新 4 条的紧凑视图,当前后端返回的最近事件超过 4 条时可展开查看全部返回记录,确保同一 run 的六类语义事件不会因 UI 硬截断而无法检查。
|
||||
- 2026-07-11 补充:开发单 Agent 聊天页继续使用整页纵向滚动,不把 Runtime 锁进固定视口;聊天消息区使用固定响应式高度并在内部滚动,避免历史消息持续撑高聊天面板。可选的 Runtime 恢复确认区始终占据独立布局行,不能与 Runtime 详情或聊天消息重叠。Runtime 状态面板支持折叠详情,折叠时只卸载目标、计划、事件、动作和任务等详情 DOM,仍保留状态标题与取消、重试、确认、拒绝、刷新操作;等待 LLM 时在消息区显示连接 / 等待首包 / 接收中的动态状态。流式聊天的连续 delta 通过 `requestAnimationFrame` 合并为每帧最多一次消息更新,delta 不重复提交未变化的 Runtime state;OpenAI Chat SSE 的 `choices: []` usage 事件按非内容事件跳过,事件监听不可用或流式请求在首个文本片段前失败时自动降级普通回复。
|
||||
- 2026-07-11 补充:开发单 Agent 聊天页继续使用整页纵向滚动,不把 Runtime 锁进固定视口;聊天消息区使用固定响应式高度并在内部滚动,避免历史消息持续撑高聊天面板。可选的 Runtime 恢复确认区始终占据独立布局行,不能与 Runtime 详情或聊天消息重叠。Runtime 状态面板支持折叠详情,折叠时只卸载目标、计划、事件、动作和任务等详情 DOM,仍保留状态标题与取消、重试、确认、拒绝、刷新操作;等待 LLM 时在消息区显示连接 / 等待首包 / 接收中的动态状态。流式聊天的连续 delta 通过 `requestAnimationFrame` 合并为每帧最多一次消息更新,delta 不重复提交未变化的 Runtime state;OpenAI Chat SSE 的 usage-only 事件会回填最终 token usage,finish-only 事件会把结束原因送入状态流,上游 error 保留真实消息,`[DONE]` 立即结束读取;持久事件订阅失败时显示非致命 Runtime 错误,聊天事件监听不可用或流式请求在首个文本片段前失败时自动降级普通回复。
|
||||
- 2026-07-11 补充:后台单 Agent 新增 Codex 风格的代码导航与局部编辑闭环。`project.search` 接受 `query / path / maxResults / caseSensitive`,在项目边界内做字面量搜索并返回 `path:line`,最多扫描 500 个、单个不超过 512 KiB 的文本文件,跳过 `.agent`、`.git`、`node_modules`、`dist`、`build`、`target`、`.next`、`coverage` 和 `.env*`;该工具映射到 `file.read` 权限。`file.read` 接受 `startLine / maxLines`,返回带行号的指定片段、总行数和下一页提示,单次最多 240 行、8,000 字符。`file.patch` 接受 `path / oldText / newText / expectedReplacements`,只在实际匹配数与预期一致时持锁写入,目标文件和修改后文件最大 2 MiB,成功后写 `agent.runtime.file.patch` 审计;该工具映射到 `file.write` 权限。Agent planning prompt 明确要求批量修改前创建 checkpoint,并可在修改后再次 `file.read` 验证;本轮不开放任意 shell 命令。
|
||||
- 2026-07-11 补充:后台工具规划与最终回复的 LLM 请求新增空 content 恢复,只对 `LlmError::EmptyResponse` 原样自动重试最多 3 次,其他错误不重试。重试发生在任何工具动作执行前或已有 observation 后的下一次规划请求,因此不会因空响应重复执行已经落盘的工具副作用;重试耗尽后仍写入原有 `error / turn.failed` 事件并把失败消息追加到当前 Agent 会话。
|
||||
- 2026-07-11 调整:后台单 Agent planning loop 上限提升为 6 轮、每轮最多 3 个工具动作,`loopIteration / maxLoopIterations / toolActionBudget` 继续向 UI 和 `agent.run_status` 暴露真实进度。待确认动作在第 N 轮暂停时,确认或重启恢复后从下一轮继续,最多使用剩余的 `6-N` 轮完成修改后复读和自检;6 轮仍返回非空 actions 时终态保持 `failed / budget-exhausted`。这只调整后台单 Agent Runtime;游戏草案 Generator/Evaluator 仍保持独立的 3 轮修复预算。旧摘要中“后台最多 3 轮、耗尽后生成最终回复”的描述不再有效,以本条和 budget-exhausted 决策为准。
|
||||
|
||||
+1
-1
@@ -136,7 +136,7 @@
|
||||
"ai-game-creator-shell:agent-run": "npm --prefix apps/ai-game-creator-shell run agent-run --",
|
||||
"ai-game-creator-shell:agent-run:smoke": "npm --prefix apps/ai-game-creator-shell run agent-run:smoke",
|
||||
"ai-game-creator-shell:typecheck": "npm --prefix apps/ai-game-creator-shell run typecheck",
|
||||
"ai-game-creator-shell:check": "npm run ai-game-creator-shell:typecheck && npm run test -- apps/ai-game-creator-shell/tests && cargo test -p platform-agent --manifest-path server-rs/Cargo.toml game_creation && cargo test -p shared-contracts --manifest-path server-rs/Cargo.toml game_creation_app && cargo test --manifest-path apps/ai-game-creator-shell/src-tauri/Cargo.toml && npm run ai-game-creator-shell:agent-run:smoke",
|
||||
"ai-game-creator-shell:check": "npm run ai-game-creator-shell:typecheck && npm run test -- apps/ai-game-creator-shell/tests && cargo test -p platform-llm --manifest-path server-rs/Cargo.toml && cargo test -p platform-agent --manifest-path server-rs/Cargo.toml game_creation && cargo test -p shared-contracts --manifest-path server-rs/Cargo.toml game_creation_app && cargo test --manifest-path apps/ai-game-creator-shell/src-tauri/Cargo.toml && npm run ai-game-creator-shell:agent-run:smoke",
|
||||
"check:native-shells": "node scripts/check-native-shells.mjs"
|
||||
},
|
||||
"dependencies": {
|
||||
|
||||
@@ -443,12 +443,15 @@ struct OpenAiCompatibleSseParser {
|
||||
buffer: String,
|
||||
raw_text: String,
|
||||
api_kind: LlmApiKind,
|
||||
terminated: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct ParsedStreamEvent {
|
||||
delta_text: Option<String>,
|
||||
finish_reason: Option<String>,
|
||||
usage: Option<LlmTokenUsage>,
|
||||
is_terminal: bool,
|
||||
}
|
||||
|
||||
impl LlmProvider {
|
||||
@@ -935,7 +938,10 @@ impl LlmClient {
|
||||
let mut parser = OpenAiCompatibleSseParser::new(request.api_kind);
|
||||
let mut accumulated_text = String::new();
|
||||
let mut finish_reason = None;
|
||||
let mut usage = None;
|
||||
let mut undecoded_chunk_bytes = Vec::new();
|
||||
let emit_finish_only_delta = request.api_kind == LlmApiKind::OpenAiChat;
|
||||
let mut stream_terminated = false;
|
||||
|
||||
loop {
|
||||
let next_chunk = response.chunk().await.map_err(|error| {
|
||||
@@ -983,26 +989,20 @@ impl LlmClient {
|
||||
);
|
||||
error
|
||||
})?;
|
||||
for event in stream_events {
|
||||
if let Some(delta_text) = event.delta_text
|
||||
&& !delta_text.is_empty()
|
||||
{
|
||||
accumulated_text.push_str(delta_text.as_str());
|
||||
let update = LlmStreamDelta {
|
||||
accumulated_text: accumulated_text.clone(),
|
||||
delta_text,
|
||||
finish_reason: event.finish_reason.clone(),
|
||||
};
|
||||
on_delta(&update);
|
||||
}
|
||||
|
||||
if event.finish_reason.is_some() {
|
||||
finish_reason = event.finish_reason;
|
||||
}
|
||||
stream_terminated = consume_stream_events(
|
||||
stream_events,
|
||||
&mut accumulated_text,
|
||||
&mut finish_reason,
|
||||
&mut usage,
|
||||
emit_finish_only_delta,
|
||||
&mut on_delta,
|
||||
);
|
||||
if stream_terminated {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if !undecoded_chunk_bytes.is_empty() {
|
||||
if !stream_terminated && !undecoded_chunk_bytes.is_empty() {
|
||||
let trailing_text =
|
||||
std_str::from_utf8(undecoded_chunk_bytes.as_slice()).map_err(|error| {
|
||||
log_llm_raw_failure(
|
||||
@@ -1027,53 +1027,37 @@ impl LlmClient {
|
||||
);
|
||||
error
|
||||
})?;
|
||||
for event in trailing_events {
|
||||
if let Some(delta_text) = event.delta_text
|
||||
&& !delta_text.is_empty()
|
||||
{
|
||||
accumulated_text.push_str(delta_text.as_str());
|
||||
let update = LlmStreamDelta {
|
||||
accumulated_text: accumulated_text.clone(),
|
||||
delta_text,
|
||||
finish_reason: event.finish_reason.clone(),
|
||||
};
|
||||
on_delta(&update);
|
||||
}
|
||||
|
||||
if event.finish_reason.is_some() {
|
||||
finish_reason = event.finish_reason;
|
||||
}
|
||||
}
|
||||
stream_terminated = consume_stream_events(
|
||||
trailing_events,
|
||||
&mut accumulated_text,
|
||||
&mut finish_reason,
|
||||
&mut usage,
|
||||
emit_finish_only_delta,
|
||||
&mut on_delta,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
let remaining_events = parser.finish().map_err(|error| {
|
||||
log_llm_raw_failure(
|
||||
&self.config,
|
||||
&request,
|
||||
true,
|
||||
1,
|
||||
"parse_stream_failed",
|
||||
parser.raw_text().as_str(),
|
||||
if !stream_terminated {
|
||||
let remaining_events = parser.finish().map_err(|error| {
|
||||
log_llm_raw_failure(
|
||||
&self.config,
|
||||
&request,
|
||||
true,
|
||||
1,
|
||||
"parse_stream_failed",
|
||||
parser.raw_text().as_str(),
|
||||
);
|
||||
error
|
||||
})?;
|
||||
consume_stream_events(
|
||||
remaining_events,
|
||||
&mut accumulated_text,
|
||||
&mut finish_reason,
|
||||
&mut usage,
|
||||
emit_finish_only_delta,
|
||||
&mut on_delta,
|
||||
);
|
||||
error
|
||||
})?;
|
||||
for event in remaining_events {
|
||||
if let Some(delta_text) = event.delta_text
|
||||
&& !delta_text.is_empty()
|
||||
{
|
||||
accumulated_text.push_str(delta_text.as_str());
|
||||
let update = LlmStreamDelta {
|
||||
accumulated_text: accumulated_text.clone(),
|
||||
delta_text,
|
||||
finish_reason: event.finish_reason.clone(),
|
||||
};
|
||||
on_delta(&update);
|
||||
}
|
||||
|
||||
if event.finish_reason.is_some() {
|
||||
finish_reason = event.finish_reason;
|
||||
}
|
||||
}
|
||||
|
||||
let content = accumulated_text.trim().to_string();
|
||||
@@ -1095,7 +1079,7 @@ impl LlmClient {
|
||||
text: content,
|
||||
finish_reason,
|
||||
response_id,
|
||||
usage: None,
|
||||
usage,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1288,11 +1272,16 @@ impl OpenAiCompatibleSseParser {
|
||||
buffer: String::new(),
|
||||
raw_text: String::new(),
|
||||
api_kind,
|
||||
terminated: false,
|
||||
}
|
||||
}
|
||||
|
||||
fn push_chunk(&mut self, chunk: &str) -> Result<Vec<ParsedStreamEvent>, LlmError> {
|
||||
self.raw_text.push_str(chunk);
|
||||
if self.terminated {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
self.buffer.push_str(chunk);
|
||||
self.buffer = self.buffer.replace("\r\n", "\n");
|
||||
self.drain_complete_events()
|
||||
@@ -1303,7 +1292,7 @@ impl OpenAiCompatibleSseParser {
|
||||
}
|
||||
|
||||
fn finish(&mut self) -> Result<Vec<ParsedStreamEvent>, LlmError> {
|
||||
if self.buffer.trim().is_empty() {
|
||||
if self.terminated || self.buffer.trim().is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
@@ -1319,7 +1308,13 @@ impl OpenAiCompatibleSseParser {
|
||||
self.buffer = self.buffer[(boundary + 2)..].to_string();
|
||||
|
||||
if let Some(event) = parse_sse_event_block(self.api_kind, block.as_str())? {
|
||||
let is_terminal = event.is_terminal;
|
||||
events.push(event);
|
||||
if is_terminal {
|
||||
self.terminated = true;
|
||||
self.buffer.clear();
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1327,6 +1322,62 @@ impl OpenAiCompatibleSseParser {
|
||||
}
|
||||
}
|
||||
|
||||
fn consume_stream_events<F>(
|
||||
events: Vec<ParsedStreamEvent>,
|
||||
accumulated_text: &mut String,
|
||||
finish_reason: &mut Option<String>,
|
||||
usage: &mut Option<LlmTokenUsage>,
|
||||
emit_finish_only_delta: bool,
|
||||
on_delta: &mut F,
|
||||
) -> bool
|
||||
where
|
||||
F: FnMut(&LlmStreamDelta),
|
||||
{
|
||||
for event in events {
|
||||
let ParsedStreamEvent {
|
||||
delta_text,
|
||||
finish_reason: event_finish_reason,
|
||||
usage: event_usage,
|
||||
is_terminal,
|
||||
} = event;
|
||||
|
||||
if let Some(event_usage) = event_usage {
|
||||
*usage = Some(event_usage);
|
||||
}
|
||||
|
||||
let delta_text = delta_text.unwrap_or_default();
|
||||
let has_delta = !delta_text.is_empty();
|
||||
if has_delta {
|
||||
accumulated_text.push_str(delta_text.as_str());
|
||||
}
|
||||
|
||||
if let Some(event_finish_reason) = event_finish_reason {
|
||||
*finish_reason = Some(event_finish_reason.clone());
|
||||
if has_delta || emit_finish_only_delta {
|
||||
let update = LlmStreamDelta {
|
||||
accumulated_text: accumulated_text.clone(),
|
||||
delta_text,
|
||||
finish_reason: Some(event_finish_reason),
|
||||
};
|
||||
on_delta(&update);
|
||||
}
|
||||
} else if has_delta {
|
||||
let update = LlmStreamDelta {
|
||||
accumulated_text: accumulated_text.clone(),
|
||||
delta_text,
|
||||
finish_reason: None,
|
||||
};
|
||||
on_delta(&update);
|
||||
}
|
||||
|
||||
if is_terminal {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
false
|
||||
}
|
||||
|
||||
fn normalize_non_empty(value: String, error_message: &str) -> Result<String, LlmError> {
|
||||
let trimmed = value.trim().to_string();
|
||||
if trimmed.is_empty() {
|
||||
@@ -1813,10 +1864,23 @@ fn parse_sse_event_block(
|
||||
}
|
||||
|
||||
let data = data_lines.join("\n");
|
||||
if data.trim().is_empty() || data.trim() == "[DONE]" {
|
||||
if data.trim().is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
if data.trim() == "[DONE]" {
|
||||
return if api_kind == LlmApiKind::OpenAiChat {
|
||||
Ok(Some(ParsedStreamEvent {
|
||||
delta_text: None,
|
||||
finish_reason: None,
|
||||
usage: None,
|
||||
is_terminal: true,
|
||||
}))
|
||||
} else {
|
||||
Ok(None)
|
||||
};
|
||||
}
|
||||
|
||||
if api_kind == LlmApiKind::OpenAiResponses {
|
||||
return parse_responses_sse_event(data.as_str());
|
||||
}
|
||||
@@ -1825,11 +1889,29 @@ fn parse_sse_event_block(
|
||||
return parse_anthropic_sse_event(data.as_str());
|
||||
}
|
||||
|
||||
let parsed: ChatCompletionsResponseEnvelope = serde_json::from_str(data.as_str())
|
||||
let parsed: serde_json::Value = serde_json::from_str(data.as_str())
|
||||
.map_err(|error| LlmError::Deserialize(format!("解析 LLM SSE 事件失败:{error}")))?;
|
||||
if let Some(error) = parsed.get("error").filter(|error| error.is_object()) {
|
||||
return Err(LlmError::Upstream {
|
||||
status_code: 502,
|
||||
message: error
|
||||
.get("message")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.unwrap_or("LLM Chat SSE 返回失败事件")
|
||||
.to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
let parsed: ChatCompletionsResponseEnvelope = serde_json::from_value(parsed)
|
||||
.map_err(|error| LlmError::Deserialize(format!("解析 LLM SSE 事件失败:{error}")))?;
|
||||
let Some(first_choice) = parsed.choices.first() else {
|
||||
return if parsed.usage.is_some() {
|
||||
Ok(None)
|
||||
return if let Some(usage) = parsed.usage {
|
||||
Ok(Some(ParsedStreamEvent {
|
||||
delta_text: None,
|
||||
finish_reason: None,
|
||||
usage: Some(usage),
|
||||
is_terminal: false,
|
||||
}))
|
||||
} else {
|
||||
Err(LlmError::Deserialize(
|
||||
"LLM SSE 响应缺少 choices[0]".to_string(),
|
||||
@@ -1840,6 +1922,8 @@ fn parse_sse_event_block(
|
||||
Ok(Some(ParsedStreamEvent {
|
||||
delta_text: extract_message_text(first_choice),
|
||||
finish_reason: first_choice.finish_reason.clone(),
|
||||
usage: parsed.usage,
|
||||
is_terminal: false,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -1859,10 +1943,14 @@ fn parse_responses_sse_event(data: &str) -> Result<Option<ParsedStreamEvent>, Ll
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::to_string),
|
||||
finish_reason: None,
|
||||
usage: None,
|
||||
is_terminal: false,
|
||||
})),
|
||||
"response.completed" => Ok(Some(ParsedStreamEvent {
|
||||
delta_text: None,
|
||||
finish_reason: Some("completed".to_string()),
|
||||
usage: None,
|
||||
is_terminal: false,
|
||||
})),
|
||||
"response.failed" | "error" => {
|
||||
let message = parsed
|
||||
@@ -1907,6 +1995,8 @@ fn parse_anthropic_sse_event(data: &str) -> Result<Option<ParsedStreamEvent>, Ll
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::to_string),
|
||||
finish_reason: None,
|
||||
usage: None,
|
||||
is_terminal: false,
|
||||
}))
|
||||
}
|
||||
"message_delta" => Ok(Some(ParsedStreamEvent {
|
||||
@@ -1916,6 +2006,8 @@ fn parse_anthropic_sse_event(data: &str) -> Result<Option<ParsedStreamEvent>, Ll
|
||||
.and_then(|value| value.get("stop_reason"))
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::to_string),
|
||||
usage: None,
|
||||
is_terminal: false,
|
||||
})),
|
||||
// message_stop 只是流终止信号;真正的 stop_reason 已由 message_delta 提供,
|
||||
// 这里不要伪造 finish_reason,否则会覆盖掉 end_turn 等真实值。
|
||||
@@ -2148,9 +2240,11 @@ mod tests {
|
||||
|
||||
assert_eq!(events_a.len(), 1);
|
||||
assert_eq!(events_a[0].delta_text.as_deref(), Some("你"));
|
||||
assert_eq!(events_b.len(), 1);
|
||||
assert_eq!(events_b.len(), 2);
|
||||
assert_eq!(events_b[0].delta_text.as_deref(), Some("好"));
|
||||
assert_eq!(events_b[0].finish_reason.as_deref(), Some("stop"));
|
||||
assert!(!events_b[0].is_terminal);
|
||||
assert!(events_b[1].is_terminal);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -2569,10 +2663,117 @@ mod tests {
|
||||
.await
|
||||
.expect("stream_run should succeed");
|
||||
|
||||
assert_eq!(updates, vec!["你".to_string(), "你好".to_string()]);
|
||||
assert_eq!(
|
||||
updates,
|
||||
vec!["你".to_string(), "你好".to_string(), "你好".to_string()]
|
||||
);
|
||||
assert_eq!(response.text, "你好");
|
||||
assert_eq!(response.finish_reason.as_deref(), Some("stop"));
|
||||
assert_eq!(response.response_id.as_deref(), Some("req_stream_01"));
|
||||
assert_eq!(
|
||||
response.usage,
|
||||
Some(LlmTokenUsage {
|
||||
prompt_tokens: 2,
|
||||
completion_tokens: 2,
|
||||
total_tokens: 4,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_emits_chat_finish_only_delta_without_repeating_text() {
|
||||
let server_url = spawn_mock_server(vec![MockResponse {
|
||||
status_line: "200 OK",
|
||||
content_type: "text/event-stream; charset=utf-8",
|
||||
body: concat!(
|
||||
"data: {\"choices\":[{\"delta\":{\"content\":\"你好\"}}]}\n\n",
|
||||
"data: {\"choices\":[{\"finish_reason\":\"stop\"}]}\n\n",
|
||||
"data: [DONE]\n\n"
|
||||
)
|
||||
.to_string(),
|
||||
extra_headers: Vec::new(),
|
||||
}]);
|
||||
|
||||
let client = build_test_client(server_url, 0);
|
||||
let mut updates = Vec::new();
|
||||
let response = client
|
||||
.stream_run(
|
||||
LlmRunRequest::single_turn("系统", "用户").with_openai_chat(),
|
||||
|delta| updates.push(delta.clone()),
|
||||
)
|
||||
.await
|
||||
.expect("stream_run should succeed");
|
||||
|
||||
assert_eq!(updates.len(), 2);
|
||||
assert_eq!(updates[0].accumulated_text, "你好");
|
||||
assert_eq!(updates[0].delta_text, "你好");
|
||||
assert_eq!(updates[0].finish_reason, None);
|
||||
assert_eq!(updates[1].accumulated_text, "你好");
|
||||
assert_eq!(updates[1].delta_text, "");
|
||||
assert_eq!(updates[1].finish_reason.as_deref(), Some("stop"));
|
||||
assert_eq!(response.text, "你好");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_stops_at_chat_done_without_waiting_for_eof() {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").expect("listener should bind");
|
||||
let address = listener.local_addr().expect("listener should have addr");
|
||||
let (stream_done_sender, stream_done_receiver) = std::sync::mpsc::channel();
|
||||
let server_handle = thread::spawn(move || {
|
||||
let (mut stream, _) = listener.accept().expect("request should connect");
|
||||
read_request(&mut stream);
|
||||
stream
|
||||
.write_all(
|
||||
concat!(
|
||||
"HTTP/1.1 200 OK\r\n",
|
||||
"Content-Type: text/event-stream; charset=utf-8\r\n",
|
||||
"Connection: keep-alive\r\n\r\n",
|
||||
"data: {\"choices\":[{\"delta\":{\"content\":\"你好\"}}]}\n\n",
|
||||
"data: [DONE]\n\n"
|
||||
)
|
||||
.as_bytes(),
|
||||
)
|
||||
.expect("stream response should be written");
|
||||
stream.flush().expect("stream response should flush");
|
||||
stream_done_receiver
|
||||
.recv_timeout(StdDuration::from_secs(1))
|
||||
.expect("stream_run should complete before upstream closes");
|
||||
});
|
||||
|
||||
let client = build_test_client(format!("http://{address}"), 0);
|
||||
let response = tokio::time::timeout(
|
||||
StdDuration::from_secs(1),
|
||||
client.stream_run(
|
||||
LlmRunRequest::single_turn("系统", "用户").with_openai_chat(),
|
||||
|_| {},
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("chat done should finish before upstream closes")
|
||||
.expect("stream_run should succeed");
|
||||
|
||||
assert_eq!(response.text, "你好");
|
||||
stream_done_sender
|
||||
.send(())
|
||||
.expect("server should still keep the stream open");
|
||||
server_handle.join().expect("server thread should join");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn chat_sse_error_preserves_upstream_message() {
|
||||
let error = parse_sse_event_block(
|
||||
LlmApiKind::OpenAiChat,
|
||||
"data: {\"error\":{\"message\":\"上游余额不足\"}}",
|
||||
)
|
||||
.expect_err("chat error event should fail");
|
||||
|
||||
assert_eq!(
|
||||
error,
|
||||
LlmError::Upstream {
|
||||
status_code: 502,
|
||||
message: "上游余额不足".to_string(),
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user