-
Notifications
You must be signed in to change notification settings - Fork 45
/
workerpool.go
103 lines (87 loc) · 1.87 KB
/
workerpool.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
96
97
98
99
100
101
102
103
package work
import (
"context"
"time"
"go.uber.org/zap"
"github.com/streamingfast/substreams/reqctx"
)
type WorkerPool struct {
workers []*WorkerStatus
started *time.Time
}
type WorkerState int
const (
WorkerFree WorkerState = iota
WorkerWorking
WorkerInitialWait
)
type WorkerStatus struct {
State WorkerState
Worker Worker
}
func NewWorkerPool(ctx context.Context, workerCount int, workerFactory WorkerFactory) *WorkerPool {
logger := reqctx.Logger(ctx)
logger.Debug("initializing worker pool", zap.Int("worker_count", workerCount))
workers := make([]*WorkerStatus, workerCount)
for i := 0; i < workerCount; i++ {
state := WorkerFree
if i > 0 {
state = WorkerInitialWait
}
workers[i] = &WorkerStatus{
Worker: workerFactory(logger),
State: state,
}
}
now := time.Now()
return &WorkerPool{
workers: workers,
started: &now,
}
}
func (p *WorkerPool) rampupWorkers() {
if time.Since(*p.started) < time.Second*4 {
// no rampup yet
return
}
for _, w := range p.workers {
if w.State == WorkerInitialWait {
w.State = WorkerFree
}
}
p.started = nil
}
func (p *WorkerPool) inRampupPhase() bool {
return p.started != nil
}
func (p *WorkerPool) WorkerAvailable() (avail bool, shouldRetry bool) {
if p.inRampupPhase() {
p.rampupWorkers()
}
for _, w := range p.workers {
if w.State == WorkerFree {
return true, false
}
}
return false, p.inRampupPhase()
}
func (p *WorkerPool) Borrow() Worker {
for _, status := range p.workers {
if status.State == WorkerFree {
status.State = WorkerWorking
return status.Worker
}
}
panic("no free workers, call WorkerAvailable() first")
}
func (p *WorkerPool) Return(worker Worker) {
for _, status := range p.workers {
if status.Worker == worker {
if status.State != WorkerWorking {
panic("returned worker was already free")
}
status.State = WorkerFree
return
}
}
}