Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 7 additions & 9 deletions backend/modules/evaluation/application/wire_gen.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 12 additions & 6 deletions backend/modules/evaluation/domain/service/expt_result_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -1722,20 +1722,26 @@ func (b *PayloadBuilder) fillExptTurnResultFilters(ctx context.Context, createdD
}
}
}
// ★ 三层 nil 守卫不可省: buildTargetOutput 末尾会为「turn_result 有 target_result_id 但
// BatchGetRecordByIDs 未命中 record」的行构造 stub(仅 ID/SpaceID/ItemID/TurnID/Status,
// 无 EvalTargetOutputData)。只判 map 命中(ok)会在 stub 行直接 NPE 打挂整个调度器 goroutine
// → 实验被 HandleEventErr 置 Failed。stub 行本就没有 target 数据可填, 跳过即可。
evalTargetOutput, ok := exptResultBuilder.turnResultID2TargetOutput[exptTurnResult.ID]
if ok {
for outputFieldKey, outputFieldValue := range evalTargetOutput.EvalTargetRecord.EvalTargetOutputData.OutputFields {
if ok && evalTargetOutput != nil && evalTargetOutput.EvalTargetRecord != nil &&
evalTargetOutput.EvalTargetRecord.EvalTargetOutputData != nil {
outputData := evalTargetOutput.EvalTargetRecord.EvalTargetOutputData
for outputFieldKey, outputFieldValue := range outputData.OutputFields {
exptTurnResultFilter.EvalTargetData[outputFieldKey] = outputFieldValue.GetText()
}
// 填充 eval_target_metrics
if evalTargetOutput.EvalTargetRecord.EvalTargetOutputData.EvalTargetUsage != nil {
usage := evalTargetOutput.EvalTargetRecord.EvalTargetOutputData.EvalTargetUsage
if outputData.EvalTargetUsage != nil {
usage := outputData.EvalTargetUsage
exptTurnResultFilter.EvalTargetMetrics["input_tokens"] = usage.InputTokens
exptTurnResultFilter.EvalTargetMetrics["output_tokens"] = usage.OutputTokens
exptTurnResultFilter.EvalTargetMetrics["total_tokens"] = usage.TotalTokens
}
if evalTargetOutput.EvalTargetRecord.EvalTargetOutputData.TimeConsumingMS != nil {
exptTurnResultFilter.EvalTargetMetrics["total_latency"] = *evalTargetOutput.EvalTargetRecord.EvalTargetOutputData.TimeConsumingMS
if outputData.TimeConsumingMS != nil {
exptTurnResultFilter.EvalTargetMetrics["total_latency"] = *outputData.TimeConsumingMS
}
}
evaluatorScoreCorrected, ok := exptResultBuilder.turnResultID2ScoreCorrected[exptTurnResult.ID]
Expand Down
95 changes: 95 additions & 0 deletions backend/modules/evaluation/domain/service/expt_result_impl_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6093,6 +6093,101 @@ func TestExptResultBuilder_FillExptTurnResultFilters_RecalculateWeightedScore(t
})
}

// TestExptResultBuilder_FillExptTurnResultFilters_TargetOutputNilGuard 回归: buildTargetOutput 为
// 「turn_result 有 target_result_id 但 record 未命中」的行构造的 stub 只带 ID/SpaceID/ItemID/TurnID/Status,
// 没有 EvalTargetOutputData。fillExptTurnResultFilters 若只判 map 命中就解引用 .EvalTargetOutputData.OutputFields
// 会 NPE 打挂调度器 goroutine → 实验被置 Failed(线上 expt 7590117239516698626 即此故障)。
func TestExptResultBuilder_FillExptTurnResultFilters_TargetOutputNilGuard(t *testing.T) {
ctx := context.Background()

newBuilder := func(targetOutput *entity.TurnTargetOutput) *PayloadBuilder {
return &PayloadBuilder{
SpaceID: 100,
BaselineExptID: 1,
ScoreCalculator: NewEvaluatorScoreCalculator(nil, nil),
BaseExptItemResultDO: []*entity.ExptItemResult{
{ItemID: 1, ItemIdx: 1, Status: entity.ItemRunState_Success},
},
BaseExptTurnResultDO: []*entity.ExptTurnResult{
{ID: 1, ItemID: 1, TurnID: 0},
},
ExptResultBuilders: []*ExptResultBuilder{
{
exptDO: &entity.Experiment{ID: 1},
turnResultID2TargetOutput: map[int64]*entity.TurnTargetOutput{
1: targetOutput,
},
},
},
}
}

t.Run("stub record 无 EvalTargetOutputData 时不 panic 且跳过 target 数据", func(t *testing.T) {
builder := newBuilder(&entity.TurnTargetOutput{
EvalTargetRecord: &entity.EvalTargetRecord{
ID: 999,
SpaceID: 100,
ItemID: 1,
TurnID: 0,
Status: gptr.Of(entity.EvalTargetRunStatusAsyncInvoking),
},
})

assert.NotPanics(t, func() {
assert.NoError(t, builder.fillExptTurnResultFilters(ctx, nil, 0, 1))
})
assert.Len(t, builder.ExptTurnResultFilters, 1)
assert.Empty(t, builder.ExptTurnResultFilters[0].EvalTargetData)
assert.Empty(t, builder.ExptTurnResultFilters[0].EvalTargetMetrics)
})

t.Run("TurnTargetOutput 或 EvalTargetRecord 为 nil 时不 panic", func(t *testing.T) {
for name, output := range map[string]*entity.TurnTargetOutput{
"nil TurnTargetOutput": nil,
"nil EvalTargetRecord": {EvalTargetRecord: nil},
} {
builder := newBuilder(output)
assert.NotPanics(t, func() {
assert.NoError(t, builder.fillExptTurnResultFilters(ctx, nil, 0, 1))
}, name)
assert.Len(t, builder.ExptTurnResultFilters, 1, name)
assert.Empty(t, builder.ExptTurnResultFilters[0].EvalTargetData, name)
}
})

t.Run("正常 record 仍正确填充 target 数据与 metrics", func(t *testing.T) {
builder := newBuilder(&entity.TurnTargetOutput{
EvalTargetRecord: &entity.EvalTargetRecord{
ID: 999,
SpaceID: 100,
EvalTargetOutputData: &entity.EvalTargetOutputData{
OutputFields: map[string]*entity.Content{
"actual_output": {
ContentType: gptr.Of(entity.ContentTypeText),
Text: gptr.Of("hello"),
},
},
EvalTargetUsage: &entity.EvalTargetUsage{
InputTokens: 3,
OutputTokens: 5,
TotalTokens: 8,
},
TimeConsumingMS: gptr.Of(int64(120)),
},
},
})

assert.NoError(t, builder.fillExptTurnResultFilters(ctx, nil, 0, 1))
assert.Len(t, builder.ExptTurnResultFilters, 1)
got := builder.ExptTurnResultFilters[0]
assert.Equal(t, "hello", got.EvalTargetData["actual_output"])
assert.Equal(t, int64(3), got.EvalTargetMetrics["input_tokens"])
assert.Equal(t, int64(5), got.EvalTargetMetrics["output_tokens"])
assert.Equal(t, int64(8), got.EvalTargetMetrics["total_tokens"])
assert.Equal(t, int64(120), got.EvalTargetMetrics["total_latency"])
})
}

// TestExptResultServiceImpl_RecalculateWeightedScore 测试 RecalculateWeightedScore 函数(2707-2802行)
func TestExptResultServiceImpl_RecalculateWeightedScore(t *testing.T) {
ctrl := gomock.NewController(t)
Expand Down
51 changes: 28 additions & 23 deletions backend/modules/evaluation/domain/service/expt_run_item_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -458,25 +458,9 @@ func (e *ExptItemEvalCtxExecutor) CompleteItemRun(ctx context.Context, eiec *ent
}
}

// 仅 item 评测成功才推送 item-complete;失败不发(下游只消费成功行)。
if evalErr == nil && e.itemCompletePublisher != nil {
completeEvent := buildItemCompleteEvent(eiec)
// 记录 enable_analysis 的固化取值与来源链路:下游 gate 依赖此值,
// 便于排查"评测对象已开分析但 item-complete 未发"(EnableAnalysis 固化断点)。
var hasTarget, hasVersion, hasSandbox bool
if expt := eiec.Expt; expt != nil && expt.Target != nil {
hasTarget = true
if ver := expt.Target.EvalTargetVersion; ver != nil {
hasVersion = true
hasSandbox = ver.SandboxAgent != nil
}
}
logs.CtxInfo(ctx, "[ExptTurnEval] item complete enable_analysis resolved, expt_id: %v, item_id: %v, enable_analysis: %v, has_target: %v, has_version: %v, has_sandbox_agent: %v",
event.ExptID, event.EvalSetItemID, completeEvent.EnableAnalysis, hasTarget, hasVersion, hasSandbox)
if err := e.itemCompletePublisher.PublishItemComplete(ctx, completeEvent); err != nil {
logs.CtxWarn(ctx, "[ExptTurnEval] publish item complete event failed, expt_id: %v, item_id: %v, err: %v", event.ExptID, event.EvalSetItemID, err)
}
}
// item-complete(success) MQ 发送点已后移到链路B(scheduler daemon 的 recordEvalItemRunLogs),
// 在 RecordItemRunLogs 写完读侧三张表后才发,消除"下游收到 success 反查读侧却未就绪"的竞态。
// 此处仅写 result_state=Logged,不再发 MQ。

if e.evalErrNeedTerminateExpt(persistCtx, event.SpaceID, evalErr) {
logs.CtxWarn(ctx, "[ExptTurnEval] found error which should terminate expt, expt_id: %v, expt_run_id: %v, item_id: %v, err: %v", event.ExptID, event.ExptRunID, event.EvalSetItemID, evalErr)
Expand Down Expand Up @@ -551,7 +535,7 @@ func buildItemCompleteEvent(eiec *entity.ExptItemEvalCtx) *component.ItemComplet
}

// version_name / dataset_key: 按 item 归属集从内存查找(GetDetail 已批量拉全所有集详情)。
if es := findEvalSetForItem(eiec, datasetID); es != nil {
if es := findEvalSetForItem(eiec.Expt, datasetID); es != nil {
ev.DatasetKey = es.DatasetKey
if ver := es.EvaluationSetVersion; ver != nil {
if ev.DatasetVersionID == "" {
Expand All @@ -564,13 +548,34 @@ func buildItemCompleteEvent(eiec *entity.ExptItemEvalCtx) *component.ItemComplet
return ev
}

// findEvalSetForItem 从实验详情里找 item 归属集的 EvaluationSet(含 version/dataset_key)。
// buildItemCompleteEventFromScheduler 供链路B(scheduler daemon)组装 item-complete 事件。
// 链路B 循环里只有 event(*ExptScheduleEvent) + item(*ExptEvalItem, 仅 ItemID 可信) + expt(全量详情),
// 缺 ItemKey/归属集/per-item 版本; 调用方须在循环外用 expt_item_ref + BatchGetEvaluationSetItems 批量补出:
// - evalSetItem: 提供 ItemKey / SpaceID / EvaluationSetID(归属集);
// - evalSetVersionID: 该 item 归属集的 per-item 版本(来自 expt_item_ref, 多集非主集也正确)。
// 切勿用 ExptEvalItem.EvalSetVersionID —— 那是 scanIncompleteAndComplete 硬编码的主集版本(张冠李戴)。
//
// 组装逻辑复用 buildItemCompleteEvent(构造最小 ExptItemEvalCtx), 与链路A 逐字段等价、单一实现不漂移。
func buildItemCompleteEventFromScheduler(spaceID, exptID, exptRunID int64, expt *entity.Experiment, item *entity.ExptEvalItem, evalSetItem *entity.EvaluationSetItem, evalSetVersionID int64) *component.ItemCompleteEvent {
eiec := &entity.ExptItemEvalCtx{
Event: &entity.ExptItemEvalEvent{
SpaceID: spaceID,
ExptID: exptID,
ExptRunID: exptRunID,
EvalSetItemID: item.ItemID,
},
Expt: expt,
EvalSetItem: evalSetItem,
EvalSetVersionID: evalSetVersionID,
}
return buildItemCompleteEvent(eiec)
}

// 按 experiment.eval_set_source_type 显式分流(权威分流开关,DB not null default 1):
// - MultiSetConfig(2) 新实验: 从 EvalSetDetails 按 datasetID 匹配归属集(GetDetail 已批量填充所有集详情);
// 匹配不到即返回 nil,不回退主集,避免把主集版本误安到非主集 item 上(张冠李戴)。
// - SingleSet(1) 老实验/单评测集: 直接用主集单数字段 eiec.Expt.EvalSet。
func findEvalSetForItem(eiec *entity.ExptItemEvalCtx, datasetID int64) *entity.EvaluationSet {
expt := eiec.Expt
func findEvalSetForItem(expt *entity.Experiment, datasetID int64) *entity.EvaluationSet {
if expt == nil {
return nil
}
Expand Down
Loading
Loading