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
78 changes: 64 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,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 是 <expt_id>+"-macvm" (见 commercial operator macVMTaskID)。
// Destroy 收尾据此把 execution 分回各自 task。两者字面量相同但语义不同 (一个是 execute id 后缀、
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,保持向前兼容。
Expand All @@ -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)
}
Expand All @@ -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)
}
})
}
Expand Down
Loading
Loading