From 177e84f35a83432018fcb243915cc0152918f1dc Mon Sep 17 00:00:00 2001 From: Olaoluwa Osuntokun Date: Tue, 1 Jul 2025 17:58:57 -0700 Subject: [PATCH] funding: fix stream handling issues in consumePendingOpenChannels In this commit, we fix a few issues with the streaming handling: 1. If it got an error on Recv() it would just spin for ever, generating an endless stream of spammy logs. 2. It never listened on the Done() context channel for the stream. So if the `lnd` had actually hung up, it wouldn't detect that and would just keep spinning. To make things more graceful, we also add back off when we try to reconnect. --- funding/manager.go | 179 ++++++++++++++++++++++++++++++++++++++-- funding/manager_test.go | 4 + 2 files changed, 175 insertions(+), 8 deletions(-) diff --git a/funding/manager.go b/funding/manager.go index 917d9ff..f8e01c3 100644 --- a/funding/manager.go +++ b/funding/manager.go @@ -40,6 +40,12 @@ var ( rpcCodeFundingFailed = auctioneerrpc.OrderReject_CHANNEL_FUNDING_FAILED ) +const ( + // maxStreamRecreateAttempts is the maximum number of attempts to + // recreate a channel event stream before giving up. + maxStreamRecreateAttempts = 3 +) + // MatchRejectErr is an error type that is returned from the funding manager if // the trader rejects certain orders instead of the whole batch. type MatchRejectErr struct { @@ -176,10 +182,12 @@ func (m *Manager) Start() error { streamCtx, &lnrpc.ChannelEventSubscription{}, ) if err != nil { + streamCancel() return err } if err := m.pendingOpenChanServer.Start(); err != nil { + streamCancel() return fmt.Errorf("error starting pending chan subscription "+ "server: %v", err) } @@ -189,6 +197,7 @@ func (m *Manager) Start() error { // updates, that's why we are a client to our own server. m.pendingOpenChanClient, err = m.SubscribePendingOpenChan() if err != nil { + streamCancel() return fmt.Errorf("error subscribing to pending open "+ "channel events: %v", err) } @@ -227,6 +236,36 @@ func (m *Manager) Stop() error { return nil } +// streamBackoff manages exponential backoff for stream errors. +type streamBackoff struct { + attempts int + maxDelay time.Duration +} + +// newStreamBackoff creates a new stream backoff manager. +func newStreamBackoff(maxDelay time.Duration) *streamBackoff { + return &streamBackoff{ + maxDelay: maxDelay, + } +} + +// nextDelay returns the next backoff delay based on the number of attempts. +func (s *streamBackoff) nextDelay() time.Duration { + s.attempts++ + + delay := time.Duration(s.attempts) * time.Second + if delay > s.maxDelay { + delay = s.maxDelay + } + + return delay +} + +// reset resets the backoff counter. +func (s *streamBackoff) reset() { + s.attempts = 0 +} + // consumePendingOpenChannels consumes pending open channel events from the // stream and notifies them if the trader currently has an ongoing batch. func (m *Manager) consumePendingOpenChannels( @@ -234,33 +273,68 @@ func (m *Manager) consumePendingOpenChannels( defer m.wg.Done() + currentStream := subStream + streamCtx, streamCancel := context.WithCancel(context.Background()) + defer streamCancel() + + // Initialize backoff manager with 30 second max delay. + backoff := newStreamBackoff(30 * time.Second) + for { select { case <-m.quit: return + case <-currentStream.Context().Done(): + // The stream context was canceled, we need to establish + // a new stream. + log.Warnf("Channel event stream context " + + "canceled, creating new stream") + + newStream := m.recreateChannelEventStream(streamCtx) + if newStream == nil { + return + } + currentStream = newStream + + backoff.reset() + + continue default: } - msg, err := subStream.Recv() + msg, err := currentStream.Recv() if err != nil { - select { - case <-m.quit: + if m.shouldExitOnError(err) { return - default: } - log.Errorf("Unable to read channel event: %v", err) + if m.shouldRecreateStream(err) { + newStream := m.recreateChannelEventStream( + streamCtx, + ) - // If the lnd node shut down, there's no use continuing. - if err == io.EOF || err == io.ErrUnexpectedEOF || - status.Code(err) == codes.Unavailable { + if newStream == nil { + return + } + currentStream = newStream + + backoff.reset() + continue + } + + // For other errors, wait with a backoff and retry with + // the same stream. + if !m.waitOrQuit(backoff.nextDelay()) { return } continue } + // Reset backoff on successful receive. + backoff.reset() + // Skip any events other than the pending open channel one. channel, ok := msg.Channel.(*lnrpc.ChannelEventUpdate_PendingOpenChannel) if !ok { @@ -274,6 +348,95 @@ func (m *Manager) consumePendingOpenChannels( } } +// recreateChannelEventStream attempts to create a new channel event stream. It +// returns nil if the manager is shutting down or if stream creation fails after +// retries. +func (m *Manager) recreateChannelEventStream(ctx context.Context, +) lnrpc.Lightning_SubscribeChannelEventsClient { + + for retries := 0; retries < maxStreamRecreateAttempts; retries++ { + newStream, err := m.cfg.BaseClient.SubscribeChannelEvents( + ctx, &lnrpc.ChannelEventSubscription{}, + ) + if err == nil { + return newStream + } + + log.Errorf("Unable to establish channel event "+ + "stream (attempt %d/%d): %v", + retries+1, maxStreamRecreateAttempts, err) + + // Check if we're shutting down before retrying. + backoff := time.Duration(retries+1) * time.Second + if !m.waitOrQuit(backoff) { + return nil + } + } + + log.Errorf("Failed to re-establish channel event stream "+ + "after %d attempts", maxStreamRecreateAttempts) + + return nil +} + +// shouldExitOnError determines if the given error should cause the consumer to +// exit completely. +func (m *Manager) shouldExitOnError(err error) bool { + log.Errorf("Unable to read channel event: %v", err) + + // Check if we're shutting down first. + select { + case <-m.quit: + return true + default: + } + + // If the lnd node is unavailable, we should exit. + if status.Code(err) == codes.Unavailable { + log.Errorf("lnd node unavailable, stopping " + + "channel event consumption") + return true + } + + return false +} + +// shouldRecreateStream determines if the given error indicates that a new +// stream should be created. +func (m *Manager) shouldRecreateStream(err error) bool { + switch { + case err == io.EOF || err == io.ErrUnexpectedEOF: + log.Infof("Channel event stream ended (EOF), " + + "creating new stream") + return true + + case status.Code(err) == codes.Canceled: + log.Infof("Channel event stream canceled, " + + "creating new stream") + + return true + + case status.Code(err) == codes.DeadlineExceeded: + log.Infof("Channel event stream deadline " + + "exceeded, creating new stream") + + return true + default: + return false + } +} + +// waitOrQuit waits for the specified duration or until the manager is +// shutting down. Returns false if the manager is shutting down. +func (m *Manager) waitOrQuit(duration time.Duration) bool { + select { + case <-time.After(duration): + return true + case <-m.quit: + return false + } +} + // SubscribePendingOpenChan creates a new subscription client to receive events // for pending open channels from lnd. func (m *Manager) SubscribePendingOpenChan() (*subscribe.Client, error) { diff --git a/funding/manager_test.go b/funding/manager_test.go index df2c8a0..bfb4acf 100644 --- a/funding/manager_test.go +++ b/funding/manager_test.go @@ -99,6 +99,10 @@ type channelEventStream struct { quit chan struct{} } +func (c *channelEventStream) Context() context.Context { + return c.ctx +} + func (c *channelEventStream) Recv() (*lnrpc.ChannelEventUpdate, error) { select { case msg := <-c.updateChan: