dfi.rs 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. #[macro_use]
  2. extern crate clap;
  3. use async_channel::unbounded;
  4. use async_dup::Arc;
  5. use async_executor::Executor;
  6. use async_std::sync::Mutex;
  7. use easy_parallel::Parallel;
  8. use log::*;
  9. use std::collections::HashMap;
  10. use std::net::{IpAddr, Ipv4Addr, SocketAddr};
  11. use sapvi::{ClientProtocol, Result, SeedProtocol, ServerProtocol};
  12. async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
  13. let connections = Arc::new(Mutex::new(HashMap::new()));
  14. let stored_addrs = Arc::new(Mutex::new(Vec::new()));
  15. let executor2 = executor.clone();
  16. let stored_addrs2 = stored_addrs.clone();
  17. let mut server_task = None;
  18. if let Some(accept_addr) = options.accept_addr {
  19. let accept_addr = accept_addr.clone();
  20. let mut protocol = ServerProtocol::new(connections.clone());
  21. server_task = Some(executor.spawn(async move {
  22. protocol
  23. .start(accept_addr, stored_addrs2, executor2)
  24. .await?;
  25. Ok::<(), sapvi::Error>(())
  26. }));
  27. }
  28. let mut seed_protocols = Vec::with_capacity(options.seed_addrs.len());
  29. // Normally we query this from a server
  30. let local_addr = options.accept_addr.clone();
  31. for seed_addr in options.seed_addrs.iter() {
  32. let mut protocol = SeedProtocol::new();
  33. protocol
  34. .start(
  35. seed_addr.clone(),
  36. local_addr,
  37. stored_addrs.clone(),
  38. executor.clone(),
  39. )
  40. .await;
  41. seed_protocols.push(protocol);
  42. }
  43. debug!("Waiting for seed node queries to finish...");
  44. for seed_protocol in seed_protocols {
  45. seed_protocol.await_finish().await;
  46. }
  47. debug!("Seed nodes queried.");
  48. let accept_addr = options.accept_addr.clone();
  49. let mut client_slots: Vec<ClientProtocol> = vec![];
  50. for i in 0..options.connection_slots {
  51. debug!("Starting connection slot {}", i);
  52. let mut client = ClientProtocol::new(connections.clone());
  53. client
  54. .start(accept_addr.clone(), stored_addrs.clone(), executor.clone())
  55. .await;
  56. client_slots.push(client);
  57. }
  58. for remote_addr in options.manual_connects {
  59. debug!("Starting connection (manual) to {}", remote_addr);
  60. let mut client = ClientProtocol::new(connections.clone());
  61. client
  62. .start_manual(
  63. remote_addr,
  64. accept_addr.clone(),
  65. stored_addrs.clone(),
  66. executor.clone(),
  67. )
  68. .await;
  69. client_slots.push(client);
  70. }
  71. loop {
  72. sapvi::sleep(2).await;
  73. }
  74. //server_task.cancel().await;
  75. //Ok(())
  76. }
  77. struct ProgramOptions {
  78. accept_addr: Option<SocketAddr>,
  79. seed_addrs: Vec<SocketAddr>,
  80. manual_connects: Vec<SocketAddr>,
  81. connection_slots: u32,
  82. }
  83. impl ProgramOptions {
  84. fn load() -> Result<ProgramOptions> {
  85. let app = clap_app!(dfi =>
  86. (version: "0.1.0")
  87. (author: "Amir Taaki <amir@dyne.org>")
  88. (about: "Dark node")
  89. (@arg ACCEPT: -a --accept +takes_value "Accept address")
  90. (@arg SEED_NODES: -s --seeds ... "Seed nodes")
  91. (@arg CONNECTS: -c --connect ... "Manual connections")
  92. (@arg CONNECT_SLOTS: --slots +takes_value "Connection slots")
  93. )
  94. .get_matches();
  95. let accept_addr = if let Some(accept_addr) = app.value_of("ACCEPT") {
  96. Some(accept_addr.parse()?)
  97. } else {
  98. None
  99. };
  100. let mut seed_addrs: Vec<SocketAddr> = vec![];
  101. if let Some(seeds) = app.values_of("SEED_NODES") {
  102. for seed in seeds {
  103. seed_addrs.push(seed.parse()?);
  104. }
  105. }
  106. let mut manual_connects: Vec<SocketAddr> = vec![];
  107. if let Some(connections) = app.values_of("CONNECTS") {
  108. for connect in connections {
  109. manual_connects.push(connect.parse()?);
  110. }
  111. }
  112. let connection_slots = if let Some(connection_slots) = app.value_of("CONNECT_SLOTS") {
  113. connection_slots.parse()?
  114. } else {
  115. 0
  116. };
  117. Ok(ProgramOptions {
  118. accept_addr,
  119. seed_addrs,
  120. manual_connects,
  121. connection_slots,
  122. })
  123. }
  124. }
  125. fn main() -> Result<()> {
  126. use simplelog::*;
  127. CombinedLogger::init(vec![TermLogger::new(
  128. LevelFilter::Debug,
  129. Config::default(),
  130. TerminalMode::Mixed,
  131. )
  132. .unwrap()])
  133. .unwrap();
  134. let options = ProgramOptions::load()?;
  135. let ex = Arc::new(Executor::new());
  136. let (signal, shutdown) = unbounded::<()>();
  137. let ex2 = ex.clone();
  138. let (_, result) = Parallel::new()
  139. // Run four executor threads.
  140. .each(0..3, |_| smol::future::block_on(ex.run(shutdown.recv())))
  141. // Run the main future on the current thread.
  142. .finish(|| {
  143. smol::future::block_on(async move {
  144. start(ex2, options).await?;
  145. drop(signal);
  146. Ok::<(), sapvi::Error>(())
  147. })
  148. });
  149. result
  150. }