diff --git a/src-tauri/src/commands/workflow.rs b/src-tauri/src/commands/workflow.rs index 4d6dc46..4d9789c 100644 --- a/src-tauri/src/commands/workflow.rs +++ b/src-tauri/src/commands/workflow.rs @@ -67,7 +67,13 @@ pub async fn run_workflow( let mut rx = state.event_bus.subscribe(); let forward_app = app.clone(); let forward_exec_id = execution_id.clone(); + // B-260615-35: forward 任务内需查 DB 终态兜底 Lagged 丢终态事件,故 move 进 db 句柄 + let forward_db = state.db.clone(); tauri::async_runtime::spawn(async move { + // Lagged 累计次数(用于"累计达阈值"触发兜底查询) + let mut lagged_total: u64 = 0; + // Lagged 单次/累计触发 DB 终态查询兜底的阈值 + const LAGGED_PROBE_THRESHOLD: u64 = 8; loop { match rx.recv().await { Ok(event) => { @@ -88,22 +94,81 @@ pub async fn run_workflow( } } // R-PD-13: Lagged 表示消费过慢,broadcast 滑动窗口丢 n 条最旧事件。 - // broadcast 不暴露被丢事件的具体类型,无法在此精确判断是否为关键终态事件 - // (NodeCompleted/NodeFailed/HumanApprovalRequest/WorkflowCompleted/WorkflowFailed)。 - // 风险:若关键终态事件被丢,下面 emit 不到,finished 判定不触发,循环不会 break, - // 前端永久收不到工作流结束信号(依赖 DB 状态轮询兜底,但实时性差)。 - // 最小兜底:warn 记录丢失条数 + execution_id,便于事后从 tracing 追溯。 - // 后续优化方向(不在此 todo 范围): - // 1. 提升广播容量(目前常量值偏小,事件突发易满); - // 2. 关键终态事件经持久化队列/DB 重发,receiver 启动时回放未消费事件; - // 3. forward 循环加 watchdog 超时(长期无事件则主动查 DB 终态并退出)。 + // broadcast 不暴露被丢事件的具体类型,无法在此精确判断是否为关键终态事件。 + // 风险:若关键终态事件(WorkflowCompleted/WorkflowFailed)被丢,finished 判定不触发, + // 循环不会 break,前端永久收不到工作流结束信号。 + // B-260615-35 兜底:收到 Lagged 时,累计次数或单次 n 达阈值则查 DB 终态, + // 若工作流已终态则补发对应 workflow-event 并 break 退出 forward。 + // 避免依赖实时性差的 DB 轮询兜底导致前端审批/完成弹窗卡死。 Err(RecvError::Lagged(n)) => { + lagged_total = lagged_total.saturating_add(n); tracing::warn!( execution_id = %forward_exec_id, lagged = n, - "工作流事件消费滞后,丢失 {} 条事件(broadcast 不暴露被丢事件类型,关键终态事件可能永久丢失,依赖 DB 轮询兜底)", + lagged_total = lagged_total, + "工作流事件消费滞后,丢失 {} 条事件(broadcast 不暴露被丢事件类型,关键终态事件可能已丢失)", n ); + if lagged_total < LAGGED_PROBE_THRESHOLD { + continue; + } + // 累计/单次达阈值:查 DB 终态兜底 + let workflows = WorkflowRepo::new(&forward_db); + match workflows.get_by_id(&forward_exec_id).await { + Ok(Some(record)) => { + let terminal = match record.status.as_str() { + "completed" => Some(WorkflowEvent::WorkflowCompleted { + total_duration_ms: 0, + }), + "failed" => Some(WorkflowEvent::WorkflowFailed { + error: "工作流执行失败(Lagged 兜底补发,详情见 DB)" + .to_string(), + failed_node: String::new(), + }), + "cancelled" => Some(WorkflowEvent::WorkflowFailed { + error: "工作流被取消(Lagged 兜底补发)".to_string(), + failed_node: String::new(), + }), + _ => None, + }; + if let Some(synth_event) = terminal { + let payload = WorkflowEventPayload { + execution_id: forward_exec_id.clone(), + event: synth_event.clone(), + }; + if let Err(e) = forward_app.emit("workflow-event", &payload) { + tracing::warn!("Lagged 终态兜底事件转发失败: {}", e); + } + tracing::info!( + execution_id = %forward_exec_id, + status = %record.status, + "Lagged 兜底命中终态,补发 {} 并退出 forward", + match synth_event { + WorkflowEvent::WorkflowCompleted { .. } => "WorkflowCompleted", + _ => "WorkflowFailed", + } + ); + break; + } + // DB 仍非终态(running),重置累计计数继续等待后续事件 + lagged_total = 0; + } + Ok(None) => { + tracing::warn!( + execution_id = %forward_exec_id, + "Lagged 兜底查询未找到执行记录,继续等待事件" + ); + lagged_total = 0; + } + Err(e) => { + tracing::warn!( + execution_id = %forward_exec_id, + error = %e, + "Lagged 兜底查询 DB 失败,继续等待事件" + ); + lagged_total = 0; + } + } } Err(RecvError::Closed) => break, }