Files
DevFlow/src-tauri/src/commands/ai/agentic/mod.rs
绝尘 eda2a5f887 优化: aichat L2状态机双轨清理+读侧迁移+L3事件总线emit双写
L2 双轨清理:删 guard 本地 state 冗余字段(ConvState 已落 PerConvState
持久化作读侧真相源)。读侧迁移:ai_is_generating/ai_chat_stop 加
CONV_STATE_ENABLED 门控读 conv_state.is_active()。L3 emit 双写:mod.rs
关键事件 AiCompleted/AiError/AiAgentRound 经 ai_event_bus.publish 双写
(EVENT_BUS_ENABLED 门控,与 app.emit 并存非替换,高频 delta 不双写)。

双轨过渡收尾,可一键回退(CONV_STATE_ENABLED off 回退 bool 真相源)。
2026-06-22 01:03:49 +08:00

1748 lines
102 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, Manager};
use tokio::sync::Mutex;
use df_ai::ai_tools::AiToolRegistry;
use df_ai::context::TokenEstimator;
// 阶段2 占位配对完整性:agentic 出口第二道防线断言(深度防御)。
use df_ai::context::ContextManager;
// 改进3 B: 压缩失败兜底关键词摘要(纯函数 extract_keyword_summary)。
// 改进4: 工具结果 view-only 摘要(should_summarize_tool_result / extract_key_info)。
use df_ai::context_helpers::{
extract_key_info, extract_keyword_summary, should_summarize_tool_result,
PLACEHOLDER_INTEGRITY_ENABLED,
};
// 改进2 B:意图收敛工具(LLM 可见 tool_defs 按 intent 过滤,执行路径仍走完整 registry)
// B 路线 Phase 1:plan_hint 接入主 loop——filter_tool_defs_planned 在 filter_tool_defs
// 收敛的扁平子集之上叠加 plan_hint 编排(并行组同批聚拢/顺序依赖源在前),供 LLM 看到
// 一份按编排意图排序的工具列表。feature flag PLANNING_ENABLED(false 默认关)门控接入。
use df_ai::intent::{filter_tool_defs, filter_tool_defs_planned, IntentRecognizer};
use df_ai::provider::{ChatMessage, CompletionRequest, LlmProvider, MessageRole};
// 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 crate::commands::ai::event_bus::AiBusEvent;
use super::audit::process_tool_calls;
use super::compress::compress_via_llm;
use super::conversation::{save_conversation, TokenAccumulator};
use super::knowledge_inject::{inject_knowledge_into_prompt, 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, SessionState};
/// L1 补丁:run_agentic_loop 入口 provider 解析超时保护的内部错误类型。
///
/// 用于把 provider 解析块(list_all + select + resolve + ensure + build)包入
/// `tokio::time::timeout` 后的内部 Err 路由——区分 ensure_resolved_key 失败(走 Auth
/// 错误)与整体超时(走 Unknown 错误)。超时由外层 timeout 的 Err(Elapsed) 单独匹配。
enum ProviderResolveError {
/// ensure_resolved_key 失败(key 缺失/钥匙串损坏):走 Auth 错误分支(对齐原 :412 处理)。
EnsureKeyFailed(String),
}
/// 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;
/// 改进3 B 常量开关:压缩失败时是否启用关键词摘要兜底(默认 true)。
///
/// true(默认):LLM 压缩失败 → compress_old_messages 标 compressed(释放 token)
/// + insert_at(0, system 关键词摘要),保留用户反复提及的主题词作续接锚点。
/// false(回退):回退原裸裁剪行为(不插关键词摘要,仅 build_for_request 裁剪保最近 6 条)。
/// 排障/对比用:置 false 即可观察无兜底时的裁剪效果。
pub const KEYWORD_FALLBACK_ENABLED: bool = true;
/// 改进4 常量开关:tool_result 大输出是否做 view-only 摘要压缩(默认 true)。
///
/// true(默认):build_for_request 后送 stream 前,遍历 history 中 tool_result,
/// 超阈值(content >2KB 或占历史 token >40%)的 content 应用 extract_key_info
/// 替换(保留错误行 + 首尾各 5 行)。**仅 messages clone 视图,不改 ContextManager
/// 持久化**(对齐 sanitize_messages:DB 原始 tool_result 完整保留)。
/// false(关闭):tool_result 原样送 LLM(旧行为)。排障/对比用。
pub const TOOL_RESULT_COMPRESS_ENABLED: bool = true;
/// 改进5 常量开关:是否启用主题切换系统标记(默认 true)。
///
/// true(默认):push 时末两条 user 消息 topic 都非 None 且不同(双高置信)→
/// pending_topic_marker 置位 → agentic loop 顶部读取并 insert 一条 system 软提示
/// `"── 用户已切换话题(从「{old}」到「{new}」),请以新话题为准 ──"`(软提示,不强制 LLM)。
/// false(关闭):loop 顶部跳过读取/insert(置 false 即可观察无主题标记效果)。
/// 保守:双高置信才标(任一 topic None 不标),不强制 LLM(软提示非硬约束)。
pub const TOPIC_MARKER_ENABLED: bool = true;
/// L1 断路器:连续同类工具失败熔断阈值(治 kms 会话 53 轮 0 产出死循环)。
///
/// 背景:agent 无止损,某工具反复同类失败(权限拒绝/路径错误等)仍每轮重试,
/// 耗尽 max_iterations 前 0 产出。机制(非 prompt 教 AI):每轮 process_tool_calls
/// 后取末尾连续 Tool 消息,失败内容前 40 字符归一为 key 计数,同一 key 累计达此阈值 →
/// guard.reset + emit AiError + return 强制熔断,逼用户换思路或人工介入。
///
/// 阈值 3:同类失败 3 次足以判死循环(去重后仍累加,不同错误各自计数互不干扰)。
pub const CIRCUIT_BREAKER_THRESHOLD: u32 = 3;
/// L1 断路器总开关(默认 true)。false → 跳过断路器检查,降级为纯 max_iterations
/// 旧行为(排障/对比/临时关闭用)。机制优先 prompt 说教,每改配开关 + 兜底(关降级旧行为)。
pub const CIRCUIT_BREAKER_ENABLED: bool = true;
/// L1 断路器熔断时是否发结构化求助(aichat 体验与 agent 能力系统化重构 §2.3,2026-06-21)。
///
/// true(默认):熔断 emit AiHelpRequired(结构化求助卡:reason + context + options),
/// 引导用户换思路/授权路径/人工接管(机制优先 prompt 说教,非教 AI 自己止损)。
/// false(兜底回退):熔断仍 emit AiError(旧行为,前端错误气泡),用于求助卡未就绪/
/// 排障/对比。两路保留 guard.reset + return 强制熔断语义不变,仅换前端呈现形态。
/// 配合 CIRCUIT_BREAKER_ENABLED:CIRCUIT_BREAKER_ENABLED=false 时断路器整段跳过,
/// 本开关无意义;CIRCUIT_BREAKER_ENABLED=true 时本开关决定呈现形态。
pub const CIRCUIT_BREAKER_HELP_EVENT: bool = true;
// 阶段2(path_auth 审批链重构):占位配对完整性开关(解 400 orphan)。
//
// 单一真相源:`df_ai::context_helpers::PLACEHOLDER_INTEGRITY_ENABLED`(本模块顶部已 use)。
// 删除本地副本(B 路线改进:避免与 df-ai 同名 const 双源,改一处易漏改另一处)。
//
// 根因:审批挂起占位 tool_result(内容 audit/cache.rs:pending_placeholder_for,带
// `__PENDING__:tc_id` 标记)与其 tool_call 头经 sanitize/裁剪后可能丢配对头 → orphan
// tool_result → deepseek-v4-pro 等端点 400。df-ai 的 sanitize_messages step3.5(反向
// orphan)已豁免保留占位,build_for_request 出口已用 assert_placeholder_pairing 自愈补头。
//
// 该开关控制 agentic loop 末尾(build_for_request + tool_result 压缩后)的**第二道防线**
// 出口断言:在 messages 送 stream 前再过一次 assert_placeholder_pairing(depth-defense,
// 防 build_for_request 到送 stream 之间的转换引入新 orphan)。
//
// true(默认):agentic 出口再断言一次占位配对,失败自愈补头。
// false(回退):agentic 出口不断言(仅依赖 df-ai build_for_request 内部一次自愈,旧行为)。
// 兜底:flag 关→等价改动前(仅 df-ai 内部自愈);view-only 不改持久化。
// ============================================================
// 重构第一批(2026-06-19):GeneratingGuard 抽离到 guard.rs(纯结构搬迁,行为零变更)。
// run_agentic_loop 内仍 `GeneratingGuard::new(...)`,路径从本模块改 super::guard。
// ============================================================
mod guard;
use guard::GeneratingGuard;
// ============================================================
// L2 统一状态机(渐进第一步,2026-06-21):ConvState enum + 转换守卫。
// 设计:generating状态机加固-2026-06-15.md §3 + aichat体验与agent能力系统化重构-2026-06-21.md §3。
//
// 本批范围(渐进):
// - 新建 conv_state.rs(纯逻辑 enum + 守卫 + 单测)。
// - run_agentic_loop 入口桥接:guard 置 generating=true 时同步迁移 ConvState
// (Idle→Generating / Error→Generating),作视图层。CONV_STATE_ENABLED 开关门控。
// - guard.reset 同步 ConvState→Idle(正常退出路径)。
// - **不替换** PerConvState.generating bool(渐进不一次全换,批2+ 逐步迁读侧)。
//
// 开关 CONV_STATE_ENABLED(默认 on):off 降级纯旧 bool 行为(可回退)。
// 兜底:状态机层迁移失败(非法转换)记 warn 不 panic,核心 generating 复位仍走旧 bool。
// ============================================================
pub mod conv_state;
// ============================================================
// 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 重算)。
/// 携带 `error`(最后一条诊断文本),供外层「全 candidate 耗尽」时 emit 最终 AiError。
/// - `Fatal`:流前 `InitFailed{retryable=false}`(4xx 非429/鉴权/参数错)。立即放弃
/// 整个 fallback(对齐 retry.rs Fatal 分类 + provider_pool 文档「不可重试错误
/// 立即放弃不浪费备用 provider」)。携带 `error`,stream_one_provider 内 Fatal 分支
/// 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>,
},
/// UX-260618-15: error 字段供外层全 candidate 耗尽时 emit 最终 AiError(单气泡聚合)。
InitFailedExhausted { error: String },
/// UX-260618-15: error 字段供 stream_one_provider 内 Fatal 分支 emit AiError。
Fatal { error: String },
}
/// 单 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) => {
// UX-260618-15: Fatal 分支 emit 移到外层(stream_one_provider 调用方 Fatal 分支统一 emit)。
// 此处仅返回 error 文本,避免 stream_one_provider 内 emit 与外层 emit 重复(Fatal 终态单 emit)。
return StreamOutcome::Fatal { error: msg };
}
};
// 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 && m.enabled) =>
{
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, error } => {
// 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,
);
// UX-260618-15: 携带 error 文本,外层 Fatal 分支 emit AiError(终态单气泡)。
return StreamOutcome::Fatal { error };
}
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。
// UX-260618-15: stream_llm 不再 emit AiError(retryable 重试过程 emit 权归此处)。
// 此处不 emit(可能切下一 candidate 成功,emit 会留残留气泡);
// 携带 error 交外层「全 candidate 耗尽」时 emit 最后一条 AiError。
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 { error };
}
// 退避复用 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 {
error: "AI 调用失败:重试预算耗尽(防御兜底)".to_string(),
}
}
/// 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 兜底)
// F-260616-09 B 批2:guard 持 conv_id,复位改 per-conv.generating(设计 §4.3)。
// L2 批2 1b:guard 持 app_handle,ConvState 迁移后 emit AiConvStateChanged 推前端。
let mut guard = GeneratingGuard::new(session_arc.clone(), conv_id.clone(), app_handle.clone());
// F-260616-09 B 批2 入口桥接:loop 启动前确保 per_conv 存在(已存在则保留累积,不存在则建)。
//
// 批2 把所有调用方(IPC commands.rs 写路径 + agentic/mod.rs loop + audit.rs process_tool_calls +
// conversation.rs save + title.rs + knowledge_inject.rs)迁移到 per_conv 真相源。IPC 在 spawn
// loop 前已通过 `session.conv(active_conversation_id).*` 建立 per_conv 并写入初始状态
// (messages push user / generating=true / stop_flag=false / iteration_used=0 等),故 loop
// 入口只需确保 per_conv 存在(防御性 conv() 惰性建,正常路径下已存在)。
//
// 已存在的 per_conv(同 conv 上一轮 loop 留下,如审批暂停后续跑)**保留累积状态**:loop 上一轮
// 在 per_conv.messages 累积的 assistant tool_calls / tool_result 不会因桥接抹掉。审批拒绝/通过
// 时 IPC 直接写 per_conv.messages.replace_tool_result_content,续跑 loop 读 per_conv 拿到正确状态。
//
// guard 语义:loop 启动占用生成态,per_conv.generating=true(IPC spawn 前已置 true,此处幂等确认)。
//
// 单 loop 安全性:批2 阶段是单 active_conversation_id(决策 e 真并发批3+ 落地),无并发 loop 抢
// per_conv 覆盖。批3+ 多 loop 并发时,每 conv 各自 per_conv 条目,互不干扰(本桥接无需改)。
//
// L2 状态机视图层:ConvState 的写收敛已收敛到 `guard`(GeneratingGuard::new 已在上方 L399 创建,
// new 内部 Idle→Generating 迁移 + reset/drop Generating→Idle 迁移,CONV_STATE_ENABLED 门控)。
// 此处不再重复迁移 enum,仅写核心 generating bool(唯一真相源,enum 是其视图)。批2+ 把 enum
// 提升到 PerConvState.conv_state 字段后,此 bool 写入与 enum 迁移由 guard 统一收敛。
{
let mut session = session_arc.lock().await;
let conv = session.conv(&conv_id);
conv.generating = true;
// L2 批2:ConvState 持久化(Idle→Generating 写收敛,CONV_STATE_ENABLED 门控)。
// enum 与 generating bool 双轨(bool 核心复位,enum 供读侧/前端);off 降级跳过。
if conv_state::CONV_STATE_ENABLED {
match conv.conv_state.transition_to(conv_state::ConvState::Generating) {
Ok(ns) => conv.conv_state = ns,
Err(e) => tracing::warn!(
conv_id = %conv_id,
error = %e,
"[ai] 入口 ConvState→Generating 非法(状态机持久化层,不阻断核心生成)"
),
}
}
}
// 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);
// L1 补丁:provider 解析(list_all + select + resolve + ensure + build)整体包 30s timeout。
// 原实现无超时,数据库/keyring 卡死时 run_agentic_loop 入口卡住,generating 永真 + 前端看门狗
// 超时静默吞消息。超时走 AiError 分支(对齐 :412 ensure_resolved_key 失败处理),guard.reset 复位。
// 内部 list_all 失败仍容忍(空池兜底,行为不变);仅整体超时(如 DB 挂死无响应)才走 Err 分支。
let provider_resolve = tokio::time::timeout(
std::time::Duration::from_secs(30),
async {
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())
}
};
// 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(&primary_provider);
let key_len = resolved_key.len();
let provider: Box<dyn LlmProvider> = match super::secret::ensure_resolved_key(
&primary_provider.name, &resolved_key,
) {
Ok(()) => df_ai::build_provider(
&primary_provider.provider_type,
&primary_provider.base_url,
&resolved_key,
&primary_provider.default_model,
),
Err(msg) => return Err(ProviderResolveError::EnsureKeyFailed(msg)),
};
Ok::<_, ProviderResolveError>((primary_provider, candidates, provider, key_len))
},
).await;
let (primary_provider, candidates, provider, key_len) = match provider_resolve {
Ok(Ok((pc, cands, prov, kl))) => (pc, cands, prov, kl),
Ok(Err(ProviderResolveError::EnsureKeyFailed(msg))) => {
guard.reset().await;
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: msg.clone(),
// ensure_resolved_key 失败 = key 缺失/钥匙串损坏,归 Auth
error_type: Some(ErrorType::Auth),
conversation_id: Some(conv_id.clone()),
});
// L3 emit 双写(2026-06-22):关键 AiError publish 到事件总线,供跨模块订阅。
// EVENT_BUS_ENABLED 门控在 publish 内部(false 静默丢弃),无消费者时空转不报错。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Error {
error: msg,
conversation_id: Some(conv_id.clone()),
});
return;
}
Err(_elapsed) => {
// L1 补丁:provider 解析 30s 超时(DB list_all / keyring resolve 卡死)。
// 走 AiError 分支复位 generating,对齐 ensure_resolved_key 失败处理口径。
guard.reset().await;
tracing::error!(conv_id = %conv_id, "[ai] provider 解析超时(30s),可能 DB/keyring 卡死");
let err_msg = "Provider 解析超时(30s),请检查数据库/钥匙串状态后重试".to_string();
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: err_msg.clone(),
error_type: Some(ErrorType::Unknown),
conversation_id: Some(conv_id.clone()),
});
// L3 emit 双写:超时 AiError publish 到事件总线(同上,门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Error {
error: err_msg,
conversation_id: Some(conv_id.clone()),
});
return;
}
};
// 用主候选覆盖入参 provider_config(下游 build_provider / 路由 / 日志均用此)。
// mut:F-04b 切换 candidate 后更新为实际成功所用 provider(供后续 push/save/标题 spawn)。
let mut provider_config = primary_provider;
// 诊断日志: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 && m.enabled) =>
{
id.to_string()
}
_ => resolved_model,
};
// 改进2 B:意图收敛工具(LLM 可见 tool_defs 按 intent 过滤,执行路径仍走完整 registry)
//
// 机制化收敛跑题:取末条 active user 消息 → IntentRecognizer 识别 → 置信 ≥ 阈值
// 则按 intent domain 过滤工具子集(减少 LLM 在无关工具上分心/误用)。
// 三重 fallback(改进2 B §可靠):
// 1. 置信 < INTENT_CONF_THRESHOLD(0.7)→ 全量
// 2. subset 空(Chat/Unknown/Conversation)→ filter_tool_defs 内部回全量
// 3. 过滤后 < 3 条(疑似漂移/误收敛)→ 回全量
// 关键安全:filter 仅改 LLM 可见 tool_defs,**不改执行**——audit 走 tools_arc.get/execute
// 完整 registry,LLM 即使幻觉一个被滤掉的工具名,audit 也能查到/拒绝。
//
// INTENT_CONF_THRESHOLD(常量开关):阈值,低于此值不收敛(回全量)。
// 注意:conf 截断到 1.0(intent.rs),单关键词 weight=1.0 即达 conf=1.0,故阈值=1.0
// 时 conf>=1.0 仍过滤(非关闭)。真关闭收敛:置 >1.0(如 1.1)。调低 = 更激进收敛。
let user_text: String = {
let session = session_arc.lock().await;
session
.conv_read(&conv_id)
.and_then(|c| {
let msgs = c.messages.all_messages_clone();
msgs.iter()
.rev()
.find(|m| matches!(m.role, df_ai::provider::MessageRole::User))
.map(|m| m.content.clone())
})
.unwrap_or_default()
};
let (intent, conf) = IntentRecognizer::recognize(&user_text);
const INTENT_CONF_THRESHOLD: f32 = 0.7;
let all_defs = tools_arc.tool_definitions();
let total = all_defs.len(); // 提前记录全量数(all_defs 将 move 进 tool_defs)
// B 路线 Phase 1 接入:PLANNING_ENABLED(false 默认关)门控 plan_hint 编排。
//
// **零行为变更(关时)**:flag 关时走 filter_tool_defs(intent 收敛扁平子集),
// 与 Phase 1 接入前完全一致——现有 intent/agentic 测试全绿无回归。
//
// **开启时**:调 filter_tool_defs_planned,它内部先 filter_tool_defs 收敛再叠加
// plan_hint 编排(并行组同批聚拢/顺序依赖源在前)。三重 fallback 与 filter_tool_defs
// 同语义(plan_hint 空/非法/registry 漂移均退 filter_tool_defs 扁平结果)。
//
// 关键安全(与 filter_tool_defs 同):本段只改 LLM 可见 tool_defs 的可见性/顺序,
// 不改执行(audit 走 tools_arc 完整 registry)。PLAN_HINT_ENABLED(plan_hint 函数开关,
// Phase0a 就绪 true)与 PLANNING_ENABLED(planner.rs 主 loop 规划开关,本批仍是 false)
// 分离:即使将来 PLANNING_ENABLED 翻 true,plan_hint 内部 PLAN_HINT_ENABLED 关闭时
// filter_tool_defs_planned 仍退扁平(双层开关,任一关闭均退旧行为)。
let tool_defs = if conf >= INTENT_CONF_THRESHOLD {
let filtered = if df_ai::planner::PLANNING_ENABLED {
// Phase 1:plan_hint 编排排序。intent_label 供 plan_hint 备用(当前规则纯关键词驱动)。
filter_tool_defs_planned(&all_defs, &intent, intent.as_str(), &user_text)
} else {
// 旧行为:intent 收敛扁平子集(零行为变更,flag 默认关走此路)。
filter_tool_defs(&all_defs, &intent)
};
if filtered.len() < 3 {
all_defs // 兜底:过滤<3(漂移/误收敛)回全量
} else {
filtered
}
} else {
all_defs // 低置信 fallback 全量
};
tracing::info!(
conv_id = %conv_id,
intent = intent.as_str(),
conf,
filtered = tool_defs.len(),
total,
planning_enabled = df_ai::planner::PLANNING_ENABLED,
"[ai] 意图收敛工具"
);
// 停止信号副本stream_llm 与每轮迭代共享读取,避免重复加锁
// notify 同取一份 Arc 引用B-260615-14stream_llm select! 监听 notified() 即时唤醒
// F-260616-09 B 批2:取 per_conv 的 stop_flag/notify(设计 §4.2 :446)。
// conv_id 来源:run_agentic_loop 入参(loop 启动快照,与 guard 一致)。
// 批2 桥接期 per_conv.stop_flag/notify 与顶层共享同一 Arc(入口桥接时顶层 Arc 直接搬入),
// 故 ai_chat_stop 写顶层 stop_flag 与 loop 读 per_conv stop_flag 互见(同一 AtomicBool 实例)。
let (stop_flag, notify) = {
let session = session_arc.lock().await;
let conv = session.conv_read(&conv_id).expect("[F-09 B 批2] loop 入口桥接后 per_conv 必存在");
(conv.stop_flag.clone(), conv.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;
// L1 断路器:连续同类工具失败计数器(key=失败内容前 40 字符,value=累计次数)。
// loop 生命周期内累加,每轮 process_tool_calls 后检查。达 CIRCUIT_BREAKER_THRESHOLD → 熔断退出。
let mut fail_counts: std::collections::HashMap<String, u32> = std::collections::HashMap::new();
// 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);
// F-09 batch5 修正(用户决策「不设并发会话上限」):删除 loop 入口 acquire_global。
// 原决策 c-1 把 global 改「会话数上限」(loop 入口持整 loop,N=3 排队第 4 个),
// 现用户取消会话数上限 → 多对话 loop 并发不限,第 N+1 个不再 await 阻塞。
// global 字段回归原义「LLM 调用并发限流」:仅由 stream_llm 每轮 + 标题/压缩/提炼/项目分析
// 等单次 LLM 调用点 acquire/drop(防 provider 429),不再绑 loop 生命周期。
// per_conv(每对话 permits=2)保留 — 单对话内主循环 stream_llm + 后台标题/压缩/提炼共享该 conv 的
// permits=2,防单对话内并发 LLM 调用失控(单对话内限流,非会话数限制,符合「不设上限」)。
// permit 绑 guard(函数返回)Drop 自动释放——各 return 点退出即释放槽位。
// retry 同 loop 内,持 per_conv 合理(F-260616-12 核验)。
let _conv_per_conv_permit = llm_concurrency.acquire_per_conv(&conv_id).await;
// 0 = 不限:effective_max=usize::MAX,for 到不了上界,靠 stop_flag/收敛/审批退出(下方达上限暂停分支不触发)
let effective_max = if max_iterations == 0 { usize::MAX } else { max_iterations };
for iteration in start_iteration..effective_max {
// 用户请求停止 → 收尾退出(已生成文本已在上一轮入库)
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, true).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()) });
// L3 emit 双写:入口 stop 的 AiCompleted publish 到事件总线(EVENT_BUS_ENABLED 门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Completed {
total_tokens: usage.total_tokens,
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
conversation_id: Some(conv_id.clone()),
});
return;
}
// 改进5: 主题切换系统标记(保守,双高置信才标,软提示非强制)。
//
// push 时若末两条 user 消息的 topic 都非 None 且不同(Intent 高置信推断的双 topic),
// ContextManager 已置位 pending_topic_marker("old|new")。loop 顶部读取并消费
// (take 一次性清空,防重复 insert),insert 一条 system 软提示告知 LLM 用户已切换话题。
//
// 保守设计:
// - 双高置信:两条 topic 都非 None(都达 0.7 阈值)才标,任一 None(低置信未标)不标。
// - 软提示:仅 insert 一条 system 消息,不强制 LLM 行为(LLM 仍可按自己理解响应)。
// - TOPIC_MARKER_ENABLED(常量开关)false → 跳过(排障/对比用)。
// topic 不参与裁剪/压缩(只供检测),insert_at(0) 同压缩摘要定位(首位 system)。
if TOPIC_MARKER_ENABLED {
let topic_marker_raw: Option<String> = {
let mut session = session_arc.lock().await;
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
"[ai] conv 已删除,loop 退出(主题标记段入口)"
);
return;
}
let conv = session.conv(&conv_id);
conv.messages.take_topic_marker()
};
if let Some((old_topic, new_topic)) = topic_marker_raw.and_then(|s| {
// 解析 "old|new" 格式;splitn 防 topic 名内含 '|' 误切(仅切首 '|' 一次)。
let mut parts = s.splitn(2, '|');
let old = parts.next()?.to_string();
let new = parts.next()?.to_string();
Some((old, new))
}) {
let marker_text = format!(
"── 用户已切换话题(从「{}」到「{}」),请以新话题为准 ──",
old_topic, new_topic
);
let mut session = session_arc.lock().await;
if session.per_conv.contains_key(&conv_id) {
let conv = session.conv(&conv_id);
conv.messages.insert_at(0, ChatMessage::system(&marker_text));
tracing::info!(
conv_id = %conv_id,
iteration,
old_topic = %old_topic,
new_topic = %new_topic,
"[ai] 主题切换标记已 insert(软提示,保守双高置信)"
);
}
}
}
// B-260615-11: 旧 loop 污染防护——每轮开始校验对话一致性。
// 用户新建/切换对话后 active_conversation_id 变更,本 loop(conv_id 快照)成陈旧,
// 继续跑会往新对话 push 消息/pending 造成污染。检测到即退出(guard Drop 复位 generating)。
//
// F-260616-09 B 批2(设计 §4.1):退出判据从 `active_conversation_id != conv_id` 改为
// `!per_conv.contains_key(conv_id)`(conv 存在性)。决策 e 真并发下旧 loop 跑自己的 conv
// 不污染他人,active_conversation_id 切换不应让旧 loop 退出;仅当 conv 被删(clear/delete)
// 才退出。批2 阶段单 loop,conv 不会被删,此校验主要保留语义对齐(批3+ 真并发场景生效)。
// conv_id 来源:run_agentic_loop 入参(与 guard/stop_flag 取用同源)。
{
let mut session = session_arc.lock().await;
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
"[ai] conv 已删除,旧 loop 退出(F-260616-09 B 批2/决策 e)"
);
return;
}
// F-260616-11 决策 a: 累计 iteration 计数(「下一轮起算值」= 当前轮+1)。
// 审批等待/达 max 暂停退出时,本字段即「当前轮+1」;审批续跑 ai_approve 读此值作
// start_iteration 透传,实现 iteration 累计不重置(防多次审批反复跑满 max 致 token 失控)。
// F-09 B 批4:per_conv.iteration_used 唯一真相源,删顶层双写。
let conv = session.conv(&conv_id);
conv.iteration_used = iteration + 1;
}
// 新一轮通知前端(第二轮起),前端需新建 assistant 消息
if iteration > 0 {
let round_n = (iteration + 1) as u32;
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiAgentRound {
round: round_n,
conversation_id: Some(conv_id.clone()),
});
// L3 emit 双写:AiAgentRound publish 到事件总线(EVENT_BUS_ENABLED 门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Round {
round: round_n,
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)不触发,零行为变化。
// F-260616-09 B 批2:messages 操作改 per_conv(设计 §4.2)。
// conv_id 来源:run_agentic_loop 入参。
let prev_compressing = session_arc.lock().await.conv_read(&conv_id).map(|c| c.messages.is_compressing()).unwrap_or(false);
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 = match session.conv_read(&conv_id) {
Some(c) => c,
None => {
tracing::warn!(stale_conv = %conv_id, "[ai] conv 已删除,loop 退出(压缩段入口)");
return;
}
};
let mgr = &mgr.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;
let conv = session.conv(&conv_id);
conv.messages.set_compressing(true);
// 读 active 克隆(不改 status):filter is_active,LLM 失败则消息状态完全不变。
let active_msgs: Vec<ChatMessage> = conv.messages.messages_mut()
[..protect_start]
.iter()
.filter(|t| t.message.is_active())
.map(|t| t.message.clone())
.collect();
let lang = conv.agent_language.clone()
.unwrap_or_else(|| "zh-CN".to_string());
(active_msgs, lang)
};
// 改进3 B:在 active_msgs move 进 compress_via_llm 前,先算关键词摘要兜底文本。
// LLM 压缩失败时仍想保留用户反复提及的主题词(续接锚点),避免裸裁剪丢主题。
// KEYWORD_FALLBACK_ENABLED=false → 跳过(回退原裸裁剪行为,排障/对比用)。
let keyword_fallback: String = if KEYWORD_FALLBACK_ENABLED {
extract_keyword_summary(&active_msgs)
} else {
String::new()
};
// 压缩调用(复用 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,
&conv_id,
&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 conv = session.conv(&conv_id);
let _compressed = conv.messages.compress_old_messages(protect_start);
conv.messages.insert_at(0, ChatMessage::system(&summary));
conv.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.conv(&conv_id).messages.set_compressing(false);
}
Err(e) => {
// LLM 失败 → 改进3 B:仍标 compressed 释放 token + 关键词摘要塞回首条(非裸裁剪)。
//
// 旧行为:消息状态完全不变,降级走 build_for_request 裁剪(丢主题)。
// 新行为(KEYWORD_FALLBACK_ENABLED=true 默认):
// - compress_old_messages 标 [..protect_start] active 为 compressed(释放 token,
// 与成功路径一致,后续 build_for_request 不再把它们进 LLM 上下文);
// - keyword_fallback 非空 → insert_at(0, system 关键词摘要)作续接锚点;
// - keyword_fallback 空(无 user 消息/无可提取词)→ 不插,等价旧裁剪(保底)。
// 持久化语义不变:compressed 仍软删可追溯(DB 全量保留),与成功路径一致。
// KEYWORD_FALLBACK_ENABLED=false → 跳过兜底,等价旧行为(set_compressing(false) +
// 消息状态不变,降级 build_for_request 裁剪)。
session_arc.lock().await.conv(&conv_id).messages.set_compressing(false);
tracing::warn!(
conv_id = %conv_id,
error = %e,
keyword_fallback_len = keyword_fallback.len(),
"[ai] 自动压缩失败,降级走关键词摘要兜底(KEYWORD_FALLBACK_ENABLED={})",
KEYWORD_FALLBACK_ENABLED,
);
if KEYWORD_FALLBACK_ENABLED {
// 标 compressed 释放 token + 关键词摘要塞首位(若非空)。
let inserted = {
let mut session = session_arc.lock().await;
let conv = session.conv(&conv_id);
let _compressed = conv.messages.compress_old_messages(protect_start);
if !keyword_fallback.is_empty() {
conv.messages.insert_at(0, ChatMessage::system(&keyword_fallback));
true
} else {
false
}
};
if inserted {
tracing::info!(
conv_id = %conv_id,
pre_tokens = pre_compress_tokens,
"[ai] 压缩失败兜底:关键词摘要已插首位(compressed 标记已扣 token)"
);
}
}
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 发生,重建等价于复用)。
// F-260616-09 B 批2:build_for_request 改 per_conv.messages(设计 §4.2 :631)。
// conv_id 来源:run_agentic_loop 入参。
let messages = {
let session = session_arc.lock().await;
let conv = match session.conv_read(&conv_id) {
Some(c) => c,
None => {
tracing::warn!(stale_conv = %conv_id, "[ai] conv 已删除,loop 退出(build_for_request 入口)");
return;
}
};
let (history_msgs, _trimmed) = conv.messages.build_for_request(sys_tokens);
let mut msgs = vec![ChatMessage::system(&system_prompt)];
msgs.extend(history_msgs);
msgs
};
// 改进4: tool_result view-only 摘要压缩(build_for_request 后,送 stream 前)。
//
// 遍历 messages(history clone,已含 system prompt + sanitize 后历史)中的 role=Tool 消息,
// 超阈值(content >2KB 或占历史 token >40%)的 content 应用 extract_key_info 替换:
// 保留错误行(error/panic/失败/.rs:N)+ 首/尾各 5 行,中间省略。
//
// **view-only**:messages 是 build_for_request 返回的 clone,改它只影响本轮 LLM 请求视图,
// 不改 ContextManager 持久化(对齐 sanitize_messages:DB 原始 tool_result 完整保留)。
// 故即使摘要有误/过度压缩,下次 build_for_request 仍从 DB 全量重建,可自愈。
//
// TOOL_RESULT_COMPRESS_ENABLED=false(常量开关)→ 跳过(原样送 LLM,排障/对比用)。
let messages: Vec<ChatMessage> = if TOOL_RESULT_COMPRESS_ENABLED {
// 历史总 token 用于占比判定(history_tokens 即 build_for_request 前的快照,
// 此处 messages 已裁剪过,用 estimated_prompt 不准;改用 history_tokens 快照更稳)。
// 注:用 conv.messages.history_tokens() 快照做占比基准(裁剪前的真相),避免循环依赖。
let history_tokens_snapshot: u32 = {
let session = session_arc.lock().await;
session
.conv_read(&conv_id)
.map(|c| c.messages.history_tokens())
.unwrap_or(0)
};
let est = TokenEstimator::default();
let mut compressed_bytes: usize = 0;
let mut original_bytes: usize = 0;
let mut summarized_count: usize = 0;
let msgs: Vec<ChatMessage> = messages
.into_iter()
.map(|mut m| {
if !matches!(m.role, df_ai::provider::MessageRole::Tool) {
return m;
}
let content_len = m.content.len();
let content_tokens = est.estimate_text(&m.content);
if !should_summarize_tool_result(content_len, history_tokens_snapshot, content_tokens) {
return m;
}
original_bytes += content_len;
// tool_name:无 tool_calls 关联(本消息是 tool_result,无 name 字段),用 call_id 或占位。
let tool_name = m.tool_call_id.clone().unwrap_or_else(|| "tool".to_string());
let compressed = extract_key_info(&m.content, &tool_name);
compressed_bytes += compressed.len();
summarized_count += 1;
m.content = compressed;
m
})
.collect();
if summarized_count > 0 {
tracing::info!(
conv_id = %conv_id,
iteration,
summarized_count,
original_bytes,
compressed_bytes,
"[ai] tool_result view-only 摘要压缩(不改持久化)"
);
}
msgs
} else {
messages
};
// 阶段2 占位配对完整性(第二道防线,depth-defense):build_for_request 已在 df-ai 内部
// 自愈一次,此处 tool_result 压缩/系统提示插入后再断言一次,防转换引入新 orphan。
// view-only:messages 是 clone,assert_placeholder_pairing 仅改本 Vec,不改 ContextManager 持久化。
let messages = ContextManager::assert_placeholder_pairing(messages, PLACEHOLDER_INTEGRITY_ENABLED);
// 预估输入 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 并发限流 — F-09 batch5 修正后:
// per_conv 由 loop 入口(L521 _conv_per_conv_permit)整 loop 持有(含工具执行/审批等待/重试),
// 防单对话内并发 LLM 调用失控(单对话内 permits=2,非会话数限制)。
// global 已不再由 loop 入口持有(用户决策「不设并发会话上限」),回归原义「LLM 调用并发限流」:
// 由各单次 LLM 调用点(stream_llm 重试循环内/标题/压缩/提炼/项目分析)各自 acquire/drop 防 429。
//
// 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 总挂钟预算。重试期间持有 per_conv permit 不释放(防新请求挤占)。
//
// F-260614-04 / F-260614-04b: per-provider permit 仍在 candidate 循环内取
// (切换 candidate 时释放旧取新,避免占用未用 provider 的槽);per_conv 由 loop 入口持有。
// 重试总预算(挂钟,含 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;
// UX-260618-15: 追踪最后一个 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(可选):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 { error } => {
// 本 candidate 重试耗尽(retryable)。切下一 candidate 继续尝试。
// candidates 空(单 provider)时此即「耗尽 return」语义——
// 跳出循环后 outcome 仍 None,落入下方「全耗尽 return」兜底。
// UX-260618-15: 不在此 emit AiError(可能切下一 candidate 成功留残留气泡),
// 仅记录 error 供全耗尽兜底 emit 最后一条。
tracing::warn!(
conv_id = %conv_id,
provider = %candidate.name,
"[ai] 候选 {} 流前失败重试耗尽,F-04b 切换下一 provider",
candidate.name,
);
last_exhausted_error = Some(error);
continue 'candidate;
}
StreamOutcome::Fatal { error } => {
// Fatal(4xx 非429/鉴权/参数错):立即放弃整轮 fallback。
// UX-260618-15: stream_llm 不再 emit AiError,此处统一 emit 最终错误气泡(单气泡聚合)。
guard.reset().await;
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: error.clone(),
error_type: Some(ErrorType::Network),
conversation_id: Some(conv_id.clone()),
});
// L3 emit 双写:Fatal AiError publish 到事件总线(门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Error {
error,
conversation_id: Some(conv_id.clone()),
});
return;
}
}
}
match outcome {
Some(r) => r,
None => {
// 全 candidate 耗尽(含单 provider 场景):沿用原单 provider 失败语义。
// UX-260618-15: stream_llm 不再 emit AiError,此处 emit 最后一个 Exhausted candidate 的
// 诊断文本作为最终错误气泡(单气泡聚合,替代原各 candidate emit 多气泡)。
tracing::warn!(
conv_id = %conv_id,
candidates_tried = candidate_chain.len(),
"[ai] 全 provider 候选流前失败重试耗尽,放弃本轮",
);
guard.reset().await;
let err_msg = last_exhausted_error
.unwrap_or_else(|| "AI 调用失败:所有候选 provider 重试耗尽".to_string());
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: err_msg.clone(),
error_type: Some(ErrorType::Network),
conversation_id: Some(conv_id.clone()),
});
// L3 emit 双写:全 candidate 耗尽 AiError publish 到事件总线(门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Error {
error: err_msg,
conversation_id: Some(conv_id.clone()),
});
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 且有文本)
// F-260616-09 B 批2(设计 §4.1 :786):退出校验改 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(F-260616-09 B 批2)"
);
return;
}
if !full_text.is_empty() {
let conv = session.conv(&conv_id);
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();
conv.messages.push(msg);
// 追加系统提示消息:响应因网络中断不完整(对齐决策 a1 系统提示机制)
let mut notice = ChatMessage::system("⚠ 响应因网络中断不完整,以上为已接收的部分内容。可重新发送以获取完整回复。");
notice.model = Some(resolved_model.clone());
conv.messages.push(notice);
}
}
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model), true).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),
});
// L3 emit 双写:MidStream 保文 AiCompleted publish 到事件总线(门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Completed {
total_tokens: usage.total_tokens,
prompt_tokens: tokens.prompt(),
completion_tokens: tokens.completion(),
conversation_id: Some(conv_id.clone()),
});
return;
}
// F-260616-09 B 批5: global/per_conv permit 已上移 loop 入口(L520-521)整 loop 持有,
// stream 后不再立即释放(会话级并发语义:工具执行期间也占槽)。per-provider permit
// (_provider_permit)仍在 candidate 循环内随作用域 Drop 自动释放(F-04b 语义不变)。
// 累加本轮 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();
// F-260616-09 B 批2(设计 §4.1 :837):退出校验改 conv 存在性 + push 改 per_conv.messages。
// conv_id 来源:run_agentic_loop 入参。
{
let mut session = session_arc.lock().await;
// 决策 e 真并发下:push 前再校验 conv 是否仍存在(被删则退出)。
// 读端读到被 clear 的空历史不致命,但 push 写回已删 conv 是污染,必须挡。
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
"[ai] stream 后 conv 已删除,丢弃本轮 push(F-260616-09 B 批2)"
);
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.conv(&conv_id).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.conv(&conv_id).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), true).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()) });
// L3 emit 双写:stream 后 stop 的 AiCompleted publish 到事件总线(门控在 publish 内)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Completed {
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
};
// F-260620 卡死已根治(DIRAUTH 审批链已闭环)。原 eprintln 诊断降级为 tracing::debug,
// 避免污染 stderr(用户可见),保留排障能力(RUST_LOG=debug 可见)。
tracing::debug!(
conv_id = %conv_id,
pending_count,
"[AI-DIRAUTH-DIAG] agentic loop 收到 pending"
);
// L1 断路器:连续同类工具失败熔断(治 agent 无止损死循环,机制非 prompt 说教)。
// CIRCUIT_BREAKER_ENABLED=false → 整段跳过降级 max_iterations 旧行为(开关 + 兜底)。
// 仅检查自动执行(Low)的工具结果——pending_count>0(待审批)交给下方审批分支,
// 此处只看已回填的 Tool 消息。取末尾连续 Tool 消息(倒序 take_while role==Tool),
// 失败内容前 40 字符归一 key 计数,同 key 累计达阈值 → guard.reset + AiError + return。
if CIRCUIT_BREAKER_ENABLED {
let (max_count, sample_key) = {
let session = session_arc.lock().await;
// conv 可能已被删除(stop/新对话),get 不到 → 无消息可判,跳过本轮断路器检查。
let messages = match session.conv_read(&conv_id) {
Some(conv) => conv.messages.all_messages_clone(),
None => Vec::new(),
};
// 倒序取末尾连续 role==Tool 消息(本轮工具回填结果;非 Tool 即停)。
// MessageRole 未派生 PartialEq,用 matches! 宏判变体(不改共享类型 df-ai-core)。
let recent_tool_results: Vec<&ChatMessage> = messages
.iter()
.rev()
.take_while(|m| matches!(m.role, MessageRole::Tool))
.collect();
for m in recent_tool_results {
let content = m.content.as_str();
// 失败判定:禁止/跳过重试/失败/Error 关键词(覆盖权限拒绝/路径错误/异常等)。
let is_failure = content.starts_with("禁止")
|| content.starts_with("已跳过重试")
|| content.contains("失败")
|| content.contains("Error")
|| content.contains("error");
if is_failure {
// 前 40 字符归一 key:同类失败(同前缀)累加,不同错误各自计数互不干扰。
let key: String = content.chars().take(40).collect();
*fail_counts.entry(key).or_insert(0) += 1;
}
}
// 取当前最大计数及其 key(无失败 → max_count=0,不触发)。
fail_counts
.iter()
.max_by_key(|(_, &v)| v)
.map(|(k, &v)| (v, k.clone()))
.unwrap_or((0u32, String::new()))
};
// 锁已随作用域 drop,可安全 await/emit(避免持锁 await 死锁)。
if max_count >= CIRCUIT_BREAKER_THRESHOLD {
tracing::warn!(
conv_id = %conv_id,
max_count,
sample_key = %sample_key,
"[ai] L1 断路器熔断:连续同类失败 {} 次,疑似死循环停止", max_count
);
guard.reset().await;
// L1 求助协议(§2.3,2026-06-21):熔断改发结构化 AiHelpRequired(机制优先 prompt 说教)。
// 开关 CIRCUIT_BREAKER_HELP_EVENT=true(默认)→ AiHelpRequired(求助卡,显 reason/options
// 供用户选);false(兜底回退)→ AiError(旧错误气泡)。两路均 guard.reset + return 强制熔断,
// 仅前端呈现形态不同,语义不变(均终止 loop,逼用户介入)。
if CIRCUIT_BREAKER_HELP_EVENT {
let _ = app_handle.emit(
"ai-chat-event",
AiChatEvent::AiHelpRequired {
reason: format!(
"连续同类失败 {} 次,疑似死循环已停止",
max_count
),
context: format!("最近错误: {}", sample_key),
options: vec![
"换思路".into(),
"授权路径".into(),
"人工接管".into(),
],
conversation_id: Some(conv_id.clone()),
},
);
} else {
// 兜底回退:求助卡未就绪/排障/对比时,沿用旧 AiError 错误气泡呈现。
let _ = app_handle.emit(
"ai-chat-event",
AiChatEvent::AiError {
error: format!(
"连续同类失败 {} 次,疑似死循环已停止。请换思路或人工介入。最近错误: {}",
max_count, sample_key
),
error_type: Some(ErrorType::Unknown),
conversation_id: Some(conv_id.clone()),
},
);
}
return;
}
}
// 有待审批 → 暂停循环,等待用户审批后通过 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), true).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)。
// 注:max_iterations=0(不限)时 effective_max=usize::MAX,for 不会正常结束至此,故不触发暂停(靠 stop/收敛/审批退出)
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), true).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), true).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()) });
// L3 emit 双写:正常完成 AiCompleted publish 到事件总线(EVENT_BUS_ENABLED 门控在 publish 内,
// 无消费者空转留批3 真实消费者接入)。AiTextDelta/AiToolCall* 高频事件不双写(无消费者空转)。
let _ = app_handle.state::<AppState>().ai_event_bus.publish(AiBusEvent::Completed {
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,
conv_id: &str,
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 后中途插入的直接收敛退出。
//
// F-260616-09 B 批2(设计 §4.4):conv_id 改显式入参(替代从 pending/active 推断),
// has_pending 按 conv_id 过滤(决策 e 真并发准备);generating/agent_language/model_override
// 读 per_conv(批2 桥接期与顶层等价,per_conv 未建 fallback 顶层)。
// conv_id 来源:调用方 ai_approve 传 approval.conversation_id;ai_continue_loop/ai_stop_loop
// 传 IPC 参数 conversation_id。
let snap = {
let session = state.ai_session.lock().await;
// path_auth 审批链阶段1:has_pending 改调 session_state(conv_id) 收敛状态机判定,
// 替代手写 path_auth+risk 两表 any 合并(mod.rs 已统一封装)。
// 阶段3a 单真相源合并后:两表合一进 pending_approvals,session_state 单表 any 判定。
// 两表语义不变——任一类挂起都阻塞续跑(agentic loop 在 path_auth 挂起时也已 return 等待,
// 漏任一会致 loop 误续跑空转)。
//
// 兜底/快速回退:若需切回手写 has_pending,原双表组合保留如下(改一行即可):
// let has_pending = session.pending_approvals.values()
// .any(|a| a.conversation_id.as_deref() == Some(conv_id));
let has_pending = session.session_state(conv_id) == SessionState::AwaitingApproval;
// pending_conv_id 保留(should_continue=false 路径的 emit conv_id 回退逻辑)。
// 阶段3a:单表 find_map(原两表合一)。
let pending_conv_id = session.pending_approvals.values()
.find_map(|a| a.conversation_id.clone());
let conv = session.conv_read(conv_id);
// F-09 B 批4:per_conv 唯一真相源,删顶层 fallback(conv 不存在则各字段默认值)。
// L2 读侧迁移(2026-06-22):CONV_STATE_ENABLED on 时读 conv_state.is_active()(Generating/Compressed),
// off 回退 generating bool。双轨过渡:off 等价旧行为。
let is_generating = conv.map(|c| {
if conv_state::CONV_STATE_ENABLED {
c.conv_state.is_active()
} else {
c.generating
}
}).unwrap_or(false);
let agent_language = conv.and_then(|c| c.agent_language.clone());
let model_override = conv.and_then(|c| c.model_override.clone());
ContinueSnapshot {
is_generating,
has_pending,
pending_conv_id,
agent_language,
model_override,
}
};
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!(conv_id = %conv_id, "[ai] try_continue 跳过:仍有待审批,转审批等待态");
} else {
// generating 已复位(用户 stop 或前序循环已 emit Completed):补发 AiCompleted 防前端卡住
tracing::info!(conv_id = %conv_id, "[ai] try_continue 跳过:generating 已复位(被 stop/已结束),补发 AiCompleted 清前端 streaming");
// R-PD-6: 优先用审批所属 conversation_id(审批等待态被 stop 触发,审批仍在 pending_approvals),
// 仅当无任何审批(has_pending=false 且 generating=false)时回退入参 conv_id(批2 显式参数)。
let emit_conv_id = match snap.pending_conv_id.clone() {
Some(cid) => cid,
None => conv_id.to_string(),
};
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompleted {
total_tokens: 0,
prompt_tokens: 0,
completion_tokens: 0,
incomplete: None,
conversation_id: Some(emit_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 复位,语义独立。
// F-09 B 批4:per_conv 唯一真相源,删顶层双写复位。
let mut session = state.ai_session.lock().await;
session.conv(conv_id).generating = false;
drop(session);
tracing::warn!(conv_id = %conv_id, error = %e, "[ai] try_continue 失败:无可用 provider");
let _ = app.emit("ai-chat-event", AiChatEvent::AiError {
error: e.clone(),
// 无可用 provider(配置丢失/全删):用户需在 Settings 设 provider,归 ProviderConfig
error_type: Some(ErrorType::ProviderConfig),
conversation_id: Some(conv_id.to_string()),
});
// L3 emit 双写:try_continue provider-Err AiError publish 到事件总线(门控在 publish 内)。
let _ = app.state::<AppState>().ai_event_bus.publish(AiBusEvent::Error {
error: e,
conversation_id: Some(conv_id.to_string()),
});
return;
}
};
// F-09 B 批2:续生成路径 conv_id 用入参(显式,不再从 active 推断),语言/override 读 per_conv 快照。
let lang = snap.agent_language.clone().unwrap_or_else(|| "zh-CN".to_string());
let conv_id_owned = conv_id.to_string();
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();
// 知识注入:DRY(B):收敛至 inject_knowledge_into_prompt 单一入口。
// P1 修复(审批恢复路径缺知识注入):try_continue 续跑轮此前用裸 build_system_prompt,
// 不调 build_knowledge_context 致续跑轮丢知识库上下文。现与 chat.rs 四处同款走 helper。
// 同消息取 text+id(②口径修复):原 last_user_text 过滤 is_active / user_message_id 走
// last_user_message_id 不过滤 is_active,末条 user 压缩后两值取自不同消息;helper 单次
// 反向扫描同一条消息取两值。复用上方已 clone 的 knowledge_config 快照(避免重复加锁)。
let system_prompt = inject_knowledge_into_prompt(state, conv_id, system_prompt, &knowledge_config).await;
// 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。)
// F-09 B 批4:重检 per_conv.generating(唯一真相源;conv 不存在则 false)。
// L2 读侧迁移(2026-06-22):CONV_STATE_ENABLED on 时读 conv_state.is_active(),off 回退 generating bool。
let still_generating = {
let session = state.ai_session.lock().await;
session.conv_read(conv_id).map(|c| {
if conv_state::CONV_STATE_ENABLED {
c.conv_state.is_active()
} else {
c.generating
}
}).unwrap_or(false)
};
if !still_generating {
tracing::info!(conv_id = %conv_id, "[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_owned.clone()),
});
return;
}
// 恢复循环前通知前端新建 assistant 消息:审批(通过/拒绝)后新一轮文本
// 不应追加到发起工具调用的旧消息,用 AiAgentRound 隔开
let _ = app.emit("ai-chat-event", AiChatEvent::AiAgentRound {
round: 0,
conversation_id: Some(conv_id_owned.clone()),
});
// L3 emit 双写:try_continue 续跑 AiAgentRound publish 到事件总线(门控在 publish 内)。
let _ = app.state::<AppState>().ai_event_bus.publish(AiBusEvent::Round {
round: 0,
conversation_id: Some(conv_id_owned.clone()),
});
tauri::async_runtime::spawn(async move {
run_agentic_loop(session_arc, tools_arc, db, app_handle, provider_config, system_prompt, conv_id_owned, 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 竞态。
///
/// F-260616-09 B 批2:active_conversation_id 字段移除(改用入参 conv_id,不再从 snap 读)。
struct ContinueSnapshot {
is_generating: bool,
has_pending: bool,
pending_conv_id: Option<String>,
agent_language: Option<String>,
model_override: Option<String>,
}