Skip to main content

tuwunel_service/sending/
sender.rs

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/// In-flight bookkeeping for one `Destination`. Cross-attempt backoff lives
69/// in `peer_status` (federation only); appservice/push paths keep their own
70/// status because they are not server-keyed.
71#[derive(Debug)]
72enum TransactionStatus {
73	Running,
74	RunningForceRetry,
75	Failed(u32, Instant), // push backoff: tries, last failure
76	Retrying(u32),        // number of times failed
77}
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
90/// MSC3202 `device_one_time_keys_count`: unclaimed one-time-key counts per
91/// algorithm, keyed by user then device. Matches the ruma request field type.
92type OtkCounts =
93	BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, BTreeMap<OneTimeKeyAlgorithm, UInt>>>;
94
95/// MSC3202 `device_unused_fallback_key_types`: algorithms with an unused
96/// fallback key, keyed by user then device.
97type FallbackTypes = BTreeMap<OwnedUserId, BTreeMap<OwnedDeviceId, Vec<OneTimeKeyAlgorithm>>>;
98
99/// The MSC3202-interesting devices of one transaction: the appservice
100/// sender's plus matched PDU senders' devices and the to-device recipients.
101type Devices = SmallVec<[(OwnedUserId, OwnedDeviceId); 1]>;
102
103/// Per-worker retry timer keyed by earliest-retry deadline (tokio time). Every
104/// recorded federation failure arms one entry; heap size is bounded by the
105/// concurrently-sad destinations plus transient re-arms.
106type WakeQueue = BinaryHeap<Reverse<(tokio::time::Instant, OwnedServerName)>>;
107
108/// Per-(room, user) bucket of `ReceiptData`. MSC3771 allows one receipt
109/// per thread context per user per EDU window; the dominant case is
110/// still a single receipt, so inline-1 fits without a heap touch.
111type UserReceipts = SmallVec<[ReceiptData; 1]>;
112
113/// Per-rank slice of receipt EDU output. Each entry becomes one
114/// `Edu::Receipt` buffer; rank 0 carries each user's earliest receipt
115/// in the window, rank 1 the next, and so on. Most windows produce a
116/// single rank.
117type RankedReceipts = SmallVec<[ReceiptMap; 1]>;
118
119/// Per-room ranked receipts gathered for one federation EDU window. The
120/// common case is a single room, so inline-1 avoids a heap touch.
121type RoomReceipts = SmallVec<[(OwnedRoomId, RankedReceipts); 1]>;
122
123/// Output of one EDU selector. `shipped` rides the current transaction up to
124/// the shared budget; `overflow` past the budget is written as queued rows for
125/// later transactions to drain.
126#[derive(Default)]
127struct Selected {
128	shipped: EduVec,
129	overflow: Vec<EduBuf>,
130}
131
132/// The appservice-injected recipient fields of a queued to-device event
133/// (MSC4203), parsed to scope MSC3202 one-time-key counts to the addressed
134/// devices.
135#[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						// Arm a one-shot retry at the destination's earliest-retry time.
231						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		// Push backs off locally; federation defers to peer_status, appservice retries.
272		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		// Find events that have been added since starting the last request
305		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		// Top up with EDUs that accrued while the transaction was in flight.
322		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		// Active transaction generations must own their queued successors before
511		// queued-only badge destinations are woken.
512		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>, // Events we want to send: event and full key
547		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		// Nothing can be done for this remote, bail out.
564		if !allow {
565			return Ok(None);
566		}
567
568		let mut events = Vec::new();
569
570		// Must retry any previous transaction for this remote.
571		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		// Compose the next transaction
581		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		// Add EDU's into the transaction
593		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		// peer_status gates federation only; appservice and push fall through.
614		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; // already running
632					if matches!(retry_action, RetryAction::Force) {
633						*e = TransactionStatus::RunningForceRetry;
634					}
635				},
636				| TransactionStatus::Failed(tries, time) => {
637					// Push backoff: hold off until the exponential window elapses.
638					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; // push retry already in flight
649				},
650				| TransactionStatus::Retrying(_) => {
651					// Promote to Running so a concurrent select does not double-send.
652					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		// selection window
664		let since = self.db.get_latest_educount(server_name).await;
665		let since_upper = self.services.globals.current_count();
666
667		// Nothing new since the last window: skip the scan and the watermark.
668		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		// Presence rides last and is excluded from the durable prefix because
705		// its content is compose-time-relative and regenerates fresh.
706		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		// EDUs past the budget become queued rows drained by later transactions.
715		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		// Persist the durable prefix so a failed or restarted transaction
729		// replays it; the ACK deletes these active rows.
730		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	/// Look for device changes
745	#[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				// Empty prev id forces synapse to resync; because synapse resyncs,
782				// we can just insert placeholder data
783				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				// Past the budget these rows overflow to the queue; replay is
798				// benign because the placeholder content is user-id-only.
799				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	/// Look for read receipts in this room
813	///
814	/// MSC3771 lets a user emit multiple receipts in the same EDU window, one
815	/// per thread context. The federation EDU shape allows only one
816	/// `ReceiptData` per `(room, user)` slot, so a user with N parallel
817	/// thread receipts ships across N parallel `Edu::Receipt` buffers within
818	/// the same transaction. Each buffer is shape-compliant; receivers
819	/// process them as independent receipt EDUs and our storage keeps each
820	/// thread distinct.
821	#[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		// Ranks reserve from the shared budget in order; those past the cap
881		// overflow to the queue instead of truncating the tail.
882		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	/// Look for read receipts in this room.
897	///
898	/// Returns a per-rank vector of [`ReceiptMap`]s. Each user's receipts in
899	/// the window (one per thread context, count-ordered) are placed into
900	/// successive ranks, so rank 0 carries each user's earliest receipt,
901	/// rank 1 the next, and so on. The receipt-limit budget bounds distinct
902	/// users only; subsequent thread receipts for an already-counted user do
903	/// not consume additional budget.
904	#[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		// Pivot per-user count-ordered receipts into rank-major
971		// `RankedReceipts`. Rank 0 carries each user's earliest receipt in
972		// the window, rank 1 the next, and so on.
973		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	/// Look for presence
991	#[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		// A budget trip drops presence, which self-heals on the next transition.
1062		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			//TODO: appservice queue cleanup.
1112			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		// MSC3202 one-time-key scope: the appservice sender plus (below) the
1140		// namespace-matched PDU senders and to-device recipients of this txn.
1141		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	/// MSC3202 one-time-key counts and unused fallback key types over every
1256	/// device of `users` plus the specific `recipients`. Recomputed per build
1257	/// rather than snapshotted, so a retry ships fresh counts.
1258	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		// Reconciliation, not an alert: a suppressed drop strands a stale badge.
1364		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	/// Schedule a flush of the pushes suppressed for one pushkey.
1421	///
1422	/// The flush runs as a task this service owns, so the caller never waits on
1423	/// the push gateway.
1424	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	/// Schedule a flush of the pushes suppressed for every pushkey a user owns.
1440	///
1441	/// The flush runs as a task this service owns, so the caller never waits on
1442	/// the push gateway.
1443	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		// A flush scheduled during shutdown is dropped, not spawned.
1458		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	// optional suppression: heuristic combining presence age and recent sync
1675	// activity.
1676	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	// Floor the delay at 1s so clock steps and past deadlines wake promptly.
1817	let delay = earliest_retry
1818		.duration_since(SystemTime::now())
1819		.unwrap_or_default()
1820		.max(Duration::from_secs(1));
1821
1822	// Spread the wake over another delay-width (3s minimum), so destinations
1823	// sharing a backoff tier trickle back rather than retrying in one burst.
1824	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}