protocol_privmsg.rs 2.7 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283
  1. use async_executor::Executor;
  2. use async_trait::async_trait;
  3. use darkfi::{net, Result};
  4. use log::debug;
  5. use std::sync::Arc;
  6. use crate::privmsg::{PrivMsg, 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 init(
  16. channel: net::ChannelPtr,
  17. notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
  18. seen_privmsg_ids: SeenPrivMsgIdsPtr,
  19. p2p: net::P2pPtr,
  20. ) -> net::ProtocolBasePtr {
  21. let message_subsytem = channel.get_message_subsystem();
  22. message_subsytem.add_dispatch::<PrivMsg>().await;
  23. let privmsg_sub =
  24. channel.subscribe_msg::<PrivMsg>().await.expect("Missing PrivMsg dispatcher!");
  25. Arc::new(Self {
  26. notify_queue_sender,
  27. privmsg_sub,
  28. jobsman: net::ProtocolJobsManager::new("PrivMsgProtocol", channel),
  29. seen_privmsg_ids,
  30. p2p,
  31. })
  32. }
  33. async fn handle_receive_privmsg(self: Arc<Self>) -> Result<()> {
  34. debug!(target: "ircd", "ProtocolPrivMsg::handle_receive_privmsg() [START]");
  35. loop {
  36. let privmsg = self.privmsg_sub.receive().await?;
  37. debug!(
  38. target: "ircd",
  39. "ProtocolPrivMsg::handle_receive_privmsg() received {:?}",
  40. privmsg
  41. );
  42. // Do we already have this message?
  43. if self.seen_privmsg_ids.is_seen(privmsg.id).await {
  44. continue
  45. }
  46. self.seen_privmsg_ids.add_seen(privmsg.id).await;
  47. // If not then broadcast to everybody else
  48. let privmsg_copy = (*privmsg).clone();
  49. self.p2p.broadcast(privmsg_copy).await?;
  50. self.notify_queue_sender.send(privmsg).await.expect("notify_queue_sender send failed!");
  51. }
  52. }
  53. }
  54. #[async_trait]
  55. impl net::ProtocolBase for ProtocolPrivMsg {
  56. /// Starts ping-pong keep-alive messages exchange. Runs ping-pong in the
  57. /// protocol task manager, then queues the reply. Sends out a ping and
  58. /// waits for pong reply. Waits for ping and replies with a pong.
  59. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  60. debug!(target: "ircd", "ProtocolPrivMsg::start() [START]");
  61. self.jobsman.clone().start(executor.clone());
  62. self.jobsman.clone().spawn(self.clone().handle_receive_privmsg(), executor.clone()).await;
  63. debug!(target: "ircd", "ProtocolPrivMsg::start() [END]");
  64. Ok(())
  65. }
  66. fn name(&self) -> &'static str {
  67. "ProtocolPrivMsg"
  68. }
  69. }