//! 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}; use df_storage::db::Database; use df_storage::models::AiProviderRecord; use crate::state::{AppState, LlmConcurrency}; 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; use super::title::{ensure_conversation_title, spawn_ensure_title}; use super::audit::process_tool_calls; use super::{AiChatEvent, AiSession}; /// Agentic 循环最大迭代次数 /// /// 默认 10 轮。未来可配置接入点:接入 AppState(新增 `agent_config` 字段)/df-storage settings KV /// 表后,改为从配置读(默认值仍为 10)。当前项目无 agent 配置位(AppState/df-storage config 表 /// 均无 agent 配置槽),故暂以常量承载,避免引入 AppState 新字段等大改(违反 P2 零行为变边界)。 /// 接入路径:run_agentic_loop 签名增 `max_iterations: usize` 参数,调用方从 AppState 读取透传。 pub const MAX_AGENT_ITERATIONS: usize = 10; // ============================================================ // 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>, done: bool, } impl GeneratingGuard { fn new(session: Arc>) -> 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; }); } } } /// Agentic 循环:流式接收 → 工具执行 → 结果回传 LLM → 循环 /// /// 退出条件: /// - LLM 只返回文本(无 tool_calls)→ 正常结束 /// - 有工具需要审批 → 暂停循环(generating 保持 true),等 ai_approve 恢复 /// - 达到最大迭代次数 → 正常结束 pub(crate) async fn run_agentic_loop( session_arc: Arc>, tools_arc: Arc, db: Arc, app_handle: AppHandle, provider_config: AiProviderRecord, system_prompt: String, conv_id: String, knowledge_config: crate::state::KnowledgeConfig, llm_concurrency: LlmConcurrency, ) { // B-260615-09: generating 状态由 RAII guard 收敛复位(正常 exit 显式 reset;panic/异常 Drop 兜底) let mut guard = GeneratingGuard::new(session_arc.clone()); // 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 = 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, 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 请求" ); let tool_defs = tools_arc.tool_definitions(); // 停止信号副本:stream_llm 与每轮迭代共享读取,避免重复加锁 // notify 同取一份 Arc 引用(B-260615-14):stream_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; for iteration in 0..MAX_AGENT_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(), conversation_id: Some(conv_id.clone()) }); return; } // B-260615-11: 旧 loop 污染防护——每轮开始校验对话一致性。 // 用户新建/切换对话后 active_conversation_id 变更,本 loop(conv_id 快照)成陈旧, // 继续跑会往新对话 push 消息/pending 造成污染。检测到即退出(guard Drop 复位 generating)。 { let 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; } } // 新一轮通知前端(第二轮起),前端需新建 assistant 消息 if iteration > 0 { let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiAgentRound { round: (iteration + 1) as u32, conversation_id: Some(conv_id.clone()), }); } // 构建请求消息(超预算时自动裁剪旧消息,保护工具调用三元组 + 最近 6 条) let messages = { let session = session_arc.lock().await; let sys_tokens = TokenEstimator::default().estimate_text(&system_prompt); 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,后段用它补) let estimated_prompt: u32 = { let est = TokenEstimator::default(); messages.iter().map(|m| est.estimate_message(m)).sum() }; let request = CompletionRequest { model: provider_config.default_model.clone(), messages, 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, }; // LLM 并发限流(全局 + 单对话双层),仅覆盖 stream_llm 调用本身; // 工具执行(process_tool_calls)是本地操作无 RPM 成本,permit 在 stream 后立即释放避免占槽 let _global_permit = llm_concurrency.acquire_global().await; let _per_conv_permit = llm_concurrency.acquire_per_conv().await; // 流式接收(内部处理 idle timeout / 断连检测 / 停止信号) let (full_text, tool_calls_acc, round_usage) = match stream_llm(&*provider, request, &app_handle, &stop_flag, ¬ify, &conv_id).await { Some(result) => result, None => { // 错误已在 stream_llm 中 emit,直接结束 guard.reset().await; return; } }; // stream 结束立即释放 permit,后续工具执行不受限流(本地操作无 RPM 成本) 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 = tool_calls_acc.keys().copied().collect(); order.sort_unstable(); let ai_tool_calls: Vec = 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(provider_config.default_model.clone()); session.messages.push(msg); } else if !full_text.is_empty() { let mut msg = ChatMessage::assistant(&full_text); msg.model = Some(provider_config.default_model.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(&provider_config.default_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(), 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(&provider_config.default_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):异常中断,提示用户 // 与 break 正常收敛(break→converged=true)区分:这里仍走入库+Completed,但前置发 AiError 警示 if !converged { tracing::warn!( conv_id = %conv_id, max_iter = MAX_AGENT_ITERATIONS, "[ai] agentic 循环达最大轮次(MAX_AGENT_ITERATIONS={})仍未收敛,可能未完成", MAX_AGENT_ITERATIONS, ); let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError { error: format!("达到最大轮次({} 轮),Agent 可能未完成(末轮工具结果未回传模型)", MAX_AGENT_ITERATIONS), conversation_id: Some(conv_id.clone()), }); } // 正常完成 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(); tauri::async_runtime::spawn(async move { save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&provider_config.default_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(), 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) { let (is_generating, has_pending, pending_conv_id) = { 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()); (session.generating, !session.pending_approvals.is_empty(), pending_conv_id) }; let should_continue = is_generating && !has_pending; if !should_continue { // generating=false(被 stop)或仍有审批(pending_approvals 非空): // 统一 emit AiCompleted 标当前轮收敛,清前端 streaming。 // 轮 token 已在前序 AiCompleted/AiApprovalResult 流程落库,此处零 token 上报仅作收敛信号。 if 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 pending_conv_id { Some(cid) => cid, None => { let session = state.ai_session.lock().await; session.active_conversation_id.clone().unwrap_or_default() } }; let _ = app.emit("ai-chat-event", AiChatEvent::AiCompleted { total_tokens: 0, prompt_tokens: 0, completion_tokens: 0, 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, 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, conv_id) = { let session = state.ai_session.lock().await; let lang = session.agent_language.clone().unwrap_or_else(|| "zh-CN".to_string()); let conv_id = session.active_conversation_id.clone().unwrap_or_default(); (lang, conv_id) }; 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(); // 恢复循环前通知前端新建 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).await; }); }