diff --git a/e2e-tests/src/lib.rs b/e2e-tests/src/lib.rs index d31723f7..53d65b5a 100644 --- a/e2e-tests/src/lib.rs +++ b/e2e-tests/src/lib.rs @@ -11,6 +11,7 @@ use std::io::{BufRead, BufReader, Write}; use std::net::TcpListener; use std::path::{Path, PathBuf}; use std::process::{Child, Command, Stdio}; +use std::sync::atomic::{AtomicU32, Ordering}; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use corepc_node::Node; @@ -115,7 +116,7 @@ impl TestBitcoind { /// Handle to a running ldk-server child process. pub struct LdkServerHandle { - child: Option, + child: ServerProcess, pub grpc_port: u16, pub p2p_port: u16, pub storage_dir: PathBuf, @@ -126,6 +127,24 @@ pub struct LdkServerHandle { client: LdkServerClient, } +/// Reap the process even if startup fails before a server handle can be constructed. +struct ServerProcess(Child); + +impl ServerProcess { + fn assert_running(&mut self) { + if let Some(status) = self.0.try_wait().expect("Failed to check ldk-server process") { + panic!("ldk-server exited during startup with {status}; see server output above"); + } + } +} + +impl Drop for ServerProcess { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + #[derive(Default)] pub struct LdkServerConfig { pub metrics_auth: Option<(String, String)>, @@ -385,8 +404,9 @@ impl LdkServerHandle { pub async fn start_with_config( config_bitcoind: &TestBitcoind, config: impl FnOnce(&TestServerParams) -> String, ) -> Self { - let (mut child, params, config_path) = spawn_server(config_bitcoind, config); - forward_server_output(&mut child); + let (child, params, config_path) = spawn_server(config_bitcoind, config); + let mut child = ServerProcess(child); + forward_server_output(&mut child.0); let TestServerParams { grpc_port, p2p_port, storage_dir, .. } = params; // Wait for the admin macaroon and TLS certificate files to appear. @@ -394,8 +414,15 @@ impl LdkServerHandle { let macaroon_path = network_dir.join("macaroons").join("admin.macaroon"); let tls_cert_path = storage_dir.join("tls.crt"); - wait_for_file(&macaroon_path, Duration::from_secs(30)).await; - wait_for_file(&tls_cert_path, Duration::from_secs(30)).await; + let start = Instant::now(); + while !macaroon_path.exists() || !tls_cert_path.exists() { + child.assert_running(); + assert!( + start.elapsed() < Duration::from_secs(30), + "Timed out waiting for credentials in {storage_dir:?}" + ); + tokio::time::sleep(Duration::from_millis(100)).await; + } let macaroon = std::fs::read_to_string(&macaroon_path).unwrap().trim().to_string(); @@ -406,7 +433,7 @@ impl LdkServerHandle { let client = LdkServerClient::new(base_url, macaroon.clone(), &tls_cert_pem).unwrap(); let mut handle = Self { - child: Some(child), + child, grpc_port, p2p_port, storage_dir, @@ -418,7 +445,7 @@ impl LdkServerHandle { }; // Wait for server to be ready and get node info - let node_info = wait_for_server_ready(&handle, Duration::from_secs(60)).await; + let node_info = wait_for_server_ready(&mut handle, Duration::from_secs(60)).await; handle.node_id = node_info.node_id; handle @@ -426,12 +453,10 @@ impl LdkServerHandle { /// Kill and restart the server with the same config and storage to test crash recovery. pub async fn restart(&mut self) { - let mut child = self.child.take().expect("Server is not running"); - child.kill().expect("Failed to kill ldk-server"); - child.wait().expect("Failed to reap ldk-server"); - let mut child = spawn_server_process(&self.config_path); - forward_server_output(&mut child); - self.child = Some(child); + self.child.0.kill().expect("Failed to kill ldk-server"); + self.child.0.wait().expect("Failed to reap ldk-server"); + self.child = ServerProcess(spawn_server_process(&self.config_path)); + forward_server_output(&mut self.child.0); let info = wait_for_server_ready(self, Duration::from_secs(60)).await; assert_eq!(info.node_id, self.node_id, "Node identity changed after restart"); } @@ -449,23 +474,16 @@ impl LdkServerHandle { } } -impl Drop for LdkServerHandle { - fn drop(&mut self) { - if let Some(mut child) = self.child.take() { - let _ = child.kill(); - let _ = child.wait(); - } - } -} - /// Prepare test server params and spawn the ldk-server process. fn spawn_server( bitcoind: &TestBitcoind, config_fn: impl FnOnce(&TestServerParams) -> String, ) -> (Child, TestServerParams, PathBuf) { #[allow(deprecated)] let storage_dir = tempfile::tempdir().unwrap().into_path(); - let grpc_port = find_available_port(); - let p2p_port = find_available_port(); + let grpc_listener = reserve_port(); + let p2p_listener = reserve_port(); + let grpc_port = grpc_listener.local_addr().unwrap().port(); + let p2p_port = p2p_listener.local_addr().unwrap().port(); let (rpc_host, rpc_port_num, rpc_user, rpc_password) = bitcoind.rpc_details(); let rpc_address = format!("{rpc_host}:{rpc_port_num}"); @@ -478,6 +496,9 @@ fn spawn_server( let config_path = params.storage_dir.join("config.toml"); std::fs::write(&config_path, &config_content).unwrap(); + // Reserve both ports through config preparation, then release them before spawning. + // Another process can still claim them before the child binds. + drop((grpc_listener, p2p_listener)); let child = spawn_server_process(&config_path); (child, params, config_path) } @@ -531,6 +552,7 @@ pub fn start_expect_failure( Ok(None) => { if start.elapsed() > timeout { let _ = child.kill(); + let _ = child.wait(); panic!( "Server did not exit within {:?} — it may have started successfully \ instead of failing", @@ -555,9 +577,24 @@ pub fn start_expect_failure( String::from_utf8_lossy(&output.stderr).to_string() } -/// Find an available TCP port by binding to port 0. +/// Allocate distinct ports outside the default Linux/macOS ephemeral ranges. +/// Allocate each port only once per test process, so stopped servers can restart. +fn reserve_port() -> TcpListener { + static NEXT_PORT: AtomicU32 = AtomicU32::new(10_000); + loop { + let port = NEXT_PORT.fetch_add(1, Ordering::Relaxed); + assert!(port < 30_000, "Exhausted test server ports"); + match TcpListener::bind(("127.0.0.1", port as u16)) { + Ok(listener) => return listener, + Err(error) if error.kind() == std::io::ErrorKind::AddrInUse => continue, + Err(error) => panic!("Failed to reserve test port {port}: {error}"), + } + } +} + +/// Find an available port that will not be allocated to another server in this test process. pub fn find_available_port() -> u16 { - let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let listener = reserve_port(); listener.local_addr().unwrap().port() } @@ -600,9 +637,12 @@ pub async fn splice_txid(events: &mut EventStream) -> String { } /// Poll get_node_info until the server responds successfully. -async fn wait_for_server_ready(handle: &LdkServerHandle, timeout: Duration) -> GetNodeInfoResponse { +async fn wait_for_server_ready( + handle: &mut LdkServerHandle, timeout: Duration, +) -> GetNodeInfoResponse { let start = std::time::Instant::now(); loop { + handle.child.assert_running(); match handle.client().get_node_info(GetNodeInfoRequest {}).await { Ok(info) => return info, Err(_) => { @@ -938,11 +978,48 @@ pub async fn setup_funded_channel( .await .unwrap(); - // Mine blocks to confirm the channel and wait for servers to sync + // Opening returns before funding is broadcast. Mining immediately can leave the + // funding transaction with fewer than the six confirmations required for gossip. + let funding_txo = tokio::time::timeout(Duration::from_secs(30), async { + loop { + let channels = server_a.client().list_channels(ListChannelsRequest {}).await.unwrap(); + if let Some(txo) = channels + .channels + .iter() + .find(|channel| channel.user_channel_id == open_resp.user_channel_id) + .and_then(|channel| channel.funding_txo.clone()) + { + break txo; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + }) + .await + .expect("Opened channel did not acquire a funding outpoint"); + wait_for_transaction(bitcoind, &funding_txo.txid).await; mine_and_sync(bitcoind, &[server_a, server_b], 6).await; - // Wait for channel to become usable (mines blocks periodically to trigger chain sync) - wait_for_usable_channel(server_a.client(), bitcoind, Duration::from_secs(60)).await; + // User channel IDs are local to each node. Match the shared funding outpoint + // instead, and never let an older usable channel satisfy this wait. + for server in [server_a, server_b] { + let start = Instant::now(); + loop { + let channels = server.client().list_channels(ListChannelsRequest {}).await.unwrap(); + if channels.channels.iter().any(|channel| { + channel.funding_txo.as_ref() == Some(&funding_txo) + && channel.is_usable + && channel.confirmations.is_some_and(|count| count >= 6) + }) { + break; + } + assert!( + start.elapsed() < Duration::from_secs(60), + "Waiting for confirmed channel {funding_txo:?} on {}: {channels:?}", + server.node_id() + ); + tokio::time::sleep(Duration::from_millis(100)).await; + } + } open_resp.user_channel_id } @@ -1085,16 +1162,16 @@ pub async fn wait_for_settled_balance( } } -/// Wait for exactly `count` announced channels with both routing directions enabled. +/// Wait for exactly `count` announced channels that are ready for routing. +/// Local channels use live channel state, as in LDK's first-hop selection; remote +/// channels require both routing directions in the graph to be enabled. pub async fn wait_for_gossip(server: &LdkServerHandle, count: usize, timeout: Duration) { let start = Instant::now(); loop { let graph = server.client().graph_list_channels(GraphListChannelsRequest {}).await.unwrap(); - let mut ready = graph.short_channel_ids.len() == count; + let local_channels = server.client().list_channels(ListChannelsRequest {}).await.unwrap(); + let mut channels = Vec::new(); for short_channel_id in graph.short_channel_ids { - if !ready { - break; - } let channel = server .client() .graph_get_channel(GraphGetChannelRequest { short_channel_id }) @@ -1102,15 +1179,26 @@ pub async fn wait_for_gossip(server: &LdkServerHandle, count: usize, timeout: Du .unwrap() .channel .unwrap(); - ready = channel.one_to_two.is_some_and(|update| update.enabled) - && channel.two_to_one.is_some_and(|update| update.enabled); + channels.push((short_channel_id, channel)); } + let ready = channels.len() == count + && channels.iter().all(|(short_channel_id, channel)| { + if channel.node_one == server.node_id() || channel.node_two == server.node_id() { + // A peer's update can arrive before our own announcement and be dropped. + // The router overrides these graph entries with usable local channels. + return local_channels.channels.iter().any(|local| { + local.short_channel_id == Some(*short_channel_id) && local.is_usable + }); + } + channel.one_to_two.as_ref().is_some_and(|update| update.enabled) + && channel.two_to_one.as_ref().is_some_and(|update| update.enabled) + }); if ready { return; } assert!( start.elapsed() < timeout, - "Timed out waiting for {count} channel announcements with enabled routing updates" + "Timed out waiting for {count} routable channel announcements on {}: graph={channels:?}, local={local_channels:?}", server.node_id() ); tokio::time::sleep(Duration::from_millis(200)).await; } diff --git a/e2e-tests/tests/e2e.rs b/e2e-tests/tests/e2e.rs index be22d12a..7d361a81 100644 --- a/e2e-tests/tests/e2e.rs +++ b/e2e-tests/tests/e2e.rs @@ -870,6 +870,9 @@ async fn test_subscribe_events_channel_state_lifecycle_pending_ready_force_close let bitcoind = TestBitcoind::new(); let server_a = LdkServerHandle::start(&bitcoind).await; let server_b = LdkServerHandle::start(&bitcoind).await; + // Keep the peer connection alive when the target channel closes. LDK Node otherwise + // reconnects after closing the last channel, racing delivery of the force-close error. + let remaining_channel_id = setup_funded_channel(&bitcoind, &server_a, &server_b, 100_000).await; let addr_a = server_a.client().onchain_receive(OnchainReceiveRequest {}).await.unwrap().address; let addr_b = server_b.client().onchain_receive(OnchainReceiveRequest {}).await.unwrap().address; @@ -939,8 +942,11 @@ async fn test_subscribe_events_channel_state_lifecycle_pending_ready_force_close assert!(pending_b.reason.is_none()); assert_eq!(pending_b.closure_initiator, ChannelClosureInitiator::Unspecified as i32); + let (funding_txid, _) = pending_a.funding_txo.as_ref().unwrap().split_once(':').unwrap(); + wait_for_transaction(&bitcoind, funding_txid).await; mine_and_sync(&bitcoind, &[&server_a, &server_b], 6).await; - wait_for_usable_channel(server_a.client(), &bitcoind, Duration::from_secs(60)).await; + wait_for_channels(&server_a, 2, Duration::from_secs(60)).await; + wait_for_channels(&server_b, 2, Duration::from_secs(60)).await; let ready_a = wait_for_event(&mut events_a, |e| { matches!( @@ -982,7 +988,7 @@ async fn test_subscribe_events_channel_state_lifecycle_pending_ready_force_close assert_eq!(ready_b.closure_initiator, ChannelClosureInitiator::Unspecified as i32); run_cli(&server_a, &["force-close-channel", &open_resp.user_channel_id, server_b.node_id()]); - mine_and_sync(&bitcoind, &[&server_a, &server_b], 6).await; + // Observe the peer notification before mining, so an on-chain close cannot win the race. let closed_a = wait_for_event(&mut events_a, |e| { matches!( @@ -1030,6 +1036,11 @@ async fn test_subscribe_events_channel_state_lifecycle_pending_ready_force_close Some(ChannelStateChangeReasonKind::CounterpartyForceClosed) ); assert_eq!(closed_b.closure_initiator, ChannelClosureInitiator::Remote as i32); + mine_and_sync(&bitcoind, &[&server_a, &server_b], 6).await; + let remaining_a = wait_for_channels(&server_a, 1, Duration::from_secs(60)).await; + let remaining_b = wait_for_channels(&server_b, 1, Duration::from_secs(60)).await; + assert_eq!(remaining_a[0].user_channel_id, remaining_channel_id); + assert_eq!(remaining_a[0].channel_id, remaining_b[0].channel_id); } #[tokio::test] diff --git a/e2e-tests/tests/postgres.rs b/e2e-tests/tests/postgres.rs index fdcbcdfe..6626530b 100644 --- a/e2e-tests/tests/postgres.rs +++ b/e2e-tests/tests/postgres.rs @@ -7,7 +7,8 @@ // You may not use this file except in accordance with one or both of these // licenses. -use std::time::Duration; +use std::collections::BTreeMap; +use std::time::{Duration, Instant}; use e2e_tests::{ assert_recovered_balance, close_channel, expected_onchain_balance, list_payments, @@ -18,17 +19,96 @@ use e2e_tests::{ }; use ldk_server_grpc::api::{ open_channel_request, ConnectPeerRequest, DisconnectPeerRequest, ForceCloseChannelRequest, - GetBalancesRequest, GetNodeInfoRequest, ListChannelForwardingStatsRequest, + GetBalancesRequest, GetBalancesResponse, GetNodeInfoRequest, ListChannelForwardingStatsRequest, OnchainReceiveRequest, OpenChannelRequest, }; use ldk_server_grpc::types::{lightning_balance, payment_kind, BalanceSource, PaymentStatus}; const TIMEOUT: Duration = Duration::from_secs(60); +/// Wait until all funds are claimable on the expected open channels, without outstanding +/// outgoing HTLC claims or dust rounding. An inbound HTLC with a known preimage can still +/// be in the commitment at this point; its value is already included in the claimable amount. +async fn open_channel_balances(server: &LdkServerHandle, count: usize) -> GetBalancesResponse { + let start = Instant::now(); + loop { + let balances = server.client().get_balances(GetBalancesRequest {}).await.unwrap(); + if balances.lightning_balances.len() == count + && balances.pending_balances_from_channel_closures.is_empty() + && balances.lightning_balances.iter().all(|balance| { + matches!( + balance.balance_type.as_ref(), + Some(lightning_balance::BalanceType::ClaimableOnChannelClose(claim)) + if claim.outbound_payment_htlc_rounded_msat == 0 + && claim.outbound_forwarded_htlc_rounded_msat == 0 + && claim.inbound_claiming_htlc_rounded_msat == 0 + && claim.inbound_htlc_rounded_msat == 0 + ) + }) { + return balances; + } + assert!( + start.elapsed() < TIMEOUT, + "Waiting for {count} open channel balances: {balances:?}" + ); + tokio::time::sleep(Duration::from_millis(100)).await; + } +} + +/// Compare funds, including commitment fees, independently for each channel. Removing an +/// already-claimed inbound HTLC reduces the commitment fee and increases the claimable +/// amount by the same value. That may finish during restart without changing our funds. +fn channel_funds(balances: &GetBalancesResponse) -> BTreeMap { + let mut total_claimable = 0; + let mut funds = BTreeMap::new(); + for balance in &balances.lightning_balances { + let Some(lightning_balance::BalanceType::ClaimableOnChannelClose(claim)) = + &balance.balance_type + else { + panic!("Expected open channel balance: {balance:?}"); + }; + total_claimable += claim.amount_satoshis; + assert!(funds + .insert( + claim.channel_id.clone(), + claim.amount_satoshis + claim.transaction_fee_satoshis + ) + .is_none()); + } + assert_eq!(balances.total_lightning_balance_sats, total_claimable); + funds +} + +#[test] +fn test_channel_funds_preserve_exact_value_across_fee_changes() { + let balances = |amount_satoshis, transaction_fee_satoshis| GetBalancesResponse { + total_lightning_balance_sats: amount_satoshis, + lightning_balances: vec![ldk_server_grpc::types::LightningBalance { + balance_type: Some(lightning_balance::BalanceType::ClaimableOnChannelClose( + ldk_server_grpc::types::ClaimableOnChannelClose { + channel_id: "channel".to_string(), + amount_satoshis, + transaction_fee_satoshis, + ..Default::default() + }, + )), + }], + ..Default::default() + }; + // Reproduce the observed 43-sat change when the final HTLC leaves the commitment. + let before = channel_funds(&balances(919_014, 327)); + assert_eq!(before, BTreeMap::from([("channel".to_string(), 919_341)])); + assert_eq!(before, channel_funds(&balances(919_057, 284))); + // A real one-satoshi loss must still fail the persistence comparison. + assert_ne!(before, channel_funds(&balances(919_056, 284))); +} + async fn start_postgres(bitcoind: &TestBitcoind, connection_string: &str) -> LdkServerHandle { let server = LdkServerHandle::start_with_config(bitcoind, |params| { // Each server gets its own table, even when sharing the same test database. - let table_name = format!("node_{}", params.grpc_port); + // Ports may be reused by a later test process against the same database. + let directory = params.storage_dir.file_name().unwrap().to_str().unwrap(); + let table_name = format!("node_{}", directory.replace('.', "_")); TestConfigBuilder::new(params) .postgres(connection_string, &table_name) .forwarded_payment_tracking_mode("detailed") @@ -96,9 +176,6 @@ async fn test_postgres_persistence_and_sqlite_interoperability() { // SQLite A -> PostgreSQL B -> SQLite C. B both accepts and initiates a channel. let channel_ab = setup_funded_channel(&bitcoind, &sqlite_a, &postgres, 1_000_000).await; let channel_bc = setup_funded_channel(&bitcoind, &postgres, &sqlite_c, 1_000_000).await; - // The shared helper waits for any usable channel on the funder; B already has A-B. - // Keep mining until C's only channel is confirmed, too. - wait_for_usable_channel(sqlite_c.client(), &bitcoind, TIMEOUT).await; wait_for_channels(&sqlite_a, 1, TIMEOUT).await; wait_for_channels(&postgres, 2, TIMEOUT).await; wait_for_channels(&sqlite_c, 1, TIMEOUT).await; @@ -139,7 +216,8 @@ async fn test_postgres_persistence_and_sqlite_interoperability() { 4 ); let saved_channels = wait_for_channels(&postgres, 2, TIMEOUT).await; - let saved_balances = postgres.client().get_balances(GetBalancesRequest {}).await.unwrap(); + let saved_balances = open_channel_balances(&postgres, 2).await; + let saved_channel_funds = channel_funds(&saved_balances); assert!(saved_balances.total_onchain_balance_sats > 0); assert!(saved_balances.total_lightning_balance_sats > 0); let saved_stats = postgres @@ -170,15 +248,13 @@ async fn test_postgres_persistence_and_sqlite_interoperability() { assert_eq!(restored.funding_txo, saved.funding_txo); assert_eq!(restored.channel_value_sats, saved.channel_value_sats); } - let restored_balances = postgres.client().get_balances(GetBalancesRequest {}).await.unwrap(); + let restored_balances = open_channel_balances(&postgres, 2).await; assert_eq!( restored_balances.total_onchain_balance_sats, saved_balances.total_onchain_balance_sats ); - assert_eq!( - restored_balances.total_lightning_balance_sats, - saved_balances.total_lightning_balance_sats - ); + assert_eq!(channel_funds(&restored_balances), saved_channel_funds, + "Channel funds changed across restart: before={saved_balances:?}, after={restored_balances:?}"); assert_eq!(list_payments(&postgres).await, saved_payments); assert_eq!(wait_for_forwarded_payments(&postgres, 2, TIMEOUT).await, saved_forwards); let restored_stats = postgres