main.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330
  1. use async_std::{
  2. net::TcpListener,
  3. sync::{Arc, Mutex},
  4. };
  5. use std::{fs::File, net::SocketAddr};
  6. use async_channel::Receiver;
  7. use async_executor::Executor;
  8. use futures::{io::BufReader, AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWrite, FutureExt};
  9. use futures_rustls::{rustls, TlsAcceptor};
  10. use fxhash::FxHashMap;
  11. use log::{error, info, warn};
  12. use rand::rngs::OsRng;
  13. use smol::future;
  14. use structopt_toml::StructOptToml;
  15. use darkfi::{
  16. async_daemonize, net,
  17. rpc::server::listen_and_serve,
  18. system::{Subscriber, SubscriberPtr},
  19. util::{
  20. cli::{get_log_config, get_log_level, spawn_config},
  21. expand_path,
  22. path::get_config_path,
  23. },
  24. Error, Result,
  25. };
  26. pub mod crypto;
  27. pub mod privmsg;
  28. pub mod protocol_privmsg;
  29. pub mod rpc;
  30. pub mod server;
  31. pub mod settings;
  32. use crate::{
  33. privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
  34. protocol_privmsg::ProtocolPrivmsg,
  35. rpc::JsonRpcInterface,
  36. server::IrcServerConnection,
  37. settings::{
  38. parse_configured_channels, parse_configured_contacts, Args, ChannelInfo, CONFIG_FILE,
  39. CONFIG_FILE_CONTENTS,
  40. },
  41. };
  42. const SIZE_OF_MSG_IDSS_BUFFER: usize = 65536;
  43. const SIZE_OF_MSGS_BUFFER: usize = 4096;
  44. pub const MAXIMUM_LENGTH_OF_MESSAGE: usize = 1024;
  45. pub const MAXIMUM_LENGTH_OF_NICKNAME: usize = 32;
  46. struct Ircd {
  47. // msgs
  48. seen_msg_ids: SeenMsgIds,
  49. privmsgs_buffer: PrivmsgsBuffer,
  50. // channels
  51. autojoin_chans: Vec<String>,
  52. configured_chans: FxHashMap<String, ChannelInfo>,
  53. configured_contacts: FxHashMap<String, crypto_box::Box>,
  54. // p2p
  55. p2p: net::P2pPtr,
  56. senders: SubscriberPtr<Privmsg>,
  57. }
  58. impl Ircd {
  59. fn new(
  60. seen_msg_ids: SeenMsgIds,
  61. privmsgs_buffer: PrivmsgsBuffer,
  62. autojoin_chans: Vec<String>,
  63. configured_chans: FxHashMap<String, ChannelInfo>,
  64. configured_contacts: FxHashMap<String, crypto_box::Box>,
  65. p2p: net::P2pPtr,
  66. ) -> Self {
  67. let senders = Subscriber::new();
  68. Self {
  69. seen_msg_ids,
  70. privmsgs_buffer,
  71. autojoin_chans,
  72. configured_chans,
  73. configured_contacts,
  74. p2p,
  75. senders,
  76. }
  77. }
  78. fn start_p2p_receive_loop(&self, executor: Arc<Executor<'_>>, p2p_receiver: Receiver<Privmsg>) {
  79. let senders = self.senders.clone();
  80. executor
  81. .spawn(async move {
  82. while let Ok(msg) = p2p_receiver.recv().await {
  83. senders.notify(msg).await;
  84. }
  85. })
  86. .detach();
  87. }
  88. async fn process_new_connection<C: AsyncRead + AsyncWrite + Send + Unpin + 'static>(
  89. &self,
  90. executor: Arc<Executor<'_>>,
  91. stream: C,
  92. peer_addr: SocketAddr,
  93. ) -> Result<()> {
  94. let (reader, writer) = stream.split();
  95. let mut reader = BufReader::new(reader);
  96. // New subscriber
  97. let receiver = self.senders.clone().subscribe().await;
  98. // New irc connection
  99. let mut conn = IrcServerConnection::new(
  100. writer,
  101. peer_addr,
  102. self.seen_msg_ids.clone(),
  103. self.privmsgs_buffer.clone(),
  104. self.autojoin_chans.clone(),
  105. self.configured_chans.clone(),
  106. self.configured_contacts.clone(),
  107. self.p2p.clone(),
  108. self.senders.clone(),
  109. receiver.get_id(),
  110. );
  111. executor
  112. .spawn(async move {
  113. loop {
  114. let mut line = String::new();
  115. let result: Result<()> = futures::select! {
  116. msg = receiver.receive().fuse() => {
  117. match conn.process_msg_from_p2p(&msg).await {
  118. Ok(_) => Ok(()),
  119. Err(e) => {
  120. error!("Process msg from p2p failed {}: {}", peer_addr, e);
  121. Err(Error::ChannelStopped)
  122. }
  123. }
  124. }
  125. err = reader.read_line(&mut line).fuse() => {
  126. match conn.process_line_from_client(err, line).await {
  127. Ok(_) => Ok(()),
  128. Err(e) => {
  129. error!("Process line from client failed {}: {}", peer_addr, e);
  130. Err(Error::ChannelStopped)
  131. }
  132. }
  133. }
  134. };
  135. if let Err(e) = result {
  136. warn!("Close connection for clinet {}: {}", peer_addr, e);
  137. receiver.unsubscribe().await;
  138. break
  139. }
  140. }
  141. })
  142. .detach();
  143. Ok(())
  144. }
  145. }
  146. async_daemonize!(realmain);
  147. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  148. let seen_msg_ids =
  149. Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSG_IDSS_BUFFER)));
  150. let privmsgs_buffer: PrivmsgsBuffer =
  151. Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSGS_BUFFER)));
  152. if settings.gen_secret {
  153. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  154. let encoded = bs58::encode(secret_key.as_bytes());
  155. println!("{}", encoded.into_string());
  156. return Ok(())
  157. }
  158. // Pick up channel settings from the TOML configuration
  159. let cfg_path = get_config_path(settings.config, CONFIG_FILE)?;
  160. let toml_contents = std::fs::read_to_string(cfg_path)?;
  161. let configured_chans = parse_configured_channels(&toml_contents)?;
  162. let configured_contacts = parse_configured_contacts(&toml_contents)?;
  163. //
  164. // P2p setup
  165. //
  166. let net_settings = settings.net;
  167. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<Privmsg>();
  168. let p2p = net::P2p::new(net_settings.into()).await;
  169. let p2p2 = p2p.clone();
  170. let registry = p2p.protocol_registry();
  171. let seen_msg_ids_cloned = seen_msg_ids.clone();
  172. let privmsgs_buffer_cloned = privmsgs_buffer.clone();
  173. registry
  174. .register(net::SESSION_ALL, move |channel, p2p| {
  175. let sender = p2p_send_channel.clone();
  176. let seen_msg_ids_cloned = seen_msg_ids_cloned.clone();
  177. let privmsgs_buffer_cloned = privmsgs_buffer_cloned.clone();
  178. async move {
  179. ProtocolPrivmsg::init(
  180. channel,
  181. sender,
  182. p2p,
  183. seen_msg_ids_cloned,
  184. privmsgs_buffer_cloned,
  185. )
  186. .await
  187. }
  188. })
  189. .await;
  190. p2p.clone().start(executor.clone()).await?;
  191. let executor_cloned = executor.clone();
  192. executor_cloned.spawn(p2p.clone().run(executor.clone())).detach();
  193. //
  194. // RPC interface
  195. //
  196. let rpc_listen_addr = settings.rpc_listen.clone();
  197. let rpc_interface =
  198. Arc::new(JsonRpcInterface { addr: rpc_listen_addr.clone(), p2p: p2p.clone() });
  199. executor.spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface).await }).detach();
  200. //
  201. // IRC instance
  202. //
  203. let listenaddr = settings.irc_listen.socket_addrs(|| None)?[0];
  204. let listener = TcpListener::bind(listenaddr).await?;
  205. let acceptor = match settings.irc_listen.scheme() {
  206. "tls" => {
  207. // openssl genpkey -algorithm ED25519 > example.com.key
  208. // openssl req -new -out example.com.csr -key example.com.key
  209. // openssl x509 -req -days 700 -in example.com.csr -signkey example.com.key -out example.com.crt
  210. if settings.irc_tls_secret.is_none() || settings.irc_tls_cert.is_none() {
  211. error!("To listen using TLS, please set irc_tls_secret and irc_tls_cert in your config file.");
  212. return Err(Error::KeypairPathNotFound)
  213. }
  214. let file = File::open(expand_path(&settings.irc_tls_secret.unwrap())?)?;
  215. let mut reader = std::io::BufReader::new(file);
  216. let secret = &rustls_pemfile::pkcs8_private_keys(&mut reader)?[0];
  217. let secret = rustls::PrivateKey(secret.clone());
  218. let file = File::open(expand_path(&settings.irc_tls_cert.unwrap())?)?;
  219. let mut reader = std::io::BufReader::new(file);
  220. let certificate = &rustls_pemfile::certs(&mut reader)?[0];
  221. let certificate = rustls::Certificate(certificate.clone());
  222. let config = rustls::ServerConfig::builder()
  223. .with_safe_defaults()
  224. .with_no_client_auth()
  225. .with_single_cert(vec![certificate], secret)?;
  226. let acceptor = TlsAcceptor::from(Arc::new(config));
  227. Some(acceptor)
  228. }
  229. _ => None,
  230. };
  231. info!("IRC listening on {}", settings.irc_listen);
  232. let executor_cloned = executor.clone();
  233. executor
  234. .spawn(async move {
  235. let ircd = Ircd::new(
  236. seen_msg_ids.clone(),
  237. privmsgs_buffer.clone(),
  238. settings.autojoin.clone(),
  239. configured_chans.clone(),
  240. configured_contacts.clone(),
  241. p2p.clone(),
  242. );
  243. ircd.start_p2p_receive_loop(executor_cloned.clone(), p2p_recv_channel);
  244. loop {
  245. let (stream, peer_addr) = match listener.accept().await {
  246. Ok((s, a)) => (s, a),
  247. Err(e) => {
  248. error!("failed accepting new connections: {}", e);
  249. continue
  250. }
  251. };
  252. let result = if let Some(acceptor) = acceptor.clone() {
  253. let stream = match acceptor.accept(stream).await {
  254. Ok(s) => s,
  255. Err(e) => {
  256. error!("Failed accepting TLS connection: {}", e);
  257. continue
  258. }
  259. };
  260. ircd.process_new_connection(executor_cloned.clone(), stream, peer_addr).await
  261. } else {
  262. ircd.process_new_connection(executor_cloned.clone(), stream, peer_addr).await
  263. };
  264. if let Err(e) = result {
  265. error!("Failed processing connection {}: {}", peer_addr, e);
  266. continue
  267. };
  268. info!("IRC Accepted new client: {}", peer_addr);
  269. }
  270. })
  271. .detach();
  272. // Run once receive exit signal
  273. let (signal, shutdown) = async_channel::bounded::<()>(1);
  274. ctrlc_async::set_async_handler(async move {
  275. warn!(target: "ircd", "ircd start Exit Signal");
  276. // cleaning up tasks running in the background
  277. signal.send(()).await.unwrap();
  278. })
  279. .unwrap();
  280. // Wait for SIGINT
  281. shutdown.recv().await?;
  282. print!("\r");
  283. info!("Caught termination signal, cleaning up and exiting...");
  284. p2p2.stop().await;
  285. Ok(())
  286. }