main.rs 11 KB

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