Skip to content

Commit f53cb19

Browse files
committed
[fix][evaluation] clean up owned runs rejected before dispatch
1 parent 71c4cb4 commit f53cb19

4 files changed

Lines changed: 562 additions & 11 deletions

File tree

‎backend/modules/evaluation/domain/service/expt_manage_execution_impl.go‎

Lines changed: 59 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -303,12 +303,48 @@ func withRetryYieldExt(ext map[string]string, enabled bool) map[string]string {
303303
return ext
304304
}
305305

306-
func (e *ExptMangerImpl) Run(ctx context.Context, exptID, runID, spaceID int64, itemRetryNum int, session *entity.Session, runMode entity.ExptRunMode, ext map[string]string) error {
307-
if err := NewQuotaService(e.quotaRepo, e.configer).AllowExptRun(ctx, exptID, spaceID, session); err != nil {
308-
return err
306+
func (e *ExptMangerImpl) prepareRun(ctx context.Context, exptID, runID, spaceID int64, session *entity.Session) (expt *entity.Experiment, err error) {
307+
defer func() {
308+
if err != nil {
309+
e.cleanupUnscheduledRun(ctx, exptID, runID, err)
310+
}
311+
}()
312+
313+
// Read dependencies before reserving quota so read failures cannot leak a quota slot.
314+
expt, err = e.GetDetail(ctx, exptID, spaceID, session)
315+
if err != nil {
316+
return nil, err
309317
}
318+
err = NewQuotaService(e.quotaRepo, e.configer).AllowExptRun(ctx, exptID, spaceID, session)
319+
return expt, err
320+
}
310321

311-
expt, err := e.GetDetail(ctx, exptID, spaceID, session)
322+
func (e *ExptMangerImpl) cleanupUnscheduledRun(ctx context.Context, exptID, runID int64, cause error) {
323+
if runID <= 0 {
324+
logs.CtxError(ctx, "[ExptEval][RunLock] invalid cleanup run [expt_id=%v run_id=%v]", exptID, runID)
325+
return
326+
}
327+
328+
// Only the owned run before MQ publication may be unlocked here.
329+
// Attempt this bare DEL synchronously once; a delayed retry could delete a new run's lock.
330+
ctx = context.WithoutCancel(ctx)
331+
stateCtx, cancelState := context.WithTimeout(ctx, exptRunLogPersistTimeout)
332+
stateErr := e.runLogRepo.Update(stateCtx, exptID, runID, map[string]any{
333+
"status": int64(entity.ExptStatus_Failed), "status_message": []byte(cause.Error()),
334+
})
335+
cancelState()
336+
337+
unlockCtx, cancelUnlock := context.WithTimeout(ctx, 2*time.Second)
338+
released, unlockErr := e.mutex.UnlockForce(unlockCtx, e.makeExptMutexLockKey(exptID))
339+
cancelUnlock()
340+
if stateErr != nil || unlockErr != nil {
341+
logs.CtxError(ctx, "[ExptEval][RunLock] cleanup failed [expt_id=%v run_id=%v state_err=%v unlock_err=%v]", exptID, runID, stateErr, unlockErr)
342+
}
343+
logs.CtxInfo(ctx, "[ExptEval][RunLock] cleanup result [expt_id=%v run_id=%v released=%v]", exptID, runID, released)
344+
}
345+
346+
func (e *ExptMangerImpl) Run(ctx context.Context, exptID, runID, spaceID int64, itemRetryNum int, session *entity.Session, runMode entity.ExptRunMode, ext map[string]string) error {
347+
expt, err := e.prepareRun(ctx, exptID, runID, spaceID, session)
312348
if err != nil {
313349
return err
314350
}
@@ -414,11 +450,7 @@ func buildExptNotifyParam(expt *entity.Experiment, toStatus entity.ExptStatus) (
414450
}
415451

416452
func (e *ExptMangerImpl) RetryItems(ctx context.Context, exptID, runID, spaceID int64, itemRetryNum int, itemIDs []int64, session *entity.Session, ext map[string]string) error {
417-
if err := NewQuotaService(e.quotaRepo, e.configer).AllowExptRun(ctx, exptID, spaceID, session); err != nil {
418-
return err
419-
}
420-
421-
expt, err := e.GetDetail(ctx, exptID, spaceID, session)
453+
expt, err := e.prepareRun(ctx, exptID, runID, spaceID, session)
422454
if err != nil {
423455
return err
424456
}
@@ -1679,15 +1711,22 @@ func (e *ExptMangerImpl) unlockCompletingRun(ctx context.Context, exptID, exptRu
16791711
return err
16801712
}
16811713

1682-
func (e *ExptMangerImpl) LogRun(ctx context.Context, exptID, exptRunID int64, mode entity.ExptRunMode, spaceID int64, itemIDs []int64, session *entity.Session) error {
1714+
func (e *ExptMangerImpl) LogRun(ctx context.Context, exptID, exptRunID int64, mode entity.ExptRunMode, spaceID int64, itemIDs []int64, session *entity.Session) (err error) {
16831715
duration := time.Duration(e.configer.GetExptExecConf(ctx, spaceID).GetZombieIntervalSecond()) * time.Second
16841716
locked, err := e.mutex.LockBackoff(ctx, e.makeExptMutexLockKey(exptID), duration, time.Second)
16851717
if err != nil {
16861718
return err
16871719
}
16881720
if !locked {
1721+
logs.CtxInfo(ctx, "[ExptEval][RunLock] lock occupied [expt_id=%v run_id=%v]", exptID, exptRunID)
16891722
return errorx.NewByCode(errno.ExperimentRunningExistedCode)
16901723
}
1724+
ownedRunID := exptRunID
1725+
defer func(runID int64) {
1726+
if err != nil {
1727+
e.cleanupUnscheduledRun(ctx, exptID, runID, err)
1728+
}
1729+
}(ownedRunID)
16911730

16921731
defer e.mtr.EmitExptExecRun(spaceID, int64(mode))
16931732

@@ -1733,16 +1772,25 @@ func (e *ExptMangerImpl) LogRetryItemsRun(ctx context.Context, exptID int64, mod
17331772
if err != nil {
17341773
return 0, false, err
17351774
}
1775+
if locked {
1776+
ownedRunID := runID
1777+
defer func(runID int64) {
1778+
if err != nil {
1779+
e.cleanupUnscheduledRun(ctx, exptID, runID, err)
1780+
}
1781+
}(ownedRunID)
1782+
}
17361783

17371784
var rl *entity.ExptRunLog
17381785
retried = !locked
17391786

17401787
if retried {
17411788
runID, err = strconv.ParseInt(existedRunID, 10, 64)
17421789
if err != nil {
1743-
logs.CtxError(ctx, "parsing expt run lock value to runid failed, raw: %v", existedRunID)
1790+
logs.CtxInfo(ctx, "[ExptEval][RunLock] lock occupied [expt_id=%v holder=%s]", exptID, existedRunID)
17441791
return 0, false, errorx.NewByCode(errno.ExperimentRunningExistedCode)
17451792
}
1793+
logs.CtxDebug(ctx, "[ExptEval][RunLock] joining existing run [expt_id=%v run_id=%v]", exptID, runID)
17461794

17471795
completing, err := e.ExistCompletingRunLock(ctx, exptID, runID, spaceID)
17481796
if err != nil {

‎backend/modules/evaluation/domain/service/expt_manage_execution_impl_test.go‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -153,6 +153,8 @@ func TestExptMangerImpl_Run(t *testing.T) {
153153
EXPECT().
154154
CreateOrUpdate(ctx, int64(789), gomock.Any(), session).
155155
Return(errors.New("quota exceeded"))
156+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Update(gomock.Any(), int64(123), int64(456), gomock.Any()).Return(nil)
157+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().UnlockForce(gomock.Any(), mgr.makeExptMutexLockKey(123)).Return(true, nil)
156158
},
157159
wantErr: true,
158160
},
@@ -622,6 +624,8 @@ func TestExptMangerImpl_LogRun(t *testing.T) {
622624
EXPECT().
623625
Create(ctx, gomock.Any()).
624626
Return(errors.New("create failed"))
627+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Update(gomock.Any(), int64(123), int64(456), gomock.Any()).Return(nil)
628+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().UnlockForce(gomock.Any(), mgr.makeExptMutexLockKey(123)).Return(true, nil)
625629
},
626630
wantErr: true,
627631
},
@@ -802,6 +806,8 @@ func TestExptMangerImpl_LogRetryItemsRun(t *testing.T) {
802806
Return(true, "1006", nil)
803807
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).
804808
EXPECT().Save(ctx, gomock.Any()).Return(errors.New("save failed"))
809+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Update(gomock.Any(), exptID, int64(1006), gomock.Any()).Return(nil)
810+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().UnlockForce(gomock.Any(), mgr.makeExptMutexLockKey(exptID)).Return(true, nil)
805811
},
806812
wantErr: true,
807813
},
@@ -898,6 +904,8 @@ func TestExptMangerImpl_RetryItems(t *testing.T) {
898904
setup: func() {
899905
mgr.quotaRepo.(*repoMocks.MockQuotaRepo).
900906
EXPECT().CreateOrUpdate(ctx, spaceID, gomock.Any(), session).Return(errors.New("quota exceeded"))
907+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Update(gomock.Any(), exptID, runID, gomock.Any()).Return(nil)
908+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().UnlockForce(gomock.Any(), mgr.makeExptMutexLockKey(exptID)).Return(true, nil)
901909
mgr.configer.(*componentMocks.MockIConfiger).
902910
EXPECT().GetExptExecConf(ctx, spaceID).AnyTimes().
903911
Return(&entity.ExptExecConf{SpaceExptConcurLimit: 10})
Lines changed: 232 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,232 @@
1+
// Copyright (c) 2025 coze-dev Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
package service
5+
6+
import (
7+
"context"
8+
"errors"
9+
"strconv"
10+
"testing"
11+
"time"
12+
13+
"github.com/stretchr/testify/require"
14+
"go.uber.org/mock/gomock"
15+
16+
idgenMocks "github.com/coze-dev/coze-loop/backend/infra/idgen/mocks"
17+
lockMocks "github.com/coze-dev/coze-loop/backend/infra/lock/mocks"
18+
lwtMocks "github.com/coze-dev/coze-loop/backend/infra/platestwrite/mocks"
19+
metricsMocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/component/metrics/mocks"
20+
componentMocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/component/mocks"
21+
"github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/entity"
22+
eventsMocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/events/mocks"
23+
repoMocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/repo/mocks"
24+
svcMocks "github.com/coze-dev/coze-loop/backend/modules/evaluation/domain/service/mocks"
25+
"github.com/coze-dev/coze-loop/backend/modules/evaluation/pkg/errno"
26+
"github.com/coze-dev/coze-loop/backend/pkg/errorx"
27+
)
28+
29+
func TestRunAdmissionCleanup(t *testing.T) {
30+
for _, tc := range []struct {
31+
name string
32+
mode entity.ExptRunMode
33+
publishFail bool
34+
}{
35+
{"submit", entity.EvaluationModeSubmit, false},
36+
{"retry_all", entity.EvaluationModeRetryAll, false},
37+
{"retry_failure", entity.EvaluationModeFailRetry, false},
38+
{"retry_items", entity.EvaluationModeRetryItems, false},
39+
{"publish_failure_preserves_run", entity.EvaluationModeFailRetry, true},
40+
{"publish_failure_preserves_retry_items_run", entity.EvaluationModeRetryItems, true},
41+
} {
42+
t.Run(tc.name, func(t *testing.T) {
43+
ctrl := gomock.NewController(t)
44+
mgr := newTestExptManager(ctrl)
45+
ctx := context.Background()
46+
session := &entity.Session{UserID: "u"}
47+
const exptID, runID, spaceID = int64(1), int64(2), int64(3)
48+
limit, logCalls := 2, 2
49+
publishErr := errors.New("publish timed out")
50+
if tc.publishFail {
51+
limit, logCalls = 3, 1
52+
mgr.publisher.(*eventsMocks.MockExptEventPublisher).EXPECT().PublishExptScheduleEvent(ctx, gomock.Any(), gomock.Any()).Return(publishErr)
53+
}
54+
configer := componentMocks.NewMockIConfiger(ctrl)
55+
configer.EXPECT().GetExptExecConf(ctx, spaceID).Return(&entity.ExptExecConf{
56+
SpaceExptConcurLimit: limit, ZombieIntervalSecond: 72 * 60 * 60,
57+
}).AnyTimes()
58+
configer.EXPECT().GetRetryYieldEnabled(ctx, spaceID).Return(false).AnyTimes()
59+
mgr.configer = configer
60+
quota := &entity.QuotaSpaceExpt{ExptID2RunTime: map[int64]int64{10: time.Now().Unix(), 11: time.Now().Unix()}}
61+
mgr.quotaRepo.(*repoMocks.MockQuotaRepo).EXPECT().CreateOrUpdate(ctx, spaceID, gomock.Any(), session).
62+
DoAndReturn(func(_ context.Context, _ int64, update func(*entity.QuotaSpaceExpt) (*entity.QuotaSpaceExpt, bool, error), _ *entity.Session) error {
63+
_, changed, err := update(quota)
64+
require.Equal(t, tc.publishFail, changed)
65+
return err
66+
})
67+
mgr.lwt.(*lwtMocks.MockILatestWriteTracker).EXPECT().CheckWriteFlagByID(ctx, gomock.Any(), exptID).Return(false).AnyTimes()
68+
mgr.exptRepo.(*repoMocks.MockIExperimentRepo).EXPECT().MGetByID(ctx, []int64{exptID}, spaceID).
69+
Return([]*entity.Experiment{{ID: exptID, SpaceID: spaceID, Status: entity.ExptStatus_Failed}}, nil).AnyTimes()
70+
mgr.evaluationSetService.(*svcMocks.MockIEvaluationSetService).EXPECT().
71+
GetEvaluationSet(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any(), gomock.Nil()).Return(&entity.EvaluationSet{}, nil).AnyTimes()
72+
mgr.exptResultService.(*svcMocks.MockExptResultService).EXPECT().
73+
MGetStats(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil, nil).AnyTimes()
74+
mgr.exptAggrResultService.(*svcMocks.MockExptAggrResultService).EXPECT().
75+
BatchGetExptAggrResultByExperimentIDs(gomock.Any(), gomock.Any(), gomock.Any()).Return(nil, nil).AnyTimes()
76+
77+
var runLog *entity.ExptRunLog
78+
locked := false
79+
lockKey := mgr.makeExptMutexLockKey(exptID)
80+
mgr.mtr.(*metricsMocks.MockExptMetric).EXPECT().EmitExptExecRun(spaceID, int64(tc.mode)).Times(logCalls)
81+
var latestRunID int64
82+
mgr.exptRepo.(*repoMocks.MockIExperimentRepo).EXPECT().Update(ctx, gomock.Any()).
83+
DoAndReturn(func(_ context.Context, expt *entity.Experiment) error { latestRunID = expt.LatestRunID; return nil }).Times(logCalls)
84+
if tc.mode == entity.EvaluationModeRetryItems {
85+
mgr.idgenerator.(*idgenMocks.MockIIDGenerator).EXPECT().GenID(ctx).Return(runID, nil).Times(logCalls)
86+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().BackoffLockWithValue(ctx, lockKey, strconv.FormatInt(runID, 10), gomock.Any(), gomock.Any()).
87+
DoAndReturn(func(context.Context, string, string, time.Duration, time.Duration) (bool, string, error) {
88+
if locked {
89+
return false, strconv.FormatInt(runID, 10), nil
90+
}
91+
locked = true
92+
return true, "", nil
93+
}).Times(logCalls)
94+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Save(gomock.Any(), gomock.Any()).
95+
DoAndReturn(func(_ context.Context, got *entity.ExptRunLog) error { runLog = got; return nil }).Times(logCalls)
96+
} else {
97+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().LockBackoff(ctx, lockKey, gomock.Any(), gomock.Any()).
98+
DoAndReturn(func(context.Context, string, time.Duration, time.Duration) (bool, error) {
99+
if locked {
100+
return false, nil
101+
}
102+
locked = true
103+
return true, nil
104+
}).Times(logCalls)
105+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Create(ctx, gomock.Any()).
106+
DoAndReturn(func(_ context.Context, got *entity.ExptRunLog) error { runLog = got; return nil }).Times(logCalls)
107+
}
108+
if !tc.publishFail {
109+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Update(gomock.Any(), exptID, runID, gomock.Any()).
110+
DoAndReturn(func(_ context.Context, _, _ int64, fields map[string]any) error {
111+
runLog.Status = fields["status"].(int64)
112+
runLog.StatusMessage = fields["status_message"].([]byte)
113+
return nil
114+
})
115+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().UnlockForce(gomock.Any(), lockKey).
116+
DoAndReturn(func(context.Context, string) (bool, error) {
117+
require.Equal(t, int64(entity.ExptStatus_Failed), runLog.Status)
118+
locked = false
119+
return true, nil
120+
})
121+
}
122+
123+
var err error
124+
if tc.mode == entity.EvaluationModeRetryItems {
125+
_, retried, logErr := mgr.LogRetryItemsRun(ctx, exptID, tc.mode, spaceID, []int64{4}, session)
126+
require.NoError(t, logErr)
127+
require.False(t, retried)
128+
err = mgr.RetryItems(ctx, exptID, runID, spaceID, 0, []int64{4}, session, nil)
129+
} else {
130+
require.NoError(t, mgr.LogRun(ctx, exptID, runID, tc.mode, spaceID, nil, session))
131+
err = mgr.Run(ctx, exptID, runID, spaceID, 0, session, tc.mode, nil)
132+
}
133+
if tc.publishFail {
134+
require.ErrorIs(t, err, publishErr)
135+
require.True(t, locked)
136+
require.Equal(t, int64(entity.ExptStatus_Pending), runLog.Status)
137+
return
138+
}
139+
status, ok := errorx.FromStatusError(err)
140+
require.True(t, ok)
141+
require.Equal(t, int32(errno.ExperimentRunningCountLimitCode), status.Code())
142+
require.Len(t, quota.ExptID2RunTime, 2)
143+
require.NotContains(t, quota.ExptID2RunTime, exptID)
144+
require.False(t, locked)
145+
require.Equal(t, int64(entity.ExptStatus_Failed), runLog.Status)
146+
require.Contains(t, string(runLog.StatusMessage), "max limit: 2")
147+
require.Equal(t, runID, latestRunID, "admission cleanup must leave the latest run pointing at the failed attempt")
148+
if tc.mode == entity.EvaluationModeRetryItems {
149+
_, retried, logErr := mgr.LogRetryItemsRun(ctx, exptID, tc.mode, spaceID, []int64{4}, session)
150+
require.NoError(t, logErr)
151+
require.False(t, retried)
152+
} else {
153+
require.NoError(t, mgr.LogRun(ctx, exptID, runID+1, tc.mode, spaceID, nil, session))
154+
}
155+
require.True(t, locked)
156+
})
157+
}
158+
}
159+
160+
func TestRunPreparationReadFailure(t *testing.T) {
161+
for _, tc := range []struct {
162+
name string
163+
updateFails, updateTimesOut, updateCommittedError bool
164+
lockPresent, released bool
165+
unlockErr error
166+
}{
167+
{name: "success", lockPresent: true, released: true},
168+
{name: "update_failure_still_unlocks", updateFails: true, lockPresent: true, released: true},
169+
{name: "update_committed_then_error_still_unlocks", updateCommittedError: true, lockPresent: true, released: true},
170+
{name: "update_timeout_leaves_full_unlock_budget", updateTimesOut: true, lockPresent: true, released: true},
171+
{name: "lock_already_absent"},
172+
{name: "unlock_error_without_delete_preserves_original", lockPresent: true, unlockErr: errors.New("redis unavailable")},
173+
{name: "update_and_unlock_fail_preserves_original", updateFails: true, lockPresent: true, unlockErr: errors.New("redis unavailable")},
174+
} {
175+
t.Run(tc.name, func(t *testing.T) {
176+
ctrl := gomock.NewController(t)
177+
mgr := newTestExptManager(ctrl)
178+
ctx, cancel := context.WithCancel(context.Background())
179+
defer cancel()
180+
const exptID, runID, spaceID = int64(1), int64(2), int64(3)
181+
mgr.lwt.(*lwtMocks.MockILatestWriteTracker).EXPECT().CheckWriteFlagByID(ctx, gomock.Any(), exptID).Return(false)
182+
mgr.exptRepo.(*repoMocks.MockIExperimentRepo).EXPECT().MGetByID(ctx, []int64{exptID}, spaceID).
183+
DoAndReturn(func(context.Context, []int64, int64) ([]*entity.Experiment, error) {
184+
cancel()
185+
return nil, context.Canceled
186+
})
187+
runLog := &entity.ExptRunLog{ExptID: exptID, ExptRunID: runID, Status: int64(entity.ExptStatus_Pending)}
188+
locked := tc.lockPresent
189+
var updateCtx context.Context
190+
mgr.runLogRepo.(*repoMocks.MockIExptRunLogRepo).EXPECT().Update(gomock.Any(), exptID, runID, gomock.Any()).
191+
DoAndReturn(func(cleanupCtx context.Context, _, _ int64, fields map[string]any) error {
192+
updateCtx = cleanupCtx
193+
require.NoError(t, cleanupCtx.Err())
194+
deadline, ok := cleanupCtx.Deadline()
195+
require.True(t, ok)
196+
require.InDelta(t, 5, time.Until(deadline).Seconds(), 1)
197+
if tc.updateTimesOut {
198+
<-cleanupCtx.Done()
199+
return cleanupCtx.Err()
200+
}
201+
if tc.updateFails {
202+
return errors.New("run log unavailable")
203+
}
204+
runLog.Status = fields["status"].(int64)
205+
runLog.StatusMessage = fields["status_message"].([]byte)
206+
if tc.updateCommittedError {
207+
return errors.New("run log response lost")
208+
}
209+
return nil
210+
})
211+
mgr.mutex.(*lockMocks.MockILocker).EXPECT().UnlockForce(gomock.Any(), mgr.makeExptMutexLockKey(exptID)).
212+
DoAndReturn(func(unlockCtx context.Context, _ string) (bool, error) {
213+
require.NotSame(t, updateCtx, unlockCtx)
214+
require.NoError(t, unlockCtx.Err())
215+
deadline, ok := unlockCtx.Deadline()
216+
require.True(t, ok)
217+
require.InDelta(t, 2, time.Until(deadline).Seconds(), 0.5)
218+
if tc.unlockErr == nil {
219+
locked = false
220+
}
221+
return tc.released, tc.unlockErr
222+
})
223+
err := mgr.Run(ctx, exptID, runID, spaceID, 0, &entity.Session{UserID: "u"}, entity.EvaluationModeFailRetry, nil)
224+
require.ErrorIs(t, err, context.Canceled)
225+
require.Equal(t, tc.lockPresent && tc.unlockErr != nil, locked)
226+
if !tc.updateFails && !tc.updateTimesOut {
227+
require.Equal(t, int64(entity.ExptStatus_Failed), runLog.Status)
228+
require.Equal(t, context.Canceled.Error(), string(runLog.StatusMessage))
229+
}
230+
})
231+
}
232+
}

0 commit comments

Comments
 (0)