fix: 表名常量对齐短名约定 + segment_index 去重守卫防视频批量节点续跑静默损坏

This commit is contained in:
2026-08-25 10:59:09 +08:00
parent 01a73e5153
commit 855d0cca72
5 changed files with 55 additions and 5 deletions
+1 -1
View File
@@ -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"
)
+5 -1
View File
@@ -53,7 +53,11 @@ func (d *flowSegmentResultDao) ListByNode(ctx context.Context, execId int64, nod
// 而 Save 的 ON CONFLICT DO UPDATE SET 不含 deleted_atOmitNil 丢弃 nil),
// 软删后同键重存无法复活该行 → 续跑复用静默失效(每次重执行都全量重生成)。
func (d *flowSegmentResultDao) DeleteByExecution(ctx context.Context, execId int64) error {
// 表名须用物理全名:raw Exec 不经过 GoFrame 的 config prefixblack_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
}
+7 -2
View File
@@ -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
+20 -1
View File
@@ -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
}
}
@@ -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{