main.rs 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212
  1. use std::{net::SocketAddr, sync::Arc};
  2. use async_executor::Executor;
  3. use clap::Parser;
  4. use easy_parallel::Parallel;
  5. use simplelog::{ColorChoice, TermLogger, TerminalMode};
  6. use rand::{rngs::OsRng, Rng, RngCore};
  7. use darkfi::{
  8. cli_desc, net,
  9. util::{cli::log_config, sleep},
  10. Result,
  11. };
  12. pub(crate) mod proto;
  13. use crate::proto::debugmsg::{Debugmsg, ProtocolDebugmsg, SeenDebugmsgIds};
  14. #[derive(Parser)]
  15. #[clap(name = "p2pdebugging", about = cli_desc!(), version)]
  16. struct Args {
  17. /// Verbosity level
  18. #[clap(short, parse(from_occurrences))]
  19. verbose: u8,
  20. /// node number:
  21. /// 0-2 is for seed nodes
  22. /// 3-20 is for inbound connections nodes
  23. /// 21- is for outbound connections nodes
  24. #[clap(short, long, default_value = "0")]
  25. node: u8,
  26. /// broadcast messages
  27. #[clap(short, long)]
  28. broadcast: bool,
  29. }
  30. #[derive(Debug, Clone)]
  31. enum State {
  32. Seed,
  33. Inbound,
  34. Outbound,
  35. }
  36. struct MockP2p {
  37. node_number: u8,
  38. state: State,
  39. p2p: net::P2pPtr,
  40. broadcast: bool,
  41. address: Option<SocketAddr>,
  42. }
  43. impl MockP2p {
  44. async fn new(node_number: u8, _broadcast: bool) -> Result<Self> {
  45. let seed_addrs: Vec<SocketAddr> = vec![
  46. "127.0.0.1:11001".parse()?,
  47. "127.0.0.1:11002".parse()?,
  48. "127.0.0.1:11003".parse()?,
  49. ];
  50. let state: State;
  51. let address: Option<SocketAddr>;
  52. let mut broadcast = _broadcast;
  53. let p2p = match node_number {
  54. 0..=2 => {
  55. address = Some(seed_addrs[node_number as usize]);
  56. let net_settings = net::Settings { inbound: address, ..Default::default() };
  57. let p2p = net::P2p::new(net_settings).await;
  58. broadcast = false;
  59. state = State::Seed;
  60. p2p
  61. }
  62. 3..=20 => {
  63. let random_port: u32 = rand::thread_rng().gen_range(11007..49151);
  64. address = Some(format!("127.0.0.1:{}", random_port).parse()?);
  65. let net_settings = net::Settings {
  66. inbound: address,
  67. external_addr: address,
  68. seeds: seed_addrs,
  69. ..Default::default()
  70. };
  71. let p2p = net::P2p::new(net_settings).await;
  72. state = State::Inbound;
  73. p2p
  74. }
  75. _ => {
  76. address = None;
  77. let net_settings = net::Settings {
  78. outbound_connections: 3,
  79. seeds: seed_addrs,
  80. ..Default::default()
  81. };
  82. let p2p = net::P2p::new(net_settings).await;
  83. state = State::Outbound;
  84. p2p
  85. }
  86. };
  87. println!("start {:?} node #{} address {:?}", state, node_number, address);
  88. Ok(Self { node_number, state, p2p, broadcast, address })
  89. }
  90. async fn run(&self, executor: Arc<Executor<'_>>) -> Result<()> {
  91. let p2p = self.p2p.clone();
  92. let state = self.state.clone();
  93. let node_number = self.node_number;
  94. let address = self.address;
  95. let (sender, receiver) = async_channel::unbounded();
  96. let sender_clone = sender.clone();
  97. let seen_debugmsg_ids = SeenDebugmsgIds::new();
  98. let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
  99. let registry = p2p.protocol_registry();
  100. registry
  101. .register(net::SESSION_ALL, move |channel, p2p| {
  102. let sender = sender_clone.clone();
  103. let seen_debugmsg_ids = seen_debugmsg_ids_clone.clone();
  104. async move { ProtocolDebugmsg::init(channel, sender, seen_debugmsg_ids, p2p).await }
  105. })
  106. .await;
  107. if self.broadcast {
  108. println!("start broadcast {:?} node #{} address {:?}", state, node_number, address);
  109. let sleep_time = 10;
  110. let p2p_clone = p2p.clone();
  111. let executor_clone = executor.clone();
  112. executor_clone
  113. .spawn(async move {
  114. loop {
  115. sleep(sleep_time).await;
  116. println!(
  117. "broadcast sleep for {} {:?} node #{} address {:?}",
  118. sleep_time, state, node_number, address
  119. );
  120. let random_id = OsRng.next_u32();
  121. let msg = Debugmsg { id: random_id, message: "hello".to_string() };
  122. println!(
  123. "send {:?} {:?} node #{} address {:?}",
  124. msg, state, node_number, address
  125. );
  126. p2p_clone.broadcast(msg).await.unwrap();
  127. }
  128. })
  129. .detach();
  130. }
  131. let state = self.state.clone();
  132. let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
  133. executor
  134. .spawn(async move {
  135. loop {
  136. let msg = receiver.recv().await.unwrap();
  137. println!(
  138. "receive {:?} {:?} node #{} address {:?}",
  139. msg, state, node_number, address
  140. );
  141. seen_debugmsg_ids_clone.add_seen(msg.id).await;
  142. }
  143. })
  144. .detach();
  145. p2p.clone().start(executor.clone()).await?;
  146. p2p.run(executor).await
  147. }
  148. }
  149. async fn start(executor: Arc<Executor<'_>>, args: Args) -> Result<()> {
  150. let mock_p2p = MockP2p::new(args.node, args.broadcast).await?;
  151. mock_p2p.run(executor).await
  152. }
  153. fn main() -> Result<()> {
  154. let args = Args::parse();
  155. let (lvl, conf) = log_config(args.verbose.into())?;
  156. TermLogger::init(lvl, conf, TerminalMode::Mixed, ColorChoice::Auto)?;
  157. let ex = Arc::new(Executor::new());
  158. let ex_clone = ex.clone();
  159. let (signal, shutdown) = async_channel::unbounded::<()>();
  160. let (_, result) = Parallel::new()
  161. .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
  162. // Run the main future on the current thread.
  163. .finish(|| {
  164. smol::future::block_on(async move {
  165. start(ex_clone.clone(), args).await?;
  166. drop(signal);
  167. Ok::<(), darkfi::Error>(())
  168. })
  169. });
  170. result
  171. }