main.rs 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214
  1. #[macro_use]
  2. extern crate clap;
  3. use async_trait::async_trait;
  4. use std::{
  5. net::{SocketAddr, TcpListener, TcpStream},
  6. sync::Arc,
  7. };
  8. use async_executor::Executor;
  9. use async_std::io::BufReader;
  10. use futures::{
  11. AsyncBufReadExt, AsyncReadExt, FutureExt,
  12. };
  13. use serde_json::{json, Value};
  14. use log::{debug, error, info, warn};
  15. use simplelog::{ColorChoice, LevelFilter, TermLogger, TerminalMode};
  16. use smol::Async;
  17. use drk::{
  18. net,
  19. rpc::{
  20. jsonrpc::{
  21. error as jsonerr, request as jsonreq, response as jsonresp, send_raw_request,
  22. ErrorCode::*, JsonRequest, JsonResult,
  23. },
  24. rpcserver::{listen_and_serve, RequestHandler, RpcServerConfig},
  25. },
  26. Error, Result,
  27. };
  28. mod privmsg;
  29. mod program_options;
  30. mod protocol_privmsg;
  31. mod irc_server;
  32. use crate::privmsg::PrivMsg;
  33. use crate::program_options::ProgramOptions;
  34. use crate::protocol_privmsg::ProtocolPrivMsg;
  35. use crate::irc_server::IrcServerConnection;
  36. async fn process(
  37. recvr: async_channel::Receiver<Arc<PrivMsg>>,
  38. stream: Async<TcpStream>,
  39. peer_addr: SocketAddr,
  40. p2p: net::P2pPtr,
  41. _executor: Arc<Executor<'_>>,
  42. ) -> Result<()> {
  43. let (reader, writer) = stream.split();
  44. let mut reader = BufReader::new(reader);
  45. let mut connection = IrcServerConnection::new(writer);
  46. loop {
  47. let mut line = String::new();
  48. futures::select! {
  49. privmsg = recvr.recv().fuse() => {
  50. let privmsg = privmsg.expect("internal message queue error");
  51. debug!("ABOUT TO SEND {:?}", privmsg);
  52. let irc_msg = format!(
  53. ":{}!darkfi@127.0.0.1 PRIVMSG {} :{}\n",
  54. privmsg.nickname,
  55. privmsg.channel,
  56. privmsg.message
  57. );
  58. connection.reply(&irc_msg).await?;
  59. }
  60. err = reader.read_line(&mut line).fuse() => {
  61. if let Err(err) = err {
  62. warn!("Read line error. Closing stream for {}: {}", peer_addr, err);
  63. return Ok(())
  64. }
  65. process_user_input(line, peer_addr, &mut connection, p2p.clone()).await;
  66. }
  67. };
  68. }
  69. }
  70. async fn process_user_input(
  71. mut line: String,
  72. peer_addr: SocketAddr,
  73. connection: &mut IrcServerConnection,
  74. p2p: net::P2pPtr,
  75. ) {
  76. if line.len() == 0 {
  77. warn!("Received empty line from {}. Closing connection.", peer_addr);
  78. return
  79. }
  80. assert!(&line[(line.len() - 1)..] == "\n");
  81. // Remove the \n character
  82. line.pop();
  83. debug!("Received '{}' from {}", line, peer_addr);
  84. if let Err(err) = connection.update(line, p2p.clone()).await {
  85. warn!("Connection error: {} for {}", err, peer_addr);
  86. return
  87. }
  88. }
  89. async fn channel_loop(
  90. p2p: net::P2pPtr,
  91. sender: async_channel::Sender<Arc<PrivMsg>>,
  92. executor: Arc<Executor<'_>>,
  93. ) -> Result<()> {
  94. debug!("CHANNEL SUBS LOOP");
  95. let new_channel_sub = p2p.subscribe_channel().await;
  96. loop {
  97. let channel = new_channel_sub.receive().await?;
  98. debug!("NEWCHANNEL");
  99. let protocol_privmsg = ProtocolPrivMsg::new(channel, sender.clone(), p2p.clone()).await;
  100. protocol_privmsg.start(executor.clone()).await;
  101. }
  102. }
  103. async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
  104. let listener = match Async::<TcpListener>::bind(options.irc_accept_addr) {
  105. Ok(listener) => listener,
  106. Err(err) => {
  107. error!("Bind listener failed: {}", err);
  108. return Err(Error::OperationFailed)
  109. }
  110. };
  111. let local_addr = match listener.get_ref().local_addr() {
  112. Ok(addr) => addr,
  113. Err(err) => {
  114. error!("Failed to get local address: {}", err);
  115. return Err(Error::OperationFailed)
  116. }
  117. };
  118. info!("Listening on {}", local_addr);
  119. let p2p = net::P2p::new(options.network_settings);
  120. // Performs seed session
  121. p2p.clone().start(executor.clone()).await?;
  122. // Actual main p2p session
  123. let ex2 = executor.clone();
  124. let p2p2 = p2p.clone();
  125. executor
  126. .spawn(async move {
  127. if let Err(err) = p2p2.run(ex2).await {
  128. error!("Error: p2p run failed {}", err);
  129. }
  130. })
  131. .detach();
  132. let (sender, recvr) = async_channel::unbounded();
  133. // for now the p2p and channel sub sessions just run forever
  134. // so detach them as background processes.
  135. executor.spawn(channel_loop(p2p.clone(), sender, executor.clone())).detach();
  136. let rpc_interface = Arc::new(JsonRpcInterface {});
  137. executor.spawn(async move {
  138. listen_and_serve(server_config, rpc_interface, executor).await
  139. }).detach();
  140. loop {
  141. let (stream, peer_addr) = match listener.accept().await {
  142. Ok((s, a)) => (s, a),
  143. Err(err) => {
  144. error!("Error listening for connections: {}", err);
  145. return Err(Error::ServiceStopped)
  146. }
  147. };
  148. info!("Accepted client: {}", peer_addr);
  149. let p2p2 = p2p.clone();
  150. let ex2 = executor.clone();
  151. executor.spawn(process(recvr.clone(), stream, peer_addr, p2p2, ex2)).detach();
  152. }
  153. }
  154. struct JsonRpcInterface {
  155. }
  156. #[async_trait]
  157. impl RequestHandler for JsonRpcInterface {
  158. async fn handle_request(&self, req: JsonRequest, _executor: Arc<Executor<'_>>) -> JsonResult {
  159. if req.params.as_array().is_none() {
  160. return JsonResult::Err(jsonerr(InvalidParams, None, req.id))
  161. }
  162. debug!(target: "RPC", "--> {}", serde_json::to_string(&req).unwrap());
  163. match req.method.as_str() {
  164. Some("say_hello") => return self.say_hello(req.id, req.params).await,
  165. Some(_) | None => return JsonResult::Err(jsonerr(MethodNotFound, None, req.id)),
  166. }
  167. }
  168. }
  169. impl JsonRpcInterface {
  170. // --> {"method": "say_hello", "params": []}
  171. // <-- {"result": "hello world"}
  172. async fn say_hello(&self, id: Value, _params: Value) -> JsonResult {
  173. JsonResult::Resp(jsonresp(json!("hello world"), id))
  174. }
  175. }
  176. fn main() -> Result<()> {
  177. TermLogger::init(
  178. LevelFilter::Debug,
  179. simplelog::Config::default(),
  180. TerminalMode::Mixed,
  181. ColorChoice::Auto,
  182. )?;
  183. let options = ProgramOptions::load()?;
  184. let ex = Arc::new(Executor::new());
  185. smol::block_on(ex.run(start(ex.clone(), options)))
  186. }