main.rs 6.6 KB

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