/
sender.go
70 lines (62 loc) · 1.5 KB
/
sender.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
package unregistration
import (
"os"
"time"
"code.cloudfoundry.org/clock"
"code.cloudfoundry.org/lager/v3"
"code.cloudfoundry.org/route-emitter/emitter"
"code.cloudfoundry.org/route-emitter/routingtable"
)
type Sender struct {
logger lager.Logger
clock clock.Clock
cache Cache
natsEmitter emitter.NATSEmitter
interval time.Duration
sendCount int
}
func NewSender(
logger lager.Logger,
clock clock.Clock,
cache Cache,
natsEmitter emitter.NATSEmitter,
interval time.Duration,
sendCount int,
) Sender {
return Sender{
logger: logger.Session("unregistration-sender"),
clock: clock,
cache: cache,
natsEmitter: natsEmitter,
interval: interval,
sendCount: sendCount,
}
}
func (s Sender) Run(signals <-chan os.Signal, ready chan<- struct{}) error {
s.logger.Info("starting")
close(ready)
defer s.logger.Info("exiting")
sendTicker := s.clock.NewTicker(s.interval)
defer sendTicker.Stop()
for {
select {
case <-signals:
s.logger.Info("stopping")
return nil
case <-sendTicker.C():
messages := s.cache.List()
if len(messages) > 0 {
s.logger.Debug("messages", lager.Data{"cache": messages})
}
for _, message := range messages {
s.natsEmitter.Emit(routingtable.MessagesToEmit{
UnregistrationMessages: []routingtable.RegistryMessage{message.RegistryMessage},
})
message.SentCount++
if message.SentCount == s.sendCount {
s.cache.Remove([]routingtable.RegistryMessage{message.RegistryMessage})
}
}
}
}
}