diff --git a/src-tauri/src/commands/ai/agentic/mod.rs b/src-tauri/src/commands/ai/agentic/mod.rs index 0b0e0dc..3601e1b 100644 --- a/src-tauri/src/commands/ai/agentic/mod.rs +++ b/src-tauri/src/commands/ai/agentic/mod.rs @@ -545,6 +545,363 @@ async fn stream_with_retry( } } +// ============================================================ +// stream_round: 单轮流式调用(candidate 链 + stream_one_provider + MidStream 保文) +// ============================================================ + +/// `run_agentic_loop` for 循环体「单轮流式调用」自包含块的返回结果。 +/// +/// 抽自 L1505-L1756(retry_deadline → candidate_chain 构建/降级 → candidate 循环 +/// (stream_one_provider + match 三态) → epoch 后校验 → MidStream 保文)。 +/// 表达本轮的 4 种结束路径,供 loop 主框架据以 `return` 或继续后续处理。 +/// +/// **不在本函数范围**(loop 内其它段): +/// - 停止/收敛/审批等待/熔断(L1815-L1917):这些在 process_tool_calls 之后, +/// 属「工具执行后路径」,与流式调用本身解耦,留 loop 主体。 +/// - 正常完成(loop 外 L1957+):loop 退出后由 run_agentic_loop 直收尾。 +enum RoundOutcome { + /// 本轮 stream 成功(Complete,无 incomplete)。携带本轮输出供 loop 后续 + /// push_assistant_message / save / 工具执行处理。 + /// + /// 注:reasoning_content 不在此 variant 携带——stream_round 内已通过 + /// `*last_reasoning_content = rc` 写回 caller 的 loop 级 mut 变量,loop 后续 + /// push_assistant_message 读 last_reasoning_content 即可(零行为变更)。 + Continue { + full_text: String, + tool_calls: std::collections::HashMap, + round_usage: df_ai::provider::TokenUsage, + }, + /// 流返回后 stale(epoch 已变,被新 loop 接管):guard 已 disarm, + /// loop 应直接 return(不再 push/save/执行工具)。 + StaleAfterStream, + /// MidStream 保文路径已完整收尾(tokens.add + push partial + 系统提示 + + /// finish_round_exit save/reset/emit)→ loop 直接 return。 + MidStreamSaved, + /// Fatal(候选 Fatal 或全候选耗尽):guard.reset + epoch 检查 + emit AiError + + /// save 已完成 → loop 直接 return。 + Fatal, +} + +/// 单轮流式调用:抽自 `run_agentic_loop` for 循环体的「单轮流式」自包含块。 +/// +/// 范围 = retry_deadline 计算 → candidate_chain 构建(G4.1 防抖降级) +/// → candidate 循环(try_acquire_for_provider + stream_one_provider + match 三态) +/// → epoch 后 stale 检查 → MidStream 保文路径。返回 [`RoundOutcome`] 表达本轮结果。 +/// +/// **逐字节等价**:原 L1505-L1756 内联块的 save/emit/guard.reset/permit 顺序、 +/// token 累加、push_assistant_message 调用、finish_round_exit 退出序列全部保留。 +/// 抽取仅改控制流(原各路径 `return` → `return RoundOutcome::Xxx`,loop 据返回值决定 return)。 +/// +/// **借用处理**: +/// - `tokens: &mut TokenAccumulator` / `saved_token_snapshot: &mut TokenUsage` / +/// `last_round_estimated: &mut bool`:loop 级累计状态,&mut 借用,函数内修改 caller 可见。 +/// - `resolved_model: &mut String` / `provider_config: &mut AiProviderRecord` / +/// `last_reasoning_content: &mut Option`:loop 级 mut 变量(被 candidate +/// 切换覆盖),&mut 借用,函数内赋新值后 loop 后续 push/save 用最新。 +/// - `guard: &mut GeneratingGuard`:RAII 复位,&mut 借用,函数内 guard.reset()/disarm() +/// 直接作用于 caller 的 guard。 +/// - 其它(session_arc/db/app_handle/conv_id/...):& 引用。 +/// +/// **permit**:per-provider permit 仍在 candidate 循环内取/释放(随作用域 Drop), +/// 与原实现一致;per_conv permit 由 loop 入口整 loop 持有,本函数不涉及。 +#[allow(clippy::too_many_arguments)] +async fn stream_round( + session_arc: &Arc>, + db: &Arc, + conv_id: &str, + messages: &[ChatMessage], + tool_defs: &[df_ai::provider::ToolDefinition], + app_handle: &AppHandle, + stop_flag: &std::sync::atomic::AtomicBool, + notify: &tokio::sync::Notify, + max_retries: usize, + model_override: &Option, + agentic_req: &TaskRequirements, + estimated_prompt: u32, + // &mut loop 级状态(函数内修改,caller 可见) + last_round_estimated: &mut bool, + last_reasoning_content: &mut Option, + resolved_model: &mut String, + provider_config: &mut AiProviderRecord, + provider_saturation: &mut HashMap, + tokens: &mut TokenAccumulator, + saved_token_snapshot: &mut df_ai::provider::TokenUsage, + guard: &mut GeneratingGuard, + // & 引用 + candidates: &[AiProviderRecord], + llm_concurrency: &LlmConcurrency, + pinned_goals_snapshot: &[super::GoalEntry], + loop_epoch_arc: &std::sync::Arc, + my_epoch: u64, +) -> RoundOutcome { + // 重试总预算(挂钟,含 sleep + 各次请求耗时),对齐 retry::MAX_TOTAL_BUDGET 30s。 + // 本轮各 candidate 共享一个 30s 预算(切换 provider 不重置预算, + // 防多 provider 串行重试累加超过单轮总预算)。超预算直接放弃重试交最终错误/保文路径。 + let retry_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); + + // 外层 for candidate 包裹 stream_one_provider。 + // 顺序:[primary, ...candidates](空池兜底时仅 primary)。 + // - Success → 用结果覆盖迭代级 resolved_model/provider_config(供 push/save/标题); + // 成功即 break。 + // - InitFailedExhausted(retryable 耗尽)→ continue 切下一 candidate。 + // - Fatal(4xx 非429/鉴权)→ stream_llm 已 emit AiError,guard.reset + return + // 放弃整轮(对齐 retry.rs Fatal + provider_pool 文档)。 + // 全 candidate 耗尽 → 沿用单 provider 失败语义(guard.reset + return,AiError 已在 + // 各 candidate 的 stream_llm 内 emit 最后一条)。 + // 候选链(拥有 Vec):[primary, ...candidates]。预先 clone primary 入链, + // 避免借用 provider_config(切换成功后需写回 provider_config = candidate.clone())。 + // candidates 空(单 provider)→ candidate_chain 仅 primary,循环跑一次,耗尽即 return。 + let mut candidate_chain: Vec = Vec::with_capacity(1 + candidates.len()); + candidate_chain.push(provider_config.clone()); + candidate_chain.extend(candidates.iter().cloned()); + + // 防瞬时抖动(G4.1):连续 PROVIDER_SATURATION_DEMOTE_THRESHOLD 次饱和才把主 provider + // 降级到链尾。计数器在候选循环内维护(Exhausted +1 / Acquired 清零),跨迭代累计。 + // 链首(当前主 provider)连续 N 次 Exhausted → 移到链尾让 fallback 先试; + // 恢复后 Acquired 清零 → 下轮回到链首。单次/瞬时饱和(N 次内)不降级,防跳过主 provider。 + // first():链构建总是 push 至少 1 个(provider_config),但防御性取首防空链 panic。 + if let Some(streak) = candidate_chain + .first() + .and_then(|c| provider_saturation.get(&c.id)) + .copied() + .filter(|s| *s >= PROVIDER_SATURATION_DEMOTE_THRESHOLD) + { + let saturated = candidate_chain.remove(0); + candidate_chain.push(saturated.clone()); + tracing::info!( + conv_id = %conv_id, + provider = %saturated.name, + streak = streak, + threshold = PROVIDER_SATURATION_DEMOTE_THRESHOLD, + "[ai] 主 provider 连续 {} 次 per-provider 饱和,降级到候选链尾(防抖动跳过)", + streak, + ); + } + + let (full_text, tool_calls_acc, round_usage, incomplete, round_reasoning_content) = { + // outcome 累积成功结果(Complete/Partial),InitFailed 不写入。 + let mut outcome: Option<(String, std::collections::HashMap, df_ai::provider::TokenUsage, bool, Option)> = None; + // 追踪最后一个 Exhausted candidate 的诊断文本,供「全 candidate 耗尽」时 emit 最终 AiError。 + // Fatal 分支即时 emit(终态不切候选),Exhausted 分支仅记录 error 不 emit(可能切下一 candidate 成功, + // emit 会留残留气泡)。 + let mut last_exhausted_error: Option = None; + + 'candidate: for candidate in &candidate_chain { + // per-provider permit(非阻塞三态,G4.1): + // NotConfigured(未配置,无限流)→ 无 permit 直接 proceed(单 provider 零变化); + // Acquired(permit)→ 持 permit stream,切换 candidate 时随作用域 Drop 释放; + // Exhausted(占满)→ 连续饱和计数 +1,跳下一 candidate(降级,不阻塞)。 + // 旧 acquire_for_provider(阻塞 acquire_owned)在主 provider 信号量占满时整条 + // fallback 链卡死(降级失效),已由本三态非阻塞版本取代。 + let _provider_permit = match llm_concurrency.try_acquire_for_provider(&candidate.id).await { + ProviderAcquire::NotConfigured => None, + ProviderAcquire::Acquired(permit) => { + // 成功拿到 permit:该 provider 已恢复,连续饱和计数清零(跨迭代防抖)。 + provider_saturation.insert(candidate.id.clone(), 0); + Some(permit) + } + ProviderAcquire::Exhausted => { + // 占满:连续饱和计数 +1;本次跳下一 candidate(非阻塞,不等待)。 + // 跨迭代防抖:连续 N 次饱和后链首降级(见候选链构建处),防瞬时抖动跳过主 provider。 + *provider_saturation.entry(candidate.id.clone()).or_insert(0) += 1; + tracing::debug!( + conv_id = %conv_id, + provider = %candidate.name, + streak = provider_saturation.get(&candidate.id).copied().unwrap_or(0), + "[ai] 候选 {} per-provider 信号量占满,跳过(非阻塞,防单 provider 429)", + candidate.name, + ); + continue 'candidate; + } + }; + + let stream_outcome = stream_one_provider( + candidate, + messages, + tool_defs, + app_handle, + stop_flag, + notify, + conv_id, + max_retries, + retry_deadline, + model_override, + agentic_req, + last_reasoning_content, + ).await; + + match stream_outcome { + StreamOutcome::Success { text, tool_calls, usage, incomplete, resolved_model: m, reasoning_content: rc } => { + // 成功:更新迭代级 resolved_model + provider_config(供后续 push/save/标题)。 + // mut 解构(已在顶部声明 mut):本轮后续 save/push 用「实际成功所用 provider」 + // 而非主 candidate。compress 已在本轮 stream 之前用过 primary,不受影响。 + *resolved_model = m; + *provider_config = candidate.clone(); + *last_reasoning_content = rc.clone(); + outcome = Some((text, tool_calls, usage, incomplete, last_reasoning_content.clone())); + break 'candidate; + } + StreamOutcome::InitFailedExhausted { error } => { + // 本 candidate 重试耗尽(retryable)。切下一 candidate 继续尝试。 + // candidates 空(单 provider)时此即「耗尽 return」语义—— + // 跳出循环后 outcome 仍 None,落入下方「全耗尽 return」兜底。 + // 不在此 emit AiError(可能切下一 candidate 成功留残留气泡), + // 仅记录 error 供全耗尽兜底 emit 最后一条。 + tracing::warn!( + conv_id = %conv_id, + provider = %candidate.name, + "[ai] 候选 {} 流前失败重试耗尽,切换下一 provider", + candidate.name, + ); + last_exhausted_error = Some(error); + continue 'candidate; + } + StreamOutcome::Fatal { error } => { + // Fatal(4xx 非429/鉴权/参数错):立即放弃整轮 fallback。 + // 统一 emit 最终错误气泡(单气泡聚合)。 + guard.reset().await; + // F1:旧 loop(被新 loop 接管)不 emit 错误气泡(owner loop 负责呈现)。 + if loop_epoch_arc.load(Ordering::SeqCst) != my_epoch { + return RoundOutcome::Fatal; + } + emit_fatal_error(app_handle, conv_id, &error).await; + // G1.1(2026-08-05):Fatal 退出也落库 user 消息(镜像下方 Exhausted 分支), + // 否则重启/切会话从 DB 恢复时末条 user 丢失(与 Exhausted 同根因不对称)。 + // epoch 校验通过后、return 前补 save(owner loop 才写,防 owner 已转移仍写)。 + save_conversation(session_arc, db, conv_id, None, Some(&*resolved_model), true).await; + return RoundOutcome::Fatal; + } + } + } + + match outcome { + Some(r) => r, + None => { + // 全 candidate 耗尽(含单 provider 场景):沿用原单 provider 失败语义。 + tracing::warn!( + conv_id = %conv_id, + candidates_tried = candidate_chain.len(), + "[ai] 全 provider 候选流前失败重试耗尽,放弃本轮", + ); + guard.reset().await; + if loop_epoch_arc.load(Ordering::SeqCst) != my_epoch { + return RoundOutcome::Fatal; + } + let err_msg = last_exhausted_error + .unwrap_or_else(|| "AI 调用失败:所有候选 provider 重试耗尽".to_string()); + emit_fatal_error(app_handle, conv_id, &err_msg).await; + // 全 provider 失败也需落库 user 消息(已 push 内存), + // 否则切换/重载从 DB 恢复时末条 user 丢失(实测压缩触发失败场景:DB 末条 tool 无后续 user)。 + save_conversation(session_arc, db, conv_id, None, Some(&*resolved_model), true).await; + return RoundOutcome::Fatal; + } + } + }; + + // G4.3:本轮 token 用量是否估算值(provider 未报 prompt_tokens → estimated_prompt 兜底), + // 供 loop 内/loop 后各退出路径透传 AiCompleted(is_estimated)仅作展示标注。 + // 语义修正:整 loop 只要任一轮估算即标估算(累计总量含估算成分),非仅末轮。 + *last_round_estimated |= round_usage.prompt_tokens == 0; + + // F1 并发 epoch:流返回后若已被新 loop 接管(force_send 等),旧 loop 不再 push 消息 / + // 保文 / 执行工具,立即退出(guard disarm 跳过复位,防 clobber 新 loop 的 Generating)。 + if loop_epoch_arc.load(Ordering::SeqCst) != my_epoch { + guard.disarm(); + return RoundOutcome::StaleAfterStream; + } + + // CR-30-2 / UX-2025-04 / 决策 a1: MidStream 保文路径——partial_text 已接收, + // 入库为正常 assistant 消息 + emit AiCompleted(incomplete=true) + 追加系统提示消息。 + // 不走 AiError(非异常中断,已有可用文本),不重试(决策 a1)。 + // 并发时用户停优先:网络断同时用户点停 → stop_flag true 时不走保文路径, + // 落入下方 ~1508 stop_flag 检查走停止路径(用户意图优先)。 + // partial 文本仍由下方 push_assistant_message(~1485)保文不丢。 + if incomplete && !stop_flag.load(Ordering::SeqCst) { + let usage = df_ai::provider::TokenUsage { + prompt_tokens: if round_usage.prompt_tokens == 0 { estimated_prompt } else { round_usage.prompt_tokens }, + completion_tokens: round_usage.completion_tokens, + total_tokens: if round_usage.prompt_tokens == 0 { estimated_prompt.saturating_add(round_usage.completion_tokens) } else { round_usage.total_tokens }, + // 分项 token(2026-08-02):cache/reasoning 透传自 round_usage,落库 + 累加器都需 + prompt_cache_hit_tokens: round_usage.prompt_cache_hit_tokens, + prompt_cache_miss_tokens: round_usage.prompt_cache_miss_tokens, + reasoning_tokens: round_usage.reasoning_tokens, + }; + // 累加本轮全量 usage(含 cache/reasoning 分项)到 tokens 累加器 + tokens.add_usage(&usage); + + // 追加 partial assistant 消息(若无 tool_calls 且有文本) + // 退出校验改 conv 存在性 + push 改 per_conv.messages。 + // conv_id 来源:run_agentic_loop 入参。 + { + let mut session = session_arc.lock().await; + if !session.per_conv.contains_key(conv_id) { + tracing::warn!( + stale_conv = %conv_id, + "[ai] MidStream 保文后 conv 已删除,丢弃本轮 push" + ); + return RoundOutcome::MidStreamSaved; + } + // 与 push_assistant_message 一致:trim 归一化,纯空白保文(网络中断只有空白)不落库 + let full_text = full_text.trim(); + if !full_text.is_empty() { + let conv = session.conv(conv_id); + let mut msg = ChatMessage::assistant(full_text); + msg.model = Some(resolved_model.clone()); + // MidStream 保文也回填 reasoning_content + msg.reasoning_content = round_reasoning_content.clone(); + // 消息级 token(对齐 push_assistant_message 双轨持久化):本轮 partial usage + // (prompt=round 或 estimated 兜底,completion=round)。系统提示消息无 token,不设。 + // 分项 token(2026-08-02):cache/reasoning 透传自 round_usage。 + msg.prompt_tokens = Some(usage.prompt_tokens); + msg.completion_tokens = Some(usage.completion_tokens); + msg.prompt_cache_hit_tokens = Some(usage.prompt_cache_hit_tokens); + msg.prompt_cache_miss_tokens = Some(usage.prompt_cache_miss_tokens); + msg.reasoning_tokens = Some(usage.reasoning_tokens); + // 消息级估算标记:本轮 prompt 是否 estimated 兜底,reload 逐条回显对齐 live 态 + msg.is_estimated = Some(round_usage.prompt_tokens == 0); + conv.messages.push(msg); + // 追加系统提示消息:响应因网络中断不完整(对齐决策 a1 系统提示机制) + let mut notice = ChatMessage::system("⚠ 响应因网络中断不完整,以上为已接收的部分内容。可重新发送以获取完整回复。"); + notice.model = Some(resolved_model.clone()); + conv.messages.push(notice); + } + } + + // 统一走 finish_round_exit 收尾(save + spawn_title + reset + emit)。 + // 注意:partial 文本+系统提示已先 push(上方 block),此 save 落库含本轮 partial,幂等覆盖。 + // save_usage 用增量(tokens 已累加本轮,减上次快照);emit_usage 用 tokens 快照累计。 + // MidStream 分叉:emit_incomplete=Some(true)(前端标不完整),publish_incomplete=None(总线消费方), + // do_publish=true(publish 走总线)。spawn_title=true(后台标题,失败 extract 兜底)。 + let save_usage = usage_delta_since(tokens, saved_token_snapshot); + let emit_usage = tokens_snapshot(tokens); + finish_round_exit( + session_arc, db, conv_id, + Some(&save_usage), Some(&*resolved_model), + true, + provider_config, llm_concurrency, + guard, + &emit_usage, + *last_round_estimated, + Some(true), None, true, + pinned_goals_snapshot, + app_handle, + loop_epoch_arc, my_epoch, + ).await; + return RoundOutcome::MidStreamSaved; + } + + // 正常完成本轮 stream(无 incomplete,或 incomplete 但用户停优先走 stop 路径)。 + // 携带本轮输出交 loop 主体处理后续 push/save/工具执行。 + // round_reasoning_content 已通过上方 `*last_reasoning_content = rc` 写回 caller, + // loop 后续 push_assistant_message 读 last_reasoning_content 即可。 + let _ = round_reasoning_content; // 防未用 warning(MidStream 分支已消费) + RoundOutcome::Continue { + full_text, + tool_calls: tool_calls_acc, + round_usage, + } +} + // ── check_retry_exhausted: 判断 InitFailed Retryable 是否已耗尽重试预算 ── // 返回 Some(error) 表示已耗尽(调用方应返回 InitFailedExhausted);None 表示可继续重试。 fn check_retry_exhausted( @@ -1499,261 +1856,40 @@ pub(crate) async fn run_agentic_loop( // per-provider permit 仍在 candidate 循环内取(切换 candidate 时释放旧取新,避免占用未用 provider 的槽); // per_conv 由 loop 入口持有。 - // 重试总预算(挂钟,含 sleep + 各次请求耗时),对齐 retry::MAX_TOTAL_BUDGET 30s。 - // 本轮各 candidate 共享一个 30s 预算(切换 provider 不重置预算, - // 防多 provider 串行重试累加超过单轮总预算)。超预算直接放弃重试交最终错误/保文路径。 - let retry_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); - - // 外层 for candidate 包裹 stream_one_provider。 - // 顺序:[primary, ...candidates](空池兜底时仅 primary)。 - // - Success → 用结果覆盖迭代级 resolved_model/provider_config(供 push/save/标题); - // 成功即 break。 - // - InitFailedExhausted(retryable 耗尽)→ continue 切下一 candidate。 - // - Fatal(4xx 非429/鉴权)→ stream_llm 已 emit AiError,guard.reset + return - // 放弃整轮(对齐 retry.rs Fatal + provider_pool 文档)。 - // 全 candidate 耗尽 → 沿用单 provider 失败语义(guard.reset + return,AiError 已在 - // 各 candidate 的 stream_llm 内 emit 最后一条)。 - // 候选链(拥有 Vec):[primary, ...candidates]。预先 clone primary 入链, - // 避免借用 provider_config(切换成功后需写回 provider_config = candidate.clone())。 - // candidates 空(单 provider)→ candidate_chain 仅 primary,循环跑一次,耗尽即 return。 - let mut candidate_chain: Vec = Vec::with_capacity(1 + candidates.len()); - candidate_chain.push(provider_config.clone()); - candidate_chain.extend(candidates.iter().cloned()); - - // 防瞬时抖动(G4.1):连续 PROVIDER_SATURATION_DEMOTE_THRESHOLD 次饱和才把主 provider - // 降级到链尾。计数器在候选循环内维护(Exhausted +1 / Acquired 清零),跨迭代累计。 - // 链首(当前主 provider)连续 N 次 Exhausted → 移到链尾让 fallback 先试; - // 恢复后 Acquired 清零 → 下轮回到链首。单次/瞬时饱和(N 次内)不降级,防跳过主 provider。 - // first():链构建总是 push 至少 1 个(provider_config),但防御性取首防空链 panic。 - if let Some(streak) = candidate_chain - .first() - .and_then(|c| provider_saturation.get(&c.id)) - .copied() - .filter(|s| *s >= PROVIDER_SATURATION_DEMOTE_THRESHOLD) - { - let saturated = candidate_chain.remove(0); - candidate_chain.push(saturated.clone()); - tracing::info!( - conv_id = %conv_id, - provider = %saturated.name, - streak = streak, - threshold = PROVIDER_SATURATION_DEMOTE_THRESHOLD, - "[ai] 主 provider 连续 {} 次 per-provider 饱和,降级到候选链尾(防抖动跳过)", - streak, - ); - } - - let (full_text, tool_calls_acc, round_usage, incomplete, round_reasoning_content) = { - // outcome 累积成功结果(Complete/Partial),InitFailed 不写入。 - let mut outcome: Option<(String, std::collections::HashMap, df_ai::provider::TokenUsage, bool, Option)> = None; - // 追踪最后一个 Exhausted candidate 的诊断文本,供「全 candidate 耗尽」时 emit 最终 AiError。 - // Fatal 分支即时 emit(终态不切候选),Exhausted 分支仅记录 error 不 emit(可能切下一 candidate 成功, - // emit 会留残留气泡)。 - let mut last_exhausted_error: Option = None; - - 'candidate: for candidate in &candidate_chain { - // per-provider permit(非阻塞三态,G4.1): - // NotConfigured(未配置,无限流)→ 无 permit 直接 proceed(单 provider 零变化); - // Acquired(permit)→ 持 permit stream,切换 candidate 时随作用域 Drop 释放; - // Exhausted(占满)→ 连续饱和计数 +1,跳下一 candidate(降级,不阻塞)。 - // 旧 acquire_for_provider(阻塞 acquire_owned)在主 provider 信号量占满时整条 - // fallback 链卡死(降级失效),已由本三态非阻塞版本取代。 - let _provider_permit = match llm_concurrency.try_acquire_for_provider(&candidate.id).await { - ProviderAcquire::NotConfigured => None, - ProviderAcquire::Acquired(permit) => { - // 成功拿到 permit:该 provider 已恢复,连续饱和计数清零(跨迭代防抖)。 - provider_saturation.insert(candidate.id.clone(), 0); - Some(permit) - } - ProviderAcquire::Exhausted => { - // 占满:连续饱和计数 +1;本次跳下一 candidate(非阻塞,不等待)。 - // 跨迭代防抖:连续 N 次饱和后链首降级(见候选链构建处),防瞬时抖动跳过主 provider。 - *provider_saturation.entry(candidate.id.clone()).or_insert(0) += 1; - tracing::debug!( - conv_id = %conv_id, - provider = %candidate.name, - streak = provider_saturation.get(&candidate.id).copied().unwrap_or(0), - "[ai] 候选 {} per-provider 信号量占满,跳过(非阻塞,防单 provider 429)", - candidate.name, - ); - continue 'candidate; - } - }; - - let stream_outcome = stream_one_provider( - candidate, - &messages, - &tool_defs, - &app_handle, - &stop_flag, - ¬ify, - &conv_id, - max_retries, - retry_deadline, - &model_override, - &agentic_req, - &last_reasoning_content, - ).await; - - match stream_outcome { - StreamOutcome::Success { text, tool_calls, usage, incomplete, resolved_model: m, reasoning_content: rc } => { - // 成功:更新迭代级 resolved_model + provider_config(供后续 push/save/标题)。 - // mut 解构(已在顶部声明 mut):本轮后续 save/push 用「实际成功所用 provider」 - // 而非主 candidate。compress 已在本轮 stream 之前用过 primary,不受影响。 - resolved_model = m; - provider_config = candidate.clone(); - last_reasoning_content = rc; - outcome = Some((text, tool_calls, usage, incomplete, last_reasoning_content.clone())); - break 'candidate; - } - StreamOutcome::InitFailedExhausted { error } => { - // 本 candidate 重试耗尽(retryable)。切下一 candidate 继续尝试。 - // candidates 空(单 provider)时此即「耗尽 return」语义—— - // 跳出循环后 outcome 仍 None,落入下方「全耗尽 return」兜底。 - // 不在此 emit AiError(可能切下一 candidate 成功留残留气泡), - // 仅记录 error 供全耗尽兜底 emit 最后一条。 - tracing::warn!( - conv_id = %conv_id, - provider = %candidate.name, - "[ai] 候选 {} 流前失败重试耗尽,切换下一 provider", - candidate.name, - ); - last_exhausted_error = Some(error); - continue 'candidate; - } - StreamOutcome::Fatal { error } => { - // Fatal(4xx 非429/鉴权/参数错):立即放弃整轮 fallback。 - // 统一 emit 最终错误气泡(单气泡聚合)。 - guard.reset().await; - // F1:旧 loop(被新 loop 接管)不 emit 错误气泡(owner loop 负责呈现)。 - if loop_epoch_arc.load(Ordering::SeqCst) != my_epoch { - return; - } - emit_fatal_error(&app_handle, &conv_id, &error).await; - // G1.1(2026-08-05):Fatal 退出也落库 user 消息(镜像下方 Exhausted 分支), - // 否则重启/切会话从 DB 恢复时末条 user 丢失(与 Exhausted 同根因不对称)。 - // epoch 校验通过后、return 前补 save(owner loop 才写,防 owner 已转移仍写)。 - save_conversation(&session_arc, &db, &conv_id, None, Some(&resolved_model), true).await; - return; - } - } - } - - match outcome { - Some(r) => r, - None => { - // 全 candidate 耗尽(含单 provider 场景):沿用原单 provider 失败语义。 - tracing::warn!( - conv_id = %conv_id, - candidates_tried = candidate_chain.len(), - "[ai] 全 provider 候选流前失败重试耗尽,放弃本轮", - ); - guard.reset().await; - if loop_epoch_arc.load(Ordering::SeqCst) != my_epoch { - return; - } - let err_msg = last_exhausted_error - .unwrap_or_else(|| "AI 调用失败:所有候选 provider 重试耗尽".to_string()); - emit_fatal_error(&app_handle, &conv_id, &err_msg).await; - // 全 provider 失败也需落库 user 消息(已 push 内存), - // 否则切换/重载从 DB 恢复时末条 user 丢失(实测压缩触发失败场景:DB 末条 tool 无后续 user)。 - save_conversation(&session_arc, &db, &conv_id, None, Some(&resolved_model), true).await; - return; - } + // 单轮流式调用(candidate 链 + stream_one_provider + MidStream 保文)抽至 stream_round。 + // + // 范围:retry_deadline 计算 → candidate_chain 构建/降级 → candidate 循环 + // (try_acquire_for_provider + stream_one_provider + match 三态)→ epoch 后 stale 检查 + // → MidStream 保文路径。返回 RoundOutcome 表达本轮结果: + // - Continue:正常完成,携输出交后续 push/save/工具执行; + // - StaleAfterStream:流后 epoch 变,guard 已 disarm,直接 return; + // - MidStreamSaved:保文路径已完整收尾,直接 return; + // - Fatal:候选 Fatal / 全耗尽,guard.reset + emit + save 已完成,直接 return。 + // + // 借用:tokens/saved_token_snapshot/last_round_estimated/last_reasoning_content/ + // resolved_model/provider_config/provider_saturation/guard 用 &mut(stream_round 内 + // 修改 caller 可见);其它 & 引用。permit 责任不变(per_conv loop 入口持有, + // per-provider candidate 循环内取/释放,均在 stream_round 内)。 + let (full_text, tool_calls_acc, round_usage) = match stream_round( + &session_arc, &db, &conv_id, + &messages, &tool_defs, &app_handle, + &stop_flag, ¬ify, + max_retries, &model_override, &agentic_req, estimated_prompt, + &mut last_round_estimated, &mut last_reasoning_content, + &mut resolved_model, &mut provider_config, + &mut provider_saturation, + &mut tokens, &mut saved_token_snapshot, &mut guard, + &candidates, &llm_concurrency, + &pinned_goals_snapshot, + &loop_epoch_arc, my_epoch, + ).await { + RoundOutcome::Continue { full_text, tool_calls, round_usage } => { + (full_text, tool_calls, round_usage) } + RoundOutcome::StaleAfterStream => return, + RoundOutcome::MidStreamSaved => return, + RoundOutcome::Fatal => return, }; - - // G4.3:本轮 token 用量是否估算值(provider 未报 prompt_tokens → estimated_prompt 兜底), - // 供 loop 内/loop 后各退出路径透传 AiCompleted(is_estimated)仅作展示标注。 - // 语义修正:整 loop 只要任一轮估算即标估算(累计总量含估算成分),非仅末轮。 - last_round_estimated |= round_usage.prompt_tokens == 0; - - // F1 并发 epoch:流返回后若已被新 loop 接管(force_send 等),旧 loop 不再 push 消息 / - // 保文 / 执行工具,立即退出(guard disarm 跳过复位,防 clobber 新 loop 的 Generating)。 - if loop_epoch_arc.load(Ordering::SeqCst) != my_epoch { - guard.disarm(); - return; - } - - // CR-30-2 / UX-2025-04 / 决策 a1: MidStream 保文路径——partial_text 已接收, - // 入库为正常 assistant 消息 + emit AiCompleted(incomplete=true) + 追加系统提示消息。 - // 不走 AiError(非异常中断,已有可用文本),不重试(决策 a1)。 - // 并发时用户停优先:网络断同时用户点停 → stop_flag true 时不走保文路径, - // 落入下方 ~1508 stop_flag 检查走停止路径(用户意图优先)。 - // partial 文本仍由下方 push_assistant_message(~1485)保文不丢。 - if incomplete && !stop_flag.load(Ordering::SeqCst) { - let usage = df_ai::provider::TokenUsage { - prompt_tokens: if round_usage.prompt_tokens == 0 { estimated_prompt } else { round_usage.prompt_tokens }, - completion_tokens: round_usage.completion_tokens, - total_tokens: if round_usage.prompt_tokens == 0 { estimated_prompt.saturating_add(round_usage.completion_tokens) } else { round_usage.total_tokens }, - // 分项 token(2026-08-02):cache/reasoning 透传自 round_usage,落库 + 累加器都需 - prompt_cache_hit_tokens: round_usage.prompt_cache_hit_tokens, - prompt_cache_miss_tokens: round_usage.prompt_cache_miss_tokens, - reasoning_tokens: round_usage.reasoning_tokens, - }; - // 累加本轮全量 usage(含 cache/reasoning 分项)到 tokens 累加器 - tokens.add_usage(&usage); - - // 追加 partial assistant 消息(若无 tool_calls 且有文本) - // 退出校验改 conv 存在性 + push 改 per_conv.messages。 - // conv_id 来源:run_agentic_loop 入参。 - { - let mut session = session_arc.lock().await; - if !session.per_conv.contains_key(&conv_id) { - tracing::warn!( - stale_conv = %conv_id, - "[ai] MidStream 保文后 conv 已删除,丢弃本轮 push" - ); - return; - } - // 与 push_assistant_message 一致:trim 归一化,纯空白保文(网络中断只有空白)不落库 - let full_text = full_text.trim(); - if !full_text.is_empty() { - let conv = session.conv(&conv_id); - let mut msg = ChatMessage::assistant(full_text); - msg.model = Some(resolved_model.clone()); - // MidStream 保文也回填 reasoning_content - msg.reasoning_content = round_reasoning_content.clone(); - // 消息级 token(对齐 push_assistant_message 双轨持久化):本轮 partial usage - // (prompt=round 或 estimated 兜底,completion=round)。系统提示消息无 token,不设。 - // 分项 token(2026-08-02):cache/reasoning 透传自 round_usage。 - msg.prompt_tokens = Some(usage.prompt_tokens); - msg.completion_tokens = Some(usage.completion_tokens); - msg.prompt_cache_hit_tokens = Some(usage.prompt_cache_hit_tokens); - msg.prompt_cache_miss_tokens = Some(usage.prompt_cache_miss_tokens); - msg.reasoning_tokens = Some(usage.reasoning_tokens); - // 消息级估算标记:本轮 prompt 是否 estimated 兜底,reload 逐条回显对齐 live 态 - msg.is_estimated = Some(round_usage.prompt_tokens == 0); - conv.messages.push(msg); - // 追加系统提示消息:响应因网络中断不完整(对齐决策 a1 系统提示机制) - let mut notice = ChatMessage::system("⚠ 响应因网络中断不完整,以上为已接收的部分内容。可重新发送以获取完整回复。"); - notice.model = Some(resolved_model.clone()); - conv.messages.push(notice); - } - } - - // 统一走 finish_round_exit 收尾(save + spawn_title + reset + emit)。 - // 注意:partial 文本+系统提示已先 push(上方 block),此 save 落库含本轮 partial,幂等覆盖。 - // save_usage 用增量(tokens 已累加本轮,减上次快照);emit_usage 用 tokens 快照累计。 - // MidStream 分叉:emit_incomplete=Some(true)(前端标不完整),publish_incomplete=None(总线消费方), - // do_publish=true(publish 走总线)。spawn_title=true(后台标题,失败 extract 兜底)。 - let save_usage = usage_delta_since(&tokens, &mut saved_token_snapshot); - let emit_usage = tokens_snapshot(&tokens); - finish_round_exit( - &session_arc, &db, &conv_id, - Some(&save_usage), Some(&resolved_model), - true, - &provider_config, &llm_concurrency, - &mut guard, - &emit_usage, - last_round_estimated, - Some(true), None, true, - &pinned_goals_snapshot, - &app_handle, - &loop_epoch_arc, my_epoch, - ).await; - return; - } // global/per_conv permit 已上移 loop 入口整 loop 持有, // stream 后不再立即释放(会话级并发语义:工具执行期间也占槽)。per-provider permit // (_provider_permit)仍在 candidate 循环内随作用域 Drop 自动释放。