messenger: add retry intents on connect to peer

If the `connect_to_peer` task fails somehow, it won't retry the connection
probably missing the connection to the peer. Therefore, this commit adds
a exponential back off strategy for the retries of `connect_to_peer`.
This commit is contained in:
f3r10 2025-09-18 06:18:33 -05:00
parent 3ddf4aaeb2
commit e47bfa5b14
4 changed files with 150 additions and 52 deletions

View file

@ -255,6 +255,10 @@ pub(super) mod tests {
mock! {
pub TestPeerConnector{}
impl Clone for TestPeerConnector {
fn clone(&self) -> Self;
}
#[async_trait]
impl PeerConnector for TestPeerConnector {
async fn list_peers(&mut self) -> Result<tonic_lnd::lnrpc::ListPeersResponse, Status>;

View file

@ -26,11 +26,13 @@ use lightning::{
sign::EntropySource,
types::payment::PaymentHash,
};
use log::{debug, error, trace};
use log::{debug, error, info, trace};
use tokio::{select, time::sleep};
use tonic_lnd::{
lnrpc::{ChanInfoRequest, GetInfoRequest, Payment},
Client,
};
use triggered::Listener;
use crate::{
lnd::{
@ -43,6 +45,10 @@ use crate::{
use super::{validate_amount, OfferError};
const INITIAL_DELAY_MS: u64 = 500;
const MAX_DELAY_MS: u64 = 60_000;
const MAX_ATTEMPTS: u32 = 10;
#[derive(Debug)]
pub struct LndkBolt12InvoiceInfo {
pub payment_hash: PaymentHash,
@ -375,6 +381,47 @@ pub async fn send_invoice_request(
Ok((contents, send_instructions))
}
pub(crate) async fn connect_to_peer_with_retry(
connector: impl PeerConnector + Clone,
node_id: PublicKey,
shutdown_listener: Listener,
) -> Result<(), OfferError> {
let mut attempts = 0;
loop {
select! {
biased;
_ = shutdown_listener.clone() => {
info!("Received shutdown signal, exiting connect to peer loop.");
return Ok(())
}
_ = async {
if attempts > 0 {
let mut delay = INITIAL_DELAY_MS * (2u64.pow(attempts - 1));
delay = delay.min(MAX_DELAY_MS);
info!("Attempt {attempts} to connect to peer failed. Retrying in {}ms...", delay);
sleep(Duration::from_millis(delay)).await;
}
} => {
match connect_to_peer(connector.clone(), node_id).await {
Ok(_) => {
debug!("Connect to peer with node_id {node_id} successful.");
return Ok(())
},
Err(e) => {
error!("Connect to peer produced an error: {e}.");
attempts += 1;
if attempts > MAX_ATTEMPTS {
error!("Max retry attempts reached for peer {node_id}. Giving up.");
return Err(e);
}
}
}
}
}
}
}
pub(crate) async fn connect_to_peer(
mut connector: impl PeerConnector,
node_id: PublicKey,

View file

@ -14,7 +14,7 @@ pub mod handler;
mod lnd_requests;
mod parse;
pub(crate) use lnd_requests::connect_to_peer;
pub(crate) use lnd_requests::connect_to_peer_with_retry;
pub use lnd_requests::create_reply_path;
pub use parse::{decode, get_destination, validate_amount};

View file

@ -1,7 +1,7 @@
use crate::clock::TokioClock;
use crate::grpc::Retryable;
use crate::lnd::{features_support_onion_messages, PeerConnector, ONION_MESSAGES_OPTIONAL};
use crate::offers::connect_to_peer;
use crate::offers::connect_to_peer_with_retry;
use crate::rate_limit::{RateLimiter, RateLimiterCfg, TokenLimiter};
use crate::{LifecycleSignals, LndkOnionMessenger, LDK_LOGGER_NAME};
use async_trait::async_trait;
@ -343,6 +343,7 @@ impl LndkOnionMessenger {
};
let event_handler = LndkEventHandler {
lnd_client: ln_client.clone(),
shutdown_listener: signals.listener.clone(),
};
let consume_result = consume_messenger_events(
onion_messenger,
@ -1032,6 +1033,7 @@ async fn relay_outgoing_msg_event(
struct LndkEventHandler<T: PeerConnector + Send + 'static + Clone> {
lnd_client: T,
shutdown_listener: Listener,
}
impl<T: PeerConnector + Send + 'static + Clone> EventHandler for LndkEventHandler<T> {
@ -1043,11 +1045,16 @@ impl<T: PeerConnector + Send + 'static + Clone> EventHandler for LndkEventHandle
} => {
debug!("ConnectionNeeded event received for node: {}", node_id);
let lnd_client = self.lnd_client.clone();
let shutdown = self.shutdown_listener.clone();
// TODO: we probably want to retry this connect_to_peer call if it fails
tokio::spawn(async move {
if let Err(e) = connect_to_peer(lnd_client, node_id).await {
error!("Failed to connect to peer: {}", e);
match connect_to_peer_with_retry(lnd_client, node_id, shutdown).await {
Ok(_) => {
debug!("Connect to peer with node_id {node_id} successful.");
}
Err(e) => {
error!("Failed to connect to peer: {}", e);
}
}
});
Ok(())
@ -1374,6 +1381,7 @@ mod tests {
drop(sender);
let (_, listener) = triggered::trigger();
let consume_resp = consume_messenger_events(
mock,
receiver,
@ -1381,6 +1389,7 @@ mod tests {
&mut rate_limiter,
LndkEventHandler {
lnd_client: MockTestPeerConnector::new(),
shutdown_listener: listener.clone(),
},
Network::Regtest,
)
@ -1406,6 +1415,7 @@ mod tests {
let mut sender_mock = MockSendCustomMessenger::new();
let (_, listener) = triggered::trigger();
let consume_err = consume_messenger_events(
mock,
receiver,
@ -1413,6 +1423,7 @@ mod tests {
&mut rate_limiter,
LndkEventHandler {
lnd_client: MockTestPeerConnector::new(),
shutdown_listener: listener.clone(),
},
Network::Regtest,
)
@ -1429,6 +1440,7 @@ mod tests {
drop(sender_done);
let mut sender_mock = MockSendCustomMessenger::new();
let mut rate_limiter = MockRateLimiter::new();
let (_, listener) = triggered::trigger();
assert!(consume_messenger_events(
MockOnionHandler::new(),
@ -1436,7 +1448,8 @@ mod tests {
&mut sender_mock,
&mut rate_limiter,
LndkEventHandler {
lnd_client: MockTestPeerConnector::new()
lnd_client: MockTestPeerConnector::new(),
shutdown_listener: listener.clone()
},
Network::Regtest,
)
@ -1455,34 +1468,47 @@ mod tests {
connector_mock.expect_clone().times(1).returning(move || {
let notify = notify_clone.clone();
let mut mock = MockTestPeerConnector::new();
mock.expect_clone().times(1).returning(move || {
let notify = notify.clone();
let mut mock_clone = MockTestPeerConnector::new();
mock_clone
.expect_list_peers()
.times(..)
.returning(|| Ok(tonic_lnd::lnrpc::ListPeersResponse { peers: vec![] }));
mock_clone
.expect_get_node_info()
.times(1)
.returning(|_, _| {
let node_addr = tonic_lnd::lnrpc::NodeAddress {
network: String::from("regtest"),
addr: String::from("127.0.0.1:9735"),
};
let node = tonic_lnd::lnrpc::LightningNode {
addresses: vec![node_addr],
..Default::default()
};
Ok(tonic_lnd::lnrpc::NodeInfo {
node: Some(node),
..Default::default()
})
});
mock.expect_list_peers()
.returning(|| Ok(tonic_lnd::lnrpc::ListPeersResponse { peers: vec![] }));
mock.expect_get_node_info().returning(|_, _| {
let node_addr = tonic_lnd::lnrpc::NodeAddress {
network: String::from("regtest"),
addr: String::from("127.0.0.1:9735"),
};
let node = tonic_lnd::lnrpc::LightningNode {
addresses: vec![node_addr],
..Default::default()
};
Ok(tonic_lnd::lnrpc::NodeInfo {
node: Some(node),
..Default::default()
})
});
mock.expect_connect_peer().times(1).returning(move |_, _| {
notify.notify_one(); // Signal completion
Ok(())
mock_clone
.expect_connect_peer()
.times(1)
.returning(move |_, _| {
notify.notify_one(); // Signal completion
Ok(())
});
mock_clone
});
mock
});
let (_, listener) = triggered::trigger();
let handler = LndkEventHandler {
lnd_client: connector_mock,
shutdown_listener: listener.clone(),
};
let event = Event::ConnectionNeeded {
@ -1505,8 +1531,10 @@ mod tests {
#[tokio::test]
async fn test_lndk_event_handler_other_events() {
let connector_mock = MockTestPeerConnector::new();
let (_, listener) = triggered::trigger();
let handler = LndkEventHandler {
lnd_client: connector_mock,
shutdown_listener: listener.clone(),
};
// Test that other events just return ok.
@ -1527,44 +1555,63 @@ mod tests {
}
#[tokio::test]
async fn test_lndk_event_handler_connection_failed() {
async fn test_lndk_event_handler_connection_needed_success_with_retry() {
let mut connector_mock = MockTestPeerConnector::new();
let expected_node_id = pubkey(1);
let notify = Arc::new(tokio::sync::Notify::new());
let notify_clone = notify.clone();
// Define a mutable counter for the mock function.
let mut counter = 0;
connector_mock.expect_clone().times(1).returning(move || {
let notify = notify_clone.clone();
let mut mock = MockTestPeerConnector::new();
mock.expect_list_peers()
.returning(|| Ok(tonic_lnd::lnrpc::ListPeersResponse { peers: vec![] }));
mock.expect_get_node_info().returning(|_, _| {
let node_addr = tonic_lnd::lnrpc::NodeAddress {
network: String::from("regtest"),
addr: String::from("127.0.0.1:9735"),
};
let node = tonic_lnd::lnrpc::LightningNode {
addresses: vec![node_addr],
..Default::default()
};
Ok(tonic_lnd::lnrpc::NodeInfo {
node: Some(node),
..Default::default()
})
});
mock.expect_connect_peer().times(1).returning(move |_, _| {
notify.notify_one();
Err(tonic_lnd::tonic::Status::unavailable("Connection failed"))
mock.expect_clone().times(..).returning(move || {
counter += 1;
let notify = notify.clone();
let mut mock_clone = MockTestPeerConnector::new();
mock_clone
.expect_list_peers()
.times(..)
.returning(|| Ok(tonic_lnd::lnrpc::ListPeersResponse { peers: vec![] }));
mock_clone
.expect_get_node_info()
.times(..)
.returning(|_, _| {
let node_addr = tonic_lnd::lnrpc::NodeAddress {
network: String::from("regtest"),
addr: String::from("127.0.0.1:9735"),
};
let node = tonic_lnd::lnrpc::LightningNode {
addresses: vec![node_addr],
..Default::default()
};
Ok(tonic_lnd::lnrpc::NodeInfo {
node: Some(node),
..Default::default()
})
});
mock_clone
.expect_connect_peer()
.times(..)
.returning(move |_, _| {
if counter <= 5 {
Err(tonic_lnd::tonic::Status::unavailable("Connection failed"))
} else {
notify.notify_one();
Ok(())
}
});
mock_clone
});
mock
});
let (_, listener) = triggered::trigger();
let handler = LndkEventHandler {
lnd_client: connector_mock,
shutdown_listener: listener.clone(),
};
let event = Event::ConnectionNeeded {
@ -1580,7 +1627,7 @@ mod tests {
_ = notify.notified() => {
// Successfully called
}
_ = tokio::time::sleep(tokio::time::Duration::from_millis(200)) => {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(60)) => {
panic!("Spawned connection should have completed within timeout");
}
}