From 17f27965115b1bda08d02e5f884466387b8e368b Mon Sep 17 00:00:00 2001 From: DEV Date: Wed, 27 Mar 2024 22:31:18 +0100 Subject: [PATCH] terminando de crear el formateador de texto de errores --- core/errs/fmt_text.go | 122 ------------------ core/errs/formatter.go | 111 ----------------- core/syncs/pool.go | 6 +- core/syncs/pool_ops.go | 10 +- core/syncs/tasks/completion.go | 12 +- core/syncs/tasks/completion_ops.go | 12 +- core/syncs/tasks/gogo.go | 2 +- core/syncs/tasks/worker_pool.go | 8 +- core/timers/active.go | 71 +++++++++++ core/timers/chrono.go | 108 ++++++++++++++++ core/timers/chronometer.go | 59 +++++++++ core/timers/expiration.go | 182 +++++++++++++++++++++++++++ core/timers/expiration_source.go | 62 ++++++++++ core/timers/expirer.go | 57 +++++++++ core/timers/go.mod | 4 + core/timers/pulsar.go | 174 ++++++++++++++++++++++++++ core/timers/time_task.go | 48 ++++++++ core/timers/timer_queue.go | 190 +++++++++++++++++++++++++++++ core/timers/timer_state.go | 81 ++++++++++++ core/timers/timer_task.go | 60 +++++++++ core/timers/timer_task_heap.go | 89 ++++++++++++++ go.work | 1 + 22 files changed, 1211 insertions(+), 258 deletions(-) create mode 100644 core/timers/active.go create mode 100644 core/timers/chrono.go create mode 100644 core/timers/chronometer.go create mode 100644 core/timers/expiration.go create mode 100644 core/timers/expiration_source.go create mode 100644 core/timers/expirer.go create mode 100644 core/timers/go.mod create mode 100644 core/timers/pulsar.go create mode 100644 core/timers/time_task.go create mode 100644 core/timers/timer_queue.go create mode 100644 core/timers/timer_state.go create mode 100644 core/timers/timer_task.go create mode 100644 core/timers/timer_task_heap.go diff --git a/core/errs/fmt_text.go b/core/errs/fmt_text.go index 9cc4eb0..c79e388 100644 --- a/core/errs/fmt_text.go +++ b/core/errs/fmt_text.go @@ -215,125 +215,3 @@ func formatStackTrace(s *FmtTextState, t StackTrace) { } } - -/* - - -func (f *FmtText) SetErrorFormatter (errType reflect.Type, formatFn func (*bytes.Buffer,error,int)) { - if formatFn == nil { - formatFn = textFormatterDefault - } - f.providers.Store(errType,formatFn) -} - -func (f *FmtText) SetFrameFormatter (formatFn func (*bytes.Buffer, StackTrace, int)) { - if formatFn == nil { - formatFn = textFormatterTrace - } - f.providers.Store(G.stackTraceRType,formatFn) -} - - - -func (f *FmtText) Format (err error) []byte { - if err == nil { - return nil - } - d,e := (&FmtTextState{}).format(err) - if e != nil { - return nil - } - return d -} - - -func (f *FmtText) fmtError (b *bytes.Buffer, e error, l int) { - tp := reflect.TypeOf(e) - if l == 0 { - writeFieldString(b, "error", e.Error(), l) - } else { - writeFieldString(b, "prev error", e.Error(), l) - } - writeFieldString(b,"error-type",tp.String(),l) - - ls := UnwrapAll(e) - for _,e := range ls { - - switch tv := e.(type) { - case *ErrCode : f.fmtErrCode (b,*tv, l) - case ErrCode : f.fmtErrCode (b, tv, l) - case *ErrInfo : f.fmtErrInfo (b,*tv, l) - case ErrInfo : f.fmtErrInfo (b, tv, l) - case *ErrCause : f.fmtErrCause(b,*tv, l) - case ErrCause : f.fmtErrCause(b, tv, l) - case *ErrTrace : f.fmtErrTrace(b,*tv, l) - case ErrTrace : f.fmtErrTrace(b, tv, l) - - default: - fm,_ := f.providers.LoadOrStore(tp, textFormatterDefault) - fm.(func (b *bytes.Buffer, e error, l int))(b,e,l) - } - } - - if c := UnwrapCause(e); c != nil { - b.WriteString("\n") - //writeString(b,"previous error",l) - f.fmtError(b,c.Cause(),l+1) - } - - if l == 0 { - if st,is := e.(StackTracer); is { - f.fmtTrace (b,st.StackTrace(),l) - } - } - -} - - -func (f *FmtText) getFormatter (err error) TextFormatterFnc { - -} - - -func textFormatterDefault(b *bytes.Buffer, e error, l int) { - fmt.Printf("Por defect0*******\n") -} - -func textFormatterTrace (b *bytes.Buffer, s StackTrace, l int) { - -} - - -func (f *FmtText) fmtErrCode (b *bytes.Buffer, e ErrCode, l int) {} - -func (f *FmtText) fmtErrInfo (b *bytes.Buffer, e ErrInfo, l int) { - writeFieldString(b,"message" ,e.Message(),l) - writeFieldString(b,"info" ,e.Info (),l) - if l == 0 { - writeFieldString(b, "date-time", e.Date().Format(time.RFC822Z), l) - e.MapMetas(func(n string, v any) bool { - writeFieldString(b, n,fmt.Sprintf("%v",v), l) - return true - }) - } -} - -func (f *FmtText) fmtErrCause(b *bytes.Buffer, e ErrCause,l int) {} - -func (f *FmtText) fmtErrTrace(b *bytes.Buffer, e ErrTrace,l int) { - -} - -func (f *FmtText) fmtTrace (b *bytes.Buffer, s StackTrace, l int) { - fm,_ := f.providers.LoadOrStore(G.stackTraceRType,textFormatterTrace) - fm.(func (b *bytes.Buffer, s StackTrace, l int))(b,s,l) -} - - - - - - - - -*/ \ No newline at end of file diff --git a/core/errs/formatter.go b/core/errs/formatter.go index 421e5d2..397dd37 100644 --- a/core/errs/formatter.go +++ b/core/errs/formatter.go @@ -23,14 +23,6 @@ package errs -import ( - "bytes" - "fmt" - "reflect" - "strings" - "time" -) - type ( Formatter interface { @@ -47,106 +39,3 @@ func (f FormatFnc) Format (err error) []byte { return f(err) } - - - - - - -func ErrorFormat(err error) []byte { - bf := &bytes.Buffer{} - formatError(bf,err,0) - return bf.Bytes() -} - - -func formatError(bf *bytes.Buffer, err error, level int) { - switch tp := err.(type) { - case ErrInfo : formatInfoErr(bf, tp,level) - case *ErrInfo : formatInfoErr(bf,*tp,level) - case ErrCode : formatCodeErr (bf, tp, level) - case *ErrCode : formatCodeErr (bf,*tp, level) - case ErrCause : formatCauseErr(bf, tp, level) - case *ErrCause: formatCauseErr(bf,*tp, level) - case *ErrTrace: formatTraceErr(bf,*tp, level) - case ErrTrace: formatTraceErr(bf, tp, level) - default: - writeString(bf,reflect.TypeOf(err).String(),level) - writeFieldString(bf,"error",err.Error(),level) - } - if ce := Unwrap(err); ce != nil { - //formatCauseErr(bf, ce, level) - } -} - - - -func formatInfoErr(bf *bytes.Buffer, err ErrInfo, level int) { - er := err.Unwrap() - formatError(bf,er,level) - - - writeFieldString(bf,"message" ,err.Message(),level) - writeFieldString(bf,"info" ,err.Info (),level) - if level == 0 { - writeFieldString(bf, "date-time", err.Date().Format(time.RFC822Z), level) - err.MapMetas(func(n string, v any) bool { - writeFieldString(bf, n,fmt.Sprintf("%v",v), level) - return true - }) - } - - - if ce := UnwrapCause(err); ce != nil { - bf.WriteString("\n") - writeString(bf,fmt.Sprintf("**** UnwrapCause **** (%d)",level+1) ,level+1) - formatError(bf,ce.Cause(),level+1) - } - - - //if st := err.StackTrace(); st != nil && level == 0 { - if st := UnwrapStackTrace(err); st != nil && level == 0 { - bf.WriteString("\nStack Trace\n") - formatTrace(bf,st.StackTrace(),level) - } -} - - -func formatCodeErr (bf *bytes.Buffer, err ErrCode , level int) { - writeFieldString(bf,"error-code",err.Code(),level) -} - -func formatCauseErr(bf *bytes.Buffer, err error , level int) { - formatError (bf,Unwrap(err),level) - //bf.WriteString("\n") - //writeString(bf,fmt.Sprintf("**** UnwrapCause **** (%d)",level+1) ,level+1) - //formatError (bf,err.(Causer).UnwrapCause(),level+1) - -} - -func formatTraceErr(bf *bytes.Buffer, err ErrTrace, level int) { - writeFieldString(bf,"error-stack-trace",err.Error(),level) - if level == 0 { - formatTrace(bf, err.StackTrace(), level) - } -} - -func formatTrace (bf *bytes.Buffer, trace StackTrace, level int) { - for _,f := range trace.Frames() { - sl,_ := f.SourceLine() - writeFieldString(bf,f.FuncName,sl,level) - } -} - - -func writeString (bf *bytes.Buffer, value string, level int) { - bf.WriteString(fmt.Sprintf("%s%s\n",strings.Repeat(" ",level),value)) -} - -func writeFieldString (bf *bytes.Buffer, name, value string, level int) { - if len(value) > 0 { - bf.WriteString(fmt.Sprintf("%s%-10s: %s\n",strings.Repeat(" ",level),name,value)) - } -} - - diff --git a/core/syncs/pool.go b/core/syncs/pool.go index c3ad97a..67709cb 100644 --- a/core/syncs/pool.go +++ b/core/syncs/pool.go @@ -24,7 +24,7 @@ package syncs import ( - "activething.com/go/gogo/core/opts" + "opts" "sync" "sync/atomic" ) @@ -43,8 +43,8 @@ type ( } ) -func NewPool[T any](ops ...opts.OptionFnc[Pool[T]]) *Pool[T] { - return opts.Apply(&Pool[T]{}, ops...) } +func NewPool[T any](ops ...opts.GOptionFnc[Pool[T]]) *Pool[T] { + return opts.GApply(&Pool[T]{}, ops...) } func (p *Pool[T]) Len() int32 { diff --git a/core/syncs/pool_ops.go b/core/syncs/pool_ops.go index 7902a28..98b7d54 100644 --- a/core/syncs/pool_ops.go +++ b/core/syncs/pool_ops.go @@ -24,25 +24,25 @@ package syncs import ( - "activething.com/go/gogo/core/opts" + "opts" "sync/atomic" ) -func PoolWithNew[T any](nfn func() *T) opts.OptionFnc[Pool[T]] { +func PoolWithNew[T any](nfn func() *T) opts.GOptionFnc[Pool[T]] { return func(p *Pool[T]) { p.createFn = nfn } } -func PoolWithCap[T any](cap int32) opts.OptionFnc[Pool[T]] { +func PoolWithCap[T any](cap int32) opts.GOptionFnc[Pool[T]] { return func(p *Pool[T]) { atomic.StoreInt32(&p.cap, cap) } } -func PoolWithLen[T any](len int32) opts.OptionFnc[Pool[T]] { +func PoolWithLen[T any](len int32) opts.GOptionFnc[Pool[T]] { return func(p *Pool[T]) { if p.createFn == nil { return @@ -58,7 +58,7 @@ func PoolWithLen[T any](len int32) opts.OptionFnc[Pool[T]] { } -func PoolWithSet[T any](els ...*T) opts.OptionFnc[Pool[T]] { +func PoolWithSet[T any](els ...*T) opts.GOptionFnc[Pool[T]] { return func(p *Pool[T]) { for _, e := range els { if e != nil { diff --git a/core/syncs/tasks/completion.go b/core/syncs/tasks/completion.go index 3ac3a17..42a1639 100644 --- a/core/syncs/tasks/completion.go +++ b/core/syncs/tasks/completion.go @@ -24,11 +24,11 @@ package tasks import ( - "activething.com/go/gogo/core/opts" - "activething.com/go/gogo/core/timers" "context" + "opts" "sync" "sync/atomic" + "timers" ) type ( @@ -61,8 +61,8 @@ type ( ) -func NewCompletion (options ...opts.OptionErrFnc[Completion]) (*Completion,error) { - c,e := opts.ApplyAndCheck[Completion](&Completion { +func NewCompletion (options ...opts.GOptionErrFnc[Completion]) (*Completion,error) { + c,e := opts.GApplyAndCheck[Completion](&Completion { mode : CompletionModeAll, cid : CompletionFromSelf, dispatchFn: dispatcher, @@ -160,7 +160,7 @@ func (c *Completion) complete (ctx context.Context, cctx CompletionContext, cid if ctx.Err() != nil { c.mux.Unlock() c.notify(CompletionFromContext,CompletionStateCanceled,nil,G.ErrCancelByContext) - return G.ErrCancelByContext.WithCause(ctx.Err()) + return G.ErrCancelByContext.ToErrCause(ctx.Err()) } ch := make(chan struct{}) @@ -169,7 +169,7 @@ func (c *Completion) complete (ctx context.Context, cctx CompletionContext, cid select { case <- ch: case <- cc.Context.Done() : - er := G.ErrCancelByContext.WithCause(cc.Context.Err()) + er := G.ErrCancelByContext.ToErrCause(cc.Context.Err()) cc.notify(CompletionFromContext,CompletionStateCanceled,nil,er) } }(c,ch) diff --git a/core/syncs/tasks/completion_ops.go b/core/syncs/tasks/completion_ops.go index 93a077d..0cc2ec6 100644 --- a/core/syncs/tasks/completion_ops.go +++ b/core/syncs/tasks/completion_ops.go @@ -24,9 +24,9 @@ package tasks import ( - "activething.com/go/gogo/core/opts" - "activething.com/go/gogog/core/timers" + "opts" "time" + "timers" ) @@ -34,7 +34,7 @@ import ( // Desc: Función para el constructor de Completion, // establece el mode de como se llegara a la completion // -func CompletionWithMode (mode CompletionMode) opts.OptionErrFnc[Completion] { +func CompletionWithMode (mode CompletionMode) opts.GOptionErrFnc[Completion] { return func(c *Completion) error { c.mode = mode return nil @@ -46,7 +46,7 @@ func CompletionWithMode (mode CompletionMode) opts.OptionErrFnc[Completion] { // Desc: Función para el constructor de Completion // Establece el tiempo máximo de la completion // -func CompletionWithTimeout (tmo time.Duration) opts.OptionErrFnc[Completion] { +func CompletionWithTimeout (tmo time.Duration) opts.GOptionErrFnc[Completion] { return func(c *Completion) error { c.expiration = timers.NewExpiration(tmo,func(){ c.notify(CompletionFromTimer,CompletionStateCanceled,nil,G.ErrCancelByTimer) @@ -60,7 +60,7 @@ func CompletionWithTimeout (tmo time.Duration) opts.OptionErrFnc[Completion] { // Desc: Función para el constructor de Completion, // añade los objeto a completar en la completion // -func CompletionWithCompleters (completers ... Completer) opts.OptionErrFnc[Completion] { +func CompletionWithCompleters (completers ... Completer) opts.GOptionErrFnc[Completion] { return func(c *Completion) error { if c.State() != CompletionStateCreated { return G.ErrInvalidState @@ -78,7 +78,7 @@ func CompletionWithCompleters (completers ... Completer) opts.OptionErrFnc[Compl // utiliza un pool de trabajo para completar // las tareas. // -func CompletionWithWorkPool (pool *WorkerPool) opts.OptionErrFnc[Completion] { +func CompletionWithWorkPool (pool *WorkerPool) opts.GOptionErrFnc[Completion] { return func(c *Completion) error { c.dispatchFn = dispatchFnc(func (cm *Completion, cr Completer, ci int) error { return pool.DoWork(c,func(){ diff --git a/core/syncs/tasks/gogo.go b/core/syncs/tasks/gogo.go index 8aa92cc..90f1f4e 100644 --- a/core/syncs/tasks/gogo.go +++ b/core/syncs/tasks/gogo.go @@ -23,7 +23,7 @@ package tasks -import "activething.com/go/gogo/core/errs" +import "errs" diff --git a/core/syncs/tasks/worker_pool.go b/core/syncs/tasks/worker_pool.go index 7bf355a..a870959 100644 --- a/core/syncs/tasks/worker_pool.go +++ b/core/syncs/tasks/worker_pool.go @@ -24,12 +24,12 @@ package tasks import ( - "activething.com/go/gogo/core/opts" - "activething.com/go/gogo/core/syncs" "context" + "opts" "runtime" "sync" "sync/atomic" + "syncs" "time" ) @@ -64,12 +64,12 @@ type ( ) -func NewWorkerPool (options ...opts.OptionFnc[WorkerPool]) *WorkerPool { +func NewWorkerPool (options ...opts.GOptionFnc[WorkerPool]) *WorkerPool { c := int32(1) if runtime.GOMAXPROCS(0) == 1 { c = 0 } - wp := opts.Apply[WorkerPool](&WorkerPool{ + wp := opts.GApply[WorkerPool](&WorkerPool{ stats : syncs.NewPoolStatsInfo(), maxIdle: int64(maxIdle), cap : c, diff --git a/core/timers/active.go b/core/timers/active.go new file mode 100644 index 0000000..1e62e8f --- /dev/null +++ b/core/timers/active.go @@ -0,0 +1,71 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "activething.com/go/gogo/core/errs" +) + + + + +var ( + + + G = struct { + ErrExpired errs.ErrCode + ErrTimeout errs.ErrCode + ErrIntervalInvalid errs.ErrCode + ErrNotStarted errs.ErrCode + ErrNotPaused errs.ErrCode + ErrAlreadyStarted errs.ErrCode + ErrTransitionInvalid errs.ErrCode + ErrTimerTaskInvalid errs.ErrCode + } { + + ErrExpired: errs.ErrCode("expired"), + ErrTimeout: errs.ErrCode("timeout"), + ErrIntervalInvalid: errs.ErrCode("interval_invalid"), + ErrNotStarted: errs.ErrCode("not_started"), + ErrNotPaused: errs.ErrCode("not_paused"), + ErrAlreadyStarted: errs.ErrCode("already_started"), + ErrTransitionInvalid:errs.ErrCode("transition_invalid"), + ErrTimerTaskInvalid: errs.ErrCode("timer_task_invalid"), + } + + +) + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/chrono.go b/core/timers/chrono.go new file mode 100644 index 0000000..9164614 --- /dev/null +++ b/core/timers/chrono.go @@ -0,0 +1,108 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "sync/atomic" + "time" + +) + +type ( + + + Chrono struct { + started int64 + stopped int64 + } +) + + +func NewChrono(str bool) *Chrono { + ch := &Chrono{} + if str { + ch.Start() + } + return ch +} + + +func (c *Chrono) Started() time.Time { + return time.Unix(0,atomic.LoadInt64(&c.started)) } + + +func (c *Chrono) Stopped() time.Time { + return time.Unix(0,atomic.LoadInt64(&c.stopped)) } + + +func (c *Chrono) Elapsed() time.Duration { + st := atomic.LoadInt64(&c.started) + if st == 0 { + return 0 + } + sp := atomic.LoadInt64(&c.stopped) + if sp == 0 { + return time.Now().Sub(time.Unix(0, int64(st))) + } + return time.Duration(sp - st) +} + + +func (c *Chrono) Reset() { + atomic.StoreInt64(&c.started,0) + atomic.StoreInt64(&c.stopped,0) +} + + +func (c *Chrono) Start() (tm time.Time) { + tm = time.Now() + atomic.CompareAndSwapInt64(&c.started,0,tm.UnixNano()) + return +} + + +func (c *Chrono) Stop() (tm time.Time) { + tm = time.Now() + atomic.CompareAndSwapInt64(&c.stopped,0,tm.UnixNano()) + return +} + + +func (c *Chrono) SetTimes(str, stp time.Time) { + atomic.StoreInt64(&c.started,str.UnixNano()) + atomic.StoreInt64(&c.stopped,str.UnixNano()) +} + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/chronometer.go b/core/timers/chronometer.go new file mode 100644 index 0000000..da30045 --- /dev/null +++ b/core/timers/chronometer.go @@ -0,0 +1,59 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "time" +) + +type ( + + + Chronometer interface { + + Started() time.Time + + Stopped() time.Time + + Elapsed() time.Duration + + } + +) + + + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/expiration.go b/core/timers/expiration.go new file mode 100644 index 0000000..66e303b --- /dev/null +++ b/core/timers/expiration.go @@ -0,0 +1,182 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "sync/atomic" + "time" + + +) + +type ( + + + Expiration struct { + timer *time.Timer + times Chrono + + state uint32 + source uint32 + timeout int64 + deadline int64 + + expireFn func() + } +) + + +func NewExpiration(duration time.Duration, expirer func()) *Expiration { + ex := &Expiration{ + times: *NewChrono(false), + state: uint32(TimerStateCreated), + source: uint32(ExpirationSourceSelf), + timeout: int64(duration), + expireFn: expirer, + } + if duration > 0 { + atomic.StoreUint32(&ex.source,uint32(ExpirationSourceTimer)) + } + return ex +} + + + +func (e *Expiration) Reset() { + e.times.Reset() + atomic.StoreUint32(&e.state ,uint32(TimerStateCreated)) + atomic.StoreUint32(&e.source,uint32(ExpirationSourceSelf)) + atomic.StoreInt64 (&e.timeout ,0) + atomic.StoreInt64 (&e.deadline,0) + if e.timer != nil { + if !e.timer.Stop() { + select { + case <- e.timer.C: + default: + } + } + } +} + + + +func (e *Expiration) State() TimerState { + return TimerState(atomic.LoadUint32(&e.state)) } + + + +func (e *Expiration) Chono() Chronometer { + return &e.times } + + +func (e *Expiration) Source() ExpirationSource { + return ExpirationSource(atomic.LoadUint32(&e.source)) } + + + +func (e *Expiration) Timeout() time.Duration { + return time.Duration(atomic.LoadInt64(&e.timeout)) } + + +func (e *Expiration) Deadline() time.Time { + return time.Unix(0,atomic.LoadInt64(&e.deadline)) } + + + +func (e *Expiration) Start(expirer Expirer) error { + if !atomic.CompareAndSwapUint32(&e.state,uint32(TimerStateCreated),uint32(TimerStateStarted)) { + return G.ErrAlreadyStarted + } + + st := e.times.Start() + to := time.Duration(atomic.LoadInt64(&e.timeout)) + dl := st.Add(to) + + if expirer != nil { + pd, ph := expirer.Deadline() + if ph { + if pd.Before(st) { + atomic.StoreUint32(&e.state ,uint32(TimerStateExpired)) + atomic.StoreUint32(&e.source,uint32(ExpirationSourceParent)) + return G.ErrExpired + } + if pd.Before(dl) { + atomic.StoreInt64 (&e.deadline,pd.UnixNano()) + atomic.StoreUint32(&e.source ,uint32(ExpirationSourceParent)) + return nil + } + } + } + + if atomic.LoadUint32(&e.source) == uint32(ExpirationSourceTimer) { + atomic.StoreInt64 (&e.deadline,dl.UnixNano()) + if e.timer == nil { + e.timer = time.NewTimer(to) + } else { + e.timer.Reset(to) + } + go func(ex *Expiration) { + select { + case <-ex.timer.C: + atomic.StoreUint32(&ex.state,uint32(TimerStateExpired)) + ex.times.Stop() + if ex.expireFn != nil { + ex.expireFn() + } + return + } + }(e) + + + } + + return nil +} + + +func (e *Expiration) Stop() error { + if !atomic.CompareAndSwapUint32(&e.state,uint32(TimerStateStarted), uint32(TimerStateStopped)) { + return G.ErrNotStarted + } + e.times.Stop() + if e.timer != nil && !e.timer.Stop() { + select { + case <-e.timer.C: + } + } + return nil +} + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/expiration_source.go b/core/timers/expiration_source.go new file mode 100644 index 0000000..da24175 --- /dev/null +++ b/core/timers/expiration_source.go @@ -0,0 +1,62 @@ +// ------------------------------------------------------------------------ +// 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 timers + +const ( + + //ExpirationSource + + ExpirationSourceSelf ExpirationSource = iota + ExpirationSourceParent + ExpirationSourceTimer +) + +type ( + + ExpirationSource uint32 +) + + +func (c ExpirationSource) String() string { + switch c { + case ExpirationSourceSelf : return "self" + case ExpirationSourceParent : return "parent" + case ExpirationSourceTimer : return "timer" + default: + return "unknown" + } +} + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/expirer.go b/core/timers/expirer.go new file mode 100644 index 0000000..7162bf1 --- /dev/null +++ b/core/timers/expirer.go @@ -0,0 +1,57 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "time" +) + +type ( + + + Expirer interface { + + Deadline() (time.Time, bool) + } + + + DeadlineFnc func() (time.Time, bool) +) + + +func (f DeadlineFnc) Deadline() (time.Time, bool) { + return f() } + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/go.mod b/core/timers/go.mod new file mode 100644 index 0000000..8bad306 --- /dev/null +++ b/core/timers/go.mod @@ -0,0 +1,4 @@ +module timers + + +go 1.20 \ No newline at end of file diff --git a/core/timers/pulsar.go b/core/timers/pulsar.go new file mode 100644 index 0000000..dcc6789 --- /dev/null +++ b/core/timers/pulsar.go @@ -0,0 +1,174 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "activething.com/go/gogo/core/contexts" + "context" + "sync/atomic" + "time" + + + + + +) + +type ( + + + Pulsar struct { + state uint32 + pulse int64 + err atomic.Value + cancel contexts.Cancel + workFn PulsarWorkFnc + } + + + PulsarWorkFnc func(ctx context.Context) error + +) + + +func NewPulsar(pulse time.Duration, work PulsarWorkFnc) *Pulsar { + return &Pulsar{ + state: uint32(TimerStateCreated), + cancel: *contexts.NewCancel(), + pulse: int64(pulse), + workFn: work, + } +} + + +func (p *Pulsar) State() TimerState { + return TimerState(atomic.LoadUint32(&p.state)) } + + +func (p *Pulsar) Pulse() time.Duration { + return time.Duration(atomic.LoadInt64(&p.pulse)) } + + +func (p *Pulsar) Err() error { + if v := p.err.Load(); v != nil { + return v.(error) + } + return nil +} + + +func (p *Pulsar) Start(ctx context.Context) error { + if !atomic.CompareAndSwapUint32(&p.state,uint32(TimerStateCreated),uint32(TimerStateStarted)) { + return G.ErrAlreadyStarted + } + to := atomic.LoadInt64(&p.pulse) + if to <= 0 { + atomic.StoreUint32(&p.state,uint32(TimerStateExpired)) + return G.ErrIntervalInvalid + } + + cx := ctx + if cx == nil { + cx = context.Background() + } else if ctx.Err() != nil { + atomic.StoreUint32(&p.state,uint32(TimerStateCanceled)) + return context.Canceled + } + + go func(pl *Pulsar, tk *time.Ticker, cx context.Context) { + EXIT: + for { + select { + case <-cx.Done(): + atomic.StoreUint32(&pl.state,uint32(TimerStateCanceled)) + break EXIT + + case <-pl.cancel.Done(): + break EXIT + + case <-tk.C: + if TimerState(atomic.LoadUint32(&pl.state)) == TimerStateStarted { + if er := pl.workFn(cx); er != nil { + pl.err.Store(er) + atomic.StoreUint32(&pl.state,uint32(TimerStateCanceled)) + break EXIT + } + } + } + } + tk.Stop() + pl.cancel.Cancel() + + }(p, time.NewTicker(time.Duration(to)), cx) + + return nil +} + + +func (p *Pulsar) Stop() error { + if atomic.CompareAndSwapUint32(&p.state,uint32(TimerStateStarted), uint32(TimerStateStopped)) || + atomic.CompareAndSwapUint32(&p.state,uint32(TimerStatePaused), uint32(TimerStateStopped)) { + p.cancel.Cancel() + return nil + } + return G.ErrNotStarted +} + + +func (p *Pulsar) Pause() error { + if atomic.CompareAndSwapUint32(&p.state,uint32(TimerStateStarted), uint32(TimerStatePaused)) { + return nil + } + return G.ErrNotStarted +} + + +func (p *Pulsar) Continue() error { + if atomic.CompareAndSwapUint32(&p.state,uint32(TimerStatePaused), uint32(TimerStateStarted)) { + return nil + } + return G.ErrNotPaused +} + + +func (p *Pulsar) Done() <-chan struct{} { + return p.cancel.Done() } + + +func (p *Pulsar) Wait() { + <-p.cancel.Done() } + + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/time_task.go b/core/timers/time_task.go new file mode 100644 index 0000000..5bcc147 --- /dev/null +++ b/core/timers/time_task.go @@ -0,0 +1,48 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "context" + "time" +) + +type ( + + + TimeTaskFnc func(now time.Time, ctx context.Context) +) + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/timer_queue.go b/core/timers/timer_queue.go new file mode 100644 index 0000000..7f5b3d7 --- /dev/null +++ b/core/timers/timer_queue.go @@ -0,0 +1,190 @@ +// ------------------------------------------------------------------------ +// 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 timers + +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 - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/timer_state.go b/core/timers/timer_state.go new file mode 100644 index 0000000..00c88b2 --- /dev/null +++ b/core/timers/timer_state.go @@ -0,0 +1,81 @@ +// ------------------------------------------------------------------------ +// 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 timers + +const ( + + //Enum values for TimerState + TimerStateCreated TimerState = 1 << iota + TimerStateStarted + TimerStatePaused + TimerStateExpired + TimerStateStopped + TimerStateCanceled + + timerStateActive = TimerStateStarted + TimerStatePaused +) + +type ( + + + TimerState uint32 +) + + + +func (t TimerState) String() string { + switch t { + case TimerStateCreated: + return "created" + case TimerStateStarted: + return "started" + case TimerStatePaused: + return "paused" + case TimerStateExpired: + return "expired" + case TimerStateStopped: + return "stopped" + case TimerStateCanceled: + return "canceled" + default: + return "unknown" + } +} + + +func (t TimerState) IsActive() bool { + return t != timerStateActive } + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/timer_task.go b/core/timers/timer_task.go new file mode 100644 index 0000000..e52ef44 --- /dev/null +++ b/core/timers/timer_task.go @@ -0,0 +1,60 @@ +// ------------------------------------------------------------------------ +// 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 timers + +import ( + "time" +) + +type ( + + + TimerTask struct { + + RunAt time.Time + TaskFn TimeTaskFnc + } +) + + +func NewTimerTask(runAt time.Time, task TimeTaskFnc) *TimerTask { + return &TimerTask{ + RunAt: runAt, + TaskFn: task, + } +} + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/core/timers/timer_task_heap.go b/core/timers/timer_task_heap.go new file mode 100644 index 0000000..13f53dd --- /dev/null +++ b/core/timers/timer_task_heap.go @@ -0,0 +1,89 @@ +// ------------------------------------------------------------------------ +// 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 timers + +type ( + + + TimerTaskHeap []*TimerTask +) + + +func (h *TimerTaskHeap) Reset() { + *h = (*h)[:0] } + + +func (h *TimerTaskHeap) PushTimerTask(ele *TimerTask) { + if ele != nil { + *h = append(*h, ele) + } +} + +func (h *TimerTaskHeap) PopTimerTask() *TimerTask { + th := *h + if ln := len(th); ln > 0 { + tt := th[ln-1] + *h = th[:ln-1] + return tt + } + return nil +} + + +func (h *TimerTaskHeap) Push(ele interface{}) { + h.PushTimerTask(ele.(*TimerTask)) } + + +func (h *TimerTaskHeap) Pop() interface{} { + return h.PopTimerTask() } + + +func (h TimerTaskHeap) Len() int { + return len(h) } + + +func (h TimerTaskHeap) Less(i, j int) bool { + return h[i].RunAt.Before(h[j].RunAt) } + + +func (h TimerTaskHeap) Swap(i, j int) { + h[i], h[j] = h[j], h[i] } + + + + + + +/* + :::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + + I HAVE NO DESIRE TO WALK ON WATER," SAID SIDDHARTA. + "LET THE OLD SHRAMANAS SATISFY THEMSELVES WITH SUCH SKILLS. + + SIDDHARTA + - HERMANN HESSE - + + ::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::::: + +*/ diff --git a/go.work b/go.work index db9f427..43b243d 100644 --- a/go.work +++ b/go.work @@ -8,6 +8,7 @@ use ( ./core/bins ./core/opts ./core/syncs + ./core/timers ./app