Kaynağa Gözat

bin: Apply changes for the P2P API modification.

parazyd 3 yıl önce
ebeveyn
işleme
3b0e9ecb59

+ 7 - 7
bin/darkfid/src/main.rs

@@ -267,7 +267,7 @@ impl Darkfid {
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
+async fn realmain(args: Args, ex: Arc<smol::Executor<'static>>) -> Result<()> {
     if args.consensus && args.clock_sync {
         // We verify that if peer/seed nodes are configured, their rpc config also exists
         if ((!args.consensus_p2p_peer.is_empty() && args.consensus_peer_rpc.is_empty()) ||
@@ -359,7 +359,7 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
             ..Default::default()
         };
 
-        let p2p = net::P2p::new(sync_network_settings).await;
+        let p2p = net::P2p::new(sync_network_settings, ex.clone()).await;
         let registry = p2p.protocol_registry();
 
         let _state = state.clone();
@@ -401,7 +401,7 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
                 localnet: args.localnet,
                 ..Default::default()
             };
-            let p2p = net::P2p::new(consensus_network_settings).await;
+            let p2p = net::P2p::new(consensus_network_settings, ex.clone()).await;
             let registry = p2p.protocol_registry();
 
             let _state = state.clone();
@@ -445,9 +445,9 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
     );
 
     info!("Starting sync P2P network");
-    sync_p2p.clone().unwrap().start(ex.clone()).await?;
+    sync_p2p.clone().unwrap().start().await?;
     StoppableTask::new().start(
-        sync_p2p.clone().unwrap().run(ex.clone()),
+        sync_p2p.clone().unwrap().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }
@@ -471,9 +471,9 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
     let proposal_task = if args.consensus && *darkfid.synced.lock().await {
         info!("Starting consensus P2P network");
         let consensus_p2p = consensus_p2p.clone().unwrap();
-        consensus_p2p.clone().start(ex.clone()).await?;
+        consensus_p2p.clone().start().await?;
         StoppableTask::new().start(
-            consensus_p2p.clone().run(ex.clone()),
+            consensus_p2p.clone().run(),
             |res| async {
                 match res {
                     Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 11 - 7
bin/darkfid2/src/main.rs

@@ -126,7 +126,7 @@ impl Darkfid {
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
+async fn realmain(args: Args, ex: Arc<smol::Executor<'static>>) -> Result<()> {
     info!(target: "darkfid", "Initializing DarkFi node...");
 
     if args.testing_mode {
@@ -165,11 +165,15 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
     }
 
     // Initialize syncing P2P network
-    let sync_p2p = spawn_sync_p2p(&args.sync_net.into(), &validator, &subscribers).await;
+    let sync_p2p =
+        spawn_sync_p2p(&args.sync_net.into(), &validator, &subscribers, ex.clone()).await;
 
     // Initialize consensus P2P network
     let consensus_p2p = if args.consensus {
-        Some(spawn_consensus_p2p(&args.consensus_net.into(), &validator, &subscribers).await)
+        Some(
+            spawn_consensus_p2p(&args.consensus_net.into(), &validator, &subscribers, ex.clone())
+                .await,
+        )
     } else {
         None
     };
@@ -200,9 +204,9 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
     );
 
     info!(target: "darkfid", "Starting sync P2P network");
-    sync_p2p.clone().start(ex.clone()).await?;
+    sync_p2p.clone().start().await?;
     StoppableTask::new().start(
-        sync_p2p.clone().run(ex.clone()),
+        sync_p2p.clone().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }
@@ -217,9 +221,9 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
     if args.consensus {
         info!("Starting consensus P2P network");
         let consensus_p2p = consensus_p2p.clone().unwrap();
-        consensus_p2p.clone().start(ex.clone()).await?;
+        consensus_p2p.clone().start().await?;
         StoppableTask::new().start(
-            consensus_p2p.run(ex.clone()),
+            consensus_p2p.run(),
             |res| async {
                 match res {
                     Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 1 - 4
bin/darkfid2/src/task/sync.rs

@@ -16,10 +16,7 @@
  * along with this program.  If not, see <https://www.gnu.org/licenses/>.
  */
 
-use darkfi::{
-    util::{async_util::sleep, encoding::base64},
-    Result,
-};
+use darkfi::{system::sleep, util::encoding::base64, Result};
 use darkfi_serial::serialize;
 use log::{debug, info, warn};
 use tinyjson::JsonValue;

+ 8 - 8
bin/darkfid2/src/tests/harness.rs

@@ -61,7 +61,7 @@ pub struct Harness {
 }
 
 impl Harness {
-    pub async fn new(config: HarnessConfig, ex: &Arc<smol::Executor<'_>>) -> Result<Self> {
+    pub async fn new(config: HarnessConfig, ex: &Arc<smol::Executor<'static>>) -> Result<Self> {
         // Use test harness to generate genesis transactions
         let mut th = TestHarness::new(&["money".to_string(), "consensus".to_string()]).await?;
         let (genesis_stake_tx, _) = th.genesis_stake(&Holder::Alice, config.alice_initial)?;
@@ -217,7 +217,7 @@ pub async fn generate_node(
     config: &ValidatorConfig,
     sync_settings: &Settings,
     consensus_settings: Option<&Settings>,
-    ex: &Arc<smol::Executor<'_>>,
+    ex: &Arc<smol::Executor<'static>>,
     skip_sync: bool,
 ) -> Result<Darkfid> {
     let sled_db = sled::Config::new().temporary(true).open()?;
@@ -232,17 +232,17 @@ pub async fn generate_node(
         subscribers.insert("proposals", JsonSubscriber::new("blockchain.subscribe_proposals"));
     }
 
-    let sync_p2p = spawn_sync_p2p(&sync_settings, &validator, &subscribers).await;
+    let sync_p2p = spawn_sync_p2p(&sync_settings, &validator, &subscribers, ex.clone()).await;
     let consensus_p2p = if let Some(settings) = consensus_settings {
-        Some(spawn_consensus_p2p(settings, &validator, &subscribers).await)
+        Some(spawn_consensus_p2p(settings, &validator, &subscribers, ex.clone()).await)
     } else {
         None
     };
     let node = Darkfid::new(sync_p2p.clone(), consensus_p2p.clone(), validator, subscribers).await;
 
-    sync_p2p.clone().start(ex.clone()).await?;
+    sync_p2p.clone().start().await?;
     StoppableTask::new().start(
-        sync_p2p.run(ex.clone()),
+        sync_p2p.run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }
@@ -255,9 +255,9 @@ pub async fn generate_node(
 
     if consensus_settings.is_some() {
         let consensus_p2p = consensus_p2p.unwrap();
-        consensus_p2p.clone().start(ex.clone()).await?;
+        consensus_p2p.clone().start().await?;
         StoppableTask::new().start(
-            consensus_p2p.run(ex.clone()),
+            consensus_p2p.run(),
             |res| async {
                 match res {
                     Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 1 - 1
bin/darkfid2/src/tests/mod.rs

@@ -27,7 +27,7 @@ use harness::{generate_node, Harness, HarnessConfig};
 
 mod forks;
 
-async fn sync_blocks_real(ex: Arc<Executor<'_>>) -> Result<()> {
+async fn sync_blocks_real(ex: Arc<Executor<'static>>) -> Result<()> {
     init_logger();
 
     // Initialize harness in testing mode

+ 6 - 3
bin/darkfid2/src/utils.rs

@@ -16,9 +16,10 @@
  * along with this program.  If not, see <https://www.gnu.org/licenses/>.
  */
 
-use std::collections::HashMap;
+use std::{collections::HashMap, sync::Arc};
 
 use log::info;
+use smol::Executor;
 
 use darkfi::{
     error::TxVerifyFailed,
@@ -76,9 +77,10 @@ pub async fn spawn_sync_p2p(
     settings: &Settings,
     validator: &ValidatorPtr,
     subscribers: &HashMap<&'static str, JsonSubscriber>,
+    executor: Arc<Executor<'static>>,
 ) -> P2pPtr {
     info!(target: "darkfid", "Registering sync network P2P protocols...");
-    let p2p = P2p::new(settings.clone()).await;
+    let p2p = P2p::new(settings.clone(), executor.clone()).await;
     let registry = p2p.protocol_registry();
 
     let _validator = validator.clone();
@@ -117,9 +119,10 @@ pub async fn spawn_consensus_p2p(
     settings: &Settings,
     validator: &ValidatorPtr,
     subscribers: &HashMap<&'static str, JsonSubscriber>,
+    executor: Arc<Executor<'static>>,
 ) -> P2pPtr {
     info!(target: "darkfid", "Registering consensus network P2P protocols...");
-    let p2p = P2p::new(settings.clone()).await;
+    let p2p = P2p::new(settings.clone(), executor.clone()).await;
     let registry = p2p.protocol_registry();
 
     let _validator = validator.clone();

+ 6 - 6
bin/darkirc/src/main.rs

@@ -41,8 +41,8 @@ use darkfi::{
     },
     net,
     rpc::{jsonrpc::JsonSubscriber, server::listen_and_serve},
-    system::{StoppableTask, Subscriber, SubscriberPtr},
-    util::{async_util::sleep, file::save_json_file, path::expand_path, time::Timestamp},
+    system::{sleep, StoppableTask, Subscriber, SubscriberPtr},
+    util::{file::save_json_file, path::expand_path, time::Timestamp},
     Error, Result,
 };
 
@@ -95,7 +95,7 @@ async fn remove_old_events(model: ModelPtr<PrivMsgEvent>) -> Result<()> {
 }
 
 async_daemonize!(realmain);
-async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
+async fn realmain(settings: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
     let datastore_path = expand_path(&settings.datastore)?;
 
     // mkdir datastore_path if not exists
@@ -194,7 +194,7 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
     let net_settings = settings.net.clone();
 
     // New p2p
-    let p2p = net::P2p::new(net_settings.into()).await;
+    let p2p = net::P2p::new(net_settings.into(), executor.clone()).await;
 
     // Register the protocol_event
     let registry = p2p.protocol_registry();
@@ -263,9 +263,9 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
     // Start P2P network
     ////////////////////
     info!(target: "darkirc", "Starting P2P network");
-    p2p.clone().start(executor.clone()).await?;
+    p2p.clone().start().await?;
     StoppableTask::new().start(
-        p2p.clone().run(executor.clone()),
+        p2p.clone().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 6 - 6
bin/faucetd/src/main.rs

@@ -71,9 +71,9 @@ use darkfi::{
         server::{listen_and_serve, RequestHandler},
     },
     runtime::vm_runtime::SMART_CONTRACT_ZKAS_DB_NAME,
-    system::StoppableTask,
+    system::{sleep, StoppableTask},
     tx::Transaction,
-    util::{async_util::sleep, parse::decode_base10, path::expand_path},
+    util::{parse::decode_base10, path::expand_path},
     wallet::{WalletDb, WalletPtr},
     zk::{proof::ProvingKey, vm::ZkCircuit, vm_heap::empty_witnesses},
     zkas::ZkBinary,
@@ -627,7 +627,7 @@ async fn prune_airdrop_maps(
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
+async fn realmain(args: Args, ex: Arc<smol::Executor<'static>>) -> Result<()> {
     // Initialize or load wallet
     let wallet = WalletDb::new(Some(expand_path(&args.wallet_path)?), Some(&args.wallet_pass))?;
 
@@ -696,7 +696,7 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
         ..Default::default()
     };
 
-    let sync_p2p = net::P2p::new(network_settings).await;
+    let sync_p2p = net::P2p::new(network_settings, ex.clone()).await;
     let registry = sync_p2p.protocol_registry();
 
     info!(target: "faucetd", "Registering block sync P2P protocols...");
@@ -764,9 +764,9 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
     );
 
     info!(target: "faucetd", "Starting sync P2P network");
-    sync_p2p.clone().start(ex.clone()).await?;
+    sync_p2p.clone().start().await?;
     StoppableTask::new().start(
-        sync_p2p.clone().run(ex.clone()),
+        sync_p2p.clone().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 4 - 4
bin/fud/fud/src/main.rs

@@ -519,7 +519,7 @@ async fn fetch_chunk_task(fud: Arc<Fud>, executor: Arc<Executor<'_>>) -> Result<
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, ex: Arc<Executor<'_>>) -> Result<()> {
+async fn realmain(args: Args, ex: Arc<Executor<'static>>) -> Result<()> {
     // The working directory for this daemon and geode.
     let basedir = expand_path(&args.base_dir)?;
 
@@ -531,7 +531,7 @@ async fn realmain(args: Args, ex: Arc<Executor<'_>>) -> Result<()> {
     let geode = Geode::new(&basedir.into()).await?;
 
     info!("Instantiating P2P network");
-    let p2p = P2p::new(args.net.into()).await;
+    let p2p = P2p::new(args.net.into(), ex.clone()).await;
 
     // Daemon instantiation
     let (file_fetch_tx, file_fetch_rx) = smol::channel::unbounded();
@@ -598,9 +598,9 @@ async fn realmain(args: Args, ex: Arc<Executor<'_>>) -> Result<()> {
             async move { ProtocolFud::init(fud_, channel, p2p).await.unwrap() }
         })
         .await;
-    p2p.clone().start(ex.clone()).await?;
+    p2p.clone().start().await?;
     StoppableTask::new().start(
-        p2p.clone().run(ex.clone()),
+        p2p.clone().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 4 - 7
bin/genev/genevd/src/main.rs

@@ -85,7 +85,7 @@ async fn start_sync_loop(
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
+async fn realmain(args: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
     ////////////////////
     // Initialize the base structures
     ////////////////////
@@ -105,7 +105,7 @@ async fn realmain(args: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
     let net_settings = args.net.clone();
 
     // New p2p
-    let p2p = net::P2p::new(net_settings.into()).await;
+    let p2p = net::P2p::new(net_settings.into(), executor.clone()).await;
     let p2p2 = p2p.clone();
 
     // Register the protocol_event
@@ -119,14 +119,11 @@ async fn realmain(args: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
         })
         .await;
 
-    // Start
-    p2p.clone().start(executor.clone()).await?;
-
     // Run
     info!(target: "genevd", "Starting P2P network");
-    p2p.clone().start(executor.clone()).await?;
+    p2p.clone().start().await?;
     StoppableTask::new().start(
-        p2p.clone().run(executor.clone()),
+        p2p.clone().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }

+ 6 - 7
bin/lilith/src/main.rs

@@ -41,9 +41,8 @@ use darkfi::{
         jsonrpc::*,
         server::{listen_and_serve, RequestHandler},
     },
-    system::StoppableTask,
+    system::{sleep, StoppableTask},
     util::{
-        async_util::sleep,
         file::{load_file, save_file},
         path::{expand_path, get_config_path},
     },
@@ -352,7 +351,7 @@ async fn spawn_net(
     info: &NetInfo,
     accept_addrs: &[Url],
     saved_hosts: &HashSet<Url>,
-    ex: Arc<Executor<'_>>,
+    ex: Arc<Executor<'static>>,
 ) -> Result<Spawn> {
     let mut listen_urls = vec![];
 
@@ -384,7 +383,7 @@ async fn spawn_net(
     };
 
     // Create P2P instance
-    let p2p = P2p::new(settings).await;
+    let p2p = P2p::new(settings, ex.clone()).await;
 
     // Fill db with cached hosts
     let hosts: Vec<Url> = saved_hosts.iter().cloned().collect();
@@ -392,10 +391,10 @@ async fn spawn_net(
 
     let addrs_str: Vec<&str> = listen_urls.iter().map(|x| x.as_str()).collect();
     info!(target: "lilith", "Starting seed network node for \"{}\" on {:?}", name, addrs_str);
-    p2p.clone().start(ex.clone()).await?;
+    p2p.clone().start().await?;
     let name_ = name.clone();
     StoppableTask::new().start(
-        p2p.clone().run(ex.clone()),
+        p2p.clone().run(),
         |res| async move {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }
@@ -411,7 +410,7 @@ async fn spawn_net(
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, ex: Arc<Executor<'_>>) -> Result<()> {
+async fn realmain(args: Args, ex: Arc<Executor<'static>>) -> Result<()> {
     // Pick up network settings from the TOML config
     let cfg_path = get_config_path(args.config, CONFIG_FILE)?;
     let toml_contents = std::fs::read_to_string(cfg_path)?;

+ 4 - 6
bin/tau/taud/src/main.rs

@@ -253,7 +253,7 @@ async fn on_receive_task(
 }
 
 async_daemonize!(realmain);
-async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
+async fn realmain(settings: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
     let datastore_path = expand_path(&settings.datastore)?;
 
     let nickname =
@@ -348,7 +348,7 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
     //
     let net_settings = settings.net.clone();
 
-    let p2p = net::P2p::new(net_settings.into()).await;
+    let p2p = net::P2p::new(net_settings.into(), executor.clone()).await;
     let registry = p2p.protocol_registry();
 
     registry
@@ -360,12 +360,10 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
         })
         .await;
 
-    p2p.clone().start(executor.clone()).await?;
-
     info!(target: "taud", "Starting P2P network");
-    p2p.clone().start(executor.clone()).await?;
+    p2p.clone().start().await?;
     StoppableTask::new().start(
-        p2p.clone().run(executor.clone()),
+        p2p.clone().run(),
         |res| async {
             match res {
                 Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }