Skip to content

Commit b1c11b2

Browse files
fix daemon polling after task framework reenable (#26077)
fix daemon polling after task framework reenable Approved by: @XuPeng-SH
1 parent 1ac64a8 commit b1c11b2

2 files changed

Lines changed: 67 additions & 8 deletions

File tree

pkg/taskservice/daemon_task.go

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -547,18 +547,27 @@ func (r *taskRunner) startDaemonTaskWorker() error {
547547
func (r *taskRunner) poll(ctx context.Context) {
548548
timer := time.NewTimer(r.options.fetchInterval)
549549
defer timer.Stop()
550+
r.pollWithTimer(ctx, timer.C, func() {
551+
timer.Reset(r.options.fetchInterval)
552+
})
553+
}
554+
555+
func (r *taskRunner) pollWithTimer(
556+
ctx context.Context,
557+
timerC <-chan time.Time,
558+
resetTimer func(),
559+
) {
550560
for {
551561
select {
552562
case <-ctx.Done():
553563
r.logger.Info("daemon task poll worker stopped")
554564
return
555565

556-
case <-timer.C:
557-
if taskFrameworkDisabled() {
558-
continue
566+
case <-timerC:
567+
if !taskFrameworkDisabled() {
568+
r.dispatchTaskHandle(ctx)
559569
}
560-
r.dispatchTaskHandle(ctx)
561-
timer.Reset(r.options.fetchInterval)
570+
resetTimer()
562571
}
563572
}
564573
}

pkg/taskservice/daemon_task_test.go

Lines changed: 53 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -94,12 +94,14 @@ func (r *mockErrActiveRoutine) Restart() error { return r.restartErr }
9494

9595
type serviceWithDaemonHook struct {
9696
TaskService
97-
mu sync.RWMutex
98-
queryErr error
99-
updateErr error
97+
mu sync.RWMutex
98+
queryErr error
99+
updateErr error
100+
queryCalls atomic.Int64
100101
}
101102

102103
func (s *serviceWithDaemonHook) QueryDaemonTask(ctx context.Context, conds ...Condition) ([]task.DaemonTask, error) {
104+
s.queryCalls.Add(1)
103105
s.mu.RLock()
104106
queryErr := s.queryErr
105107
s.mu.RUnlock()
@@ -131,6 +133,54 @@ func (s *serviceWithDaemonHook) setUpdateErr(err error) {
131133
s.updateErr = err
132134
}
133135

136+
func TestDaemonTaskPollResumesAfterTaskFrameworkReenabled(t *testing.T) {
137+
wasDisabled := taskFrameworkDisabled()
138+
DebugCtlTaskFramework(true)
139+
t.Cleanup(func() {
140+
DebugCtlTaskFramework(wasDisabled)
141+
})
142+
143+
r, _ := newDaemonHandleTestRunner(t)
144+
hook := &serviceWithDaemonHook{TaskService: r.service}
145+
r.service = hook
146+
147+
ctx, cancel := context.WithCancel(context.Background())
148+
done := make(chan struct{})
149+
timerC := make(chan time.Time)
150+
resetC := make(chan struct{}, 2)
151+
go func() {
152+
defer close(done)
153+
r.pollWithTimer(ctx, timerC, func() {
154+
resetC <- struct{}{}
155+
})
156+
}()
157+
defer func() {
158+
cancel()
159+
select {
160+
case <-done:
161+
case <-time.After(time.Second):
162+
t.Error("daemon poll did not stop after cancellation")
163+
}
164+
}()
165+
166+
timerC <- time.Now()
167+
select {
168+
case <-resetC:
169+
case <-time.After(time.Second):
170+
t.Fatal("disabled daemon poll did not reset its timer")
171+
}
172+
require.Zero(t, hook.queryCalls.Load())
173+
174+
DebugCtlTaskFramework(false)
175+
timerC <- time.Now()
176+
select {
177+
case <-resetC:
178+
case <-time.After(time.Second):
179+
t.Fatal("enabled daemon poll did not reset its timer")
180+
}
181+
require.Positive(t, hook.queryCalls.Load())
182+
}
183+
134184
func daemonTaskMetadata() task.TaskMetadata {
135185
return task.TaskMetadata{
136186
ID: "-",

0 commit comments

Comments
 (0)