From 5df53cd81a1bd5e5381ded026143c71fe3d2fa8e Mon Sep 17 00:00:00 2001 From: sputn1ck Date: Tue, 10 Sep 2024 18:18:46 +0200 Subject: [PATCH] notifications: add notification manager This commit adds a generic notification manager that can be used to subscribe to different types of notifications. --- notifications/log.go | 26 +++++ notifications/manager.go | 212 ++++++++++++++++++++++++++++++++++ notifications/manager_test.go | 173 +++++++++++++++++++++++++++ 3 files changed, 411 insertions(+) create mode 100644 notifications/log.go create mode 100644 notifications/manager.go create mode 100644 notifications/manager_test.go diff --git a/notifications/log.go b/notifications/log.go new file mode 100644 index 00000000..5fb30f10 --- /dev/null +++ b/notifications/log.go @@ -0,0 +1,26 @@ +package notifications + +import ( + "github.com/btcsuite/btclog" + "github.com/lightningnetwork/lnd/build" +) + +// Subsystem defines the sub system name of this package. +const Subsystem = "NTFNS" + +// log is a logger that is initialized with no output filters. This +// means the package will not perform any logging by default until the caller +// requests it. +var log btclog.Logger + +// The default amount of logging is none. +func init() { + UseLogger(build.NewSubLogger(Subsystem, nil)) +} + +// UseLogger uses a specified Logger to output package logging info. +// This should be used in preference to SetLogWriter if the caller is also +// using btclog. +func UseLogger(logger btclog.Logger) { + log = logger +} diff --git a/notifications/manager.go b/notifications/manager.go new file mode 100644 index 00000000..ac6e16a8 --- /dev/null +++ b/notifications/manager.go @@ -0,0 +1,212 @@ +package notifications + +import ( + "context" + "sync" + "time" + + "github.com/lightninglabs/loop/swapserverrpc" + "google.golang.org/grpc" +) + +// NotificationType is the type of notification that the manager can handle. +type NotificationType int + +const ( + // NotificationTypeUnknown is the default notification type. + NotificationTypeUnknown NotificationType = iota + + // NotificationTypeReservation is the notification type for reservation + // notifications. + NotificationTypeReservation +) + +// Client is the interface that the notification manager needs to implement in +// order to be able to subscribe to notifications. +type Client interface { + // SubscribeNotifications subscribes to the notifications from the server. + SubscribeNotifications(ctx context.Context, + in *swapserverrpc.SubscribeNotificationsRequest, + opts ...grpc.CallOption) ( + swapserverrpc.SwapServer_SubscribeNotificationsClient, error) +} + +// Config contains all the services that the notification manager needs to +// operate. +type Config struct { + // Client is the client used to communicate with the swap server. + Client Client + + // FetchL402 is the function used to fetch the l402 token. + FetchL402 func(context.Context) error +} + +// Manager is a manager for notifications that the swap server sends to the +// client. +type Manager struct { + cfg *Config + + hasL402 bool + + subscribers map[NotificationType][]subscriber + sync.Mutex +} + +// NewManager creates a new notification manager. +func NewManager(cfg *Config) *Manager { + return &Manager{ + cfg: cfg, + subscribers: make(map[NotificationType][]subscriber), + } +} + +type subscriber struct { + subCtx context.Context + recvChan interface{} +} + +// SubscribeReservations subscribes to the reservation notifications. +func (m *Manager) SubscribeReservations(ctx context.Context, +) <-chan *swapserverrpc.ServerReservationNotification { + + notifChan := make(chan *swapserverrpc.ServerReservationNotification, 1) + sub := subscriber{ + subCtx: ctx, + recvChan: notifChan, + } + + m.addSubscriber(NotificationTypeReservation, sub) + + // Start a goroutine to remove the subscriber when the context is canceled + go func() { + <-ctx.Done() + m.removeSubscriber(NotificationTypeReservation, sub) + close(notifChan) + }() + + return notifChan +} + +// Run starts the notification manager. It will keep on running until the +// context is canceled. It will subscribe to notifications and forward them to +// the subscribers. On a first successful connection to the server, it will +// close the readyChan to signal that the manager is ready. +func (m *Manager) Run(ctx context.Context) error { + // Initially we want to immediately try to connect to the server. + waitTime := time.Duration(0) + + // Start the notification runloop. + for { + timer := time.NewTimer(waitTime) + // Increase the wait time for the next iteration. + waitTime += time.Second * 1 + + // Return if the context has been canceled. + select { + case <-ctx.Done(): + return nil + + case <-timer.C: + } + + // In order to create a valid l402 we first are going to call + // the FetchL402 method. As a client might not have outbound capacity + // yet, we'll retry until we get a valid response. + if !m.hasL402 { + err := m.cfg.FetchL402(ctx) + if err != nil { + log.Errorf("Error fetching L402: %v", err) + continue + } + m.hasL402 = true + } + + connectedFunc := func() { + // Reset the wait time to 10 seconds. + waitTime = time.Second * 10 + } + + err := m.subscribeNotifications(ctx, connectedFunc) + if err != nil { + log.Errorf("Error subscribing to notifications: %v", err) + } + } +} + +// subscribeNotifications subscribes to the notifications from the server. +func (m *Manager) subscribeNotifications(ctx context.Context, + connectedFunc func()) error { + + callCtx, cancel := context.WithCancel(ctx) + defer cancel() + + notifStream, err := m.cfg.Client.SubscribeNotifications( + callCtx, &swapserverrpc.SubscribeNotificationsRequest{}, + ) + if err != nil { + return err + } + + // Signal that we're connected to the server. + connectedFunc() + log.Debugf("Successfully subscribed to server notifications") + + for { + notification, err := notifStream.Recv() + if err == nil && notification != nil { + log.Debugf("Received notification: %v", notification) + m.handleNotification(notification) + continue + } + + log.Errorf("Error receiving notification: %v", err) + + return err + } +} + +// handleNotification handles an incoming notification from the server, +// forwarding it to the appropriate subscribers. +func (m *Manager) handleNotification(notification *swapserverrpc. + SubscribeNotificationsResponse) { + + switch notification.Notification.(type) { + case *swapserverrpc.SubscribeNotificationsResponse_ReservationNotification: + // We'll forward the reservation notification to all subscribers. + reservationNtfn := notification.GetReservationNotification() + m.Lock() + defer m.Unlock() + + for _, sub := range m.subscribers[NotificationTypeReservation] { + recvChan := sub.recvChan.(chan *swapserverrpc. + ServerReservationNotification) + + recvChan <- reservationNtfn + } + + default: + log.Warnf("Received unknown notification type: %v", + notification) + } +} + +// addSubscriber adds a subscriber to the manager. +func (m *Manager) addSubscriber(notifType NotificationType, sub subscriber) { + m.Lock() + defer m.Unlock() + m.subscribers[notifType] = append(m.subscribers[notifType], sub) +} + +// removeSubscriber removes a subscriber from the manager. +func (m *Manager) removeSubscriber(notifType NotificationType, sub subscriber) { + m.Lock() + defer m.Unlock() + subs := m.subscribers[notifType] + newSubs := make([]subscriber, 0, len(subs)) + for _, s := range subs { + if s != sub { + newSubs = append(newSubs, s) + } + } + m.subscribers[notifType] = newSubs +} diff --git a/notifications/manager_test.go b/notifications/manager_test.go new file mode 100644 index 00000000..ea756ff6 --- /dev/null +++ b/notifications/manager_test.go @@ -0,0 +1,173 @@ +package notifications + +import ( + "context" + "io" + "sync" + "testing" + "time" + + "github.com/lightninglabs/loop/swapserverrpc" + "github.com/stretchr/testify/require" + "google.golang.org/grpc" + "google.golang.org/grpc/metadata" +) + +var ( + testReservationId = []byte{0x01, 0x02} + testReservationId2 = []byte{0x01, 0x02} +) + +// mockNotificationsClient implements the NotificationsClient interface for testing. +type mockNotificationsClient struct { + mockStream swapserverrpc.SwapServer_SubscribeNotificationsClient + subscribeErr error + timesCalled int + sync.Mutex +} + +func (m *mockNotificationsClient) SubscribeNotifications(ctx context.Context, + in *swapserverrpc.SubscribeNotificationsRequest, + opts ...grpc.CallOption) ( + swapserverrpc.SwapServer_SubscribeNotificationsClient, error) { + + m.Lock() + defer m.Unlock() + + m.timesCalled++ + if m.subscribeErr != nil { + return nil, m.subscribeErr + } + return m.mockStream, nil +} + +// mockSubscribeNotificationsClient simulates the server stream. +type mockSubscribeNotificationsClient struct { + grpc.ClientStream + recvChan chan *swapserverrpc.SubscribeNotificationsResponse + recvErrChan chan error +} + +func (m *mockSubscribeNotificationsClient) Recv() ( + *swapserverrpc.SubscribeNotificationsResponse, error) { + + select { + case err := <-m.recvErrChan: + return nil, err + case notif, ok := <-m.recvChan: + if !ok { + return nil, io.EOF + } + return notif, nil + } +} + +func (m *mockSubscribeNotificationsClient) Header() (metadata.MD, error) { + return nil, nil +} + +func (m *mockSubscribeNotificationsClient) Trailer() metadata.MD { + return nil +} + +func (m *mockSubscribeNotificationsClient) CloseSend() error { + return nil +} + +func (m *mockSubscribeNotificationsClient) Context() context.Context { + return context.TODO() +} + +func (m *mockSubscribeNotificationsClient) SendMsg(interface{}) error { + return nil +} + +func (m *mockSubscribeNotificationsClient) RecvMsg(interface{}) error { + return nil +} + +func TestManager_ReservationNotification(t *testing.T) { + // Create a mock notification client + recvChan := make(chan *swapserverrpc.SubscribeNotificationsResponse, 1) + errChan := make(chan error, 1) + mockStream := &mockSubscribeNotificationsClient{ + recvChan: recvChan, + recvErrChan: errChan, + } + mockClient := &mockNotificationsClient{ + mockStream: mockStream, + } + + // Create a Manager with the mock client + mgr := NewManager(&Config{ + Client: mockClient, + FetchL402: func(ctx context.Context) error { + // Simulate successful fetching of L402 + return nil + }, + }) + + // Subscribe to reservation notifications. + subCtx, subCancel := context.WithCancel(context.Background()) + subChan := mgr.SubscribeReservations(subCtx) + + // Run the manager. + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + go func() { + err := mgr.Run(ctx) + require.NoError(t, err) + }() + + // Wait a bit to ensure manager is running and has subscribed + require.Eventually(t, func() bool { + mgr.Lock() + defer mgr.Unlock() + return len(mgr.subscribers[NotificationTypeReservation]) > 0 + }, time.Second*5, 10*time.Millisecond) + + mockClient.Lock() + require.Equal(t, 1, mockClient.timesCalled) + mockClient.Unlock() + + // Send a test notification + testNotif := getTestNotification(testReservationId) + + // Send the notification to the recvChan + recvChan <- testNotif + + // Collect the notification in the callback + receivedNotification := <-subChan + + // Now, check that the notification received in the callback matches the one sent + require.NotNil(t, receivedNotification) + require.Equal(t, testReservationId, receivedNotification.ReservationId) + + // Cancel the subscription + subCancel() + + // Send another test notification` + testNotif2 := getTestNotification(testReservationId2) + recvChan <- testNotif2 + + // Check that the subChan is eventually closed. + require.Eventually(t, func() bool { + select { + case _, ok := <-subChan: + return !ok + default: + return false + } + }, time.Second*5, 10*time.Millisecond) +} + +func getTestNotification(resId []byte) *swapserverrpc.SubscribeNotificationsResponse { + return &swapserverrpc.SubscribeNotificationsResponse{ + Notification: &swapserverrpc.SubscribeNotificationsResponse_ReservationNotification{ + ReservationNotification: &swapserverrpc.ServerReservationNotification{ + ReservationId: resId, + }, + }, + } +}