diff --git a/src/builder.rs b/src/builder.rs index 1158044e47..3d7af99b51 100644 --- a/src/builder.rs +++ b/src/builder.rs @@ -522,7 +522,7 @@ impl NodeBuilder { /// 0-confirmation channels opened by this LSP. If `false`, 0-confirmation /// acceptance for this peer falls back to [`Config::trusted_peers_0conf`]. /// - /// May be called multiple times to register several LSPs. Duplicate `node_id`s are ignored. + /// May be called multiple times to register several LSPs. Re-adding an existing `node_id` updates its address, token, and 0conf trust settings. /// /// [bLIP-50 / LSPS0]: https://github.com/lightning/blips/blob/master/blip-0050.md pub fn add_liquidity_source( @@ -532,7 +532,12 @@ impl NodeBuilder { let liquidity_source_config = self.liquidity_source_config.get_or_insert(LiquiditySourceConfig::default()); - if liquidity_source_config.lsp_nodes.iter().any(|n| n.node_id == node_id) { + if let Some(existing) = + liquidity_source_config.lsp_nodes.iter_mut().find(|n| n.node_id == node_id) + { + existing.address = address; + existing.token = token; + existing.trust_peer_0conf = trust_peer_0conf; return self; } @@ -1176,7 +1181,7 @@ impl Builder { /// 0-confirmation channels opened by this LSP. If `false`, 0-confirmation /// acceptance for this peer falls back to [`Config::trusted_peers_0conf`]. /// - /// May be called multiple times to register several LSPs. Duplicate `node_id`s are ignored. + /// May be called multiple times to register several LSPs. Re-adding an existing `node_id` updates its address, token, and 0conf trust settings. /// /// [bLIP-50 / LSPS0]: https://github.com/lightning/blips/blob/master/blip-0050.md pub fn add_liquidity_source( diff --git a/src/liquidity/mod.rs b/src/liquidity/mod.rs index ffc1f878bc..651081cabc 100644 --- a/src/liquidity/mod.rs +++ b/src/liquidity/mod.rs @@ -154,33 +154,56 @@ impl Liquidity { /// The given `token` will be used by the LSP to authenticate the user. /// `trust_peer_0conf` controls whether the node will accept 0-confirmation channels opened by this /// LSP. Note this supersedes [`Config::trusted_peers_0conf`] for this peer. - /// Duplicate `node_id`s are ignored. + /// Re-adding an existing `node_id` updates its address/token/0conf settings and reconnects. pub fn add_liquidity_source( &self, node_id: PublicKey, address: SocketAddress, token: Option, trust_peer_0conf: bool, ) -> Result<(), Error> { + let mut previous: Option = None; { let mut lsp_nodes = self.liquidity_source.lsp_nodes.write().expect("lock"); - if lsp_nodes.iter().any(|n| n.node_id == node_id) { - log_info!(self.logger, "LSP node {} already added, skipping.", node_id); - return Ok(()); + if let Some(existing) = lsp_nodes.iter_mut().find(|n| n.node_id == node_id) { + if existing.address == address + && existing.token == token + && existing.trust_peer_0conf == trust_peer_0conf + { + log_info!(self.logger, "LSP node {} already added, skipping.", node_id); + return Ok(()); + } + log_info!( + self.logger, + "Updating existing LSP node {} address/config and reconnecting.", + node_id + ); + previous = Some(existing.clone()); + existing.address = address.clone(); + existing.token = token.clone(); + existing.trust_peer_0conf = trust_peer_0conf; + // Force rediscovery after config/address change. + existing.supported_protocols = None; + } else { + lsp_nodes.push(LspNode { + node_id, + address: address.clone(), + token: token.clone(), + trust_peer_0conf, + supported_protocols: None, + }); } - - lsp_nodes.push(LspNode { - node_id, - address: address.clone(), - token: token.clone(), - trust_peer_0conf, - supported_protocols: None, - }); } - // If anything below fails, drop the half-initialized entry so the user can retry cleanly. + // On failure: remove half-initialized inserts; restore prior config for updates. let lsp_nodes = Arc::clone(&self.liquidity_source.lsp_nodes); let cleanup = move || { - lsp_nodes.write().expect("lock").retain(|n| n.node_id != node_id); + let mut nodes = lsp_nodes.write().expect("lock"); + if let Some(prev) = previous { + if let Some(existing) = nodes.iter_mut().find(|n| n.node_id == node_id) { + *existing = prev; + } + } else { + nodes.retain(|n| n.node_id != node_id); + } }; - let con_cm = Arc::clone(&self.connection_manager); let connect_addr = address.clone(); if let Err(e) = self @@ -225,6 +248,7 @@ pub(crate) struct LspConfig { pub trust_peer_0conf: bool, } +#[derive(Clone)] pub(crate) struct LspNode { node_id: PublicKey, address: SocketAddress, diff --git a/src/peer_store.rs b/src/peer_store.rs index 8345bf7111..b8175409b5 100644 --- a/src/peer_store.rs +++ b/src/peer_store.rs @@ -42,17 +42,42 @@ where Self { peers, mutation_lock, kv_store, logger } } + /// Inserts or updates a peer entry. + /// + /// If the peer is already known with the same address, this is a no-op. If the peer is new or + /// the stored address changed (e.g. an LSP migrated hosts), the entry is updated and persisted. + /// In-memory state is only mutated after a successful store write, matching [`Self::remove_peer`]. pub(crate) async fn add_peer(&self, peer_info: PeerInfo) -> Result<(), Error> { let _guard = self.mutation_lock.lock().await; let data = { let mut locked_peers = self.peers.write().expect("lock"); - if locked_peers.contains_key(&peer_info.node_id) { + if locked_peers.get(&peer_info.node_id) == Some(&peer_info) { return Ok(()); } - locked_peers.insert(peer_info.node_id, peer_info); - PeerStoreSerWrapper(&locked_peers).encode() + + // Temporarily apply the update so it is included in the + // serialized representation. + let previous = locked_peers.insert(peer_info.node_id, peer_info.clone()); + let data = PeerStoreSerWrapper(&locked_peers).encode(); + + // Restore the original in-memory state before persistence. + match previous { + Some(previous) => { + locked_peers.insert(peer_info.node_id, previous); + }, + None => { + locked_peers.remove(&peer_info.node_id); + }, + } + + data }; - self.persist_peers(data).await + + self.persist_peers(data).await?; + + // Persistence succeeded, so apply the update permanently. + self.peers.write().expect("lock").insert(peer_info.node_id, peer_info); + Ok(()) } pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> { @@ -277,4 +302,93 @@ mod tests { assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await); assert_eq!(Some(peer_info), peer_store.get_peer(&node_id)); } + + #[tokio::test] + async fn peer_address_updated_on_readd() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let logger = Arc::new(TestLogger::new()); + let peer_store = PeerStore::new(Arc::clone(&store), Arc::clone(&logger)); + + let node_id = PublicKey::from_str( + "0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993", + ) + .unwrap(); + let old_address = SocketAddress::from_str("127.0.0.1:9738").unwrap(); + let new_address = SocketAddress::from_str("127.0.0.1:9739").unwrap(); + + peer_store.add_peer(PeerInfo { node_id, address: old_address.clone() }).await.unwrap(); + assert_eq!(peer_store.get_peer(&node_id), Some(PeerInfo { node_id, address: old_address })); + + // Re-adding the same peer with a new socket address must refresh the stored entry + // (regression for https://github.com/lightningdevkit/ldk-node/issues/700). + let updated = PeerInfo { node_id, address: new_address.clone() }; + peer_store.add_peer(updated.clone()).await.unwrap(); + assert_eq!(peer_store.get_peer(&node_id), Some(updated.clone())); + + let persisted_bytes = KVStore::read( + &*store, + PEER_INFO_PERSISTENCE_PRIMARY_NAMESPACE, + PEER_INFO_PERSISTENCE_SECONDARY_NAMESPACE, + PEER_INFO_PERSISTENCE_KEY, + ) + .await + .unwrap(); + let deser_peer_store = + PeerStore::read(&mut &persisted_bytes[..], (Arc::clone(&store), logger)).unwrap(); + assert_eq!(deser_peer_store.get_peer(&node_id), Some(updated)); + } + + #[tokio::test] + async fn peer_same_address_skips_persist() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let logger = Arc::new(TestLogger::new()); + let peer_store = PeerStore::new(Arc::clone(&store), Arc::clone(&logger)); + + let node_id = PublicKey::from_str( + "0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993", + ) + .unwrap(); + let address = SocketAddress::from_str("127.0.0.1:9738").unwrap(); + let peer_info = PeerInfo { node_id, address }; + + peer_store.add_peer(peer_info.clone()).await.unwrap(); + let first_bytes = KVStore::read( + &*store, + PEER_INFO_PERSISTENCE_PRIMARY_NAMESPACE, + PEER_INFO_PERSISTENCE_SECONDARY_NAMESPACE, + PEER_INFO_PERSISTENCE_KEY, + ) + .await + .unwrap(); + + // Identical re-add is a no-op for the store payload. + peer_store.add_peer(peer_info.clone()).await.unwrap(); + let second_bytes = KVStore::read( + &*store, + PEER_INFO_PERSISTENCE_PRIMARY_NAMESPACE, + PEER_INFO_PERSISTENCE_SECONDARY_NAMESPACE, + PEER_INFO_PERSISTENCE_KEY, + ) + .await + .unwrap(); + assert_eq!(first_bytes, second_bytes); + assert_eq!(peer_store.get_peer(&node_id), Some(peer_info)); + } + + #[tokio::test] + async fn add_peer_does_not_mutate_memory_if_persist_fails() { + let store: Arc = Arc::new(DynStoreWrapper(FailingStore)); + let logger = Arc::new(TestLogger::new()); + let peer_store = PeerStore::new(store, logger); + + let node_id = PublicKey::from_str( + "0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993", + ) + .unwrap(); + let peer_info = + PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() }; + + assert_eq!(Err(Error::PersistenceFailed), peer_store.add_peer(peer_info.clone()).await); + assert_eq!(None, peer_store.get_peer(&node_id)); + } }