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.

288 lines
5.2 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 tasks
import (
"activething.com/go/gogo/core/opts"
"activething.com/go/gogo/core/syncs"
"context"
"runtime"
"sync"
"sync/atomic"
"time"
)
const (
minWorkers = 10
maxWorkers = 1000
minIdle = time.Millisecond * 100
maxIdle = time.Second * 10
)
type (
WorkerPool struct {
stats *syncs.PoolStatsInfo
pool sync.Pool
mux sync.Mutex
cap int32
max int32
count int32
state uint32
maxIdle int64
workers []*worker
doneCh chan struct{}
}
)
func NewWorkerPool (options ...opts.OptionFnc[WorkerPool]) *WorkerPool {
c := int32(1)
if runtime.GOMAXPROCS(0) == 1 {
c = 0
}
wp := opts.Apply[WorkerPool](&WorkerPool{
stats : syncs.NewPoolStatsInfo(),
maxIdle: int64(maxIdle),
cap : c,
max : maxWorkers,
state : uint32(WorkerPoolStateCreated),
}, options ...)
wp.pool.New = wp.newWorker
return wp
}
func (p *WorkerPool) newWorker () any {
p.stats.AddCreated()
return &worker {
jobCh: make(chan func(), atomic.LoadInt32(&p.cap)) }}
func (p *WorkerPool) Stats () syncs.PoolStats {
return p.stats }
func (p *WorkerPool) State () WorkerPoolState {
return WorkerPoolState(atomic.LoadUint32(&p.state)) }
func (p *WorkerPool) Count () int32 {
return atomic.LoadInt32(&p.count) }
func (p *WorkerPool) MaxWorkers () int32 {
return atomic.LoadInt32(&p.max) }
func (p *WorkerPool) SetMaxWorkers (max int32) {
if max <= 0 {
max = maxWorkers
}
atomic.StoreInt32(&p.max,max)
}
func (p *WorkerPool) WorkerCap () int32 {
return atomic.LoadInt32(&p.cap) }
func (p *WorkerPool) SetWorkerCap (cap int32) {
if cap <= 0 {
cap = 0
}
atomic.StoreInt32(&p.cap,cap)
}
func (p *WorkerPool) MaxWorkerIdle () time.Duration {
return time.Duration(atomic.LoadInt64(&p.maxIdle)) }
func (p *WorkerPool) SetMaxWorkerIdle (idle time.Duration) {
if idle < minIdle {
idle = minIdle
}
atomic.StoreInt64(&p.maxIdle,int64(idle))
}
func (p *WorkerPool) Done () <- chan struct{} {
return p.doneCh }
func (p *WorkerPool) Start () error {
p.mux.Lock()
st := p.State()
if st == WorkerPoolStateRunning {
p.mux.Unlock()
return G.ErrPoolStarted
}
if st == WorkerPoolStateStopped {
atomic.StoreInt32(&p.count,0)
}
p.doneCh = make(chan struct{})
p.mux.Unlock()
go func(p *WorkerPool) {
atomic.StoreUint32(&p.state,uint32(WorkerPoolStateRunning))
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 *WorkerPool) Stop () error {
p.mux.Lock()
if p.State() != WorkerPoolStateRunning {
p.mux.Unlock()
return G.ErrPoolStopped
}
atomic.StoreUint32(&p.state,uint32(WorkerPoolStateStopped))
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 *WorkerPool) DoWork(ctx context.Context, job func()) error {
p.mux.Lock()
if p.State() != WorkerPoolStateRunning {
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 *WorkerPool, 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() != WorkerPoolStateStopped {
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 *WorkerPool) 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
}
}

Powered by TurnKey Linux.