protocol_privmsg.rs 3.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105
  1. use async_std::sync::Arc;
  2. use async_executor::Executor;
  3. use async_trait::async_trait;
  4. use log::debug;
  5. use ringbuffer::{RingBufferExt, RingBufferWrite};
  6. use darkfi::{net, Result};
  7. use crate::privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds};
  8. pub struct ProtocolPrivmsg {
  9. jobsman: net::ProtocolJobsManagerPtr,
  10. notify_queue_sender: async_channel::Sender<Privmsg>,
  11. msg_sub: net::MessageSubscription<Privmsg>,
  12. p2p: net::P2pPtr,
  13. msg_ids: SeenMsgIds,
  14. msgs: PrivmsgsBuffer,
  15. channel: net::ChannelPtr,
  16. }
  17. impl ProtocolPrivmsg {
  18. pub async fn init(
  19. channel: net::ChannelPtr,
  20. notify_queue_sender: async_channel::Sender<Privmsg>,
  21. p2p: net::P2pPtr,
  22. msg_ids: SeenMsgIds,
  23. msgs: PrivmsgsBuffer,
  24. ) -> net::ProtocolBasePtr {
  25. let message_subsytem = channel.get_message_subsystem();
  26. message_subsytem.add_dispatch::<Privmsg>().await;
  27. let msg_sub =
  28. channel.subscribe_msg::<Privmsg>().await.expect("Missing Privmsg dispatcher!");
  29. Arc::new(Self {
  30. notify_queue_sender,
  31. msg_sub,
  32. jobsman: net::ProtocolJobsManager::new("ProtocolPrivmsg", channel.clone()),
  33. p2p,
  34. msg_ids,
  35. msgs,
  36. channel,
  37. })
  38. }
  39. async fn handle_receive_msg(self: Arc<Self>) -> Result<()> {
  40. debug!(target: "ircd", "ProtocolPrivmsg::handle_receive_msg() [START]");
  41. let exclude_list = vec![self.channel.address().clone()];
  42. // once a channel get started
  43. let msgs_buffer = self.msgs.lock().await;
  44. let msgs = msgs_buffer.to_vec();
  45. drop(msgs_buffer);
  46. for m in msgs {
  47. self.channel.send(m.clone()).await?;
  48. }
  49. loop {
  50. let msg = self.msg_sub.receive().await?;
  51. if msg.nickname.len() > 10 {
  52. continue
  53. }
  54. if self.msg_ids.lock().await.contains(&msg.id) {
  55. continue
  56. }
  57. self.msg_ids.lock().await.push(msg.id);
  58. let msg = (*msg).clone();
  59. // add the msg to the buffer
  60. self.msgs.lock().await.push(msg.clone());
  61. self.notify_queue_sender.send(msg.clone()).await?;
  62. self.p2p.broadcast_with_exclude(msg.clone(), &exclude_list).await?;
  63. }
  64. }
  65. }
  66. #[async_trait]
  67. impl net::ProtocolBase for ProtocolPrivmsg {
  68. /// Starts ping-pong keep-alive messages exchange. Runs ping-pong in the
  69. /// protocol task manager, then queues the reply. Sends out a ping and
  70. /// waits for pong reply. Waits for ping and replies with a pong.
  71. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  72. debug!(target: "ircd", "ProtocolPrivmsg::start() [START]");
  73. self.jobsman.clone().start(executor.clone());
  74. self.jobsman.clone().spawn(self.clone().handle_receive_msg(), executor.clone()).await;
  75. debug!(target: "ircd", "ProtocolPrivmsg::start() [END]");
  76. Ok(())
  77. }
  78. fn name(&self) -> &'static str {
  79. "ProtocolPrivmsg"
  80. }
  81. }
  82. impl net::Message for Privmsg {
  83. fn name() -> &'static str {
  84. "privmsg"
  85. }
  86. }