2018-05-19 07:28:03 +03:00
|
|
|
package ants
|
|
|
|
|
2018-05-19 14:08:31 +03:00
|
|
|
import (
|
|
|
|
"sync/atomic"
|
|
|
|
"container/list"
|
|
|
|
"sync"
|
|
|
|
)
|
2018-05-19 07:28:03 +03:00
|
|
|
|
|
|
|
type Worker struct {
|
|
|
|
pool *Pool
|
|
|
|
task chan f
|
|
|
|
exit chan sig
|
|
|
|
}
|
|
|
|
|
|
|
|
func (w *Worker) run() {
|
|
|
|
go func() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case f := <-w.task:
|
|
|
|
f()
|
2018-05-19 20:38:50 +03:00
|
|
|
w.pool.putWorker(w)
|
2018-05-19 14:51:33 +03:00
|
|
|
w.pool.wg.Done()
|
2018-05-19 07:28:03 +03:00
|
|
|
case <-w.exit:
|
2018-05-19 13:24:36 +03:00
|
|
|
atomic.AddInt32(&w.pool.running, -1)
|
2018-05-19 07:28:03 +03:00
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (w *Worker) stop() {
|
|
|
|
w.exit <- sig{}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (w *Worker) sendTask(task f) {
|
|
|
|
w.task <- task
|
|
|
|
}
|
2018-05-19 14:08:31 +03:00
|
|
|
|
|
|
|
//--------------------------------------------------------------------------------
|
|
|
|
|
|
|
|
type ConcurrentQueue struct {
|
|
|
|
queue *list.List
|
|
|
|
m sync.Mutex
|
|
|
|
}
|
|
|
|
|
|
|
|
func NewConcurrentQueue() *ConcurrentQueue {
|
|
|
|
q := new(ConcurrentQueue)
|
|
|
|
q.queue = list.New()
|
|
|
|
return q
|
|
|
|
}
|
|
|
|
|
|
|
|
func (q *ConcurrentQueue) push(v interface{}) {
|
|
|
|
defer q.m.Unlock()
|
|
|
|
q.m.Lock()
|
|
|
|
q.queue.PushFront(v)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (q *ConcurrentQueue) pop() interface{} {
|
|
|
|
defer q.m.Unlock()
|
|
|
|
q.m.Lock()
|
2018-05-19 20:38:50 +03:00
|
|
|
if elem := q.queue.Back(); elem != nil {
|
2018-05-19 14:08:31 +03:00
|
|
|
return q.queue.Remove(elem)
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|