-
Notifications
You must be signed in to change notification settings - Fork 0
/
subscriber.go
86 lines (73 loc) · 1.8 KB
/
subscriber.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
package broadcaster
import (
"sync"
)
// Subscriber is is registered to a broadcast agent and will receive messages
// through a channel
type Subscriber struct {
Comm chan *Broadcast
ErrChan chan error
workersAmmount uint32
frequency string
handler SubscriberHandler
wg sync.WaitGroup
agentWg *sync.WaitGroup
agentWgBusy *sync.WaitGroup
}
func newSubscriber(frequency string, handler SubscriberHandler, workersAmmount uint32, wgBusy, wg *sync.WaitGroup) *Subscriber {
// ErrChan needs to be buffered to be non blocking in case it's not used
errChanBufferSize := 1
if workersAmmount > 0 {
errChanBufferSize = int(workersAmmount)
}
return &Subscriber{
Comm: make(chan *Broadcast, workersAmmount),
ErrChan: make(chan error, errChanBufferSize),
workersAmmount: workersAmmount,
frequency: frequency,
handler: handler,
agentWg: wg,
agentWgBusy: wgBusy,
wg: sync.WaitGroup{},
}
}
func (sub *Subscriber) stop() {
close(sub.Comm)
sub.wg.Wait()
close(sub.ErrChan)
}
func (sub *Subscriber) startWorker() {
defer sub.agentWg.Done()
defer sub.wg.Done()
for {
broadcast := <-sub.Comm
if broadcast == nil {
return
}
sub.agentWgBusy.Add(1)
if err := sub.handler(broadcast); err != nil {
// for non blocking error buffering in the chan
// avoids deadlocks if ErrChan is never read
select {
case sub.ErrChan <- err:
break
}
}
sub.agentWgBusy.Done()
}
}
// Start will start listening to the primary Comm channel for message
// broadcasted on its agent
func (sub *Subscriber) Start() {
if sub.workersAmmount == 0 {
sub.agentWg.Add(1)
sub.wg.Add(1)
go sub.startWorker()
return
}
for i := uint32(0); i < sub.workersAmmount; i++ {
sub.agentWg.Add(1)
sub.wg.Add(1)
go sub.startWorker()
}
}