main.rs 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270
  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. util::{
  18. cli::{get_log_config, get_log_level, spawn_config},
  19. path::get_config_path,
  20. },
  21. Error, Result,
  22. };
  23. pub mod crypto;
  24. pub mod privmsg;
  25. pub mod protocol_privmsg;
  26. pub mod rpc;
  27. pub mod server;
  28. pub mod settings;
  29. pub mod util;
  30. use crate::{
  31. crypto::try_decrypt_message,
  32. privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
  33. protocol_privmsg::ProtocolPrivmsg,
  34. rpc::JsonRpcInterface,
  35. server::IrcServerConnection,
  36. settings::{parse_configured_channels, Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  37. util::clean_input,
  38. };
  39. const SIZE_OF_MSGS_BUFFER: usize = 4096;
  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. p2p_receiver: Receiver<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. p2p_receiver: Receiver<Privmsg>,
  59. ) -> Self {
  60. Self { seen_msg_ids, privmsgs_buffer, autojoin_chans, configured_chans, p2p, p2p_receiver }
  61. }
  62. async fn process(
  63. &self,
  64. executor: Arc<Executor<'_>>,
  65. stream: TcpStream,
  66. peer_addr: SocketAddr,
  67. ) -> Result<()> {
  68. let (reader, writer) = stream.split();
  69. let mut reader = BufReader::new(reader);
  70. let mut conn = IrcServerConnection::new(
  71. writer,
  72. self.seen_msg_ids.clone(),
  73. self.privmsgs_buffer.clone(),
  74. self.autojoin_chans.clone(),
  75. self.configured_chans.clone(),
  76. self.p2p.clone(),
  77. );
  78. let p2p_receiver = self.p2p_receiver.clone();
  79. executor
  80. .spawn(async move {
  81. loop {
  82. let mut line = String::new();
  83. futures::select! {
  84. privmsg = p2p_receiver.recv().fuse() => {
  85. let mut msg = privmsg?;
  86. info!("Received msg from P2p network: {:?}", msg);
  87. // Try to potentially decrypt the incoming message.
  88. if conn.configured_chans.contains_key(&msg.channel) {
  89. let chan_info = conn.configured_chans.get(&msg.channel).unwrap();
  90. if !chan_info.joined {
  91. continue
  92. }
  93. let salt_box = chan_info.salt_box.clone();
  94. if salt_box.is_some() {
  95. let decrypted_msg =
  96. try_decrypt_message(&salt_box.unwrap(), &msg.message);
  97. if decrypted_msg.is_none() {
  98. continue
  99. }
  100. msg.message = decrypted_msg.unwrap();
  101. info!("Decrypted received message: {:?}", msg);
  102. }
  103. }
  104. conn.reply(&msg.to_irc_msg()).await?;
  105. }
  106. err = reader.read_line(&mut line).fuse() => {
  107. if let Err(e) = err {
  108. warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
  109. return Ok(())
  110. }
  111. info!("Received msg from IRC client: {:?}", line);
  112. let irc_msg = match clean_input(line, &peer_addr) {
  113. Ok(m) => m,
  114. Err(e) => return Err(e)
  115. };
  116. info!("Send msg to IRC client '{}' from {}", irc_msg, peer_addr);
  117. if let Err(e) = conn.update(irc_msg).await {
  118. warn!("Connection error: {} for {}", e, peer_addr);
  119. return Err(Error::ChannelStopped)
  120. }
  121. }
  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. p2p_recv_channel.clone(),
  197. );
  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 process the {} connections: {}", 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. }