// ------------------------------------------------------------------------ // Project atila // Active Thing (activething.com) git.activething.com/go // // File name rkey.go // Created by DEV // Modified 05/02/2024 // // Copyright 2024 activething.com // ------------------------------------------------------------------------ // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. // ------------------------------------------------------------------------ package times import ( "activething.com/go/gogo/core/opts" "container/heap" "context" "sync/atomic" "time" ) const ( timerQueueLen = 4096 timerQueueTime = time.Hour * 86400 ) type ( TimerQueue struct { heap TimerTaskHeap state uint32 pushCh chan *TimerTask ctx context.Context ctxCancelFn context.CancelFunc nextTask *time.Time timer *time.Timer } ) func NewTimerQueue(options ... opts.OptionFnc[TimerQueue]) *TimerQueue { tq := &TimerQueue{ heap : TimerTaskHeap{}, state: uint32(TimerStateCreated), } heap.Init(&tq.heap) opts.Apply(tq, options ...) if tq.pushCh == nil { tq.pushCh = make(chan *TimerTask, timerQueueLen) } return tq } func (q *TimerQueue) State() TimerState { return TimerState(atomic.LoadUint32(&q.state)) } func (q *TimerQueue) Start(ctx context.Context) error { if !atomic.CompareAndSwapUint32(&q.state,uint32(TimerStateCreated),uint32(TimerStateStarted)) { // error already started return G.ErrAlreadyStarted } if q.timer == nil { q.timer = time.NewTimer(timerQueueTime) } else { q.timer.Reset(timerQueueTime) } q.ctx, q.ctxCancelFn = context.WithCancel(ctx) go func(qu *TimerQueue, cx context.Context) { ts := make([]*TimerTask, 0, 10) EXIT: for { select { //-- case <-cx.Done(): atomic.CompareAndSwapUint32(&qu.state,uint32(TimerStateStarted),uint32(TimerStateCanceled)) break EXIT //-- case nw := <-qu.timer.C: ts = ts[:0] for ix, ln := 0, qu.heap.Len(); ix < ln; ix++ { pk := qu.heap[0] if pk.RunAt.Before(nw) { pp := heap.Pop(&qu.heap).(*TimerTask) ts = append(ts, pp) continue } break } if len(ts) > 0 { go func(tm time.Time, tk []*TimerTask) { for _, tt := range tk { tt.TaskFn(tm, ctx) } }(nw, ts) } if qu.heap.Len() > 0 { nr := qu.heap[0].RunAt qu.nextTask = &nr qu.timer.Reset(nr.Sub(time.Now())) } else { qu.timer.Stop() qu.nextTask = nil } //-- case tt := <-q.pushCh: heap.Push(&q.heap, tt) if qu.nextTask == nil || qu.nextTask.After(tt.RunAt) { if qu.nextTask != nil && !qu.timer.Stop() { <-qu.timer.C } qu.timer.Reset(tt.RunAt.Sub(time.Now())) qu.nextTask = &tt.RunAt } } } if !q.timer.Stop() { <-qu.timer.C } qu.heap.Reset() }(q, q.ctx) return nil } func (q *TimerQueue) Stop() error { if atomic.CompareAndSwapUint32(&q.state,uint32(TimerStateStarted),uint32(TimerStateCanceled)) { q.ctxCancelFn() return nil } return G.ErrNotStarted } func (q *TimerQueue) Push(runAt time.Time, task TimeTaskFnc) error { switch { case runAt.Before(time.Now()): return G.ErrExpired case task == nil: return G.ErrTimerTaskInvalid case TimerState(atomic.LoadUint32(&q.state)) != TimerStateStarted: return G.ErrNotStarted default: select { case q.pushCh <- NewTimerTask(runAt, task): return nil case <-q.ctx.Done(): return context.Canceled } } } /* :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. SIDDHARTA - HERMANN HESSE - ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: */