main.rs 8.3 KB

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