main.rs 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263
  1. use async_std::{
  2. net::{TcpListener, TcpStream},
  3. sync::{Arc, Mutex},
  4. };
  5. use std::net::SocketAddr;
  6. use async_channel::Receiver;
  7. use async_executor::Executor;
  8. use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
  9. use fxhash::FxHashMap;
  10. use log::{error, info, warn};
  11. use rand::rngs::OsRng;
  12. use smol::future;
  13. use structopt_toml::StructOptToml;
  14. use darkfi::{
  15. async_daemonize, net,
  16. rpc::server::listen_and_serve,
  17. system::{Subscriber, SubscriberPtr},
  18. util::{
  19. cli::{get_log_config, get_log_level, spawn_config},
  20. path::get_config_path,
  21. },
  22. Error, Result,
  23. };
  24. pub mod crypto;
  25. pub mod privmsg;
  26. pub mod protocol_privmsg;
  27. pub mod rpc;
  28. pub mod server;
  29. pub mod settings;
  30. use crate::{
  31. privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
  32. protocol_privmsg::ProtocolPrivmsg,
  33. rpc::JsonRpcInterface,
  34. server::IrcServerConnection,
  35. settings::{parse_configured_channels, Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  36. };
  37. const SIZE_OF_MSGS_BUFFER: usize = 4096;
  38. struct Ircd {
  39. // msgs
  40. seen_msg_ids: SeenMsgIds,
  41. privmsgs_buffer: PrivmsgsBuffer,
  42. // channels
  43. autojoin_chans: Vec<String>,
  44. configured_chans: FxHashMap<String, ChannelInfo>,
  45. // p2p
  46. p2p: net::P2pPtr,
  47. senders: SubscriberPtr<Privmsg>,
  48. }
  49. impl Ircd {
  50. fn new(
  51. seen_msg_ids: SeenMsgIds,
  52. privmsgs_buffer: PrivmsgsBuffer,
  53. autojoin_chans: Vec<String>,
  54. configured_chans: FxHashMap<String, ChannelInfo>,
  55. p2p: net::P2pPtr,
  56. ) -> Self {
  57. let senders = Subscriber::new();
  58. Self { seen_msg_ids, privmsgs_buffer, autojoin_chans, configured_chans, p2p, senders }
  59. }
  60. fn start_p2p_receive_loop(&self, executor: Arc<Executor<'_>>, p2p_receiver: Receiver<Privmsg>) {
  61. let senders = self.senders.clone();
  62. let p2p_receiver_cloned = p2p_receiver.clone();
  63. executor
  64. .spawn(async move {
  65. while let Ok(msg) = p2p_receiver_cloned.recv().await {
  66. senders.notify(msg).await;
  67. }
  68. })
  69. .detach();
  70. }
  71. async fn process(
  72. &self,
  73. executor: Arc<Executor<'_>>,
  74. stream: TcpStream,
  75. peer_addr: SocketAddr,
  76. ) -> Result<()> {
  77. let (reader, writer) = stream.split();
  78. let mut reader = BufReader::new(reader);
  79. // New subscriber
  80. let receiver = self.senders.clone().subscribe().await;
  81. // New irc connection
  82. let mut conn = IrcServerConnection::new(
  83. writer,
  84. peer_addr,
  85. self.seen_msg_ids.clone(),
  86. self.privmsgs_buffer.clone(),
  87. self.autojoin_chans.clone(),
  88. self.configured_chans.clone(),
  89. self.p2p.clone(),
  90. self.senders.clone(),
  91. receiver.get_id(),
  92. );
  93. executor
  94. .spawn(async move {
  95. loop {
  96. let mut line = String::new();
  97. let result: Result<()> = futures::select! {
  98. msg = receiver.receive().fuse() => {
  99. info!("Received msg from P2p network: {:?}", msg);
  100. match conn.process_msg_from_p2p(&msg).await {
  101. Ok(_) => Ok(()),
  102. Err(e) => {
  103. error!("Process msg from p2p failed {}: {}", peer_addr, e);
  104. Err(Error::ChannelStopped)
  105. }
  106. }
  107. }
  108. err = reader.read_line(&mut line).fuse() => {
  109. match conn.process_line_from_client(err, line).await {
  110. Ok(_) => Ok(()),
  111. Err(e) => {
  112. error!("Process line from client failed {}: {}", peer_addr, e);
  113. Err(Error::ChannelStopped)
  114. }
  115. }
  116. }
  117. };
  118. if let Err(e) = result {
  119. warn!("Close connection for clinet {}: {}", peer_addr, e);
  120. receiver.unsubscribe().await;
  121. break
  122. }
  123. }
  124. })
  125. .detach();
  126. Ok(())
  127. }
  128. }
  129. async_daemonize!(realmain);
  130. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  131. let seen_msg_ids = Arc::new(Mutex::new(vec![]));
  132. let privmsgs_buffer: PrivmsgsBuffer =
  133. Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSGS_BUFFER)));
  134. if settings.gen_secret {
  135. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  136. let encoded = bs58::encode(secret_key.as_bytes());
  137. println!("{}", encoded.into_string());
  138. return Ok(())
  139. }
  140. // Pick up channel settings from the TOML configuration
  141. let cfg_path = get_config_path(settings.config, CONFIG_FILE)?;
  142. let configured_chans = parse_configured_channels(&cfg_path)?;
  143. //
  144. // P2p setup
  145. //
  146. let net_settings = settings.net;
  147. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<Privmsg>();
  148. let p2p = net::P2p::new(net_settings.into()).await;
  149. let p2p = p2p.clone();
  150. let registry = p2p.protocol_registry();
  151. let seen_msg_ids_cloned = seen_msg_ids.clone();
  152. let privmsgs_buffer_cloned = privmsgs_buffer.clone();
  153. registry
  154. .register(net::SESSION_ALL, move |channel, p2p| {
  155. let sender = p2p_send_channel.clone();
  156. let seen_msg_ids_cloned = seen_msg_ids_cloned.clone();
  157. let privmsgs_buffer_cloned = privmsgs_buffer_cloned.clone();
  158. async move {
  159. ProtocolPrivmsg::init(
  160. channel,
  161. sender,
  162. p2p,
  163. seen_msg_ids_cloned,
  164. privmsgs_buffer_cloned,
  165. )
  166. .await
  167. }
  168. })
  169. .await;
  170. p2p.clone().start(executor.clone()).await?;
  171. let executor_cloned = executor.clone();
  172. executor_cloned.spawn(p2p.clone().run(executor.clone())).detach();
  173. //
  174. // RPC interface
  175. //
  176. let rpc_listen_addr = settings.rpc_listen.clone();
  177. let rpc_interface =
  178. Arc::new(JsonRpcInterface { addr: rpc_listen_addr.clone(), p2p: p2p.clone() });
  179. executor.spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface).await }).detach();
  180. //
  181. // IRC instance
  182. //
  183. let irc_listen_addr = settings.irc_listen.socket_addrs(|| None)?[0];
  184. let listener = TcpListener::bind(irc_listen_addr).await?;
  185. let local_addr = listener.local_addr()?;
  186. info!("IRC listening on {}", local_addr);
  187. let executor_cloned = executor.clone();
  188. executor
  189. .spawn(async move {
  190. let ircd = Ircd::new(
  191. seen_msg_ids.clone(),
  192. privmsgs_buffer.clone(),
  193. settings.autojoin.clone(),
  194. configured_chans.clone(),
  195. p2p.clone(),
  196. );
  197. ircd.start_p2p_receive_loop(executor_cloned.clone(), p2p_recv_channel);
  198. loop {
  199. let (stream, peer_addr) = match listener.accept().await {
  200. Ok((s, a)) => (s, a),
  201. Err(e) => {
  202. error!("failed accepting new connections: {}", e);
  203. continue
  204. }
  205. };
  206. let result = ircd.process(executor_cloned.clone(), stream, peer_addr).await;
  207. if let Err(e) = result {
  208. error!("Failed processing connection {}: {}", peer_addr, e);
  209. continue
  210. };
  211. info!("IRC Accepted new client: {}", peer_addr);
  212. }
  213. })
  214. .detach();
  215. // Run once receive exit signal
  216. let (signal, shutdown) = async_channel::bounded::<()>(1);
  217. ctrlc_async::set_async_handler(async move {
  218. warn!(target: "ircd", "ircd start Exit Signal");
  219. // cleaning up tasks running in the background
  220. signal.send(()).await.unwrap();
  221. })
  222. .unwrap();
  223. // Wait for SIGINT
  224. shutdown.recv().await?;
  225. print!("\r");
  226. info!("Caught termination signal, cleaning up and exiting...");
  227. Ok(())
  228. }