client.rs 18 KB

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