main.rs 6.7 KB

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