Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 58 additions & 3 deletions src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -525,7 +525,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(
Expand All @@ -535,7 +535,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;
}

Expand Down Expand Up @@ -1185,7 +1190,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(
Expand Down Expand Up @@ -2805,4 +2810,54 @@ mod tests {
builder.set_runtime(multi_thread_runtime.handle().clone()).unwrap();
assert!(builder.setup_runtime(&logger).is_ok(), "a multi-threaded runtime should be used");
}

#[test]
fn add_liquidity_source_updates_existing_lsp_in_place() {
use std::str::FromStr;

use bitcoin::secp256k1::PublicKey;
use lightning::ln::msgs::SocketAddress;

let mut builder = NodeBuilder::new();
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let first_address = SocketAddress::from_str("127.0.0.1:9735").unwrap();
let second_address = SocketAddress::from_str("127.0.0.2:9735").unwrap();

// Register an LSP, then re-add the same node_id with a new address and
// different 0conf/trust settings. Re-adding must update the existing
// entry in place rather than silently dropping the new values.
builder.add_liquidity_source(node_id, first_address.clone(), None, true);
builder.add_liquidity_source(
node_id,
second_address.clone(),
Some("token".to_owned()),
false,
);

let config = builder
.liquidity_source_config
.as_ref()
.expect("liquidity source config should be initialized");
assert_eq!(
config.lsp_nodes.len(),
1,
"re-adding the same node_id must not create a duplicate entry",
);
assert_eq!(
config.lsp_nodes[0].address, second_address,
"re-adding must update the address in place",
);
assert_eq!(
config.lsp_nodes[0].token,
Some("token".to_owned()),
"re-adding must update the token in place",
);
assert!(
!config.lsp_nodes[0].trust_peer_0conf,
"re-adding must update the 0conf trust setting in place",
);
}
}
129 changes: 98 additions & 31 deletions src/liquidity/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -154,50 +154,113 @@ 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. A changed address
/// disconnects any live session, dials the new address, and rediscovers protocols. Token or
/// 0conf-only updates are applied in place. A failed update restores the previous config only
/// if no newer update for the same `node_id` has landed.
pub fn add_liquidity_source(
&self, node_id: PublicKey, address: SocketAddress, token: Option<String>,
trust_peer_0conf: bool,
) -> Result<(), Error> {
let mut previous: Option<(SocketAddress, Option<String>, bool, Option<Vec<u16>>)> = None;
let mut addr_changed = true;
// Generation of this write. Cleanup restores or removes only if it is still current,
// so a concurrent re-add for the same node_id is not rolled back.
// New entries start at generation 1; existing entries inherit the current generation.
let mut update_generation = 1u64;
{
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) {
addr_changed = existing.address != address;
let changed = addr_changed
|| existing.token != token
|| existing.trust_peer_0conf != trust_peer_0conf;
if !changed {
log_info!(self.logger, "LSP node {} already added, skipping.", node_id);
return Ok(());
}
if addr_changed {
log_info!(
self.logger,
"Updating existing LSP node {} address/config; disconnecting and reconnecting.",
node_id
);
} else {
log_info!(
self.logger,
"Updating existing LSP node {} token/0conf without reconnecting.",
node_id
);
}
let prev_address = std::mem::replace(&mut existing.address, address.clone());
let prev_token = std::mem::replace(&mut existing.token, token.clone());
let prev_trust =
std::mem::replace(&mut existing.trust_peer_0conf, trust_peer_0conf);
// Keep supported_protocols until rediscovery overwrites them. Clearing here
// makes get_lsp_config miss this LSP for the whole connect+discover window.
let prev_protocols = existing.supported_protocols.clone();
existing.config_generation = existing.config_generation.wrapping_add(1);
update_generation = existing.config_generation;
previous = Some((prev_address, prev_token, prev_trust, prev_protocols));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, the take() removal (L192-194) resolves the JIT-window concern.

One thing still open: the snapshot rollback can race with a concurrent update for the same node_id. The lock is released at L205, then cleanup (L209-221) restores the snapshot taken at L195:

  1. A snapshots addr0, sets addr1
  2. B snapshots addr1, sets addr2
  3. A fails → L213-216 write addr0 (and stale supported_protocols) over B's in-flight addr2
  4. B succeeds → Ok(()), but the stored config is addr0

Suggested fix: serialize updates per LSP (lock held across the whole call), or connect + discover first and commit to lsp_nodes only on success, which makes rollback unnecessary. If concurrent re-adds aren't a supported use, a doc note would also be acceptable.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch — the lock is dropped before connect/discover, so a failed older update could write its snapshot over a newer re-add.

0499288 stamps each write with config_generation. Cleanup restores the snapshot (or removes a fresh insert) only when that generation is still current; a later add for the same node_id is left alone. Combined with the disconnect-then-dial fix so an address change actually hits the new endpoint before we treat the update as committed.

} else {
lsp_nodes.push(LspNode {
node_id,
address: address.clone(),
token: token.clone(),
trust_peer_0conf,
supported_protocols: None,
config_generation: update_generation,
});
}

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.
// Skip if a newer add for this node_id already committed a later generation.
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");
let still_current = nodes
.iter()
.find(|n| n.node_id == node_id)
.is_some_and(|n| n.config_generation == update_generation);
if !still_current {
return;
}
if let Some((prev_address, prev_token, prev_trust, prev_protocols)) = previous {
if let Some(existing) = nodes.iter_mut().find(|n| n.node_id == node_id) {
existing.address = prev_address;
existing.token = prev_token;
existing.trust_peer_0conf = prev_trust;
existing.supported_protocols = prev_protocols;
}
} 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
.runtime
.block_on(async move { con_cm.connect_peer_if_necessary(node_id, connect_addr).await })
{
cleanup();
return Err(e);
}
log_info!(self.logger, "Connected to LSP {}@{}.", node_id, address);

if let Err(e) = self
.runtime
.block_on(async { self.liquidity_source.discover_lsp_protocols(&node_id).await })
{
cleanup();
return Err(e);
if addr_changed {
let con_cm = Arc::clone(&self.connection_manager);
let connect_addr = address.clone();
// connect_peer_if_necessary no-ops while the peer is up, so a live LSP would keep
// the old session and never dial the new address. Drop it first, then dial.
if let Err(e) = self.runtime.block_on(async move {
con_cm.disconnect_peer(node_id);
con_cm.do_connect_peer(node_id, connect_addr).await
}) {
cleanup();
return Err(e);
}
log_info!(self.logger, "Connected to LSP {}@{}.", node_id, address);

if let Err(e) = self
.runtime
.block_on(async { self.liquidity_source.discover_lsp_protocols(&node_id).await })
{
cleanup();
return Err(e);
}
} else {
log_info!(self.logger, "Updated LSP {} config (address unchanged).", node_id);
}

Ok(())
Expand Down Expand Up @@ -232,6 +295,9 @@ pub(crate) struct LspNode {
trust_peer_0conf: bool,
// Protocol numbers discovered via LSPS0 (e.g., 1 = LSPS1, 2 = LSPS2, 5 = LSPS5).
supported_protocols: Option<Vec<u16>>,
// Bumped on every in-place config write. Failure cleanup matches this so a newer
// concurrent update is not overwritten by an older rollback.
config_generation: u64,
}

pub(crate) struct LiquiditySourceBuilder<L: Deref>
Expand Down Expand Up @@ -331,6 +397,7 @@ where
token: cfg.token,
trust_peer_0conf: cfg.trust_peer_0conf,
supported_protocols: None,
config_generation: 0,
})
.collect(),
));
Expand Down
Loading