broker.gno
3.00 Kb · 123 lines
1package message
2
3import (
4 "errors"
5 "strings"
6
7 "gno.land/p/moul/ulist/v0"
8 "gno.land/p/nt/bptree/v0"
9)
10
11var (
12 // ErrInvalidTopic is triggered when an invalid topic is used.
13 ErrInvalidTopic = errors.New("invalid topic")
14
15 // ErrRequiredCallback is triggered when subscribing without a callback.
16 ErrRequiredCallback = errors.New("message callback is required")
17
18 // ErrRequiredSubscriptionID is triggered when unsubscribing without an ID.
19 ErrRequiredSubscriptionID = errors.New("message sibscription ID is required")
20
21 // ErrRequiredTopic is triggered when (un)subscribing without a topic.
22 ErrRequiredTopic = errors.New("message topic is required")
23)
24
25// NewBroker creates a new message broker.
26func NewBroker() *Broker {
27 return &Broker{
28 callbacks: bptree.NewBPTree32(),
29 }
30}
31
32// Broker is a message broker that handles subscriptions and message publishing.
33type Broker struct {
34 callbacks *bptree.BPTree // string(topic) -> *ulist.List(Callback)
35}
36
37// Topics returns the list of current subscription topics.
38func (b *Broker) Topics() []Topic {
39 var topics []Topic
40 b.callbacks.Iterate("", "", func(k string, _ any) bool {
41 topic := Topic(k)
42 if topic == TopicAll {
43 // Skip catchall topic from the list
44 return false
45 }
46
47 topics = append(topics, topic)
48 return false
49 })
50 return topics
51}
52
53// Subscribe subscribes to messages published for a topic.
54// It returns the callback ID within the topic.
55func (b *Broker) Subscribe(topic Topic, cb Callback) (id int, _ error) {
56 key := strings.TrimSpace(string(topic))
57 if key == "" {
58 return 0, ErrRequiredTopic
59 }
60
61 if cb == nil {
62 return 0, ErrRequiredCallback
63 }
64
65 callbacks, _ := b.callbacks.Get(key).(*ulist.List)
66 if callbacks == nil {
67 callbacks = ulist.New()
68 }
69
70 callbacks.Append(cb)
71 b.callbacks.Set(key, callbacks)
72 return callbacks.TotalSize(), nil
73}
74
75// Unsubscribe unsubscribes a callback from a message topic.
76// ID is the callback ID within the topic, returned on subscription.
77func (b *Broker) Unsubscribe(topic Topic, id int) (unsubscribed bool, _ error) {
78 key := strings.TrimSpace(string(topic))
79 if key == "" {
80 return false, ErrRequiredTopic
81 }
82
83 if id == 0 {
84 return false, ErrRequiredSubscriptionID
85 }
86
87 callbacks, ok := b.callbacks.Get(key).(*ulist.List)
88 if !ok {
89 return false, errors.New("message topic not found: " + key)
90 }
91
92 i := id - 1
93 return callbacks.Delete(i) == nil, nil
94}
95
96// Publish publishes a message for a topic.
97func (b *Broker) Publish(topic Topic, data any) error {
98 if topic == TopicAll {
99 return ErrInvalidTopic
100 }
101
102 key := strings.TrimSpace(string(topic))
103 if key == "" {
104 return ErrRequiredTopic
105 }
106
107 iterCb := func(_ int, v any) bool {
108 cb := v.(Callback)
109 cb(Message{topic, data})
110 return false
111 }
112
113 // Trigger callbacks subscribed to current topic
114 if callbacks, ok := b.callbacks.Get(key).(*ulist.List); ok {
115 callbacks.Iterator(0, callbacks.Size(), iterCb)
116 }
117
118 // Trigger callbacks subscribed to all topics
119 if callbacks, ok := b.callbacks.Get(string(TopicAll)).(*ulist.List); ok {
120 callbacks.Iterator(0, callbacks.Size(), iterCb)
121 }
122 return nil
123}