Search Apps Documentation Source Content File Folder Download Copy Actions Download State String Boolean Number Struct Map Slice Pointer Function Closure Reference Nil Package Type Interface Unknown

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}