Prechádzať zdrojové kódy

make example/p2pdebug crate for debugging/testing the p2p network code

ghassmo 4 rokov pred
rodič
commit
394daad1f7

+ 2 - 0
example/p2pdebug/.gitignore

@@ -0,0 +1,2 @@
+target/*
+Cargo.lock

+ 32 - 0
example/p2pdebug/Cargo.toml

@@ -0,0 +1,32 @@
+[package]
+name = "p2pdebug"
+version = "0.3.0"
+homepage = "https://dark.fi"
+repository = "https://github.com/darkrenaissance/darkfi"
+license = "AGPL-3.0-only"
+edition = "2021"
+
+[workspace]
+
+[dependencies]
+darkfi = {path = "../../", features = ["net"]}
+# Async
+smol = "1.2.5"
+futures = "0.3.21"
+async-std = "1.11.0"
+async-trait = "0.1.53"
+async-channel = "1.6.1"
+async-executor = "1.4.1"
+easy-parallel = "3.2.0"
+
+# Crypto
+rand = "0.8.5"
+
+# Misc
+clap = {version = "3.1.8", features = ["derive"]}
+log = "0.4.16"
+simplelog = "0.12.0-alpha1"
+fxhash = "0.2.1"
+
+# Encoding and parsing
+serde_json = "1.0.79"

+ 214 - 0
example/p2pdebug/src/main.rs

@@ -0,0 +1,214 @@
+use std::{net::SocketAddr, sync::Arc};
+
+use async_executor::Executor;
+use clap::Parser;
+use easy_parallel::Parallel;
+use simplelog::{ColorChoice, TermLogger, TerminalMode};
+
+use rand::{rngs::OsRng, Rng, RngCore};
+
+use darkfi::{
+    cli_desc, net,
+    util::{cli::log_config, sleep},
+    Result,
+};
+
+pub(crate) mod proto;
+
+use crate::proto::debugmsg::{Debugmsg, ProtocolDebugmsg, SeenDebugmsgIds};
+
+#[derive(Parser)]
+#[clap(name = "p2pdebugging", about = cli_desc!(), version)]
+struct Args {
+    /// Verbosity level
+    #[clap(short, parse(from_occurrences))]
+    verbose: u8,
+    /// node number:  
+    /// 0-2 is for seed nodes
+    /// 3-20 is for inbound connections nodes
+    /// 21- is for outbound connections nodes
+    #[clap(short, long, default_value = "0")]
+    node: u8,
+    /// broadcast messages
+    #[clap(short, long)]
+    broadcast: bool,
+}
+
+#[derive(Debug, Clone)]
+enum State {
+    Seed,
+    Inbound,
+    Outbound,
+}
+
+struct MockP2p {
+    node_number: u8,
+    state: State,
+    p2p: net::P2pPtr,
+    broadcast: bool,
+    address: Option<SocketAddr>,
+}
+
+impl MockP2p {
+    async fn new(node_number: u8, _broadcast: bool) -> Result<Self> {
+        let seed_addrs: Vec<SocketAddr> = vec![
+            "127.0.0.1:11001".parse()?,
+            "127.0.0.1:11002".parse()?,
+            "127.0.0.1:11003".parse()?,
+        ];
+
+        let state: State;
+        let address: Option<SocketAddr>;
+
+        let mut broadcast = _broadcast;
+
+        let p2p = match node_number {
+            0..=2 => {
+                address = Some(seed_addrs[node_number as usize]);
+
+                let net_settings = net::Settings { inbound: address, ..Default::default() };
+                let p2p = net::P2p::new(net_settings).await;
+
+                broadcast = false;
+                state = State::Seed;
+
+                p2p
+            }
+            3..=20 => {
+                let random_port: u32 = rand::thread_rng().gen_range(11007..49151);
+                address = Some(format!("127.0.0.1:{}", random_port).parse()?);
+
+                let net_settings = net::Settings {
+                    inbound: address,
+                    external_addr: address,
+                    seeds: seed_addrs,
+                    ..Default::default()
+                };
+
+                let p2p = net::P2p::new(net_settings).await;
+
+                state = State::Inbound;
+
+                p2p
+            }
+            _ => {
+                address = None;
+
+                let net_settings = net::Settings {
+                    outbound_connections: 3,
+                    seeds: seed_addrs,
+                    ..Default::default()
+                };
+
+                let p2p = net::P2p::new(net_settings).await;
+                state = State::Outbound;
+
+                p2p
+            }
+        };
+
+        println!("start {:?} node #{} address {:?}", state, node_number, address);
+
+        Ok(Self { node_number, state, p2p, broadcast, address })
+    }
+
+    async fn run(&self, executor: Arc<Executor<'_>>) -> Result<()> {
+        let p2p = self.p2p.clone();
+        let state = self.state.clone();
+        let node_number = self.node_number;
+        let address = self.address;
+
+        let (sender, receiver) = async_channel::unbounded();
+        let sender_clone = sender.clone();
+
+        let seen_debugmsg_ids = SeenDebugmsgIds::new();
+        let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
+
+        let registry = p2p.protocol_registry();
+        registry
+            .register(!net::SESSION_SEED, move |channel, p2p| {
+                let sender = sender_clone.clone();
+                let seen_debugmsg_ids = seen_debugmsg_ids_clone.clone();
+                async move { ProtocolDebugmsg::init(channel, sender, seen_debugmsg_ids, p2p).await }
+            })
+            .await;
+
+        if self.broadcast {
+            println!("start broadcast {:?} node #{} address {:?}", state, node_number, address);
+            let sleep_time = 10;
+            let p2p_clone = p2p.clone();
+            let executor_clone = executor.clone();
+            let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
+            executor_clone
+                .spawn(async move {
+                    loop {
+                        sleep(sleep_time).await;
+
+                        println!(
+                            "broadcast sleep for {} {:?} node #{} address {:?}",
+                            sleep_time, state, node_number, address
+                        );
+
+                        let random_id = OsRng.next_u32();
+                        seen_debugmsg_ids_clone.add_seen(random_id).await;
+
+                        let msg = Debugmsg { id: random_id, message: "hello".to_string() };
+
+                        println!(
+                            "send {:?} {:?} node #{} address {:?}",
+                            msg, state, node_number, address
+                        );
+
+                        p2p_clone.broadcast(msg).await.unwrap();
+                    }
+                })
+                .detach();
+        }
+
+        let state = self.state.clone();
+        let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
+        executor
+            .spawn(async move {
+                loop {
+                    let msg = receiver.recv().await.unwrap();
+                    println!(
+                        "receive {:?} {:?} node #{} address {:?}",
+                        msg, state, node_number, address
+                    );
+                    seen_debugmsg_ids_clone.add_seen(msg.id).await;
+                }
+            })
+            .detach();
+
+        p2p.clone().start(executor.clone()).await?;
+        p2p.run(executor).await
+    }
+}
+
+async fn start(executor: Arc<Executor<'_>>, args: Args) -> Result<()> {
+    let mock_p2p = MockP2p::new(args.node, args.broadcast).await?;
+    mock_p2p.run(executor).await
+}
+
+fn main() -> Result<()> {
+    let args = Args::parse();
+
+    let (lvl, conf) = log_config(args.verbose.into())?;
+    TermLogger::init(lvl, conf, TerminalMode::Mixed, ColorChoice::Auto)?;
+
+    let ex = Arc::new(Executor::new());
+    let ex_clone = ex.clone();
+    let (signal, shutdown) = async_channel::unbounded::<()>();
+    let (_, result) = Parallel::new()
+        .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
+        // Run the main future on the current thread.
+        .finish(|| {
+            smol::future::block_on(async move {
+                start(ex_clone.clone(), args).await?;
+                drop(signal);
+                Ok::<(), darkfi::Error>(())
+            })
+        });
+
+    result
+}

+ 122 - 0
example/p2pdebug/src/proto/debugmsg.rs

@@ -0,0 +1,122 @@
+use std::sync::Arc;
+
+use async_channel::Sender;
+use async_executor::Executor;
+use async_std::sync::Mutex;
+use async_trait::async_trait;
+use fxhash::FxHashSet;
+use log::debug;
+
+use darkfi::{
+    net,
+    util::serial::{SerialDecodable, SerialEncodable},
+    Result,
+};
+
+pub type DebugmsgId = u32;
+
+#[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
+pub struct Debugmsg {
+    pub id: DebugmsgId,
+    pub message: String,
+}
+
+impl net::Message for Debugmsg {
+    fn name() -> &'static str {
+        "debugmsg"
+    }
+}
+
+pub struct SeenDebugmsgIds {
+    ids: Mutex<FxHashSet<DebugmsgId>>,
+}
+
+pub type SeenDebugmsgIdsPtr = Arc<SeenDebugmsgIds>;
+
+impl SeenDebugmsgIds {
+    pub fn new() -> Arc<Self> {
+        Arc::new(Self { ids: Mutex::new(FxHashSet::default()) })
+    }
+
+    pub async fn add_seen(&self, id: u32) {
+        self.ids.lock().await.insert(id);
+    }
+
+    pub async fn is_seen(&self, id: u32) -> bool {
+        self.ids.lock().await.contains(&id)
+    }
+}
+
+pub struct ProtocolDebugmsg {
+    notify_queue_sender: Sender<Arc<Debugmsg>>,
+    debugmsg_sub: net::MessageSubscription<Debugmsg>,
+    jobsman: net::ProtocolJobsManagerPtr,
+    seen_ids: SeenDebugmsgIdsPtr,
+    p2p: net::P2pPtr,
+}
+
+#[async_trait]
+impl net::ProtocolBase for ProtocolDebugmsg {
+    /// Starts ping-pong keep-alive messages exchange. Runs ping-pong in the
+    /// protocol task manager, then queues the reply. Sends out a ping and
+    /// waits for pong reply. Waits for ping and replies with a pong.
+    async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
+        debug!(target: "ircd", "Protocoldebugmsg::start() [START]");
+        self.jobsman.clone().start(executor.clone());
+        self.jobsman.clone().spawn(self.clone().handle_receive_debugmsg(), executor.clone()).await;
+        debug!(target: "ircd", "ProtocolDebugmsg::start() [END]");
+        Ok(())
+    }
+
+    fn name(&self) -> &'static str {
+        "Protocoldebugmsg"
+    }
+}
+
+impl ProtocolDebugmsg {
+    pub async fn init(
+        channel: net::ChannelPtr,
+        notify_queue_sender: Sender<Arc<Debugmsg>>,
+        seen_ids: SeenDebugmsgIdsPtr,
+        p2p: net::P2pPtr,
+    ) -> net::ProtocolBasePtr {
+        let message_subsystem = channel.get_message_subsystem();
+        message_subsystem.add_dispatch::<Debugmsg>().await;
+
+        let sub = channel.subscribe_msg::<Debugmsg>().await.expect("Missing Debugmsg dispatcher!");
+
+        Arc::new(Self {
+            notify_queue_sender,
+            debugmsg_sub: sub,
+            jobsman: net::ProtocolJobsManager::new("DebugmsgProtocol", channel),
+            seen_ids,
+            p2p,
+        })
+    }
+
+    async fn handle_receive_debugmsg(self: Arc<Self>) -> Result<()> {
+        debug!(target: "ircd", "ProtocolDebugmsg::handle_receive_debugmsg() [START]");
+
+        loop {
+            let debugmsg = self.debugmsg_sub.receive().await?;
+
+            debug!(target: "ircd", "ProtocolDebugmsg::handle_receive_debugmsg() received {:?}", debugmsg);
+
+            // Do we already have this message?
+            if self.seen_ids.is_seen(debugmsg.id).await {
+                continue
+            }
+
+            self.seen_ids.add_seen(debugmsg.id).await;
+
+            // If not, then broadcast to network.
+            let debugmsg_copy = (*debugmsg).clone();
+            self.p2p.broadcast(debugmsg_copy).await?;
+
+            self.notify_queue_sender
+                .send(debugmsg)
+                .await
+                .expect("notify_queue_sender send failed!");
+        }
+    }
+}

+ 1 - 0
example/p2pdebug/src/proto/mod.rs

@@ -0,0 +1 @@
+pub(crate) mod debugmsg;

+ 6 - 0
example/p2pdebug/tmux_sessions.sh

@@ -0,0 +1,6 @@
+#!/bin/sh
+
+tmux new-session -d './target/release/p2pdebug'
+tmux split-window -v './target/release/p2pdebug -n 3 '
+tmux split-window -h './target/release/p2pdebug -n 21 -b'
+tmux attach