client.rs 18 KB

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