Files
DevFlow/src-tauri/src/commands/ai/agentic.rs

1101 lines
61 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Agentic 循环 — 流式接收 → 工具执行 → 结果回传 LLM → 循环
use std::sync::Arc;
use std::sync::atomic::Ordering;
use tauri::{AppHandle, Emitter};
use tokio::sync::Mutex;
use df_ai::ai_tools::AiToolRegistry;
use df_ai::context::TokenEstimator;
use df_ai::provider::{ChatMessage, CompletionRequest, LlmProvider};
// CR-30-1: 复用 retry::backoff_delay(jitter 1s→2s→4s) + is_status_retryable(Fatal 分类)
// 实现流前失败重试退避对齐(决策 F-260616-07 a1),避免重写退避逻辑。
use df_ai::retry;
// F-01 阶段5: 智能路由 helper + TaskRequirements + 维度枚举。
// 调用点构造 TaskRequirements(主对话:needs_tool_use=true,
// 当前仅 Text 模态;后续多模态接入时检测消息内 Part/Image 追加 Vision),
// 经 select_model_id 在 provider.model_configs 池中选最优;池空/无匹配兜底 default_model。
// 注:路由已解耦(B-260618-03),cost_tier/intelligence 不再参与硬路由。
use df_ai::router::{
select_model_id, Modality, TaskRequirements,
};
use df_storage::db::Database;
use df_storage::models::AiProviderRecord;
use crate::state::{AppState, LlmConcurrency};
use super::audit::process_tool_calls;
use super::compress::compress_via_llm;
use super::conversation::{save_conversation, TokenAccumulator};
use super::knowledge_inject::maybe_spawn_extraction;
use super::prompt::{build_system_prompt, get_active_provider};
use super::stream_recv::{stream_llm, StreamResult};
use super::title::{ensure_conversation_title, spawn_ensure_title};
use super::{AiChatEvent, AiSession, ErrorType};
/// Agentic 循环默认最大迭代次数(可配置项的默认值)
///
/// 默认 10 轮。F-260616-01 已接入配置:AppState.agent_max_iterations(Arc<AtomicUsize>) +
/// ai_set_agent_max_iterations command + Settings.vue 数字配置。调用方在 loop 入口
/// load AtomicUsize 快照后透传 `max_iterations: usize` 形参,当前 loop 锁定边界,
/// 热改下次发消息生效(与 llm_concurrency 传 Arc 实时反映的区别)。
pub const DEFAULT_MAX_AGENT_ITERATIONS: usize = 10;
/// 压缩保护区条数(对齐 clear/compress IPC 的 PROTECT_COUNT=6)。
/// protect_start = len.saturating_sub(PROTECT_COUNT):保护区内的最近 6 条(含本轮三元组)
/// 不参与压缩,避免压缩正在使用的活跃消息。
const PROTECT_COUNT: usize = 6;
/// 流式对话失败自动重试默认次数F-260616-07 / 决策 a1
///
/// 默认 3 次(初次 + 2 次重试)。复用 retry::backoff_delay 退避(1s→2s→4s+jitter) +
/// retry::is_status_retryable Fatal 分类(4xx 非429 立即放弃) + 30s 总预算。
/// 只重试流前失败Init Err未输出任何 token流中途失败MidStream已输出
/// partial_text不重试——保文入库 + AiCompleted(incomplete=true) + 系统提示网络中断。
pub const DEFAULT_MAX_AGENT_RETRIES: usize = 3;
// ============================================================
// B-260615-09: generating 状态 RAII guard
// ============================================================
/// generating 复位 RAII guard,取代散布的手动 `session.generating = false`。
///
/// 两路复位:
/// - 正常路径:exit 点显式 `reset().await` 即时复位(emit 前调,保证"复位→emit"顺序,
/// 前端收事件时后端已可接下一条)。
/// - 异常路径(panic/未走正常 return):Drop 兜底 spawn 复位,防 generating 永真卡死前端。
///
/// 注:try_continue_agent_loop 不用 guard——其 should_continue=false 路径需保持
/// generating=true(审批等待态),全函数 guard 会误复位;该函数单点 provider-Err 复位保持手动。
struct GeneratingGuard {
session: Arc<Mutex<AiSession>>,
done: bool,
}
impl GeneratingGuard {
fn new(session: Arc<Mutex<AiSession>>) -> Self {
Self { session, done: false }
}
/// 显式复位 generating=false。emit 前调用保证顺序。幂等。
async fn reset(&mut self) {
if !self.done {
self.session.lock().await.generating = false;
self.done = true;
}
}
/// 解除 Drop 兜底复位但不复位 generating。审批等待 return 路径调用:
/// 保持 generating=true 留 try_continue 续生成,同时 Drop 因 done=true 跳过复位 spawn。
/// (B-260615-26: 修复审批执行后对话不续生成回归)
fn disarm(&mut self) {
self.done = true;
}
}
impl Drop for GeneratingGuard {
fn drop(&mut self) {
if !self.done {
let session = self.session.clone();
tauri::async_runtime::spawn(async move {
session.lock().await.generating = false;
});
}
}
}
// ============================================================
// 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,
/// DeepSeek thinking 模式推理内容(需回传到下一轮请求)
reasoning_content: Option<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,
last_reasoning_content: &Option<String>,
) -> 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,
reasoning_content: last_reasoning_content.clone(),
};
match stream_llm(&*provider, retry_request, app_handle, stop_flag, notify, conv_id).await {
StreamResult::Complete { text, tool_calls, usage, reasoning_content } => {
return StreamOutcome::Success {
text, tool_calls, usage,
incomplete: false,
resolved_model,
reasoning_content,
};
}
StreamResult::Partial { text, tool_calls, usage, reasoning_content } => {
// 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,
reasoning_content,
};
}
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 → 循环
///
/// 退出条件:
/// - LLM 只返回文本(无 tool_calls→ 正常结束
/// - 有工具需要审批 → 暂停循环generating 保持 true等 ai_approve 恢复
/// - 达到最大迭代次数 → 正常结束
pub(crate) async fn run_agentic_loop(
session_arc: Arc<Mutex<AiSession>>,
tools_arc: Arc<AiToolRegistry>,
db: Arc<Database>,
app_handle: AppHandle,
provider_config: AiProviderRecord,
system_prompt: String,
conv_id: String,
knowledge_config: crate::state::KnowledgeConfig,
llm_concurrency: LlmConcurrency,
max_iterations: usize,
max_retries: usize,
start_iteration: usize,
model_override: Option<String>,
) {
// B-260615-09: generating 状态由 RAII guard 收敛复位(正常 exit 显式 reset;panic/异常 Drop 兜底)
let mut guard = GeneratingGuard::new(session_arc.clone());
// F-260614-04 / F-260614-04b: 多 Provider 负载均衡池 — 选主 + fallback 候选列表。
//
// 流程:list_all → ProviderPool::select(按 模型亲和 > weight > is_default 排序)→ 有序 Vec。
// 单 provider 场景:池仅 1 enabled provider → select 返回单元素 Vec → 首位 = 唯一 provider,
// 行为同 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。
// compress_via_llm / 标题 / 后台 spawn 沿用主 candidate(非 fallback 范围,见各调用点注释)。
//
// AiProviderRepo::new 仅 clone Arc<Database>(廉价),不复用 AppState.ai_providers
// (run_agentic_loop 签名只传 Arc<Database>,改签名会牵动 3 调用点 + try_continue)。
let provider_repo = df_storage::crud::AiProviderRepo::new(&db);
let pool_providers: Vec<AiProviderRecord> = match provider_repo.list_all().await {
Ok(v) => v,
Err(e) => {
tracing::warn!(error = %e, "[ai] list_all providers 失败,负载均衡池退化为入参默认 provider(空池兜底)");
Vec::new()
}
};
let ranked_candidates: Vec<AiProviderRecord> = super::provider_pool::ProviderPool::select(
&pool_providers,
// specify 模式(用户指定 model):传 override 作亲和键,ProviderPool 优先选
// 「池中含该 model 的 provider」作 primary,打破下方「router 选模型需 provider_config」
// 的鸡生蛋——override 此时已知(入参 ← session.model_override),无须等 router。
// auto 模式(override=None)→ 全亲和纯权重排序,行为不变(向后兼容)。
model_override.as_deref(),
);
let (primary_provider, candidates): (AiProviderRecord, Vec<AiProviderRecord>) =
match ranked_candidates.split_first() {
Some((first, rest)) => (first.clone(), rest.to_vec()),
None => {
// 空池兜底:用调用方传入的 provider_config 作唯一候选(启动行为不变)。
// candidates 空 → fallback 循环仅跑 primary 一次,等同单 provider 路径。
(provider_config.clone(), Vec::new())
}
};
// 用主候选覆盖入参 provider_config(下游 build_provider / 路由 / 日志均用此)。
// mut:F-04b 切换 candidate 后更新为实际成功所用 provider(供后续 push/save/标题 spawn)。
let mut provider_config = primary_provider;
// FR-S1: resolve→ensure_resolved_key(空 key 早失败)→build_provider 三步统一走工厂
// 空 key 早失败(逻辑见 secret::ensure_resolved_key 单测):避免空 key 发请求吃 401,错误伪装成"API Key 无效"
//
// B-260615-17:resolve 一次复用——原实现 build_provider_for 成功后又独立调 resolve_provider_secret
// 取 key_len(重复 keyring resolve)。现 resolve 一次:既供 key_len 诊断日志,又供 build_provider,
// 去重复 keyring resolve 调用。逻辑等价于 secret::build_provider_for(resolve→ensure→build 三步),
// 仅因 build_provider_for 隐藏 resolved key 无法复用而在此内联(未改 secret.rs 锁边界)。
let resolved_key = super::secret::resolve_provider_secret(&provider_config);
let key_len = resolved_key.len();
let provider: Box<dyn LlmProvider> = match super::secret::ensure_resolved_key(
&provider_config.name, &resolved_key,
) {
Ok(()) => df_ai::build_provider(
&provider_config.provider_type,
&provider_config.base_url,
&resolved_key,
&provider_config.default_model,
),
Err(msg) => {
guard.reset().await;
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: msg,
// ensure_resolved_key 失败 = key 缺失/钥匙串损坏,归 Auth
error_type: Some(ErrorType::Auth),
conversation_id: Some(conv_id.clone()),
});
return;
}
};
// 诊断日志:401/错误时据此定位是 url/type/model/key 哪项问题(只记长度不记明文)
tracing::info!(
provider = %provider_config.name,
provider_type = %provider_config.provider_type,
base_url = %provider_config.base_url,
model = %provider_config.default_model,
key_len = key_len,
"[ai] 发起 LLM 请求"
);
// F-01 阶段5: 主对话路由 — TaskRequirements(needs_tool_use=true)。
// 模态当前仅 Text(图像消息类型未实现,后续多模态接入时检测 Part/Image 追加 Vision)。
// select_model_id None(池空/无匹配)→ 兜底 default_model,行为与接入前一致。
let agentic_req = TaskRequirements {
modalities: vec![Modality::Text],
needs_tool_use: true,
estimated_context: 0,
};
let resolved_model = select_model_id(&agentic_req, &provider_config.model_configs)
.unwrap_or_else(|| provider_config.default_model.clone());
// F-01 阶段6: 用户指定模型 override 穿透(仅主对话生效,标题/扫描/灵感仍走路由)。
// 兜底原则:override 非空且在该 provider model_configs 池中 → 用 override;否则用 resolved_model。
// 绝不让 override 导致无模型(空/不在池 → 落回路由结果,行为不变)。
// 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()
&& provider_config.model_configs.iter().any(|m| m.model_id == id) =>
{
id.to_string()
}
_ => resolved_model,
};
let tool_defs = tools_arc.tool_definitions();
// 停止信号副本stream_llm 与每轮迭代共享读取,避免重复加锁
// notify 同取一份 Arc 引用B-260615-14stream_llm select! 监听 notified() 即时唤醒
let (stop_flag, notify) = {
let session = session_arc.lock().await;
(session.stop_flag.clone(), session.notify.clone())
};
// token 累加器:loop 生命周期内各轮叠加,退出时传 save_conversation(累加模式落库)
let mut tokens = TokenAccumulator::default();
// 收敛标志:仅当 LLM 末轮无 tool_calls 自行 break(正常收敛)时置 true;
// 区分"正常收敛退出"与"达 MAX 被截断退出"——后者末轮 tool_calls 仍非空(tool_result 不再回传 LLM),属异常
let mut converged = false;
// BUG-260617-12: DeepSeek thinking 模式推理内容跨轮透传
let mut last_reasoning_content: Option<String> = None;
// F-260616-13: system_prompt 是 run_agentic_loop 的不变参数(整个 loop 期间文本不变),
// 其 token 估算在 loop 外算一次缓存复用,避免每轮/每次重试重复 estimate_text(低收益优化,行为不变)。
let sys_tokens = TokenEstimator::default().estimate_text(&system_prompt);
for iteration in start_iteration..max_iterations {
// 用户请求停止 → 收尾退出(已生成文本已在上一轮入库)
if stop_flag.load(Ordering::SeqCst) {
let usage = df_ai::provider::TokenUsage {
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
// 入口 stop:本轮可能尚未 stream(首轮即停),不记 model——避免把未实际生成的 model 写入 models 数组
save_conversation(&session_arc, &db, &conv_id, Some(&usage), None).await;
// 标题生成后台化:不阻塞 Completed emit失败有 extract_title 兜底)
spawn_ensure_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, &llm_concurrency);
guard.reset().await;
// generating 复位后再 emit Completed保证前端收事件时后端已可接下一条(发送队列续发不被"正在生成中"拒绝)
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiCompleted { total_tokens: usage.total_tokens, prompt_tokens: tokens.prompt(), completion_tokens: tokens.completion(), incomplete: None, conversation_id: Some(conv_id.clone()) });
return;
}
// B-260615-11: 旧 loop 污染防护——每轮开始校验对话一致性。
// 用户新建/切换对话后 active_conversation_id 变更,本 loop(conv_id 快照)成陈旧,
// 继续跑会往新对话 push 消息/pending 造成污染。检测到即退出(guard Drop 复位 generating)。
{
let mut session = session_arc.lock().await;
if session.active_conversation_id.as_deref() != Some(conv_id.as_str()) {
tracing::warn!(
stale_conv = %conv_id,
active_conv = ?session.active_conversation_id,
"[ai] 对话已切换,旧 loop 退出(B-260615-11)避免污染新对话"
);
return;
}
// F-260616-11 决策 a: 累计 iteration 计数(「下一轮起算值」= 当前轮+1)。
// 审批等待/达 max 暂停退出时,本字段即「当前轮+1」;审批续跑 ai_approve 读此值作
// start_iteration 透传,实现 iteration 累计不重置(防多次审批反复跑满 max 致 token 失控)。
session.iteration_used = iteration + 1;
}
// 新一轮通知前端(第二轮起),前端需新建 assistant 消息
if iteration > 0 {
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiAgentRound {
round: (iteration + 1) as u32,
conversation_id: Some(conv_id.clone()),
});
}
// F-15 阶段3: 自动压缩(智能裁剪)——在 build_for_request 之前预处理。
//
// 触发条件:history_tokens > budget*0.6 且 保护区外有可压缩消息 且 未在压缩中。
//
// 流程(对齐阶段2 ai_chat_compress_context IPC 的 read-but-don't-mutate 模式):
// ① 读 active 克隆(不改 status / 不扣 token)→ 喂 LLM 出摘要;
// ② LLM 成功 → compress_old_messages(标 compressed + 扣 token)+ insert_at(摘要 system);
// ③ LLM 失败 → 消息状态完全不变(未改 status / 未扣 token),降级走原 build_for_request 裁剪。
//
// 口径决策(任务让"你判断"):选**延迟 mutate(成功才改)**而非"失败回滚 status"。
// 理由:ContextManager.history_tokens 字段私有、无 set_history_tokens 公开接口;
// 若先 compress_old_messages(扣 token)再 LLM,失败回滚需精确恢复 history_tokens,
// 但 ChatMessage 克隆不含 token_count,无法等量加回——回滚 token 不精确。
// 延迟 mutate 则失败时零副作用(消息状态/token 完全不变),语义最干净。
// 注:延迟 mutate 的窗口(active_msgs 读出→LLM 出摘要期间)不持锁,但 loop 串行无并发
// (本函数独占 session_arc,工具执行/审批分支在 stream 之后),故此窗口内 messages 不变。
//
// 安全(FR-S1):复用 loop 顶部已 build+验证 的 provider(不再 build_provider_for 重复 resolve
// keyring),api_key 经 df_storage::secret 闭环;summary/error payload/日志均不含 api_key。
// is_compressing 防重入:set_compressing(true/false) 成对(LLM 调用前后均复位)。
// 单轮问答(history_tokens 未超 0.6*budget)不触发,零行为变化。
let prev_compressing = session_arc.lock().await.messages.is_compressing();
if !prev_compressing {
// 读触发条件(history_tokens / budget / has_compressible_messages),持锁快照判断。
let (should_compress, protect_start, pre_compress_tokens) = {
let session = session_arc.lock().await;
let mgr = &session.messages;
let protect_start = mgr.len().saturating_sub(PROTECT_COUNT);
let history_tokens = mgr.history_tokens();
let budget = mgr.budget_limit();
// 触发阈值 0.6*budget(整数比避免浮点):budget*6/10 < history_tokens
let should = protect_start > 0
&& (budget as u64) * 6 / 10 < history_tokens as u64
&& mgr.has_compressible_messages(protect_start);
(should, protect_start, history_tokens)
};
if should_compress {
// emit 压缩开始 + 置位防重入 + 读 active 克隆(不改 status)。
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiCompressing {
conversation_id: Some(conv_id.clone()),
});
let (active_msgs, lang) = {
let mut session = session_arc.lock().await;
session.messages.set_compressing(true);
// 读 active 克隆(不改 status):filter is_active,LLM 失败则消息状态完全不变。
let active_msgs: Vec<ChatMessage> = session.messages.messages_mut()
[..protect_start]
.iter()
.filter(|t| t.message.is_active())
.map(|t| t.message.clone())
.collect();
let lang = session.agent_language.clone()
.unwrap_or_else(|| "zh-CN".to_string());
(active_msgs, lang)
};
// 压缩调用(复用 loop 顶部已 build 的 provider,api_key 经 secret 闭环)。
// 成功 → Some(summary);失败 → Err;无 active 可压缩(active_msgs 空)→ 视为 noop。
let compress_outcome: Result<Option<String>, String> = if active_msgs.is_empty() {
Ok(None)
} else {
compress_via_llm(
provider.as_ref(),
&provider_config,
active_msgs,
&lang,
&llm_concurrency,
).await.map(Some)
};
match compress_outcome {
Ok(Some(summary)) => {
// LLM 成功 → 标 compressed(扣 token)+ 摘要 system 插首位 + set_compressing(false)。
// compress_old_messages 幂等:此时 status 仍是 active(本流程未先标),它会把
// [..protect_start] 内 active 标 compressed 并扣 token。返回的 cloned 与之前读的
// active_msgs 等价(LLM 调用期间 messages 不变,见上方口径决策注)。
{
let mut session = session_arc.lock().await;
let _compressed = session.messages.compress_old_messages(protect_start);
session.messages.insert_at(0, ChatMessage::system(&summary));
session.messages.set_compressing(false);
}
tracing::info!(
conv_id = %conv_id,
iteration,
pre_tokens = pre_compress_tokens,
"[ai] 自动压缩成功,摘要已插首位"
);
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiCompressed {
conversation_id: Some(conv_id.clone()),
summary,
});
}
Ok(None) => {
// 保护区外无 active 可压缩(已全 compressed/archived)→ noop,仅复位 is_compressing。
session_arc.lock().await.messages.set_compressing(false);
}
Err(e) => {
// LLM 失败 → 消息状态完全不变(未改 status / 未扣 token)。
// set_compressing(false) 复位 + emit AiError(message 不含 api_key)。
// 不阻塞 loop:继续走下方 build_for_request 原裁剪路径(保最近 6 条)。
session_arc.lock().await.messages.set_compressing(false);
tracing::warn!(
conv_id = %conv_id,
error = %e,
"[ai] 自动压缩失败,降级走原裁剪(build_for_request)"
);
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: format!("自动上下文压缩失败,已降级为普通裁剪: {}", e),
error_type: Some(ErrorType::Unknown),
conversation_id: Some(conv_id.clone()),
});
}
}
}
}
// 构建请求消息(超预算时自动裁剪旧消息,保护工具调用三元组 + 最近 6 条)
// F-260616-13: sys_tokens 已在 loop 外缓存;本块构建的 messages 在本轮重试循环中复用
// (本轮 stream_llm 不持 session_arc、不改 messages,重试无 push 发生,重建等价于复用)。
let messages = {
let session = session_arc.lock().await;
let (history_msgs, _trimmed) = session.messages.build_for_request(sys_tokens);
let mut msgs = vec![ChatMessage::system(&system_prompt)];
msgs.extend(history_msgs);
msgs
};
// 预估输入 token(兜底:部分 provider 如 GLM 流式 usage 不报 prompt_tokens,后段用它补)
// 注:stream_one_provider 内每次重试重建 request(因 provider.stream 消费 body),
// 此处不再预构建 request(旧 request 变量已废弃),仅保留 messages 供 estimated_prompt。
let estimated_prompt: u32 = {
let est = TokenEstimator::default();
messages.iter().map(|m| est.estimate_message(m)).sum()
};
// LLM 并发限流(全局 + 单对话双层),仅覆盖 stream_llm 调用本身;
// 工具执行(process_tool_calls)是本地操作无 RPM 成本,permit 在 stream 后立即释放避免占槽
//
// CR-30-1 / F-260616-07 / 决策 a1: 流前失败(Init Err)重试,流中途失败(MidStream
// Partial)不重试保文。重试退避复用 retry::backoff_delay(1s→2s→4s±20% jitter) +
// retry::is_status_retryable Fatal 分类(stream_recv classify_status_or_class 镜像,
// 4xx 非429 立即放弃) + 30s 总挂钟预算。重试期间持有 permit 不释放(防新请求挤占)。
//
// F-260614-04 / F-260614-04b: global/per_conv permit 整轮迭代共享(不随 candidate 切换
// 释放——并发语义是「全局/单对话」级,与具体 provider 无关);per-provider permit 改为
// 在 candidate 循环内取(切换 candidate 时释放旧取新,避免占用未用 provider 的槽)。
// 两 permit 均 Drop 释放(L803-804 显式 drop _global/_per_conv)。
let _global_permit = llm_concurrency.acquire_global().await;
let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
// 重试总预算(挂钟,含 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);
// 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, 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;
'candidate: for candidate in &candidate_chain {
// per-provider permit(可选):set_provider_caps 未配置时返回 None(单 provider 零变化);
// 配置后取额外 permit 防单 provider 被打满(限流 429)。切换 candidate 时上一 permit
// 随 _provider_permit 绑定作用域 Drop 释放。
let _provider_permit = llm_concurrency.acquire_for_provider(&candidate.id).await;
match 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 {
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 => {
// 本 candidate 重试耗尽(retryable)。切下一 candidate 继续尝试。
// candidates 空(单 provider)时此即「耗尽 return」语义——
// 跳出循环后 outcome 仍 None,落入下方「全耗尽 return」兜底。
tracing::warn!(
conv_id = %conv_id,
provider = %candidate.name,
"[ai] 候选 {} 流前失败重试耗尽,F-04b 切换下一 provider",
candidate.name,
);
continue 'candidate;
}
StreamOutcome::Fatal => {
// Fatal(4xx 非429/鉴权/参数错):立即放弃整轮 fallback。
// stream_llm 已 emit AiError,仅做 guard.reset + return。
guard.reset().await;
return;
}
}
}
match outcome {
Some(r) => r,
None => {
// 全 candidate 耗尽(含单 provider 场景):沿用原单 provider 失败语义。
// 各 candidate 的 stream_llm 已 emit 最后一条 AiError,此处仅复位退出。
tracing::warn!(
conv_id = %conv_id,
candidates_tried = candidate_chain.len(),
"[ai] 全 provider 候选流前失败重试耗尽,放弃本轮",
);
guard.reset().await;
return;
}
}
};
// CR-30-2 / UX-2025-04 / 决策 a1: MidStream 保文路径——partial_text 已接收,
// 入库为正常 assistant 消息 + emit AiCompleted(incomplete=true) + 追加系统提示消息。
// 不走 AiError(非异常中断,已有可用文本),不重试(决策 a1)。
if incomplete {
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 + round_usage.completion_tokens } else { round_usage.total_tokens },
};
tokens.add(usage.prompt_tokens, usage.completion_tokens);
// 追加 partial assistant 消息(若无 tool_calls 且有文本)
{
let mut session = session_arc.lock().await;
if session.active_conversation_id.as_deref() != Some(conv_id.as_str()) {
tracing::warn!(
stale_conv = %conv_id,
active_conv = ?session.active_conversation_id,
"[ai] MidStream 保文后对话已切换,丢弃本轮 push(B-260615-11)"
);
return;
}
if !full_text.is_empty() {
let mut msg = ChatMessage::assistant(&full_text);
msg.model = Some(resolved_model.clone());
// BUG-260617-12: MidStream 保文也回填 reasoning_content
msg.reasoning_content = round_reasoning_content.clone();
session.messages.push(msg);
// 追加系统提示消息:响应因网络中断不完整(对齐决策 a1 系统提示机制)
let mut notice = ChatMessage::system("⚠ 响应因网络中断不完整,以上为已接收的部分内容。可重新发送以获取完整回复。");
notice.model = Some(resolved_model.clone());
session.messages.push(notice);
}
}
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
// 标题生成后台化(失败有 extract_title 兜底)
spawn_ensure_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, &llm_concurrency);
guard.reset().await;
// generating 复位后再 emit Completed(incomplete=true):前端据此标记消息为不完整
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiCompleted {
total_tokens: usage.total_tokens,
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
conversation_id: Some(conv_id.clone()),
incomplete: Some(true),
});
return;
}
// stream 结束立即释放 permit,后续工具执行不受限流(本地操作无 RPM 成本)
// per-provider permit(_provider_permit)在 candidate 循环内随作用域 Drop 自动释放,
// 此处仅显式释放 global/per_conv(F-04b:per-provider 改为循环内取)。
drop(_global_permit);
drop(_per_conv_permit);
// 累加本轮 token:provider 流式 usage 的 prompt_tokens 为 0 时(GLM 等),用预估输入兜底
let round_prompt = if round_usage.prompt_tokens == 0 { estimated_prompt } else { round_usage.prompt_tokens };
tokens.add(round_prompt, round_usage.completion_tokens);
// 追加 assistant 消息到历史
let has_tool_calls = !tool_calls_acc.is_empty();
{
let mut session = session_arc.lock().await;
// B-260615-11: push 前再校验(stream_llm 期间用户可能新建对话)。
// 读端读到被 clear 的空历史不致命,但 push 写回新对话是污染,必须挡。
if session.active_conversation_id.as_deref() != Some(conv_id.as_str()) {
tracing::warn!(
stale_conv = %conv_id,
active_conv = ?session.active_conversation_id,
"[ai] stream 后对话已切换,丢弃本轮 push(B-260615-11)避免污染新对话"
);
return;
}
if has_tool_calls {
let mut order: Vec<u32> = tool_calls_acc.keys().copied().collect();
order.sort_unstable();
let ai_tool_calls: Vec<df_ai::provider::ToolCall> = order.iter()
.map(|i| {
let draft = &tool_calls_acc[i];
df_ai::provider::ToolCall::new(&draft.id, &draft.name, &draft.args)
})
.collect();
let mut msg = ChatMessage::assistant_with_tools(&full_text, ai_tool_calls);
msg.model = Some(resolved_model.clone());
// BUG-260617-12: 回填 reasoning_content 供下一轮请求透传
msg.reasoning_content = last_reasoning_content.clone();
session.messages.push(msg);
} else if !full_text.is_empty() {
let mut msg = ChatMessage::assistant(&full_text);
msg.model = Some(resolved_model.clone());
msg.reasoning_content = last_reasoning_content.clone();
session.messages.push(msg);
}
}
// 停止信号:已生成文本入库后退出,不再执行后续工具调用
if stop_flag.load(Ordering::SeqCst) {
let usage = df_ai::provider::TokenUsage {
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
// 标题生成后台化:不阻塞 Completed emit失败有 extract_title 兜底)
spawn_ensure_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, &llm_concurrency);
guard.reset().await;
// generating 复位后再 emit Completed保证前端收事件时后端已可接下一条(发送队列续发不被"正在生成中"拒绝)
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiCompleted { total_tokens: usage.total_tokens, prompt_tokens: tokens.prompt(), completion_tokens: tokens.completion(), incomplete: None, conversation_id: Some(conv_id.clone()) });
return;
}
// 无工具调用 → 最终文本响应,正常收敛退出
if !has_tool_calls { converged = true; break; }
// 处理工具调用Low 自动执行 / Medium+High 待审批)
let pending_count = {
let mut session = session_arc.lock().await;
process_tool_calls(&mut session, tool_calls_acc, &tools_arc, &db, &app_handle, &conv_id).await
};
// 有待审批 → 暂停循环,等待用户审批后通过 ai_approve → try_continue_agent_loop 恢复
if pending_count > 0 {
let usage = df_ai::provider::TokenUsage {
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
// B-260615-26: 审批等待 return 前 disarm guard——保持 generating=true 留 try_continue 续生成,
// 同时 Drop 因 done=true 跳过复位 spawn(避免误复位审批态 generating 致 ai_approve→try_continue 不续)
guard.disarm();
return; // generating 保持 true
}
// 全部自动执行完成 → 继续下一轮
}
// 达 MAX 未收敛(LLM 末轮仍想调工具被截断,末轮 tool_result 不再回传 LLM):转入暂停态询问用户
// F-260616-03:不再 emit AiError + 走完成流程,改为 emit AiMaxRoundsReached + 保持 generating=true
// (仿审批等待 L313-316),等用户点继续(ai_continue_loop → try_continue_agent_loop 再跑 max_iterations 轮)
// 或点停止(ai_stop_loop → 走完成流程)。try_continue 续跑 iteration 由调用方传 start_iteration 决定:
// 审批续跑累计(F-260616-11 决策 a,防多次审批反复跑满 max 致 token 失控,传 session.iteration_used),
// 达 max 续跑重计(F-260616-03 决策 a,用户点继续=授权重来,传 0 + 重置 iteration_used)。
if !converged {
tracing::warn!(
conv_id = %conv_id,
max_iter = max_iterations,
"[ai] agentic 循环达最大轮次(max_iterations={})仍未收敛,转暂停态询问用户(F-260616-03)",
max_iterations,
);
// 轮 token 落库(保留末轮已生成内容,续跑/停止都据此累加)
let usage = df_ai::provider::TokenUsage {
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
// 暂停态保持 generating=true(防其他 send 抢占,仿审批),disarm guard 跳过 Drop 兜底复位
guard.disarm();
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiMaxRoundsReached {
conversation_id: Some(conv_id.clone()),
});
return; // generating 保持 true,等 ai_continue_loop / ai_stop_loop
}
// 正常完成
let usage = df_ai::provider::TokenUsage {
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
// 落库 + 标题 + 知识提炼打包后台化:不阻塞 generating 复位与 Completed 事件
// save 先行(extract/title 都读已落库消息);extract 内部 fire-and-forget,与 title 可能并发
// (均受 per_conv 信号量约束,读写不同字段互不干扰)
// 并发取舍:与新对话新 loop 的 save 存在低概率并发 upsert最多丢少量 token 累加(非功能错误,可接受)
let usage_total = usage.total_tokens;
{
let session_arc = session_arc.clone();
let db = db.clone();
let conv_id = conv_id.clone();
let provider_config = provider_config.clone();
let knowledge_config = knowledge_config.clone();
let app_handle = app_handle.clone();
let llm_concurrency = llm_concurrency.clone();
let resolved_model = resolved_model.clone();
tauri::async_runtime::spawn(async move {
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
// 知识提炼:需读已落库的对话消息,故在 save 之后
if let Err(e) = maybe_spawn_extraction(&session_arc, &db, &conv_id, &provider_config, &knowledge_config, llm_concurrency.clone()).await {
tracing::warn!("知识提炼触发失败(非阻断): {}", e);
}
ensure_conversation_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, llm_concurrency).await;
});
}
guard.reset().await;
// generating 复位后再 emit Completed落库/标题/提炼已在后台,前端立即感知完成
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiCompleted { total_tokens: usage_total, prompt_tokens: tokens.prompt(), completion_tokens: tokens.completion(), incomplete: None, conversation_id: Some(conv_id.clone()) });
}
/// 检查是否所有待审批已处理,如果是则恢复 agentic 循环
///
/// B-260615-08:所有静默 return 点显式 emit 收尾事件,避免前端 streaming=true 永久卡。
/// 各 return 点的语义判断:
/// 1) should_continue=false(generating 已复位 / pending_approvals 非空):
/// - generating=false → 用户点了停止(ai_chat_stop 复位)或会话已结束,emit AiCompleted 标当前轮收敛
/// (streaming=true 由 AiCompleted 清理)
/// - pending_approvals 非空 → 转入审批等待态(其他审批未决),emit AiCompleted 标当前轮结束
/// (前端审批态 watchdog 已 clear,不卡)
/// 2) get_active_provider Err → 无可用 provider(配置丢失/全删),无法续生成,emit AiError
/// (语义:配置错误,用户需设 provider;非 generating 复位可恢复)
///
/// R-PD-6: conv_id 来源从全局 active_conversation_id 解耦到审批所属会话。
/// 触发本函数的 ai_approve 已 remove 触发审批,但 pending_approvals 内剩余审批(若 has_pending)
/// 仍各自携带 conversation_id(审批产生时由 process_tool_calls 写入,业务真相源)。
/// 故 has_pending=true 分支(审批等待态)直接取剩余审批的 conversation_id 做 conv_id,
/// 不读 active_conversation_id 全局单例——该字段在审批等待态(非 generating-only 期)可被
/// ai_chat_stop/clear/switch 并发改写,属竞态耦合。has_pending=false(全部审批已处理,续生成)
/// 分支:审批已被 remove,改为以剩余 pending_approvals 任一 conversation_id 做一致性校验
/// (此处空,校验通过即沿用全局值,该期 generating=true 且 switch 为 readonly 不并发)。
pub(crate) async fn try_continue_agent_loop(app: &AppHandle, state: &AppState, start_iteration: usize) {
// BUG-260617-05: 原 5 次独立 lock().await 造成 TOCTOU 竞态——should_continue=true 判出后、
// spawn 前用户点 stop(ai_chat_stop 复位 generating=false),续跑仍按过时快照继续 spawn。
// 修复:单次 lock 取结构化快照(所有续跑判定所需字段),无锁态判定;spawn 前单次 lock 原子
// 重检 generating 仍为 true 才续跑,stop 后中途插入的直接收敛退出。
let snap = {
let session = state.ai_session.lock().await;
// pending_approvals 中任一审批的 conversation_id:审批等待态(has_pending)下作为 conv_id 来源,
// 取第一个非空值(同一对话的审批 conversation_id 一致,见 process_tool_calls 写入路径)。
let pending_conv_id = session.pending_approvals.values()
.find_map(|a| a.conversation_id.clone());
ContinueSnapshot {
is_generating: session.generating,
has_pending: !session.pending_approvals.is_empty(),
pending_conv_id,
active_conversation_id: session.active_conversation_id.clone(),
agent_language: session.agent_language.clone(),
model_override: session.model_override.clone(),
}
};
let should_continue = snap.is_generating && !snap.has_pending;
if !should_continue {
// generating=false(被 stop)或仍有审批(pending_approvals 非空):
// 统一 emit AiCompleted 标当前轮收敛,清前端 streaming。
// 轮 token 已在前序 AiCompleted/AiApprovalResult 流程落库,此处零 token 上报仅作收敛信号。
if snap.is_generating {
// pending_approvals 非空但 generating 仍 true:转审批态,前端审批态 watchdog 已 clear,不卡
tracing::info!("[ai] try_continue 跳过:仍有待审批,转审批等待态");
} else {
// generating 已复位(用户 stop 或前序循环已 emit Completed):补发 AiCompleted 防前端卡住
tracing::info!("[ai] try_continue 跳过:generating 已复位(被 stop/已结束),补发 AiCompleted 清前端 streaming");
// R-PD-6: 优先用审批所属 conversation_id(审批等待态被 stop 触发,审批仍在 pending_approvals),
// 仅当无任何审批(has_pending=false 且 generating=false)时回退 active_conversation_id。
let conv_id = match snap.pending_conv_id.clone() {
Some(cid) => cid,
None => snap.active_conversation_id.clone().unwrap_or_default(),
};
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompleted {
total_tokens: 0,
prompt_tokens: 0,
completion_tokens: 0,
incomplete: None,
conversation_id: Some(conv_id),
});
}
return;
}
let provider_config = match get_active_provider(state).await {
Ok(p) => p,
Err(e) => {
// 无可用 provider(配置丢失/全删):无法续生成,emit AiError。
// 语义:配置错误,用户需在 Settings 设 provider;generating 复位由 run_agentic_loop 内
// build_provider_for Err 分支处理(同样 emit AiError),此处与之一致。
// 不用 GeneratingGuard:try_continue 的 should_continue=false 路径需保 generating=true(审批等待态),
// 全函数 guard 会误复位。此点单点 provider-Err 复位,语义独立。
let mut session = state.ai_session.lock().await;
session.generating = false;
// R-PD-6: 续生成被拒(provider 缺失)回退 conv_id 优先审批所属;无审批再读全局。
let conv_id = session.pending_approvals.values()
.find_map(|a| a.conversation_id.clone())
.or_else(|| session.active_conversation_id.clone())
.unwrap_or_default();
drop(session);
tracing::warn!(error = %e, "[ai] try_continue 失败:无可用 provider");
let _ = app.emit("ai-chat-event", AiChatEvent::AiError {
error: e,
// 无可用 provider(配置丢失/全删):用户需在 Settings 设 provider,归 ProviderConfig
error_type: Some(ErrorType::ProviderConfig),
conversation_id: Some(conv_id),
});
return;
}
};
// R-PD-6: 续生成路径 conv_id 解耦——has_pending=false 时审批已 remove,无审批 conversation_id 可取;
// 此期 generating=true 且 switchConversation 为 readonly 不并发改 active_conversation_id,
// 故读快照值安全(非竞态期);若 has_pending=true 已在上面 return,不会到此。
let lang = snap.agent_language.clone().unwrap_or_else(|| "zh-CN".to_string());
let conv_id = snap.active_conversation_id.clone().unwrap_or_default();
let system_prompt = build_system_prompt(state, &lang).await;
let session_arc = state.ai_session.clone();
let tools_arc = state.ai_tools.clone();
let db = state.db.clone();
let app_handle = app.clone();
let knowledge_config = state.knowledge_config.lock().await.clone();
let llm_concurrency = state.llm_concurrency.clone();
// F-260616-01: loop 入口 load 快照,当前续生成 loop 锁定边界(热改下次发消息生效)
let max_iterations = state.agent_max_iterations.load(Ordering::SeqCst);
// F-260616-07: 流式失败重试次数快照
let max_retries = state.agent_max_retries.load(Ordering::SeqCst);
// F-01 阶段6: 续跑沿用同一主对话的 model_override(审批续跑/达 max 续跑保持一致)。
let model_override = snap.model_override.clone();
// BUG-260617-05 续: provider 解析/build_system_prompt 期间用户可能点 stop。
// spawn 前单次 lock 原子重检 generating——若已被 stop 复位,收敛退出而非覆盖用户的 stop。
// (run_agentic_loop 入口 GeneratingGuard 会再次置 generating=true,若不重检会抹掉 stop。)
if !state.ai_session.lock().await.generating {
tracing::info!("[ai] try_continue 终止:spawn 前重检 generating 已被 stop 复位,补发 AiCompleted");
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompleted {
total_tokens: 0,
prompt_tokens: 0,
completion_tokens: 0,
incomplete: None,
conversation_id: Some(conv_id.clone()),
});
return;
}
// 恢复循环前通知前端新建 assistant 消息:审批(通过/拒绝)后新一轮文本
// 不应追加到发起工具调用的旧消息,用 AiAgentRound 隔开
let _ = app.emit("ai-chat-event", AiChatEvent::AiAgentRound {
round: 0,
conversation_id: Some(conv_id.clone()),
});
tauri::async_runtime::spawn(async move {
run_agentic_loop(session_arc, tools_arc, db, app_handle, provider_config, system_prompt, conv_id, knowledge_config, llm_concurrency, max_iterations, max_retries, start_iteration, model_override).await;
});
}
/// BUG-260617-05: try_continue_agent_loop 续跑判定所需 session 字段的一次性快照。
/// 单次 lock 取出后无锁态判定,消除多 lock 间其他 IPC(ai_chat_stop/clear/switch)改写 session
/// 致续跑判断基于过时快照的 TOCTOU 竞态。
struct ContinueSnapshot {
is_generating: bool,
has_pending: bool,
pending_conv_id: Option<String>,
active_conversation_id: Option<String>,
agent_language: Option<String>,
model_override: Option<String>,
}