1use std::{
2 cmp::Reverse,
3 collections::{BTreeMap, BTreeSet, BinaryHeap, HashMap, HashSet, btree_map::Entry},
4 fmt::Debug,
5 iter::once,
6 str::from_utf8,
7 sync::{
8 Arc,
9 atomic::{AtomicU64, AtomicUsize, Ordering},
10 },
11 time::{Duration, Instant, SystemTime},
12};
13
14use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
15use futures::{
16 FutureExt, StreamExt, TryFutureExt,
17 future::{BoxFuture, join, join3, try_join3},
18 pin_mut,
19 stream::FuturesUnordered,
20};
21use ruma::{
22 MilliSecondsSinceUnixEpoch, OneTimeKeyAlgorithm, OwnedDeviceId, OwnedRoomId, OwnedServerName,
23 OwnedUserId, RoomId, ServerName, UInt, UserId,
24 api::{
25 appservice::event::push_events::v1::{
26 DeviceLists, EphemeralData, Request as PushEventsRequest,
27 },
28 client::push::Pusher,
29 federation::transactions::{
30 edu::{
31 DeviceListUpdateContent, Edu, PresenceContent, PresenceUpdate, ReceiptContent,
32 ReceiptData, ReceiptMap,
33 },
34 send_transaction_message,
35 },
36 },
37 device_id,
38 events::{
39 AnySyncEphemeralRoomEvent, GlobalAccountDataEventType, push_rules::PushRulesEvent,
40 receipt::ReceiptType,
41 },
42 presence::PresenceState,
43 push::Ruleset,
44 serde::Raw,
45 uint,
46};
47use serde::Deserialize;
48use tuwunel_core::{
49 Error, Event, Result, debug, debug_warn, err, error, extract_variant,
50 result::LogErr,
51 smallvec::SmallVec,
52 trace,
53 utils::{
54 BoolExt, ReadyExt, calculate_hash, continue_exponential_backoff_secs,
55 future::TryExtExt,
56 rand::secs as rand_secs,
57 stream::{BroadbandExt, IterStream, WidebandExt},
58 },
59 warn,
60};
61
62use super::{
63 Destination, EduBuf, EduVec, Msg, SendingEvent, Service, TAG_PREFIX_LEN, data::QueueItem,
64 reap_flushes,
65};
66use crate::{federation::ShouldAttempt, rooms::timeline::RawPduId};
67
68#[derive(Debug)]
72enum TransactionStatus {
73 Running,
74 RunningForceRetry,
75 Failed(u32, Instant), Retrying(u32), }
78
79enum RetryAction {
80 None,
81 Force,
82}
83
84type SendingError = (Destination, Error);
85type SendingResult = Result<Destination, SendingError>;
86type SendingFuture<'a> = BoxFuture<'a, SendingResult>;
87type SendingFutures<'a> = FuturesUnordered<SendingFuture<'a>>;
88type CurTransactionStatus = HashMap<Destination, TransactionStatus>;
89
90type OtkCounts =
93 BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, BTreeMap<OneTimeKeyAlgorithm, UInt>>>;
94
95type FallbackTypes = BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, Vec<OneTimeKeyAlgorithm>>>;
98
99type Devices = SmallVec<[(OwnedUserId, OwnedDeviceId); 1]>;
102
103type WakeQueue = BinaryHeap<Reverse<(tokio::time::Instant, OwnedServerName)>>;
107
108type UserReceipts = SmallVec<[ReceiptData; 1]>;
112
113type RankedReceipts = SmallVec<[ReceiptMap; 1]>;
118
119type RoomReceipts = SmallVec<[(OwnedRoomId, RankedReceipts); 1]>;
122
123#[derive(Default)]
127struct Selected {
128 shipped: EduVec,
129 overflow: Vec<EduBuf>,
130}
131
132#[derive(Deserialize)]
136struct ToDeviceRecipient {
137 to_user_id: OwnedUserId,
138 to_device_id: OwnedDeviceId,
139}
140
141const SELECT_PRESENCE_LIMIT: usize = 256;
142const SELECT_RECEIPT_LIMIT: usize = 256;
143const DEQUEUE_LIMIT: usize = 48;
144
145pub const PDU_LIMIT: usize = 50;
146pub const EDU_LIMIT: usize = 100;
147
148impl Service {
149 #[tracing::instrument(skip(self), level = "debug")]
150 pub(super) async fn sender(self: Arc<Self>, id: usize) -> Result {
151 let mut statuses: CurTransactionStatus = CurTransactionStatus::new();
152 let mut futures: SendingFutures<'_> = FuturesUnordered::new();
153 let mut wakes: WakeQueue = WakeQueue::new();
154
155 self.startup_netburst(id, &mut futures, &mut statuses)
156 .boxed()
157 .await;
158
159 self.work_loop(id, &mut futures, &mut statuses, &mut wakes)
160 .await;
161
162 if !futures.is_empty() {
163 self.finish_responses(&mut futures).boxed().await;
164 }
165
166 Ok(())
167 }
168
169 #[tracing::instrument(
170 name = "work",
171 level = "trace",
172 skip_all,
173 fields(
174 futures = %futures.len(),
175 statuses = %statuses.len(),
176 ),
177 )]
178 async fn work_loop<'a>(
179 &'a self,
180 id: usize,
181 futures: &mut SendingFutures<'a>,
182 statuses: &mut CurTransactionStatus,
183 wakes: &mut WakeQueue,
184 ) {
185 use tokio::time::{Instant, sleep_until};
186
187 let receiver = self
188 .channels
189 .get(id)
190 .map(|(_, receiver)| receiver.clone())
191 .expect("Missing channel for sender worker");
192
193 while !receiver.is_closed() {
194 let next_due = wakes
195 .peek()
196 .map_or_else(Instant::now, |Reverse((instant, _))| *instant);
197
198 tokio::select! {
199 Some(response) = futures.next() => {
200 self.handle_response(response, futures, statuses, wakes).await;
201 },
202 request = receiver.recv_async() => match request {
203 Ok(request) => self.handle_request(request, futures, statuses).await,
204 Err(_) => return,
205 },
206 () = sleep_until(next_due), if !wakes.is_empty() => {
207 self.drain_due_wakes(futures, statuses, wakes).await;
208 },
209 }
210 }
211 }
212
213 #[tracing::instrument(name = "response", level = "debug", skip_all)]
214 async fn handle_response<'a>(
215 &'a self,
216 response: SendingResult,
217 futures: &mut SendingFutures<'a>,
218 statuses: &mut CurTransactionStatus,
219 wakes: &mut WakeQueue,
220 ) {
221 match response {
222 | Ok(dest) =>
223 self.handle_response_ok(&dest, futures, statuses)
224 .await,
225 | Err((dest, e)) => {
226 let retry_action = Self::handle_response_err(&dest, statuses, &e);
227
228 match dest {
229 | Destination::Federation(server) => {
230 if let ShouldAttempt::No { earliest_retry } = self
232 .services
233 .federation
234 .should_attempt(&server)
235 .await
236 {
237 arm_wake(wakes, server, earliest_retry);
238 }
239 },
240 | dest if matches!(retry_action, RetryAction::Force) =>
241 self.handle_force_retry(dest, futures, statuses)
242 .await,
243 | _ => {},
244 }
245 },
246 }
247 }
248
249 async fn handle_force_retry<'a>(
250 &'a self,
251 dest: Destination,
252 futures: &mut SendingFutures<'a>,
253 statuses: &mut CurTransactionStatus,
254 ) {
255 let Ok(Some(events)) = self
256 .select_events(&dest, Vec::new(), statuses)
257 .await
258 else {
259 return;
260 };
261
262 self.schedule_events(dest, events, futures, statuses);
263 }
264
265 fn handle_response_err(
266 dest: &Destination,
267 statuses: &mut CurTransactionStatus,
268 e: &Error,
269 ) -> RetryAction {
270 debug!(?dest, "{e:?}");
271 let push = matches!(dest, Destination::Push(..));
273
274 let Some(status) = statuses.get_mut(dest) else {
275 return RetryAction::None;
276 };
277
278 let (tries, retry_action) = match status {
279 | TransactionStatus::Running => (1, RetryAction::None),
280 | TransactionStatus::RunningForceRetry => (1, RetryAction::Force),
281 | TransactionStatus::Failed(n, _) | TransactionStatus::Retrying(n) =>
282 (n.saturating_add(1), RetryAction::None),
283 };
284
285 *status = if push {
286 TransactionStatus::Failed(tries, Instant::now())
287 } else {
288 TransactionStatus::Retrying(tries)
289 };
290
291 retry_action
292 }
293
294 #[expect(clippy::needless_pass_by_ref_mut)]
295 async fn handle_response_ok<'a>(
296 &'a self,
297 dest: &Destination,
298 futures: &mut SendingFutures<'a>,
299 statuses: &mut CurTransactionStatus,
300 ) {
301 let _cork = self.db.db.cork();
302 self.db.delete_all_active_requests_for(dest).await;
303
304 let new_events = self
306 .db
307 .queued_requests(dest)
308 .take(DEQUEUE_LIMIT)
309 .collect::<Vec<_>>()
310 .await;
311
312 if !new_events.is_empty() {
313 self.db.mark_as_active(new_events.iter());
314 }
315
316 let mut events: Vec<SendingEvent> = new_events
317 .into_iter()
318 .map(|(_, event)| event)
319 .collect();
320
321 if let Destination::Federation(server_name) = dest {
323 let budget_used = events
324 .iter()
325 .filter(|event| matches!(event, SendingEvent::Edu(_)))
326 .count();
327
328 if let Ok(select_edus) = self.select_edus(server_name, budget_used).await {
329 events.extend(select_edus.into_iter().map(SendingEvent::Edu));
330 }
331 }
332
333 if events.is_empty() {
334 statuses.remove(dest);
335 } else {
336 if let Some(status) = statuses.get_mut(dest) {
337 *status = TransactionStatus::Running;
338 }
339
340 futures.push(self.send_events(dest.clone(), events));
341 }
342 }
343
344 #[expect(
345 clippy::needless_pass_by_ref_mut,
346 reason = "mutable reference avoids requiring SendingFutures to be Sync"
347 )]
348 fn schedule_events<'a>(
349 &'a self,
350 dest: Destination,
351 events: Vec<SendingEvent>,
352 futures: &mut SendingFutures<'a>,
353 statuses: &mut CurTransactionStatus,
354 ) {
355 if events.is_empty() {
356 statuses.remove(&dest);
357 } else {
358 futures.push(self.send_events(dest, events));
359 }
360 }
361
362 #[tracing::instrument(name = "request", level = "debug", skip_all)]
363 async fn handle_request<'a>(
364 &'a self,
365 msg: Msg,
366 futures: &mut SendingFutures<'a>,
367 statuses: &mut CurTransactionStatus,
368 ) {
369 let synthetic_badge =
370 msg.queue_id.is_empty() && matches!(&msg.event, SendingEvent::BadgeRefresh);
371
372 let new_events = match (synthetic_badge, statuses.contains_key(&msg.dest)) {
373 | (false, _) => vec![(msg.queue_id, msg.event)],
374 | (true, true) => Vec::new(),
375 | (true, false) =>
376 self.db
377 .queued_requests(&msg.dest)
378 .take(DEQUEUE_LIMIT)
379 .collect()
380 .await,
381 };
382
383 if let Ok(Some(events)) = self
384 .select_events(&msg.dest, new_events, statuses)
385 .await
386 {
387 self.schedule_events(msg.dest, events, futures, statuses);
388 }
389 }
390
391 async fn drain_due_wakes<'a>(
392 &'a self,
393 futures: &mut SendingFutures<'a>,
394 statuses: &mut CurTransactionStatus,
395 wakes: &mut WakeQueue,
396 ) {
397 use tokio::time::Instant;
398
399 let now = Instant::now();
400 while wakes
401 .peek()
402 .is_some_and(|Reverse((due, _))| *due <= now)
403 {
404 let Reverse((_, server)) = wakes.pop().expect("peeked entry");
405 self.handle_wake(server, futures, statuses, wakes)
406 .await;
407 }
408 }
409
410 async fn handle_wake<'a>(
411 &'a self,
412 server: OwnedServerName,
413 futures: &mut SendingFutures<'a>,
414 statuses: &mut CurTransactionStatus,
415 wakes: &mut WakeQueue,
416 ) {
417 let dest = Destination::Federation(server.clone());
418 if matches!(statuses.get(&dest), Some(TransactionStatus::Running)) {
419 return;
420 }
421
422 match self
423 .services
424 .federation
425 .should_attempt(&server)
426 .await
427 {
428 | ShouldAttempt::No { earliest_retry } => arm_wake(wakes, server, earliest_retry),
429 | _ => {
430 let msg = Msg {
431 dest,
432 event: SendingEvent::Flush,
433 queue_id: Vec::new(),
434 };
435
436 self.handle_request(msg, futures, statuses).await;
437 },
438 }
439 }
440
441 #[tracing::instrument(
442 name = "finish",
443 level = "info",
444 skip_all,
445 fields(futures = %futures.len()),
446 )]
447 async fn finish_responses<'a>(&'a self, futures: &mut SendingFutures<'a>) {
448 use tokio::{
449 select,
450 time::{Instant, sleep_until},
451 };
452
453 let timeout = self.server.config.sender_shutdown_timeout;
454 let timeout = Duration::from_secs(timeout);
455 let now = Instant::now();
456 let deadline = now.checked_add(timeout).unwrap_or(now);
457 loop {
458 trace!("Waiting for {} requests to complete...", futures.len());
459 select! {
460 () = sleep_until(deadline) => return,
461 response = futures.next() => match response {
462 Some(Ok(dest)) => self.db.delete_all_active_requests_for(&dest).await,
463 Some(_) => {},
464 None => return,
465 },
466 }
467 }
468 }
469
470 #[tracing::instrument(
471 name = "netburst",
472 level = "debug",
473 skip_all,
474 fields(futures = %futures.len()),
475 )]
476 async fn startup_netburst<'a>(
477 &'a self,
478 id: usize,
479 futures: &mut SendingFutures<'a>,
480 statuses: &mut CurTransactionStatus,
481 ) {
482 let keep =
483 usize::try_from(self.server.config.startup_netburst_keep).unwrap_or(usize::MAX);
484
485 let mut txns = HashMap::<Destination, Vec<SendingEvent>>::new();
486 let active = self.db.active_requests();
487
488 pin_mut!(active);
489 while let Some((key, event, dest)) = active.next().await {
490 if self.shard_id(&dest) != id {
491 continue;
492 }
493
494 let entry = txns.entry(dest.clone()).or_default();
495 if self.server.config.startup_netburst_keep >= 0 && entry.len() >= keep {
496 warn!("Dropping unsent event {dest:?} {:?}", String::from_utf8_lossy(&key));
497 self.db.delete_active_request(&key);
498 } else {
499 entry.push(event);
500 }
501 }
502
503 for (dest, events) in txns {
504 if self.server.config.startup_netburst && !events.is_empty() {
505 statuses.insert(dest.clone(), TransactionStatus::Running);
506 futures.push(self.send_events(dest.clone(), events));
507 }
508 }
509
510 if !self.server.config.startup_netburst || keep == 0 {
513 return;
514 }
515
516 let destinations = self
517 .db
518 .queued_badge_refresh_destinations()
519 .ready_filter(|dest| self.shard_id(dest) == id)
520 .collect::<HashSet<_>>()
521 .await;
522
523 for dest in destinations {
524 let msg = Msg {
525 dest,
526 event: SendingEvent::BadgeRefresh,
527 queue_id: Vec::new(),
528 };
529
530 self.handle_request(msg, futures, statuses).await;
531 }
532 }
533
534 #[tracing::instrument(
535 name = "select",
536 level = "debug",
537 skip_all,
538 fields(
539 ?dest,
540 new_events = %new_events.len(),
541 ),
542 )]
543 async fn select_events(
544 &self,
545 dest: &Destination,
546 new_events: Vec<QueueItem>, statuses: &mut CurTransactionStatus,
548 ) -> Result<Option<Vec<SendingEvent>>> {
549 let retry_action = if matches!(dest, Destination::Appservice(_))
550 && new_events
551 .iter()
552 .any(|(_, event)| matches!(event, SendingEvent::Flush))
553 {
554 RetryAction::Force
555 } else {
556 RetryAction::None
557 };
558
559 let (allow, retry) = self
560 .select_events_current(dest, statuses, retry_action)
561 .await?;
562
563 if !allow {
565 return Ok(None);
566 }
567
568 let mut events = Vec::new();
569
570 if retry {
572 self.db
573 .active_requests_for(dest)
574 .ready_for_each(|(_, e)| events.push(e))
575 .await;
576
577 return Ok(Some(events));
578 }
579
580 let _cork = self.db.db.cork();
582 self.db
583 .retain_queued(new_events)
584 .ready_for_each(|item| {
585 self.db.mark_as_active(once(&item));
586 if !matches!(&item.1, SendingEvent::Flush) {
587 events.push(item.1);
588 }
589 })
590 .await;
591
592 if let Destination::Federation(server_name) = dest {
594 let budget_used = events
595 .iter()
596 .filter(|event| matches!(event, SendingEvent::Edu(_)))
597 .count();
598
599 if let Ok(select_edus) = self.select_edus(server_name, budget_used).await {
600 events.extend(select_edus.into_iter().map(SendingEvent::Edu));
601 }
602 }
603
604 Ok(Some(events))
605 }
606
607 async fn select_events_current(
608 &self,
609 dest: &Destination,
610 statuses: &mut CurTransactionStatus,
611 retry_action: RetryAction,
612 ) -> Result<(bool, bool)> {
613 if let Destination::Federation(server) = dest {
615 let should_attempt = self
616 .services
617 .federation
618 .should_attempt(server)
619 .await;
620
621 if matches!(should_attempt, ShouldAttempt::No { .. }) {
622 return Ok((false, false));
623 }
624 }
625
626 let (mut allow, mut retry) = (true, false);
627 statuses
628 .entry(dest.clone())
629 .and_modify(|e| match e {
630 | TransactionStatus::Running | TransactionStatus::RunningForceRetry => {
631 allow = false; if matches!(retry_action, RetryAction::Force) {
633 *e = TransactionStatus::RunningForceRetry;
634 }
635 },
636 | TransactionStatus::Failed(tries, time) => {
637 let min = self.server.config.sender_timeout;
639 let max = self.server.config.sender_retry_backoff_limit;
640 if continue_exponential_backoff_secs(min, max, time.elapsed(), *tries) {
641 allow = false;
642 } else {
643 retry = true;
644 *e = TransactionStatus::Retrying(*tries);
645 }
646 },
647 | TransactionStatus::Retrying(_) if matches!(dest, Destination::Push(..)) => {
648 allow = false; },
650 | TransactionStatus::Retrying(_) => {
651 retry = true;
653 *e = TransactionStatus::Running;
654 },
655 })
656 .or_insert(TransactionStatus::Running);
657
658 Ok((allow, retry))
659 }
660
661 #[tracing::instrument(name = "edus", level = "debug", skip_all)]
662 async fn select_edus(&self, server_name: &ServerName, budget_used: usize) -> Result<EduVec> {
663 let since = self.db.get_latest_educount(server_name).await;
665 let since_upper = self.services.globals.current_count();
666
667 if since == since_upper {
669 return Ok(EduVec::new());
670 }
671
672 let batch = (since, since_upper);
673 debug_assert!(batch.0 <= batch.1, "since range must not be negative");
674
675 let events_len = AtomicUsize::new(budget_used);
676 let max_edu_count = AtomicU64::new(since);
677 let device_changes =
678 self.select_edus_device_changes(server_name, batch, &max_edu_count, &events_len);
679
680 let receipts = self
681 .server
682 .config
683 .allow_outgoing_read_receipts
684 .then_async(|| {
685 self.select_edus_receipts(server_name, batch, &max_edu_count, &events_len)
686 });
687
688 let presence = self
689 .server
690 .config
691 .allow_outgoing_presence
692 .then_async(|| {
693 self.select_edus_presence(server_name, batch, &max_edu_count, &events_len)
694 });
695
696 let (device_changes, receipts, presence) =
697 join3(device_changes, receipts, presence).await;
698
699 let receipts = receipts.unwrap_or_default();
700 let mut events = device_changes.shipped;
701
702 events.extend(receipts.shipped);
703
704 let durable_len = events.len();
707
708 events.extend(presence.flatten());
709 debug_assert!(
710 budget_used.saturating_add(events.len()) <= EDU_LIMIT,
711 "exceeded edus limit"
712 );
713
714 let overflow: Vec<SendingEvent> = device_changes
716 .overflow
717 .into_iter()
718 .chain(receipts.overflow)
719 .map(SendingEvent::Edu)
720 .collect();
721
722 if !overflow.is_empty() {
723 let dest = Destination::Federation(server_name.to_owned());
724 self.db
725 .queue_requests(overflow.iter().map(|event| (event, &dest)));
726 }
727
728 if durable_len > 0 {
731 self.db
732 .persist_active_edus(server_name, &events[..durable_len]);
733 }
734
735 let last_count = max_edu_count.load(Ordering::Acquire);
736 if last_count > since {
737 self.db
738 .set_latest_educount(server_name, last_count);
739 }
740
741 Ok(events)
742 }
743
744 #[tracing::instrument(
746 name = "device_changes",
747 level = "trace",
748 skip(self, server_name, max_edu_count, events_len)
749 )]
750 async fn select_edus_device_changes(
751 &self,
752 server_name: &ServerName,
753 since: (u64, u64),
754 max_edu_count: &AtomicU64,
755 events_len: &AtomicUsize,
756 ) -> Selected {
757 let mut selected = Selected::default();
758 let server_rooms = self
759 .services
760 .state_cache
761 .server_rooms(server_name);
762
763 pin_mut!(server_rooms);
764 let mut device_list_changes = HashSet::<OwnedUserId>::new();
765 while let Some(room_id) = server_rooms.next().await {
766 let keys_changed = self
767 .services
768 .users
769 .room_keys_changed(room_id, since.0, Some(since.1))
770 .ready_filter(|(user_id, _)| self.services.globals.user_is_local(user_id));
771
772 pin_mut!(keys_changed);
773 while let Some((user_id, count)) = keys_changed.next().await {
774 debug_assert!(count <= since.1, "exceeds upper-bound");
775
776 max_edu_count.fetch_max(count, Ordering::Relaxed);
777 if !device_list_changes.insert(user_id.into()) {
778 continue;
779 }
780
781 let edu = Edu::DeviceListUpdate(DeviceListUpdateContent {
784 user_id: user_id.into(),
785 device_id: device_id!("placeholder").to_owned(),
786 device_display_name: Some("Placeholder".to_owned()),
787 stream_id: uint!(1),
788 prev_id: Vec::new(),
789 deleted: None,
790 keys: None,
791 });
792
793 let mut buf = EduBuf::new();
794 serde_json::to_writer(&mut buf, &edu)
795 .expect("failed to serialize device list update to JSON");
796
797 if !selected.overflow.is_empty()
800 || events_len.fetch_add(1, Ordering::Relaxed) >= EDU_LIMIT
801 {
802 selected.overflow.push(buf);
803 } else {
804 selected.shipped.push(buf);
805 }
806 }
807 }
808
809 selected
810 }
811
812 #[tracing::instrument(
822 name = "receipts",
823 level = "trace",
824 skip(self, server_name, max_edu_count, events_len)
825 )]
826 async fn select_edus_receipts(
827 &self,
828 server_name: &ServerName,
829 since: (u64, u64),
830 max_edu_count: &AtomicU64,
831 events_len: &AtomicUsize,
832 ) -> Selected {
833 let num = AtomicUsize::new(0);
834 let by_room: RoomReceipts = self
835 .services
836 .state_cache
837 .server_rooms(server_name)
838 .map(ToOwned::to_owned)
839 .broad_filter_map(async |room_id| {
840 let ranked = self
841 .select_edus_receipts_room(&room_id, since, max_edu_count, &num)
842 .await;
843
844 ranked
845 .is_empty()
846 .is_false()
847 .then_some((room_id, ranked))
848 })
849 .collect()
850 .boxed()
851 .await;
852
853 let max_rank = by_room
854 .iter()
855 .map(|(_, maps)| maps.len())
856 .max()
857 .unwrap_or(0);
858
859 let pivot_rank = |rank: usize| -> Option<BTreeMap<OwnedRoomId, ReceiptMap>> {
860 let receipts: BTreeMap<_, _> = by_room
861 .iter()
862 .filter_map(|(room_id, maps)| {
863 maps.get(rank)
864 .cloned()
865 .map(|map| (room_id.clone(), map))
866 })
867 .collect();
868
869 receipts.is_empty().is_false().then_some(receipts)
870 };
871
872 let serialize_edu = |receipts: BTreeMap<OwnedRoomId, ReceiptMap>| -> EduBuf {
873 let mut buf = EduBuf::new();
874 serde_json::to_writer(&mut buf, &Edu::Receipt(ReceiptContent { receipts }))
875 .expect("Failed to serialize Receipt EDU to JSON vec");
876
877 buf
878 };
879
880 let mut selected = Selected::default();
883 for receipts in (0..max_rank).filter_map(pivot_rank) {
884 if !selected.overflow.is_empty()
885 || events_len.fetch_add(1, Ordering::Relaxed) >= EDU_LIMIT
886 {
887 selected.overflow.push(serialize_edu(receipts));
888 } else {
889 selected.shipped.push(serialize_edu(receipts));
890 }
891 }
892
893 selected
894 }
895
896 #[tracing::instrument(
905 name = "receipts",
906 level = "trace",
907 skip(self, since, max_edu_count)
908 )]
909 async fn select_edus_receipts_room(
910 &self,
911 room_id: &RoomId,
912 since: (u64, u64),
913 max_edu_count: &AtomicU64,
914 num: &AtomicUsize,
915 ) -> RankedReceipts {
916 let receipts =
917 self.services
918 .read_receipt
919 .readreceipts_since(room_id, since.0, Some(since.1));
920
921 pin_mut!(receipts);
922 let mut by_user = BTreeMap::<OwnedUserId, UserReceipts>::new();
923 while let Some((user_id, count, read_receipt)) = receipts.next().await {
924 debug_assert!(count <= since.1, "exceeds upper-bound");
925
926 max_edu_count.fetch_max(count, Ordering::Relaxed);
927 if !self.services.globals.user_is_local(user_id) {
928 continue;
929 }
930
931 let Ok(event) = serde_json::from_str(read_receipt.json().get()) else {
932 error!(?user_id, ?count, ?read_receipt, "Invalid edu event in read_receipts.");
933 continue;
934 };
935
936 let AnySyncEphemeralRoomEvent::Receipt(r) = event else {
937 error!(?user_id, ?count, ?event, "Invalid event type in read_receipts");
938 continue;
939 };
940
941 let (event_id, mut receipt) = r
942 .content
943 .0
944 .into_iter()
945 .next()
946 .expect("we only use one event per read receipt");
947
948 let receipt = receipt
949 .remove(&ReceiptType::Read)
950 .expect("our read receipts always set this")
951 .remove(user_id)
952 .expect("our read receipts always have the user here");
953
954 let receipt_data = ReceiptData { data: receipt, event_ids: vec![event_id] };
955
956 match by_user.entry(user_id.to_owned()) {
957 | Entry::Vacant(slot) => {
958 slot.insert(SmallVec::from_buf([receipt_data]));
959 let num = num.fetch_add(1, Ordering::Relaxed);
960 if num >= SELECT_RECEIPT_LIMIT {
961 break;
962 }
963 },
964 | Entry::Occupied(mut slot) => {
965 slot.get_mut().push(receipt_data);
966 },
967 }
968 }
969
970 by_user
974 .into_iter()
975 .fold(RankedReceipts::new(), |mut acc, (user_id, receipts)| {
976 for (rank, receipt_data) in receipts.into_iter().enumerate() {
977 if rank >= acc.len() {
978 acc.push(ReceiptMap { read: BTreeMap::new() });
979 }
980
981 acc[rank]
982 .read
983 .insert(user_id.clone(), receipt_data);
984 }
985
986 acc
987 })
988 }
989
990 #[tracing::instrument(
992 name = "presence",
993 level = "trace",
994 skip(self, server_name, max_edu_count, events_len)
995 )]
996 async fn select_edus_presence(
997 &self,
998 server_name: &ServerName,
999 since: (u64, u64),
1000 max_edu_count: &AtomicU64,
1001 events_len: &AtomicUsize,
1002 ) -> Option<EduBuf> {
1003 let presence_since = self
1004 .services
1005 .presence
1006 .presence_since(since.0, Some(since.1));
1007
1008 pin_mut!(presence_since);
1009 let mut presence_updates = HashMap::<OwnedUserId, PresenceUpdate>::new();
1010 while let Some((user_id, count, presence_bytes)) = presence_since.next().await {
1011 debug_assert!(count <= since.1, "exceeded upper-bound");
1012
1013 max_edu_count.fetch_max(count, Ordering::Relaxed);
1014 if !self.services.globals.user_is_local(user_id) {
1015 continue;
1016 }
1017
1018 if !self
1019 .services
1020 .state_cache
1021 .server_sees_user(server_name, user_id)
1022 .await
1023 {
1024 continue;
1025 }
1026
1027 let Ok(presence_event) = self
1028 .services
1029 .presence
1030 .from_json_bytes_to_event(presence_bytes, user_id)
1031 .await
1032 .log_err()
1033 else {
1034 continue;
1035 };
1036
1037 let update = PresenceUpdate {
1038 user_id: user_id.into(),
1039 presence: presence_event.content.presence,
1040 currently_active: presence_event
1041 .content
1042 .currently_active
1043 .unwrap_or(false),
1044 status_msg: presence_event.content.status_msg,
1045 last_active_ago: presence_event
1046 .content
1047 .last_active_ago
1048 .unwrap_or_else(|| uint!(0)),
1049 };
1050
1051 presence_updates.insert(user_id.into(), update);
1052 if presence_updates.len() >= SELECT_PRESENCE_LIMIT {
1053 break;
1054 }
1055 }
1056
1057 if presence_updates.is_empty() {
1058 return None;
1059 }
1060
1061 if events_len.fetch_add(1, Ordering::Relaxed) >= EDU_LIMIT {
1063 return None;
1064 }
1065
1066 let presence_content = Edu::Presence(PresenceContent {
1067 push: presence_updates.into_values().collect(),
1068 });
1069
1070 let mut buf = EduBuf::new();
1071 serde_json::to_writer(&mut buf, &presence_content)
1072 .expect("failed to serialize Presence EDU to JSON");
1073
1074 Some(buf)
1075 }
1076
1077 fn send_events(&self, dest: Destination, events: Vec<SendingEvent>) -> SendingFuture<'_> {
1078 debug_assert!(!events.is_empty(), "sending empty transaction");
1079 match dest {
1080 | Destination::Federation(server) => self
1081 .send_events_dest_federation(server, events)
1082 .boxed(),
1083 | Destination::Appservice(id) => self
1084 .send_events_dest_appservice(id, events)
1085 .boxed(),
1086 | Destination::Push(user_id, pushkey) => self
1087 .send_events_dest_push(user_id, pushkey, events)
1088 .boxed(),
1089 }
1090 }
1091
1092 #[tracing::instrument(
1093 name = "appservice",
1094 level = "debug",
1095 skip(self, events),
1096 fields(
1097 events = %events.len(),
1098 ),
1099 )]
1100 async fn send_events_dest_appservice(
1101 &self,
1102 id: String,
1103 events: Vec<SendingEvent>,
1104 ) -> SendingResult {
1105 let Some(info) = self
1106 .services
1107 .appservice
1108 .get_registration_info(&id)
1109 .await
1110 else {
1111 return Err((
1113 Destination::Appservice(id.clone()),
1114 err!(Database(debug_warn!(?id, "Missing appservice registration"))),
1115 ));
1116 };
1117
1118 let msc3202 = info.registration.msc3202_transaction_extensions;
1119
1120 let (pdu_count, edu_count, to_device_count, device_list_count) = events.iter().fold(
1121 (0_usize, 0_usize, 0_usize, 0_usize),
1122 |(pdus, edus, to_device, device_list), event| match event {
1123 | SendingEvent::Pdu(_) => (pdus.saturating_add(1), edus, to_device, device_list),
1124 | SendingEvent::Edu(_) => (pdus, edus.saturating_add(1), to_device, device_list),
1125 | SendingEvent::ToDevice(_) =>
1126 (pdus, edus, to_device.saturating_add(1), device_list),
1127 | SendingEvent::DeviceListChanged(_) =>
1128 (pdus, edus, to_device, device_list.saturating_add(1)),
1129 | SendingEvent::BadgeRefresh | SendingEvent::Flush =>
1130 (pdus, edus, to_device, device_list),
1131 },
1132 );
1133
1134 let mut pdu_jsons = Vec::with_capacity(pdu_count);
1135 let mut edu_jsons: Vec<Raw<EphemeralData>> = Vec::with_capacity(edu_count);
1136 let mut to_device = Vec::with_capacity(to_device_count);
1137 let mut changed = Vec::with_capacity(device_list_count);
1138
1139 let mut otk_users = BTreeSet::new();
1142 let mut otk_recipients = BTreeSet::new();
1143 if msc3202 {
1144 otk_users.insert(info.sender.clone());
1145 }
1146
1147 for event in &events {
1148 match event {
1149 | SendingEvent::Pdu(pdu_id) => {
1150 if let Ok(pdu) = self
1151 .services
1152 .timeline
1153 .get_pdu_from_id(pdu_id)
1154 .await
1155 {
1156 if msc3202 && info.is_user_match(pdu.sender()) {
1157 otk_users.insert(pdu.sender().to_owned());
1158 }
1159
1160 pdu_jsons.push(pdu.to_format());
1161 }
1162 },
1163 | SendingEvent::Edu(edu) => {
1164 if info.registration.receive_ephemeral
1165 && let Ok(edu) =
1166 serde_json::from_slice(edu).and_then(|edu| Raw::new(&edu))
1167 {
1168 edu_jsons.push(edu);
1169 }
1170 },
1171 | SendingEvent::ToDevice(buf) => {
1172 let Some(bytes) = buf.get(TAG_PREFIX_LEN..) else {
1173 debug_warn!("skipping malformed queued to-device event");
1174 continue;
1175 };
1176
1177 if msc3202
1178 && let Ok(recipient) = serde_json::from_slice::<ToDeviceRecipient>(bytes)
1179 {
1180 otk_recipients.insert((recipient.to_user_id, recipient.to_device_id));
1181 }
1182
1183 if let Ok(raw) = serde_json::from_slice(bytes) {
1184 to_device.push(raw);
1185 } else {
1186 debug_warn!("skipping malformed queued to-device event");
1187 }
1188 },
1189 | SendingEvent::DeviceListChanged(buf) => {
1190 if msc3202
1191 && let Some(bytes) = buf.get(TAG_PREFIX_LEN..)
1192 && let Ok(user) = from_utf8(bytes)
1193 && let Ok(user_id) = UserId::parse(user)
1194 {
1195 changed.push(user_id);
1196 }
1197 },
1198 | SendingEvent::BadgeRefresh | SendingEvent::Flush => {},
1199 }
1200 }
1201
1202 let txn_hash = calculate_hash(events.iter().filter_map(|e| match e {
1203 | SendingEvent::Edu(b)
1204 | SendingEvent::ToDevice(b)
1205 | SendingEvent::DeviceListChanged(b) => Some(b.as_ref()),
1206 | SendingEvent::Pdu(b) => Some(b.as_ref()),
1207 | SendingEvent::BadgeRefresh | SendingEvent::Flush => None,
1208 }));
1209
1210 let txn_id = &*URL_SAFE_NO_PAD.encode(txn_hash);
1211
1212 let (device_lists, device_one_time_keys_count, device_unused_fallback_key_types) =
1213 if msc3202 {
1214 changed.sort_unstable();
1215 changed.dedup();
1216
1217 let (counts, fallbacks) = self
1218 .msc3202_key_counts(otk_users, otk_recipients)
1219 .await;
1220
1221 (DeviceLists { changed, left: Vec::new() }, counts, fallbacks)
1222 } else {
1223 (DeviceLists::new(), OtkCounts::new(), FallbackTypes::new())
1224 };
1225
1226 if pdu_jsons.is_empty()
1227 && edu_jsons.is_empty()
1228 && to_device.is_empty()
1229 && device_lists.is_empty()
1230 && device_one_time_keys_count.is_empty()
1231 && device_unused_fallback_key_types.is_empty()
1232 {
1233 return Ok(Destination::Appservice(id));
1234 }
1235
1236 match self
1237 .services
1238 .appservice
1239 .send_request(info.registration, PushEventsRequest {
1240 txn_id: txn_id.into(),
1241 events: pdu_jsons,
1242 ephemeral: edu_jsons,
1243 to_device,
1244 device_lists,
1245 device_one_time_keys_count,
1246 device_unused_fallback_key_types,
1247 })
1248 .await
1249 {
1250 | Ok(_) => Ok(Destination::Appservice(id)),
1251 | Err(e) => Err((Destination::Appservice(id), e)),
1252 }
1253 }
1254
1255 async fn msc3202_key_counts(
1259 &self,
1260 users: BTreeSet<OwnedUserId>,
1261 recipients: BTreeSet<(OwnedUserId, OwnedDeviceId)>,
1262 ) -> (OtkCounts, FallbackTypes) {
1263 let mut devices: Devices = users
1264 .into_iter()
1265 .stream()
1266 .broad_then(async |user_id: OwnedUserId| {
1267 self.services
1268 .users
1269 .all_device_ids(&user_id)
1270 .map(|device_id| (user_id.clone(), device_id.to_owned()))
1271 .collect()
1272 .await
1273 })
1274 .flat_map(|pairs: Vec<(OwnedUserId, OwnedDeviceId)>| pairs.into_iter().stream())
1275 .chain(recipients.into_iter().stream())
1276 .collect()
1277 .await;
1278
1279 devices.sort_unstable();
1280 devices.dedup();
1281
1282 devices
1283 .into_iter()
1284 .stream()
1285 .broad_then(async |(user_id, device_id): (OwnedUserId, OwnedDeviceId)| {
1286 let counts = self
1287 .services
1288 .users
1289 .count_one_time_keys(&user_id, &device_id);
1290
1291 let fallbacks = self
1292 .services
1293 .users
1294 .unused_fallback_key_algorithms(&user_id, &device_id)
1295 .collect();
1296
1297 let (counts, fallbacks) = join(counts, fallbacks).await;
1298
1299 (user_id, device_id, counts, fallbacks)
1300 })
1301 .ready_fold(
1302 (OtkCounts::new(), FallbackTypes::new()),
1303 |(mut counts, mut fallbacks), (user_id, device_id, otk, fallback)| {
1304 counts
1305 .entry(user_id.clone())
1306 .or_default()
1307 .insert(device_id.clone(), otk);
1308
1309 fallbacks
1310 .entry(user_id)
1311 .or_default()
1312 .insert(device_id, fallback);
1313
1314 (counts, fallbacks)
1315 },
1316 )
1317 .await
1318 }
1319
1320 #[tracing::instrument(
1321 name = "push",
1322 level = "info",
1323 skip(self, events),
1324 fields(
1325 events = %events.len(),
1326 ),
1327 )]
1328 async fn send_events_dest_push(
1329 &self,
1330 user_id: OwnedUserId,
1331 pushkey: String,
1332 events: Vec<SendingEvent>,
1333 ) -> SendingResult {
1334 let has_pdu = events
1335 .iter()
1336 .any(|event| matches!(event, SendingEvent::Pdu(_)));
1337
1338 let suppressed = self.pushing_suppressed(&user_id).map(Ok);
1339 let pusher = self
1340 .services
1341 .pusher
1342 .get_pusher(&user_id, &pushkey)
1343 .map_err(|_| {
1344 (
1345 Destination::Push(user_id.clone(), pushkey.clone()),
1346 err!(Database(error!(?user_id, ?pushkey, "Missing pusher"))),
1347 )
1348 });
1349
1350 let rules_for_user = has_pdu
1351 .then_async(async || {
1352 self.services
1353 .account_data
1354 .get_global::<PushRulesEvent>(&user_id, GlobalAccountDataEventType::PushRules)
1355 .await
1356 .map_or_else(|_| Ruleset::server_default(&user_id), |ev| ev.content.global)
1357 })
1358 .map(Ok);
1359
1360 let (pusher, rules_for_user, suppressed) =
1361 try_join3(pusher, rules_for_user, suppressed).await?;
1362
1363 if events.contains(&SendingEvent::BadgeRefresh) {
1365 self.services
1366 .pusher
1367 .send_badge_notice(&user_id, &pusher)
1368 .map_err(|e| (Destination::Push(user_id.clone(), pushkey.clone()), e))
1369 .await?;
1370 }
1371
1372 if suppressed {
1373 let queued = self
1374 .enqueue_suppressed_push_events(&user_id, &pushkey, &events)
1375 .await;
1376
1377 debug!(
1378 ?user_id,
1379 pushkey,
1380 queued,
1381 events = events.len(),
1382 "Push suppressed; queued events"
1383 );
1384 return Ok(Destination::Push(user_id, pushkey));
1385 }
1386
1387 self.schedule_flush_suppressed_for_pushkey(
1388 user_id.clone(),
1389 pushkey.clone(),
1390 "non-suppressed push",
1391 );
1392
1393 if let Some(rules_for_user) = rules_for_user {
1394 let _sent = events
1395 .iter()
1396 .stream()
1397 .ready_filter_map(|event| extract_variant!(event, SendingEvent::Pdu))
1398 .wide_filter_map(|pdu_id| {
1399 self.services
1400 .timeline
1401 .get_pdu_from_id(pdu_id)
1402 .ok()
1403 })
1404 .ready_filter(|pdu| !pdu.is_redacted())
1405 .wide_filter_map(async |pdu| {
1406 self.services
1407 .pusher
1408 .send_push_notice(&user_id, &pusher, &rules_for_user, &pdu)
1409 .await
1410 .map_err(|e| (Destination::Push(user_id.clone(), pushkey.clone()), e))
1411 .ok()
1412 })
1413 .count()
1414 .await;
1415 }
1416
1417 Ok(Destination::Push(user_id, pushkey))
1418 }
1419
1420 pub fn schedule_flush_suppressed_for_pushkey(
1425 &self,
1426 user_id: OwnedUserId,
1427 pushkey: String,
1428 reason: &'static str,
1429 ) {
1430 let sending = self.services.sending.clone();
1431
1432 self.spawn_flush(async move {
1433 sending
1434 .flush_suppressed_for_pushkey(user_id, pushkey, reason)
1435 .await;
1436 });
1437 }
1438
1439 pub fn schedule_flush_suppressed_for_user(&self, user_id: OwnedUserId, reason: &'static str) {
1444 let sending = self.services.sending.clone();
1445
1446 self.spawn_flush(async move {
1447 sending
1448 .flush_suppressed_for_user(user_id, reason)
1449 .await;
1450 });
1451 }
1452
1453 fn spawn_flush<F>(&self, flush: F)
1454 where
1455 F: Future<Output = ()> + Send + 'static,
1456 {
1457 if !self.server.is_running() {
1459 return;
1460 }
1461
1462 let mut flushes = self.flushes.lock().expect("locked");
1463
1464 reap_flushes(&mut flushes);
1465 let _abort = flushes.spawn_on(flush, self.server.runtime());
1466 }
1467
1468 async fn enqueue_suppressed_push_events(
1469 &self,
1470 user_id: &UserId,
1471 pushkey: &str,
1472 events: &[SendingEvent],
1473 ) -> usize {
1474 let mut queued = 0_usize;
1475 for event in events {
1476 let SendingEvent::Pdu(pdu_id) = event else {
1477 continue;
1478 };
1479
1480 let Ok(pdu) = self
1481 .services
1482 .timeline
1483 .get_pdu_from_id(pdu_id)
1484 .await
1485 else {
1486 debug!(?user_id, ?pdu_id, "Suppressing push but PDU is missing");
1487 continue;
1488 };
1489
1490 if pdu.is_redacted() {
1491 trace!(?user_id, ?pdu_id, "Suppressing push for redacted PDU");
1492 continue;
1493 }
1494
1495 if self.services.pusher.queue_suppressed_push(
1496 user_id,
1497 pushkey,
1498 pdu.room_id(),
1499 *pdu_id,
1500 ) {
1501 queued = queued.saturating_add(1);
1502 }
1503 }
1504
1505 queued
1506 }
1507
1508 async fn flush_suppressed_rooms(
1509 &self,
1510 user_id: &UserId,
1511 pushkey: &str,
1512 pusher: &Pusher,
1513 rules_for_user: &Ruleset,
1514 rooms: Vec<(OwnedRoomId, Vec<RawPduId>)>,
1515 reason: &'static str,
1516 ) {
1517 if rooms.is_empty() {
1518 return;
1519 }
1520
1521 let mut sent = 0_usize;
1522 debug!(?user_id, pushkey, rooms = rooms.len(), "Flushing suppressed pushes ({reason})");
1523
1524 for (room_id, pdu_ids) in rooms {
1525 let unread = self
1526 .services
1527 .pusher
1528 .notification_count(user_id, &room_id)
1529 .await;
1530
1531 if unread == 0 {
1532 trace!(?user_id, ?room_id, "Skipping suppressed push flush: no unread");
1533 continue;
1534 }
1535
1536 for pdu_id in pdu_ids {
1537 let Ok(pdu) = self
1538 .services
1539 .timeline
1540 .get_pdu_from_id(&pdu_id)
1541 .await
1542 else {
1543 debug!(?user_id, ?pdu_id, "Suppressed PDU missing during flush");
1544 continue;
1545 };
1546
1547 if pdu.is_redacted() {
1548 trace!(?user_id, ?pdu_id, "Suppressed PDU redacted during flush");
1549 continue;
1550 }
1551
1552 if let Err(error) = self
1553 .services
1554 .pusher
1555 .send_push_notice(user_id, pusher, rules_for_user, &pdu)
1556 .await
1557 {
1558 let requeued = self
1559 .services
1560 .pusher
1561 .queue_suppressed_push(user_id, pushkey, &room_id, pdu_id);
1562
1563 warn!(
1564 ?user_id,
1565 ?room_id,
1566 ?error,
1567 requeued,
1568 "Failed to send suppressed push notification"
1569 );
1570 } else {
1571 sent = sent.saturating_add(1);
1572 }
1573 }
1574 }
1575
1576 debug!(?user_id, pushkey, sent, "Flushed suppressed push notifications");
1577 }
1578
1579 async fn flush_suppressed_for_pushkey(
1580 &self,
1581 user_id: OwnedUserId,
1582 pushkey: String,
1583 reason: &'static str,
1584 ) {
1585 let suppressed = self
1586 .services
1587 .pusher
1588 .take_suppressed_for_pushkey(&user_id, &pushkey);
1589
1590 if suppressed.is_empty() {
1591 return;
1592 }
1593
1594 let pusher = match self
1595 .services
1596 .pusher
1597 .get_pusher(&user_id, &pushkey)
1598 .await
1599 {
1600 | Ok(pusher) => pusher,
1601 | Err(error) => {
1602 warn!(?user_id, pushkey, ?error, "Missing pusher for suppressed flush");
1603 return;
1604 },
1605 };
1606
1607 let rules_for_user = match self
1608 .services
1609 .account_data
1610 .get_global::<PushRulesEvent>(&user_id, GlobalAccountDataEventType::PushRules)
1611 .await
1612 {
1613 | Ok(ev) => ev.content.global,
1614 | Err(_) => Ruleset::server_default(&user_id),
1615 };
1616
1617 self.flush_suppressed_rooms(
1618 &user_id,
1619 &pushkey,
1620 &pusher,
1621 &rules_for_user,
1622 suppressed,
1623 reason,
1624 )
1625 .await;
1626 }
1627
1628 pub async fn flush_suppressed_for_user(&self, user_id: OwnedUserId, reason: &'static str) {
1629 let suppressed = self
1630 .services
1631 .pusher
1632 .take_suppressed_for_user(&user_id);
1633
1634 if suppressed.is_empty() {
1635 return;
1636 }
1637
1638 let rules_for_user = match self
1639 .services
1640 .account_data
1641 .get_global::<PushRulesEvent>(&user_id, GlobalAccountDataEventType::PushRules)
1642 .await
1643 {
1644 | Ok(ev) => ev.content.global,
1645 | Err(_) => Ruleset::server_default(&user_id),
1646 };
1647
1648 for (pushkey, rooms) in suppressed {
1649 let pusher = match self
1650 .services
1651 .pusher
1652 .get_pusher(&user_id, &pushkey)
1653 .await
1654 {
1655 | Ok(pusher) => pusher,
1656 | Err(error) => {
1657 warn!(?user_id, pushkey, ?error, "Missing pusher for suppressed flush");
1658 continue;
1659 },
1660 };
1661
1662 self.flush_suppressed_rooms(
1663 &user_id,
1664 &pushkey,
1665 &pusher,
1666 &rules_for_user,
1667 rooms,
1668 reason,
1669 )
1670 .await;
1671 }
1672 }
1673
1674 async fn pushing_suppressed(&self, user_id: &UserId) -> bool {
1677 if !self.services.config.suppress_push_when_active {
1678 debug!(?user_id, "push not suppressed: suppress_push_when_active disabled");
1679 return false;
1680 }
1681
1682 let Ok(presence) = self.services.presence.get_presence(user_id).await else {
1683 debug!(?user_id, "push not suppressed: presence unavailable");
1684 return false;
1685 };
1686
1687 if presence.content.presence != PresenceState::Online {
1688 debug!(
1689 ?user_id,
1690 presence = ?presence.content.presence,
1691 "push not suppressed: presence not online"
1692 );
1693 return false;
1694 }
1695
1696 let presence_age_ms = presence
1697 .content
1698 .last_active_ago
1699 .map(u64::from)
1700 .unwrap_or(u64::MAX);
1701
1702 if presence_age_ms >= 65_000 {
1703 debug!(?user_id, presence_age_ms, "push not suppressed: presence too old");
1704 return false;
1705 }
1706
1707 let sync_gap_ms = self
1708 .services
1709 .presence
1710 .last_sync_gap_ms(user_id)
1711 .await;
1712
1713 let considered_active = sync_gap_ms.is_some_and(|gap| gap < 32_000);
1714
1715 match sync_gap_ms {
1716 | Some(gap) if gap < 32_000 => debug!(
1717 ?user_id,
1718 presence_age_ms,
1719 sync_gap_ms = gap,
1720 "suppressing push: active heuristic"
1721 ),
1722 | Some(gap) => debug!(
1723 ?user_id,
1724 presence_age_ms,
1725 sync_gap_ms = gap,
1726 "push not suppressed: sync gap too large"
1727 ),
1728 | None => debug!(?user_id, presence_age_ms, "push not suppressed: no recent sync"),
1729 }
1730
1731 considered_active
1732 }
1733
1734 async fn send_events_dest_federation(
1735 &self,
1736 server: OwnedServerName,
1737 events: Vec<SendingEvent>,
1738 ) -> SendingResult {
1739 let pdus: Vec<_> = events
1740 .iter()
1741 .filter_map(|event| extract_variant!(event, SendingEvent::Pdu))
1742 .stream()
1743 .wide_filter_map(|pdu_id| {
1744 self.services
1745 .timeline
1746 .get_pdu_json_from_id(pdu_id)
1747 .ok()
1748 })
1749 .wide_then(|pdu| {
1750 self.services
1751 .state_accessor
1752 .erased_for_server(&server, pdu)
1753 })
1754 .wide_then(|pdu| {
1755 self.services
1756 .federation
1757 .format_pdu_into(pdu, None)
1758 })
1759 .collect()
1760 .await;
1761
1762 let edus: Vec<Raw<Edu>> = events
1763 .iter()
1764 .filter_map(|edu| match edu {
1765 | SendingEvent::Edu(edu) => Some(edu.as_ref()),
1766 | _ => None,
1767 })
1768 .map(serde_json::from_slice)
1769 .filter_map(Result::ok)
1770 .collect();
1771
1772 if pdus.is_empty() && edus.is_empty() {
1773 return Ok(Destination::Federation(server));
1774 }
1775
1776 let preimage = pdus
1777 .iter()
1778 .map(|raw| raw.get().as_bytes())
1779 .chain(edus.iter().map(|raw| raw.json().get().as_bytes()));
1780
1781 let txn_hash = calculate_hash(preimage);
1782 let txn_id = &*URL_SAFE_NO_PAD.encode(txn_hash);
1783 let request = send_transaction_message::v1::Request {
1784 transaction_id: txn_id.into(),
1785 origin: self.server.name.clone(),
1786 origin_server_ts: MilliSecondsSinceUnixEpoch::now(),
1787 pdus,
1788 edus,
1789 };
1790
1791 let result = self
1792 .services
1793 .federation
1794 .execute_on(&self.services.client.sender, &server, request)
1795 .await;
1796
1797 for (event_id, result) in result.iter().flat_map(|resp| resp.pdus.iter()) {
1798 if let Err(e) = result {
1799 warn!(
1800 %txn_id, %server,
1801 "error sending PDU {event_id} to remote server: {e:?}"
1802 );
1803 }
1804 }
1805
1806 match result {
1807 | Ok(_) => Ok(Destination::Federation(server)),
1808 | Err(error) => Err((Destination::Federation(server), error)),
1809 }
1810 }
1811}
1812
1813fn arm_wake(wakes: &mut WakeQueue, server: OwnedServerName, earliest_retry: SystemTime) {
1814 use tokio::time::Instant;
1815
1816 let delay = earliest_retry
1818 .duration_since(SystemTime::now())
1819 .unwrap_or_default()
1820 .max(Duration::from_secs(1));
1821
1822 let jitter = rand_secs(0..delay.as_secs().max(3));
1825 let deadline = Instant::now()
1826 .checked_add(delay.saturating_add(jitter))
1827 .unwrap_or_else(Instant::now);
1828
1829 wakes.push(Reverse((deadline, server)));
1830}