From 855d0cca72502971daf00c20e806c5fa0adb00b9 Mon Sep 17 00:00:00 2001 From: qhd <1766646056@qq.com> Date: Tue, 25 Aug 2026 10:59:09 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E8=A1=A8=E5=90=8D=E5=B8=B8=E9=87=8F?= =?UTF-8?q?=E5=AF=B9=E9=BD=90=E7=9F=AD=E5=90=8D=E7=BA=A6=E5=AE=9A=20+=20se?= =?UTF-8?q?gment=5Findex=20=E5=8E=BB=E9=87=8D=E5=AE=88=E5=8D=AB=E9=98=B2?= =?UTF-8?q?=E8=A7=86=E9=A2=91=E6=89=B9=E9=87=8F=E8=8A=82=E7=82=B9=E7=BB=AD?= =?UTF-8?q?=E8=B7=91=E9=9D=99=E9=BB=98=E6=8D=9F=E5=9D=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- workflow/consts/public/table_name.go | 2 +- workflow/dao/flow/flow_segment_result_dao.go | 6 ++++- workflow/service/flow/lambda_node.go | 9 ++++++-- .../service/flow/lambda_segment_resume.go | 21 +++++++++++++++++- .../flow/lambda_segment_resume_test.go | 22 +++++++++++++++++++ 5 files changed, 55 insertions(+), 5 deletions(-) diff --git a/workflow/consts/public/table_name.go b/workflow/consts/public/table_name.go index 15e6bff..2d78460 100644 --- a/workflow/consts/public/table_name.go +++ b/workflow/consts/public/table_name.go @@ -22,5 +22,5 @@ const ( TableNameExecChat = "exec_chat" TableNameExecWorkflow = "exec_workflow" TableNameExecWorkflowResult = "exec_workflow_result" - TableNameFlowSegmentResult = "black_deacon_flow_segment_result" + TableNameFlowSegmentResult = "flow_segment_result" ) diff --git a/workflow/dao/flow/flow_segment_result_dao.go b/workflow/dao/flow/flow_segment_result_dao.go index 175d023..b072c9c 100644 --- a/workflow/dao/flow/flow_segment_result_dao.go +++ b/workflow/dao/flow/flow_segment_result_dao.go @@ -53,7 +53,11 @@ func (d *flowSegmentResultDao) ListByNode(ctx context.Context, execId int64, nod // 而 Save 的 ON CONFLICT DO UPDATE SET 不含 deleted_at(OmitNil 丢弃 nil), // 软删后同键重存无法复活该行 → 续跑复用静默失效(每次重执行都全量重生成)。 func (d *flowSegmentResultDao) DeleteByExecution(ctx context.Context, execId int64) error { + // 表名须用物理全名:raw Exec 不经过 GoFrame 的 config prefix(black_deacon_)自动加前缀, + // 与 update.sql 物理建表名 black_deacon_flow_segment_result 保持一致;常量是短名 + // flow_segment_result(经 Model() 自动加前缀),不能在此复用。 + const physicalTable = "black_deacon_flow_segment_result" _, err := gfdb.DB(ctx, public.DbNameBlackDeacon). - Exec(ctx, fmt.Sprintf("DELETE FROM %s WHERE execution_id = ?", public.TableNameFlowSegmentResult), execId) + Exec(ctx, fmt.Sprintf("DELETE FROM %s WHERE execution_id = ?", physicalTable), execId) return err } diff --git a/workflow/service/flow/lambda_node.go b/workflow/service/flow/lambda_node.go index 9609c0f..7bb77d3 100644 --- a/workflow/service/flow/lambda_node.go +++ b/workflow/service/flow/lambda_node.go @@ -76,8 +76,13 @@ func ModelLambda(ctx context.Context, input any) (any, error) { var totalTokens int64 var totalCost float64 if len(paramsList) > 1 { - // 段级续跑仅在"多段 + 视频模型"启用;非视频分段(批量文本等)走原逻辑零影响 - segVideo := isVideoModel(ctx, nodeInput.Config.ModelConfig.ModelId) + // 段级续跑仅在"多段 + 视频模型"启用;非视频分段(批量文本等)走原逻辑零影响。 + // 段级复用/落库还要求各段带互不重复的 segment_index(目前仅 split_shots_pipeline 产出): + // 视频节点若被批量拆分(如 split_batch_model_params),paramsList>1 但无 segment_index, + // 此时不得启用段级续跑,否则续跑会拿同一段结果拼出 N 段重复视频(静默损坏)。 + // segVideo 为 false 时走既有非段级合并路径(全量生成、不落库、不复用),全新执行行为不变, + // 后续自动 concat 判断(独立的 isVideoModel 调用)仍正常执行。 + segVideo := isVideoModel(ctx, nodeInput.Config.ModelConfig.ModelId) && hasDistinctSegmentIndex(paramsList) // 续跑(!ForceNewRun)时读取该节点已成功段;全新执行不查(BuildExecution 已清旧段),saved 为 nil → 全量重生成 var saved map[int]entity.SegmentRef diff --git a/workflow/service/flow/lambda_segment_resume.go b/workflow/service/flow/lambda_segment_resume.go index b3b12ed..52d1798 100644 --- a/workflow/service/flow/lambda_segment_resume.go +++ b/workflow/service/flow/lambda_segment_resume.go @@ -11,6 +11,25 @@ import ( // segmentGenerateMaxAttempts 视频段生成最大尝试次数(失败自动重试 1 次,共 2 次尝试),参数化可调 const segmentGenerateMaxAttempts = 2 +// hasDistinctSegmentIndex 判断各段参数是否都带 segment_index 且互不重复。 +// 段级续跑依赖 segment_index 作为段的稳定身份;缺失或重复(如视频节点误用批量拆分)时 +// 不得启用复用/落库,否则续跑会拿同一段结果拼出 N 段重复视频(静默损坏)。 +func hasDistinctSegmentIndex(paramsList []map[string]any) bool { + seen := make(map[int]bool, len(paramsList)) + for _, params := range paramsList { + v, ok := params["segment_index"] + if !ok { + return false + } + idx := gconv.Int(v) + if seen[idx] { + return false + } + seen[idx] = true + } + return len(paramsList) > 0 +} + // planSegmentResume 段级续跑决策:把 paramsList 各段映射到"是否需重新生成"。 // savedMap 为该节点已成功段(段序号 → {key,url});段在表中缺失或地址为空则需重新生成。 // 返回值与 paramsList 对齐。全新执行(savedMap 为 nil/空)时全部需生成。 @@ -52,4 +71,4 @@ func mergeSegmentOutputs(idxList []int, needGen []bool, newRes [][]map[string]an merged = append(merged, o.recs...) } return merged -} \ No newline at end of file +} diff --git a/workflow/service/flow/lambda_segment_resume_test.go b/workflow/service/flow/lambda_segment_resume_test.go index 15ffacb..a8df876 100644 --- a/workflow/service/flow/lambda_segment_resume_test.go +++ b/workflow/service/flow/lambda_segment_resume_test.go @@ -6,6 +6,28 @@ import ( "ai-agent/workflow/model/entity" ) +// hasDistinctSegmentIndex:各段都带 segment_index 且互不重复 → 允许段级续跑 +func TestHasDistinctSegmentIndex(t *testing.T) { + distinct := []map[string]any{{"segment_index": 1}, {"segment_index": 2}, {"segment_index": 3}} + if !hasDistinctSegmentIndex(distinct) { + t.Fatal("互不重复的 segment_index 应返回 true") + } + missing := []map[string]any{{"segment_index": 1}, {"other": 2}} + if hasDistinctSegmentIndex(missing) { + t.Fatal("缺失 segment_index 应返回 false") + } + duplicate := []map[string]any{{"segment_index": 1}, {"segment_index": 1}} + if hasDistinctSegmentIndex(duplicate) { + t.Fatal("重复 segment_index 应返回 false") + } + if hasDistinctSegmentIndex(nil) { + t.Fatal("空列表应返回 false") + } + if hasDistinctSegmentIndex([]map[string]any{}) { + t.Fatal("空列表(非 nil)应返回 false") + } +} + // planSegmentResume:段序号从 params 的 segment_index 读出;已成功段复用,缺失/空地址需生成 func TestPlanSegmentResume(t *testing.T) { paramsList := []map[string]any{