protocol_privmsg.rs 2.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273
  1. use async_executor::Executor;
  2. use async_std::sync::Mutex;
  3. use darkfi::{net, Result};
  4. use log::debug;
  5. use std::{collections::HashSet, sync::Arc};
  6. use crate::privmsg::{PrivMsg, PrivMsgId, SeenPrivMsgIdsPtr};
  7. pub struct ProtocolPrivMsg {
  8. notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
  9. privmsg_sub: net::MessageSubscription<PrivMsg>,
  10. jobsman: net::ProtocolJobsManagerPtr,
  11. seen_privmsg_ids: SeenPrivMsgIdsPtr,
  12. p2p: net::P2pPtr,
  13. }
  14. impl ProtocolPrivMsg {
  15. pub async fn new(
  16. channel: net::ChannelPtr,
  17. notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
  18. seen_privmsg_ids: SeenPrivMsgIdsPtr,
  19. p2p: net::P2pPtr,
  20. ) -> Arc<Self> {
  21. let message_subsytem = channel.get_message_subsystem();
  22. message_subsytem.add_dispatch::<PrivMsg>().await;
  23. debug!("ADDED DISPATCH");
  24. let privmsg_sub =
  25. channel.subscribe_msg::<PrivMsg>().await.expect("Missing PrivMsg dispatcher!");
  26. Arc::new(Self {
  27. notify_queue_sender,
  28. privmsg_sub,
  29. jobsman: net::ProtocolJobsManager::new("PrivMsgProtocol", channel),
  30. seen_privmsg_ids,
  31. p2p,
  32. })
  33. }
  34. pub async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) {
  35. debug!(target: "ircd", "ProtocolPrivMsg::start() [START]");
  36. self.jobsman.clone().start(executor.clone());
  37. self.jobsman.clone().spawn(self.clone().handle_receive_privmsg(), executor.clone()).await;
  38. debug!(target: "ircd", "ProtocolPrivMsg::start() [END]");
  39. }
  40. async fn handle_receive_privmsg(self: Arc<Self>) -> Result<()> {
  41. debug!(target: "ircd", "ProtocolAddress::handle_receive_privmsg() [START]");
  42. loop {
  43. let privmsg = self.privmsg_sub.receive().await?;
  44. debug!(
  45. target: "ircd",
  46. "ProtocolPrivMsg::handle_receive_privmsg() received {:?}",
  47. privmsg
  48. );
  49. // Do we already have this message?
  50. if self.seen_privmsg_ids.is_seen(privmsg.id).await {
  51. continue
  52. }
  53. self.seen_privmsg_ids.add_seen(privmsg.id).await;
  54. // If not then broadcast to everybody else
  55. let privmsg_copy = (*privmsg).clone();
  56. self.p2p.broadcast(privmsg_copy).await?;
  57. self.notify_queue_sender.send(privmsg).await.expect("notify_queue_sender send failed!");
  58. }
  59. }
  60. }