Răsfoiți Sursa

[stakeholder/stakeholder] fix blocking block sync

mohab 3 ani în urmă
părinte
comite
356002e575
2 a modificat fișierele cu 55 adăugiri și 53 ștergeri
  1. 28 25
      example/crypsinous.rs
  2. 27 28
      src/stakeholder/stakeholder.rs

+ 28 - 25
example/crypsinous.rs

@@ -1,4 +1,9 @@
-use ::darkfi::{blockchain::EpochConsensus, net::Settings, stakeholder::Stakeholder};
+use ::darkfi::{blockchain::EpochConsensus,
+               net::Settings,
+               stakeholder::Stakeholder,
+               util::time::Timestamp,
+               
+};
 
 use clap::Parser;
 use futures::executor::block_on;
@@ -8,9 +13,14 @@ use vec;
 
 #[derive(Parser)]
 struct NetCli {
+    #[clap(long,value_parser,default_value="tls://127.0.0.1:12003")]
     addr: String,
+    #[clap(long,value_parser,default_value="/tmp/db")]
     path: String,
+    #[clap(long,value_parser,default_value="tls://127.0.0.1:12004")]
     peers: Vec<String>,
+    #[clap(long,value_parser,default_value="tls://lilith.dark.fi:25551")]
+    seeds: Vec<String>,
 }
 
 #[async_std::main]
@@ -21,11 +31,10 @@ async fn main() {
     for i in 0..args.peers.len() {
         peers.push(Url::parse(args.peers[i].as_str()).unwrap());
     }
-    let seeds = [
-        Url::parse("tls://irc0.dark.fi:11001").unwrap(),
-        Url::parse("tls://irc1.dark.fi:11001").unwrap(),
-    ]
-    .to_vec();
+    let mut seeds = vec![];
+    for i in 0..args.seeds.len() {
+        seeds.push(Url::parse(args.seeds[i].as_str()).unwrap());
+    }
     let slots = 3;
     let epochs = 3;
     let ticks = 10;
@@ -47,26 +56,20 @@ async fn main() {
     };
     //proof's number of rows
     let k: u32 = 13;
-    let mut handles = vec![];
     let path = args.path;
-    for i in 0..2 {
-        let rel_path = format!("{}{}", path, i.to_string());
+    let id = Timestamp::current_time().0;
 
-        let mut stakeholder = block_on(Stakeholder::new(
-            epoch_consensus.clone(),
-            settings.clone(),
-            &rel_path,
-            i,
-            Some(k),
-        ))
+    let mut stakeholder = block_on(Stakeholder::new(
+        epoch_consensus.clone(),
+        settings.clone(),
+        &path,
+        id,
+        Some(k),
+    ))
         .unwrap();
-
-        let handle = thread::spawn(move || {
-            block_on(stakeholder.background(Some(9)));
-        });
-        handles.push(handle);
-    }
-    for handle in handles {
-        handle.join().unwrap();
-    }
+    
+    let handle = thread::spawn(move || {
+        block_on(stakeholder.background(Some(9)));
+    });
+    handle.join().unwrap();
 }

+ 27 - 28
src/stakeholder/stakeholder.rs

@@ -128,7 +128,7 @@ pub struct Stakeholder {
     pub vk: VerifyingKey,
     pub playing: bool,
     pub workspace: SlotWorkspace,
-    pub id: u8,
+    pub id: i64,
     pub keypair: Keypair,
     //pub subscription: Subscription<Result<ChannelPtr>>,
     //pub chanptr : ChannelPtr,
@@ -140,7 +140,7 @@ impl Stakeholder {
         consensus: EpochConsensus,
         settings: Settings,
         rel_path: &str,
-        id: u8,
+        id: i64,
         k: Option<u32>,
     ) -> Result<Self> {
         let path = expand_path(rel_path).unwrap();
@@ -228,7 +228,8 @@ impl Stakeholder {
         println!("runing p2p net");
         let exec = Arc::new(Executor::new());
         self.net.clone().start(exec.clone()).await?;
-        self.net.clone().run(exec).await?;
+        //TODO (fix) await blocks
+        self.net.clone().run(exec);
         println!("p2p net running...");
         Ok(())
     }
@@ -276,31 +277,29 @@ impl Stakeholder {
     /// validate the block proof, and the transactions,
     /// if so add the proof to metadata if stakeholder isn't the lead.
     pub async fn sync_block(&self) {
-        let subscription: Subscription<Result<ChannelPtr>> = self.net.subscribe_channel().await;
-        println!("--> channel");
-        let chanptr: ChannelPtr = subscription.receive().await.unwrap();
-        println!("--> received channel");
-        //
-        let message_subsytem = chanptr.get_message_subsystem();
-        println!("--> adding dispatcher to msg subsystem");
-        message_subsytem.add_dispatch::<BlockInfo>().await;
-        println!("--> added");
-        //TODO start channel if isn't started yet
-        //let info = chanptr.get_info();
-        //println!("channel info: {}", info);
-        println!("--> subscribe msg_sub");
-        let msg_sub: MessageSubscription<BlockInfo> =
-            chanptr.subscribe_msg::<BlockInfo>().await.expect("missing blockinfo");
-        println!("--> subscribed");
-
-        let res = msg_sub.receive().await.unwrap();
-        let blk: BlockInfo = (*res).to_owned();
-        //TODO validate the block proof, and transactions.
-        if self.valid_block(blk.clone()) {
-            //TODO if valid only.
-            let _len = self.blockchain.add(&[blk]);
-        } else {
-            println!("received block is invalid!");
+        for chanptr in self.net.channels().lock().await.values() {
+            //
+            let message_subsytem = chanptr.get_message_subsystem();
+            println!("--> adding dispatcher to msg subsystem");
+            message_subsytem.add_dispatch::<BlockInfo>().await;
+            println!("--> added");
+            //TODO start channel if isn't started yet
+            //let info = chanptr.get_info();
+            //println!("channel info: {}", info);
+            println!("--> subscribe msg_sub");
+            let msg_sub: MessageSubscription<BlockInfo> =
+                chanptr.subscribe_msg::<BlockInfo>().await.expect("missing blockinfo");
+            println!("--> subscribed");
+            
+            let res = msg_sub.receive().await.unwrap();
+            let blk: BlockInfo = (*res).to_owned();
+            //TODO validate the block proof, and transactions.
+            if self.valid_block(blk.clone()) {
+                //TODO if valid only.
+                let _len = self.blockchain.add(&[blk]);
+            } else {
+                println!("received block is invalid!");
+            }
         }
     }