main.rs 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273
  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, net,
  10. rpc::server::listen_and_serve,
  11. util::cli::{get_log_config, get_log_level},
  12. Result,
  13. };
  14. pub(crate) mod proto;
  15. pub(crate) mod rpc;
  16. use crate::proto::debugmsg::{Debugmsg, ProtocolDebugmsg, SeenDebugmsgIds};
  17. #[derive(Parser)]
  18. #[clap(name = "p2pdebugging", about = cli_desc!(), version)]
  19. struct Args {
  20. /// Verbosity level
  21. #[clap(short, parse(from_occurrences))]
  22. verbose: u8,
  23. /// node number:
  24. /// 0-2 is for seed nodes
  25. /// 3-20 is for inbound connections nodes
  26. /// 21- is for outbound connections nodes
  27. #[clap(short, long, default_value = "0")]
  28. node: u8,
  29. /// broadcast messages
  30. #[clap(short, long)]
  31. broadcast: bool,
  32. /// communicate using tls protocol by default is tcp
  33. #[clap(long)]
  34. tls: bool,
  35. /// communicate using tor protocol
  36. #[clap(long)]
  37. tor: bool,
  38. /// peers url to connect to manually
  39. #[clap(long)]
  40. peers: Vec<Url>,
  41. /// communicate using tls protocol by default is tcp
  42. #[clap(long, default_value = "tcp://127.0.0.1:11055")]
  43. rpc: String,
  44. /// open manual connection
  45. #[clap(long)]
  46. connect: Option<String>,
  47. }
  48. #[derive(Debug, Clone)]
  49. enum State {
  50. Seed,
  51. Inbound,
  52. Outbound,
  53. }
  54. struct MockP2p {
  55. node_number: u8,
  56. state: State,
  57. broadcast: bool,
  58. address: Option<Url>,
  59. }
  60. impl MockP2p {
  61. async fn new(
  62. node_number: u8,
  63. _broadcast: bool,
  64. scheme: &str,
  65. connect: Option<String>,
  66. peers: Vec<Url>,
  67. ) -> Result<(net::P2pPtr, Self)> {
  68. let seed_addrs: Vec<Url> = vec![
  69. Url::parse(&format!("{}://127.0.0.1:11001", scheme))?,
  70. Url::parse(&format!("{}://127.0.0.1:11002", scheme))?,
  71. Url::parse(&format!("{}://127.0.0.1:11003", scheme))?,
  72. ];
  73. let state: State;
  74. let address: Option<Url>;
  75. let mut broadcast = _broadcast;
  76. let p2p = if connect.is_none() {
  77. match node_number {
  78. 0..=2 => {
  79. address = Some(seed_addrs[node_number as usize].clone());
  80. let net_settings =
  81. net::Settings { inbound: address.clone(), peers, ..Default::default() };
  82. let p2p = net::P2p::new(net_settings).await;
  83. broadcast = false;
  84. state = State::Seed;
  85. p2p
  86. }
  87. 3..=20 => {
  88. let random_port: u32 = rand::thread_rng().gen_range(11007..49151);
  89. address = Some(format!("{}://127.0.0.1:{}", scheme, random_port).parse()?);
  90. let net_settings = net::Settings {
  91. inbound: address.clone(),
  92. external_addr: address.clone(),
  93. seeds: seed_addrs,
  94. peers,
  95. ..Default::default()
  96. };
  97. let p2p = net::P2p::new(net_settings).await;
  98. state = State::Inbound;
  99. p2p
  100. }
  101. _ => {
  102. address = None;
  103. let net_settings = net::Settings {
  104. outbound_connections: 3,
  105. seeds: seed_addrs,
  106. peers,
  107. ..Default::default()
  108. };
  109. let p2p = net::P2p::new(net_settings).await;
  110. state = State::Outbound;
  111. p2p
  112. }
  113. }
  114. } else {
  115. address = None;
  116. let net_settings =
  117. net::Settings { peers: vec![Url::parse(&connect.unwrap())?], ..Default::default() };
  118. let p2p = net::P2p::new(net_settings).await;
  119. state = State::Outbound;
  120. p2p
  121. };
  122. println!("start {:?} node #{} address {:?}", state, node_number, address);
  123. Ok((p2p, Self { node_number, state, broadcast, address }))
  124. }
  125. async fn run(
  126. &self,
  127. p2p: net::P2pPtr,
  128. rpc_addr: Url,
  129. executor: Arc<Executor<'_>>,
  130. ) -> Result<()> {
  131. let state = self.state.clone();
  132. let node_number = self.node_number;
  133. let address = self.address.clone();
  134. let (sender, receiver) = async_channel::unbounded();
  135. let sender_clone = sender.clone();
  136. let seen_debugmsg_ids = SeenDebugmsgIds::new();
  137. let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
  138. let registry = p2p.protocol_registry();
  139. registry
  140. .register(net::SESSION_ALL, move |channel, p2p| {
  141. let sender = sender_clone.clone();
  142. let seen_debugmsg_ids = seen_debugmsg_ids_clone.clone();
  143. async move { ProtocolDebugmsg::init(channel, sender, seen_debugmsg_ids, p2p).await }
  144. })
  145. .await;
  146. if self.broadcast {
  147. println!("start broadcast {:?} node #{} address {:?}", state, node_number, address);
  148. let sleep_time = 10;
  149. let p2p_clone = p2p.clone();
  150. let executor_clone = executor.clone();
  151. let address_cloned = address.clone();
  152. let state = state.clone();
  153. executor_clone
  154. .spawn(async move {
  155. loop {
  156. darkfi::util::sleep(sleep_time).await;
  157. println!(
  158. "broadcast sleep for {} {:?} node #{} address {:?}",
  159. sleep_time,
  160. state,
  161. node_number,
  162. address_cloned.clone()
  163. );
  164. let random_id = OsRng.next_u32();
  165. let msg = Debugmsg { id: random_id, message: "hello".to_string() };
  166. println!(
  167. "send {:?} {:?} node #{} address {:?}",
  168. msg, state, node_number, address_cloned
  169. );
  170. p2p_clone.broadcast(msg).await.unwrap();
  171. }
  172. })
  173. .detach();
  174. }
  175. let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
  176. let address_cloned = address.clone();
  177. executor
  178. .spawn(async move {
  179. loop {
  180. let msg = receiver.recv().await.unwrap();
  181. println!(
  182. "receive {:?} {:?} node #{} address {:?}",
  183. msg, state, node_number, address_cloned
  184. );
  185. seen_debugmsg_ids_clone.add_seen(msg.id).await;
  186. }
  187. })
  188. .detach();
  189. // RPC
  190. let rpc_interface =
  191. Arc::new(rpc::JsonRpcInterface { addr: rpc_addr.clone(), p2p: p2p.clone() });
  192. executor.spawn(async move { listen_and_serve(rpc_addr, rpc_interface).await }).detach();
  193. p2p.clone().start(executor.clone()).await?;
  194. p2p.run(executor).await
  195. }
  196. }
  197. async fn start(executor: Arc<Executor<'_>>, args: Args) -> Result<()> {
  198. let mut scheme = (if args.tor { "tor" } else { "tcp" }).to_string();
  199. if args.tls {
  200. scheme = format!("{}+tls", scheme);
  201. }
  202. let rpc_addr = Url::parse(&args.rpc)?;
  203. let (p2p, mock_p2p) =
  204. MockP2p::new(args.node, args.broadcast, &scheme, args.connect, args.peers).await?;
  205. mock_p2p.run(p2p.clone(), rpc_addr, executor).await
  206. }
  207. fn main() -> Result<()> {
  208. let args = Args::parse();
  209. let log_level = get_log_level(args.verbose.into());
  210. let log_config = get_log_config();
  211. TermLogger::init(log_level, log_config, TerminalMode::Mixed, ColorChoice::Auto)?;
  212. let ex = Arc::new(Executor::new());
  213. let ex_clone = ex.clone();
  214. let (signal, shutdown) = async_channel::unbounded::<()>();
  215. let (_, result) = Parallel::new()
  216. .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
  217. // Run the main future on the current thread.
  218. .finish(|| {
  219. smol::future::block_on(async move {
  220. start(ex_clone.clone(), args).await?;
  221. drop(signal);
  222. Ok::<(), darkfi::Error>(())
  223. })
  224. });
  225. result
  226. }