ants/worker_stack.go

83 lines
1.4 KiB
Go
Raw Permalink Normal View History

package ants
import "time"
type workerStack struct {
items []*goWorker
expiry []*goWorker
}
func newWorkerStack(size int) *workerStack {
2019-10-20 13:22:13 +03:00
return &workerStack{
items: make([]*goWorker, 0, size),
}
}
func (wq *workerStack) len() int {
return len(wq.items)
}
func (wq *workerStack) isEmpty() bool {
return len(wq.items) == 0
}
func (wq *workerStack) insert(worker *goWorker) error {
wq.items = append(wq.items, worker)
return nil
}
func (wq *workerStack) detach() *goWorker {
l := wq.len()
if l == 0 {
return nil
}
w := wq.items[l-1]
2020-08-29 13:51:56 +03:00
wq.items[l-1] = nil // avoid memory leaks
wq.items = wq.items[:l-1]
return w
}
2019-10-20 13:22:13 +03:00
func (wq *workerStack) retrieveExpiry(duration time.Duration) []*goWorker {
n := wq.len()
if n == 0 {
return nil
}
expiryTime := time.Now().Add(-duration)
2019-10-20 13:22:13 +03:00
index := wq.binarySearch(0, n-1, expiryTime)
wq.expiry = wq.expiry[:0]
if index != -1 {
wq.expiry = append(wq.expiry, wq.items[:index+1]...)
m := copy(wq.items, wq.items[index+1:])
2020-08-29 13:51:56 +03:00
for i := m; i < n; i++ {
wq.items[i] = nil
}
wq.items = wq.items[:m]
}
return wq.expiry
}
2019-10-20 13:22:13 +03:00
func (wq *workerStack) binarySearch(l, r int, expiryTime time.Time) int {
var mid int
for l <= r {
mid = (l + r) / 2
if expiryTime.Before(wq.items[mid].recycleTime) {
r = mid - 1
} else {
l = mid + 1
}
}
return r
}
2019-10-20 13:22:13 +03:00
func (wq *workerStack) reset() {
for i := 0; i < wq.len(); i++ {
wq.items[i].task <- nil
2020-08-29 13:51:56 +03:00
wq.items[i] = nil
}
wq.items = wq.items[:0]
}