main.rs 6.5 KB

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