diff --git a/backend/modules/evaluation/domain/service/target_impl.go b/backend/modules/evaluation/domain/service/target_impl.go index c40a72cc1..ca84a4a65 100644 --- a/backend/modules/evaluation/domain/service/target_impl.go +++ b/backend/modules/evaluation/domain/service/target_impl.go @@ -53,6 +53,9 @@ type EvalTargetServiceImpl struct { // 若按线上预算跑就要真等十秒,压到毫秒级才能把"重试到用尽"确定性地断言完。 // 生产路径(wire 注入)从不设置它,走默认值。 sandboxDestroyRetryBudget time.Duration + // trajectoryRetryInterval 仅用于把单测中的最终一致性等待压到毫秒级。 + // 生产路径不设置,使用 defaultTrajectoryRetryInterval。 + trajectoryRetryInterval time.Duration } const evalTargetRecordPersistTimeout = 5 * time.Second @@ -65,6 +68,23 @@ const defaultSandboxDestroyRetryBudget = 10 * time.Second // 用于吸收请求发起时间与实际 span 上报时间之间可能的时钟/延迟误差,避免漏掉最早的 span。 const trajectoryStartTimeBufferMS = int64(60 * 1000) +const ( + // trajectoryExtractAttempts / defaultTrajectoryRetryInterval 决定抽取轨迹时对"最终一致未就绪"的坚持窗口。 + // observability 侧 span 落库最终一致, 顶层 root span(唯一 ParentID 为空/"0" 的 span, 见 BuildTrajectoryFromSpans) + // 与子 span 走不同上报/落库路径, 落库延迟有长尾: 子 span 已可见时 root span 仍未落, ListTrajectory 只返回 + // {id, agent_steps} 而 RootStep=nil, IsValid()=false。此时无法用任一 agent step 顶替 root(它们都有指向未落 + // root 的真实 ParentID), 只能等。PPE 实测该长尾可达首次抽取后 180s~540s 才对 ListTrajectory 可见, 窗口过短 + // (旧值 6 次 x 10s ≈ 60s)仍会在 root span 落库前判 incomplete 放弃且不再补抽 → 轨迹永久缺失。抽取全程在 + // 后台 goroutine(脱离请求路径), 只是把 trajectory 的 UpdateEvalTargetRecord 延后, 放宽窗口无用户侧延迟代价。 + // 取 14 次 x 30s(重试期约 7min)覆盖实测长尾, 兼顾唤醒次数不过多。 + trajectoryExtractAttempts = 14 + defaultTrajectoryRetryInterval = 30 * time.Second + // trajectoryListPerAttemptTimeout 是单次 ListTrajectory RPC 的时间预算。计算后台抽取 ctx 总预算时 + // 必须按 attempts 次 RPC 预留, 否则 trace 未最终一致触发多次重试时, 最后一次 RPC 会拿到近乎 0 的 + // 剩余预算, 以 timeout=0s 立即失败并丢掉 trajectory。observability 侧实测单次约 1.5s, 取 3s 留余量。 + trajectoryListPerAttemptTimeout = 3 * time.Second +) + // sandbox mac_vm_plus_sandbox 链路: operator 给 mac_vm 那台的 execution id 加 "-macvm" 后缀 // (invokeID+"-macvm"), 且其 task 是 +"-macvm" (见 commercial operator macVMTaskID)。 // Destroy 收尾据此把 execution 分回各自 task。两者字面量相同但语义不同 (一个是 execute id 后缀、 @@ -516,14 +536,29 @@ func (e *EvalTargetServiceImpl) ExtractTrajectory(ctx context.Context, spaceID i if startTimeMS != nil { startTimeMS = gptr.Of(*startTimeMS - trajectoryStartTimeBufferMS) } - trajectories, err := e.trajectoryAdapter.ListTrajectory(ctx, spaceID, []string{traceID}, startTimeMS) - if err != nil { - return nil, err + retryInterval := e.trajectoryRetryInterval + if retryInterval <= 0 { + retryInterval = defaultTrajectoryRetryInterval } - if len(trajectories) == 0 { - return nil, nil + for attempt := 1; attempt <= trajectoryExtractAttempts; attempt++ { + trajectories, err := e.trajectoryAdapter.ListTrajectory(ctx, spaceID, []string{traceID}, startTimeMS) + if err != nil { + return nil, err + } + if len(trajectories) > 0 && trajectories[0].IsValid() { + return trajectories[0], nil + } + if attempt < trajectoryExtractAttempts { + timer := time.NewTimer(retryInterval) + select { + case <-ctx.Done(): + timer.Stop() + return nil, ctx.Err() + case <-timer.C: + } + } } - return trajectories[0], nil + return nil, errorx.New("trajectory is incomplete after %d attempts, traceID=%s", trajectoryExtractAttempts, traceID) } func (e *EvalTargetServiceImpl) AsyncExecuteTarget(ctx context.Context, spaceID, targetID, targetVersionID int64, @@ -1357,7 +1392,7 @@ func (e *EvalTargetServiceImpl) ReportInvokeRecords(ctx context.Context, param * // record.TargetID, record.TargetVersionID, record.ID, err) // } - recordTrajectory := func() error { + recordTrajectory := func(extractCtx context.Context) error { var sms *int64 // 优先用「请求发起时间」作为抽取 trajectory 的时间下界;它比 record.BaseInfo.CreatedAt(异步返回后才 stamp) // 更早,避免漏掉请求发起到返回之间的 span。为 0(未透传)时回退到 CreatedAt,保持向前兼容。 @@ -1366,7 +1401,7 @@ func (e *EvalTargetServiceImpl) ReportInvokeRecords(ctx context.Context, param * } else if record.BaseInfo != nil { sms = record.BaseInfo.CreatedAt } - trajectory, err := e.ExtractTrajectory(ctx, param.SpaceID, record.TraceID, sms) + trajectory, err := e.ExtractTrajectory(extractCtx, param.SpaceID, record.TraceID, sms) if err != nil { return errorx.Wrapf(err, "ExtractTrajectory fail, space_id: %v, trace_id: %v", param.SpaceID, record.TraceID) } @@ -1380,20 +1415,35 @@ func (e *EvalTargetServiceImpl) ReportInvokeRecords(ctx context.Context, param * if od.OutputFields == nil { od.OutputFields = map[string]*entity.Content{} } - od.OutputFields[consts.EvalTargetOutputFieldKeyTrajectory] = trajectory.ToContent(ctx) + od.OutputFields[consts.EvalTargetOutputFieldKeyTrajectory] = trajectory.ToContent(extractCtx) updateRec := &entity.EvalTargetRecord{ ID: record.ID, TraceID: record.TraceID, EvalTargetOutputData: od, } - return e.evalTargetRepo.UpdateEvalTargetRecord(ctx, updateRec, nil) + return e.evalTargetRepo.UpdateEvalTargetRecord(extractCtx, updateRec, nil) } if param.EnableExtractTrajectory == nil || *param.EnableExtractTrajectory { - goroutine.Go(ctx, func() { - time.Sleep(e.configer.GetTargetTrajectoryConf(ctx).GetExtractInterval(param.SpaceID)) - if err := recordTrajectory(); err != nil { - logs.CtxError(ctx, "extract and record trajectory fail, record_id: %v, err: %v", record.ID, err) + backgroundCtx := context.WithoutCancel(ctx) + goroutine.Go(backgroundCtx, func() { + extractInterval := e.configer.GetTargetTrajectoryConf(backgroundCtx).GetExtractInterval(param.SpaceID) + retryInterval := e.trajectoryRetryInterval + if retryInterval <= 0 { + retryInterval = defaultTrajectoryRetryInterval + } + // 等待期(extractInterval)在 backgroundCtx 上先睡完, 不占用抽取+落库的工作预算; 否则 extractInterval + // 较大时会把 extractCtx 预算耗掉, 使重试的最后一次 ListTrajectory 拿到近乎 0 的剩余而 timeout=0s。 + timer := time.NewTimer(extractInterval) + <-timer.C + // 工作预算独立于等待期, 且必须覆盖每一次 RPC(attempts 次)+ 重试间隔(attempts-1 次)+ 一次落库, + // 不能只算重试间隔——那样没给 RPC 本身留时间。 + workBudget := time.Duration(trajectoryExtractAttempts)*trajectoryListPerAttemptTimeout + + time.Duration(trajectoryExtractAttempts-1)*retryInterval + evalTargetRecordPersistTimeout + extractCtx, cancel := context.WithTimeout(backgroundCtx, workBudget) + defer cancel() + if err := recordTrajectory(extractCtx); err != nil { + logs.CtxError(backgroundCtx, "extract and record trajectory fail, record_id: %v, err: %v", record.ID, err) } }) } diff --git a/backend/modules/evaluation/domain/service/target_impl_test.go b/backend/modules/evaluation/domain/service/target_impl_test.go index b96da1e16..854c46c0d 100755 --- a/backend/modules/evaluation/domain/service/target_impl_test.go +++ b/backend/modules/evaluation/domain/service/target_impl_test.go @@ -21,6 +21,7 @@ import ( idgenmocks "github.com/coze-dev/coze-loop/backend/infra/idgen/mocks" "github.com/coze-dev/coze-loop/backend/infra/looptracer" looptracermocks "github.com/coze-dev/coze-loop/backend/infra/looptracer/mocks" + kitextrajectory "github.com/coze-dev/coze-loop/backend/kitex_gen/coze/loop/trajectory" "github.com/coze-dev/coze-loop/backend/modules/evaluation/consts" metricsmocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/component/metrics/mocks" componentmocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/component/mocks" @@ -867,7 +868,8 @@ func TestEvalTargetServiceImpl_ExecuteTarget_TrajectoryExtraction(t *testing.T) name: "trajectory extracted successfully - field added", trajectories: []*entity.Trajectory{ { - ID: gptr.Of("traj-id"), + ID: gptr.Of("traj-id"), + RootStep: &kitextrajectory.RootStep{}, }, }, expectHasField: true, @@ -1207,7 +1209,8 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_Trajectory(t *testing.T) { name: "extract trajectory success - trajectory field added", trajectories: []*entity.Trajectory{ { - ID: gptr.Of("traj-id"), + ID: gptr.Of("traj-id"), + RootStep: &kitextrajectory.RootStep{}, }, }, expectHasField: true, @@ -1358,6 +1361,73 @@ func TestEvalTargetServiceImpl_ExtractTrajectory_EmptyTraceID(t *testing.T) { assert.Nil(t, res) } +func TestEvalTargetServiceImpl_ExtractTrajectory_RetriesIncompleteTrajectory(t *testing.T) { + t.Parallel() + ctrl := gomock.NewController(t) + adapter := trajectorymocks.NewMockITrajectoryAdapter(ctrl) + traceID := "trace-x" + incomplete := &entity.Trajectory{ID: &traceID} + complete := &entity.Trajectory{ID: &traceID, RootStep: &kitextrajectory.RootStep{}} + + gomock.InOrder( + adapter.EXPECT().ListTrajectory(gomock.Any(), int64(1), []string{traceID}, nil). + Return([]*entity.Trajectory{incomplete}, nil), + adapter.EXPECT().ListTrajectory(gomock.Any(), int64(1), []string{traceID}, nil). + Return([]*entity.Trajectory{complete}, nil), + ) + + svc := &EvalTargetServiceImpl{trajectoryAdapter: adapter, trajectoryRetryInterval: time.Millisecond} + got, err := svc.ExtractTrajectory(context.Background(), 1, traceID, nil) + require.NoError(t, err) + assert.Same(t, complete, got) +} + +func TestEvalTargetServiceImpl_ExtractTrajectory_RejectsPersistentlyIncompleteTrajectory(t *testing.T) { + t.Parallel() + ctrl := gomock.NewController(t) + adapter := trajectorymocks.NewMockITrajectoryAdapter(ctrl) + traceID := "trace-x" + incomplete := &entity.Trajectory{ID: &traceID} + adapter.EXPECT().ListTrajectory(gomock.Any(), int64(1), []string{traceID}, nil). + Return([]*entity.Trajectory{incomplete}, nil).Times(trajectoryExtractAttempts) + + svc := &EvalTargetServiceImpl{trajectoryAdapter: adapter, trajectoryRetryInterval: time.Millisecond} + got, err := svc.ExtractTrajectory(context.Background(), 1, traceID, nil) + require.Error(t, err) + assert.Contains(t, err.Error(), "trajectory is incomplete") + assert.Nil(t, got) +} + +// TestEvalTargetServiceImpl_ExtractTrajectory_RetryWindowCoversRootSpanLanding 复现线上偶发丢轨迹的第二层根因: +// observability 侧 span 落库最终一致, 抽取发起时顶层 root span 尚未可见, ListTrajectory 只返回 {id, agent_steps} +// 而 RootStep=nil → IsValid()=false。root span 往往在数十秒后才落库, 因此重试必须坚持到 root_step 出现; +// 若窗口只有旧的 3 次, root_step 在第 5 次才出现时会被判 incomplete 丢弃。这里让前 4 次 incomplete、第 5 次 complete, +// 断言仍能抽到——等价要求 trajectoryExtractAttempts >= 5。 +func TestEvalTargetServiceImpl_ExtractTrajectory_RetryWindowCoversRootSpanLanding(t *testing.T) { + t.Parallel() + ctrl := gomock.NewController(t) + adapter := trajectorymocks.NewMockITrajectoryAdapter(ctrl) + traceID := "trace-late-root" + incomplete := &entity.Trajectory{ID: &traceID} // 仅 agent_steps, root span 未落库 + complete := &entity.Trajectory{ID: &traceID, RootStep: &kitextrajectory.RootStep{}} // root span 落库后 + const rootSpanLandsOnAttempt = 5 + + call := 0 + adapter.EXPECT().ListTrajectory(gomock.Any(), int64(1), []string{traceID}, nil). + DoAndReturn(func(_ context.Context, _ int64, _ []string, _ *int64) ([]*entity.Trajectory, error) { + call++ + if call >= rootSpanLandsOnAttempt { + return []*entity.Trajectory{complete}, nil + } + return []*entity.Trajectory{incomplete}, nil + }).MinTimes(rootSpanLandsOnAttempt) + + svc := &EvalTargetServiceImpl{trajectoryAdapter: adapter, trajectoryRetryInterval: time.Millisecond} + got, err := svc.ExtractTrajectory(context.Background(), 1, traceID, nil) + require.NoError(t, err) + assert.Same(t, complete, got) +} + // TestEvalTargetServiceImpl_ExtractTrajectory_StartTimeBuffer 验证抽取 trajectory 时下界额外向前预留 1 分钟 buffer; // startTimeMS 为 nil 时保持 nil 不做偏移。 func TestEvalTargetServiceImpl_ExtractTrajectory_StartTimeBuffer(t *testing.T) { @@ -1389,7 +1459,8 @@ func TestEvalTargetServiceImpl_ExtractTrajectory_StartTimeBuffer(t *testing.T) { require.NotNil(t, got) assert.Equal(t, *tt.wantOut, *got) } - return nil, nil + traceID := "trace-x" + return []*entity.Trajectory{{ID: &traceID, RootStep: &kitextrajectory.RootStep{}}}, nil }) svc := &EvalTargetServiceImpl{trajectoryAdapter: adapter} _, err := svc.ExtractTrajectory(ctx, spaceID, "trace-x", tt.in) @@ -1421,6 +1492,7 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryStartTime(t *testin for _, tt := range tests { tt := tt t.Run(tt.name, func(t *testing.T) { + requestCtx, cancelRequest := context.WithCancel(ctx) ctrl := gomock.NewController(t) defer ctrl.Finish() @@ -1445,7 +1517,7 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryStartTime(t *testin AsyncUnixMS: tt.asyncUnixMS, } - repo.EXPECT().GetEvalTargetRecordByIDAndSpaceID(ctx, param.SpaceID, param.RecordID).Return(record, nil) + repo.EXPECT().GetEvalTargetRecordByIDAndSpaceID(requestCtx, param.SpaceID, param.RecordID).Return(record, nil) repo.EXPECT().SaveEvalTargetRecord(gomock.Any(), gomock.Any(), gomock.Any()).Return(nil) repo.EXPECT().UpdateEvalTargetRecord(gomock.Any(), gomock.Any(), gomock.Any()).AnyTimes().Return(nil) configer.EXPECT().GetErrCtrl(gomock.Any()).Return(&entity.ExptErrCtrl{}).AnyTimes() @@ -1456,7 +1528,8 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryStartTime(t *testin gotStartCh := make(chan int64, 1) trajectoryAdapter.EXPECT(). ListTrajectory(gomock.Any(), spaceID, gomock.Any(), gomock.Any()). - DoAndReturn(func(_ context.Context, _ int64, _ []string, startMS *int64) ([]*entity.Trajectory, error) { + DoAndReturn(func(extractCtx context.Context, _ int64, _ []string, startMS *int64) ([]*entity.Trajectory, error) { + require.NoError(t, extractCtx.Err(), "后台轨迹抽取不应继承已取消的请求 context") var v int64 if startMS != nil { v = *startMS @@ -1465,7 +1538,7 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryStartTime(t *testin case gotStartCh <- v: default: } - return []*entity.Trajectory{{ID: gptr.Of("traj")}}, nil + return []*entity.Trajectory{{ID: gptr.Of("traj"), RootStep: &kitextrajectory.RootStep{}}}, nil }) svc := &EvalTargetServiceImpl{ @@ -1474,8 +1547,9 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryStartTime(t *testin configer: configer, } - err := svc.ReportInvokeRecords(ctx, param) + err := svc.ReportInvokeRecords(requestCtx, param) require.NoError(t, err) + cancelRequest() // 等异步抽取 goroutine(sleep 1s interval)完成 time.Sleep(1200 * time.Millisecond) @@ -1489,6 +1563,106 @@ func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryStartTime(t *testin } } +// TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryBudgetCoversAllAttempts 验证: +// 后台轨迹抽取的超时预算(extractCtx)必须覆盖"每一次 ListTrajectory RPC + 重试间隔 + 落库", +// 且不被等待期(extractInterval)蚕食。线上根因是: +// - 等待期 extractInterval 与工作期共用一个 extractCtx 预算; +// - 工作期预算里只算了 (attempts-1)*retryInterval + persistTimeout, 完全没给 attempts 次 RPC 留时间。 +// +// 于是当 trace 尚未最终一致、需要多次重试、每次 ListTrajectory RPC 又各耗一定时间时, 最后一次 RPC +// 拿到的 ctx 剩余会 <= 0, 以 timeout=0s 失败(实测 actual≈1.48s), trajectory 丢失。 +// +// 本测试模拟前若干次抽取到"不完整"轨迹(触发重试)且每次 RPC 各耗 800ms; 断言最后一次 RPC 被调用时, +// 其 ctx 剩余预算仍能容纳一次落库(>= persistTimeout)。修复前: 最后一次 RPC 要么根本没机会跑(ctx 已 Done), +// 要么剩余远不足 persistTimeout → 断言失败。 +func TestEvalTargetServiceImpl_ReportInvokeRecords_TrajectoryBudgetCoversAllAttempts(t *testing.T) { + // do not run in parallel: involves real time for the async trajectory goroutine + spaceID := int64(1) + // extractInterval 取一个较大值, 复现"等待期与工作期共用预算";用 3s 保证 UT 快速。 + const extractIntervalSec = int64(3) + const perRPCCost = 800 * time.Millisecond // 模拟单次 ListTrajectory RPC 耗时 + + requestCtx, cancelRequest := context.WithCancel(context.Background()) + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + repo := repomocks.NewMockIEvalTargetRepo(ctrl) + configer := componentmocks.NewMockIConfiger(ctrl) + trajectoryAdapter := trajectorymocks.NewMockITrajectoryAdapter(ctrl) + + record := &entity.EvalTargetRecord{ + ID: 10, + SpaceID: spaceID, + Status: gptr.Of(entity.EvalTargetRunStatusAsyncInvoking), + EvalTargetOutputData: &entity.EvalTargetOutputData{}, + TraceID: "trace-budget", + BaseInfo: &entity.BaseInfo{CreatedAt: gptr.Of(int64(2_000_000))}, + } + param := &entity.ReportTargetRecordParam{ + SpaceID: spaceID, + RecordID: record.ID, + Status: entity.EvalTargetRunStatusSuccess, + OutputData: &entity.EvalTargetOutputData{}, + } + + repo.EXPECT().GetEvalTargetRecordByIDAndSpaceID(requestCtx, param.SpaceID, param.RecordID).Return(record, nil) + repo.EXPECT().SaveEvalTargetRecord(gomock.Any(), gomock.Any(), gomock.Any()).Return(nil) + repo.EXPECT().UpdateEvalTargetRecord(gomock.Any(), gomock.Any(), gomock.Any()).AnyTimes().Return(nil) + configer.EXPECT().GetErrCtrl(gomock.Any()).Return(&entity.ExptErrCtrl{}).AnyTimes() + configer.EXPECT().GetTargetTrajectoryConf(gomock.Any()).AnyTimes().Return(&entity.TargetTrajectoryConf{ + SpaceExtractIntervalSecond: map[int64]int64{spaceID: extractIntervalSec}, + }) + + traceID := record.TraceID + incomplete := &entity.Trajectory{ID: &traceID} // 无 RootStep → IsValid()=false, 触发重试 + complete := &entity.Trajectory{ID: &traceID, RootStep: &kitextrajectory.RootStep{}} // 最后一次抽到完整 + + lastRemainingCh := make(chan time.Duration, 1) + call := 0 + trajectoryAdapter.EXPECT(). + ListTrajectory(gomock.Any(), spaceID, gomock.Any(), gomock.Any()). + DoAndReturn(func(extractCtx context.Context, _ int64, _ []string, _ *int64) ([]*entity.Trajectory, error) { + call++ + time.Sleep(perRPCCost) // 模拟 RPC 真实耗时 + if call >= trajectoryExtractAttempts { + var remaining time.Duration + if dl, ok := extractCtx.Deadline(); ok { + remaining = time.Until(dl) + } else { + remaining = time.Hour + } + select { + case lastRemainingCh <- remaining: + default: + } + return []*entity.Trajectory{complete}, nil + } + return []*entity.Trajectory{incomplete}, nil + }).Times(trajectoryExtractAttempts) + + svc := &EvalTargetServiceImpl{ + evalTargetRepo: repo, + trajectoryAdapter: trajectoryAdapter, + configer: configer, + // 用毫秒级 retryInterval 把重试等待压到可忽略: 本测试只验"工作预算覆盖 attempts 次 RPC", + // 与重试间隔长短无关; 用生产默认(10s)会让 UT 空等数十秒。 + trajectoryRetryInterval: time.Millisecond, + } + + err := svc.ReportInvokeRecords(requestCtx, param) + require.NoError(t, err) + cancelRequest() // 模拟请求返回后 ctx 被取消 + + select { + case remaining := <-lastRemainingCh: + // 最后一次 ListTrajectory 返回后, 还要走 UpdateEvalTargetRecord 落库(persistTimeout)。 + assert.GreaterOrEqual(t, remaining, evalTargetRecordPersistTimeout, + "最后一次抽取时 ctx 剩余预算不足以落库, 会因 timeout=0s 丢 trajectory") + case <-time.After(time.Duration(extractIntervalSec)*time.Second + time.Duration(trajectoryExtractAttempts)*perRPCCost + 8*time.Second): + t.Fatal("最后一次 ListTrajectory 未被调用(预算被提前耗尽)") + } +} + func TestEvalTargetServiceImpl_ValidateRuntimeParam(t *testing.T) { t.Parallel()