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
164 changes: 126 additions & 38 deletions e2e-tests/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -115,7 +116,7 @@ impl TestBitcoind {

/// Handle to a running ldk-server child process.
pub struct LdkServerHandle {
child: Option<Child>,
child: ServerProcess,
pub grpc_port: u16,
pub p2p_port: u16,
pub storage_dir: PathBuf,
Expand All @@ -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)>,
Expand Down Expand Up @@ -385,17 +404,25 @@ 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.
let network_dir = storage_dir.join("regtest");
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();

Expand All @@ -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,
Expand All @@ -418,20 +445,18 @@ 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
}

/// 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");
}
Expand All @@ -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}");
Expand All @@ -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));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

I'm not sure the comment above matches the code.

It says both ports stay reserved until the child is ready to take them over. The listeners are dropped first, and only then does spawn_server_process start ldk-server. After that drop, is anything still holding the ports, or are they free until the child binds them?

I ran the end-to-end tests and they passed, so I can't show that this gap causes a failure. I just want to check whether the comment describes what the code does.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Thanks fixed

let child = spawn_server_process(&config_path);
(child, params, config_path)
}
Expand Down Expand Up @@ -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",
Expand All @@ -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()
}

Expand Down Expand Up @@ -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(_) => {
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -1085,32 +1162,43 @@ 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 })
.await
.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;
}
Expand Down
15 changes: 13 additions & 2 deletions e2e-tests/tests/e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -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]
Expand Down
Loading
Loading