修复: B-260617-01 run_workflow审批执行(抽取run_workflow_inner+ai_approve后端分支)
This commit is contained in:
@@ -302,7 +302,23 @@ pub async fn ai_approve(
|
||||
}
|
||||
drop(session); // 释放锁后再执行
|
||||
|
||||
let exec_result = state.ai_tools.execute(&approval.tool_name, args.clone()).await;
|
||||
// B-260617-01 决策 a(方案 a2 后端分支):run_workflow 特殊处理。
|
||||
// run_workflow handler(tool_registry.rs:544)仅持 db:Arc<Database>,无 AppHandle/State
|
||||
// (registry/event_bus/workflows Repo),经 ai_tools.execute 调用必然 Err。原行为:审批通过 →
|
||||
// handler Err → tool_result=错误提示 → LLM 重试 → 又进审批 → 死循环(工具未标 no_retry)。
|
||||
// 修复:识别 run_workflow → 不走 ai_tools.execute → 直接调工作流执行函数(持有完整 State,
|
||||
// 同 invoke('run_workflow') IPC 入口)→ 结果回填 tool_result(让 LLM 收成功结果停止重试)。
|
||||
// 选 a2 而非 a1(前端拦截):ai_approve 已持 state+app,后端统一更简,无需新增前端 IPC +
|
||||
// 不把工作流知识泄漏到前端。失败(tool 未注册/参数缺失)按原 Err 路径处理(failed tool_result)。
|
||||
//
|
||||
// 仅分支 run_workflow,其余工具审批路径完全不变(含 AE-04 trust_key_for 写 trust / audit 终态
|
||||
// 回填 / emit_data_changed / save / try_continue)。workflow 联动任务推进由 run_workflow 内部
|
||||
// 按 task_id+target_status 自动处理(workflow.rs:270-308),此处不重复推进。
|
||||
let exec_result: anyhow::Result<serde_json::Value> = if approval.tool_name == "run_workflow" {
|
||||
super::execute_run_workflow_for_tool(&app, &state, &args).await
|
||||
} else {
|
||||
state.ai_tools.execute(&approval.tool_name, args.clone()).await
|
||||
};
|
||||
// 工具失败不 return Err:把错误包成 tool_result,落库 + emit completed + 续循环全走通。
|
||||
// 否则前端 approveToolCall 的 catch 会回滚 pending_approval,审批按钮卡死无法消除。
|
||||
let (audit_status, result_val) = match &exec_result {
|
||||
|
||||
@@ -422,3 +422,72 @@ pub(crate) struct ToolCallDraft {
|
||||
pub(crate) name: String,
|
||||
pub(crate) args: String,
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// B-260617-01 决策 a(方案 a2):run_workflow 工具经审批后的真实执行入口
|
||||
//
|
||||
// 背景:run_workflow handler(tool_registry.rs)仅持 db:Arc<Database>,无 AppHandle/State
|
||||
// (registry/event_bus/workflows Repo),经 ai_tools.execute 必然 Err。ai_approve 识别 run_workflow
|
||||
// 后改调本函数(同 invoke('run_workflow') IPC 入口 workflow.rs::run_workflow,持有完整 State),
|
||||
// 真正触发工作流引擎执行 + 联动任务推进,把 execution_id 回填 tool_result 让 LLM 收成功结果停止重试。
|
||||
//
|
||||
// 选后端分支(a2)而非前端拦截(a1):ai_approve 已持 state+app,后端统一更简,无需新增前端 IPC +
|
||||
// 不把工作流知识泄漏到前端。本函数仅在 ai_approve 内被调用(单一调用方)。
|
||||
// ============================================================
|
||||
|
||||
/// 把 run_workflow 工具调用的 AI args 转调工作流执行函数,返回 execution_id(供回填 tool_result)。
|
||||
///
|
||||
/// 从 args 取 task_id + target_status(两者必填),传空 DagDef(由 workflow.rs 按 target_status
|
||||
/// 自动选推进链模板,见 task_workflow_templates::template_for),config 传空 Object。
|
||||
/// 任一参数缺失返 Err(回填 failed tool_result,LLM 据此修参重发)。
|
||||
///
|
||||
/// 不在此推进任务状态 —— workflow.rs::run_workflow 内部已按 task_id+target_status 自动联动
|
||||
/// (成功推进到 target / 失败按 regression_target 退回),此处不重复。
|
||||
pub(crate) async fn execute_run_workflow_for_tool(
|
||||
app: &tauri::AppHandle,
|
||||
state: &crate::state::AppState,
|
||||
args: &serde_json::Value,
|
||||
) -> anyhow::Result<serde_json::Value> {
|
||||
let task_id = args
|
||||
.get("task_id")
|
||||
.and_then(|v| v.as_str())
|
||||
.ok_or_else(|| anyhow::anyhow!("缺少 task_id 参数"))?
|
||||
.to_string();
|
||||
let target_status = args
|
||||
.get("target_status")
|
||||
.and_then(|v| v.as_str())
|
||||
.ok_or_else(|| anyhow::anyhow!("缺少 target_status 参数"))?
|
||||
.to_string();
|
||||
|
||||
// 空 DagDef + target_status:workflow.rs:84-92 自动按 target_status 选推进链模板
|
||||
// (AiNode 自审 / HumanNode 核对闸门),无需本函数感知模板知识。
|
||||
let dag = df_workflow::dag_def::DagDef::new();
|
||||
let config = serde_json::Value::Object(serde_json::Map::new());
|
||||
let name = format!("AI 推进任务 {} → {}", task_id, target_status);
|
||||
|
||||
// 转调工作流执行核心(B-260617-01:run_workflow_inner 是 run_workflow 命令的共用核心,
|
||||
// 同 invoke('run_workflow') IPC 入口的执行逻辑,持 &AppState 非 tauri::State,绕开
|
||||
// tauri::State 在非命令上下文不可构造的限制)。
|
||||
// 返回 execution_id 字符串;失败(模板不存在/DAG 校验失败/DB 写入失败)上抛 Err。
|
||||
let execution_id = crate::commands::workflow::run_workflow_inner(
|
||||
app,
|
||||
state,
|
||||
name,
|
||||
dag,
|
||||
config,
|
||||
Some(task_id.clone()),
|
||||
Some(target_status.clone()),
|
||||
)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("工作流执行失败: {}", e))?;
|
||||
|
||||
// 回填 tool_result:返回 execution_id + 联动语义说明,让 LLM 知工作流已触发(后台异步执行,
|
||||
// 完成后自动推进任务状态),不再重试同 tool_call。
|
||||
Ok(serde_json::json!({
|
||||
"task_id": task_id,
|
||||
"target_status": target_status,
|
||||
"execution_id": execution_id,
|
||||
"status": "running",
|
||||
"note": "工作流已触发(后台异步执行,完成后自动推进任务状态,失败按退回态回滚)。前端 workflow-event 跟踪进度。"
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -529,30 +529,34 @@ pub fn build_ai_tool_registry(db: &Arc<Database>) -> AiToolRegistry {
|
||||
// 实施路径文档 §三 列为阶段3 必做项(tool_registry.rs 此前无此工具连空壳都没有)。
|
||||
// 描述明确按任务 target_status 推进对应工作流(含 AiNode 自审 / HumanNode 核对闸门)。
|
||||
//
|
||||
// handler 约束说明:run_workflow 真正执行需要 AppHandle(转发 workflow-event 到前端) +
|
||||
// AppState(registry 构建 DAG / event_bus 订阅 / workflows Repo 落库 / workflow_state_registry
|
||||
// 注销)等 Tauri 注入态,这些在 tool handler(仅持有 db: Arc<Database>)中无法构造。
|
||||
// 现阶段仅注册 ToolDefinition(schema + risk + 描述)让 LLM 知晓此能力并产出 tool_call;
|
||||
// 真正触发须走 Tauri IPC(经 invoke_handler 注册的 run_workflow 命令,持有完整 State)。
|
||||
// handler 显式报错引导走 IPC,避免在 handler 内重放 DAG 执行引擎(违反单一执行路径原则)。
|
||||
// 后续若需 AI 直驱完整工作流,需扩展 build_ai_tool_registry 注入 AppState 句柄(改 state.rs,
|
||||
// 留待推进链后续批次)。
|
||||
// handler 约束说明(B-260617-01 更新):run_workflow 真正执行需要 AppHandle(转发 workflow-event
|
||||
// 到前端) + AppState(registry 构建 DAG / event_bus 订阅 / workflows Repo 落库 /
|
||||
// workflow_state_registry 注销),这些在 tool handler(仅持 db: Arc<Database>)中无法构造。
|
||||
//
|
||||
// **执行路径(单一)**:run_workflow 是 High risk → 始终经 audit.rs:process_tool_calls 进 pending
|
||||
// → ai_approve 审批 → commands.rs ai_approve 内识别 run_workflow 分支调
|
||||
// execute_run_workflow_for_tool(mod.rs)→ workflow.rs::run_workflow_inner(持完整 State)真正执行。
|
||||
// 故本 handler 经 ai_tools.execute 调用的路径在正常流程下不可达(dead code 防御):
|
||||
// 仅当未来出现"不经 ai_approve 直接 execute run_workflow"的异常调用方时,本 Err 作为防御兜底
|
||||
// 返回明确错误,而非静默 panic。LLM 误判重试同 tool_call 时,High risk 去重缓存
|
||||
// (audit.rs:find_cached_high_risk_result)会复用旧 tool_result 跳过审批,断重试循环。
|
||||
registry.register(
|
||||
"run_workflow", "按任务 target_status 推进对应工作流(含 AiNode 自审 / HumanNode 核对闸门)。参数 task_id + target_status 同时提供才联动任务推进(完成后按 target_status 推进任务,失败按退回态回滚)。属高风险操作(触发工作流引擎执行),须人工批准",
|
||||
"run_workflow", "按任务 target_status 推进对应工作流(含 AiNode 自审 / HumanNode 核对闸门)。参数 task_id + target_status 同时提供才联动任务推进(完成后按 target_status 推进任务,失败按退回态回滚)。属高风险操作(触发工作流引擎执行),须人工批准。审批通过后由后端直接执行工作流引擎并联动推进任务,返回 execution_id",
|
||||
df_ai::ai_tools::object_schema(vec![("task_id", "string", true), ("target_status", "string", true)]),
|
||||
RiskLevel::High,
|
||||
{ let db = db.clone(); Box::new(move |args: serde_json::Value| {
|
||||
let _db = db.clone();
|
||||
{ let _db = db.clone(); Box::new(move |args: serde_json::Value| {
|
||||
Box::pin(async move {
|
||||
let task_id = args["task_id"].as_str().ok_or_else(|| anyhow::anyhow!("缺少 task_id"))?;
|
||||
let target_status = args["target_status"].as_str()
|
||||
.ok_or_else(|| anyhow::anyhow!("缺少 target_status"))?;
|
||||
// handler 无法访问 AppHandle/State(registry/event_bus/workflows Repo),
|
||||
// 真正执行须走 run_workflow Tauri IPC(持有完整 State)。
|
||||
// 这里返回明确错误引导前端走 IPC,而非在 handler 内重放 DAG 引擎。
|
||||
// 防御兜底(正常流程不可达):run_workflow 经 ai_approve 分支执行,不经此 handler。
|
||||
// 误入此路径说明调用方异常(非 ai_approve 直接 execute),返回明确错误勿盲目重试
|
||||
// (对齐 run_command 超时标注模式 L700-712),High risk 去重缓存会断 LLM 重试循环。
|
||||
Err(anyhow::anyhow!(
|
||||
"run_workflow 工具需经 Tauri IPC 执行(持有 AppHandle/State),handler 无 State 句柄。\
|
||||
请前端收到此 tool_call 后转调 invoke('run_workflow', {{ task_id: {}, target_status: {} }})。",
|
||||
"run_workflow 须经人工审批后由后端 ai_approve 分支执行(转调 run_workflow_inner,\
|
||||
持完整 State)。本 handler 经 ai_tools.execute 调用属异常路径(无 AppHandle/State),\
|
||||
勿盲目重试同调用(task_id={}, target_status={});\
|
||||
若需推进任务,重新发起 run_workflow tool_call 走审批流程。",
|
||||
task_id, target_status
|
||||
))
|
||||
})
|
||||
|
||||
@@ -73,6 +73,25 @@ pub async fn run_workflow(
|
||||
// F-260616-06 ②-2: 工作流联动任务的关联 ID 与目标态(同时 Some 才联动)
|
||||
task_id: Option<String>,
|
||||
target_status: Option<String>,
|
||||
) -> Result<String, String> {
|
||||
// 瘦转发到 run_workflow_inner(B-260617-01:抽取共用核心供 ai_approve 经
|
||||
// execute_run_workflow_for_tool 直接调,无需构造 tauri::State)。
|
||||
run_workflow_inner(&app, &state, name, dag, config, task_id, target_status).await
|
||||
}
|
||||
|
||||
/// 工作流执行核心逻辑(B-260617-01 抽取,命令层瘦转发)。
|
||||
///
|
||||
/// 原 `run_workflow` 命令体拆出,使 ai_approve 经 execute_run_workflow_for_tool 直接调用
|
||||
/// (持 &AppState 而非 tauri::State,绕开 tauri::State 在非命令上下文不可构造的限制)。
|
||||
/// 语义与原命令完全一致:模板选 DAG → 落执行记录 → spawn 事件转发 + DAG 执行 + 任务联动推进。
|
||||
pub async fn run_workflow_inner(
|
||||
app: &AppHandle,
|
||||
state: &AppState,
|
||||
name: String,
|
||||
dag: DagDef,
|
||||
config: serde_json::Value,
|
||||
task_id: Option<String>,
|
||||
target_status: Option<String>,
|
||||
) -> Result<String, String> {
|
||||
// 0. F-260616-06 ①-1: DagDef 来源选模板(方案A·最小)。
|
||||
// 规则:
|
||||
@@ -102,6 +121,10 @@ pub async fn run_workflow(
|
||||
name,
|
||||
dag_json,
|
||||
status: "running".to_string(),
|
||||
// B-260617-01: AI 工具经 ai_approve 触发时标 triggered_by=ai(区分人工手动触发)。
|
||||
// 判据:有 task_id 联动(target_status Some)且当前调用来自 execute_run_workflow_for_tool
|
||||
// → 由调用方在 name 前缀已含 "AI 推进任务" 字样,此处统一标 manual(向后兼容,不破坏旧记录语义)。
|
||||
// 精细化 triggered_by=ai 留待后续(需加参数透传,本次最小改动)。
|
||||
triggered_by: Some("manual".to_string()),
|
||||
project_id: None,
|
||||
task_id: task_id.clone(),
|
||||
|
||||
Reference in New Issue
Block a user