Files
DevFlow/src-tauri/src/commands/ai/audit.rs
绝尘 2a4d745c5b 重构: F-09批2 调用方迁移per-conv(双写桥接共存期)
- agentic.rs: GeneratingGuard加conv_id + 三处退出校验改conv存在性(决策e) + stop_flag/notify/messages/iteration per-conv + try_continue加conv_id参数
- commands.rs: IPC写路径双写桥接(per_conv新真相源+顶层双写,批4删) + ai_is_generating读per_conv fallback顶层
- audit/conversation/title/knowledge_inject/lib.rs: per_conv读写迁移
- conv()/conv_read()去allow(批2有调用方)
共存期双写: per_conv主+顶层双写兼容批4前IPC,批4签名改后删顶层
主代兜底: cargo check --workspace 0 + test 96 passed + grep三处退出/GeneratingGuard/双写印证
2026-06-19 02:01:06 +08:00

831 lines
41 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.
//! 工具调用审计 + pending 审批恢复 + 工具调用处理
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use serde::Serialize;
use tauri::{AppHandle, Emitter, State};
use df_ai::ai_tools::{AiToolRegistry, RiskLevel};
use df_ai::provider::ChatMessage;
use df_storage::crud::{AiToolExecutionRepo, ProjectRepo, TaskRepo};
use df_storage::db::Database;
use df_storage::models::AiToolExecutionRecord;
use df_types::types::new_id;
use crate::state::AppState;
use crate::commands::{err_str, now_millis};
use super::{AiChatEvent, AiSession, PendingApproval, ToolCallDraft};
use super::tool_registry::generate_diff;
/// Med/High 工具挂起审批时回填的占位 tool_result 内容。
///
/// 单一常量避免「写入处」(process_tool_calls) 与「去重排查处」
/// (find_cached_high_risk_result) 各持一份字面量致耦合——若两者漂移,
/// 去重会把 pending 占位误判为已落定结果命中缓存,污染 LLM 上下文。
const PENDING_APPROVAL_PLACEHOLDER: &str = "需要用户审批,等待确认";
/// RiskLevel → 审计记录字符串low/medium/high
pub(crate) fn risk_str(r: RiskLevel) -> &'static str {
match r {
RiskLevel::Low => "low",
RiskLevel::Medium => "medium",
RiskLevel::High => "high",
}
}
/// 审计记录字符串 → RiskLevel启动重建 pending_approvals 用,未知串返回 None 跳过)
pub(crate) fn risk_from_str(s: &str) -> Option<RiskLevel> {
match s {
"low" => Some(RiskLevel::Low),
"medium" => Some(RiskLevel::Medium),
"high" => Some(RiskLevel::High),
_ => None,
}
}
// ============================================================
// 审批历史查询 IPCAE-2025-08
// ============================================================
/// 审批历史 DTO(传给前端的精简视图,敏感字段截断防泄露)
///
/// arguments/result 在落库时是完整 JSON(可能含项目名/路径/长结果),审计面板只展示摘要,
/// 故截断到固定长度(参数 120 / 结果 160),既保留可读性又不泄露全量数据到前端 DOM。
#[derive(Debug, Clone, Serialize)]
pub struct ToolExecutionDto {
pub id: String,
pub conversation_id: Option<String>,
pub tool_call_id: String,
pub tool_name: String,
/// 参数摘要(截断 120 字符,完整原值仍留库)
pub arguments_brief: String,
/// 结果摘要(截断 160 字符,None → 空串便于前端展示)
pub result_brief: Option<String>,
/// pending/approved/rejected/executing/completed/failed
pub status: String,
/// low/medium/high
pub risk_level: String,
pub requested_at: String,
pub executed_at: Option<String>,
/// human/auto,None 表示尚未决策
pub decided_by: Option<String>,
}
/// 截断字符串到 max 字符(按 char_indices 边界切,避免切坏中文/emoji)
fn truncate_chars(s: &str, max: usize) -> String {
if s.chars().count() <= max {
return s.to_string();
}
let mut out: String = s.chars().take(max).collect();
out.push('…');
out
}
/// 审批历史面板查询:按 requested_at 倒序(最新在前)分页返回工具调用审计记录。
///
/// 默认 limit=50 / offset=0(第一页)。limit 在 storage 层钳制 ≤200 防滥用。
/// 敏感字段(arguments/result)截断成摘要返回,完整原值仍留库。
#[tauri::command]
pub async fn list_tool_executions(
state: State<'_, AppState>,
limit: Option<u32>,
offset: Option<u32>,
) -> Result<Vec<ToolExecutionDto>, String> {
let limit = limit.unwrap_or(50);
let offset = offset.unwrap_or(0);
let records = state
.ai_tool_executions
.list_recent(limit, offset)
.await
.map_err(err_str)?;
Ok(records
.into_iter()
.map(|r| ToolExecutionDto {
id: r.id,
conversation_id: r.conversation_id,
tool_call_id: r.tool_call_id,
tool_name: r.tool_name,
arguments_brief: truncate_chars(&r.arguments, 120),
result_brief: r.result.map(|s| truncate_chars(&s, 160)),
status: r.status,
risk_level: r.risk_level,
requested_at: r.requested_at,
executed_at: r.executed_at,
decided_by: r.decided_by,
})
.collect())
}
/// 查项目可读标签id → "「项目名」(id=x)",查不到回退友好提示,空 id 返回空串。
///
/// AR-3审批卡片需显示对象名而非裸 id用户反馈"只返回 ID 不知道是什么数据")。
/// 三臂区分:
/// - Ok(Some) → 项目名标签
/// - Ok(None) → "项目已不存在"(真不存在,项目已彻底清除/外部 id
/// - Err → 裸 id 回退 + warn 日志DB 故障/锁/连接断,不误报"已不存在"误导用户)
async fn resolve_project_label(db: &Arc<Database>, id: &str) -> String {
if id.is_empty() {
return String::new();
}
let repo = ProjectRepo::new(db);
match repo.get_by_id(id).await {
Ok(Some(p)) => format!("{}」(id={})", p.name, id),
Ok(None) => format!("(项目已不存在, id={})", id),
Err(e) => {
tracing::warn!("resolve_project_label: 查询项目 id={} 失败,回退裸 id: {}", id, e);
format!("(id={})", id)
}
}
}
/// 查任务可读标签:id → "「任务标题」(id=x)",对齐 resolve_project_label 三臂语义。
///
/// UX-260618-14:advance_task 审批卡的 id 是 task_id,原 build_approval_reason 把 "id"
/// 统一走 resolve_project_label(查 projects 表),误把任务 id 当项目 id 解析,永远落到
/// "项目已不存在"。本方法改查 tasks 表,与 resolve_project_label 同 Ok(None)/Err 分流。
async fn resolve_task_label(db: &Arc<Database>, id: &str) -> String {
if id.is_empty() {
return String::new();
}
let repo = TaskRepo::new(db);
match repo.get_by_id(id).await {
Ok(Some(t)) => format!("{}」(id={})", t.title, id),
Ok(None) => format!("(任务已不存在, id={})", id),
Err(e) => {
tracing::warn!("resolve_task_label: 查询任务 id={} 失败,回退裸 id: {}", id, e);
format!("(id={})", id)
}
}
}
/// 拼审批 reason按工具名 + args 取可读字段,让用户知道"审批要做什么"。
///
/// 主要审批工具特化(带对象名/标识):
/// - delete_project/restore_project/purge_project → 查项目名拼"删除项目「名」(id=x)"
/// - update_project → 查项目名 + field
/// - bind_directory → path + 查项目名
/// - create_task → title + 查 project_id 项目名create_project → namecreate_idea → title
/// - run_workflow → name
/// args 无可读字段或工具未特化时fallback 原 risk 模板(含风险等级提示)。
async fn build_approval_reason(
tool_name: &str,
args: &serde_json::Value,
risk_level: RiskLevel,
db: &Arc<Database>,
) -> String {
let s = |key: &str| args.get(key).and_then(|v| v.as_str()).unwrap_or("");
// 优先级tool_display_hint(轻量标签) → display_hint_for_tool(模板填充) → 硬编码 fallback
let detail = if let Some(hint) = super::tool_registry::tool_display_hint(tool_name) {
// 轻量命中:直接用中文动作前缀 + 风险后缀(不含参数细节)
hint.to_string()
} else if let Some((template, keys)) = super::tool_registry::display_hint_for_tool(tool_name) {
// 按 keys 列表取参数值;含 "id"/"project_id" 的值走 resolve_project_label 解析
let mut values = Vec::with_capacity(keys.len());
for &key in keys {
let val = s(key);
if val.is_empty() { values.clear(); break; }
match key {
// advance_task 的 id 是 task_id(非项目 id),改查 tasks 表,避免误报"项目已不存在"
"id" if tool_name == "advance_task" => values.push(resolve_task_label(db, &val).await),
"id" | "project_id" => values.push(resolve_project_label(db, &val).await),
_ => values.push(val.to_string()),
}
}
// 任一关键参数为空则整条 fallback 为空串(与原行为一致)
if values.is_empty() {
String::new()
} else {
// 模板 {} 占位符按位置填充
let mut result = template.to_string();
for v in &values {
result = result.replacen("{}", v, 1);
}
result
}
} else {
// fallback原有完整硬编码逻辑保持行为不变
match tool_name {
"delete_project" => {
let id = s("id");
if !id.is_empty() { format!("删除项目{}", resolve_project_label(db, id).await) } else { String::new() }
}
"restore_project" => {
let id = s("id");
if !id.is_empty() { format!("从回收站恢复项目{}", resolve_project_label(db, id).await) } else { String::new() }
}
"purge_project" => {
let id = s("id");
if !id.is_empty() { format!("永久删除项目及关联数据,不可恢复{}", resolve_project_label(db, id).await) } else { String::new() }
}
"update_project" => {
let id = s("id");
let field = s("field");
if !id.is_empty() {
format!("修改项目{}字段「{}", resolve_project_label(db, id).await, field)
} else if !field.is_empty() {
format!("修改项目字段「{}", field)
} else { String::new() }
}
"bind_directory" => {
let id = s("id");
let path = s("path");
if !path.is_empty() && !id.is_empty() {
format!("绑定目录:{}(项目{}", path, resolve_project_label(db, id).await)
} else if !path.is_empty() {
format!("绑定目录:{}", path)
} else { String::new() }
}
"create_task" => {
let title = s("title");
let pid = s("project_id");
if !title.is_empty() && !pid.is_empty() {
format!("创建任务:{}(项目{}", title, resolve_project_label(db, pid).await)
} else if !title.is_empty() {
format!("创建任务:{}", title)
} else { String::new() }
}
"create_project" => {
let name = s("name");
if !name.is_empty() { format!("创建项目:「{}", name) } else { String::new() }
}
"create_idea" => {
let title = s("title");
if !title.is_empty() { format!("捕获灵感:{}", title) } else { String::new() }
}
_ => String::new(),
}
};
if detail.is_empty() {
// fallback无可读字段保留风险等级提示模板
match risk_level {
RiskLevel::High => "高风险操作,必须人工批准".to_string(),
_ => "创建操作,请确认是否执行".to_string(),
}
} else {
// 拼上风险等级后缀(高风险标注,便于用户权衡)
match risk_level {
RiskLevel::High => format!("{}(高风险,需人工批准)", detail),
_ => format!("{}(请确认是否执行)", detail),
}
}
}
/// 启动恢复:从审计表重建 pending_approvals(重启前未审批的工具调用,内存态已丢)
///
/// pending_approvals 是 AiSession 内存 HashMap,重启必丢。ai_tool_executions 表已存
/// status=pending 的行(持久化真相源),此处读回重建内存态,使重启后待审批不丢。
/// 前端经 ai_pending_tool_calls 查询 + switchConversation 恢复 toolCard 的 pending_approval 态。
pub async fn restore_pending_approvals(state: &AppState) {
let pending = match state.ai_tool_executions.list_pending().await {
Ok(v) => v,
Err(e) => {
tracing::warn!("启动恢复 pending 审批失败(非阻断): {}", e);
return;
}
};
if pending.is_empty() {
return;
}
let mut session = state.ai_session.lock().await;
for rec in pending {
let args: serde_json::Value = serde_json::from_str(&rec.arguments).unwrap_or_default();
// 过滤 risk_level 解析失败的损坏记录(语义保留:数据异常不恢复)·不再绑定 risk(PendingApproval.risk_level 已删)
if risk_from_str(&rec.risk_level).is_none() {
continue;
}
session.pending_approvals.insert(
rec.tool_call_id.clone(),
PendingApproval {
tool_call_id: rec.tool_call_id,
tool_name: rec.tool_name,
arguments: args,
conversation_id: rec.conversation_id,
recovered: true,
// AE-2025-03: 重启恢复的审批不重读旧文件——审批可能跨重启,
// 期间文件可能已被外部改动,重读生成 diff 反映的不是当初决策时的状态,
// 且恢复路径在 session.lock 内做 async IO 复杂度高,预览价值低。
// 前端见 diff=None 时回退显新 content。
diff: None,
},
);
}
tracing::info!("启动恢复: {} 条 pending 工具审批重建到内存", session.pending_approvals.len());
}
/// 写一条工具执行审计记录insert 失败不阻断主流程,故 `let _ =`
///
/// `decided_by` 有值auto/human= 已决策执行 → 记 executed_at
/// `None`pending 待审批)→ executed_at 留空,待 audit_finalize 回填。
pub(crate) async fn audit_tool_call(
repo: &AiToolExecutionRepo,
conv_id: &str,
tool_call_id: &str,
tool_name: &str,
arguments: &str,
status: &str,
risk_level: RiskLevel,
result: Option<String>,
decided_by: Option<&str>,
) {
let executed_at = if decided_by.is_some() { Some(now_millis()) } else { None };
let _ = repo
.insert(AiToolExecutionRecord {
id: new_id(),
conversation_id: Some(conv_id.to_string()),
tool_call_id: tool_call_id.to_string(),
tool_name: tool_name.to_string(),
arguments: arguments.to_string(),
result,
status: status.to_string(),
risk_level: risk_str(risk_level).to_string(),
requested_at: now_millis(),
executed_at,
decided_by: decided_by.map(|s| s.to_string()),
})
.await;
}
/// 审批后更新审计记录状态(按 tool_call_id 定位 pending 记录,回填 status/decided_by=human/executed_at/result
///
/// 走专用 find_by_tool_call_id —— 通用 query 宏硬编码 ORDER BY created_at DESC
/// 而 ai_tool_executions 无该列,调用会报 "no such column: created_at" 被 unwrap_or_default 吞掉,
/// 导致审批后审计记录永久卡 pending。
pub(crate) async fn audit_finalize(state: &AppState, tool_call_id: &str, status: &str, result: Option<String>) {
// 区分 Err(DB 故障)与 Ok(None)(真无记录):原 unwrap_or_default 把 Err 压成 None,
// DB 故障被「未找到」日志掩盖,审批后审计记录卡 pending 无确诊线索(B-260617-17 同款吞错)。
let mut rec = match state.ai_tool_executions.find_by_tool_call_id(tool_call_id).await {
Ok(Some(rec)) => rec,
Ok(None) => {
tracing::warn!("audit_finalize: 未找到 tool_call_id={} 的审计记录", tool_call_id);
return;
}
Err(e) => {
tracing::error!("audit_finalize: 查询 tool_call_id={} 审计记录失败(DB 故障,状态未回填): {}", tool_call_id, e);
return;
}
};
rec.status = status.to_string();
rec.decided_by = Some("human".to_string());
rec.executed_at = Some(now_millis());
if let Some(r) = result {
rec.result = Some(r);
}
if let Err(e) = state.ai_tool_executions.update_full(&rec).await {
tracing::error!("audit_finalize: 回填 tool_call_id={} 审计状态失败: {}", tool_call_id, e);
}
}
/// AR-11方案A工具名 → (entity, action) 映射,用于数据变更联动刷新。
///
/// 仅数据变更类工具返回 Some其余工具返回 None不 emit df-data-changed
/// 映射规则:
/// - create_project/create_idea/create_task → (project/idea/task, create)
/// - update_project/update_task/update_idea → (..., update)
/// - delete_project/delete_task/delete_idea/purge_project → (..., delete)
/// - restore_project → (project, update)(恢复等同状态变更,前端走 update 刷新)
/// - bind_directory → (project, update)(绑定目录属项目字段更新)
fn data_change_for_tool(name: &str) -> Option<(&'static str, &'static str)> {
let (entity, action) = match name {
"create_project" => ("project", "create"),
"create_task" => ("task", "create"),
"create_idea" => ("idea", "create"),
"update_project" => ("project", "update"),
"update_task" => ("task", "update"),
"update_idea" => ("idea", "update"),
"delete_project" | "purge_project" => ("project", "delete"),
"delete_task" => ("task", "delete"),
"delete_idea" => ("idea", "delete"),
"restore_project" => ("project", "update"),
"bind_directory" => ("project", "update"),
_ => return None,
};
Some((entity, action))
}
/// AR-11方案A工具执行成功后按映射 emit "df-data-changed" { entity, action }。
///
/// 前端 store listen 该事件,按 entity 调 loadProjects/loadTasks/loadIdeas 刷新列表,
/// 取代工具执行后需手动刷新。仅工具名命中 data_change_for_tool 映射才 emit
/// 其余工具run_workflow/list_tasks 等)不 emit。emit 失败不阻断let _ =)。
pub(crate) fn emit_data_changed(app_handle: &AppHandle, tool_name: &str) {
if let Some((entity, action)) = data_change_for_tool(tool_name) {
let _ = app_handle.emit(
"df-data-changed",
serde_json::json!({ "entity": entity, "action": action }),
);
}
}
/// F-260616-05高危工具去重根治 run_command 超时→重试→重新审批循环)。
///
/// ⚠ 性能注记(CR-260618-11#5):本函数在 session lock 持有期间对每个 high risk 工具
/// 串行查 AiToolExecutionRepo::find_by_tool_call_id,工具数多时锁持有线性增长。
/// 批量预取(改签名传 tool_call_ids 批量查)属架构改暂未做,后续 High risk 工具增多时优先处理。
///
/// **根因链**F-04 batch42 已做超时标注LLM 调 run_command 超时 → F-04 把超时标注成
/// tool_result 回传 LLM → LLM不可靠仍重试同命令 → 新 tool_call_idprovider 每轮新 id
/// → `process_tool_calls` 重新 insert pending(audit.rs:418) → 用户被迫重新审批,循环。
///
/// **治本F-05**:即使 LLM 仍重试,同 (tool_name, args) 不重复审批。本函数扫描会话历史,
/// 若发现**已落定**executed/failed/rejected非 pending的同类同参高危调用返回其
/// tool_result 内容,调用方把缓存结果作为新 tool_call_id 的 tool_result 回传 LLM跳过审批。
///
/// **安全边界**
/// - 仅 High risk循环源头Low/Med 不去重Low 直行无审批Med 重试场景少且去重易误伤)。
/// - 仅匹配**已落定**结果pending 的不命中——pending 已有独立流程,不会触发循环;
/// 且避免两个 pending 互相吞掉审批。LLM 只有在收到 tool_result 后才会重试,
/// 故循环必然是「上一条已落定 → 重试」形态pending 不命中不影响治本。
/// - args 走 JSON 规范化比较(键序无关),避免 LLM 两次生成键序不同误判为不同命令。
/// - 返回 None 表示无缓存命中(走原审批流程)。
///
/// `session` 只读扫描 messages不写调用方据返回值决定是否跳过 insert pending。
///
/// F-260616-09 B 批2:经 `conv_id` 索引 per_conv.messages(顶层 messages 批2 后是死字段)。
/// conv_id 来源:process_tool_calls 入参 → 由 agentic.rs run_agentic_loop 入参透传。
async fn find_cached_high_risk_result(
session: &AiSession,
conv_id: &str,
audit_repo: &AiToolExecutionRepo,
tool_name: &str,
args: &serde_json::Value,
) -> Option<(String, String)> {
use df_ai::provider::MessageRole;
// 规范化新调用的 args 为可比字符串(排序键,键序无关)
let new_args_key = canonical_args_key(args);
// F-260616-09 B 批2:读 per_conv.messages。process_tool_calls 调用前 loop 入口已桥接建立 per_conv,
// 故 conv_read 必命中;防御性 None 时返 None(无缓存命中,走原审批流程)。
let conv = session.conv_read(conv_id)?;
// ContextManager::iter 返回 impl Iterator非 DoubleEndedcollect 成 Vec 再反向遍历。
// 单对话消息量小百级collect 开销可忽略。
let msgs: Vec<&ChatMessage> = conv.messages.iter().collect();
// 1) 反向扫描 assistant tool_calls找最近一条同名同参的 High 工具调用 → 拿到旧 tool_call_id
// 反向:循环是「最近一次超时→重试」,命中通常是末尾附近,反向先停省全扫。
let mut prev_tool_call_id: Option<String> = None;
for msg in msgs.iter().rev() {
if !matches!(msg.role, MessageRole::Assistant) {
continue;
}
let Some(tcs) = msg.tool_calls.as_ref() else { continue };
for tc in tcs {
if tc.function.name != tool_name {
continue;
}
// 旧调用的 args 是流式拼接的 JSON 字符串,解析失败跳过(不误判为命中)
let Ok(old_args) = serde_json::from_str::<serde_json::Value>(&tc.function.arguments) else {
continue;
};
if canonical_args_key(&old_args) == new_args_key {
prev_tool_call_id = Some(tc.id.clone());
break;
}
}
if prev_tool_call_id.is_some() {
break;
}
}
// 2) 用旧 tool_call_id 找对应 tool_result。注意审批拒绝/超时失败也属「已落定」,
// 其 tool_result 内容同样回传LLM 看到原反馈自行决定,不再逼用户二次审批)。
// pending 占位(「需要用户审批,等待确认」)不命中——仍在审批中,走原流程。
let old_id = prev_tool_call_id?;
for msg in msgs.iter().rev() {
if !matches!(msg.role, MessageRole::Tool) {
continue;
}
if msg.tool_call_id.as_deref() != Some(old_id.as_str()) {
continue;
}
// 命中旧 tool_result排除 pending 占位(内容固定为 PENDING_APPROVAL_PLACEHOLDER
if msg.content == PENDING_APPROVAL_PLACEHOLDER {
return None;
}
// SW-260618-16: 查审计表拿缓存来源真实 status(completed/rejected/failed),透传给
// audit_tool_call 而非固定 completed(审计语义与结果内容一致,防"rejected/failed 结果
// 记 completed"误导安全追溯)。审计记录缺失/查询失败 fallback completed(不阻塞去重,降级原行为)。
let status = audit_repo
.find_by_tool_call_id(old_id.as_str())
.await
.ok()
.and_then(|opt| opt.map(|rec| rec.status))
.unwrap_or_else(|| "completed".to_string());
return Some((msg.content.clone(), status));
}
None
}
/// 把 JSON args 规范化为可比字符串:对象键按字典序排序后序列化,
/// 键序不同的等价参数生成同一 key防 LLM 两次生成键序不同误判为不同命令)。
fn canonical_args_key(args: &serde_json::Value) -> String {
let mut v = args.clone();
sort_object_keys(&mut v);
// 紧凑序列化(无空白),保证稳定可比
serde_json::to_string(&v).unwrap_or_default()
}
/// 递归对 JSON 对象的键做字典序排序(就地),数组成员也递归排序。
fn sort_object_keys(v: &mut serde_json::Value) {
match v {
serde_json::Value::Object(map) => {
// BTreeMap 按键排序,重建 Object
let mut entries: Vec<(String, serde_json::Value)> = map
.iter()
.map(|(k, val)| (k.clone(), val.clone()))
.collect();
entries.sort_by(|a, b| a.0.cmp(&b.0));
map.clear();
for (k, mut val) in entries {
sort_object_keys(&mut val);
map.insert(k, val);
}
}
serde_json::Value::Array(arr) => {
for item in arr {
sort_object_keys(item);
}
}
_ => {}
}
}
/// AE-2025-03路径 Bwrite_file 审批预览 diff 生成。
///
/// 从 write_file args 取 path旧文件路径+ content新内容
/// 预读旧文件 → 复用 `generate_diff`F-260615-10 LCS 行级 diff注入审批事件。
/// 旧文件不存在(新建)/ 读失败 / args 缺字段 → None前端回退显新 content
///
/// **仅读不改**审批未通过前不动文件读路径不校验write_file handler 自身会校验
/// validate_path审批拒绝则 handler 不执行,恶意路径读失败仅得到 None 不产生危害)。
async fn build_write_file_diff(args: &serde_json::Value) -> Option<String> {
let path = args.get("path")?.as_str()?;
let new_content = args.get("content")?.as_str()?;
match tokio::fs::read_to_string(path).await {
Ok(old_content) => {
if old_content == new_content {
return None; // 无变化不生成 diff理论上 write_file 不会如此,容错)
}
Some(generate_diff(&old_content, new_content))
}
Err(_) => None, // 文件不存在(新建)/ 无读权限 → 无 diff前端显新 content
}
}
/// 处理流式接收的工具调用Low 风险并行执行join_allMed/High 进审批门控
/// 返回待审批的工具数量0 = 全部自动执行完成)
pub(crate) async fn process_tool_calls(
session: &mut AiSession,
tool_calls_acc: HashMap<u32, ToolCallDraft>,
tools_arc: &Arc<AiToolRegistry>,
db: &Arc<Database>,
app_handle: &AppHandle,
conv_id: &str,
) -> usize {
let mut tc_list: Vec<_> = tool_calls_acc.into_iter().collect();
tc_list.sort_unstable_by_key(|(i, _)| *i);
// B-260616-21 治本兜底:LLM 异常复用同 tool_use.id(stream_recv 按 content_block index 分桶,
// 同 id 不同 index draft 可并存 → 每 draft emit AiToolCallStarted 致同 id emit 两次 → 前端 push 两卡,
// Completed 按 id 只 update 首张 → 次张残留 running 0行)。process 层按 id 去重——同 id 保留
// 最小 index 的首个,丢弃后续,保证 emit Started 的 id 唯一。前端 useAiEvents.ts:205 findToolCall
// 守卫双保险。详 docs/02-架构设计/B-260616-21排查方案-2026-06-16.md。
let mut seen_ids: HashSet<String> = HashSet::new();
tc_list.retain(|(_, draft)| seen_ids.insert(draft.id.clone()));
let mut pending_count = 0usize;
let audit_repo = AiToolExecutionRepo::new(db);
// 解析 args + 批量发 Started前端骨架按原始 index 顺序展示)
let drafts: Vec<(u32, ToolCallDraft, serde_json::Value)> = tc_list.into_iter()
.map(|(idx, draft)| {
let args = serde_json::from_str(&draft.args).unwrap_or(serde_json::Value::Object(Default::default()));
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiToolCallStarted {
id: draft.id.clone(),
name: draft.name.clone(),
args: args.clone(),
conversation_id: Some(conv_id.to_string()),
});
(idx, draft, args)
})
.collect();
// 分类Low 收集并行执行Med/High 立即进审批门控push 占位 tool_result
//
// F-260616-05High risk 在进审批门前先查去重缓存find_cached_high_risk_result
// 若 LLM 重试同命令(同 tool_name + 同 args键序无关命中已落定的旧 tool_result
// 把缓存结果作为新 tool_call_id 的 tool_result 回传 LLM跳过 insert pending + 跳过审批,
// 断「超时→重试→重新审批」循环。Med 不去重去重易误伤Low 无审批本就不进此分支。
let mut low_risk: Vec<(ToolCallDraft, serde_json::Value)> = Vec::new();
// trust_hits 收集AE-04 会话信任命中(同会话已批准同类操作),真实执行但移到锁外 spawn
// (对齐 Low risk L673-720 不持锁模式)。命中时此处只 emit toast + 收集,不 .await execute。
// 元组携带 risk_level: 回填审计需原等级(Med/High 才进此分支),避免闭包内再查 registry
let mut trust_hits: Vec<(ToolCallDraft, serde_json::Value, String, RiskLevel)> = Vec::new();
for (_, draft, args) in drafts {
let risk_level = tools_arc.get(&draft.name).map(|t| t.risk_level).unwrap_or(RiskLevel::High);
match risk_level {
RiskLevel::Low => low_risk.push((draft, args)),
RiskLevel::Medium | RiskLevel::High => {
// AE-2025-04 会话级信任Session Trust首批 write_file / run_command
// 同会话已批准过同工具+同目录 → TrustKey 命中 → 自动放行(跳过 pending + 二次确认)。
// 命中后走与 F-05 去重命中相似的「直接执行 + Completed + 审计 decided_by=auto_trust」路径
// 但与 F-05 不同F-05 复用缓存 tool_result 跳过执行trust 放行**真实执行工具**
// (用户信任同目录同类操作,但仍要看每次的真实结果)。
let trust_hit = super::trust_key_for(&draft.name, &args)
.and_then(|key| {
// F-260616-09 B 批2:读 per_conv.session_trust(conv_id 来源:本函数入参)。
if session.conv_read(conv_id).map(|c| c.session_trust.contains(&key)).unwrap_or(false) {
Some(key)
} else {
None
}
});
if let Some(key) = trust_hit {
let dir_label = match &key {
super::TrustKey::Write { dir } | super::TrustKey::Execute { dir } => dir.clone(),
};
tracing::info!(
tool = %draft.name,
dir = %dir_label,
new_tool_call_id = %draft.id,
"[AE-2025-04] 会话信任命中: 同会话已批准同类操作,自动放行(跳过审批+二次确认)"
);
// emit 轻量 toast 事件(前端 AiChat.vue 显示"🔓 自动放行: tool(dir)")
// toast 即时反馈,锁内 emit 不阻塞
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiToolAutoApproved {
id: draft.id.clone(),
tool: draft.name.clone(),
dir: dir_label.clone(),
conversation_id: Some(conv_id.to_string()),
});
// 收集后循环外 spawn 执行,对齐 Low risk 不持锁(原注释 L593 声称"锁外"但代码持锁,
// run_command 慢命令会阻塞同会话所有触 state.ai_session 的 IPC,CR-51 修此)
trust_hits.push((draft, args, dir_label, risk_level));
continue;
}
// F-05仅 High 查去重缓存Med 保持原审批流程
if matches!(risk_level, RiskLevel::High) {
if let Some((cached, status)) = find_cached_high_risk_result(session, conv_id, &audit_repo, &draft.name, &args).await {
// 命中:把缓存结果作为新 tool_call_id 的 tool_result 回传,跳过审批
tracing::info!(
tool = %draft.name,
new_tool_call_id = %draft.id,
"[F-05] 高危工具去重命中:LLM 重试同命令,复用缓存结果跳过审批(断循环)"
);
// F-260616-09 B 批2:写 per_conv.messages(conv_id 来源:本函数入参)。
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, &cached));
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
id: draft.id.clone(),
result: serde_json::Value::String(cached.clone()),
conversation_id: Some(conv_id.to_string()),
});
// 审计:去重命中记一条(status 透传缓存来源 completed/rejected/failed,SW-260618-16;decided_by=auto_dedup),不进 pending
audit_tool_call(&audit_repo, conv_id, &draft.id, &draft.name, &draft.args, &status, risk_level, Some(cached), Some("auto_dedup")).await;
continue;
}
}
pending_count += 1;
// AE-2025-03路径 Bwrite_file 挂起审批前预读旧文件生成 diff。
// 仅 write_file覆盖整文件有完整新旧内容可对比其他工具 diff=None。
// 旧文件不存在(新建)→ diff=None前端回退显新 content。
// 读失败不阻断审批(容错:文件无读权限等极端情况降级为无 diff 预览)。
let approval_diff: Option<String> = if draft.name == "write_file" {
build_write_file_diff(&args).await
} else {
None
};
session.pending_approvals.insert(draft.id.clone(), PendingApproval {
tool_call_id: draft.id.clone(),
tool_name: draft.name.clone(),
arguments: args.clone(),
conversation_id: Some(conv_id.to_string()),
recovered: false,
diff: approval_diff.clone(),
});
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, PENDING_APPROVAL_PLACEHOLDER));
let reason = build_approval_reason(&draft.name, &args, risk_level, db).await;
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiApprovalRequired {
id: draft.id.clone(),
name: draft.name.clone(),
args: args.clone(),
reason,
diff: approval_diff,
conversation_id: Some(conv_id.to_string()),
});
audit_tool_call(&audit_repo, conv_id, &draft.id, &draft.name, &draft.args, "pending", risk_level, None, None).await;
}
}
}
// AE-04 trust-hit 并行执行execute + 即时 emit 在闭包内(闭包不访问 session非"锁已释放"——
// session 锁仍由调用方 agentic.rs 持有至 process_tool_calls 返回CR-53 审查纠正原"移锁外"误述),
// push tool_result / audit 在 join_all 后串行回填(持锁)。对齐 Low risk 并行模式。
// CR-51 修:原 inline .await execute 串行执行每个工具阻塞期间锁被持有run_command 慢命令
// 阻塞同会话 IPC改 join_all 并行多工具减少总阻塞时间(锁持有时长不变,并行化降阻塞)。
// join_all 保序——结果顺序 = trust_hits 输入顺序 = tc_list 原始 index 顺序,不额外 sort
if !trust_hits.is_empty() {
let results: Vec<(ToolCallDraft, RiskLevel, Result<String, String>)> =
futures::future::join_all(trust_hits.into_iter().map(|(draft, args, _dir_label, risk_level)| {
let tools = tools_arc.clone();
let app_clone = app_handle.clone();
let conv_clone = conv_id.to_string();
async move {
let exec_result = tools.execute(&draft.name, args).await;
match exec_result {
Ok(val) => {
let content = val.to_string();
let _ = app_clone.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
id: draft.id.clone(),
result: serde_json::Value::String(content.clone()),
conversation_id: Some(conv_clone),
});
// AR-11数据变更类工具执行成功后 emit df-data-changed 联动刷新
// (write_file/run_command 不在 data_change_for_tool 映射内,emit_data_changed 内部 None 即 noop)
emit_data_changed(&app_clone, &draft.name);
(draft, risk_level, Ok(content))
}
Err(e) => {
let content = e.to_string();
let _ = app_clone.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
id: draft.id.clone(),
result: serde_json::Value::String(content.clone()),
conversation_id: Some(conv_clone),
});
(draft, risk_level, Err(content))
}
}
}
})).await;
// 串行回填 tool_result + 审计(持 session 锁)
for (draft, risk_level, outcome) in results {
let (status, content) = match outcome {
Ok(c) => ("completed", c),
Err(c) => ("failed", c),
};
// F-260616-09 B 批2:写 per_conv.messages。
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, content.clone()));
// 审计:trust 放行仍记一条(decided_by=auto_trust),留痕可追溯
audit_tool_call(&audit_repo, conv_id, &draft.id, &draft.name, &draft.args, status, risk_level, Some(content), Some("auto_trust")).await;
}
}
// Low 风险并行执行execute + 即时 emit 在闭包内(不持 session 锁),
// push tool_result / audit 在 join_all 后串行回填(持锁,与 Med/High 占位拼接)。
// join_all 保序——结果顺序 = low_risk 输入顺序 = tc_list 原始 index 顺序,不额外 sort
if !low_risk.is_empty() {
let results: Vec<(ToolCallDraft, Result<String, String>)> =
futures::future::join_all(low_risk.into_iter().map(|(draft, args)| {
let tools = tools_arc.clone();
let app_clone = app_handle.clone();
let conv_clone = conv_id.to_string();
async move {
let result = tools.execute(&draft.name, args).await;
match result {
Ok(val) => {
let _ = app_clone.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
id: draft.id.clone(),
result: val.clone(),
conversation_id: Some(conv_clone),
});
// AR-11方案A数据变更工具执行成功后 emit df-data-changed,
// 前端 store listen 刷新列表(仅命中映射的工具 emit见 emit_data_changed
emit_data_changed(&app_clone, &draft.name);
(draft, Ok(val.to_string()))
}
Err(e) => {
// AR-6定向BLow 工具失败不 emit AiError。
// 原逻辑 emit AiError 会令前端置 streaming=false + 错误气泡,但 process_tool_calls
// 仍返回 pending_count=0agentic loop 续下一轮 → 前端 false 后端跑,状态紊乱。
// 现改为正常 emit AiToolCallCompletedresult=错误信息),错误包进 tool_result
// 让 LLM 看到工具失败自行决定下一步loop 正常续,前端状态一致。
let err_msg = format!("工具 {} 执行失败: {}", draft.name, e);
let _ = app_clone.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
id: draft.id.clone(),
result: serde_json::Value::String(err_msg.clone()),
conversation_id: Some(conv_clone),
});
(draft, Err(err_msg))
}
}
}
})).await;
// 串行回填 tool_result + 审计(持 session 锁)
for (draft, outcome) in results {
let (status, content) = match outcome {
Ok(c) => ("completed", c),
Err(c) => ("failed", c),
};
// F-260616-09 B 批2:写 per_conv.messages。
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, content.clone()));
audit_tool_call(&audit_repo, conv_id, &draft.id, &draft.name, &draft.args, status, RiskLevel::Low, Some(content), Some("auto")).await;
}
}
pending_count
}