main.rs 7.1 KB

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