server.rs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523
  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 rand::{rngs::OsRng, RngCore};
  6. use ringbuffer::{RingBufferExt, RingBufferWrite};
  7. use darkfi::{net::P2pPtr, system::SubscriberPtr, Error, Result};
  8. use crate::{
  9. crypto::{encrypt_message, try_decrypt_message},
  10. privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
  11. ChannelInfo, MAXIMUM_LENGTH_OF_MESSAGE, MAXIMUM_LENGTH_OF_NICKNAME,
  12. };
  13. const RPL_NOTOPIC: u32 = 331;
  14. const RPL_TOPIC: u32 = 332;
  15. const RPL_NAMEREPLY: u32 = 353;
  16. const RPL_ENDOFNAMES: u32 = 366;
  17. pub struct IrcServerConnection<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> {
  18. // server stream
  19. write_stream: WriteHalf<C>,
  20. peer_address: SocketAddr,
  21. // msg ids
  22. seen_msg_ids: SeenMsgIds,
  23. privmsgs_buffer: PrivmsgsBuffer,
  24. // user & channels
  25. is_nick_init: bool,
  26. is_user_init: bool,
  27. is_registered: bool,
  28. is_cap_end: bool,
  29. nickname: String,
  30. auto_channels: Vec<String>,
  31. pub configured_chans: FxHashMap<String, ChannelInfo>,
  32. pub configured_contacts: FxHashMap<String, crypto_box::Box>,
  33. capabilities: FxHashMap<String, bool>,
  34. // p2p
  35. p2p: P2pPtr,
  36. senders: SubscriberPtr<Privmsg>,
  37. subscriber_id: u64,
  38. }
  39. impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcServerConnection<C> {
  40. #[allow(clippy::too_many_arguments)]
  41. pub fn new(
  42. write_stream: WriteHalf<C>,
  43. peer_address: SocketAddr,
  44. seen_msg_ids: SeenMsgIds,
  45. privmsgs_buffer: PrivmsgsBuffer,
  46. auto_channels: Vec<String>,
  47. configured_chans: FxHashMap<String, ChannelInfo>,
  48. configured_contacts: FxHashMap<String, crypto_box::Box>,
  49. p2p: P2pPtr,
  50. senders: SubscriberPtr<Privmsg>,
  51. subscriber_id: u64,
  52. ) -> Self {
  53. let mut capabilities = FxHashMap::default();
  54. capabilities.insert("no-history".to_string(), false);
  55. Self {
  56. write_stream,
  57. peer_address,
  58. seen_msg_ids,
  59. privmsgs_buffer,
  60. is_nick_init: false,
  61. is_user_init: false,
  62. is_registered: false,
  63. is_cap_end: true,
  64. nickname: "anon".to_string(),
  65. auto_channels,
  66. configured_chans,
  67. configured_contacts,
  68. capabilities,
  69. p2p,
  70. senders,
  71. subscriber_id,
  72. }
  73. }
  74. async fn update(&mut self, line: String) -> Result<()> {
  75. if line.len() > MAXIMUM_LENGTH_OF_MESSAGE {
  76. return Err(Error::MalformedPacket)
  77. }
  78. let mut tokens = line.split_ascii_whitespace();
  79. // Commands can begin with :garbage but we will reject clients doing
  80. // that for now to keep the protocol simple and focused.
  81. let command = tokens.next().ok_or(Error::MalformedPacket)?;
  82. info!("IRC server received command: {}", command.to_uppercase());
  83. match command.to_uppercase().as_str() {
  84. "USER" => {
  85. // We can stuff any extra things like public keys in here.
  86. // Ignore it for now.
  87. self.is_user_init = true;
  88. }
  89. "NAMES" => {
  90. let channels = tokens.next().ok_or(Error::MalformedPacket)?;
  91. for chan in channels.split(',') {
  92. if !chan.starts_with('#') {
  93. warn!("{} is not a valid name for channel", chan);
  94. continue
  95. }
  96. self.on_receive_names(chan).await?;
  97. }
  98. }
  99. "NICK" => {
  100. let nickname = tokens.next().ok_or(Error::MalformedPacket)?;
  101. if nickname.len() > MAXIMUM_LENGTH_OF_NICKNAME {
  102. return Ok(())
  103. }
  104. self.is_nick_init = true;
  105. let old_nick = std::mem::replace(&mut self.nickname, nickname.to_string());
  106. let nick_reply = format!(":{}!anon@dark.fi NICK {}\r\n", old_nick, self.nickname);
  107. self.reply(&nick_reply).await?;
  108. }
  109. "JOIN" => {
  110. let channels = tokens.next().ok_or(Error::MalformedPacket)?;
  111. for chan in channels.split(',') {
  112. if !chan.starts_with('#') {
  113. warn!("{} is not a valid name for channel", chan);
  114. continue
  115. }
  116. self.on_join(chan).await?;
  117. }
  118. }
  119. "PART" => {
  120. let channels = tokens.next().ok_or(Error::MalformedPacket)?;
  121. for chan in channels.split(',') {
  122. let part_reply = format!(":{}!anon@dark.fi PART {}\r\n", self.nickname, chan);
  123. self.reply(&part_reply).await?;
  124. if self.configured_chans.contains_key(chan) {
  125. let chan_info = self.configured_chans.get_mut(chan).unwrap();
  126. chan_info.joined = false;
  127. }
  128. }
  129. }
  130. "TOPIC" => {
  131. let channel = tokens.next().ok_or(Error::MalformedPacket)?;
  132. if let Some(substr_idx) = line.find(':') {
  133. // Client is setting the topic
  134. if substr_idx >= line.len() {
  135. return Err(Error::MalformedPacket)
  136. }
  137. let topic = &line[substr_idx + 1..];
  138. let chan_info = self.configured_chans.get_mut(channel).unwrap();
  139. chan_info.topic = Some(topic.to_string());
  140. let topic_reply =
  141. format!(":{}!anon@dark.fi TOPIC {} :{}\r\n", self.nickname, channel, topic);
  142. self.reply(&topic_reply).await?;
  143. } else {
  144. // Client is asking or the topic
  145. let chan_info = self.configured_chans.get(channel).unwrap();
  146. let topic_reply = if let Some(topic) = &chan_info.topic {
  147. format!("{} {} {} :{}\r\n", RPL_TOPIC, self.nickname, channel, topic)
  148. } else {
  149. const TOPIC: &str = "No topic is set";
  150. format!("{} {} {} :{}\r\n", RPL_NOTOPIC, self.nickname, channel, TOPIC)
  151. };
  152. self.reply(&topic_reply).await?;
  153. }
  154. }
  155. "PING" => {
  156. let pong = tokens.next().ok_or(Error::MalformedPacket)?;
  157. let pong = format!("PONG {}\r\n", pong);
  158. self.reply(&pong).await?;
  159. }
  160. "PRIVMSG" => {
  161. let target = tokens.next().ok_or(Error::MalformedPacket)?;
  162. let substr_idx = line.find(':').ok_or(Error::MalformedPacket)?;
  163. if substr_idx >= line.len() {
  164. return Err(Error::MalformedPacket)
  165. }
  166. let mut message = line[substr_idx + 1..].to_string();
  167. info!("(Plain) PRIVMSG {} :{}", target, message);
  168. if target.starts_with('#') {
  169. if !self.configured_chans.contains_key(target) {
  170. return Ok(())
  171. }
  172. let channel_info = self.configured_chans.get(target).unwrap();
  173. if !channel_info.joined {
  174. return Ok(())
  175. }
  176. message = if let Some(salt_box) = &channel_info.salt_box {
  177. let encrypted = encrypt_message(salt_box, &message);
  178. info!("(Encrypted) PRIVMSG {} :{}", target, encrypted);
  179. encrypted
  180. } else {
  181. message.to_string()
  182. };
  183. } else {
  184. // If we have a configured secret for this nick, we encrypt the message.
  185. if let Some(salt_box) = self.configured_contacts.get(target) {
  186. message = encrypt_message(salt_box, &message);
  187. info!("(Encrypted) PRIVMSG {} :{}", target, message);
  188. }
  189. }
  190. self.on_receive_privmsg(&message, target).await?;
  191. }
  192. "CAP" => {
  193. self.is_cap_end = false;
  194. let subcommand = tokens.next().ok_or(Error::MalformedPacket)?.to_uppercase();
  195. let capabilities_keys: Vec<String> = self.capabilities.keys().cloned().collect();
  196. if subcommand == "LS" {
  197. let cap_ls_reply = format!(
  198. ":{}!anon@dark.fi CAP * LS :{}\r\n",
  199. self.nickname,
  200. capabilities_keys.join(" ")
  201. );
  202. self.reply(&cap_ls_reply).await?;
  203. }
  204. if subcommand == "REQ" {
  205. let substr_idx = line.find(':').ok_or(Error::MalformedPacket)?;
  206. if substr_idx >= line.len() {
  207. return Err(Error::MalformedPacket)
  208. }
  209. let cap: Vec<&str> = line[substr_idx + 1..].split(' ').collect();
  210. let mut ack_list = vec![];
  211. let mut nak_list = vec![];
  212. for c in cap {
  213. if self.capabilities.contains_key(c) {
  214. self.capabilities.insert(c.to_string(), true);
  215. ack_list.push(c);
  216. } else {
  217. nak_list.push(c);
  218. }
  219. }
  220. let cap_ack_reply = format!(
  221. ":{}!anon@dark.fi CAP * ACK :{}\r\n",
  222. self.nickname,
  223. ack_list.join(" ")
  224. );
  225. let cap_nak_reply = format!(
  226. ":{}!anon@dark.fi CAP * NAK :{}\r\n",
  227. self.nickname,
  228. nak_list.join(" ")
  229. );
  230. self.reply(&cap_ack_reply).await?;
  231. self.reply(&cap_nak_reply).await?;
  232. }
  233. if subcommand == "LIST" {
  234. let enabled_capabilities: Vec<String> = self
  235. .capabilities
  236. .clone()
  237. .into_iter()
  238. .filter(|(_, v)| *v)
  239. .map(|(k, _)| k)
  240. .collect();
  241. let cap_list_reply = format!(
  242. ":{}!anon@dark.fi CAP * LIST :{}\r\n",
  243. self.nickname,
  244. enabled_capabilities.join(" ")
  245. );
  246. self.reply(&cap_list_reply).await?;
  247. }
  248. if subcommand == "END" {
  249. self.is_cap_end = true;
  250. }
  251. }
  252. "QUIT" => {
  253. // Close the connection
  254. return Err(Error::NetworkServiceStopped)
  255. }
  256. _ => {
  257. warn!("Unimplemented `{}` command", command);
  258. }
  259. }
  260. // on registration
  261. if !self.is_registered && self.is_cap_end && self.is_nick_init && self.is_user_init {
  262. debug!("Initializing peer connection");
  263. let register_reply = format!(":darkfi 001 {} :Let there be dark\r\n", self.nickname);
  264. self.reply(&register_reply).await?;
  265. self.is_registered = true;
  266. for chan in self.auto_channels.clone() {
  267. self.on_join(&chan).await?;
  268. }
  269. // Send dm messages in buffer
  270. if *self.capabilities.get("no-history").unwrap() {
  271. return Ok(())
  272. }
  273. for msg in self.privmsgs_buffer.lock().await.to_vec() {
  274. if msg.target == self.nickname ||
  275. (msg.nickname == self.nickname && !msg.target.starts_with('#'))
  276. {
  277. self.senders.notify_by_id(msg, self.subscriber_id).await;
  278. }
  279. }
  280. }
  281. Ok(())
  282. }
  283. async fn reply(&mut self, message: &str) -> Result<()> {
  284. self.write_stream.write_all(message.as_bytes()).await?;
  285. debug!("Sent {}", message);
  286. Ok(())
  287. }
  288. async fn on_receive_names(&mut self, chan: &str) -> Result<()> {
  289. if self.configured_chans.contains_key(chan) {
  290. let chan_info = self.configured_chans.get(chan).unwrap();
  291. if chan_info.names.is_empty() {
  292. return Ok(())
  293. }
  294. let names_reply = format!(
  295. ":{}!anon@dark.fi {} = {} : {}\r\n",
  296. self.nickname,
  297. RPL_NAMEREPLY,
  298. chan,
  299. chan_info.names.join(" ")
  300. );
  301. self.reply(&names_reply).await?;
  302. let end_of_names = format!(
  303. ":DarkFi {:03} {} {} :End of NAMES list\r\n",
  304. RPL_ENDOFNAMES, self.nickname, chan
  305. );
  306. self.reply(&end_of_names).await?;
  307. }
  308. Ok(())
  309. }
  310. async fn on_receive_privmsg(&mut self, message: &str, target: &str) -> Result<()> {
  311. let random_id = OsRng.next_u64();
  312. let protocol_msg = Privmsg {
  313. id: random_id,
  314. nickname: self.nickname.clone(),
  315. target: target.to_string(),
  316. message: message.to_string(),
  317. };
  318. {
  319. (*self.seen_msg_ids.lock().await).push(random_id);
  320. (*self.privmsgs_buffer.lock().await).push(protocol_msg.clone())
  321. }
  322. self.senders.notify_with_exclude(protocol_msg.clone(), &[self.subscriber_id]).await;
  323. debug!(target: "ircd", "PRIVMSG to be sent: {:?}", protocol_msg);
  324. self.p2p.broadcast(protocol_msg).await?;
  325. Ok(())
  326. }
  327. async fn on_join(&mut self, chan: &str) -> Result<()> {
  328. if !self.configured_chans.contains_key(chan) {
  329. let mut chan_info = ChannelInfo::new()?;
  330. chan_info.topic = Some("n/a".to_string());
  331. self.configured_chans.insert(chan.to_string(), chan_info);
  332. }
  333. let chan_info = self.configured_chans.get_mut(chan).unwrap();
  334. if chan_info.joined {
  335. return Ok(())
  336. }
  337. chan_info.joined = true;
  338. let topic =
  339. if let Some(topic) = chan_info.topic.clone() { topic } else { "n/a".to_string() };
  340. chan_info.topic = Some(topic.to_string());
  341. {
  342. let j = format!(":{}!anon@dark.fi JOIN {}\r\n", self.nickname, chan);
  343. let t = format!(":DarkFi TOPIC {} :{}\r\n", chan, topic);
  344. self.reply(&j).await?;
  345. self.reply(&t).await?;
  346. }
  347. // Send messages in buffer
  348. if !self.capabilities.get("no-history").unwrap() {
  349. for msg in self.privmsgs_buffer.lock().await.to_vec() {
  350. if msg.target == chan {
  351. self.senders.notify_by_id(msg, self.subscriber_id).await;
  352. }
  353. }
  354. }
  355. self.on_receive_names(chan).await?;
  356. Ok(())
  357. }
  358. pub async fn process_msg_from_p2p(&mut self, msg: &Privmsg) -> Result<()> {
  359. info!("Received msg from P2p network: {:?}", msg);
  360. let mut msg = msg.clone();
  361. if msg.target.starts_with('#') {
  362. // Try to potentially decrypt the incoming message.
  363. if !self.configured_chans.contains_key(&msg.target) {
  364. return Ok(())
  365. }
  366. let chan_info = self.configured_chans.get_mut(&msg.target).unwrap();
  367. if !chan_info.joined {
  368. return Ok(())
  369. }
  370. let salt_box = chan_info.salt_box.clone();
  371. if salt_box.is_some() {
  372. let decrypted_msg = try_decrypt_message(&salt_box.unwrap(), &msg.message);
  373. if decrypted_msg.is_none() {
  374. return Ok(())
  375. }
  376. msg.message = decrypted_msg.unwrap();
  377. info!("Decrypted received message: {:?}", msg);
  378. }
  379. // add the nickname to the channel's names
  380. if !chan_info.names.contains(&msg.nickname) {
  381. chan_info.names.push(msg.nickname.clone());
  382. }
  383. self.reply(&msg.to_irc_msg()).await?;
  384. return Ok(())
  385. } else {
  386. if self.is_cap_end &&
  387. self.is_nick_init &&
  388. (self.nickname == msg.target || self.nickname == msg.nickname)
  389. {
  390. if self.configured_contacts.contains_key(&msg.target) {
  391. let salt_box = self.configured_contacts.get(&msg.target).unwrap();
  392. if let Some(decrypted) = try_decrypt_message(&salt_box, &msg.message) {
  393. msg.message = decrypted;
  394. info!("Decrypted received message: {:?}", msg);
  395. }
  396. }
  397. self.reply(&msg.to_irc_msg()).await?;
  398. }
  399. }
  400. Ok(())
  401. }
  402. pub async fn process_line_from_client(
  403. &mut self,
  404. err: std::result::Result<usize, std::io::Error>,
  405. line: String,
  406. ) -> Result<()> {
  407. if let Err(e) = err {
  408. warn!("Read line error {}: {}", self.peer_address, e);
  409. return Err(Error::ChannelStopped)
  410. }
  411. info!("Received msg from IRC client: {:?}", line);
  412. let irc_msg = self.clean_input_line(line)?;
  413. if let Err(e) = self.update(irc_msg).await {
  414. warn!("Connection error: {} for {}", e, self.peer_address);
  415. return Err(Error::ChannelStopped)
  416. }
  417. Ok(())
  418. }
  419. fn clean_input_line(&self, mut line: String) -> Result<String> {
  420. if line.is_empty() {
  421. warn!("Received empty line from {}. ", self.peer_address);
  422. warn!("Closing connection.");
  423. return Err(Error::ChannelStopped)
  424. }
  425. if &line[(line.len() - 2)..] == "\r\n" {
  426. // Remove CRLF
  427. line.pop();
  428. line.pop();
  429. } else if &line[(line.len() - 1)..] == "\n" {
  430. line.pop();
  431. } else {
  432. warn!("Closing connection.");
  433. return Err(Error::ChannelStopped)
  434. }
  435. if line == "\n" {
  436. warn!("Closing connection.");
  437. return Err(Error::ChannelStopped)
  438. }
  439. Ok(line.clone())
  440. }
  441. }