From 0de29988b4590a482865fe64ee42dd585f6259c3 Mon Sep 17 00:00:00 2001 From: qhd <1766646056@qq.com> Date: Tue, 25 Aug 2026 10:09:58 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20FlowExecutionInput=20=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=20ForceNewRun=EF=BC=8CBuildExecution=20=E5=85=A8=E6=96=B0?= =?UTF-8?q?=E7=94=9F=E5=89=8D=E6=B8=85=E6=97=A7=E6=AE=B5/=E6=88=90?= =?UTF-8?q?=E5=8A=9F=E5=90=8E=E6=B8=85=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.7 --- workflow/model/dto/flow/flow_execution_dto.go | 1 + workflow/service/flow/flow_graph_util.go | 3 +++ workflow/service/flow/flow_ws_exec.go | 10 ++++++++++ 3 files changed, 14 insertions(+) diff --git a/workflow/model/dto/flow/flow_execution_dto.go b/workflow/model/dto/flow/flow_execution_dto.go index 98f4eb1..5330cb0 100644 --- a/workflow/model/dto/flow/flow_execution_dto.go +++ b/workflow/model/dto/flow/flow_execution_dto.go @@ -33,6 +33,7 @@ type FlowExecutionInput struct { ConfigMap map[string]*entity.FlowNode `json:"configMap"` SessionId string `json:"sessionId" dc:"会话ID"` ExecutedNodes []ExecutedNode `json:"executedNodes"` // 已执行节点列表,包含执行状态 + ForceNewRun bool `json:"forceNewRun" dc:"是否全新执行(false=断点续跑,视频节点段级复用已成功段)"` } // ExecutedNode 已执行节点记录,包含节点ID和执行状态 diff --git a/workflow/service/flow/flow_graph_util.go b/workflow/service/flow/flow_graph_util.go index 01bf692..d752443 100644 --- a/workflow/service/flow/flow_graph_util.go +++ b/workflow/service/flow/flow_graph_util.go @@ -39,6 +39,9 @@ func BuildNodeExecutionInput(ctx context.Context, input any, flowNode entity.Flo return nil, nil, fmt.Errorf("节点:%v 进程状态为空,节点入参参数为空", flowNode.Name) } } + // 续跑必定非全新执行:checkpoint 恢复的 SavedFlowInput 里 ForceNewRun 是上次(fresh)执行留下的 true, + // 不清则 ModelLambda 误走"清段重生成"而非复用已成功段 + execInput.ForceNewRun = false } else { var ok bool execInput, ok = input.(*flowDto.FlowExecutionInput) diff --git a/workflow/service/flow/flow_ws_exec.go b/workflow/service/flow/flow_ws_exec.go index 81c86d5..0a3cb76 100644 --- a/workflow/service/flow/flow_ws_exec.go +++ b/workflow/service/flow/flow_ws_exec.go @@ -421,6 +421,7 @@ func BuildExecution(ctx context.Context, forceNewRun bool, flowId, executionId i FlowId: flowId, ConfigMap: configMap, SessionId: sessionId, + ForceNewRun: forceNewRun, } var opts []compose.Option @@ -428,6 +429,13 @@ func BuildExecution(ctx context.Context, forceNewRun bool, flowId, executionId i if forceNewRun { opts = append(opts, compose.WithForceNewRun()) } + // 全新执行前清理该执行残留段结果:forceNewRun 复用同一条 exec 记录时可能留有旧参数生成的段, + // 不清则续跑会误复用。只按 execution_id 清理(放在图启动前,避免多视频节点互相误删) + if forceNewRun { + if err := flowDao.FlowSegmentResultDao.DeleteByExecution(ctx, executionId); err != nil { + return fmt.Errorf("清理段结果失败: %v", err) + } + } _, err = runGraph.Invoke(ctx, execInput, opts...) if err != nil { // 图执行被 ctx 取消(WS 断连/用户终止):返回 context.Canceled 语义,让 recordWorkflow @@ -458,5 +466,7 @@ func BuildExecution(ctx context.Context, forceNewRun bool, flowId, executionId i } // 清理断点数据 _ = flowDao.FlowCheckpointDao.Delete(ctx, gconv.String(executionId)) + // 清理该执行段结果,下次执行无残留(失败则保留,供 reExecute 复用) + _ = flowDao.FlowSegmentResultDao.DeleteByExecution(ctx, executionId) return }