protocol_slab.rs 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187
  1. use std::sync::Arc;
  2. use log::*;
  3. use smol::Executor;
  4. use crate::error::Result as NetResult;
  5. use crate::serial::deserialize;
  6. use crate::net::{message_subscriber::MessageSubscription, protocols::ProtocolJobsManager,
  7. protocols::ProtocolJobsManagerPtr ,ChannelPtr};
  8. use crate::darkpulse::{aes_decrypt, messages, ControlCommand, ControlMessage, SlabsManagerSafe, CiphertextHash};
  9. pub struct ProtocolSlab {
  10. channel: ChannelPtr,
  11. slabman: SlabsManagerSafe,
  12. sync_sub: MessageSubscription<messages::SyncMessage>,
  13. inv_sub: MessageSubscription<messages::InvMessage>,
  14. get_slabs_sub: MessageSubscription<messages::GetSlabsMessage>,
  15. slab_sub: MessageSubscription<messages::SlabMessage>,
  16. jobsman: ProtocolJobsManagerPtr,
  17. }
  18. impl ProtocolSlab {
  19. pub async fn new(slabman: SlabsManagerSafe, channel: ChannelPtr) -> Arc<Self> {
  20. let sync_sub = channel
  21. .clone()
  22. .subscribe_msg::<messages::SyncMessage>()
  23. .await
  24. .expect("Missing sync dispatcher!");
  25. let inv_sub = channel
  26. .clone()
  27. .subscribe_msg::<messages::InvMessage>()
  28. .await
  29. .expect("Missing inv dispatcher!");
  30. let get_slabs_sub = channel
  31. .clone()
  32. .subscribe_msg::<messages::GetSlabsMessage>()
  33. .await
  34. .expect("Missing getslabs dispatcher!");
  35. let slab_sub = channel
  36. .clone()
  37. .subscribe_msg::<messages::SlabMessage>()
  38. .await
  39. .expect("Missing slab dispatcher!");
  40. Arc::new(Self {
  41. channel: channel.clone(),
  42. slabman,
  43. sync_sub,
  44. inv_sub,
  45. get_slabs_sub,
  46. slab_sub,
  47. jobsman: ProtocolJobsManager::new("ProtocolSlab", channel),
  48. })
  49. }
  50. pub async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) {
  51. debug!(target: "net", "ProtocolSlab::start() [START]");
  52. self.jobsman.clone().start(executor.clone());
  53. self.jobsman
  54. .clone()
  55. .spawn(self.clone().handle_receive_sync(), executor.clone())
  56. .await;
  57. self.jobsman
  58. .clone()
  59. .spawn(self.clone().handle_receive_inv(), executor.clone())
  60. .await;
  61. self.jobsman
  62. .clone()
  63. .spawn(self.clone().handle_receive_get_slabs(), executor.clone())
  64. .await;
  65. self.jobsman
  66. .clone()
  67. .spawn(self.clone().handle_receive_slab(), executor)
  68. .await;
  69. let _ = self.channel.send(messages::SyncMessage {}).await;
  70. debug!(target: "net", "ProtocolSlab::start() [END]");
  71. }
  72. async fn handle_receive_sync(self: Arc<Self>) -> NetResult<()> {
  73. debug!(target: "net", "ProtocolSlab::handle_receive_sync() [START]");
  74. loop {
  75. let _sync_msg = self.sync_sub.receive().await?;
  76. let slab_hashs = self.slabman.lock().await.get_slabs_hash();
  77. let inv_msg = messages::InvMessage {
  78. slabs_hash: slab_hashs.clone(),
  79. };
  80. self.channel.send(inv_msg).await?;
  81. info!("receive sync message!");
  82. }
  83. }
  84. async fn handle_receive_inv(self: Arc<Self>) -> NetResult<()> {
  85. debug!(target: "net", "ProtocolSlab::handle_receive_inv() [START]");
  86. loop {
  87. let inv_msg = self.inv_sub.receive().await?;
  88. let mut list_of_hash: Vec<CiphertextHash> = vec![];
  89. let slabs_hash = self.slabman.lock().await.get_slabs_hash();
  90. for slab in inv_msg.slabs_hash.iter() {
  91. if !slabs_hash.contains(slab) {
  92. list_of_hash.push(slab.clone());
  93. }
  94. }
  95. let getslabs_msg = messages::GetSlabsMessage {
  96. slabs_hash: list_of_hash,
  97. };
  98. self.channel.send(getslabs_msg).await?;
  99. info!("receive inv message!");
  100. }
  101. }
  102. async fn handle_receive_get_slabs(self: Arc<Self>) -> NetResult<()> {
  103. debug!(target: "net", "ProtocolSlab::handle_receive_get_slabs() [START]");
  104. loop {
  105. let get_slabs_msg = self.get_slabs_sub.receive().await?;
  106. for slab_hash in get_slabs_msg.slabs_hash.iter() {
  107. let slabman = self.slabman.lock().await;
  108. let slab = slabman.get_slab(&slab_hash);
  109. if let Some(slab) = slab {
  110. self.channel.send(slab.clone()).await?;
  111. }
  112. }
  113. info!("receive getslabs message!");
  114. }
  115. }
  116. async fn handle_receive_slab(self: Arc<Self>) -> NetResult<()> {
  117. debug!(target: "net", "ProtocolSlab::handle_receive_slab() [START]");
  118. loop {
  119. let slab_msg = self.slab_sub.receive().await?;
  120. info!("receive slab message!");
  121. let channels = self.slabman.lock().await.get_channels().unwrap_or(vec![]);
  122. let slab = messages::SlabMessage {
  123. nonce: slab_msg.nonce,
  124. ciphertext: slab_msg.ciphertext.clone(),
  125. };
  126. for channel in channels.iter() {
  127. match aes_decrypt(
  128. &channel.get_channel_secret(),
  129. &slab_msg.nonce,
  130. &slab_msg.ciphertext,
  131. ) {
  132. Some(plaintext) => {
  133. self.slabman
  134. .lock()
  135. .await
  136. .add_new_slab(slab.clone())
  137. .await
  138. .expect("error during adding new slab to database");
  139. let des_plaintext: ControlMessage = deserialize(&plaintext[..])
  140. .expect("error during deserializing the message");
  141. match des_plaintext.control {
  142. ControlCommand::Join => {
  143. info!("{} joined the group", des_plaintext.payload.nickname);
  144. }
  145. ControlCommand::Leave => {
  146. info!("{} left the group", des_plaintext.payload.nickname);
  147. }
  148. ControlCommand::Message => {
  149. info!(
  150. "{} -> {}: {}",
  151. des_plaintext.payload.timestamp,
  152. des_plaintext.payload.nickname,
  153. des_plaintext.payload.text
  154. );
  155. }
  156. }
  157. }
  158. None => {}
  159. }
  160. }
  161. }
  162. }
  163. }