From 670fe61033853da7d2fd8aa4f5919b8a7b322801 Mon Sep 17 00:00:00 2001 From: qhd <1766646056@qq.com> Date: Sat, 22 Aug 2026 08:46:57 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E5=B7=A5=E4=BD=9C?= =?UTF-8?q?=E6=B5=81=E6=89=A7=E8=A1=8C=E8=AE=B0=E5=BD=95=E5=8D=A1=E5=9C=A8?= =?UTF-8?q?=E8=BF=90=E8=A1=8C=E4=B8=AD=E7=9A=84=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- workflow/service/flow/flow_ws_exec.go | 33 +++++++++++++++++++++++---- 1 file changed, 29 insertions(+), 4 deletions(-) diff --git a/workflow/service/flow/flow_ws_exec.go b/workflow/service/flow/flow_ws_exec.go index 54d3770..9b7b818 100644 --- a/workflow/service/flow/flow_ws_exec.go +++ b/workflow/service/flow/flow_ws_exec.go @@ -128,10 +128,11 @@ func handleExecute(ctx context.Context, conn *wsCommon.WsConnection, payload int if r := recover(); r != nil { glog.Errorf(execCtx, "workflow panic: %v", r) execErr = fmt.Errorf("工作流异常: %v", r) - // panic 被 recover 吞掉时,正常路径的 recordWorkflow 不会执行, - // 在此兜底把失败状态落库,避免 exec_workflow 记录卡在 Running - if !recorded && !g.IsEmpty(execId) { - recordWorkflow(saveCtx, execId, time.Since(start), execErr) + // panic 发生在 executeOrResume 内部(如节点 lambda panic)时, + // 多返回值赋值不会完成,此处 execId 可能为 0,需走兜底按会话+流程查最近"运行中"记录标记失败, + // 避免 exec_workflow 记录卡在 Running + if !recorded { + recordExecutionFailure(saveCtx, conn.SessionId, execPayload.FlowId, execId, execErr) recorded = true } _ = writeJSON(conn, &wsCommon.WsPushMsg{Type: "error", Message: "工作流异常", Error: fmt.Sprintf("%v", r)}) @@ -159,6 +160,10 @@ func handleExecute(ctx context.Context, conn *wsCommon.WsConnection, payload int glog.Infof(saveCtx, "工作流执行完成,execId: %v", execId) recordWorkflow(saveCtx, execId, time.Since(start), execErr) recorded = true + } else if execErr != nil { + // 查询/创建执行记录失败(拿不到 execId)时,兜底把该会话+工作流最近一条"运行中"记录标记为失败 + recordExecutionFailure(saveCtx, conn.SessionId, execPayload.FlowId, execId, execErr) + recorded = true } if execErr != nil { _ = writeJSON(conn, &wsCommon.WsPushMsg{Type: "error", Message: "工作流执行失败", Error: execErr.Error()}) @@ -175,6 +180,26 @@ func handleExecute(ctx context.Context, conn *wsCommon.WsConnection, payload int }() } +// recordExecutionFailure 记录一次失败状态。有 execId 直接更新该记录;拿不到 execId +// (executeOrResume 在创建记录后、返回前 panic,或查询/创建执行记录失败)时, +// 兜底按会话+工作流查最近一条仍处于"运行中"的记录标记为失败,避免前端已报错但记录卡在 Running。 +// 最近记录已是成功/失败状态则不处理(可能是上一次执行的结果,不应误改)。 +func recordExecutionFailure(ctx context.Context, sessionId string, flowId int64, execId int64, runErr error) { + if !g.IsEmpty(execId) { + recordWorkflow(ctx, execId, 0, runErr) + return + } + lastExec, err := sessionDao.ExecWorkflowDao.GetLatestBySessionAndFlow(ctx, sessionId, flowId) + if err != nil || lastExec == nil { + glog.Errorf(ctx, "兜底标记失败状态失败: sessionId=%s flowId=%d err=%v", sessionId, flowId, err) + return + } + if lastExec.Status == nil || *lastExec.Status != *flow.FlowExecutionStatusRunning.Code() { + return + } + recordWorkflow(ctx, lastExec.Id, 0, runErr) +} + // recordWorkflow 把一次工作流执行写入 exec_workflow/exec_workflow_result:运行记录 + 输出文件结果 func recordWorkflow(ctx context.Context, id int64, duration time.Duration, runErr error) { // exec_workflow 状态沿用 1-运行中,2-成功,3-失败;前端结果卡片也只识别 1/2/3