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
Original file line number Diff line number Diff line change
Expand Up @@ -517,16 +517,16 @@ func TestRecordEvalItemRunLogs_RejectsNonTerminalState(t *testing.T) {
assert.ErrorContains(t, err, "invalid item run state")
}

// TestExptMangerImpl_CompleteExpt_TerminatedNotFailed 终态推导双向覆盖(tasks 6.6,design D5):
// - 有 terminated 行、无 fail 行 → Success(不能因为用户主动终止就把实验判失败)
// - 有 fail 行 → 仍 Failed(不误伤既有失败判定)
// TestExptMangerImpl_CompleteExpt_TerminatedNotFailed 终态推导覆盖:
// 只要实验仍存在 terminated 行(用户终止后重试个别 item 的收尾场景),
// 实验一律收敛为 Terminated,terminated 优先级最高、盖过 fail/success 推导。
func TestExptMangerImpl_CompleteExpt_TerminatedNotFailed(t *testing.T) {
ctx := context.Background()
session := &entity.Session{UserID: "test_user"}
const exptID, spaceID = int64(123), int64(789)
runID := int64(456)

t.Run("terminated_without_fail_is_success", func(t *testing.T) {
t.Run("terminated_without_fail_is_terminated", func(t *testing.T) {
ctrl := gomock.NewController(t)
defer ctrl.Finish()
mgr := newTestExptManager(ctrl)
Expand All @@ -552,11 +552,11 @@ func TestExptMangerImpl_CompleteExpt_TerminatedNotFailed(t *testing.T) {

require.NoError(t, mgr.CompleteExpt(ctx, exptID, &runID, spaceID, session))
require.NotNil(t, captured)
assert.Equal(t, entity.ExptStatus_Success, captured.Status,
"D5: terminated 是用户主动放弃,不参与 Failed 推导")
assert.Equal(t, entity.ExptStatus_Terminated, captured.Status,
"仍有 terminated 行时实验一律收敛为 Terminated")
})

t.Run("fail_still_failed", func(t *testing.T) {
t.Run("terminated_with_fail_is_terminated", func(t *testing.T) {
ctrl := gomock.NewController(t)
defer ctrl.Finish()
mgr := newTestExptManager(ctrl)
Expand All @@ -582,7 +582,8 @@ func TestExptMangerImpl_CompleteExpt_TerminatedNotFailed(t *testing.T) {

require.NoError(t, mgr.CompleteExpt(ctx, exptID, &runID, spaceID, session))
require.NotNil(t, captured)
assert.Equal(t, entity.ExptStatus_Failed, captured.Status, "有 fail 行仍必须判 Failed")
assert.Equal(t, entity.ExptStatus_Terminated, captured.Status,
"terminated 优先级最高,即便同时有 fail 行也判 Terminated")
})
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -613,22 +613,15 @@ func (e *ExptMangerImpl) CompleteExpt(ctx context.Context, exptID int64, exptRun
return err
}

if err := e.statsRepo.UpdateByExptID(ctx, exptID, spaceID, &entity.ExptStats{
SuccessItemCnt: int32(stats.SuccessItemCnt),
PendingItemCnt: int32(stats.PendingItemCnt),
FailItemCnt: int32(stats.FailItemCnt),
ProcessingItemCnt: int32(stats.ProcessingItemCnt),
TerminatedItemCnt: int32(stats.TerminatedItemCnt),
}); err != nil {
return err
}

status := opt.Status
if !entity.IsExptFinished(status) {
// ⚠️ 判据中**不含** stats.TerminatedItemCnt(design D5):terminated 是用户主动放弃的行,不是跑失败。
// 保留它会让「用户终止了 2 行、其余全绿」的实验渲染成红色 Failed 并触发失败侧 webhook。
// 实验级 Kill 不受影响:它走 opt.Status=Terminated,IsExptFinished 为真、根本不进这个分支。
if stats.FailItemCnt > 0 || stats.ProcessingItemCnt > 0 || stats.PendingItemCnt > 0 {
// 只要实验里仍存在终止的 item,实验就必须收敛为 Terminated,而不是塌成 Success/Failed。
// 场景:实验被终止后手动重试个别 item,重试的 item 跑完再次进入 CompleteExpt,
// 此时其余 item 仍处终止态;若忽略 TerminatedItemCnt 会误判成功/失败,
// 前端也就只显示成功/失败数、不再展示终止/执行中/待执行。terminated 优先级最高。
if stats.TerminatedItemCnt > 0 {
status = entity.ExptStatus_Terminated
} else if stats.FailItemCnt > 0 || stats.ProcessingItemCnt > 0 || stats.PendingItemCnt > 0 {
status = entity.ExptStatus_Failed
} else {
status = entity.ExptStatus_Success
Expand Down Expand Up @@ -663,6 +656,9 @@ func (e *ExptMangerImpl) CompleteExpt(ctx context.Context, exptID int64, exptRun
// 先置终态它就一条都查不到(静默不释放)。理由见 terminateIncompleteItemRunLogs。
e.terminateIncompleteItemRunLogs(ctx, got, exptRunID)

// statsDirty 标记本次是否真的把 in-flight 行改写成了终态;只有为真时才在收口后
// 重算并回写 stats 表,避免正常完成/Failed 路径多做一次全量扫描。
statsDirty := false
if !opt.NoCompleteItemTurn {
incompleteTurnIDs, err := e.exptResultService.GetIncompleteTurns(ctx, exptID, spaceID, session)
if err != nil {
Expand All @@ -685,6 +681,7 @@ func (e *ExptMangerImpl) CompleteExpt(ctx context.Context, exptID int64, exptRun
}
// 在实验行状态更新完成后,更新 ExptTurnResultFilter
if len(terminatedItemIDSet) > 0 {
statsDirty = true
terminatedItemIDs := maps.ToSlice(terminatedItemIDSet, func(k int64, v bool) int64 {
return k
})
Expand Down Expand Up @@ -717,6 +714,27 @@ func (e *ExptMangerImpl) CompleteExpt(ctx context.Context, exptID int64, exptRun
}
}

// stats 回写放到终止收口之后:terminateIncompleteItemRunLogs / terminateItemTurns
// 已把 in-flight 行由 Processing 改写成 Terminated,此处必须按 item 主表现值重算,
// 否则 stats 表会残留虚高的 processing_cnt(≈终止瞬间并发数)、少算 terminated_cnt,
// 前端就显示成「执行中 = 并发数 + 1」。只有真正终止了行(statsDirty)才重算,
// 正常完成 / Failed 路径直接复用首次快照,不多做一次全量扫描。
if statsDirty {
stats, err = e.exptResultService.CalculateStats(ctx, exptID, spaceID, session)
if err != nil {
return err
}
}
if err := e.statsRepo.UpdateByExptID(ctx, exptID, spaceID, &entity.ExptStats{
SuccessItemCnt: int32(stats.SuccessItemCnt),
PendingItemCnt: int32(stats.PendingItemCnt),
FailItemCnt: int32(stats.FailItemCnt),
ProcessingItemCnt: int32(stats.ProcessingItemCnt),
TerminatedItemCnt: int32(stats.TerminatedItemCnt),
}); err != nil {
return err
}

exptDo := &entity.Experiment{
ID: exptID,
SpaceID: spaceID,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1623,15 +1623,31 @@ func TestExptMangerImpl_CompleteExpt(t *testing.T) {
}, nil)

// Mock stats calculation
// 走 Terminated 且有真实终止行 → 收口后按新语义重算一次:
// 第一次(终止前)返回带 processing 的旧快照,第二次(重算)返回终止后的正确计数。
// 断言写入 stats 表的必须是第二次(重算)的值。
calcCall := 0
mgr.exptResultService.(*svcMocks.MockExptResultService).
EXPECT().
CalculateStats(ctx, int64(123), int64(789), session).
Return(&entity.ExptCalculateStats{
SuccessItemCnt: 3,
FailItemCnt: 1,
ProcessingItemCnt: 0,
TerminatedItemCnt: 0,
}, nil)
DoAndReturn(func(_ context.Context, _, _ int64, _ *entity.Session) (*entity.ExptCalculateStats, error) {
calcCall++
if calcCall == 1 {
return &entity.ExptCalculateStats{
SuccessItemCnt: 3,
FailItemCnt: 1,
ProcessingItemCnt: 2,
TerminatedItemCnt: 0,
}, nil
}
return &entity.ExptCalculateStats{
SuccessItemCnt: 3,
FailItemCnt: 1,
ProcessingItemCnt: 0,
TerminatedItemCnt: 2,
}, nil
}).
Times(2)

// Mock incomplete turns retrieval
mgr.exptResultService.(*svcMocks.MockExptResultService).
Expand Down Expand Up @@ -1682,11 +1698,20 @@ func TestExptMangerImpl_CompleteExpt(t *testing.T) {
return nil
})

// Mock stats update
// Mock stats update:断言写入的是重算后(第二次 CalculateStats)的计数,
// processing 归零、terminated=2,正是 bug#1 的核心不变量。
mgr.statsRepo.(*repoMocks.MockIExptStatsRepo).
EXPECT().
UpdateByExptID(ctx, int64(123), int64(789), gomock.Any()).
Return(nil)
DoAndReturn(func(_ context.Context, _, _ int64, s *entity.ExptStats) error {
if s.ProcessingItemCnt != 0 {
return fmt.Errorf("expected recalculated ProcessingItemCnt=0, got %d", s.ProcessingItemCnt)
}
if s.TerminatedItemCnt != 2 {
return fmt.Errorf("expected recalculated TerminatedItemCnt=2, got %d", s.TerminatedItemCnt)
}
return nil
})

// Mock experiment update
mgr.exptRepo.(*repoMocks.MockIExperimentRepo).
Expand Down
Loading