1use super::*;
4
5impl Client {
6 #[cfg_attr(
12 feature = "tracing",
13 tracing::instrument(name = "wa.recv.request_retry", level = "debug", skip_all, err(Debug))
14 )]
15 pub async fn request_message_retry(
16 self: &Arc<Self>,
17 stanza: &NodeRef<'_>,
18 options: crate::features::RetryRequestOptions,
19 ) -> Result<crate::features::RetryRequestOutcome, crate::features::RetryRequestError> {
20 if stanza.tag.as_ref() != "message" {
21 return Err(crate::features::RetryRequestError::UnsupportedStanzaClass);
22 }
23 if stanza.get_attr("id").is_none() {
24 return Err(crate::features::RetryRequestError::MissingAttribute("id"));
25 }
26 if stanza.get_attr("from").is_none() {
27 return Err(crate::features::RetryRequestError::MissingAttribute("from"));
28 }
29 if !self.is_connected() {
30 return Err(crate::client::ClientError::NotConnected.into());
31 }
32
33 let device = self.persistence_manager.get_device_snapshot();
34 let own_pn = device
35 .pn
36 .as_ref()
37 .ok_or(crate::features::RetryRequestError::MissingLocalIdentity)?;
38 let info = wacore::messages::parse_message_info(stanza, own_pn, device.lid.as_ref())
39 .map_err(crate::features::RetryRequestError::InvalidStanza)?;
40 let max_sender_retry_count = message_enc_nodes_for_device(stanza, Some(own_pn))
41 .map(sender_retry_count)
42 .max()
43 .unwrap_or(0);
44 let info = Arc::new(info);
45 drop(device);
46
47 self.request_retry_for_info(
48 &info,
49 options,
50 (max_sender_retry_count > 0).then_some(max_sender_retry_count),
51 )
52 .await
53 }
54
55 #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.recv.undecryptable", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), msg_id = %info.id)))]
64 pub(crate) async fn dispatch_undecryptable_event(
65 &self,
66 info: Arc<MessageInfo>,
67 is_unavailable: bool,
68 unavailable_type: crate::types::events::UnavailableType,
69 decrypt_fail_mode: crate::types::events::DecryptFailMode,
70 ) -> bool {
71 let dedup_key =
72 wacore::types::message::ChatMessageId::new(info.source.chat.clone(), info.id.clone());
73 let fresh = Arc::new(std::sync::atomic::AtomicBool::new(false));
76 let fresh_clone = fresh.clone();
77 self.undecryptable_dispatched
78 .get_with(dedup_key, async move {
79 fresh_clone.store(true, Ordering::Release);
80 })
81 .await;
82 let was_fresh = fresh.load(Ordering::Acquire);
83 if was_fresh {
84 wacore::telemetry::recv("undecryptable");
85 self.core.event_bus.dispatch(Event::UndecryptableMessage(
86 crate::types::events::UndecryptableMessage::builder()
87 .info(info)
88 .is_unavailable(is_unavailable)
89 .unavailable_type(unavailable_type)
90 .decrypt_fail_mode(decrypt_fail_mode)
91 .build(),
92 ));
93 } else {
94 log::debug!(
95 "[msg:{}] UndecryptableMessage already dispatched for this id; skipping duplicate event",
96 info.id,
97 );
98 }
99 was_fresh
100 }
101
102 #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.recv.decrypt_failure", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), sender = %info.source.sender.observe(), msg_id = %info.id, reason = ?reason)))]
116 pub(crate) async fn handle_decrypt_failure(
117 self: &Arc<Self>,
118 info: &Arc<MessageInfo>,
119 reason: RetryReason,
120 decrypt_fail_mode: crate::types::events::DecryptFailMode,
121 ) -> bool {
122 self.dispatch_undecryptable_event(
123 Arc::clone(info),
124 false,
125 crate::types::events::UnavailableType::Unknown,
126 decrypt_fail_mode,
127 )
128 .await;
129 let client = Arc::clone(self);
130 let info = Arc::clone(info);
131 self.outbound_flush.spawn(&*self.runtime, async move {
132 if info.source.is_self_fanout()
140 && !info.source.is_bot_authored_non_bot_chat()
141 && Self::should_send_delivery_receipt(&info)
142 {
143 client.send_delivery_receipt(&info).await;
144 return;
145 }
146 let resend_sent = client.run_retry_receipt(&info, reason).await;
149 if resend_sent {
150 client.send_transport_ack(&info).await;
151 }
152 });
153 true
154 }
155
156 #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.recv.plaintext_failure", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), msg_id = %info.id)))]
157 pub(crate) async fn handle_plaintext_failure(
158 self: &Arc<Self>,
159 info: &Arc<MessageInfo>,
160 decrypt_fail_mode: crate::types::events::DecryptFailMode,
161 ) -> bool {
162 let dispatched = self
163 .dispatch_undecryptable_event(
164 Arc::clone(info),
165 false,
166 crate::types::events::UnavailableType::Unknown,
167 decrypt_fail_mode,
168 )
169 .await;
170 self.spawn_nack(info, NackReason::InvalidProtobuf, None);
171 dispatched
172 }
173
174 pub(crate) async fn increment_retry_count(
178 &self,
179 cache_key: &str,
180 reason: RetryReason,
181 ) -> Option<u8> {
182 self.message_retry_counts
183 .upsert_with_by_ref(cache_key, |current| {
184 let count = match current {
185 Some((count, _)) if *count >= MAX_DECRYPT_RETRIES => return (None, None),
186 Some((count, _)) => *count + 1,
187 None => 1,
188 };
189 (Some((count, Some(reason))), Some(count))
190 })
191 .await
192 }
193
194 pub(crate) async fn preseed_retry_count(&self, cache_key: &str, sender_count: u8) {
197 self.message_retry_counts
198 .upsert_with_by_ref(cache_key, |current| match current {
199 Some((count, _)) if *count >= sender_count => (None, ()),
200 Some((_, reason)) => (Some((sender_count, *reason)), ()),
201 None => (Some((sender_count, None)), ()),
202 })
203 .await;
204 }
205
206 pub(crate) async fn make_retry_cache_key(
208 &self,
209 chat: &Jid,
210 msg_id: &str,
211 sender: &Jid,
212 ) -> String {
213 let (chat, sender) = futures::join!(
215 self.resolve_encryption_jid(chat),
216 self.resolve_encryption_jid(sender),
217 );
218 let mut key =
220 String::with_capacity(chat.user.len() + msg_id.len() + sender.user.len() + 40);
221 chat.push_to(&mut key);
222 key.push(':');
223 key.push_str(msg_id);
224 key.push(':');
225 sender.push_to(&mut key);
226 key
227 }
228
229 #[cfg(test)]
254 pub(crate) fn spawn_retry_receipt(
255 self: &Arc<Self>,
256 info: &Arc<MessageInfo>,
257 reason: RetryReason,
258 ) {
259 let client = Arc::clone(self);
260 let info = Arc::clone(info);
261 self.outbound_flush.spawn(&*self.runtime, async move {
262 client.run_retry_receipt(&info, reason).await;
263 });
264 }
265
266 async fn request_retry_for_info(
270 self: &Arc<Self>,
271 info: &Arc<MessageInfo>,
272 options: crate::features::RetryRequestOptions,
273 sender_retry_count: Option<u8>,
274 ) -> Result<crate::features::RetryRequestOutcome, crate::features::RetryRequestError> {
275 let reason = options.reason();
276 let cache_key = self
277 .make_retry_cache_key(&info.source.chat, &info.id, &info.source.sender)
278 .await;
279
280 if let Some(sender_retry_count) = sender_retry_count {
281 self.preseed_retry_count(&cache_key, sender_retry_count)
282 .await;
283 }
284
285 let Some(retry_count) = self.increment_retry_count(&cache_key, reason).await else {
286 log::debug!(
287 "Max retries ({}) reached for message {} from {} [{:?}]. Requesting PDO fallback.",
288 MAX_DECRYPT_RETRIES,
289 info.id,
290 info.source.sender.observe(),
291 reason
292 );
293 self.run_pdo_request(info).await;
294 return Ok(crate::features::RetryRequestOutcome::LimitReached);
295 };
296
297 if retry_count > HIGH_RETRY_COUNT_THRESHOLD {
298 log::warn!(
299 "High retry count ({}) for message {} in chat {} from {} [{:?}]",
300 retry_count,
301 info.id,
302 info.source.chat.observe(),
303 info.source.sender.observe(),
304 reason
305 );
306 }
307
308 let send_result = self
309 .send_retry_receipt(info, retry_count, reason, options.force_include_keys())
310 .await;
311
312 if retry_count == 1 {
317 self.run_pdo_request(info).await;
318 }
319
320 let send_outcome = send_result?;
321
322 let outcome = match send_outcome {
323 crate::retry::RetryReceiptSendOutcome::Sent { included_keys } => {
324 wacore::telemetry::retry_receipt(reason.as_str());
325 if retry_count >= MAX_DECRYPT_RETRIES {
326 wacore::telemetry::high_retry(reason.as_str());
327 }
328 debug!(
329 "Sent retry receipt #{} for message {} in chat {} from {} [{:?}]",
330 retry_count,
331 info.id,
332 info.source.chat.observe(),
333 info.source.sender.observe(),
334 reason
335 );
336 crate::features::RetryRequestOutcome::Sent {
337 retry_count,
338 included_keys,
339 }
340 }
341 crate::retry::RetryReceiptSendOutcome::Suppressed => {
342 crate::features::RetryRequestOutcome::Suppressed { retry_count }
343 }
344 };
345
346 Ok(outcome)
347 }
348
349 #[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.recv.retry_receipt", level = "debug", skip_all, fields(chat = %info.source.chat.observe(), sender = %info.source.sender.observe(), msg_id = %info.id, reason = ?reason)))]
355 async fn run_retry_receipt(
356 self: &Arc<Self>,
357 info: &Arc<MessageInfo>,
358 reason: RetryReason,
359 ) -> bool {
360 match self
361 .request_retry_for_info(
362 info,
363 crate::features::RetryRequestOptions::new().with_reason(reason),
364 None,
365 )
366 .await
367 {
368 Ok(_) => true,
369 Err(error) => {
370 log::error!(
371 "Failed to send retry receipt for message {} [{:?}]: {error:?}",
372 info.id,
373 reason
374 );
375 false
376 }
377 }
378 }
379}