main.rs 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::sync::Arc;
  19. use async_executor::Executor;
  20. use clap::Parser;
  21. use easy_parallel::Parallel;
  22. use rand::{rngs::OsRng, Rng, RngCore};
  23. use simplelog::{ColorChoice, TermLogger, TerminalMode};
  24. use url::Url;
  25. use darkfi::{
  26. cli_desc, net,
  27. rpc::server::listen_and_serve,
  28. util::cli::{get_log_config, get_log_level},
  29. Result,
  30. };
  31. pub(crate) mod proto;
  32. pub(crate) mod rpc;
  33. use crate::proto::debugmsg::{Debugmsg, ProtocolDebugmsg, SeenDebugmsgIds};
  34. #[derive(Parser)]
  35. #[clap(name = "p2pdebugging", about = cli_desc!(), version)]
  36. struct Args {
  37. /// Verbosity level
  38. #[clap(short, parse(from_occurrences))]
  39. verbose: u8,
  40. /// node number:
  41. /// 0-2 is for seed nodes
  42. /// 3-20 is for inbound connections nodes
  43. /// 21- is for outbound connections nodes
  44. #[clap(short, long, default_value = "0")]
  45. node: u8,
  46. /// broadcast messages
  47. #[clap(short, long)]
  48. broadcast: bool,
  49. /// communicate using tls protocol by default is tcp
  50. #[clap(long)]
  51. tls: bool,
  52. /// communicate using tor protocol
  53. #[clap(long)]
  54. tor: bool,
  55. /// peers url to connect to manually
  56. #[clap(long)]
  57. peers: Vec<Url>,
  58. /// communicate using tls protocol by default is tcp
  59. #[clap(long, default_value = "tcp://127.0.0.1:11055")]
  60. rpc: String,
  61. /// open manual connection
  62. #[clap(long)]
  63. connect: Option<String>,
  64. }
  65. #[derive(Debug, Clone)]
  66. enum State {
  67. Seed,
  68. Inbound,
  69. Outbound,
  70. }
  71. struct MockP2p {
  72. node_number: u8,
  73. state: State,
  74. broadcast: bool,
  75. address: Option<Url>,
  76. }
  77. impl MockP2p {
  78. async fn new(
  79. node_number: u8,
  80. _broadcast: bool,
  81. scheme: &str,
  82. connect: Option<String>,
  83. peers: Vec<Url>,
  84. ) -> Result<(net::P2pPtr, Self)> {
  85. let seed_addrs: Vec<Url> = vec![
  86. Url::parse(&format!("{}://127.0.0.1:11001", scheme))?,
  87. Url::parse(&format!("{}://127.0.0.1:11002", scheme))?,
  88. Url::parse(&format!("{}://127.0.0.1:11003", scheme))?,
  89. ];
  90. let state: State;
  91. let address: Option<Url>;
  92. let mut broadcast = _broadcast;
  93. let p2p = if connect.is_none() {
  94. match node_number {
  95. 0..=2 => {
  96. address = Some(seed_addrs[node_number as usize].clone());
  97. let net_settings =
  98. net::Settings { inbound: address.clone(), peers, ..Default::default() };
  99. let p2p = net::P2p::new(net_settings).await;
  100. broadcast = false;
  101. state = State::Seed;
  102. p2p
  103. }
  104. 3..=20 => {
  105. let random_port: u32 = rand::thread_rng().gen_range(11007..49151);
  106. address = Some(format!("{}://127.0.0.1:{}", scheme, random_port).parse()?);
  107. let net_settings = net::Settings {
  108. inbound: address.clone(),
  109. external_addr: address.clone(),
  110. seeds: seed_addrs,
  111. peers,
  112. ..Default::default()
  113. };
  114. let p2p = net::P2p::new(net_settings).await;
  115. state = State::Inbound;
  116. p2p
  117. }
  118. _ => {
  119. address = None;
  120. let net_settings = net::Settings {
  121. outbound_connections: 3,
  122. seeds: seed_addrs,
  123. peers,
  124. ..Default::default()
  125. };
  126. let p2p = net::P2p::new(net_settings).await;
  127. state = State::Outbound;
  128. p2p
  129. }
  130. }
  131. } else {
  132. address = None;
  133. let net_settings =
  134. net::Settings { peers: vec![Url::parse(&connect.unwrap())?], ..Default::default() };
  135. let p2p = net::P2p::new(net_settings).await;
  136. state = State::Outbound;
  137. p2p
  138. };
  139. println!("start {:?} node #{} address {:?}", state, node_number, address);
  140. Ok((p2p, Self { node_number, state, broadcast, address }))
  141. }
  142. async fn run(
  143. &self,
  144. p2p: net::P2pPtr,
  145. rpc_addr: Url,
  146. executor: Arc<Executor<'_>>,
  147. ) -> Result<()> {
  148. let state = self.state.clone();
  149. let node_number = self.node_number;
  150. let address = self.address.clone();
  151. let (sender, receiver) = async_channel::unbounded();
  152. let sender_clone = sender.clone();
  153. let seen_debugmsg_ids = SeenDebugmsgIds::new();
  154. let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
  155. let registry = p2p.protocol_registry();
  156. registry
  157. .register(net::SESSION_ALL, move |channel, p2p| {
  158. let sender = sender_clone.clone();
  159. let seen_debugmsg_ids = seen_debugmsg_ids_clone.clone();
  160. async move { ProtocolDebugmsg::init(channel, sender, seen_debugmsg_ids, p2p).await }
  161. })
  162. .await;
  163. if self.broadcast {
  164. println!("start broadcast {:?} node #{} address {:?}", state, node_number, address);
  165. let sleep_time = 10;
  166. let p2p_clone = p2p.clone();
  167. let executor_clone = executor.clone();
  168. let address_cloned = address.clone();
  169. let state = state.clone();
  170. executor_clone
  171. .spawn(async move {
  172. loop {
  173. darkfi::util::sleep(sleep_time).await;
  174. println!(
  175. "broadcast sleep for {} {:?} node #{} address {:?}",
  176. sleep_time,
  177. state,
  178. node_number,
  179. address_cloned.clone()
  180. );
  181. let random_id = OsRng.next_u32();
  182. let msg = Debugmsg { id: random_id, message: "hello".to_string() };
  183. println!(
  184. "send {:?} {:?} node #{} address {:?}",
  185. msg, state, node_number, address_cloned
  186. );
  187. p2p_clone.broadcast(msg).await.unwrap();
  188. }
  189. })
  190. .detach();
  191. }
  192. let seen_debugmsg_ids_clone = seen_debugmsg_ids.clone();
  193. let address_cloned = address.clone();
  194. executor
  195. .spawn(async move {
  196. loop {
  197. let msg = receiver.recv().await.unwrap();
  198. println!(
  199. "receive {:?} {:?} node #{} address {:?}",
  200. msg, state, node_number, address_cloned
  201. );
  202. seen_debugmsg_ids_clone.add_seen(msg.id).await;
  203. }
  204. })
  205. .detach();
  206. // RPC
  207. let rpc_interface =
  208. Arc::new(rpc::JsonRpcInterface { addr: rpc_addr.clone(), p2p: p2p.clone() });
  209. let _ex = executor.clone();
  210. executor.spawn(async move { listen_and_serve(rpc_addr, rpc_interface, _ex).await }).detach();
  211. p2p.clone().start(executor.clone()).await?;
  212. p2p.run(executor).await
  213. }
  214. }
  215. async fn start(executor: Arc<Executor<'_>>, args: Args) -> Result<()> {
  216. let mut scheme = (if args.tor { "tor" } else { "tcp" }).to_string();
  217. if args.tls {
  218. scheme = format!("{}+tls", scheme);
  219. }
  220. let rpc_addr = Url::parse(&args.rpc)?;
  221. let (p2p, mock_p2p) =
  222. MockP2p::new(args.node, args.broadcast, &scheme, args.connect, args.peers).await?;
  223. mock_p2p.run(p2p.clone(), rpc_addr, executor).await
  224. }
  225. fn main() -> Result<()> {
  226. let args = Args::parse();
  227. let log_level = get_log_level(args.verbose.into());
  228. let log_config = get_log_config();
  229. TermLogger::init(log_level, log_config, TerminalMode::Mixed, ColorChoice::Auto)?;
  230. let ex = Arc::new(Executor::new());
  231. let ex_clone = ex.clone();
  232. let (signal, shutdown) = async_channel::unbounded::<()>();
  233. let (_, result) = Parallel::new()
  234. .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
  235. // Run the main future on the current thread.
  236. .finish(|| {
  237. smol::future::block_on(async move {
  238. start(ex_clone.clone(), args).await?;
  239. drop(signal);
  240. Ok::<(), darkfi::Error>(())
  241. })
  242. });
  243. result
  244. }