main.rs 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261
  1. use std::{net::SocketAddr, sync::atomic::Ordering};
  2. use async_channel::{Receiver, Sender};
  3. use async_executor::Executor;
  4. use async_std::{
  5. net::{TcpListener, TcpStream},
  6. sync::{Arc, Mutex},
  7. };
  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 smol::future;
  13. use structopt_toml::StructOptToml;
  14. use darkfi::{
  15. async_daemonize, net,
  16. raft::{NetMsg, ProtocolRaft, Raft},
  17. rpc::server::listen_and_serve,
  18. util::{
  19. cli::{get_log_config, get_log_level, spawn_config},
  20. path::{expand_path, get_config_path},
  21. },
  22. Error, Result,
  23. };
  24. pub mod crypto;
  25. pub mod privmsg;
  26. pub mod rpc;
  27. pub mod server;
  28. pub mod settings;
  29. use crate::{
  30. crypto::try_decrypt_message,
  31. privmsg::Privmsg,
  32. rpc::JsonRpcInterface,
  33. server::IrcServerConnection,
  34. settings::{parse_configured_channels, Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  35. };
  36. pub type SeenMsgIds = Arc<Mutex<Vec<u32>>>;
  37. fn build_irc_msg(msg: &Privmsg) -> String {
  38. debug!("ABOUT TO SEND: {:?}", msg);
  39. let irc_msg =
  40. format!(":{}!anon@dark.fi PRIVMSG {} :{}\r\n", msg.nickname, msg.channel, msg.message);
  41. irc_msg
  42. }
  43. fn clean_input(mut line: String, peer_addr: &SocketAddr) -> Result<String> {
  44. if line.is_empty() {
  45. warn!("Received empty line from {}. ", peer_addr);
  46. warn!("Closing connection.");
  47. return Err(Error::ChannelStopped)
  48. }
  49. if &line[(line.len() - 2)..] != "\r\n" {
  50. warn!("Closing connection.");
  51. return Err(Error::ChannelStopped)
  52. }
  53. // Remove CRLF
  54. line.pop();
  55. line.pop();
  56. Ok(line)
  57. }
  58. async fn broadcast_msg(
  59. irc_msg: String,
  60. peer_addr: SocketAddr,
  61. conn: &mut IrcServerConnection,
  62. ) -> Result<()> {
  63. info!("Send msg to IRC client '{}' from {}", irc_msg, peer_addr);
  64. if let Err(e) = conn.update(irc_msg).await {
  65. warn!("Connection error: {} for {}", e, peer_addr);
  66. return Err(Error::ChannelStopped)
  67. }
  68. Ok(())
  69. }
  70. async fn process(
  71. raft_receiver: Receiver<Privmsg>,
  72. stream: TcpStream,
  73. peer_addr: SocketAddr,
  74. raft_sender: Sender<Privmsg>,
  75. seen_msg_id: SeenMsgIds,
  76. autojoin_chans: Vec<String>,
  77. configured_chans: FxHashMap<String, ChannelInfo>,
  78. ) -> Result<()> {
  79. let (reader, writer) = stream.split();
  80. let mut reader = BufReader::new(reader);
  81. let mut conn = IrcServerConnection::new(
  82. writer,
  83. seen_msg_id.clone(),
  84. raft_sender,
  85. autojoin_chans,
  86. configured_chans,
  87. );
  88. loop {
  89. let mut line = String::new();
  90. futures::select! {
  91. privmsg = raft_receiver.recv().fuse() => {
  92. let mut msg = privmsg?;
  93. info!("Received msg from Raft: {:?}", msg);
  94. let mut smi = seen_msg_id.lock().await;
  95. if smi.contains(&msg.id) {
  96. continue
  97. }
  98. smi.push(msg.id);
  99. drop(smi);
  100. // Try to potentially decrypt the incoming message.
  101. if conn.configured_chans.contains_key(&msg.channel) {
  102. let chan_info = conn.configured_chans.get(&msg.channel).unwrap();
  103. if !chan_info.joined.load(Ordering::Relaxed) {
  104. continue
  105. }
  106. if let Some(salt_box) = &chan_info.salt_box {
  107. if let Some(decrypted_msg) = try_decrypt_message(salt_box, &msg.message) {
  108. msg.message = decrypted_msg;
  109. info!("Decrypted received message: {:?}", msg);
  110. }
  111. }
  112. }
  113. let irc_msg = build_irc_msg(&msg);
  114. conn.reply(&irc_msg).await?;
  115. }
  116. err = reader.read_line(&mut line).fuse() => {
  117. if let Err(e) = err {
  118. warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
  119. return Ok(())
  120. }
  121. info!("Received msg from IRC client: {:?}", line);
  122. let irc_msg = match clean_input(line, &peer_addr) {
  123. Ok(m) => m,
  124. Err(e) => return Err(e)
  125. };
  126. broadcast_msg(irc_msg, peer_addr,&mut conn).await?;
  127. }
  128. };
  129. }
  130. }
  131. async_daemonize!(realmain);
  132. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  133. if settings.gen_secret {
  134. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  135. let encoded = bs58::encode(secret_key.as_bytes());
  136. println!("{}", encoded.into_string());
  137. return Ok(())
  138. }
  139. let seen_msg_id: SeenMsgIds = Arc::new(Mutex::new(vec![]));
  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. //Raft
  145. //
  146. let datastore_path = expand_path(&settings.datastore)?;
  147. let net_settings = settings.net;
  148. let datastore_raft = datastore_path.join("ircd.db");
  149. let mut raft = Raft::<Privmsg>::new(net_settings.inbound.clone(), datastore_raft)?;
  150. let raft_sender = raft.get_msgs_channel();
  151. let raft_receiver = raft.get_commits_channel();
  152. // P2p setup
  153. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<NetMsg>();
  154. let p2p = net::P2p::new(net_settings.into()).await;
  155. let p2p = p2p.clone();
  156. let registry = p2p.protocol_registry();
  157. let seen_net_msg = Arc::new(Mutex::new(vec![]));
  158. let raft_node_id = raft.id.clone();
  159. registry
  160. .register(net::SESSION_ALL, move |channel, p2p| {
  161. let raft_node_id = raft_node_id.clone();
  162. let sender = p2p_send_channel.clone();
  163. let seen_net_msg_cloned = seen_net_msg.clone();
  164. async move {
  165. ProtocolRaft::init(raft_node_id, channel, sender, p2p, seen_net_msg_cloned).await
  166. }
  167. })
  168. .await;
  169. p2p.clone().start(executor.clone()).await?;
  170. let executor_cloned = executor.clone();
  171. let p2p_run_task = executor_cloned.spawn(p2p.clone().run(executor.clone()));
  172. //
  173. // RPC interface
  174. //
  175. let rpc_listen_addr = settings.rpc_listen.clone();
  176. let rpc_interface =
  177. Arc::new(JsonRpcInterface { addr: rpc_listen_addr.clone(), p2p: p2p.clone() });
  178. let rpc_task =
  179. executor.spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface).await });
  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. let raft_receiver_cloned = raft_receiver.clone();
  189. let irc_task: smol::Task<Result<()>> = executor.spawn(async move {
  190. loop {
  191. let (stream, peer_addr) = match listener.accept().await {
  192. Ok((s, a)) => (s, a),
  193. Err(e) => {
  194. error!("Failed listening for connections: {}", e);
  195. return Err(Error::NetworkServiceStopped)
  196. }
  197. };
  198. info!("IRC Accepted client: {}", peer_addr);
  199. executor_cloned
  200. .spawn(process(
  201. raft_receiver_cloned.clone(),
  202. stream,
  203. peer_addr,
  204. raft_sender.clone(),
  205. seen_msg_id.clone(),
  206. settings.autojoin.clone(),
  207. configured_chans.clone(),
  208. ))
  209. .detach();
  210. }
  211. });
  212. // Run once receive exit signal
  213. let (signal, shutdown) = async_channel::bounded::<()>(1);
  214. ctrlc::set_handler(move || {
  215. warn!(target: "ircd", "ircd start Exit Signal");
  216. // cleaning up tasks running in the background
  217. async_std::task::block_on(signal.send(())).unwrap();
  218. async_std::task::block_on(rpc_task.cancel());
  219. async_std::task::block_on(irc_task.cancel());
  220. async_std::task::block_on(p2p_run_task.cancel());
  221. })
  222. .unwrap();
  223. // blocking
  224. raft.start(p2p.clone(), p2p_recv_channel.clone(), executor.clone(), shutdown.clone()).await?;
  225. Ok(())
  226. }