diff --git a/Cargo.lock b/Cargo.lock index 70f3192a..5f556ff6 100755 --- a/Cargo.lock +++ b/Cargo.lock @@ -224,6 +224,15 @@ dependencies = [ "windows-targets 0.52.6", ] +[[package]] +name = "base58ck" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f1cba749a07c1efb1f4d87518f39cea4aec25d0991fb97e80459c057238f0d2" +dependencies = [ + "bitcoin_hashes 0.14.1", +] + [[package]] name = "base64" version = "0.13.1" @@ -248,20 +257,60 @@ version = "0.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d86b93f97252c47b41663388e6d155714a9d0c398b99f1005cbc5f978b29f445" +[[package]] +name = "bech32" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32637268377fc7b10a8c6d51de3e7fba1ce5dd371a96e342b34e6078db558e7f" + [[package]] name = "bitcoin" version = "0.30.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1945a5048598e4189e239d3f809b19bdad4845c4b2ba400d304d2dcf26d2c462" dependencies = [ - "bech32", + "bech32 0.9.1", "bitcoin-private", "bitcoin_hashes 0.12.0", "hex_lit", - "secp256k1", + "secp256k1 0.27.0", + "serde", +] + +[[package]] +name = "bitcoin" +version = "0.32.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7ca0f87890d2398219d15182cb6af9681722bd891758dab4ca01896e2038a510" +dependencies = [ + "base58ck", + "bech32 0.11.1", + "bitcoin-io", + "bitcoin-units", + "bitcoin_hashes 0.14.1", + "hex-conservative 0.2.2", + "hex_lit", + "secp256k1 0.29.1", + "serde", +] + +[[package]] +name = "bitcoin-consensus-encoding" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6712f9c6fd6785b3b270884e57c441c403dc5d7e19ca45368c97c7a1de3000ec" +dependencies = [ + "bitcoin-internals", + "hex-conservative 1.3.0", "serde", ] +[[package]] +name = "bitcoin-internals" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d573f4cf32996a8dce612e4348cece65a241f1882ed594047c9ba348e8869fa5" + [[package]] name = "bitcoin-io" version = "0.1.4" @@ -274,6 +323,16 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "73290177011694f38ec25e165d0387ab7ea749a4b81cd4c80dae5988229f7a57" +[[package]] +name = "bitcoin-units" +version = "0.1.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9cb95693f371d089a4b5b6fc41c6f3ea6e01ee8c15388335dfac8ea685173b51" +dependencies = [ + "bitcoin-consensus-encoding", + "serde", +] + [[package]] name = "bitcoin_hashes" version = "0.12.0" @@ -292,6 +351,7 @@ checksum = "26ec84b80c482df901772e931a9a681e26a1b9ee2302edeff23cb30328745c8b" dependencies = [ "bitcoin-io", "hex-conservative 0.2.2", + "serde", ] [[package]] @@ -402,7 +462,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d528ade112e169cbb79ceae2f235b4502fee6e939326bbaa36aafcdfd54cd91c" dependencies = [ "anyhow", - "bitcoin", + "bitcoin 0.30.2", "futures-core", "hex", "log", @@ -630,6 +690,12 @@ dependencies = [ "syn 2.0.100", ] +[[package]] +name = "dnssec-prover" +version = "0.6.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9468f1a08c50bd1e5ad91b151e11ce8e806f8fa1c1eb9b07f66c7011de45a2e" + [[package]] name = "downcast" version = "0.11.0" @@ -941,6 +1007,12 @@ version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" +[[package]] +name = "hashbrown" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "43a3c133739dddd0d2990f9a4bdf8eb4b21ef50e4851ca85ab661199821d510e" + [[package]] name = "hashbrown" version = "0.14.5" @@ -989,15 +1061,18 @@ checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" [[package]] name = "hex-conservative" -version = "0.1.2" +version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "212ab92002354b4819390025006c897e8140934349e8635c9b077f47b4dcbd20" +checksum = "fda06d18ac606267c40c04e41b9947729bf8b9efe74bd4e82b61a5f26a510b9f" +dependencies = [ + "arrayvec 0.7.6", +] [[package]] name = "hex-conservative" -version = "0.2.2" +version = "1.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fda06d18ac606267c40c04e41b9947729bf8b9efe74bd4e82b61a5f26a510b9f" +checksum = "271e0d19bcb473b6675739a2b536076b24a082316cb5199ad918edce10c599e8" dependencies = [ "arrayvec 0.7.6", ] @@ -1484,12 +1559,50 @@ checksum = "8355be11b20d696c8f18f6cc018c4e372165b1fa8126cef092399c9951984ffa" [[package]] name = "lightning" -version = "0.0.123" +version = "0.2.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a8c4111946d048fecc72245a2ecacf28cde4390bcc946b010f544c6ac2b2d8b" +dependencies = [ + "bech32 0.11.1", + "bitcoin 0.32.11", + "dnssec-prover", + "hashbrown 0.13.2", + "libm", + "lightning-invoice", + "lightning-macros", + "lightning-types", + "possiblyrandom", +] + +[[package]] +name = "lightning-invoice" +version = "0.34.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47d83bd798e04ab9eecc8bbef1fa17d3808859bcdc0406bd16c55d51c8834444" +dependencies = [ + "bech32 0.11.1", + "bitcoin 0.32.11", + "lightning-types", +] + +[[package]] +name = "lightning-macros" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d4c717494cdc2c8bb85bee7113031248f5f6c64f8802b33c1c9e2d98e594aa71" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.100", +] + +[[package]] +name = "lightning-types" +version = "0.3.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5fd92d4aa159374be430c7590e169b4a6c0fb79018f5bc4ea1bffde536384db3" +checksum = "c77c676d4a34cceb2ae3756916e446b4d17f9430a24107e099981f0f9aec77e6" dependencies = [ - "bitcoin", - "hex-conservative 0.1.2", + "bitcoin 0.32.11", ] [[package]] @@ -1887,6 +2000,15 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" +[[package]] +name = "possiblyrandom" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9c564dbf654befd49035528299f1208a40508f6e07efb11c163444e304e4484f" +dependencies = [ + "getrandom 0.2.15", +] + [[package]] name = "powerfmt" version = "0.2.0" @@ -2537,7 +2659,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "25996b82292a7a57ed3508f052cfff8640d38d32018784acd714758b43da9c8f" dependencies = [ "bitcoin_hashes 0.12.0", - "secp256k1-sys", + "secp256k1-sys 0.8.1", + "serde", +] + +[[package]] +name = "secp256k1" +version = "0.29.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9465315bc9d4566e1724f0fffcbcc446268cb522e60f9a27bcded6b19c108113" +dependencies = [ + "bitcoin_hashes 0.14.1", + "secp256k1-sys 0.10.1", "serde", ] @@ -2550,6 +2683,15 @@ dependencies = [ "cc", ] +[[package]] +name = "secp256k1-sys" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d4387882333d3aa8cb20530a17c69a3752e97837832f34f6dccc760e715001d9" +dependencies = [ + "cc", +] + [[package]] name = "security-framework" version = "2.11.1" @@ -2661,7 +2803,7 @@ name = "sim-cli" version = "0.1.1" dependencies = [ "anyhow", - "bitcoin", + "bitcoin 0.32.11", "clap", "console-subscriber", "ctrlc", @@ -2686,7 +2828,7 @@ version = "0.1.1" dependencies = [ "anyhow", "async-trait", - "bitcoin", + "bitcoin 0.32.11", "cln-grpc", "csv", "expanduser", diff --git a/sim-cli/Cargo.toml b/sim-cli/Cargo.toml index e3249418..4c6809bd 100755 --- a/sim-cli/Cargo.toml +++ b/sim-cli/Cargo.toml @@ -21,7 +21,7 @@ simple_logger = "4.2.0" # The virtual-time feature is required for the --virtual-time flag, which runs simulations on a paused runtime. simln-lib = { path = "../simln-lib", features = ["virtual-time"] } tokio = { version = "1.26.0", features = ["full"] } -bitcoin = { version = "0.30.1" } +bitcoin = { version = "0.32" } ctrlc = "3.4.0" rand = "0.8.5" hex = {version = "0.4.3"} diff --git a/sim-cli/src/parsing.rs b/sim-cli/src/parsing.rs index 350e5e39..87107525 100755 --- a/sim-cli/src/parsing.rs +++ b/sim-cli/src/parsing.rs @@ -260,7 +260,7 @@ pub async fn create_simulation_with_network( ( Simulation, Vec, - HashMap>>>, + HashMap>>, ), anyhow::Error, > { @@ -322,9 +322,9 @@ pub async fn create_simulation_with_network( // to a dyn trait and exclude any nodes that shouldn't be included in random activity // generation. let nodes = ln_node_from_graph(simulation_graph, routing_graph, clock.clone()).await?; - let mut nodes_dyn: HashMap<_, Arc>> = nodes + let mut nodes_dyn: HashMap<_, Arc> = nodes .iter() - .map(|(pk, node)| (*pk, Arc::clone(node) as Arc>)) + .map(|(pk, node)| (*pk, Arc::clone(node) as Arc)) .collect(); for pk in exclude { nodes_dyn.remove(pk); @@ -388,26 +388,25 @@ async fn get_clients( nodes: Vec, ) -> Result< ( - HashMap>>, + HashMap>, HashMap, ), LightningError, > { - let mut clients: HashMap>> = HashMap::new(); + let mut clients: HashMap> = HashMap::new(); let mut clients_info: HashMap = HashMap::new(); for connection in nodes { - // TODO: Feels like there should be a better way of doing this without having to Arc>> it at this time. - // Box sort of works, but we won't know the size of the dyn LightningNode at compile time so the compiler will - // scream at us when trying to create the Arc> later on while adding the node to the clients map - let node: Arc> = match connection { - NodeConnection::Lnd(c) => Arc::new(Mutex::new(LndNode::new(c).await?)), - NodeConnection::Cln(c) => Arc::new(Mutex::new(ClnNode::new(c).await?)), - NodeConnection::Eclair(c) => Arc::new(Mutex::new(EclairNode::new(c).await?)), - NodeConnection::LdkServer(c) => Arc::new(Mutex::new(LdkServerNode::new(c).await?)), + // We won't know the size of the dyn LightningNode at compile time, so the node is boxed into an Arc to + // store it in the clients map. + let node: Arc = match connection { + NodeConnection::Lnd(c) => Arc::new(LndNode::new(c).await?), + NodeConnection::Cln(c) => Arc::new(ClnNode::new(c).await?), + NodeConnection::Eclair(c) => Arc::new(EclairNode::new(c).await?), + NodeConnection::LdkServer(c) => Arc::new(LdkServerNode::new(c).await?), }; - let node_info = node.lock().await.get_info().clone(); + let node_info = node.get_info().clone(); clients.insert(node_info.pubkey, node); clients_info.insert(node_info.pubkey, node_info); @@ -606,7 +605,7 @@ pub fn parse_sim_params(cli: &Cli) -> anyhow::Result { } pub async fn get_validated_activities( - clients: &HashMap>>, + clients: &HashMap>, nodes_info: HashMap, activity: Vec, ) -> Result, LightningError> { @@ -614,8 +613,6 @@ pub async fn get_validated_activities( // nodes that we do not control. To do this, we can just grab the first node in our map and perform the lookup. let graph = match clients.values().next() { Some(client) => client - .lock() - .await .get_graph() .await .map_err(|e| LightningError::GetGraphError(format!("Error getting graph {:?}", e))), diff --git a/simln-lib/Cargo.toml b/simln-lib/Cargo.toml index 1b6aa6d7..ee55cec1 100755 --- a/simln-lib/Cargo.toml +++ b/simln-lib/Cargo.toml @@ -15,8 +15,8 @@ cln-grpc = "0.1.3" expanduser = "1.2.2" serde = { version="1.0.183", features=["derive"] } serde_json = "1.0.104" -bitcoin = { version = "0.30.1", features=["serde"] } -lightning = { version = "0.0.123" } +bitcoin = { version = "0.32", features=["serde"] } +lightning = { version = "0.2" } tonic_lnd = { package="fedimint-tonic-lnd", version="0.1.2", features=["lightningrpc", "routerrpc"]} tonic = { version = "0.8", features = ["tls", "transport"] } async-trait = "0.1.73" diff --git a/simln-lib/src/batched_writer.rs b/simln-lib/src/batched_writer.rs index 8a9b0934..7b4f22fb 100644 --- a/simln-lib/src/batched_writer.rs +++ b/simln-lib/src/batched_writer.rs @@ -26,9 +26,11 @@ impl BatchedWriter { let file = directory.join(file_name); let writer = WriterBuilder::new() - .from_path(file) + .from_path(&file) .map_err(SimulationError::CsvError)?; + log::info!("Writing simulation results to {}.", file.display()); + Ok(BatchedWriter { batch_size, counter: 0, diff --git a/simln-lib/src/cln.rs b/simln-lib/src/cln.rs index 5bfa1aa0..3f962851 100644 --- a/simln-lib/src/cln.rs +++ b/simln-lib/src/cln.rs @@ -8,8 +8,8 @@ use cln_grpc::pb::{ KeysendRequest, KeysendResponse, ListchannelsRequest, ListnodesRequest, ListpaysRequest, ListpaysResponse, }; -use lightning::ln::features::NodeFeatures; -use lightning::ln::PaymentHash; +use lightning::types::features::NodeFeatures; +use lightning::types::payment::PaymentHash; use serde::{Deserialize, Serialize}; use tokio::fs::File; use tokio::io::{AsyncReadExt, Error}; diff --git a/simln-lib/src/eclair.rs b/simln-lib/src/eclair.rs index bea94575..90d9b4cd 100644 --- a/simln-lib/src/eclair.rs +++ b/simln-lib/src/eclair.rs @@ -5,8 +5,8 @@ use crate::{ use async_trait::async_trait; use bitcoin::secp256k1::PublicKey; use bitcoin::Network; -use lightning::ln::features::NodeFeatures; -use lightning::ln::{PaymentHash, PaymentPreimage}; +use lightning::types::features::NodeFeatures; +use lightning::types::payment::{PaymentHash, PaymentPreimage}; use reqwest::multipart::Form; use reqwest::{Client, Method, Url}; use serde::{Deserialize, Serialize}; diff --git a/simln-lib/src/latency_interceptor.rs b/simln-lib/src/latency_interceptor.rs index 3d904972..8798428e 100644 --- a/simln-lib/src/latency_interceptor.rs +++ b/simln-lib/src/latency_interceptor.rs @@ -81,7 +81,7 @@ mod tests { use crate::sim_node::{CustomRecords, HtlcRef, InterceptRequest}; use crate::test_utils::get_random_keypair; use crate::ShortChannelID; - use lightning::ln::PaymentHash; + use lightning::types::payment::PaymentHash; use ntest::assert_true; use rand::distributions::Distribution; use rand::rngs::StdRng; diff --git a/simln-lib/src/ldk_server.rs b/simln-lib/src/ldk_server.rs index ed379865..a7ad47aa 100644 --- a/simln-lib/src/ldk_server.rs +++ b/simln-lib/src/ldk_server.rs @@ -10,8 +10,8 @@ use ldk_server_client::ldk_server_grpc::api::{ ListChannelsRequest, SpontaneousSendRequest, }; use ldk_server_client::ldk_server_grpc::types::{GraphNodeAnnouncement, PaymentStatus}; -use lightning::ln::features::NodeFeatures; -use lightning::ln::PaymentHash; +use lightning::types::features::NodeFeatures; +use lightning::types::payment::PaymentHash; use serde::{Deserialize, Serialize}; use tokio::time::{self, Duration}; use triggered::Listener; diff --git a/simln-lib/src/lib.rs b/simln-lib/src/lib.rs index d2171af8..05e7e1ac 100755 --- a/simln-lib/src/lib.rs +++ b/simln-lib/src/lib.rs @@ -5,8 +5,8 @@ use self::clock::Clock; use async_trait::async_trait; use bitcoin::secp256k1::PublicKey; use bitcoin::Network; -use lightning::ln::features::NodeFeatures; -use lightning::ln::PaymentHash; +use lightning::types::features::NodeFeatures; +use lightning::types::payment::PaymentHash; use rand::{Rng, RngCore, SeedableRng}; use rand_chacha::ChaCha8Rng; use random_activity::RandomActivityError; @@ -330,7 +330,7 @@ impl Default for Graph { /// LightningNode represents the functionality that is required to execute events on a lightning node. #[async_trait] -pub trait LightningNode: Send { +pub trait LightningNode: Send + Sync { /// Get information about the node. fn get_info(&self) -> &NodeInfo; /// Get the network this node is running at. @@ -605,7 +605,7 @@ pub struct Simulation { /// Config for the simulation itself. cfg: SimulationCfg, /// The lightning node that is being simulated. - nodes: HashMap>>, + nodes: HashMap>, /// Results logger that holds the simulation statistics. results: Arc>, /// Track all tasks spawned for use in the simulation. When used in the `run` method, it will wait for @@ -683,7 +683,7 @@ struct ProducePaymentEventsTrackers { impl Simulation { pub fn new( cfg: SimulationCfg, - nodes: HashMap>>, + nodes: HashMap>, tasks: TaskTracker, clock: Arc, shutdown_trigger: Trigger, @@ -715,7 +715,6 @@ impl Simulation { )); } else { for node in self.nodes.values() { - let node = node.lock().await; if !node.get_info().features.supports_keysend() { return Err(LightningError::ValidationError(format!( "All nodes eligible for random activity generation must support keysend, {} does not", @@ -765,7 +764,7 @@ impl Simulation { let mut running_network = Option::None; for node in self.nodes.values() { - let network = node.lock().await.get_network(); + let network = node.get_network(); if network == Network::Bitcoin { return Err(LightningError::ValidationError( "mainnet is not supported".to_string(), @@ -1015,7 +1014,7 @@ impl Simulation { // While we're at it, we get the node info and store it with capacity to create activity generators in our // second pass. for (pk, node) in self.nodes.iter() { - let chan_capacity = node.lock().await.channel_capacities().await?; + let chan_capacity = node.channel_capacities().await?; if let Err(e) = RandomPaymentActivity::validate_capacity( chan_capacity, @@ -1028,7 +1027,7 @@ impl Simulation { // Don't double count channel capacity because each channel reports the total balance between counter // parities. Track capacity separately to be used for our network generator. let capacity = chan_capacity / 2; - let node_info = node.lock().await.get_node_info(pk).await?; + let node_info = node.get_node_info(pk).await?; active_nodes.insert(node_info.pubkey, (node_info, capacity)); } @@ -1125,7 +1124,7 @@ impl Simulation { async fn produce_payment_events( mut heap: BinaryHeap>, mut payments_tracker: HashMap, - nodes: HashMap>>, + nodes: HashMap>, clock: Arc, output_sender: Sender, trackers: ProducePaymentEventsTrackers, @@ -1281,15 +1280,13 @@ async fn generate_payment( /// events that are crated for a lightning node that we can execute events on. Any output that is generated from the /// event being executed is piped into a channel to handle the result of the event. async fn send_payment( - node: Arc>, + node: Arc, sender: Sender, simulation_event: SimulationEvent, dispatch_time: SystemTime, ) -> Result<(), SimulationError> { match simulation_event { SimulationEvent::SendPayment(dest, amt_msat) => { - let node = node.lock().await; - let mut payment = Payment { source: node.get_info().pubkey, hash: None, @@ -1499,7 +1496,7 @@ async fn run_results_logger( /// out. In the multiple-producer case, a single producer shutting down does not drop *all* sending channels so the /// consumer will not exit and a trigger is required. async fn produce_simulation_results( - nodes: HashMap>>, + nodes: HashMap>, mut output_receiver: Receiver, results: Sender<(Payment, PaymentResult)>, listener: Listener, @@ -1552,15 +1549,13 @@ async fn produce_simulation_results( } async fn track_payment_result( - node: Arc>, + node: Arc, results: Sender<(Payment, PaymentResult)>, payment: Payment, listener: Listener, ) -> Result<(), SimulationError> { log::trace!("Payment result tracker starting."); - let node = node.lock().await; - let res = match payment.hash { Some(hash) => { log::debug!("Tracking payment outcome for: {}.", hex::encode(hash.0)); @@ -1626,7 +1621,6 @@ mod tests { use std::sync::Arc; use std::sync::Mutex as StdMutex; use std::time::{Duration, SystemTime}; - use tokio::sync::Mutex; use tokio_util::task::TaskTracker; #[test] @@ -1904,7 +1898,7 @@ mod tests { /// "we don't control any nodes". #[tokio::test] async fn test_validate_node_network_empty_nodes() { - let empty_nodes: HashMap>> = HashMap::new(); + let empty_nodes: HashMap> = HashMap::new(); let simulation = test_utils::create_simulation(empty_nodes); let result = simulation.validate_node_network().await; @@ -1974,13 +1968,15 @@ mod tests { assert!(result.is_ok()); } - async fn mock_send_payment( - mock_node: &mut Arc>, + /// Sets up a mock node's expectations. Takes the node before it is shared, so that it can be mutated through + /// its `Arc`. + fn mock_send_payment( + mock_node: &mut Arc, node_info: NodeInfo, payments_list: Arc>>, payment_hash: [u8; 32], ) { - let mut mock_node = mock_node.lock().await; + let mock_node = Arc::get_mut(mock_node).expect("node is not shared yet"); mock_node.expect_get_info().return_const(node_info.clone()); mock_node .expect_get_network() @@ -2001,7 +1997,7 @@ mod tests { let pl = payments_list.clone(); mock_node.expect_send_payment().returning(move |a, _| { pl.lock().unwrap().push(a); - Ok(lightning::ln::PaymentHash(payment_hash)) + Ok(lightning::types::payment::PaymentHash(payment_hash)) }); } @@ -2022,7 +2018,7 @@ mod tests { #[allow(clippy::type_complexity)] /// Helper to create and configure mock nodes for testing - async fn setup_test_nodes_for_testing_deterministic_events( + fn setup_test_nodes_for_testing_deterministic_events( fixed_pubkeys: Option>, ) -> (TestNodesResult, Arc>>) { let mut builder = LightningTestNodeBuilder::new(4); @@ -2038,32 +2034,28 @@ mod tests { network.nodes[0].clone(), payments_list.clone(), [0; 32], - ) - .await; + ); mock_send_payment( &mut network.clients[1], network.nodes[1].clone(), payments_list.clone(), [0; 32], - ) - .await; + ); mock_send_payment( &mut network.clients[2], network.nodes[2].clone(), payments_list.clone(), [0; 32], - ) - .await; + ); mock_send_payment( &mut network.clients[3], network.nodes[3].clone(), payments_list.clone(), [0; 32], - ) - .await; + ); (network, payments_list) } @@ -2073,7 +2065,7 @@ mod tests { #[tokio::test(start_paused = true)] async fn test_deterministic_payments_events_defined_activities() { let (network, payments_list) = - setup_test_nodes_for_testing_deterministic_events(Some(fixed_test_pubkeys())).await; + setup_test_nodes_for_testing_deterministic_events(Some(fixed_test_pubkeys())); // Define two activities // Activity 1: From node_1 to node_2 @@ -2156,8 +2148,7 @@ mod tests { let pks = fixed_test_pubkeys(); let (pk1, pk2, pk3, pk4) = (pks[0], pks[1], pks[2], pks[3]); - let (network, payments_list) = - setup_test_nodes_for_testing_deterministic_events(Some(pks)).await; + let (network, payments_list) = setup_test_nodes_for_testing_deterministic_events(Some(pks)); let (shutdown_trigger, shutdown_listener) = triggered::trigger(); diff --git a/simln-lib/src/lnd.rs b/simln-lib/src/lnd.rs index 4bbc9470..531abc5d 100644 --- a/simln-lib/src/lnd.rs +++ b/simln-lib/src/lnd.rs @@ -9,8 +9,8 @@ use async_trait::async_trait; use bitcoin::hashes::{sha256, Hash}; use bitcoin::secp256k1::PublicKey; use bitcoin::Network; -use lightning::ln::features::NodeFeatures; -use lightning::ln::{PaymentHash, PaymentPreimage}; +use lightning::types::features::NodeFeatures; +use lightning::types::payment::{PaymentHash, PaymentPreimage}; use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; use tonic_lnd::lnrpc::{payment::PaymentStatus, GetInfoRequest}; diff --git a/simln-lib/src/runtime.rs b/simln-lib/src/runtime.rs index 43932412..617e8694 100644 --- a/simln-lib/src/runtime.rs +++ b/simln-lib/src/runtime.rs @@ -6,6 +6,7 @@ //! advance, the library builds and owns the runtime here. use std::future::Future; +use std::num::NonZero; use std::sync::Arc; use std::time::SystemTime; @@ -55,9 +56,14 @@ where )); } + let blocking_threads = std::thread::available_parallelism() + .map(NonZero::get) + .unwrap_or(4); + let runtime = Builder::new_current_thread() .enable_all() .start_paused(true) + .max_blocking_threads(blocking_threads) .build() .map_err(|e| { SimulationError::RuntimeError(format!("could not build virtual-time runtime: {e}")) diff --git a/simln-lib/src/serializers.rs b/simln-lib/src/serializers.rs index c15d5217..8e7ee48f 100644 --- a/simln-lib/src/serializers.rs +++ b/simln-lib/src/serializers.rs @@ -2,7 +2,7 @@ use expanduser::expanduser; use serde::Deserialize; pub mod serde_option_payment_hash { - use lightning::ln::PaymentHash; + use lightning::types::payment::PaymentHash; pub fn serialize(hash: &Option, serializer: S) -> Result where diff --git a/simln-lib/src/sim_node.rs b/simln-lib/src/sim_node.rs index f25c4b92..931f937f 100755 --- a/simln-lib/src/sim_node.rs +++ b/simln-lib/src/sim_node.rs @@ -5,7 +5,7 @@ use crate::{ use async_trait::async_trait; use bitcoin::constants::ChainHash; use bitcoin::secp256k1::PublicKey; -use bitcoin::{Network, ScriptBuf, TxOut}; +use bitcoin::{Amount, Network, ScriptBuf, TxOut}; use lightning::ln::chan_utils::make_funding_redeemscript; use serde::{Deserialize, Serialize}; use std::collections::{hash_map::Entry, HashMap}; @@ -15,17 +15,17 @@ use std::time::UNIX_EPOCH; use tokio::task::JoinSet; use tokio_util::task::TaskTracker; -use lightning::ln::features::{ChannelFeatures, NodeFeatures}; use lightning::ln::msgs::{ LightningError as LdkError, UnsignedChannelAnnouncement, UnsignedChannelUpdate, }; -use lightning::ln::{PaymentHash, PaymentPreimage}; use lightning::routing::gossip::{NetworkGraph, NodeId}; use lightning::routing::router::{find_route, Path, PaymentParameters, Route, RouteParameters}; use lightning::routing::scoring::{ ProbabilisticScorer, ProbabilisticScoringDecayParameters, ScoreUpdate, }; use lightning::routing::utxo::{UtxoLookup, UtxoResult}; +use lightning::types::features::{ChannelFeatures, NodeFeatures}; +use lightning::types::payment::{PaymentHash, PaymentPreimage}; use lightning::util::logger::{Level, Logger, Record}; use thiserror::Error; use tokio::select; @@ -514,6 +514,7 @@ pub trait SimNetwork: Send + Sync { } type LdkNetworkGraph = NetworkGraph>; +type Scorer = ProbabilisticScorer, Arc>; struct InFlightPayment { /// The channel used to report payment results to. @@ -538,7 +539,7 @@ pub struct SimNode { pathfinding_graph: Arc, /// Probabilistic scorer used to rank paths through the network for routing. This is reused across /// multiple payments to maintain scoring state. - scorer: Mutex, Arc>>, + scorer: Arc>, /// Clock for tracking simulation time. clock: Arc, } @@ -566,7 +567,7 @@ impl SimNode { network: payment_network, in_flight: Mutex::new(HashMap::new()), pathfinding_graph, - scorer: Mutex::new(scorer), + scorer: Arc::new(std::sync::RwLock::new(scorer)), clock, }) } @@ -581,7 +582,7 @@ impl SimNode { /// /// **Note:** The route passed in here must contain only one path. pub async fn send_to_route( - &mut self, + &self, route: Route, payment_hash: PaymentHash, custom_records: Option, @@ -635,14 +636,13 @@ fn node_info(pubkey: PublicKey, alias: String) -> NodeInfo { /// Uses LDK's pathfinding algorithm with default parameters to find a path from source to destination, with no /// restrictions on fee budget. -async fn find_payment_route( +fn find_payment_route( source: &PublicKey, dest: PublicKey, amount_msat: u64, pathfinding_graph: &LdkNetworkGraph, - scorer: &Mutex, Arc>>, + scorer: &Scorer, ) -> Result { - let scorer_guard = scorer.lock().await; find_route( source, &RouteParameters { @@ -657,11 +657,11 @@ async fn find_payment_route( pathfinding_graph, None, &WrappedLog {}, - &scorer_guard, + scorer, &Default::default(), &[0; 32], ) - .map_err(|e| SimulationError::SimulatedNetworkError(e.err)) + .map_err(|e| SimulationError::SimulatedNetworkError(e.to_string())) } #[async_trait] @@ -687,7 +687,35 @@ impl LightningNode for SimNode { let preimage = PaymentPreimage(rand::random()); let payment_hash = preimage.into(); - // Check for payment hash collision, failing the payment if we happen to repeat one. + // Pathfinding dominates the cost of a payment and is pure CPU work, so run it on the blocking pool rather + // than on the scheduler thread. Routes for different payments are then computed in parallel, using the stored + // scorer under a shared read guard. + // + // On a paused runtime, virtual time cannot advance while a blocking task is outstanding, so the simulation + // clock still only moves once every route in progress has been computed. + let route = { + let source = self.info.pubkey; + let pathfinding_graph = self.pathfinding_graph.clone(); + let scorer = self.scorer.clone(); + + tokio::task::spawn_blocking(move || -> Result<_, LightningError> { + let scorer = scorer.read().map_err(|e| { + LightningError::SendPaymentError(format!("scorer lock poisoned: {e}")) + })?; + Ok(find_payment_route( + &source, + dest, + amount_msat, + &pathfinding_graph, + &scorer, + )) + }) + .await + .map_err(|e| { + LightningError::SendPaymentError(format!("pathfinding task failed: {e}")) + })?? + }; + let mut in_flight_guard = self.in_flight.lock().await; let entry = match in_flight_guard.entry(payment_hash) { Entry::Occupied(_) => { @@ -698,16 +726,7 @@ impl LightningNode for SimNode { Entry::Vacant(vacant) => vacant, }; - // Use the stored scorer when finding a route - let route = match find_payment_route( - &self.info.pubkey, - dest, - amount_msat, - &self.pathfinding_graph, - &self.scorer, - ) - .await - { + let route = match route { Ok(path) => path, // In the case that we can't find a route for the payment, we still report a successful payment *api call* // and report RouteNotFound to the tracking channel. This mimics the behavior of real nodes. @@ -784,10 +803,13 @@ impl LightningNode for SimNode { }; match &in_flight.path { Some(path) => { + let mut scorer = self.scorer.write().map_err(|e| { + LightningError::TrackPaymentError(format!("scorer lock poisoned: {e}")) + })?; if payment_result.payment_outcome == PaymentOutcome::Success { - self.scorer.lock().await.payment_path_successful(path, duration); + scorer.payment_path_successful(path, duration); } else if let PaymentOutcome::IndexFailure(index) = payment_result.payment_outcome { - self.scorer.lock().await.payment_path_failed(path, index as u64, duration); + scorer.payment_path_failed(path, index as u64, duration); } }, None => { @@ -1122,20 +1144,20 @@ pub async fn ln_node_from_graph( graph: Arc>, routing_graph: Arc, clock: Arc, -) -> Result>>>, LightningError> { +) -> Result>>, LightningError> { let sim_graph = graph.lock().await; - let mut nodes: HashMap>>> = + let mut nodes: HashMap>> = HashMap::with_capacity(sim_graph.nodes.len()); for node in sim_graph.nodes.iter() { nodes.insert( *node.0, - Arc::new(Mutex::new(SimNode::new( + Arc::new(SimNode::new( node.1 .0.clone(), graph.clone(), routing_graph.clone(), clock.clone(), - )?)), + )?), ); } @@ -1179,7 +1201,7 @@ pub fn populate_network_graph( &channel.node_1.policy.pubkey, &channel.node_2.policy.pubkey, ) - .to_v0_p2wsh(), + .to_p2wsh(), }; graph.update_channel_from_unsigned_announcement(&announcement, &Some(&utxo_validator))?; @@ -1189,9 +1211,11 @@ pub fn populate_network_graph( chain_hash, short_channel_id: channel.short_channel_id.into(), timestamp: now, + // Only the must_be_one bit is defined for message_flags. + message_flags: 1, // The least significant bit of the channel flag field represents the direction that the channel update // applies to. This value is interpreted as node_1 if it is zero, and node_2 otherwise. - flags: i as u8, + channel_flags: i as u8, cltv_expiry_delta: node.policy.cltv_expiry_delta as u16, htlc_minimum_msat: node.policy.min_htlc_size_msat, htlc_maximum_msat: node.policy.max_htlc_size_msat, @@ -1402,7 +1426,7 @@ async fn add_htlcs( let request = InterceptRequest { forwarding_node: hop.pubkey, payment_hash, - incoming_htlc: incoming_htlc.clone(), + incoming_htlc, incoming_custom_records, outgoing_channel_id: next_scid, incoming_amount_msat: outgoing_amount, @@ -1626,7 +1650,7 @@ struct UtxoValidator { impl UtxoLookup for UtxoValidator { fn get_utxo(&self, _genesis_hash: &ChainHash, _short_channel_id: u64) -> UtxoResult { UtxoResult::Sync(Ok(TxOut { - value: self.amount_sat, + value: Amount::from_sat(self.amount_sat), script_pubkey: self.script.clone(), })) } @@ -2026,14 +2050,14 @@ mod tests { assert!(nodes.len() == 3); - let node_1 = nodes.get(&pk1).unwrap().lock().await; + let node_1 = nodes.get(&pk1).unwrap(); let node_1_capacity = node_1.channel_capacities().await.unwrap(); // Node 1 has 2 channels but one was excluded so here we should only have the capacity of // the channel that was not excluded. assert!(node_1_capacity == capacity_1); - let node_2 = nodes.get(&pk2).unwrap().lock().await; + let node_2 = nodes.get(&pk2).unwrap(); let node_2_capacity = node_2.channel_capacities().await.unwrap(); assert!(node_2_capacity == capacity_1); @@ -2041,7 +2065,7 @@ mod tests { // present because its only channel was excluded. let node_3 = nodes.get(&pk3); assert!(node_3.is_some()); - let node_3 = node_3.unwrap().lock().await; + let node_3 = node_3.unwrap(); assert!(node_3.channel_capacities().await.unwrap() == 0); } @@ -2329,7 +2353,7 @@ mod tests { let test_kit = DispatchPaymentTestKit::new(chan_capacity, vec![], CustomRecords::default()).await; - let mut node = SimNode::new( + let node = SimNode::new( node_info(test_kit.nodes[0], String::default()), Arc::new(Mutex::new(test_kit.graph)), test_kit.routing_graph.clone(), @@ -2402,7 +2426,7 @@ mod tests { graph: SimGraph, nodes: Vec, routing_graph: Arc, - scorer: Mutex, Arc>>, + scorer: Scorer, shutdown: (Trigger, Listener), } @@ -2428,11 +2452,11 @@ mod tests { .unwrap(), ); - let scorer = Mutex::new(ProbabilisticScorer::new( + let scorer = ProbabilisticScorer::new( ProbabilisticScoringDecayParameters::default(), routing_graph.clone(), Arc::new(WrappedLog {}), - )); + ); // Collect pubkeys in-order, pushing the last node on separately because they don't have an outgoing // channel (they are not node_1 in any channel, only node_2). @@ -2498,9 +2522,8 @@ mod tests { dest: PublicKey, amt: u64, ) -> (Route, Result) { - let route = find_payment_route(&source, dest, amt, &self.routing_graph, &self.scorer) - .await - .unwrap(); + let route = + find_payment_route(&source, dest, amt, &self.routing_graph, &self.scorer).unwrap(); let (sender, receiver) = oneshot::channel(); self.graph @@ -2704,7 +2727,7 @@ mod tests { let test_kit = DispatchPaymentTestKit::new(chan_capacity, vec![], CustomRecords::default()).await; - let mut node = SimNode::new( + let node = SimNode::new( node_info(test_kit.nodes[0], String::default()), Arc::new(Mutex::new(test_kit.graph)), test_kit.routing_graph.clone(), diff --git a/simln-lib/src/test_utils.rs b/simln-lib/src/test_utils.rs index 2760ea23..4f3b72a6 100644 --- a/simln-lib/src/test_utils.rs +++ b/simln-lib/src/test_utils.rs @@ -2,14 +2,13 @@ use async_trait::async_trait; use bitcoin::secp256k1::{PublicKey, Secp256k1, SecretKey}; use bitcoin::Network; -use lightning::ln::features::Features; +use lightning::types::features::Features; use mockall::mock; use rand::distributions::Uniform; use rand::Rng; use std::collections::HashMap; use std::time::SystemTime; use std::{fmt, sync::Arc, time::Duration}; -use tokio::sync::Mutex; use tokio_util::task::TaskTracker; use crate::clock::SimulationClock; @@ -83,10 +82,10 @@ mock! { &self, dest: bitcoin::secp256k1::PublicKey, amount_msat: u64, - ) -> Result; + ) -> Result; async fn track_payment( &self, - hash: &lightning::ln::PaymentHash, + hash: &lightning::types::payment::PaymentHash, shutdown: triggered::Listener, ) -> Result; async fn get_node_info(&self, node_id: &PublicKey) -> Result; @@ -98,13 +97,13 @@ mock! { /// Type alias for the result of setup_test_nodes. pub struct TestNodesResult { pub nodes: Vec, - pub clients: Vec>>, + pub clients: Vec>, } impl TestNodesResult { // Returns a hashmap of the mocked lightning clients, cast to dyn LightningNode. - pub fn get_client_hashmap(&self) -> HashMap>> { - let mut client_map: HashMap>> = + pub fn get_client_hashmap(&self) -> HashMap> { + let mut client_map: HashMap> = HashMap::with_capacity(self.nodes.len()); for (idx, node) in self.nodes.iter().enumerate() { @@ -214,7 +213,7 @@ impl LightningTestNodeBuilder { mock_node.expect_get_network().return_const(network); } - clients.push(Arc::new(Mutex::new(mock_node))); + clients.push(Arc::new(mock_node)); nodes.push(node_info); } @@ -225,7 +224,7 @@ impl LightningTestNodeBuilder { /// Creates a new simulation with the given clients and activity definitions. /// Note: This sets a runtime for the simulation of 0, so run() will exit immediately. pub fn create_simulation( - clients: HashMap>>, + clients: HashMap>, ) -> Simulation { let (shutdown_trigger, shutdown_listener) = triggered::trigger(); Simulation::new(