mod.rs 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248
  1. use std::net::SocketAddr;
  2. use futures::{io::WriteHalf, AsyncRead, AsyncWrite, AsyncWriteExt};
  3. use fxhash::FxHashMap;
  4. use log::{debug, info, warn};
  5. use darkfi::{net::P2pPtr, system::SubscriberPtr, Error, Result};
  6. use crate::{
  7. buffers::{ArcPrivmsgsBuffer, SeenMsgIds},
  8. crypto::{decrypt_privmsg, decrypt_target},
  9. ChannelInfo, Privmsg, MAXIMUM_LENGTH_OF_MESSAGE,
  10. };
  11. mod command;
  12. pub struct IrcServerConnection<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> {
  13. // server stream
  14. write_stream: WriteHalf<C>,
  15. peer_address: SocketAddr,
  16. // msg ids
  17. seen_msg_ids: SeenMsgIds,
  18. privmsgs_buffer: ArcPrivmsgsBuffer,
  19. // user & channels
  20. is_nick_init: bool,
  21. is_user_init: bool,
  22. is_registered: bool,
  23. is_cap_end: bool,
  24. is_pass_init: bool,
  25. nickname: String,
  26. auto_channels: Vec<String>,
  27. pub configured_chans: FxHashMap<String, ChannelInfo>,
  28. pub configured_contacts: FxHashMap<String, crypto_box::SalsaBox>,
  29. capabilities: FxHashMap<String, bool>,
  30. // p2p
  31. p2p: P2pPtr,
  32. senders: SubscriberPtr<Privmsg>,
  33. subscriber_id: u64,
  34. password: String,
  35. }
  36. impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcServerConnection<C> {
  37. #[allow(clippy::too_many_arguments)]
  38. pub fn new(
  39. write_stream: WriteHalf<C>,
  40. peer_address: SocketAddr,
  41. seen_msg_ids: SeenMsgIds,
  42. privmsgs_buffer: ArcPrivmsgsBuffer,
  43. auto_channels: Vec<String>,
  44. password: String,
  45. configured_chans: FxHashMap<String, ChannelInfo>,
  46. configured_contacts: FxHashMap<String, crypto_box::SalsaBox>,
  47. p2p: P2pPtr,
  48. senders: SubscriberPtr<Privmsg>,
  49. subscriber_id: u64,
  50. ) -> Self {
  51. let mut capabilities = FxHashMap::default();
  52. capabilities.insert("no-history".to_string(), false);
  53. Self {
  54. write_stream,
  55. peer_address,
  56. seen_msg_ids,
  57. privmsgs_buffer,
  58. is_nick_init: false,
  59. is_user_init: false,
  60. is_registered: false,
  61. is_cap_end: true,
  62. is_pass_init: false,
  63. nickname: "anon".to_string(),
  64. auto_channels,
  65. password,
  66. configured_chans,
  67. configured_contacts,
  68. capabilities,
  69. p2p,
  70. senders,
  71. subscriber_id,
  72. }
  73. }
  74. pub async fn process_msg_from_p2p(&mut self, msg: &Privmsg) -> Result<()> {
  75. info!("Received msg from P2p network: {:?}", msg);
  76. let mut msg = msg.clone();
  77. decrypt_target(&mut msg, self.configured_chans.clone(), self.configured_contacts.clone());
  78. if msg.target.starts_with('#') {
  79. // Try to potentially decrypt the incoming message.
  80. if !self.configured_chans.contains_key(&msg.target) {
  81. return Ok(())
  82. }
  83. let chan_info = self.configured_chans.get_mut(&msg.target).unwrap();
  84. if !chan_info.joined {
  85. return Ok(())
  86. }
  87. if let Some(salt_box) = &chan_info.salt_box {
  88. decrypt_privmsg(salt_box, &mut msg);
  89. info!("Decrypted received message: {:?}", msg);
  90. }
  91. // add the nickname to the channel's names
  92. if !chan_info.names.contains(&msg.nickname) {
  93. chan_info.names.push(msg.nickname.clone());
  94. }
  95. self.reply(&msg.to_string()).await?;
  96. return Ok(())
  97. } else if self.is_cap_end && self.is_nick_init && self.nickname == msg.target {
  98. if self.configured_contacts.contains_key(&msg.target) {
  99. let salt_box = self.configured_contacts.get(&msg.target).unwrap();
  100. decrypt_privmsg(salt_box, &mut msg);
  101. info!("Decrypted received message: {:?}", msg);
  102. }
  103. self.reply(&msg.to_string()).await?;
  104. }
  105. Ok(())
  106. }
  107. pub async fn process_line_from_client(
  108. &mut self,
  109. err: std::result::Result<usize, std::io::Error>,
  110. line: String,
  111. ) -> Result<()> {
  112. if let Err(e) = err {
  113. warn!("Read line error {}: {}", self.peer_address, e);
  114. return Err(Error::ChannelStopped)
  115. }
  116. info!("Received msg from IRC client: {:?}", line);
  117. let irc_msg = clean_input_line(line, &self.peer_address)?;
  118. if let Err(e) = self.update(irc_msg).await {
  119. warn!("Connection error: {} for {}", e, self.peer_address);
  120. return Err(Error::ChannelStopped)
  121. }
  122. Ok(())
  123. }
  124. async fn update(&mut self, line: String) -> Result<()> {
  125. if line.len() > MAXIMUM_LENGTH_OF_MESSAGE {
  126. return Err(Error::MalformedPacket)
  127. }
  128. if self.password.is_empty() {
  129. self.is_pass_init = true
  130. }
  131. let (command, value) = parse_line(&line)?;
  132. let (command, value) = (command.as_str(), value.as_str());
  133. info!("IRC server received command: {}", command);
  134. match command {
  135. "PASS" => self.on_receive_pass(value).await?,
  136. "USER" => self.on_receive_user().await?,
  137. "NAMES" => self.on_receive_names(value.split(',').map(String::from).collect()).await?,
  138. "NICK" => self.on_receive_nick(value).await?,
  139. "JOIN" => self.on_receive_join(value.split(',').map(String::from).collect()).await?,
  140. "PART" => self.on_receive_part(value.split(',').map(String::from).collect()).await?,
  141. "TOPIC" => self.on_receive_topic(&line, value).await?,
  142. "PING" => self.on_ping(value).await?,
  143. "PRIVMSG" => self.on_receive_privmsg(&line, value).await?,
  144. "CAP" => self.on_receive_cap(&line, &value.to_uppercase()).await?,
  145. "QUIT" => self.on_quit()?,
  146. _ => warn!("Unimplemented `{}` command", command),
  147. }
  148. self.registre().await?;
  149. Ok(())
  150. }
  151. async fn registre(&mut self) -> Result<()> {
  152. if !self.is_registered && self.is_cap_end && self.is_nick_init && self.is_user_init {
  153. debug!("Initializing peer connection");
  154. let register_reply = format!(":darkfi 001 {} :Let there be dark\r\n", self.nickname);
  155. self.reply(&register_reply).await?;
  156. self.is_registered = true;
  157. self.on_receive_join(self.auto_channels.clone()).await?;
  158. if *self.capabilities.get("no-history").unwrap() {
  159. return Ok(())
  160. }
  161. // Send dm messages in buffer
  162. let mut privmsgs_buffer = self.privmsgs_buffer.lock().await;
  163. privmsgs_buffer.update();
  164. for msg in privmsgs_buffer.iter() {
  165. let is_dm = msg.target == self.nickname ||
  166. (msg.nickname == self.nickname && !msg.target.starts_with('#'));
  167. if is_dm {
  168. self.senders.notify_by_id(msg.clone(), self.subscriber_id).await;
  169. }
  170. }
  171. drop(privmsgs_buffer);
  172. }
  173. Ok(())
  174. }
  175. async fn reply(&mut self, message: &str) -> Result<()> {
  176. self.write_stream.write_all(message.as_bytes()).await?;
  177. debug!("Sent {}", message);
  178. Ok(())
  179. }
  180. }
  181. //
  182. // Helper functions
  183. //
  184. fn clean_input_line(mut line: String, peer_address: &SocketAddr) -> Result<String> {
  185. if line.is_empty() {
  186. warn!("Received empty line from {}. ", peer_address);
  187. warn!("Closing connection.");
  188. return Err(Error::ChannelStopped)
  189. }
  190. if line == "\n" || line == "\r\n" {
  191. warn!("Closing connection.");
  192. return Err(Error::ChannelStopped)
  193. }
  194. if &line[(line.len() - 2)..] == "\r\n" {
  195. // Remove CRLF
  196. line.pop();
  197. line.pop();
  198. } else if &line[(line.len() - 1)..] == "\n" {
  199. line.pop();
  200. } else {
  201. warn!("Closing connection.");
  202. return Err(Error::ChannelStopped)
  203. }
  204. Ok(line.clone())
  205. }
  206. fn parse_line(line: &str) -> Result<(String, String)> {
  207. let mut tokens = line.split_ascii_whitespace();
  208. // Commands can begin with :garbage but we will reject clients doing
  209. // that for now to keep the protocol simple and focused.
  210. let command = tokens.next().ok_or(Error::MalformedPacket)?.to_uppercase();
  211. let value = tokens.next().ok_or(Error::MalformedPacket)?;
  212. Ok((command.to_owned(), value.to_owned()))
  213. }