main.rs 6.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230
  1. use async_std::{
  2. net::{TcpListener, TcpStream},
  3. sync::{Arc, Mutex},
  4. };
  5. use std::net::SocketAddr;
  6. use async_channel::{Receiver, Sender};
  7. use async_executor::Executor;
  8. use easy_parallel::Parallel;
  9. use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
  10. use log::{debug, error, info, warn};
  11. use simplelog::{ColorChoice, TermLogger, TerminalMode};
  12. use smol::future;
  13. use structopt_toml::StructOptToml;
  14. use darkfi::{
  15. async_daemonize, net,
  16. raft::{NetMsg, ProtocolRaft, Raft},
  17. rpc::rpcserver::{listen_and_serve, RpcServerConfig},
  18. util::{
  19. cli::{log_config, spawn_config},
  20. path::{expand_path, get_config_path},
  21. },
  22. Error, Result,
  23. };
  24. pub mod privmsg;
  25. pub mod rpc;
  26. pub mod server;
  27. pub mod settings;
  28. use crate::{
  29. privmsg::Privmsg,
  30. rpc::JsonRpcInterface,
  31. server::IrcServerConnection,
  32. settings::{Args, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  33. };
  34. pub type SeenMsgIds = Arc<Mutex<Vec<u32>>>;
  35. fn build_irc_msg(msg: &Privmsg) -> String {
  36. debug!("ABOUT TO SEND: {:?}", msg);
  37. let irc_msg =
  38. format!(":{}!anon@dark.fi PRIVMSG {} :{}\r\n", msg.nickname, msg.channel, msg.message,);
  39. irc_msg
  40. }
  41. fn clean_input(mut line: String, peer_addr: &SocketAddr) -> Result<String> {
  42. if line.is_empty() {
  43. warn!("Received empty line from {}. ", peer_addr);
  44. warn!("Closing connection.");
  45. return Err(Error::ChannelStopped)
  46. }
  47. if &line[(line.len() - 2)..] != "\r\n" {
  48. warn!("Closing connection.");
  49. return Err(Error::ChannelStopped)
  50. }
  51. // Remove CRLF
  52. line.pop();
  53. line.pop();
  54. Ok(line)
  55. }
  56. async fn broadcast_msg(
  57. irc_msg: String,
  58. peer_addr: SocketAddr,
  59. conn: &mut IrcServerConnection,
  60. ) -> Result<()> {
  61. info!("Send msg to IRC server '{}' from {}", irc_msg, peer_addr);
  62. if let Err(e) = conn.update(irc_msg).await {
  63. warn!("Connection error: {} for {}", e, peer_addr);
  64. return Err(Error::ChannelStopped)
  65. }
  66. Ok(())
  67. }
  68. async fn process(
  69. raft_receiver: Receiver<Privmsg>,
  70. stream: TcpStream,
  71. peer_addr: SocketAddr,
  72. raft_sender: Sender<Privmsg>,
  73. seen_msg_id: SeenMsgIds,
  74. ) -> Result<()> {
  75. let (reader, writer) = stream.split();
  76. let mut reader = BufReader::new(reader);
  77. let mut conn = IrcServerConnection::new(writer, seen_msg_id.clone(), raft_sender);
  78. loop {
  79. let mut line = String::new();
  80. futures::select! {
  81. privmsg = raft_receiver.recv().fuse() => {
  82. info!("Receive msg from raft");
  83. let msg = privmsg?;
  84. let mut smi = seen_msg_id.lock().await;
  85. if smi.contains(&msg.id) {
  86. continue
  87. }
  88. smi.push(msg.id);
  89. drop(smi);
  90. let irc_msg = build_irc_msg(&msg);
  91. conn.reply(&irc_msg).await?;
  92. }
  93. err = reader.read_line(&mut line).fuse() => {
  94. if let Err(e) = err {
  95. warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
  96. return Ok(())
  97. }
  98. info!("Receive msg from IRC server");
  99. let irc_msg = match clean_input(line, &peer_addr) {
  100. Ok(m) => m,
  101. Err(e) => return Err(e)
  102. };
  103. broadcast_msg(irc_msg, peer_addr,&mut conn).await?;
  104. }
  105. };
  106. }
  107. }
  108. async_daemonize!(realmain);
  109. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  110. let seen_msg_id: SeenMsgIds = Arc::new(Mutex::new(vec![]));
  111. //
  112. //Raft
  113. //
  114. let datastore_path = expand_path(&settings.datastore)?;
  115. let net_settings = settings.net;
  116. let datastore_raft = datastore_path.join("ircd.db");
  117. let mut raft = Raft::<Privmsg>::new(net_settings.inbound, datastore_raft)?;
  118. let raft_sender = raft.get_broadcast();
  119. let raft_receiver = raft.get_commits();
  120. // P2p setup
  121. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<NetMsg>();
  122. let p2p = net::P2p::new(net_settings.into()).await;
  123. let p2p = p2p.clone();
  124. let registry = p2p.protocol_registry();
  125. let seen_net_msg = Arc::new(Mutex::new(vec![]));
  126. let raft_node_id = raft.id.clone();
  127. registry
  128. .register(net::SESSION_ALL, move |channel, p2p| {
  129. let raft_node_id = raft_node_id.clone();
  130. let sender = p2p_send_channel.clone();
  131. let seen_net_msg_cloned = seen_net_msg.clone();
  132. async move {
  133. ProtocolRaft::init(raft_node_id, channel, sender, p2p, seen_net_msg_cloned).await
  134. }
  135. })
  136. .await;
  137. p2p.clone().start(executor.clone()).await?;
  138. let executor_cloned = executor.clone();
  139. let p2p_run_task = executor_cloned.spawn(p2p.clone().run(executor.clone()));
  140. //
  141. // RPC interface
  142. //
  143. let rpc_config = RpcServerConfig {
  144. socket_addr: settings.rpc_listen,
  145. use_tls: false,
  146. identity_path: Default::default(),
  147. identity_pass: Default::default(),
  148. };
  149. let executor_cloned = executor.clone();
  150. let rpc_interface = Arc::new(JsonRpcInterface { addr: settings.rpc_listen, p2p: p2p.clone() });
  151. let rpc_task = executor.spawn(async move {
  152. listen_and_serve(rpc_config, rpc_interface, executor_cloned.clone()).await
  153. });
  154. //
  155. // IRC instance
  156. //
  157. let listener = TcpListener::bind(settings.irc_listen).await?;
  158. let local_addr = listener.local_addr()?;
  159. info!("IRC listening on {}", local_addr);
  160. let executor_cloned = executor.clone();
  161. let raft_receiver_cloned = raft_receiver.clone();
  162. let irc_task: smol::Task<Result<()>> = executor.spawn(async move {
  163. loop {
  164. let (stream, peer_addr) = match listener.accept().await {
  165. Ok((s, a)) => (s, a),
  166. Err(e) => {
  167. error!("Failed listening for connections: {}", e);
  168. return Err(Error::ServiceStopped)
  169. }
  170. };
  171. info!("IRC Accepted client: {}", peer_addr);
  172. executor_cloned
  173. .spawn(process(
  174. raft_receiver_cloned.clone(),
  175. stream,
  176. peer_addr,
  177. raft_sender.clone(),
  178. seen_msg_id.clone(),
  179. ))
  180. .detach();
  181. }
  182. });
  183. // Run once receive exit signal
  184. let (signal, shutdown) = async_channel::bounded::<()>(1);
  185. ctrlc_async::set_async_handler(async move {
  186. warn!(target: "ircd", "ircd start Exit Signal");
  187. // cleaning up tasks running in the background
  188. signal.send(()).await.unwrap();
  189. rpc_task.cancel().await;
  190. irc_task.cancel().await;
  191. p2p_run_task.cancel().await;
  192. })
  193. .unwrap();
  194. // blocking
  195. raft.start(p2p.clone(), p2p_recv_channel.clone(), executor.clone(), shutdown.clone()).await?;
  196. Ok(())
  197. }