main.rs 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270
  1. use async_std::{
  2. net::{TcpListener, TcpStream},
  3. sync::{Arc, Mutex},
  4. };
  5. use std::{net::SocketAddr, sync::atomic::Ordering};
  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::{debug, error, info, warn};
  11. use rand::rngs::OsRng;
  12. use ringbuffer::RingBufferWrite;
  13. use smol::future;
  14. use structopt_toml::StructOptToml;
  15. use darkfi::{
  16. async_daemonize, net,
  17. rpc::server::listen_and_serve,
  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. 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. };
  38. const SIZE_OF_MSGS_BUFFER: usize = 4096;
  39. fn clean_input(mut line: String, peer_addr: &SocketAddr) -> Result<String> {
  40. if line.is_empty() {
  41. warn!("Received empty line from {}. ", peer_addr);
  42. warn!("Closing connection.");
  43. return Err(Error::ChannelStopped)
  44. }
  45. if &line[(line.len() - 2)..] != "\r\n" {
  46. warn!("Closing connection.");
  47. return Err(Error::ChannelStopped)
  48. }
  49. // Remove CRLF
  50. line.pop();
  51. line.pop();
  52. Ok(line)
  53. }
  54. struct Ircd {
  55. // msgs
  56. seen_msg_ids: SeenMsgIds,
  57. privmsgs_buffer: PrivmsgsBuffer,
  58. // channels
  59. autojoin_chans: Vec<String>,
  60. configured_chans: FxHashMap<String, ChannelInfo>,
  61. // p2p
  62. p2p: net::P2pPtr,
  63. p2p_receiver: Receiver<Privmsg>,
  64. }
  65. impl Ircd {
  66. fn new(
  67. seen_msg_ids: SeenMsgIds,
  68. autojoin_chans: Vec<String>,
  69. configured_chans: FxHashMap<String, ChannelInfo>,
  70. p2p: net::P2pPtr,
  71. p2p_receiver: Receiver<Privmsg>,
  72. ) -> Self {
  73. let privmsgs_buffer: PrivmsgsBuffer =
  74. Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSGS_BUFFER)));
  75. Self { seen_msg_ids, privmsgs_buffer, autojoin_chans, configured_chans, p2p, p2p_receiver }
  76. }
  77. async fn process(
  78. &self,
  79. executor: Arc<Executor<'_>>,
  80. stream: TcpStream,
  81. peer_addr: SocketAddr,
  82. ) -> Result<()> {
  83. let (reader, writer) = stream.split();
  84. let mut reader = BufReader::new(reader);
  85. let mut conn = IrcServerConnection::new(
  86. writer,
  87. self.seen_msg_ids.clone(),
  88. self.privmsgs_buffer.clone(),
  89. self.autojoin_chans.clone(),
  90. self.configured_chans.clone(),
  91. self.p2p.clone(),
  92. );
  93. let p2p_receiver = self.p2p_receiver.clone();
  94. let privmsgs_buffer = self.privmsgs_buffer.clone();
  95. executor.spawn(async move {
  96. loop {
  97. let mut line = String::new();
  98. futures::select! {
  99. privmsg = p2p_receiver.recv().fuse() => {
  100. let mut msg = privmsg?;
  101. info!("Received msg from P2p network: {:?}", msg);
  102. // Try to potentially decrypt the incoming message.
  103. if conn.configured_chans.contains_key(&msg.channel) {
  104. let chan_info = conn.configured_chans.get(&msg.channel).unwrap();
  105. if !chan_info.joined.load(Ordering::Relaxed) {
  106. continue
  107. }
  108. if let Some(salt_box) = &chan_info.salt_box {
  109. if let Some(decrypted_msg) = try_decrypt_message(salt_box, &msg.message) {
  110. msg.message = decrypted_msg;
  111. info!("Decrypted received message: {:?}", msg);
  112. }
  113. }
  114. }
  115. // add the msg to buffer
  116. {
  117. (*privmsgs_buffer.lock().await).push(msg.clone());
  118. }
  119. conn.reply(&msg.to_irc_msg()).await?;
  120. }
  121. err = reader.read_line(&mut line).fuse() => {
  122. if let Err(e) = err {
  123. warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
  124. return Ok(())
  125. }
  126. info!("Received msg from IRC client: {:?}", line);
  127. let irc_msg = match clean_input(line, &peer_addr) {
  128. Ok(m) => m,
  129. Err(e) => return Err(e)
  130. };
  131. info!("Send msg to IRC client '{}' from {}", irc_msg, peer_addr);
  132. if let Err(e) = conn.update(irc_msg).await {
  133. warn!("Connection error: {} for {}", e, peer_addr);
  134. return Err(Error::ChannelStopped)
  135. }
  136. }
  137. };
  138. }
  139. }).detach();
  140. Ok(())
  141. }
  142. }
  143. async_daemonize!(realmain);
  144. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  145. let seen_msg_ids = Arc::new(Mutex::new(vec![]));
  146. if settings.gen_secret {
  147. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  148. let encoded = bs58::encode(secret_key.as_bytes());
  149. println!("{}", encoded.into_string());
  150. return Ok(())
  151. }
  152. // Pick up channel settings from the TOML configuration
  153. let cfg_path = get_config_path(settings.config, CONFIG_FILE)?;
  154. let configured_chans = parse_configured_channels(&cfg_path)?;
  155. //
  156. // P2p setup
  157. //
  158. let net_settings = settings.net;
  159. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<Privmsg>();
  160. let p2p = net::P2p::new(net_settings.into()).await;
  161. let p2p = p2p.clone();
  162. let registry = p2p.protocol_registry();
  163. let seen_msg_ids_cloned = seen_msg_ids.clone();
  164. registry
  165. .register(net::SESSION_ALL, move |channel, p2p| {
  166. let sender = p2p_send_channel.clone();
  167. let seen_msg_ids_cloned = seen_msg_ids_cloned.clone();
  168. async move { ProtocolPrivmsg::init(channel, sender, p2p, seen_msg_ids_cloned).await }
  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. 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. }