Files
DevFlow/docs/02-架构设计/已编号方案/B-03-人工审批响应机制-2026-06-14.md
绝尘 998a2f243d 文档: 架构方案文档(意图识别论证+多主题愿景/论证+文档物理分类+边界清晰化)
squash合并:
- 意图识别层论证(8维度+10业界佐证)
- 多主题上下文管理愿景+并存论证+补充论证(多轮agentic)
- 架构设计文档物理分类(四子目录+INDEX+命名规范+引用同步+边界清晰化)
- 前端架构技术债清单归档
2026-06-19 15:04:04 +08:00

18 KiB
Raw Blame History

B-03 人工审批响应机制设计

真相源(本文档唯一展开完整设计)。功能决策记录仅放摘要 + 指针。

背景B-260614-03 — df-workflow HumanNode 假实现(human_node.rs:55 注释"等待审批"但首次迭代直接 return "同意")。 状态:📐 设计完成 | 创建2026-06-14 | 来源:多代理探索 依赖B-260614-06(execution_id 硬编码)、B-260614-07(每节点全新空 StateMachine)


实施状态(2026-06-18 核对)

B-03a(响应等待 + 超时)— 已落地HumanNode.execute 完整实现 subscribe → send(Request)(.await 修复 send 缺 poll 死 bug) → select! 循环(响应/超时/取消)

  • crates/df-nodes/src/human_node.rs:70 先 subscribe:74-83 .send(HumanApprovalRequest).await(原审查报告头号 bug 已修):90-175 select! 循环(rx.recv() / sleep_until(deadline) / cancel_tick)。
  • 单测覆盖:human_node.rs:286 normal_approval_returns_decision / :339 mismatched_execution_id_filtered_then_timeout / :378 timeout_when_no_response / :388 invalid_decision_ignored_then_timeout 等。

B-03b(取消机制)— 已落地(原设计标"待做",实际已实施)

  • StateMachine::set_cancelled 已加:crates/df-workflow/src/state.rs:95(注释为"唯一受控旁路")。
  • cancel_workflow_node IPC 已加:src-tauri/src/commands/workflow.rs:486,并在 src-tauri/src/lib.rs:107 注册;含终态前置守卫(Pending/Running/Waiting 才允许 set_cancelled)。
  • 前端取消按钮已接:src/views/ProjectDetail.vue:425 + src/stores/project/workflow.ts:131(经 src/api/workflow.ts:66 invoke)。
  • B-07(共享 StateMachine)已解:NodeContext.node_statusStateMachine clone内部 Arc<Mutex<HashMap>> 共享(crates/df-workflow/src/executor.rs:39 注释、state.rs:88-95)run_workflow 把执行器状态机注册到 AppState 全局表IPC 经 execution_id 取引用直达运行中节点。
  • 端到端测试:human_node.rs:488 end_to_end_human_approval_completes_workflow / :551 end_to_end_human_approval_cancelled

超出原设计、后追加的能力

  • F-260615-01 多选审批(select_type=single|multiple + decisions 数组)human_node.rs:63-66 解析、:102-111 数量/合法性校验、src-tauri/src/commands/workflow.rs:404 approve_human_approval 签名含 decisions/select_type
  • F-260616-06 阶段2 审批拒绝语义化decision 命中拒绝关键字(human_node.rs:19-22 REJECT_KEYWORDS)→ 返 Err 触发工作流 failed(原设计拒绝与同意一样 Ok 的行为已反转)。

B-06(execution_id 下沉)— 未单独核验状态,本设计标注当时为"dummy-execution-id";现 approve_human_approval IPC 签名已显式收 execution_id: String(workflow.rs:407),由调用方传入。是否已从 run_workflow 真 ID 下沉到 executor 再到 NodeContext本次仅标注未深核。

原文以下设计正文保持不变,作为历史设计记录;落地形态以上方"实施状态"为准。


一、背景与问题

HumanNode 是工作流中唯一的阻塞节点,用于在 DAG 执行链路上插入人工确认门控(如"发布前确认""删除前确认")。当前实现 crates/df-nodes/src/human_node.rs 已正确发送 WorkflowEvent::HumanApprovalRequest 到事件总线,但紧接着直接 return NodeOutput { decision: "同意" },从不等待前端审批响应。这导致:

  1. 人工审批门控形同虚设——工作流永远按"同意"放行,无人工拦截能力。
  2. ai.rsai_approve 严谨审批链路(Low 自动 / Medium+High 暂停等审批)矛盾——同一项目两套审批机制,一严谨一形同虚设。
  3. 前端 approve_human_approval IPC、HumanApprovalResponse 事件、stores/project.ts 监听链路均已接通,却被 HumanNode 的假返回架空。

二、现状勘察:链路 90% 已通,缺口仅 1 处

经代码勘察,端到端审批响应链路的基础设施已全部就位,唯一缺口在 HumanNode 本身。

组件 位置 状态
WorkflowEvent::HumanApprovalRequest / HumanApprovalResponse 事件 df-core/src/events.rs:60-73 已定义
EventBus(tokio broadcastcapacity 256subscribe()) df-workflow/src/eventbus.rs 可用
approve_human_approval IPC(前端响应回总线) src-tauri/src/commands/workflow.rs:161 已实现
AppState.event_bus 单一全局总线 src-tauri/src/state.rs:149 run_workflow 与 NodeContext 共享同一 sender
前端监听 workflow-event + 捕获 Request + 调 IPC src/stores/project.ts:214, 228 已接通
run_workflow 转发所有事件(含 Request/Response)到前端 src-tauri/src/commands/workflow.rs:83
HumanNode.execute 订阅 Response 等待审批 crates/df-nodes/src/human_node.rs:55 缺口:发完直接 return

端到端路径验证

HumanNode.execute
  → ctx.event_bus.send(HumanApprovalRequest)          # 同一 broadcast bus
  → run_workflow 转发器(独立 receiver) emit "workflow-event" 到前端
  → 前端 stores/project.ts 捕获 Request存 pendingApproval渲染审批 UI
  → 用户点"同意/拒绝"
  → invoke('approve_human_approval', { execution_id, node_id, decision, comment })
  → workflow.rs:161 构造 HumanApprovalResponsestate.event_bus.send(Response)
  → 同一 broadcast bus
  → HumanNode 的 receiver 收到 Response ✓
  → 过滤 execution_id + node_id 命中 → 返回 NodeOutput

AppState.event_busrun_workflow(workflow.rs:100state.event_bus.clone() 传入 DagExecutor::new)与 NodeContext.event_bus(executor.rs:77 传入 self.event_bus.clone())之间共享同一 broadcast::Sender(Clone 仅复制 sender 句柄,底层通道同一)。approve_human_approval 发往 state.event_bus,即发往 HumanNode 订阅的同一通道。路径闭环成立

三、B-06 / B-07 前置依赖的真实影响

todo.md 标 B-03 依赖 B-06/B-07。核对后分级澄清,避免误解为硬阻塞:

B-06(execution_id 硬编码 "dummy-execution-id")

  • 单工作流场景B-03 照常工作。Request/Response 两端都取 ctx.execution_id(当前 = "dummy"),过滤匹配。
  • 多工作流并发场景:所有 execution_id 相同,跨工作流的 Response 会错配到同 node_id 的别的工作流实例 → 必须 B-06 修复(execution_id 从 run_workflow 已生成的真 ID 下沉到 executor 再到 NodeContext)才能正确隔离。
  • 结论B-06 是并发正确性前置非单流功能性前置。B-03 实现完成后,单工作流可用;并发安全等 B-06。

B-07(每节点全新空 StateMachine)

  • HumanNode 取消检查 ctx.node_status.is_cancelled(&ctx.node_id) 恒 false(空状态机 get() 返回 Pending)。
  • 即使 B-07 修复(共享 self.state_machine),取消仍不生效——因为 StateMachine(state.rs)set_cancelled 方法,也无 cancel_workflow_node IPC、无前端取消按钮。
  • 结论B-07 是取消机制的必要非充分条件。取消要真正端到端生效,还需另补三件(见第七节)。B-03 的响应等待 + 超时核心功能不依赖 B-07。

四、核心机制设计

human_node.rs::execute 改为:先订阅 → 发请求 → select! 循环等响应

4.1 订阅时序铁律

tokio broadcast 通道不回放历史消息——subscribe() 调用之后发送的消息才进入该 receiver 的队列。因此必须:

subscribe()  ← 必须先于 send(Request)
send(Request)
select! { rx.recv() | timeout | cancel }

若顺序颠倒(subscribe 在 send 之后)HumanNode 的 receiver 在 Response 发出时尚不存在Response 丢失HumanNode 死等到超时。

4.2 execute 实现骨架

async fn execute(&self, ctx: NodeContext) -> NodeResult {
    let config = ctx.config.as_object().cloned().unwrap_or_default();
    let title = config.get("title").and_then(|v| v.as_str()).unwrap_or("请确认");
    let description = config.get("description").and_then(|v| v.as_str()).unwrap_or("");
    let options: Vec<String> = config.get("options")
        .and_then(|v| v.as_array())
        .map(|a| a.iter().filter_map(|v| v.as_str().map(String::from)).collect())
        .unwrap_or_else(|| vec!["同意".into(), "拒绝".into()]);
    let timeout_secs = config.get("timeout_secs").and_then(|v| v.as_u64()).unwrap_or(3600);

    // 1. 先订阅再发请求broadcast 不回放历史)
    let mut rx = ctx.event_bus.subscribe();

    // 2. 发审批请求
    ctx.event_bus.send(WorkflowEvent::HumanApprovalRequest {
        execution_id: ctx.execution_id.clone(),
        node_id: ctx.node_id.clone(),
        title: title.into(),
        description: description.into(),
        options: options.clone(),
    }).await;

    // 3. select! 循环Response / 超时 / 取消
    let deadline = tokio::time::Instant::now() + Duration::from_secs(timeout_secs);
    let mut cancel_tick = tokio::time::interval(Duration::from_millis(500));
    cancel_tick.tick().await; // 丢弃首个立即触发

    loop {
        tokio::select! {
            recv = rx.recv() => match recv {
                Ok(WorkflowEvent::HumanApprovalResponse {
                    execution_id, node_id, decision, comment
                }) if execution_id == ctx.execution_id && node_id == ctx.node_id => {
                    // decision 合法性校验
                    if !decision.is_empty() && (options.is_empty() || options.contains(&decision)) {
                        return Ok(NodeOutput::from_value(serde_json::json!({
                            "decision": decision,
                            "comment": comment.unwrap_or_default(),
                        })));
                    }
                    return Err(anyhow::anyhow!("审批决策非法: {}", decision));
                }
                Ok(_) => continue, // 其他节点/类型的事件,忽略
                Err(broadcast::error::RecvError::Lagged(n)) => {
                    tracing::warn!("HumanNode {} 漏收 {} 条事件(可能错过自身响应,继续)", ctx.node_id, n);
                    continue; // 风险:若恰好漏收自身 Response本节点将等到超时
                }
                Err(broadcast::error::RecvError::Closed) => {
                    return Err(anyhow::anyhow!("事件总线关闭,审批无法完成"));
                }
            },
            _ = tokio::time::sleep_until(deadline) => {
                return Err(anyhow::anyhow!("人工审批超时({}s)", timeout_secs));
            }
            _ = cancel_tick.tick() => {
                if ctx.node_status.is_cancelled(&ctx.node_id) {
                    return Err(anyhow::anyhow!("人工审批被取消"));
                }
            }
        }
    }
}

4.3 决策点

决策 取值 原因
options 校验 空数组时不校验(允许自由文本决策);非空时强制 decision ∈ options 空数组语义 = 自由文本审批;非空 = 枚举选项,非法值应报错而非静默放行
Lagged 处理 警告日志 + continue capacity 256 + 审批低频,漏自身 Response 概率极低;丢弃则误判超时更糟
超时来源 配置 timeout_secs,默认 3600s 保留现状默认,支持节点级配置(如"删除确认"给更长超时)
取消检查频率 500ms interval 轮询 is_cancelled 当前无主动取消信号机制,轮询是 B-07 修复前的过渡B-07 + set_cancelled 后仍需轮询(除非引入 Notify)
过滤键 execution_id + node_id 双键 node_id 单键不够(跨工作流可能重复)execution_id 单键不够(同工作流同层多 HumanNode)

五、关键时序

HumanNode.execute          run_workflow 转发器        前端 store        approve_human_approval
     │                            │                       │                      │
     │ subscribe() (rx 建位)      │                       │                      │
     │ send(Request) ──broadcast──┤                       │                      │
     │                            ├─emit workflow-event──→│                      │
     │                            │                       │ pendingApproval=…    │
     │  (select! 阻塞等 rx)       │                       │  (UI 渲染审批卡片)    │
     │                            │                       │  用户点"同意"         │
     │                            │                       │──── invoke ──────────┤
     │                            │                       │                      │ send(Response)
     │                            │                       │                      │  └─broadcast─┐
     │  rx.recv() = Response ✓ ←──┼───────────────────────┼──────────────────────┼──────────────┘
     │  过滤 exec_id+node_id 命中 │                       │                      │
     │  return NodeOutput         │                       │                      │

说明run_workflow 转发器是独立的 broadcast receiver它收到 Response 后会再 emit 一次到前端(workflow.rs:83 无差别转发所有事件)。这是无害 echo——前端 store 在调 IPC 后已本地清 pendingApproval,重复的 Response 事件不影响状态。

六、并发边界

场景 处理
同层多个 HumanNode 并行 各自独立 receiver各收全量 Responsenode_id 过滤互不干扰(execution_id 同层相同,隔离靠 node_id)
跨工作流并发 HumanNode node_id 可能重复,必须 B-06 真 execution_id 隔离;未修前并发场景有错配风险
前端未渲染审批 UI 现有 store 已接通捕获 RequestUI 组件渲染属前端独立工作B-03 后端不阻塞
前端审批后转发器 echo Response 无害store 已清 pendingApproval
事件总线容量 capacity 256审批事件低频正常不触发 Lagged
审批超时无响应 select!sleep_until(deadline) 分支返回 Err节点置 Failed工作流中止后续层

七、取消机制范围界定(B-03a / B-03b 拆分)

取消要端到端生效,当前缺三件,均不在 B-03 响应等待核心内:

  1. StateMachine::set_cancelled() 方法state.rs 当前只有 is_cancelled 查询,无对应 setter(set_waiting/set_skipped 不经转换校验Cancelled 同理可加)
  2. cancel_workflow_node IPC — 前端触发取消的入口(当前无)
  3. 前端取消按钮 + 调 IPC — UI 触发点

建议拆分

子任务 范围 依赖
B-03a HumanNode 响应等待 + 超时(本设计第四节) 无硬依赖,单工作流即可用
B-03b 取消机制:set_cancelled + cancel IPC + 前端按钮 B-07(共享 StateMachine) + 上述三件

B-03a 不依赖 B-07 即可工作(取消分支恒 false等价无取消功能不残)。todo.md 原文"B-06/B-07 修了 is_cancelled 才有意义"应理解为B-07 是取消的必要前提,但取消本身需独立补全(归入 B-03b)。

八、改动清单

文件 改动 风险
crates/df-nodes/src/human_node.rs 重写 executesubscribe → send → select! 循环 低,单文件,无外部接口变更
crates/df-workflow/src/eventbus.rs(可选) try_recv_human_approval 死代码 TODO(eventbus.rs:49-53,零调用) 低,零调用方

不改动df-core/events.rs / df-workflow/state.rs / src-tauri/commands/workflow.rs IPC / 前端 store —— 全部基础设施复用,零侵入。

B-03b 额外改动(取消机制,后续)

文件 改动
crates/df-workflow/src/state.rs set_cancelled 方法(不经转换校验,同 set_waiting)
crates/df-workflow/src/executor.rs B-07NodeContext.node_statusself.state_machine.clone() 而非 StateMachine::new()
src-tauri/src/commands/workflow.rs 新增 cancel_workflow_node IPC
src/stores/project.ts 审批 UI 加"取消"按钮,调 cancel IPC

九、测试设计

用例 方法 期望
正常审批:发匹配 Response → 收到决策 构造 EventBus + NodeContextspawn execute另起 task 发匹配(exec_id+node_id) Response 返回 NodeOutput.decision = 发送的 decision
execution_id 不匹配Response 被过滤 发不匹配 execution_id 的 Response 节点继续阻塞,短超时验证 → Err "超时"
node_id 不匹配Response 被过滤 发不匹配 node_id 的 Response 同上
超时:无 Response timeout_secs=1,不发 Response Err "人工审批超时(1s)"
decision 非法options 内无该决策 配置 options=["同意","拒绝"],发 decision="随便" Err "审批决策非法: 随便"
options 空允许自由文本 配置 options=[],发任意 decision 返回该 decision(不校验)
broadcast 关闭 → Err drop 所有 sender 后 recv Err "事件总线关闭"
Lagged 容忍(可选) 构造小容量 bus 灌满跳过,验证不 panic warn 日志 + 继续

取消分支测试(B-03b):依赖 set_cancelled + B-07本阶段跳过或 mock is_cancelled 返回 true 验证分支可达。


相关

  • 功能决策记录「工作流人工审批节点(B-03)」— 设计摘要
  • docs/todo.md B-260614-03 — 任务看板
  • crates/df-nodes/src/human_node.rs — 实施位置
  • crates/df-workflow/src/eventbus.rs — EventBus 基础设施
  • src-tauri/src/commands/workflow.rs:161approve_human_approval IPC