重构: 拆audit.rs第一批reason helper(strategy)
- audit.rs → audit/mod.rs(959→812行) + 新建 audit/reason.rs(167行) - 抽离 reason helper: resolve_project_label/resolve_task_label/build_approval_reason(3 fn自成闭包, pub(super)) - re-export 调用方零变更(commands::ai::audit::* 路径透明, 4处调用方+mod re-export核验) - F-09 batch2 per_conv 保留(本批只搬 reason helper, conv_id路由不动) 主代兜底: cargo check --workspace 0 + test 98 + grep mod reason/use印证 strategy: 单批1-2文件原子; audit/mod.rs余812行(后续批restore/audit_finalize/find_cached)
This commit is contained in:
812
src-tauri/src/commands/ai/audit/mod.rs
Normal file
812
src-tauri/src/commands/ai/audit/mod.rs
Normal file
@@ -0,0 +1,812 @@
|
||||
//! 工具调用审计 + 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;
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 审批历史查询 IPC(AE-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())
|
||||
}
|
||||
|
||||
|
||||
// reason 拼装(resolve_project_label / resolve_task_label / build_approval_reason)
|
||||
// 拆至子模块 audit/reason.rs(本批 helper 抽离,行为零变更)。
|
||||
mod reason;
|
||||
|
||||
// process_tool_calls 调用 build_approval_reason,从子模块引入(super 可见,fn 为 pub(super))。
|
||||
use reason::build_approval_reason;
|
||||
|
||||
/// 启动恢复:从审计表重建 pending_approvals(重启前未审批的工具调用,内存态已丢)
|
||||
///
|
||||
/// pending_approvals 是 AiSession 内存 HashMap,重启必丢。ai_tool_executions 表已存
|
||||
/// status=pending 的行(持久化真相源),此处读回重建内存态,使重启后待审批不丢。
|
||||
/// 前端经 ai_pending_tool_calls 查询 + switchConversation 恢复 toolCard 的 pending_approval 态。
|
||||
///
|
||||
/// **F-260616-09 B 批8 多 conv 适配(设计 §3 batch8 + §5.2)**:DB pending 审批按
|
||||
/// `conversation_id` 分配到各 conv 的 per_conv state——对每个含 pending 审批的 conv
|
||||
/// **惰性建 `PerConvState`**(不依赖 `active_conversation_id` 单值),使后续 ai_approve →
|
||||
/// try_continue_agent_loop 的 `conv_read(conv_id)` 命中各自 per_conv(各归各,不串)。
|
||||
/// `pending_approvals` 本身保持单层 HashMap(设计 §2.1.2,已带 conversation_id 是业务
|
||||
/// 语义而非路由键)。conversation_id=None 的无主审批(R-9)不建 per_conv(无 conv_id 可挂),
|
||||
/// 仍进 pending_approvals 单层表,后续审批按 tool_call_id 路由,不影响正确性。
|
||||
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;
|
||||
// 多 conv 分配:先建 conv → per_conv 索引(惰性建 PerConvState::new),再 insert pending。
|
||||
// conv() 已自处理「已存在则复用」(同 conv_id 多条 pending 复用同一 PerConvState),
|
||||
// 故对每条 pending 都调一次 conv() 是幂等的(后续命中 entry().or_insert_with 跳过新建)。
|
||||
let mut convs_restored: std::collections::HashSet<String> = std::collections::HashSet::new();
|
||||
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;
|
||||
}
|
||||
// 多 conv 分配:有 conversation_id 的审批惰性建/复用对应 conv 的 per_conv(设计 §3 batch8)。
|
||||
// 无 conversation_id 的无主审批(R-9)不建 per_conv——无 conv_id 可挂,仍进 pending_approvals
|
||||
// 单层表按 tool_call_id 路由,后续 ai_approve 路径不依赖其 per_conv。
|
||||
if let Some(cid) = rec.conversation_id.as_deref() {
|
||||
if !cid.is_empty() && convs_restored.insert(cid.to_string()) {
|
||||
// 惰性建:无则建 PerConvState::new,有则 or_insert_with 跳过(幂等,无副作用)。
|
||||
// 不在此恢复 messages——messages 由 switchConversation 从 DB reload 独立路径恢复,
|
||||
// 此处仅保证 conv 在 per_conv 中存在(供 conv_read 命中)。
|
||||
let _ = session.conv(cid);
|
||||
}
|
||||
}
|
||||
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 工具审批重建到内存, 分布 {} 个 conv 的 per_conv",
|
||||
session.pending_approvals.len(),
|
||||
convs_restored.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_id(provider 每轮新 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(非 DoubleEnded),collect 成 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(路径 B):write_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_all),Med/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-05:High 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(路径 B):write_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(定向B):Low 工具失败不 emit AiError。
|
||||
// 原逻辑 emit AiError 会令前端置 streaming=false + 错误气泡,但 process_tool_calls
|
||||
// 仍返回 pending_count=0,agentic loop 续下一轮 → 前端 false 后端跑,状态紊乱。
|
||||
// 现改为正常 emit AiToolCallCompleted(result=错误信息),错误包进 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
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests_f09_batch8_restore {
|
||||
use super::*;
|
||||
use crate::commands::ai::{AiSession, PendingApproval};
|
||||
|
||||
/// F-260616-09 B 批8(设计 §3 batch8 + §5.2):验证 restore_pending_approvals 的多 conv 分配不变量。
|
||||
///
|
||||
/// restore_pending_approvals 受限于 AppState(需 DB),无法直接单测。但其核心分配逻辑
|
||||
/// 「对每个 conversation_id=Some 的恢复审批,惰性建/复用对应 conv 的 PerConvState」依赖
|
||||
/// 的契约是 AiSession::conv 的幂等惰性创建 + conv_read 可达。本测试用同一契约重现分配
|
||||
/// 逻辑,验证各 conv 各归各(不串)、无主审批(None)不建 per_conv。
|
||||
#[test]
|
||||
fn test_restore_multi_conv_distribution_invariant() {
|
||||
let mut session = AiSession::new();
|
||||
|
||||
// 模拟 DB pending 行:3 条分属 2 个 conv(conv-a 两条 + conv-b 一条)+ 1 条 None 无主。
|
||||
let pending_rows: Vec<(String, String, Option<String>)> = vec![
|
||||
("tc-1".to_string(), "write_file".to_string(), Some("conv-a".to_string())),
|
||||
("tc-2".to_string(), "run_command".to_string(), Some("conv-a".to_string())),
|
||||
("tc-3".to_string(), "delete_project".to_string(), Some("conv-b".to_string())),
|
||||
("tc-4".to_string(), "create_task".to_string(), None), // 无主(R-9)
|
||||
];
|
||||
|
||||
// 复现 restore_pending_approvals 的分配逻辑(同名变量,逐字对齐实现)。
|
||||
let mut convs_restored: HashSet<String> = HashSet::new();
|
||||
for (tool_call_id, tool_name, conversation_id) in &pending_rows {
|
||||
if let Some(cid) = conversation_id.as_deref() {
|
||||
if !cid.is_empty() && convs_restored.insert(cid.to_string()) {
|
||||
let _ = session.conv(cid); // 惰性建
|
||||
}
|
||||
}
|
||||
// pending 入单层 HashMap(不进 per_conv,设计 §2.1.2)。
|
||||
session.pending_approvals.insert(
|
||||
tool_call_id.clone(),
|
||||
PendingApproval {
|
||||
tool_call_id: tool_call_id.clone(),
|
||||
tool_name: tool_name.clone(),
|
||||
arguments: serde_json::Value::Null,
|
||||
conversation_id: conversation_id.clone(),
|
||||
recovered: true,
|
||||
diff: None,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
// 不变量 1:per_conv 含 conv-a + conv-b 各一份(惰性建,无主 None 不建)。
|
||||
assert_eq!(session.per_conv.len(), 2, "2 个 conv 各惰性建 1 份 per_conv");
|
||||
assert!(session.conv_read("conv-a").is_some(), "conv-a 恢复后 conv_read 应可达");
|
||||
assert!(session.conv_read("conv-b").is_some(), "conv-b 恢复后 conv_read 应可达");
|
||||
assert!(
|
||||
session.conv_read("never-restored").is_none(),
|
||||
"未恢复的 conv 不应被惰性建"
|
||||
);
|
||||
|
||||
// 不变量 2:pending_approvals 单层 HashMap 4 条(含无主),conversation_id 保留(业务语义)。
|
||||
assert_eq!(session.pending_approvals.len(), 4, "全部 pending 入单层 HashMap");
|
||||
let conv_a_pending: Vec<_> = session
|
||||
.pending_approvals
|
||||
.values()
|
||||
.filter(|a| a.conversation_id.as_deref() == Some("conv-a"))
|
||||
.collect();
|
||||
assert_eq!(conv_a_pending.len(), 2, "conv-a 2 条 pending 各归各");
|
||||
let conv_b_pending: Vec<_> = session
|
||||
.pending_approvals
|
||||
.values()
|
||||
.filter(|a| a.conversation_id.as_deref() == Some("conv-b"))
|
||||
.collect();
|
||||
assert_eq!(conv_b_pending.len(), 1, "conv-b 1 条 pending");
|
||||
let ownerless: Vec<_> = session
|
||||
.pending_approvals
|
||||
.values()
|
||||
.filter(|a| a.conversation_id.is_none())
|
||||
.collect();
|
||||
assert_eq!(ownerless.len(), 1, "无主审批 1 条仍入单层表(R-9)");
|
||||
|
||||
// 不变量 3:同 conv 多条 pending 复用同一 PerConvState(conv() 幂等,不重复建)。
|
||||
let conv_a_ptr = session.conv("conv-a") as *const _;
|
||||
let _ = session.conv("conv-a"); // 再次访问(模拟同 conv 二次 pending)
|
||||
let conv_a_ptr_again = session.conv("conv-a") as *const _;
|
||||
assert_eq!(conv_a_ptr, conv_a_ptr_again, "同 conv 复用同一 PerConvState(幂等)");
|
||||
assert_eq!(session.per_conv.len(), 2, "同 conv 二次访问不新增 per_conv 条目");
|
||||
}
|
||||
|
||||
/// 验证 restore 仅含无主审批(None)时不建任何 per_conv(R-9 边界)。
|
||||
#[test]
|
||||
fn test_restore_only_ownerless_no_per_conv() {
|
||||
let mut session = AiSession::new();
|
||||
let mut convs_restored: HashSet<String> = HashSet::new();
|
||||
// 全部 None 的 pending 行:无主,不应触发任何 conv() 惰性建。
|
||||
let pending_rows: Vec<Option<String>> = vec![None, None, None];
|
||||
for conversation_id in &pending_rows {
|
||||
if let Some(cid) = conversation_id.as_deref() {
|
||||
if !cid.is_empty() && convs_restored.insert(cid.to_string()) {
|
||||
let _ = session.conv(cid);
|
||||
}
|
||||
}
|
||||
}
|
||||
assert!(session.per_conv.is_empty(), "仅无主审批不应建任何 per_conv");
|
||||
assert!(convs_restored.is_empty(), "无 conv 被记录为已恢复");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user