Skip to content

Commit fca1614

Browse files
authored
Merge pull request #983 from tnull/2026-07-remove-force-close-peers-after-reconnect
Remove force-close peers after reconnect
2 parents 03812a0 + 5e241f7 commit fca1614

5 files changed

Lines changed: 162 additions & 29 deletions

File tree

src/connection.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,10 @@ where
5353
self.do_connect_peer(node_id, addr).await
5454
}
5555

56+
pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
57+
self.peer_manager.disconnect_by_node_id(node_id);
58+
}
59+
5660
pub(crate) async fn do_connect_peer(
5761
&self, node_id: PublicKey, addr: SocketAddress,
5862
) -> Result<(), Error> {

src/event.rs

Lines changed: 84 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
3434
use lightning_liquidity::lsps2::utils::compute_opening_fee;
3535
use lightning_types::payment::{PaymentHash, PaymentPreimage};
3636

37-
use crate::config::{may_announce_channel, Config};
37+
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
3838
use crate::connection::ConnectionManager;
3939
use crate::data_store::DataStoreUpdateResult;
4040
use crate::fee_estimator::ConfirmationTarget;
@@ -583,6 +583,68 @@ where
583583
}
584584
}
585585

586+
fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
587+
let channel_manager = Arc::clone(&self.channel_manager);
588+
let connection_manager = Arc::clone(&self.connection_manager);
589+
let peer_store = Arc::clone(&self.peer_store);
590+
let logger = self.logger.clone();
591+
self.runtime.spawn_cancellable_background_task(async move {
592+
let has_other_channels = || {
593+
channel_manager
594+
.list_channels_with_counterparty(&peer_info.node_id)
595+
.iter()
596+
.any(|c| c.channel_id != closed_channel_id)
597+
};
598+
599+
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
600+
return;
601+
}
602+
603+
// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
604+
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
605+
// to retransmit the force-close error before we stop persisting the peer.
606+
connection_manager.disconnect_peer(peer_info.node_id);
607+
608+
loop {
609+
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
610+
return;
611+
}
612+
613+
match connection_manager
614+
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
615+
.await
616+
{
617+
Ok(()) => {
618+
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
619+
{
620+
return;
621+
}
622+
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
623+
log_error!(
624+
logger,
625+
"Failed to remove peer {} from peer store: {}",
626+
peer_info.node_id,
627+
e
628+
);
629+
} else {
630+
return;
631+
}
632+
},
633+
Err(e) => {
634+
log_debug!(
635+
logger,
636+
"Failed to reconnect peer {} before removing from peer store: {}",
637+
peer_info.node_id,
638+
e
639+
);
640+
},
641+
}
642+
643+
tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
644+
}
645+
});
646+
}
647+
586648
async fn fail_claimable_payment(
587649
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
588650
) -> Result<(), ReplayEvent> {
@@ -1627,25 +1689,24 @@ where
16271689
let counterparty_node_id = counterparty_node_id
16281690
.expect("counterparty_node_id is always set since LDK 0.0.117");
16291691

1630-
// Drop the peer once its last channel with us has reached a terminal state
1631-
// that reconnection cannot recover. Every closure reason is terminal except
1632-
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
1633-
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
1692+
// Drop the peer once its last channel with us has reached a terminal state.
1693+
// For `HolderForceClosed`, retain it through one recovery reconnect so that
1694+
// `channel_reestablish` can retransmit the force-close error before cleanup.
16341695
// This also cleans up peers persisted for a channel that closed before funding
16351696
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
16361697
// retried forever.
16371698
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
16381699
// before removing it from its internal list.
1639-
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });
1640-
1641-
if dont_reconnect {
1642-
let has_other_channels = self
1643-
.channel_manager
1644-
.list_channels_with_counterparty(&counterparty_node_id)
1645-
.iter()
1646-
.any(|c| c.channel_id != channel_id);
1647-
1648-
if !has_other_channels {
1700+
let has_other_channels = self
1701+
.channel_manager
1702+
.list_channels_with_counterparty(&counterparty_node_id)
1703+
.iter()
1704+
.any(|c| c.channel_id != channel_id);
1705+
1706+
let peer_to_reconnect = if !has_other_channels {
1707+
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
1708+
self.peer_store.get_peer(&counterparty_node_id)
1709+
} else {
16491710
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
16501711
log_error!(
16511712
self.logger,
@@ -1655,8 +1716,11 @@ where
16551716
);
16561717
return Err(ReplayEvent());
16571718
}
1719+
None
16581720
}
1659-
}
1721+
} else {
1722+
None
1723+
};
16601724

16611725
let event = Event::ChannelClosed {
16621726
channel_id,
@@ -1672,6 +1736,10 @@ where
16721736
return Err(ReplayEvent());
16731737
},
16741738
};
1739+
1740+
if let Some(peer_info) = peer_to_reconnect {
1741+
self.remove_peer_after_reconnect(peer_info, channel_id);
1742+
}
16751743
},
16761744
LdkEvent::DiscardFunding { channel_id, funding_info } => {
16771745
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {

src/lib.rs

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -2030,12 +2030,10 @@ impl Node {
20302030
}
20312031

20322032
// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
2033-
// which drops the peer once its last channel reaches a terminal state that
2034-
// reconnection cannot recover. We intentionally do nothing here so that a
2035-
// force-closed peer is retained, letting the background reconnection task keep
2036-
// firing and drive the `channel_reestablish` recovery flow. This is especially
2037-
// important against LND peers, which don't always handle force-closure error
2038-
// messages correctly.
2033+
// which retains a force-closed peer through one recovery reconnect before
2034+
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
2035+
// especially important against LND peers that don't always handle force-closure
2036+
// error messages correctly.
20392037
}
20402038

20412039
Ok(())

src/peer_store.rs

Lines changed: 66 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -58,11 +58,14 @@ where
5858
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
5959
let _guard = self.mutation_lock.lock().await;
6060
let data = {
61-
let mut locked_peers = self.peers.write().expect("lock");
62-
locked_peers.remove(node_id);
63-
PeerStoreSerWrapper(&locked_peers).encode()
61+
let locked_peers = self.peers.read().expect("lock");
62+
let mut updated_peers = locked_peers.clone();
63+
updated_peers.remove(node_id);
64+
PeerStoreSerWrapper(&updated_peers).encode()
6465
};
65-
self.persist_peers(data).await
66+
self.persist_peers(data).await?;
67+
self.peers.write().expect("lock").remove(node_id);
68+
Ok(())
6669
}
6770

6871
/// Returns the current in-memory peer set.
@@ -170,12 +173,52 @@ mod tests {
170173
use std::str::FromStr;
171174
use std::sync::Arc;
172175

176+
use bitcoin::io;
177+
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
173178
use lightning::util::test_utils::TestLogger;
174179

175180
use super::*;
176181
use crate::io::test_utils::InMemoryStore;
177182
use crate::types::DynStoreWrapper;
178183

184+
struct FailingStore;
185+
186+
impl KVStore for FailingStore {
187+
fn read(
188+
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
189+
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
190+
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
191+
}
192+
193+
fn write(
194+
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
195+
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
196+
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
197+
}
198+
199+
fn remove(
200+
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
201+
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
202+
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
203+
}
204+
205+
fn list(
206+
&self, _primary_namespace: &str, _secondary_namespace: &str,
207+
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
208+
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
209+
}
210+
}
211+
212+
impl PaginatedKVStore for FailingStore {
213+
fn list_paginated(
214+
&self, _primary_namespace: &str, _secondary_namespace: &str,
215+
_page_token: Option<PageToken>,
216+
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
217+
{
218+
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
219+
}
220+
}
221+
179222
#[tokio::test]
180223
async fn peer_info_persistence() {
181224
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
@@ -215,4 +258,23 @@ mod tests {
215258
assert_eq!(peers[0], expected_peer_info);
216259
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
217260
}
261+
262+
#[tokio::test]
263+
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
264+
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
265+
let logger = Arc::new(TestLogger::new());
266+
let node_id = PublicKey::from_str(
267+
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
268+
)
269+
.unwrap();
270+
let peer_info =
271+
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
272+
let mut peers = HashMap::new();
273+
peers.insert(node_id, peer_info.clone());
274+
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
275+
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();
276+
277+
assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
278+
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
279+
}
218280
}

tests/common/mod.rs

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1699,10 +1699,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
16991699
}
17001700

17011701
if force_close {
1702-
// Peer retained after local force-close to allow channel_reestablish recovery.
1702+
// The recovery reconnect completed while the force-close settled, so the peer no longer
1703+
// needs to remain persisted.
17031704
assert!(
1704-
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
1705-
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
1705+
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
1706+
"node_b should be removed from node_a peer store after the recovery reconnect"
17061707
);
17071708
assert_all_nodes_have_onchain_tx_type(
17081709
&[("node_a", &node_a), ("node_b", &node_b)],

0 commit comments

Comments
 (0)