main.rs 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224
  1. use async_std::net::{TcpListener, TcpStream};
  2. use std::{net::SocketAddr, sync::Arc};
  3. use async_channel::Receiver;
  4. use async_executor::Executor;
  5. use clap::Parser;
  6. use easy_parallel::Parallel;
  7. use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
  8. use log::{debug, error, info, warn};
  9. use simplelog::{ColorChoice, TermLogger, TerminalMode};
  10. use darkfi::{
  11. cli_desc, net,
  12. raft::Raft,
  13. rpc::rpcserver::{listen_and_serve, RpcServerConfig},
  14. util::cli::log_config,
  15. Error, Result,
  16. };
  17. pub(crate) mod privmsg;
  18. pub(crate) mod rpc;
  19. pub(crate) mod server;
  20. use crate::{privmsg::Privmsg, rpc::JsonRpcInterface, server::IrcServerConnection};
  21. #[derive(Parser)]
  22. #[clap(name = "ircd", about = cli_desc!(), version)]
  23. struct Args {
  24. /// Accept address
  25. #[clap(short, long)]
  26. accept: Option<SocketAddr>,
  27. /// Seed node (repeatable)
  28. #[clap(short, long)]
  29. seed: Vec<SocketAddr>,
  30. /// Manual connection (repeatable)
  31. #[clap(short, long)]
  32. connect: Vec<SocketAddr>,
  33. /// Connection slots
  34. #[clap(long, default_value_t = 0)]
  35. slots: u32,
  36. /// External address
  37. #[clap(short, long)]
  38. external: Option<SocketAddr>,
  39. /// IRC listen address
  40. #[clap(short = 'r', long, default_value = "127.0.0.1:6667")]
  41. irc: SocketAddr,
  42. /// RPC listen address
  43. #[clap(long, default_value = "127.0.0.1:8000")]
  44. rpc: SocketAddr,
  45. /// Verbosity level
  46. #[clap(short, parse(from_occurrences))]
  47. verbose: u8,
  48. }
  49. async fn process_user_input(
  50. mut line: String,
  51. peer_addr: SocketAddr,
  52. conn: &mut IrcServerConnection,
  53. sender: async_channel::Sender<Privmsg>,
  54. ) -> Result<()> {
  55. if line.is_empty() {
  56. warn!("Received empty line from {}. Closing connection.", peer_addr);
  57. return Err(Error::ChannelStopped)
  58. }
  59. assert!(&line[(line.len() - 2)..] == "\r\n");
  60. // Remove CRLF
  61. line.pop();
  62. line.pop();
  63. debug!("Received '{}' from {}", line, peer_addr);
  64. if let Err(e) = conn.update(line, sender).await {
  65. warn!("Connection error: {} for {}", e, peer_addr);
  66. return Err(Error::ChannelStopped)
  67. }
  68. Ok(())
  69. }
  70. async fn process(
  71. receiver: Receiver<Privmsg>,
  72. stream: TcpStream,
  73. peer_addr: SocketAddr,
  74. sender: async_channel::Sender<Privmsg>,
  75. ) -> Result<()> {
  76. let (reader, writer) = stream.split();
  77. let mut reader = BufReader::new(reader);
  78. let mut conn = IrcServerConnection::new(writer);
  79. loop {
  80. let mut line = String::new();
  81. futures::select! {
  82. privmsg = receiver.recv().fuse() => {
  83. let msg = privmsg.expect("internal message queue error");
  84. debug!("ABOUT TO SEND: {:?}", msg);
  85. let irc_msg = format!(":{}!anon@dark.fi PRIVMSG {} :{}\r\n",
  86. msg.nickname,
  87. msg.channel,
  88. msg.message,
  89. );
  90. conn.reply(&irc_msg).await?;
  91. }
  92. err = reader.read_line(&mut line).fuse() => {
  93. if let Err(e) = err {
  94. warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
  95. return Ok(())
  96. }
  97. process_user_input(line, peer_addr, &mut conn, sender.clone()).await?;
  98. }
  99. };
  100. }
  101. }
  102. async fn start(executor: Arc<Executor<'_>>, args: Args, net_settings: net::Settings) -> Result<()> {
  103. let listener = TcpListener::bind(args.irc).await?;
  104. let local_addr = listener.local_addr()?;
  105. info!("Listening on {}", local_addr);
  106. //
  107. // Raft
  108. //
  109. let mut raft = Raft::<Privmsg>::new(net_settings.inbound, std::path::PathBuf::from("msgs.db"))?;
  110. let raft_sender = raft.get_broadcast();
  111. let commits = raft.get_commits();
  112. //
  113. // RPC interface
  114. let rpc_config = RpcServerConfig {
  115. socket_addr: args.rpc,
  116. // TODO: Use net/transport:
  117. use_tls: false,
  118. identity_path: Default::default(),
  119. identity_pass: Default::default(),
  120. };
  121. let executor_cloned = executor.clone();
  122. let rpc_interface = Arc::new(JsonRpcInterface { addr: args.rpc });
  123. let rpc_task = executor.spawn(async move {
  124. listen_and_serve(rpc_config, rpc_interface, executor_cloned.clone()).await
  125. });
  126. //
  127. // IRC instance
  128. //
  129. let executor_cloned = executor.clone();
  130. let irc_task: smol::Task<Result<()>> = executor.spawn(async move {
  131. loop {
  132. let (stream, peer_addr) = match listener.accept().await {
  133. Ok((s, a)) => (s, a),
  134. Err(e) => {
  135. error!("Failed listening for connections: {}", e);
  136. return Err(Error::ServiceStopped)
  137. }
  138. };
  139. info!("Accepted client: {}", peer_addr);
  140. executor_cloned
  141. .spawn(process(commits.clone(), stream, peer_addr, raft_sender.clone()))
  142. .detach();
  143. }
  144. });
  145. let stop_signal = async_channel::bounded::<()>(10);
  146. ctrlc_async::set_async_handler(async move {
  147. warn!(target: "ircd", "ircd start() Exit Signal");
  148. // cleaning up tasks running in the background
  149. stop_signal.0.send(()).await.expect("send exit signal to raft");
  150. rpc_task.cancel().await;
  151. irc_task.cancel().await;
  152. })
  153. .expect("handle exit signal");
  154. // blocking
  155. raft.start(net_settings.clone(), executor.clone(), stop_signal.1.clone()).await?;
  156. Ok(())
  157. }
  158. fn main() -> Result<()> {
  159. let args = Args::parse();
  160. let (lvl, conf) = log_config(args.verbose.into())?;
  161. TermLogger::init(lvl, conf, TerminalMode::Mixed, ColorChoice::Auto)?;
  162. let net_settings = net::Settings {
  163. inbound: args.accept,
  164. outbound_connections: args.slots,
  165. external_addr: args.external,
  166. peers: args.connect.clone(),
  167. seeds: args.seed.clone(),
  168. ..Default::default()
  169. };
  170. let ex = Arc::new(Executor::new());
  171. let ex_clone = ex.clone();
  172. let (signal, shutdown) = async_channel::unbounded::<()>();
  173. let (_, result) = Parallel::new()
  174. .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
  175. // Run the main future on the current thread.
  176. .finish(|| {
  177. smol::future::block_on(async move {
  178. start(ex_clone.clone(), args, net_settings).await?;
  179. drop(signal);
  180. Ok::<(), darkfi::Error>(())
  181. })
  182. });
  183. result
  184. }