parent
3fca3ef79d
commit
deb98e9d4e
@ -0,0 +1,65 @@
|
||||
// ------------------------------------------------------------------------
|
||||
// Project atila
|
||||
// Active Thing (activething.com) git.activething.com/go
|
||||
//
|
||||
// File name pool_stats.go
|
||||
// Created by DEV
|
||||
// Modified 17/01/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 syncs
|
||||
|
||||
import "sync/atomic"
|
||||
|
||||
type (
|
||||
|
||||
PoolStatsInfo struct {
|
||||
created uint64
|
||||
acquired uint64
|
||||
released uint64
|
||||
}
|
||||
)
|
||||
|
||||
func NewPoolStatsInfo () *PoolStatsInfo {
|
||||
return &PoolStatsInfo{} }
|
||||
|
||||
|
||||
func (p *PoolStatsInfo) Created() uint64 {
|
||||
return atomic.LoadUint64(&p.created)
|
||||
}
|
||||
|
||||
func (p *PoolStatsInfo) Acquired() uint64 {
|
||||
return atomic.LoadUint64(&p.acquired)
|
||||
}
|
||||
|
||||
func (p *PoolStatsInfo) Released() uint64 {
|
||||
return atomic.LoadUint64(&p.released)
|
||||
}
|
||||
|
||||
func (p *PoolStatsInfo) AddCreated() {
|
||||
atomic.AddUint64(&p.created, 1)
|
||||
}
|
||||
|
||||
func (p *PoolStatsInfo) AddReleased() {
|
||||
atomic.AddUint64(&p.released, 1)
|
||||
}
|
||||
|
||||
func (p *PoolStatsInfo) AddAcquired() {
|
||||
atomic.AddUint64(&p.acquired, 1)
|
||||
}
|
||||
|
||||
|
||||
@ -0,0 +1,26 @@
|
||||
// ------------------------------------------------------------------------
|
||||
// Project atila
|
||||
// Active Thing (activething.com) git.activething.com/go
|
||||
//
|
||||
// File name gogo.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
|
||||
|
||||
|
||||
@ -0,0 +1,38 @@
|
||||
// ------------------------------------------------------------------------
|
||||
// Project atila
|
||||
// Active Thing (activething.com) git.activething.com/go
|
||||
//
|
||||
// File name worker.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 "time"
|
||||
|
||||
type (
|
||||
|
||||
worker struct {
|
||||
lastTime time.Time
|
||||
jobCh chan func()
|
||||
}
|
||||
|
||||
|
||||
)
|
||||
|
||||
|
||||
@ -0,0 +1,151 @@
|
||||
// ------------------------------------------------------------------------
|
||||
// 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"
|
||||
"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 := WorkerPoolState(atomic.LoadUint32(&p.state))
|
||||
if st == WorkerPoolStateRunning {
|
||||
p.mux.Unlock()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -0,0 +1,40 @@
|
||||
// ------------------------------------------------------------------------
|
||||
// Project atila
|
||||
// Active Thing (activething.com) git.activething.com/go
|
||||
//
|
||||
// File name worker_pool_state.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
|
||||
|
||||
const (
|
||||
|
||||
WorkerPoolStateCreated WorkerPoolState = iota
|
||||
WorkerPoolStateRunning
|
||||
WorkerPoolStateStopped
|
||||
|
||||
)
|
||||
|
||||
type (
|
||||
|
||||
WorkerPoolState uint32
|
||||
)
|
||||
|
||||
|
||||
|
||||
Loading…
Reference in new issue