重构: AI执行循环拆分降低复杂度

This commit is contained in:
lxy
2026-08-10 08:04:30 +08:00
parent 6141a67b0e
commit 5a1b7a871a
+389 -253
View File
@@ -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<u32, super::ToolCallDraft>,
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<String>`: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<Mutex<AiSession>>,
db: &Arc<Database>,
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<String>,
agentic_req: &TaskRequirements,
estimated_prompt: u32,
// &mut loop 级状态(函数内修改,caller 可见)
last_round_estimated: &mut bool,
last_reasoning_content: &mut Option<String>,
resolved_model: &mut String,
provider_config: &mut AiProviderRecord,
provider_saturation: &mut HashMap<String, usize>,
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<std::sync::atomic::AtomicU64>,
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<AiProviderRecord> = 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<u32, super::ToolCallDraft>, df_ai::provider::TokenUsage, bool, Option<String>)> = None;
// 追踪最后一个 Exhausted candidate 的诊断文本,供「全 candidate 耗尽」时 emit 最终 AiError。
// Fatal 分支即时 emit(终态不切候选),Exhausted 分支仅记录 error 不 emit(可能切下一 candidate 成功,
// emit 会留残留气泡)。
let mut last_exhausted_error: Option<String> = 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 是否已耗尽重试预算 ── // ── check_retry_exhausted: 判断 InitFailed Retryable 是否已耗尽重试预算 ──
// 返回 Some(error) 表示已耗尽(调用方应返回 InitFailedExhausted);None 表示可继续重试。 // 返回 Some(error) 表示已耗尽(调用方应返回 InitFailedExhausted);None 表示可继续重试。
fn check_retry_exhausted( fn check_retry_exhausted(
@@ -1499,261 +1856,40 @@ pub(crate) async fn run_agentic_loop(
// per-provider permit 仍在 candidate 循环内取(切换 candidate 时释放旧取新,避免占用未用 provider 的槽); // per-provider permit 仍在 candidate 循环内取(切换 candidate 时释放旧取新,避免占用未用 provider 的槽);
// per_conv 由 loop 入口持有。 // per_conv 由 loop 入口持有。
// 重试总预算(挂钟,含 sleep + 各次请求耗时),对齐 retry::MAX_TOTAL_BUDGET 30s // 单轮流式调用(candidate 链 + stream_one_provider + MidStream 保文)抽至 stream_round
// 本轮各 candidate 共享一个 30s 预算(切换 provider 不重置预算, //
// 防多 provider 串行重试累加超过单轮总预算)。超预算直接放弃重试交最终错误/保文路径。 // 范围:retry_deadline 计算 → candidate_chain 构建/降级 → candidate 循环
let retry_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); // (try_acquire_for_provider + stream_one_provider + match 三态)→ epoch 后 stale 检查
// → MidStream 保文路径。返回 RoundOutcome 表达本轮结果:
// 外层 for candidate 包裹 stream_one_provider。 // - Continue:正常完成,携输出交后续 push/save/工具执行;
// 顺序:[primary, ...candidates](空池兜底时仅 primary)。 // - StaleAfterStream:流后 epoch 变,guard 已 disarm,直接 return;
// - Success → 用结果覆盖迭代级 resolved_model/provider_config(供 push/save/标题); // - MidStreamSaved:保文路径已完整收尾,直接 return;
// 成功即 break // - Fatal:候选 Fatal / 全耗尽,guard.reset + emit + save 已完成,直接 return
// - InitFailedExhausted(retryable 耗尽)→ continue 切下一 candidate。 //
// - Fatal(4xx 非429/鉴权)→ stream_llm 已 emit AiError,guard.reset + return // 借用:tokens/saved_token_snapshot/last_round_estimated/last_reasoning_content/
// 放弃整轮(对齐 retry.rs Fatal + provider_pool 文档)。 // resolved_model/provider_config/provider_saturation/guard 用 &mut(stream_round 内
// candidate 耗尽 → 沿用单 provider 失败语义(guard.reset + return,AiError 已在 // 修改 caller 可见);其它 & 引用。permit 责任不变(per_conv loop 入口持有,
// 各 candidate 的 stream_llm 内 emit 最后一条)。 // per-provider candidate 循环内取/释放,均在 stream_round 内)。
// 候选链(拥有 Vec):[primary, ...candidates]。预先 clone primary 入链, let (full_text, tool_calls_acc, round_usage) = match stream_round(
// 避免借用 provider_config(切换成功后需写回 provider_config = candidate.clone())。 &session_arc, &db, &conv_id,
// candidates 空(单 provider)→ candidate_chain 仅 primary,循环跑一次,耗尽即 return。 &messages, &tool_defs, &app_handle,
let mut candidate_chain: Vec<AiProviderRecord> = Vec::with_capacity(1 + candidates.len()); &stop_flag, &notify,
candidate_chain.push(provider_config.clone()); max_retries, &model_override, &agentic_req, estimated_prompt,
candidate_chain.extend(candidates.iter().cloned()); &mut last_round_estimated, &mut last_reasoning_content,
&mut resolved_model, &mut provider_config,
// 防瞬时抖动(G4.1):连续 PROVIDER_SATURATION_DEMOTE_THRESHOLD 次饱和才把主 provider &mut provider_saturation,
// 降级到链尾。计数器在候选循环内维护(Exhausted +1 / Acquired 清零),跨迭代累计。 &mut tokens, &mut saved_token_snapshot, &mut guard,
// 链首(当前主 provider)连续 N 次 Exhausted → 移到链尾让 fallback 先试; &candidates, &llm_concurrency,
// 恢复后 Acquired 清零 → 下轮回到链首。单次/瞬时饱和(N 次内)不降级,防跳过主 provider。 &pinned_goals_snapshot,
// first():链构建总是 push 至少 1 个(provider_config),但防御性取首防空链 panic。 &loop_epoch_arc, my_epoch,
if let Some(streak) = candidate_chain ).await {
.first() RoundOutcome::Continue { full_text, tool_calls, round_usage } => {
.and_then(|c| provider_saturation.get(&c.id)) (full_text, tool_calls, round_usage)
.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<u32, super::ToolCallDraft>, df_ai::provider::TokenUsage, bool, Option<String>)> = None;
// 追踪最后一个 Exhausted candidate 的诊断文本,供「全 candidate 耗尽」时 emit 最终 AiError。
// Fatal 分支即时 emit(终态不切候选),Exhausted 分支仅记录 error 不 emit(可能切下一 candidate 成功,
// emit 会留残留气泡)。
let mut last_exhausted_error: Option<String> = 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;
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;
}
} }
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 持有, // global/per_conv permit 已上移 loop 入口整 loop 持有,
// stream 后不再立即释放(会话级并发语义:工具执行期间也占槽)。per-provider permit // stream 后不再立即释放(会话级并发语义:工具执行期间也占槽)。per-provider permit
// (_provider_permit)仍在 candidate 循环内随作用域 Drop 自动释放。 // (_provider_permit)仍在 candidate 循环内随作用域 Drop 自动释放。