protocol_privmsg.rs 3.3 KB

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