Skip to content
Open
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
70 changes: 56 additions & 14 deletions backend/modules/evaluation/domain/service/target_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,9 @@ type EvalTargetServiceImpl struct {
// 若按线上预算跑就要真等十秒,压到毫秒级才能把"重试到用尽"确定性地断言完。
// 生产路径(wire 注入)从不设置它,走默认值。
sandboxDestroyRetryBudget time.Duration
// trajectoryRetryInterval 仅用于把单测中的最终一致性等待压到毫秒级。
// 生产路径不设置,使用 defaultTrajectoryRetryInterval。
trajectoryRetryInterval time.Duration
}

const evalTargetRecordPersistTimeout = 5 * time.Second
Expand All @@ -65,6 +68,15 @@ const defaultSandboxDestroyRetryBudget = 10 * time.Second
// 用于吸收请求发起时间与实际 span 上报时间之间可能的时钟/延迟误差,避免漏掉最早的 span。
const trajectoryStartTimeBufferMS = int64(60 * 1000)

const (
trajectoryExtractAttempts = 3
defaultTrajectoryRetryInterval = 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 是 <expt_id>+"-macvm" (见 commercial operator macVMTaskID)。
// Destroy 收尾据此把 execution 分回各自 task。两者字面量相同但语义不同 (一个是 execute id 后缀、
Expand Down Expand Up @@ -516,14 +528,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,
Expand Down Expand Up @@ -1357,7 +1384,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,保持向前兼容。
Expand All @@ -1366,7 +1393,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)
}
Expand All @@ -1380,20 +1407,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)
}
})
}
Expand Down
156 changes: 149 additions & 7 deletions backend/modules/evaluation/domain/service/target_impl_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -1358,6 +1361,43 @@ 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_StartTimeBuffer 验证抽取 trajectory 时下界额外向前预留 1 分钟 buffer;
// startTimeMS 为 nil 时保持 nil 不做偏移。
func TestEvalTargetServiceImpl_ExtractTrajectory_StartTimeBuffer(t *testing.T) {
Expand Down Expand Up @@ -1389,7 +1429,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)
Expand Down Expand Up @@ -1421,6 +1462,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()

Expand All @@ -1445,7 +1487,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()
Expand All @@ -1456,7 +1498,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
Expand All @@ -1465,7 +1508,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{
Expand All @@ -1474,8 +1517,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)
Expand All @@ -1489,6 +1533,104 @@ 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(1s), 贴近生产: (attempts-1)*retryInterval=2s。
}

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+8) * time.Second):
t.Fatal("最后一次 ListTrajectory 未被调用(预算被提前耗尽)")
}
}

func TestEvalTargetServiceImpl_ValidateRuntimeParam(t *testing.T) {
t.Parallel()

Expand Down
Loading