diff --git a/lntest/harness_assertion.go b/lntest/harness_assertion.go index 2540ef6b0..5a598033a 100644 --- a/lntest/harness_assertion.go +++ b/lntest/harness_assertion.go @@ -2878,3 +2878,65 @@ func (h *HarnessTest) AssertForceCloseAndAnchorTxnsInMempool() (*wire.MsgTx, return nil, nil } } + +// ReceiveSendToRouteUpdate waits until a message is received on the +// PeerEventsClient stream or the timeout is reached. +func (h *HarnessTest) ReceivePeerEvent( + stream rpc.PeerEventsClient) (*lnrpc.PeerEvent, error) { + + eventChan := make(chan *lnrpc.PeerEvent, 1) + errChan := make(chan error, 1) + go func() { + // Consume one message. This will block until the message is + // received. + resp, err := stream.Recv() + if err != nil { + errChan <- err + + return + } + eventChan <- resp + }() + + select { + case <-time.After(DefaultTimeout): + require.Fail(h, "timeout", "timeout waiting for peer event") + return nil, nil + + case err := <-errChan: + return nil, err + + case event := <-eventChan: + return event, nil + } +} + +// AssertPeerOnlineEvent reads an event from the PeerEventsClient stream and +// asserts it's an online event. +func (h HarnessTest) AssertPeerOnlineEvent(stream rpc.PeerEventsClient) { + event, err := h.ReceivePeerEvent(stream) + require.NoError(h, err) + + require.Equal(h, lnrpc.PeerEvent_PEER_ONLINE, event.Type) +} + +// AssertPeerOfflineEvent reads an event from the PeerEventsClient stream and +// asserts it's an offline event. +func (h HarnessTest) AssertPeerOfflineEvent(stream rpc.PeerEventsClient) { + event, err := h.ReceivePeerEvent(stream) + require.NoError(h, err) + + require.Equal(h, lnrpc.PeerEvent_PEER_OFFLINE, event.Type) +} + +// AssertPeerReconnected reads two events from the PeerEventsClient stream. The +// first event must be an offline event, and the second event must be an online +// event. This is a typical reconnection scenario, where the peer is +// disconnected then connected again. +// +// NOTE: It's important to make the subscription before the disconnection +// happens, otherwise the events can be missed. +func (h HarnessTest) AssertPeerReconnected(stream rpc.PeerEventsClient) { + h.AssertPeerOfflineEvent(stream) + h.AssertPeerOnlineEvent(stream) +} diff --git a/lntest/rpc/lnd.go b/lntest/rpc/lnd.go index 1a49cd18b..055b3f29b 100644 --- a/lntest/rpc/lnd.go +++ b/lntest/rpc/lnd.go @@ -755,3 +755,19 @@ func (h *HarnessRPC) Quiesce( return res } + +type PeerEventsClient lnrpc.Lightning_SubscribePeerEventsClient + +// SubscribePeerEvents makes a RPC call to the node's SubscribePeerEvents and +// returns the stream client. +func (h *HarnessRPC) SubscribePeerEvents( + req *lnrpc.PeerEventSubscription) PeerEventsClient { + + // SubscribePeerEvents needs to have the context alive for the entire + // test case as the returned client will be used for send and receive + // events stream. Thus we use runCtx here instead of a timeout context. + resp, err := h.LN.SubscribePeerEvents(h.runCtx, req) + h.NoError(err, "SubscribePeerEvents") + + return resp +}