forked from freehere107/go-workers
-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathmanager.go
123 lines (104 loc) · 2.05 KB
/
manager.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
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
package workers
import (
"strings"
"sync"
)
type manager struct {
queue string
fetch Fetcher
job jobFunc
concurrency int
workers []*worker
workersM *sync.Mutex
confirm chan *Msg
stop chan bool
exit chan bool
mids *Middlewares
*sync.WaitGroup
}
func (m *manager) start() {
m.Add(1)
m.loadWorkers()
go m.manage()
}
func (m *manager) prepare() {
if !m.fetch.Closed() {
m.fetch.Close()
}
}
func (m *manager) quit() {
Logger.Println("quitting queue", m.queueName(), "(waiting for", m.processing(), "/", len(m.workers), "workers).")
m.prepare()
m.workersM.Lock()
for _, worker := range m.workers {
worker.quit()
}
m.workersM.Unlock()
m.stop <- true
<-m.exit
m.reset()
m.Done()
}
func (m *manager) manage() {
Logger.Println("processing queue", m.queueName(), "with", m.concurrency, "workers.")
go m.fetch.Fetch()
for {
select {
case message := <-m.confirm:
m.fetch.Acknowledge(message)
case <-m.stop:
m.exit <- true
break
}
}
}
func (m *manager) loadWorkers() {
m.workersM.Lock()
for i := 0; i < m.concurrency; i++ {
m.workers[i] = newWorker(m)
m.workers[i].start()
}
m.workersM.Unlock()
}
func (m *manager) processing() (count int) {
m.workersM.Lock()
for _, worker := range m.workers {
if worker.processing() {
count++
}
}
m.workersM.Unlock()
return
}
func (m *manager) queueName() string {
return strings.Replace(m.queue, "queue:", "", 1)
}
func (m *manager) reset() {
m.fetch = Config.Fetch(m.queue)
}
func newManager(queue string, job jobFunc, concurrency int, mids ...Action) *manager {
var customMids *Middlewares
if len(mids) == 0 {
customMids = Middleware
} else {
customMids = NewMiddleware(Middleware.actions...)
for _, m := range mids {
customMids.Append(m)
}
}
m := &manager{
Config.Namespace + "queue:" + queue,
nil,
job,
concurrency,
make([]*worker, concurrency),
&sync.Mutex{},
make(chan *Msg),
make(chan bool),
make(chan bool),
customMids,
&sync.WaitGroup{},
}
m.fetch = Config.Fetch(m.queue)
return m
}