main.rs 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214
  1. use async_std::sync::Arc;
  2. use std::fmt;
  3. use async_channel::Receiver;
  4. use async_executor::Executor;
  5. use log::{info, warn};
  6. use rand::rngs::OsRng;
  7. use smol::future;
  8. use structopt_toml::StructOptToml;
  9. use darkfi::{
  10. async_daemonize, net,
  11. net::P2pPtr,
  12. rpc::server::listen_and_serve,
  13. system::{Subscriber, SubscriberPtr},
  14. util::{
  15. cli::{get_log_config, get_log_level, spawn_config},
  16. expand_path,
  17. file::save_json_file,
  18. path::get_config_path,
  19. sleep,
  20. },
  21. Result,
  22. };
  23. pub mod buffers;
  24. pub mod crypto;
  25. pub mod irc;
  26. pub mod privmsg;
  27. pub mod protocol_privmsg;
  28. pub mod rpc;
  29. pub mod settings;
  30. use crate::{
  31. buffers::{create_buffers, Buffers},
  32. irc::IrcServer,
  33. privmsg::Privmsg,
  34. protocol_privmsg::{LastTerm, ProtocolPrivmsg},
  35. rpc::JsonRpcInterface,
  36. settings::{Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  37. };
  38. #[derive(serde::Serialize)]
  39. struct KeyPair {
  40. private_key: String,
  41. public_key: String,
  42. }
  43. impl fmt::Display for KeyPair {
  44. fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
  45. write!(f, "Public key: {}\nPrivate key: {}", self.public_key, self.private_key)
  46. }
  47. }
  48. async fn resend_unread_msgs(p2p: P2pPtr, buffers: Buffers) -> Result<()> {
  49. loop {
  50. sleep(settings::TIMEOUT_FOR_RESEND_UNREAD_MSGS).await;
  51. for msg in buffers.unread_msgs.load().await.values() {
  52. p2p.broadcast(msg.clone()).await?;
  53. }
  54. }
  55. }
  56. async fn send_last_term(p2p: P2pPtr, buffers: Buffers) -> Result<()> {
  57. loop {
  58. sleep(settings::BROADCAST_LAST_TERM_MSG).await;
  59. let term = buffers.privmsgs.last_term().await;
  60. p2p.broadcast(LastTerm { term }).await?;
  61. }
  62. }
  63. struct Ircd {
  64. notify_clients: SubscriberPtr<Privmsg>,
  65. }
  66. impl Ircd {
  67. fn new() -> Self {
  68. let notify_clients = Subscriber::new();
  69. Self { notify_clients }
  70. }
  71. async fn start(
  72. &self,
  73. settings: &Args,
  74. buffers: Buffers,
  75. p2p: net::P2pPtr,
  76. p2p_receiver: Receiver<Privmsg>,
  77. executor: Arc<Executor<'_>>,
  78. ) -> Result<()> {
  79. let notify_clients = self.notify_clients.clone();
  80. executor
  81. .spawn(async move {
  82. while let Ok(msg) = p2p_receiver.recv().await {
  83. notify_clients.notify(msg).await;
  84. }
  85. })
  86. .detach();
  87. let irc_server = IrcServer::new(
  88. settings.clone(),
  89. buffers.clone(),
  90. p2p.clone(),
  91. self.notify_clients.clone(),
  92. )
  93. .await?;
  94. let executor_cloned = executor.clone();
  95. executor
  96. .spawn(async move {
  97. irc_server.start(executor_cloned.clone()).await.unwrap();
  98. })
  99. .detach();
  100. Ok(())
  101. }
  102. }
  103. async_daemonize!(realmain);
  104. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  105. let buffers = create_buffers();
  106. if settings.gen_secret {
  107. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  108. let encoded = bs58::encode(secret_key.as_bytes());
  109. println!("{}", encoded.into_string());
  110. return Ok(())
  111. }
  112. if settings.gen_keypair {
  113. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  114. let pub_key = secret_key.public_key();
  115. let prv_encoded = bs58::encode(secret_key.as_bytes()).into_string();
  116. let pub_encoded = bs58::encode(pub_key.as_bytes()).into_string();
  117. let kp = KeyPair { private_key: prv_encoded, public_key: pub_encoded };
  118. if settings.output.is_some() {
  119. let datastore = expand_path(&settings.output.unwrap())?;
  120. save_json_file(&datastore, &kp)?;
  121. } else {
  122. println!("Generated KeyPair:\n{}", kp);
  123. }
  124. return Ok(())
  125. }
  126. //
  127. // P2p setup
  128. //
  129. let mut net_settings = settings.net.clone();
  130. net_settings.app_version = Some(option_env!("CARGO_PKG_VERSION").unwrap_or("").to_string());
  131. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<Privmsg>();
  132. let p2p = net::P2p::new(net_settings.into()).await;
  133. let p2p2 = p2p.clone();
  134. let registry = p2p.protocol_registry();
  135. let buffers_cloned = buffers.clone();
  136. registry
  137. .register(net::SESSION_ALL, move |channel, p2p| {
  138. let sender = p2p_send_channel.clone();
  139. let buffers_cloned = buffers_cloned.clone();
  140. async move { ProtocolPrivmsg::init(channel, sender, p2p, buffers_cloned).await }
  141. })
  142. .await;
  143. p2p.clone().start(executor.clone()).await?;
  144. let executor_cloned = executor.clone();
  145. executor_cloned.spawn(p2p.clone().run(executor.clone())).detach();
  146. //
  147. // Sync tasks
  148. //
  149. executor.spawn(resend_unread_msgs(p2p.clone(), buffers.clone())).detach();
  150. executor.spawn(send_last_term(p2p.clone(), buffers.clone())).detach();
  151. //
  152. // RPC interface
  153. //
  154. let rpc_listen_addr = settings.rpc_listen.clone();
  155. let rpc_interface =
  156. Arc::new(JsonRpcInterface { addr: rpc_listen_addr.clone(), p2p: p2p.clone() });
  157. executor.spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface).await }).detach();
  158. //
  159. // IRC instance
  160. //
  161. let ircd = Ircd::new();
  162. ircd.start(&settings, buffers, p2p, p2p_recv_channel, executor.clone()).await?;
  163. // Run once receive exit signal
  164. let (signal, shutdown) = async_channel::bounded::<()>(1);
  165. ctrlc::set_handler(move || {
  166. warn!(target: "ircd", "ircd start Exit Signal");
  167. // cleaning up tasks running in the background
  168. async_std::task::block_on(signal.send(())).unwrap();
  169. })
  170. .unwrap();
  171. // Wait for SIGINT
  172. shutdown.recv().await?;
  173. print!("\r");
  174. info!("Caught termination signal, cleaning up and exiting...");
  175. p2p2.stop().await;
  176. Ok(())
  177. }