main.rs 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174
  1. use async_std::net::{TcpListener, TcpStream};
  2. use std::{net::SocketAddr, path::PathBuf, sync::Arc};
  3. use async_channel::Receiver;
  4. use async_executor::Executor;
  5. use easy_parallel::Parallel;
  6. use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
  7. use log::{debug, error, info, warn};
  8. use simplelog::{ColorChoice, TermLogger, TerminalMode};
  9. use smol::future;
  10. use structopt_toml::StructOptToml;
  11. use darkfi::{
  12. async_daemonize,
  13. raft::Raft,
  14. rpc::rpcserver::{listen_and_serve, RpcServerConfig},
  15. util::{
  16. cli::{log_config, spawn_config},
  17. path::get_config_path,
  18. },
  19. Error, Result,
  20. };
  21. pub mod privmsg;
  22. pub mod rpc;
  23. pub mod server;
  24. pub mod settings;
  25. use crate::{
  26. privmsg::Privmsg,
  27. rpc::JsonRpcInterface,
  28. server::IrcServerConnection,
  29. settings::{Args, Command, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  30. };
  31. async fn process_user_input(
  32. mut line: String,
  33. peer_addr: SocketAddr,
  34. conn: &mut IrcServerConnection,
  35. sender: async_channel::Sender<Privmsg>,
  36. ) -> Result<()> {
  37. if line.is_empty() {
  38. warn!("Received empty line from {}. Closing connection.", peer_addr);
  39. return Err(Error::ChannelStopped)
  40. }
  41. assert!(&line[(line.len() - 2)..] == "\r\n");
  42. // Remove CRLF
  43. line.pop();
  44. line.pop();
  45. debug!("Received '{}' from {}", line, peer_addr);
  46. if let Err(e) = conn.update(line, sender).await {
  47. warn!("Connection error: {} for {}", e, peer_addr);
  48. return Err(Error::ChannelStopped)
  49. }
  50. Ok(())
  51. }
  52. async fn process(
  53. receiver: Receiver<Privmsg>,
  54. stream: TcpStream,
  55. peer_addr: SocketAddr,
  56. sender: async_channel::Sender<Privmsg>,
  57. ) -> Result<()> {
  58. let (reader, writer) = stream.split();
  59. let mut reader = BufReader::new(reader);
  60. let mut conn = IrcServerConnection::new(writer);
  61. loop {
  62. let mut line = String::new();
  63. futures::select! {
  64. privmsg = receiver.recv().fuse() => {
  65. let msg = privmsg?;
  66. debug!("ABOUT TO SEND: {:?}", msg);
  67. let irc_msg = format!(":{}!anon@dark.fi PRIVMSG {} :{}\r\n",
  68. msg.nickname,
  69. msg.channel,
  70. msg.message,
  71. );
  72. conn.reply(&irc_msg).await?;
  73. }
  74. err = reader.read_line(&mut line).fuse() => {
  75. if let Err(e) = err {
  76. warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
  77. return Ok(())
  78. }
  79. process_user_input(line, peer_addr, &mut conn, sender.clone()).await?;
  80. }
  81. };
  82. }
  83. }
  84. async_daemonize!(realmain);
  85. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  86. let listener = TcpListener::bind(settings.irc_listen).await?;
  87. let local_addr = listener.local_addr()?;
  88. info!("Listening on {}", local_addr);
  89. let datastore_path = PathBuf::from(&settings.datastore);
  90. let net_settings = match settings.command {
  91. Command::Net(s) => s,
  92. };
  93. //
  94. //Raft
  95. //
  96. let datastore_raft = datastore_path.join("ircd.db");
  97. let mut raft = Raft::<Privmsg>::new(net_settings.inbound, datastore_raft)?;
  98. let raft_sender = raft.get_broadcast();
  99. let commits = raft.get_commits();
  100. //
  101. // RPC interface
  102. //
  103. let rpc_config = RpcServerConfig {
  104. socket_addr: settings.rpc_listen,
  105. // TODO: Use net/transport:
  106. use_tls: false,
  107. identity_path: Default::default(),
  108. identity_pass: Default::default(),
  109. };
  110. let executor_cloned = executor.clone();
  111. let rpc_interface = Arc::new(JsonRpcInterface { addr: settings.rpc_listen });
  112. let rpc_task = executor.spawn(async move {
  113. listen_and_serve(rpc_config, rpc_interface, executor_cloned.clone()).await
  114. });
  115. //
  116. // IRC instance
  117. //
  118. let executor_cloned = executor.clone();
  119. let irc_task: smol::Task<Result<()>> = executor.spawn(async move {
  120. loop {
  121. let (stream, peer_addr) = match listener.accept().await {
  122. Ok((s, a)) => (s, a),
  123. Err(e) => {
  124. error!("Failed listening for connections: {}", e);
  125. return Err(Error::ServiceStopped)
  126. }
  127. };
  128. info!("Accepted client: {}", peer_addr);
  129. executor_cloned
  130. .spawn(process(commits.clone(), stream, peer_addr, raft_sender.clone()))
  131. .detach();
  132. }
  133. });
  134. let (signal, shutdown) = async_channel::bounded::<()>(1);
  135. ctrlc_async::set_async_handler(async move {
  136. warn!(target: "ircd", "ircd start Exit Signal");
  137. // cleaning up tasks running in the background
  138. signal.send(()).await.unwrap();
  139. rpc_task.cancel().await;
  140. irc_task.cancel().await;
  141. })
  142. .unwrap();
  143. // blocking
  144. raft.start(net_settings, executor.clone(), shutdown.clone()).await?;
  145. Ok(())
  146. }