protocol_privmsg.rs 2.3 KB

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