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.
282 lines
4.9 KiB
282 lines
4.9 KiB
// ------------------------------------------------------------------------
|
|
// Project atila
|
|
// Active Thing (activething.com) git.activething.com/go
|
|
//
|
|
// File name worker_pool.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 works
|
|
|
|
import (
|
|
"context"
|
|
"runtime"
|
|
"sync"
|
|
"sync/atomic"
|
|
"syns/pools"
|
|
"time"
|
|
)
|
|
|
|
|
|
type (
|
|
|
|
Pool struct {
|
|
stats *pools.PoolStatsInfo
|
|
pool sync.Pool
|
|
mux sync.Mutex
|
|
|
|
cap int32
|
|
max int32
|
|
count int32
|
|
state uint32
|
|
|
|
maxIdle int64
|
|
|
|
workers []*worker
|
|
doneCh chan struct{}
|
|
}
|
|
|
|
|
|
PoolFnc func(wp *Pool)
|
|
)
|
|
|
|
|
|
func NewPool(options ...PoolFnc) *Pool {
|
|
c := int32(1)
|
|
if runtime.GOMAXPROCS(0) == 1 {
|
|
c = 0
|
|
}
|
|
wp := &Pool{
|
|
stats : pools.NewPoolStatsInfo(),
|
|
maxIdle: int64(maxIdle),
|
|
cap : c,
|
|
max : maxWorkers,
|
|
state : uint32(PoolStateCreated),
|
|
}
|
|
for _,o := range options {
|
|
o(wp)
|
|
}
|
|
wp.pool.New = wp.newWorker
|
|
return wp
|
|
}
|
|
|
|
|
|
func (p *Pool) newWorker () any {
|
|
p.stats.AddCreated()
|
|
return &worker{
|
|
jobCh: make(chan func(), atomic.LoadInt32(&p.cap)) }}
|
|
|
|
|
|
func (p *Pool) Stats () pools.PoolStats {
|
|
return p.stats }
|
|
|
|
|
|
func (p *Pool) State () PoolState {
|
|
return PoolState(atomic.LoadUint32(&p.state)) }
|
|
|
|
func (p *Pool) Count () int32 {
|
|
return atomic.LoadInt32(&p.count) }
|
|
|
|
|
|
func (p *Pool) MaxWorkers () int32 {
|
|
return atomic.LoadInt32(&p.max) }
|
|
|
|
func (p *Pool) SetMaxWorkers (max int32) {
|
|
if max <= 0 {
|
|
max = maxWorkers
|
|
}
|
|
atomic.StoreInt32(&p.max,max)
|
|
}
|
|
|
|
|
|
func (p *Pool) WorkerCap () int32 {
|
|
return atomic.LoadInt32(&p.cap) }
|
|
|
|
func (p *Pool) SetWorkerCap (cap int32) {
|
|
if cap <= 0 {
|
|
cap = 0
|
|
}
|
|
atomic.StoreInt32(&p.cap,cap)
|
|
}
|
|
|
|
|
|
func (p *Pool) MaxWorkerIdle () time.Duration {
|
|
return time.Duration(atomic.LoadInt64(&p.maxIdle)) }
|
|
|
|
func (p *Pool) SetMaxWorkerIdle (idle time.Duration) {
|
|
if idle < minIdle {
|
|
idle = minIdle
|
|
}
|
|
atomic.StoreInt64(&p.maxIdle,int64(idle))
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func (p *Pool) Done () <- chan struct{} {
|
|
return p.doneCh }
|
|
|
|
|
|
|
|
func (p *Pool) Start () error {
|
|
p.mux.Lock()
|
|
|
|
st := p.State()
|
|
if st == PoolStateRunning {
|
|
p.mux.Unlock()
|
|
return G.ErrPoolStarted
|
|
}
|
|
if st == PoolStateStopped {
|
|
atomic.StoreInt32(&p.count,0)
|
|
}
|
|
|
|
p.doneCh = make(chan struct{})
|
|
p.mux.Unlock()
|
|
|
|
go func(p *Pool) {
|
|
atomic.StoreUint32(&p.state,uint32(PoolStateRunning))
|
|
var ls []*worker
|
|
for {
|
|
p.clearWorkers(&ls)
|
|
select {
|
|
case <- p.doneCh:
|
|
return
|
|
default:
|
|
time.Sleep(time.Duration(atomic.LoadInt64(&p.maxIdle)))
|
|
}
|
|
}
|
|
}(p)
|
|
|
|
return nil
|
|
}
|
|
|
|
|
|
func (p *Pool) Stop () error {
|
|
p.mux.Lock()
|
|
if p.State() != PoolStateRunning {
|
|
p.mux.Unlock()
|
|
return G.ErrPoolStopped
|
|
}
|
|
atomic.StoreUint32(&p.state,uint32(PoolStateStopped))
|
|
|
|
close(p.doneCh)
|
|
p.doneCh = nil
|
|
ws := p.workers
|
|
|
|
for ix,wk := range ws {
|
|
wk.jobCh <- nil
|
|
ws[ix]=nil
|
|
}
|
|
|
|
p.workers = ws[:0]
|
|
|
|
p.mux.Unlock()
|
|
return nil
|
|
}
|
|
|
|
|
|
func (p *Pool) DoWork(ctx context.Context, job func()) error {
|
|
p.mux.Lock()
|
|
if p.State() != PoolStateRunning {
|
|
p.mux.Unlock()
|
|
return G.ErrPoolStopped
|
|
}
|
|
|
|
wk := (*worker)(nil)
|
|
ws := p.workers
|
|
ln := len(ws) - 1
|
|
|
|
if ln >= 0 {
|
|
wk = ws[ln]
|
|
ws[ln] = nil
|
|
p.workers = ws[:ln]
|
|
p.mux.Unlock()
|
|
} else {
|
|
p.mux.Unlock()
|
|
|
|
if atomic.LoadInt32(&p.count) >= atomic.LoadInt32(&p.max) {
|
|
return G.ErrPoolAcquire
|
|
}
|
|
atomic.AddInt32(&p.count,1)
|
|
wk = p.pool.Get().(*worker)
|
|
|
|
go func(p *Pool, w *worker, x context.Context) {
|
|
|
|
for job := range w.jobCh {
|
|
if job == nil {
|
|
break
|
|
}
|
|
job()
|
|
job = nil
|
|
|
|
p.mux.Lock()
|
|
done := false
|
|
if p.State() != PoolStateStopped {
|
|
w.lastTime = time.Now()
|
|
p.workers = append(p.workers, w)
|
|
done = true
|
|
p.stats.AddReleased()
|
|
}
|
|
p.mux.Unlock()
|
|
if !done {
|
|
break
|
|
}
|
|
}
|
|
atomic.AddInt32(&p.count,-1)
|
|
p.pool.Put(w)
|
|
|
|
}(p,wk,ctx)
|
|
|
|
}
|
|
|
|
p.stats.AddAcquired()
|
|
wk.jobCh <- job
|
|
return nil
|
|
}
|
|
|
|
|
|
|
|
|
|
func (p *Pool) clearWorkers (cls *[]*worker) {
|
|
p.mux.Lock()
|
|
tm := time.Now()
|
|
mi := time.Duration(atomic.LoadInt64(&p.maxIdle))
|
|
ws := p.workers
|
|
ln := len(ws)
|
|
ix := 0
|
|
|
|
for ix <ln && tm.Sub(ws[ix].lastTime) > mi {
|
|
ix++
|
|
}
|
|
*cls = append((*cls)[:0],ws[:ix]...)
|
|
if ix > 0 {
|
|
nt := copy(ws,ws[ix:])
|
|
for ix:= nt; ix<ln; ix++ {
|
|
ws[ix]=nil
|
|
}
|
|
p.workers=ws[:nt]
|
|
}
|
|
p.mux.Unlock()
|
|
|
|
cl := *cls
|
|
for ix,wk := range cl {
|
|
wk.jobCh <- nil
|
|
cl[ix]=nil
|
|
}
|
|
} |