ldk_event_broadcaster.go raw

   1  package ldk
   2  
   3  import (
   4  	"context"
   5  	"slices"
   6  	"time"
   7  
   8  	"github.com/getAlby/ldk-node-go/ldk_node"
   9  	// "github.com/getAlby/hub/ldk_node"
  10  	"github.com/getAlby/hub/logger"
  11  
  12  	"github.com/sirupsen/logrus"
  13  )
  14  
  15  /*
  16  *
  17  a LDK event broadcaster powered by channels.
  18  There are 3 main channels:
  19  - 1. receives the next LDK event
  20  - 2. adding a subscriber
  21  - 3. removing a subscriber
  22  Based on https://betterprogramming.pub/how-to-broadcast-messages-in-go-using-channels-b68f42bdf32e
  23  */
  24  type ldkEventBroadcastServer struct {
  25  	source         <-chan *ldk_node.Event
  26  	listeners      []chan *ldk_node.Event
  27  	addListener    chan chan *ldk_node.Event
  28  	removeListener chan (<-chan *ldk_node.Event)
  29  }
  30  
  31  type LDKEventBroadcaster interface {
  32  	Subscribe() chan *ldk_node.Event
  33  	CancelSubscription(chan *ldk_node.Event)
  34  }
  35  
  36  func NewLDKEventBroadcaster(ctx context.Context, source <-chan *ldk_node.Event) LDKEventBroadcaster {
  37  	service := &ldkEventBroadcastServer{
  38  		source:         source,
  39  		listeners:      make([]chan *ldk_node.Event, 0),
  40  		addListener:    make(chan chan *ldk_node.Event),
  41  		removeListener: make(chan (<-chan *ldk_node.Event)),
  42  	}
  43  	go service.serve(ctx)
  44  	return service
  45  }
  46  
  47  func (s *ldkEventBroadcastServer) Subscribe() chan *ldk_node.Event {
  48  	// create a new listener channel and mark it to be added to the listeners array
  49  	newListener := make(chan *ldk_node.Event)
  50  	s.addListener <- newListener
  51  	return newListener
  52  }
  53  
  54  func (s *ldkEventBroadcastServer) CancelSubscription(channel chan *ldk_node.Event) {
  55  	// close the channel - this could fail if the channel was already closed (just ignore it)
  56  	func() {
  57  		defer func() {
  58  			if r := recover(); r != nil {
  59  				logger.Logger.WithField("r", r).Error("Failed to close subscription channel")
  60  			}
  61  		}()
  62  		close(channel)
  63  	}()
  64  	// mark the channel to be removed from the listeners array
  65  	s.removeListener <- channel
  66  }
  67  
  68  func (s *ldkEventBroadcastServer) serve(ctx context.Context) {
  69  	// close down all listeners when the LDK lnclient is shut down
  70  	defer func() {
  71  		for _, listener := range s.listeners {
  72  			func() {
  73  				defer func() {
  74  					if r := recover(); r != nil {
  75  						logger.Logger.WithField("r", r).Error("Failed to close subscription channel")
  76  					}
  77  				}()
  78  				close(listener)
  79  			}()
  80  		}
  81  	}()
  82  
  83  	/**
  84  	Process events from channels in an infinite (blocking) loop.
  85  	*/
  86  	for {
  87  		select {
  88  		case <-ctx.Done():
  89  			// LDK lnclient was shut down
  90  			return
  91  		case newListener := <-s.addListener:
  92  			s.listeners = append(s.listeners, newListener)
  93  		case listenerToRemove := <-s.removeListener:
  94  			for i, listener := range s.listeners {
  95  				// delete the listener from the listeners array
  96  				if listener == listenerToRemove {
  97  					s.listeners[i] = s.listeners[len(s.listeners)-1]
  98  					s.listeners = slices.Delete(s.listeners, len(s.listeners)-1, len(s.listeners))
  99  					break
 100  				}
 101  			}
 102  		case event := <-s.source:
 103  			// got a new LDK event - send it to all listeners
 104  			logger.Logger.WithFields(logrus.Fields{
 105  				"event":         event,
 106  				"listenerCount": len(s.listeners),
 107  			}).Debug("Sending LDK event to listeners")
 108  			for _, listener := range s.listeners {
 109  				func() {
 110  					// if we fail to send the event to the listener it was probably closed
 111  					defer func() {
 112  						if r := recover(); r != nil {
 113  							logger.Logger.WithField("r", r).Error("Failed to send event to listener")
 114  						}
 115  					}()
 116  
 117  					// try to send the event to the listener
 118  					// this can fail if the listener is closed (due to unsubscribing)
 119  					// worst case scenario: it times out because the listener is stuck processing an event
 120  					select {
 121  					case listener <- event:
 122  						logger.Logger.WithFields(logrus.Fields{
 123  							"event": event,
 124  						}).Debug("Sent LDK event to listener")
 125  					case <-time.After(5 * time.Second):
 126  						logger.Logger.WithFields(logrus.Fields{
 127  							"event": event,
 128  						}).Error("Timeout sending LDK event to listener")
 129  					}
 130  				}()
 131  			}
 132  		}
 133  	}
 134  }
 135