-
-
Notifications
You must be signed in to change notification settings - Fork 27
/
broker.go
87 lines (68 loc) · 1.69 KB
/
broker.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
package services
import (
"sync"
"github.com/stjudewashere/seonaut/internal/models"
"github.com/google/uuid"
)
type (
// PubSub subscriber struct.
subscriber struct {
Id uuid.UUID
Topic string
Callback func(*models.Message) error
}
// PubSub broker service struct keeps a map of subscribers.
Broker struct {
subscribers map[string][]*subscriber
lock *sync.RWMutex
}
)
func NewPubSubBroker() *Broker {
return &Broker{
subscribers: make(map[string][]*subscriber, 0),
lock: &sync.RWMutex{},
}
}
// Returns a new subsciber to the topic.
func (b *Broker) NewSubscriber(topic string, c func(*models.Message) error) *subscriber {
b.lock.Lock()
defer b.lock.Unlock()
s := &subscriber{
Id: uuid.New(),
Topic: topic,
Callback: c,
}
b.subscribers[topic] = append(b.subscribers[topic], s)
return s
}
// Unsubscribes a subscriber.
func (b *Broker) Unsubscribe(s *subscriber) {
b.lock.Lock()
defer b.lock.Unlock()
subscribers := b.subscribers[s.Topic]
for i, v := range subscribers {
if v.Id == s.Id {
b.subscribers[s.Topic] = append(subscribers[:i], subscribers[i+1:]...)
// The topic is removed once there are no more subscribers.
if len(b.subscribers[s.Topic]) == 0 {
delete(b.subscribers, s.Topic)
}
break
}
}
}
// Publishes a message to all subscribers of a topic.
func (b *Broker) Publish(topic string, m *models.Message) {
b.lock.Lock()
defer b.lock.Unlock()
subscribers := b.subscribers[topic]
for i, v := range subscribers {
err := v.Callback(m)
if err != nil {
b.subscribers[topic] = append(subscribers[:i], subscribers[i+1:]...)
}
}
if len(b.subscribers[topic]) == 0 {
delete(b.subscribers, topic)
}
}