-
Notifications
You must be signed in to change notification settings - Fork 3
/
pool.go
95 lines (84 loc) · 1.59 KB
/
pool.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
package workerpool
import (
"context"
"sync/atomic"
)
type Worker struct {
ctx context.Context
cancel context.CancelFunc
err *errorContainer
pendingTasks chan WorkerFunc
activeOperations int64
}
type WorkerFunc func(context.Context) error
func New(ctx context.Context, workers int) *Worker {
cCtx, ctxCancel := context.WithCancel(ctx)
wrk := &Worker{
ctx: cCtx,
cancel: ctxCancel,
err: &errorContainer{},
pendingTasks: make(chan WorkerFunc),
activeOperations: 0,
}
wrk.start(workers)
return wrk
}
func (w *Worker) start(numWorkers int) {
for i := 0; i < numWorkers; i++ {
go func() {
for {
select {
case fn, ok := <-w.pendingTasks:
if !ok {
return
}
if err := fn(w.ctx); err != nil {
w.err.AssignError(err)
}
atomic.AddInt64(&w.activeOperations, -1)
case <-w.ctx.Done():
if w.ctx.Err() != nil {
w.err.AssignError(w.ctx.Err())
}
return
default:
if w.err.err != nil {
w.cancel()
}
}
}
}()
}
}
func (w *Worker) Do(fns ...WorkerFunc) {
for _, fn := range fns {
select {
case <-w.ctx.Done():
return
default:
atomic.AddInt64(&w.activeOperations, 1)
f := fn
go func() {
select {
case <-w.ctx.Done():
return
default:
w.pendingTasks <- f
}
}()
}
}
}
func (w *Worker) Wait() error {
defer close(w.pendingTasks)
for {
select {
case <-w.ctx.Done():
return w.err.err
default:
if atomic.LoadInt64(&w.activeOperations) == 0 {
return w.err.err
}
}
}
}