main.rs 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228
  1. use async_std::sync::{Arc, Mutex};
  2. use std::path::Path;
  3. use async_executor::Executor;
  4. use fxhash::FxHashMap;
  5. use log::{error, info, warn};
  6. use smol::future;
  7. use structopt::StructOpt;
  8. use url::Url;
  9. use darkfi::{
  10. net,
  11. raft::{DataStore, NetMsg, ProtocolRaft, Raft, RaftSettings},
  12. util::{
  13. cli::{get_log_config, get_log_level},
  14. expand_path,
  15. serial::{SerialDecodable, SerialEncodable},
  16. sleep,
  17. },
  18. Result,
  19. };
  20. #[derive(Clone, Debug, StructOpt)]
  21. #[structopt(name = "raft-diag")]
  22. pub struct Args {
  23. /// JSON-RPC listen URL
  24. #[structopt(long = "rpc", default_value = "tcp://127.0.0.1:12055")]
  25. pub rpc_listen: Url,
  26. /// Inbound listen URL
  27. #[structopt(long = "inbound")]
  28. pub inbound_url: Vec<Url>,
  29. /// Seed Urls
  30. #[structopt(long = "seeds")]
  31. pub seed_urls: Vec<Url>,
  32. /// Outbound connections
  33. #[structopt(long = "outbound", default_value = "0")]
  34. pub outbound_connections: u32,
  35. /// Sets Datastore Path
  36. #[structopt(long = "path", default_value = "test1.db")]
  37. pub datastore: String,
  38. /// Check if all datastore paths provided are synced
  39. #[structopt(long = "check")]
  40. pub check: Vec<String>,
  41. /// Datastore path to extract and print it
  42. #[structopt(long = "extract")]
  43. pub extract: Option<String>,
  44. /// Number of messages to broadcast
  45. #[structopt(short, default_value = "0")]
  46. pub broadcast: u32,
  47. /// Increase verbosity
  48. #[structopt(short, parse(from_occurrences))]
  49. pub verbose: u8,
  50. }
  51. #[derive(Debug, Clone, SerialEncodable, SerialDecodable, PartialEq, Eq)]
  52. pub struct Message {
  53. payload: String,
  54. }
  55. fn extract(path: &str) -> Result<()> {
  56. if !Path::new(path).exists() {
  57. return Ok(())
  58. }
  59. let db = DataStore::<Message>::new(path)?;
  60. let commits = db.commits.get_all()?;
  61. println!("{:?}", commits);
  62. Ok(())
  63. }
  64. fn check(args: Args) -> Result<()> {
  65. let mut commits_check = vec![];
  66. for path in args.check {
  67. if !Path::new(&path).exists() {
  68. continue
  69. }
  70. let db = DataStore::<Message>::new(&path)?;
  71. let commits = db.commits.get_all()?;
  72. commits_check.push(commits);
  73. }
  74. let result = commits_check.windows(2).all(|w| w[0] == w[1]);
  75. println!("Synced: {}", result);
  76. Ok(())
  77. }
  78. async fn start_broadcasting(n: u32, sender: async_channel::Sender<Message>) -> Result<()> {
  79. sleep(8).await;
  80. info!(target: "raft", "Start broadcasting...");
  81. for id in 0..n {
  82. let msg = format!("msg_test_{}", id);
  83. info!(target: "raft", "Send a message {:?}", msg);
  84. let msg = Message { payload: msg };
  85. sender.send(msg).await?;
  86. }
  87. Ok(())
  88. }
  89. async fn receive_loop(receiver: async_channel::Receiver<Message>) -> Result<()> {
  90. loop {
  91. let msg = receiver.recv().await?;
  92. info!(target: "raft", "Receive new msg {:?}", msg);
  93. }
  94. }
  95. async fn start(args: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  96. let net_settings = net::Settings {
  97. outbound_connections: args.outbound_connections,
  98. inbound: args.inbound_url.clone(),
  99. external_addr: args.inbound_url,
  100. seeds: args.seed_urls,
  101. ..net::Settings::default()
  102. };
  103. //
  104. // Raft
  105. //
  106. let datastore_raft = expand_path(&args.datastore)?;
  107. let seen_net_msgs = Arc::new(Mutex::new(FxHashMap::default()));
  108. let raft_settings = RaftSettings { datastore_path: datastore_raft, ..RaftSettings::default() };
  109. let mut raft = Raft::<Message>::new(raft_settings, seen_net_msgs.clone())?;
  110. //
  111. // P2p setup
  112. //
  113. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<NetMsg>();
  114. let p2p = net::P2p::new(net_settings).await;
  115. let p2p = p2p.clone();
  116. let registry = p2p.protocol_registry();
  117. let raft_node_id = raft.id();
  118. registry
  119. .register(net::SESSION_ALL, move |channel, p2p| {
  120. let raft_node_id = raft_node_id.clone();
  121. let sender = p2p_send_channel.clone();
  122. let seen_net_msgs_cloned = seen_net_msgs.clone();
  123. async move {
  124. ProtocolRaft::init(raft_node_id, channel, sender, p2p, seen_net_msgs_cloned).await
  125. }
  126. })
  127. .await;
  128. p2p.clone().start(executor.clone()).await?;
  129. executor.spawn(p2p.clone().run(executor.clone())).detach();
  130. //
  131. // Waiting Exit signal
  132. //
  133. let (signal, shutdown) = async_channel::bounded::<()>(1);
  134. ctrlc::set_handler(move || {
  135. warn!("Catch exit signal");
  136. // cleaning up tasks running in the background
  137. if let Err(e) = async_std::task::block_on(signal.send(())) {
  138. error!("Error on sending exit signal: {}", e);
  139. }
  140. })
  141. .unwrap();
  142. if args.broadcast != 0 {
  143. executor.spawn(start_broadcasting(args.broadcast, raft.sender())).detach();
  144. }
  145. executor.spawn(receive_loop(raft.receiver())).detach();
  146. raft.run(p2p.clone(), p2p_recv_channel.clone(), executor.clone(), shutdown.clone()).await?;
  147. Ok(())
  148. }
  149. fn main() -> Result<()> {
  150. let args = Args::from_args();
  151. let log_level = get_log_level(args.verbose.into());
  152. let log_config = get_log_config();
  153. let mut log_path = expand_path(&args.datastore)?;
  154. let log_name: String = log_path.file_name().as_ref().unwrap().to_str().unwrap().to_owned();
  155. log_path.pop();
  156. let log_path = log_path.join(&format!("{}.log", log_name));
  157. let env_log_file_path = std::fs::File::create(log_path).unwrap();
  158. simplelog::CombinedLogger::init(vec![
  159. simplelog::TermLogger::new(
  160. log_level,
  161. log_config.clone(),
  162. simplelog::TerminalMode::Mixed,
  163. simplelog::ColorChoice::Auto,
  164. ),
  165. simplelog::WriteLogger::new(log_level, log_config, env_log_file_path),
  166. ])?;
  167. if !args.check.is_empty() {
  168. return check(args)
  169. }
  170. if args.extract.is_some() {
  171. return extract(&args.extract.unwrap())
  172. }
  173. // https://docs.rs/smol/latest/smol/struct.Executor.html#examples
  174. let ex = Arc::new(async_executor::Executor::new());
  175. let (signal, shutdown) = async_channel::unbounded::<()>();
  176. let (_, result) = easy_parallel::Parallel::new()
  177. // Run four executor threads
  178. .each(0..4, |_| future::block_on(ex.run(shutdown.recv())))
  179. // Run the main future on the current thread.
  180. .finish(|| {
  181. future::block_on(async {
  182. start(args, ex.clone()).await?;
  183. drop(signal);
  184. Ok::<(), darkfi::Error>(())
  185. })
  186. });
  187. result
  188. }