修复: B-35 broadcast Lagged 终态兜底(P0,防前端永久卡死)

workflow.rs forward 任务 move state.db.clone() 进去(原仅 executor spawn 持有)+ Lagged 累计 lagged_total(LAGGED_PROBE_THRESHOLD=8)达阈值查 WorkflowRepo::get_by_id 终态,completed/failed/cancelled 补 emit workflow-event + break 退出 forward。防 Lagged 丢终态事件致 forward 永远 rx.recv().await 等 finished → 前端审批/完成弹窗永久卡死。合成事件 total_duration_ms/failed_node 占位(broadcast 不暴露原字段,DB 为准)。批7 wdntu5adj,cargo 0
This commit is contained in:
2026-06-15 05:54:15 +08:00
parent 06a2dea1d7
commit b08adcb9da

View File

@@ -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,
}