events.go raw

   1  package events
   2  
   3  import (
   4  	"context"
   5  	"slices"
   6  	"sync"
   7  
   8  	"github.com/getAlby/hub/logger"
   9  	"github.com/getAlby/hub/version"
  10  	"github.com/sirupsen/logrus"
  11  )
  12  
  13  type eventPublisher struct {
  14  	listeners        []EventSubscriber
  15  	subscriberMtx    sync.Mutex
  16  	globalProperties map[string]interface{}
  17  }
  18  
  19  func NewEventPublisher() *eventPublisher {
  20  	eventPublisher := &eventPublisher{
  21  		listeners:        []EventSubscriber{},
  22  		globalProperties: map[string]interface{}{},
  23  	}
  24  	eventPublisher.SetGlobalProperty("version", version.Tag)
  25  	return eventPublisher
  26  }
  27  
  28  func (ep *eventPublisher) RegisterSubscriber(listener EventSubscriber) {
  29  	ep.subscriberMtx.Lock()
  30  	defer ep.subscriberMtx.Unlock()
  31  	ep.listeners = append(ep.listeners, listener)
  32  }
  33  
  34  func (ep *eventPublisher) RemoveSubscriber(listenerToRemove EventSubscriber) {
  35  	ep.subscriberMtx.Lock()
  36  	defer ep.subscriberMtx.Unlock()
  37  
  38  	for i, listener := range ep.listeners {
  39  		// delete the listener from the listeners array
  40  		if listener == listenerToRemove {
  41  			ep.listeners[i] = ep.listeners[len(ep.listeners)-1]
  42  			ep.listeners = slices.Delete(ep.listeners, len(ep.listeners)-1, len(ep.listeners))
  43  			break
  44  		}
  45  	}
  46  }
  47  
  48  func (ep *eventPublisher) Publish(event *Event) {
  49  	ep.publish(event, false)
  50  }
  51  func (ep *eventPublisher) PublishSync(event *Event) {
  52  	ep.publish(event, true)
  53  }
  54  
  55  func (ep *eventPublisher) publish(event *Event, sync bool) {
  56  	ep.subscriberMtx.Lock()
  57  	defer ep.subscriberMtx.Unlock()
  58  	logger.Logger.WithFields(logrus.Fields{"event": event, "global": ep.globalProperties}).Debug("Publishing event")
  59  	for _, listener := range ep.listeners {
  60  		if sync {
  61  			listener.ConsumeEvent(context.Background(), event, ep.globalProperties)
  62  		} else {
  63  			// consume event without blocking thread
  64  			go listener.ConsumeEvent(context.Background(), event, ep.globalProperties)
  65  		}
  66  	}
  67  }
  68  
  69  func (ep *eventPublisher) SetGlobalProperty(key string, value interface{}) {
  70  	ep.globalProperties[key] = value
  71  }
  72