From 831b5b267373f6fbd3548849a3925c4e70806de2 Mon Sep 17 00:00:00 2001 From: xtaci Date: Sun, 22 Mar 2020 22:01:17 +0800 Subject: [PATCH] properly drain timer.C --- timedsched.go | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/timedsched.go b/timedsched.go index e20fc52..2db7c20 100644 --- a/timedsched.go +++ b/timedsched.go @@ -62,6 +62,7 @@ func NewTimedSched(parallel int) *TimedSched { func (ts *TimedSched) sched() { var tasks timedFuncHeap timer := time.NewTimer(0) + drained := false for { select { case task := <-ts.chTask: @@ -71,15 +72,22 @@ func (ts *TimedSched) sched() { task.execute() } else { heap.Push(&tasks, task) - // reset timer to trigger based on the top element + // properly reset timer to trigger based on the top element + stopped := timer.Stop() + if !stopped && !drained { + <-timer.C + } timer.Reset(tasks[0].ts.Sub(now)) + drained = false } case now := <-timer.C: + drained = true for tasks.Len() > 0 { if now.After(tasks[0].ts) { heap.Pop(&tasks).(timedFunc).execute() } else { timer.Reset(tasks[0].ts.Sub(now)) + drained = false break } }