新增: F-04b多provider fallback切换(stream_one_provider+候选链+resolved_model重算)

This commit is contained in:
2026-06-17 02:26:22 +08:00
parent b3684f4d1f
commit 80c0955a1a

View File

@@ -106,6 +106,197 @@ impl Drop for GeneratingGuard {
} }
} }
// ============================================================
// F-260614-04 / F-260614-04b: 单 Provider 流式结果 + fallback 辅助
// ============================================================
/// 单 candidate provider 一次迭代的流式调用结果(F-260614-04b fallback)。
///
/// 显式区分三态,驱动外层 `for candidate` fallback 循环:
/// - `Success`:`Complete` 或 `Partial`(MidStream 保文)。携带该 provider 路由出的
/// `resolved_model`(调用方据此更新迭代级 resolved_model,供 push/save)。
/// - `InitFailedExhausted`:流前 `InitFailed{retryable=true}` 在本 provider 上重试
/// 耗尽(或预算 30s 耗尽)。retryable=true 即「瞬态错误,换 provider 可能成功」→
/// 外层 `continue` 切下一 candidate(重建 build_provider + resolved_model 重算)。
/// - `Fatal`:流前 `InitFailed{retryable=false}`(4xx 非429/鉴权/参数错)。立即放弃
/// 整个 fallback(对齐 retry.rs Fatal 分类 + provider_pool 文档「不可重试错误
/// 立即放弃不浪费备用 provider」)。stream_llm 已 emit AiError,调用方仅做
/// guard.reset + return。空 key(build_provider_for Err)同样归 Fatal 语义——
/// key 缺失/钥匙串损坏非瞬态,换 provider 无意义(但实际场景:外层已对 primary
/// 做了启动 Auth 早失败,故此分支多见于 secondary 配置不一致,保守 Fatal 收敛)。
enum StreamOutcome {
Success {
text: String,
tool_calls: std::collections::HashMap<u32, super::ToolCallDraft>,
usage: df_ai::provider::TokenUsage,
incomplete: bool,
resolved_model: String,
},
InitFailedExhausted,
Fatal,
}
/// 单 provider 流式调用 + 重试(F-260614-04b)。
///
/// 封装一个 candidate 的:resolve_secret → ensure_resolved_key → build_provider
/// → select_model_id(resolved_model 在该 candidate.model_configs 上重算)
/// → 流式 stream_recv + 流前 InitFailed retryable 重试循环。
///
/// 与波 12 主链共享的退避/分类:`retry::backoff_delay`(1s→2s→4s+jitter) +
/// `StreamResult::InitFailed{retryable}`(镜像 retry::is_status_retryable) +
/// 30s 总挂钟预算(`retry_deadline` 由调用方传入,本轮 fallback 各 candidate 共享一个预算)。
///
/// **permit 责任**:本函数不取并发 permit——permit 持有期间多 provider 限流不重叠,
/// 由调用方在进 candidate 前取 global/per_conv(整轮迭代共享) + 本 candidate 的
/// per-provider permit(切换 candidate 时释放旧取新)。
///
/// 返回 [`StreamOutcome`]:Fatal 含已 emit AiError;Success 含保文或正常完成;
/// InitFailedExhausted 让调用方切下一 candidate。
async fn stream_one_provider(
candidate: &AiProviderRecord,
messages: &[ChatMessage],
tool_defs: &[df_ai::provider::ToolDefinition],
app_handle: &AppHandle,
stop_flag: &std::sync::atomic::AtomicBool,
notify: &tokio::sync::Notify,
conv_id: &str,
max_retries: usize,
retry_deadline: tokio::time::Instant,
model_override: &Option<String>,
agentic_req: &TaskRequirements,
) -> StreamOutcome {
// resolve→ensure→build 三步(复用 secret::build_provider_for,DRY)。
// key 缺失/损坏 → Err → 归 Fatal(stream_llm Fatal 也是 key 类错,语义一致)。
// 注:primary 候选的启动 Auth 早失败已在外层 run_agentic_loop 顶部处理(emit + return),
// 本函数到达时 primary 的 key 已验证过;此处 Err 多见于 secondary 配置不一致,
// 保守归 Fatal 立即放弃(不浪费预算试下一 provider,因 key 错非瞬态)。
let provider: Box<dyn LlmProvider> = match super::secret::build_provider_for(candidate) {
Ok(p) => p,
Err(msg) => {
// 不复用 stream_llm 的 AiError emit(那是 HTTP 诊断),此处补一条 Auth 分类错误。
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: msg,
error_type: Some(ErrorType::Auth),
conversation_id: Some(conv_id.to_string()),
});
return StreamOutcome::Fatal;
}
};
// resolved_model 在本 candidate.model_configs 上重算(F-260614-04b 核心:provider 切换后
// 模型池不同,必须重选;否则拿主 provider 的 model_id 去打次 provider 会吃 400/404)。
let resolved_model = select_model_id(agentic_req, &candidate.model_configs)
.unwrap_or_else(|| candidate.default_model.clone());
// model_override 穿透(对齐主链 F-01 阶段6):override 非空且在本 candidate 池中 → 用 override;
// 否则用 resolved_model(绝不因 override 致无模型)。
let resolved_model = match model_override.as_deref() {
Some(id) if !id.is_empty()
&& candidate.model_configs.iter().any(|m| m.model_id == id) =>
{
id.to_string()
}
_ => resolved_model,
};
// 流前 InitFailed 重试循环(对齐主链 CR-30-1 / F-260616-07 决策 a1)。
// outcome 累积 Complete/Partial;InitFailed retryable 未耗尽即重试,Fatal 立即返回。
for retry_attempt in 0..=max_retries {
let retry_request = CompletionRequest {
model: resolved_model.clone(),
messages: messages.to_vec(),
temperature: Some(0.7),
max_tokens: Some(8192),
stream: true,
tools: if tool_defs.is_empty() { None } else { Some(tool_defs.to_vec()) },
tool_choice: None,
};
match stream_llm(&*provider, retry_request, app_handle, stop_flag, notify, conv_id).await {
StreamResult::Complete { text, tool_calls, usage } => {
return StreamOutcome::Success {
text, tool_calls, usage,
incomplete: false,
resolved_model,
};
}
StreamResult::Partial { text, tool_calls, usage } => {
// MidStream 保文不重试(决策 a1)。携带 incomplete=true 交调用方走保文路径。
tracing::warn!(
conv_id = %conv_id,
text_len = text.len(),
"[ai] 流中途失败(候选 {}),保文不重试(incomplete=true)",
candidate.name,
);
return StreamOutcome::Success {
text, tool_calls, usage,
incomplete: true,
resolved_model,
};
}
StreamResult::InitFailed { retryable } => {
// Fatal(4xx 非429/鉴权/参数错)立即放弃整个 fallback。
if !retryable {
tracing::warn!(
conv_id = %conv_id,
provider = %candidate.name,
attempt = retry_attempt + 1,
"[ai] 候选 {} 流前失败 Fatal(4xx/鉴权),立即放弃 fallback",
candidate.name,
);
return StreamOutcome::Fatal;
}
let is_last = retry_attempt >= max_retries;
let now = tokio::time::Instant::now();
let budget_exhausted = now >= retry_deadline;
if is_last || budget_exhausted {
// 本 candidate 重试耗尽 / 预算耗尽 → 交外层切下一 candidate。
// stream_llm 已 emit AiError(每次 InitFailed 均发),故不再重复 emit。
tracing::warn!(
conv_id = %conv_id,
provider = %candidate.name,
attempt = retry_attempt + 1,
budget_exhausted = budget_exhausted,
"[ai] 候选 {} 流前失败重试{}({}次),切下一 provider",
candidate.name,
if budget_exhausted { "预算耗尽" } else { "耗尽" },
max_retries + 1,
);
return StreamOutcome::InitFailedExhausted;
}
// 退避复用 retry::backoff_delay(retry_attempt+1)(1s→2s→4s ±20% jitter)。
let delay = retry::backoff_delay((retry_attempt + 1) as u32)
.min(retry_deadline.saturating_duration_since(now));
let total_retry = retry_attempt + 1;
tracing::warn!(
conv_id = %conv_id,
provider = %candidate.name,
attempt = total_retry,
max_attempts = max_retries + 1,
delay_ms = delay.as_millis() as u64,
"[ai] 候选 {} 流前失败(Retryable),{}ms 后重试 ({}/{})",
candidate.name, delay.as_millis(), total_retry, max_retries + 1,
);
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiStreamRetry {
attempt: total_retry as u32,
max_attempts: (max_retries + 1) as u32,
conversation_id: Some(conv_id.to_string()),
});
tokio::time::sleep(delay).await;
continue;
}
}
}
// 仅 max_retries=0(无重试配置)且首试 InitFailed retryable=true 时到达此处:
// 0..=0 循环首趟即 is_last → 已在循环内 return InitFailedExhausted。此为防御兜底。
StreamOutcome::InitFailedExhausted
}
/// Agentic 循环:流式接收 → 工具执行 → 结果回传 LLM → 循环 /// Agentic 循环:流式接收 → 工具执行 → 结果回传 LLM → 循环
/// ///
/// 退出条件: /// 退出条件:
@@ -130,31 +321,40 @@ pub(crate) async fn run_agentic_loop(
// B-260615-09: generating 状态由 RAII guard 收敛复位(正常 exit 显式 reset;panic/异常 Drop 兜底) // B-260615-09: generating 状态由 RAII guard 收敛复位(正常 exit 显式 reset;panic/异常 Drop 兜底)
let mut guard = GeneratingGuard::new(session_arc.clone()); let mut guard = GeneratingGuard::new(session_arc.clone());
// F-260614-04: 多 Provider 负载均衡池 — 选主候选 provider // F-260614-04 / F-260614-04b: 多 Provider 负载均衡池 — 选主 + fallback 候选列表
// //
// 流程:list_all → ProviderPool::select(按 模型亲和 > weight > is_default 排序)→ 取首位 // 流程:list_all → ProviderPool::select(按 模型亲和 > weight > is_default 排序)→ 有序 Vec
// 单 provider 场景:池仅 1 enabled provider → select 返回单元素 Vec → 首位 = 唯一 provider, // 单 provider 场景:池仅 1 enabled provider → select 返回单元素 Vec → 首位 = 唯一 provider,
// 行为同 F-01 前(零变化)。空池(0 enabled)→ fallback 入参 provider_config(保启动行为)。 // 行为同 F-01 前(零变化)。空池(0 enabled)→ fallback 入参 provider_config(保启动行为)。
// //
// F-04b:select 返回的**完整有序 Vec** 即 fallback 序(主→备用)。本轮接入真实切换:
// 流式重试块外层包 `for candidate in &candidates`,主 candidate InitFailed{retryable=true}
// 耗尽 → continue 切下一 candidate(stream_one_provider 内重建 build_provider + resolved_model
// 在新 candidate.model_configs 上重算 select_model_id);Fatal(4xx 非429)立即放弃整轮。
//
// 主候选的 model_configs 用于路由(F-01),其 provider_config 用于 build_provider。 // 主候选的 model_configs 用于路由(F-01),其 provider_config 用于 build_provider。
// 当前批次仅选主,fallback(主失败切备用 provider)见后续批次(需重构流式重试块,单独评估)。 // compress_via_llm / 标题 / 后台 spawn 沿用主 candidate(非 fallback 范围,见各调用点注释)。
// //
// AiProviderRepo::new 仅 clone Arc<Database>(廉价),不复用 AppState.ai_providers // AiProviderRepo::new 仅 clone Arc<Database>(廉价),不复用 AppState.ai_providers
// (run_agentic_loop 签名只传 Arc<Database>,改签名会牵动 3 调用点 + try_continue)。 // (run_agentic_loop 签名只传 Arc<Database>,改签名会牵动 3 调用点 + try_continue)。
let provider_repo = df_storage::crud::AiProviderRepo::new(&db); let provider_repo = df_storage::crud::AiProviderRepo::new(&db);
let pool_providers: Vec<AiProviderRecord> = provider_repo.list_all().await.unwrap_or_default(); let pool_providers: Vec<AiProviderRecord> = provider_repo.list_all().await.unwrap_or_default();
let primary_provider: AiProviderRecord = match super::provider_pool::ProviderPool::select( let ranked_candidates: Vec<AiProviderRecord> = super::provider_pool::ProviderPool::select(
&pool_providers, &pool_providers,
None, // 模型亲和此时尚未确定(router 选模型需 provider_config,见下方) None, // 模型亲和此时尚未确定(router 选模型需 provider_config,见下方)
) );
.into_iter() let (primary_provider, candidates): (AiProviderRecord, Vec<AiProviderRecord>) =
.next() match ranked_candidates.split_first() {
{ Some((first, rest)) => (first.clone(), rest.to_vec()),
Some(p) => p, None => {
None => provider_config.clone(), // 空池兜底:用调用方传入的 provider_config(启动行为不变) // 空池兜底:用调用方传入的 provider_config 作唯一候选(启动行为不变)
}; // candidates 空 → fallback 循环仅跑 primary 一次,等同单 provider 路径。
(provider_config.clone(), Vec::new())
}
};
// 用主候选覆盖入参 provider_config(下游 build_provider / 路由 / 日志均用此)。 // 用主候选覆盖入参 provider_config(下游 build_provider / 路由 / 日志均用此)。
let provider_config = primary_provider; // mut:F-04b 切换 candidate 后更新为实际成功所用 provider(供后续 push/save/标题 spawn)。
let mut provider_config = primary_provider;
// FR-S1: resolve→ensure_resolved_key(空 key 早失败)→build_provider 三步统一走工厂 // FR-S1: resolve→ensure_resolved_key(空 key 早失败)→build_provider 三步统一走工厂
// 空 key 早失败(逻辑见 secret::ensure_resolved_key 单测):避免空 key 发请求吃 401,错误伪装成"API Key 无效" // 空 key 早失败(逻辑见 secret::ensure_resolved_key 单测):避免空 key 发请求吃 401,错误伪装成"API Key 无效"
@@ -210,7 +410,9 @@ pub(crate) async fn run_agentic_loop(
// F-01 阶段6: 用户指定模型 override 穿透(仅主对话生效,标题/扫描/灵感仍走路由)。 // F-01 阶段6: 用户指定模型 override 穿透(仅主对话生效,标题/扫描/灵感仍走路由)。
// 兜底原则:override 非空且在该 provider model_configs 池中 → 用 override;否则用 resolved_model。 // 兜底原则:override 非空且在该 provider model_configs 池中 → 用 override;否则用 resolved_model。
// 绝不让 override 导致无模型(空/不在池 → 落回路由结果,行为不变)。 // 绝不让 override 导致无模型(空/不在池 → 落回路由结果,行为不变)。
let resolved_model = match model_override.as_deref() { // mut:F-04b 切换 candidate 后由 stream_one_provider 返回的 Success.resolved_model 覆盖
// (新 candidate 模型池不同,必须重选);此后 push/save 均用迭代最新值。
let mut resolved_model = match model_override.as_deref() {
Some(id) if !id.is_empty() Some(id) if !id.is_empty()
&& provider_config.model_configs.iter().any(|m| m.model_id == id) => && provider_config.model_configs.iter().any(|m| m.model_id == id) =>
{ {
@@ -412,14 +614,14 @@ pub(crate) async fn run_agentic_loop(
}; };
// 预估输入 token(兜底:部分 provider 如 GLM 流式 usage 不报 prompt_tokens,后段用它补) // 预估输入 token(兜底:部分 provider 如 GLM 流式 usage 不报 prompt_tokens,后段用它补)
// 注:F-260616-07 重试循环内每次重建 retry_request(因 provider.stream 消费 body), // 注:stream_one_provider 内每次重试重建 request(因 provider.stream 消费 body),
// 此处不再预构建 request(旧 request 变量已废弃),仅保留 messages 供 estimated_prompt。 // 此处不再预构建 request(旧 request 变量已废弃),仅保留 messages 供 estimated_prompt。
let estimated_prompt: u32 = { let estimated_prompt: u32 = {
let est = TokenEstimator::default(); let est = TokenEstimator::default();
messages.iter().map(|m| est.estimate_message(m)).sum() messages.iter().map(|m| est.estimate_message(m)).sum()
}; };
// LLM 并发限流全局 + 单对话双层仅覆盖 stream_llm 调用本身; // LLM 并发限流(全局 + 单对话双层),仅覆盖 stream_llm 调用本身;
// 工具执行(process_tool_calls)是本地操作无 RPM 成本,permit 在 stream 后立即释放避免占槽 // 工具执行(process_tool_calls)是本地操作无 RPM 成本,permit 在 stream 后立即释放避免占槽
// //
// CR-30-1 / F-260616-07 / 决策 a1: 流前失败(Init Err)重试,流中途失败(MidStream // CR-30-1 / F-260616-07 / 决策 a1: 流前失败(Init Err)重试,流中途失败(MidStream
@@ -427,114 +629,82 @@ pub(crate) async fn run_agentic_loop(
// retry::is_status_retryable Fatal 分类(stream_recv classify_status_or_class 镜像, // retry::is_status_retryable Fatal 分类(stream_recv classify_status_or_class 镜像,
// 4xx 非429 立即放弃) + 30s 总挂钟预算。重试期间持有 permit 不释放(防新请求挤占)。 // 4xx 非429 立即放弃) + 30s 总挂钟预算。重试期间持有 permit 不释放(防新请求挤占)。
// //
// F-260614-04: per-provider permit(可选)。set_provider_caps 未配置时返回 None // F-260614-04 / F-260614-04b: global/per_conv permit 整轮迭代共享(不随 candidate 切换
// (单 provider 场景零变化);配置后取额外 permit 防单 provider 被打满(限流 429)。 // 释放——并发语义是「全局/单对话」级,与具体 provider 无关);per-provider permit 改为
// 三 permit 均 Drop 释放(L595-597 显式 drop _global/_per_conv;_provider_permit 绑块尾 Drop)。 // 在 candidate 循环内取(切换 candidate 时释放旧取新,避免占用未用 provider 的槽)。
// 两 permit 均 Drop 释放(L803-804 显式 drop _global/_per_conv)。
let _global_permit = llm_concurrency.acquire_global().await; let _global_permit = llm_concurrency.acquire_global().await;
let _per_conv_permit = llm_concurrency.acquire_per_conv().await; let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
let _provider_permit = llm_concurrency.acquire_for_provider(&provider_config.id).await;
// 重试总预算(挂钟,含 sleep + 各次请求耗时),对齐 retry::MAX_TOTAL_BUDGET 30s。 // 重试总预算(挂钟,含 sleep + 各次请求耗时),对齐 retry::MAX_TOTAL_BUDGET 30s。
// 超预算直接放弃重试交最终错误/保文路径。 // F-260614-04b:本轮各 candidate 共享一个 30s 预算(切换 provider 不重置预算,
// 防多 provider 串行重试累加超过单轮总预算)。超预算直接放弃重试交最终错误/保文路径。
let retry_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); let retry_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30);
// F-260614-04b:外层 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());
let (full_text, tool_calls_acc, round_usage, incomplete) = { let (full_text, tool_calls_acc, round_usage, incomplete) = {
// outcome 累积最后一次成功/保文结果(CompletePartial)InitFailed 不写入, // outcome 累积成功结果(Complete/Partial),InitFailed 不写入
// 走重试或耗尽 return。
let mut outcome: Option<(String, std::collections::HashMap<u32, super::ToolCallDraft>, df_ai::provider::TokenUsage, bool)> = None; let mut outcome: Option<(String, std::collections::HashMap<u32, super::ToolCallDraft>, df_ai::provider::TokenUsage, bool)> = None;
for retry_attempt in 0..=max_retries { 'candidate: for candidate in &candidate_chain {
// 每次重试重建 request(CompletionRequest 无状态,但 provider.stream() 内部可能消费 body) // per-provider permit(可选):set_provider_caps 未配置时返回 None(单 provider 零变化);
// // 配置后取额外 permit 防单 provider 被打满(限流 429)。切换 candidate 时上一 permit
// F-260616-13: messages 复用本轮外层构建的快照(不再持 session_arc.lock() 重建) // 随 _provider_permit 绑定作用域 Drop 释放
// 安全前提:本轮 stream_llm 不接收 session_arc、不 push 消息;重试期间 tool 执行 let _provider_permit = llm_concurrency.acquire_for_provider(&candidate.id).await;
// (process_tool_calls)在重试块之后,故重试内 session.messages 必然与外层构建时
// 一致,重建与复用等价。收益:省去每次重试的 lock() + build_for_request(含
// all_messages_clone 全量 clone)。
let retry_request = CompletionRequest {
model: resolved_model.clone(),
messages: messages.clone(),
temperature: Some(0.7),
max_tokens: Some(8192),
stream: true,
tools: if tool_defs.is_empty() { None } else { Some(tool_defs.clone()) },
tool_choice: None,
};
match stream_llm(&*provider, retry_request, &app_handle, &stop_flag, &notify, &conv_id).await { match stream_one_provider(
StreamResult::Complete { text, tool_calls, usage } => { candidate,
// 正常完成(流尽 + finished,或用户主动停止)。incomplete=false。 &messages,
outcome = Some((text, tool_calls, usage, false)); &tool_defs,
break; &app_handle,
&stop_flag,
&notify,
&conv_id,
max_retries,
retry_deadline,
&model_override,
&agentic_req,
).await {
StreamOutcome::Success { text, tool_calls, usage, incomplete, resolved_model: m } => {
// 成功:更新迭代级 resolved_model + provider_config(供后续 push/save/标题)。
// mut 解构(已在顶部声明 mut):本轮后续 save/push 用「实际成功所用 provider」
// 而非主 candidate。compress 已在本轮 stream 之前用过 primary,不受影响。
resolved_model = m;
provider_config = candidate.clone();
outcome = Some((text, tool_calls, usage, incomplete));
break 'candidate;
} }
StreamResult::Partial { text, tool_calls, usage } => { StreamOutcome::InitFailedExhausted => {
// CR-30-2 / UX-2025-04 / 决策 a1: 流中途失败保文,**不重试** // 本 candidate 重试耗尽(retryable)。切下一 candidate 继续尝试
// 已 emit AiTextDelta(前端 currentText 已累积),此处保文入库 + // candidates 空(单 provider)时此即「耗尽 return」语义——
// AiCompleted(incomplete=true) + 系统提示。currentText 不混乱(未重试无追加) // 跳出循环后 outcome 仍 None,落入下方「全耗尽 return」兜底
tracing::warn!( tracing::warn!(
conv_id = %conv_id, conv_id = %conv_id,
text_len = text.len(), provider = %candidate.name,
"[ai] 流中途失败,保文不重试(incomplete=true),入库 + AiCompleted + 系统提示", "[ai] 候选 {} 流前失败重试耗尽,F-04b 切换下一 provider",
candidate.name,
); );
outcome = Some((text, tool_calls, usage, true)); continue 'candidate;
break;
} }
StreamResult::InitFailed { retryable } => { StreamOutcome::Fatal => {
// stream_llm 已 emit AiError。决定是否重试(仅 Init 失败可重试,决策 a1) // Fatal(4xx 非429/鉴权/参数错):立即放弃整轮 fallback
let is_last = retry_attempt >= max_retries; // stream_llm 已 emit AiError,仅做 guard.reset + return。
let now = tokio::time::Instant::now(); guard.reset().await;
let budget_exhausted = now >= retry_deadline; return;
// Fatal(4xx 非429/鉴权/参数错)立即放弃,不浪费预算重试
if !retryable {
tracing::warn!(
conv_id = %conv_id,
attempt = retry_attempt + 1,
"[ai] 流前失败 Fatal(4xx/鉴权),立即放弃不重试",
);
guard.reset().await;
return;
}
if is_last || budget_exhausted {
// 重试耗尽或预算耗尽:错误已 emit,结束
tracing::warn!(
conv_id = %conv_id,
attempt = retry_attempt + 1,
budget_exhausted = budget_exhausted,
"[ai] 流前失败重试{}({}次),放弃",
if budget_exhausted { "预算耗尽" } else { "耗尽" },
max_retries + 1,
);
guard.reset().await;
return;
}
// CR-30-1: 退避复用 retry::backoff_delay(attempt+1)(1s→2s→4s + ±20% jitter),
// 不再使用纯指数 `1u64 << retry_attempt`。attempt 传 retry_attempt+1
// (retry::backoff_delay 是 1-based,第 1 次重试 ~1s,第 2 次 ~2s,第 3 次 ~4s)。
let delay = retry::backoff_delay((retry_attempt + 1) as u32)
.min(retry_deadline.saturating_duration_since(now));
let total_retry = retry_attempt + 1; // 1-based 当前是第几次尝试
tracing::warn!(
conv_id = %conv_id,
attempt = total_retry,
max_attempts = max_retries + 1,
delay_ms = delay.as_millis() as u64,
"[ai] 流前失败(Retryable),{}ms 后重试 ({}/{})",
delay.as_millis(), total_retry, max_retries + 1
);
// F-260616-07(d): emit 重试提示事件,前端在错误气泡内显示「重试 n/m」
// (对齐决策"错误气泡内更新")。CR-30-2: useAiEvents.ts 补 case 处理。
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiStreamRetry {
attempt: total_retry as u32,
max_attempts: (max_retries + 1) as u32,
conversation_id: Some(conv_id.clone()),
});
tokio::time::sleep(delay).await;
continue;
} }
} }
} }
@@ -542,8 +712,13 @@ pub(crate) async fn run_agentic_loop(
match outcome { match outcome {
Some(r) => r, Some(r) => r,
None => { None => {
// 不应到达:InitFailed 非 Fatal/未耗尽即重试,耗尽/Fatal 已 return; // 全 candidate 耗尽(含单 provider 场景):沿用原单 provider 失败语义。
// Complete/Partial 写入 outcome 后 break。防御兜底 // 各 candidate 的 stream_llm 已 emit 最后一条 AiError,此处仅复位退出
tracing::warn!(
conv_id = %conv_id,
candidates_tried = candidate_chain.len(),
"[ai] 全 provider 候选流前失败重试耗尽,放弃本轮",
);
guard.reset().await; guard.reset().await;
return; return;
} }
@@ -598,10 +773,10 @@ pub(crate) async fn run_agentic_loop(
return; return;
} }
// stream 结束立即释放 permit,后续工具执行不受限流(本地操作无 RPM 成本) // stream 结束立即释放 permit,后续工具执行不受限流(本地操作无 RPM 成本)
// per-provider permit(_provider_permit)在 candidate 循环内随作用域 Drop 自动释放,
// 此处仅显式释放 global/per_conv(F-04b:per-provider 改为循环内取)。
drop(_global_permit); drop(_global_permit);
drop(_per_conv_permit); drop(_per_conv_permit);
// F-260614-04: per-provider permit 为 None 时 drop None 无副作用;Some 时释放槽。
drop(_provider_permit);
// 累加本轮 token:provider 流式 usage 的 prompt_tokens 为 0 时(GLM 等),用预估输入兜底 // 累加本轮 token:provider 流式 usage 的 prompt_tokens 为 0 时(GLM 等),用预估输入兜底
let round_prompt = if round_usage.prompt_tokens == 0 { estimated_prompt } else { round_usage.prompt_tokens }; let round_prompt = if round_usage.prompt_tokens == 0 { estimated_prompt } else { round_usage.prompt_tokens };