From 71c9f9da45e7c8a74eedf7a480c43ef783b01813 Mon Sep 17 00:00:00 2001 From: Martin Saposnic Date: Mon, 2 Mar 2026 19:12:53 -0300 Subject: [PATCH] fix: LSPS4 HTLC forwarding race condition and retry mechanism forward_intercepted_htlc was returning Ok when the channel was not live (peer_disconnected flag set), causing LSPS4 to remove the HTLC from the store. LDK's send_htlc would then reject with "Cannot send while disconnected", silently losing the HTLC. Additionally, peer_connected fires before ChannelReestablish completes, so channels are never is_usable at that point. Combined with a 10s HTLC expiry, payments arriving while the peer was offline had almost no chance of being forwarded in time. Changes: - forward_intercepted_htlc: return Err when channel is not live - execute_htlc_actions: only remove HTLC from store on Ok; on Err, keep/re-insert for retry - peer_connected: skip forwarding when no channels are is_usable - New try_forward_pending_htlcs(): periodic retry for stored HTLCs, called from the 5s poll loop - New get_all_peer_node_ids() on HTLCStore - Bump HTLC_EXPIRY_THRESHOLD_SECS from 10 to 30 --- lightning-liquidity/src/lsps4/htlc_store.rs | 10 ++ lightning-liquidity/src/lsps4/service.rs | 109 ++++++++++++++++---- lightning/src/ln/channelmanager.rs | 9 +- 3 files changed, 102 insertions(+), 26 deletions(-) diff --git a/lightning-liquidity/src/lsps4/htlc_store.rs b/lightning-liquidity/src/lsps4/htlc_store.rs index 2f7963832b0..9d2790cfc8b 100644 --- a/lightning-liquidity/src/lsps4/htlc_store.rs +++ b/lightning-liquidity/src/lsps4/htlc_store.rs @@ -192,6 +192,16 @@ where L::Target: Logger, KV::Target: KVStoreSync { htlcs } + /// Get all unique peer node IDs that have stored HTLCs. + pub fn get_all_peer_node_ids(&self) -> Vec { + let locked = self.htlcs.lock().unwrap(); + let mut seen = std::collections::HashSet::new(); + for htlc in locked.values() { + seen.insert(htlc.next_node_id()); + } + seen.into_iter().collect() + } + /// Get all intercepted HTLCs that are older than the specified threshold in seconds. pub fn get_expired_htlcs(&self, now: u64, threshold_seconds: u64) -> Vec { self.list_filter(|htlc| { diff --git a/lightning-liquidity/src/lsps4/service.rs b/lightning-liquidity/src/lsps4/service.rs index 5df00c56f76..0f3c147de83 100644 --- a/lightning-liquidity/src/lsps4/service.rs +++ b/lightning-liquidity/src/lsps4/service.rs @@ -49,7 +49,7 @@ use std::string::String; use std::time::Instant; use std::vec::Vec; -const HTLC_EXPIRY_THRESHOLD_SECS: u64 = 10; +const HTLC_EXPIRY_THRESHOLD_SECS: u64 = 30; /// Action to forward a specific HTLC through a channel #[derive(Debug)] @@ -370,7 +370,20 @@ where } } - self.process_htlcs_for_peer(counterparty_node_id.clone(), htlcs); + // Only attempt forwarding if at least one channel is live (is_usable). + // At peer_connected time, ChannelReestablish hasn't completed yet, + // so channels are typically not usable. The periodic retry will pick + // these up once channels become live. + let has_usable_channel = channels.iter().any(|ch| ch.is_usable); + if has_usable_channel { + self.process_htlcs_for_peer(counterparty_node_id.clone(), htlcs); + } else if htlc_count > 0 { + log_info!( + self.logger, + "[LSPS4] peer_connected: skipping HTLC forwarding for {} - no usable channels yet. Periodic retry will handle it.", + counterparty_node_id + ); + } log_info!( self.logger, @@ -426,6 +439,47 @@ where } } + /// Try to forward all pending stored HTLCs for peers that are connected and have + /// usable channels. Called periodically from the poll loop to catch the window + /// after ChannelReestablish completes (when peer_connected couldn't forward). + pub fn try_forward_pending_htlcs(&self) { + let peer_ids = self.htlc_store.get_all_peer_node_ids(); + if peer_ids.is_empty() { + return; + } + + let connected_peers = self.connected_peers.read().unwrap().clone(); + + for peer_id in peer_ids { + if !connected_peers.contains(&peer_id) { + continue; + } + + let channels = self + .channel_manager + .get_cm() + .list_channels_with_counterparty(&peer_id); + let has_usable_channel = channels.iter().any(|ch| ch.is_usable); + if !has_usable_channel { + continue; + } + + let htlcs = self.htlc_store.get_htlcs_by_node_id(&peer_id); + if htlcs.is_empty() { + continue; + } + + log_info!( + self.logger, + "[LSPS4] try_forward_pending_htlcs: forwarding {} HTLCs for peer {} (channel now live)", + htlcs.len(), + peer_id + ); + + self.process_htlcs_for_peer(peer_id, htlcs); + } + } + fn is_peer_connected(&self, counterparty_node_id: &PublicKey) -> bool { self.connected_peers.read().unwrap().contains(counterparty_node_id) } @@ -705,38 +759,49 @@ where Ok(()) => { log_info!( self.logger, - "[LSPS4] execute_htlc_actions: forward_intercepted_htlc returned Ok for HTLC {:?} (took {}ms). NOTE: Ok does NOT guarantee delivery - channel may still be reestablishing.", + "[LSPS4] execute_htlc_actions: forward_intercepted_htlc returned Ok for HTLC {:?} (took {}ms)", forward_action.htlc.id(), fwd_start.elapsed().as_millis() ); + + // Only remove from store on successful forward + match self.htlc_store.remove(&forward_action.htlc.id()) { + Ok(()) => { + log_info!( + self.logger, + "[LSPS4] execute_htlc_actions: removed HTLC {:?} from store after forward", + forward_action.htlc.id() + ); + }, + Err(e) => { + log_info!( + self.logger, + "[LSPS4] execute_htlc_actions: HTLC {:?} was NOT in store (expected if from htlc_intercepted connected path): {}", + forward_action.htlc.id(), + e + ); + }, + } }, Err(ref e) => { log_error!( self.logger, - "[LSPS4] execute_htlc_actions: forward_intercepted_htlc FAILED for HTLC {:?} - error: {:?} (took {}ms)", + "[LSPS4] execute_htlc_actions: forward FAILED for HTLC {:?} - error: {:?} (took {}ms). Kept in store for retry.", forward_action.htlc.id(), e, fwd_start.elapsed().as_millis() ); - }, - } - // Remove from store - log whether it was actually present - match self.htlc_store.remove(&forward_action.htlc.id()) { - Ok(()) => { - log_info!( - self.logger, - "[LSPS4] execute_htlc_actions: removed HTLC {:?} from store after forward", - forward_action.htlc.id() - ); - }, - Err(e) => { - log_info!( - self.logger, - "[LSPS4] execute_htlc_actions: HTLC {:?} was NOT in store (expected if from htlc_intercepted connected path): {}", - forward_action.htlc.id(), - e - ); + // On error, ensure HTLC stays in store for retry. + // insert() is a no-op if already present. + if let Err(e) = self.htlc_store.insert(forward_action.htlc.clone()) { + log_error!( + self.logger, + "[LSPS4] execute_htlc_actions: failed to re-insert HTLC {:?} into store: {}", + forward_action.htlc.id(), + e + ); + } }, } } diff --git a/lightning/src/ln/channelmanager.rs b/lightning/src/ln/channelmanager.rs index da91b931ce5..f1c33bc0d9c 100644 --- a/lightning/src/ln/channelmanager.rs +++ b/lightning/src/ln/channelmanager.rs @@ -6746,10 +6746,11 @@ where }); } if !is_live { - log_info!(entry_logger, - "forward_intercepted_htlc: WARNING - channel {} is_usable but NOT is_live (peer_disconnected flag set). HTLC will be queued but will fail at send_htlc time.", - next_hop_channel_id - ); + return Err(APIError::ChannelUnavailable { + err: format!( + "Channel with id {next_hop_channel_id} is not live (peer_disconnected flag set)" + ), + }); } funded_chan.context.outbound_scid_alias() } else {