main.rs 12 KB

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