You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

301 lines
6.9 KiB

// ------------------------------------------------------------------------
// Project atila
// Active Thing (activething.com) git.activething.com/go
//
// File name completion.go
// Created by DEV
// Modified 20/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 tasks
import (
"context"
"opts"
"sync"
"sync/atomic"
"timers"
)
type (
Completion struct {
/*embedded*/context.Context
mux sync.RWMutex
wgr sync.WaitGroup
parent CompletionContext
expiration *timers.Expiration
mode CompletionMode
state uint32
pending uint32
cid int
result CompletionResult
dispatchFn dispatchFnc
cancelFn func()
completers []Completer
}
dispatchFnc func (*Completion, Completer, int) error
)
func NewCompletion (options ...opts.GOptionErrFnc[Completion]) (*Completion,error) {
c,e := opts.GApplyAndCheck[Completion](&Completion {
mode : CompletionModeAll,
cid : CompletionFromSelf,
dispatchFn: dispatcher,
},options...)
if e != nil {
return nil,e
}
if c.expiration == nil {
c.expiration = timers.NewExpiration(0,nil)
}
return c,nil
}
func (c *Completion) Mode () CompletionMode {
return c.mode }
func (c *Completion) State () CompletionState {
return CompletionState(atomic.LoadUint32(&c.state)) }
func (c *Completion) Pending() int {
return int(atomic.LoadUint32(&c.pending)) }
func (c *Completion) IsCompleted () bool {
return c.State() > CompletionStateStarted }
func (c *Completion) Chrono () timers.Chronometer {
return c.expiration.Chono() }
func (c *Completion) Result () CompletionResult {
c.wgr.Wait()
return c.result
}
func (c *Completion) Wait () {
c.wgr.Wait() }
func (c *Completion) AddCompleter (cms ... Completer) error {
c.mux.Lock()
if c.State() != CompletionStateCreated {
c.mux.Unlock()
return G.ErrInvalidState
}
for _,cm := range cms {
c.completers = append(c.completers, cm)
}
c.mux.Unlock()
return nil
}
func (c *Completion) Cancel () {
c.notify(CompletionFromRequest,CompletionStateCanceled,nil,G.ErrCancelByRequest) }
func (c *Completion) Start (ctx context.Context) error {
return c.complete(ctx,nil,CompletionFromSelf) }
func (c *Completion) Complete (cmp CompletionContext, cid int) {
_ = c.complete(nil,cmp,cid) }
func (c *Completion) complete (ctx context.Context, cctx CompletionContext, cid int) error {
c.mux.Lock()
if c.State() != CompletionStateCreated {
c.mux.Unlock()
return G.ErrInvalidState
}
// Nos sirve para determinar de donde va a depender la expiracion
// en tiempo de la completion en caso de que exista, o bien de la
// completion padre o bien del contexto.
//
var ex timers.Expirer
if cctx != nil {
c.parent = cctx
c.Context= cctx
c.cid = cid
if c.parent.IsCompleted() {
c.mux.Unlock()
c.notify(CompletionFromParent,CompletionStateCanceled,nil,G.ErrCancelByParent)
return G.ErrCancelByParent
}
ex = c.parent
} else {
if ctx == nil {
c.Context = context.Background()
} else {
c.Context = ctx
if ctx.Err() != nil {
c.mux.Unlock()
c.notify(CompletionFromContext,CompletionStateCanceled,nil,G.ErrCancelByContext)
return G.ErrCancelByContext.ToErrCause(ctx.Err())
}
ch := make(chan struct{})
c.cancelFn = func(){ close(ch)}
go func(cc *Completion, ch chan struct{}){
select {
case <- ch:
case <- cc.Context.Done() :
er := G.ErrCancelByContext.ToErrCause(cc.Context.Err())
cc.notify(CompletionFromContext,CompletionStateCanceled,nil,er)
}
}(c,ch)
}
ex = ctx
}
if er := c.expiration.Start(ex); er != nil {
c.mux.Unlock()
c.notify(CompletionFromSelf,CompletionStateFaulted,nil,er)
return er
}
atomic.StoreUint32(&c.state,uint32(CompletionStateStarted))
c.wgr.Add(1)
go func (c *Completion) {
ln := len(c.completers)
if ln == 0 {
c.mux.Unlock()
c.notify(CompletionFromSelf, CompletionStateCompleted, nil, nil)
return
}
atomic.StoreUint32(&c.pending,uint32(ln))
cs := c.completers
c.mux.Unlock()
for ix, cm := range cs {
if c.IsCompleted() {
return
}
if er := c.dispatchFn(c,cm,ix); er != nil {
c.notify(CompletionFromSelf,CompletionStateFaulted,nil,er)
return
}
}
}(c)
return nil
}
func (c *Completion) Notify (frm int, val interface{}, err error) {
c.notify(frm,CompletionStateCompleted,val,err) }
func (c *Completion) notify (frm int, ste CompletionState, val interface{}, err error) {
c.mux.Lock()
os := c.State()
if os > CompletionStateStarted {
c.mux.Unlock()
return
}
if frm >= CompletionFromChild {
cn := atomic.AddUint32(&c.pending,-1)
switch {
case c.mode.IsAnyOrResult() && val != nil:
atomic.StoreUint32(&c.state,uint32(CompletionStateCompleted))
c.result = NewCompletionResultValue(frm, val)
case c.mode.IsAnyOrError() && err != nil:
atomic.StoreUint32(&c.state,uint32(CompletionStateFaulted))
c.result = NewCompletionResultError(frm, err)
case c.mode == completionModeTask:
c.result = NewCompletionResult(frm, val, err)
default:
if c.mode == CompletionModeAll {
if c.result == nil {
c.result = NewCompletionResultList(c.cid, NewCompletionResult(frm, val, err))
} else {
c.result.(*CompletionResultList).Add(NewCompletionResult(frm, val, err))
}
}
if cn == 0 {
atomic.StoreUint32(&c.state,uint32(CompletionStateCompleted))
} else {
c.mux.Unlock()
return
}
}
} else {
atomic.StoreUint32(&c.state,uint32(ste))
c.result = NewCompletionResult(frm,val,err)
}
_ = c.expiration.Stop()
if c.cancelFn != nil {
c.cancelFn()
}
c.wgr.Done()
c.mux.Unlock()
// Solo notificamos a las completions hijas y al padre en caso de que se haya iniciado
if os == CompletionStateStarted {
for _, ch := range c.completers {
// Notificamos solo a los completer que a su vez implementan
// CompletionContext
if cm, is := ch.(CompletionContext); is {
if !cm.IsCompleted() {
cm.Notify(CompletionFromParent, nil, G.ErrCancelByParent)
}
}
}
// Notificamos al padre en caso de que tenga y la completion no venga de el
if frm != CompletionFromParent && c.parent != nil {
c.parent.Notify(c.cid,c.result.Value(),c.result.Error())
}
}
}
func dispatcher (completion *Completion, completer Completer, cid int) error {
go completer.Complete(completion,cid)
return nil
}

Powered by TurnKey Linux.