Skip to main content

tuwunel_service/sending/sender/dispatch/push/
suppressed.rs

1use futures::StreamExt;
2use ruma::{
3	OwnedRoomId, OwnedUserId, RoomId, UserId, api::client::push::Pusher, presence::PresenceState,
4	push::Ruleset,
5};
6use tuwunel_core::{
7	Event, debug, extract_variant, implement,
8	matrix::Pdu,
9	trace,
10	utils::{
11		IterStream, ReadyExt,
12		stream::{BroadbandExt, WidebandExt},
13	},
14	warn,
15};
16
17use super::PUSH_WIDTH;
18use crate::{
19	pusher::SuppressedRooms,
20	rooms::timeline::RawPduId,
21	sending::{SendingEvent, Service},
22};
23
24/// The pusher and ruleset one suppressed-push flush delivers with.
25///
26/// Every room and PDU of the flush is sent under the same values.
27struct Flush<'a> {
28	user_id: &'a UserId,
29	pushkey: &'a str,
30	pusher: &'a Pusher,
31	ruleset: &'a Ruleset,
32	reason: &'static str,
33}
34
35// One presence heartbeat and one sync long-poll respectively, plus margin.
36const ACTIVE_PRESENCE_AGE_MS: u64 = 65_000;
37const ACTIVE_SYNC_GAP_MS: u64 = 32_000;
38
39/// Schedule a flush of the pushes suppressed for one pushkey.
40///
41/// The flush runs as a task this service owns, so the caller never waits on
42/// the push gateway.
43#[implement(Service)]
44pub fn schedule_flush_suppressed_for_pushkey(
45	&self,
46	user_id: OwnedUserId,
47	pushkey: String,
48	reason: &'static str,
49) {
50	let sending = self.services.sending.clone();
51
52	self.spawn_flush(async move {
53		sending
54			.flush_suppressed_for_pushkey(&user_id, &pushkey, reason)
55			.await;
56	});
57}
58
59/// Schedule a flush of the pushes suppressed for every pushkey a user owns.
60///
61/// The flush runs as a task this service owns, so the caller never waits on
62/// the push gateway.
63#[implement(Service)]
64pub fn schedule_flush_suppressed_for_user(&self, user_id: OwnedUserId, reason: &'static str) {
65	let sending = self.services.sending.clone();
66
67	self.spawn_flush(async move {
68		sending
69			.flush_suppressed_for_user(&user_id, reason)
70			.await;
71	});
72}
73
74#[implement(Service)]
75async fn flush_suppressed_for_pushkey(
76	&self,
77	user_id: &UserId,
78	pushkey: &str,
79	reason: &'static str,
80) {
81	let suppressed = self
82		.services
83		.pusher
84		.take_suppressed_for_pushkey(user_id, pushkey);
85
86	if suppressed.is_empty() {
87		return;
88	}
89
90	let Ok(pusher) = self
91		.services
92		.pusher
93		.get_pusher(user_id, pushkey)
94		.await
95		.inspect_err(|error| {
96			warn!(?user_id, pushkey, ?error, "Missing pusher for suppressed flush");
97		})
98	else {
99		return;
100	};
101
102	let ruleset = self.services.pusher.ruleset(user_id).await;
103	let flush = Flush {
104		user_id,
105		pushkey,
106		pusher: &pusher,
107		ruleset: &ruleset,
108		reason,
109	};
110
111	self.flush_suppressed_rooms(&flush, suppressed)
112		.await;
113}
114
115#[implement(Service)]
116async fn flush_suppressed_for_user(&self, user_id: &UserId, reason: &'static str) {
117	let suppressed = self
118		.services
119		.pusher
120		.take_suppressed_for_user(user_id);
121
122	if suppressed.is_empty() {
123		return;
124	}
125
126	let ruleset = self.services.pusher.ruleset(user_id).await;
127
128	for (pushkey, rooms) in suppressed {
129		let Ok(pusher) = self
130			.services
131			.pusher
132			.get_pusher(user_id, &pushkey)
133			.await
134			.inspect_err(|error| {
135				warn!(?user_id, pushkey, ?error, "Missing pusher for suppressed flush");
136			})
137		else {
138			continue;
139		};
140
141		let flush = Flush {
142			user_id,
143			pushkey: &pushkey,
144			pusher: &pusher,
145			ruleset: &ruleset,
146			reason,
147		};
148
149		self.flush_suppressed_rooms(&flush, rooms).await;
150	}
151}
152
153#[implement(Service)]
154async fn flush_suppressed_rooms(&self, flush: &Flush<'_>, rooms: SuppressedRooms) {
155	if rooms.is_empty() {
156		return;
157	}
158
159	let Flush { user_id, pushkey, reason, .. } = *flush;
160
161	debug!(?user_id, pushkey, rooms = rooms.len(), reason, "Flushing suppressed pushes");
162	let sent = rooms
163		.into_iter()
164		.stream()
165		.then(|(room_id, pdu_ids)| self.flush_suppressed_room(flush, room_id, pdu_ids))
166		.ready_fold(0, usize::saturating_add)
167		.await;
168
169	debug!(?user_id, pushkey, sent, "Flushed suppressed push notifications");
170}
171
172#[implement(Service)]
173async fn flush_suppressed_room(
174	&self,
175	flush: &Flush<'_>,
176	room_id: OwnedRoomId,
177	pdu_ids: Vec<RawPduId>,
178) -> usize {
179	let unread = self
180		.services
181		.pusher
182		.notification_count(flush.user_id, &room_id)
183		.await;
184
185	if unread == 0 {
186		trace!(user_id = ?flush.user_id, ?room_id, "Skipping suppressed push flush: no unread");
187		return 0;
188	}
189
190	pdu_ids
191		.into_iter()
192		.stream()
193		.wide_filter_map(async |pdu_id| {
194			self.suppressable_pdu(flush.user_id, &pdu_id)
195				.await
196				.map(|pdu| (pdu_id, pdu))
197		})
198		.broadn_then(Some(PUSH_WIDTH), async |(pdu_id, pdu)| {
199			self.flush_suppressed_pdu(flush, &room_id, pdu_id, &pdu)
200				.await
201		})
202		.ready_filter(|&sent| sent)
203		.count()
204		.await
205}
206
207#[implement(Service)]
208async fn flush_suppressed_pdu(
209	&self,
210	flush: &Flush<'_>,
211	room_id: &RoomId,
212	pdu_id: RawPduId,
213	pdu: &Pdu,
214) -> bool {
215	let Flush { user_id, pushkey, pusher, ruleset, .. } = *flush;
216
217	let Err(error) = self
218		.services
219		.pusher
220		.send_push_notice(user_id, pusher, ruleset, pdu)
221		.await
222	else {
223		return true;
224	};
225
226	let requeued = self
227		.services
228		.pusher
229		.queue_suppressed_push(user_id, pushkey, room_id, pdu_id);
230
231	warn!(
232		?user_id,
233		?room_id,
234		?error,
235		requeued,
236		"Failed to send suppressed push notification"
237	);
238
239	false
240}
241
242#[implement(Service)]
243pub(super) async fn enqueue_suppressed_push_events(
244	&self,
245	user_id: &UserId,
246	pushkey: &str,
247	events: &[SendingEvent],
248) -> usize {
249	events
250		.iter()
251		.stream()
252		.ready_filter_map(|event| extract_variant!(event, SendingEvent::Pdu))
253		.wide_filter_map(async |pdu_id| {
254			self.suppressable_pdu(user_id, pdu_id)
255				.await
256				.map(|pdu| (pdu_id, pdu))
257		})
258		.ready_fold(0, |queued: usize, (pdu_id, pdu)| {
259			let accepted = self.services.pusher.queue_suppressed_push(
260				user_id,
261				pushkey,
262				pdu.room_id(),
263				*pdu_id,
264			);
265
266			queued.saturating_add(usize::from(accepted))
267		})
268		.await
269}
270
271/// Load a suppressed PDU if a push for it is still worth sending.
272///
273/// A missing or redacted PDU is dropped with a log line and never notified.
274#[implement(Service)]
275async fn suppressable_pdu(&self, user_id: &UserId, pdu_id: &RawPduId) -> Option<Pdu> {
276	let Ok(pdu) = self
277		.services
278		.timeline
279		.get_pdu_from_id(pdu_id)
280		.await
281	else {
282		debug!(?user_id, ?pdu_id, "Suppressed PDU is missing");
283		return None;
284	};
285
286	if pdu.is_redacted() {
287		trace!(?user_id, ?pdu_id, "Suppressed PDU is redacted");
288		return None;
289	}
290
291	Some(pdu)
292}
293
294/// Decide whether pushes for a user are suppressed as active.
295///
296/// The heuristic combines the presence age and the most recent sync gap, and
297/// only applies when `suppress_push_when_active` is enabled. An offline user
298/// is never suppressed; the two `ACTIVE_*` constants are the thresholds.
299#[implement(Service)]
300pub(super) async fn pushing_suppressed(&self, user_id: &UserId) -> bool {
301	if !self.services.config.suppress_push_when_active {
302		debug!(?user_id, "push not suppressed: suppress_push_when_active disabled");
303		return false;
304	}
305
306	let Ok(presence) = self.services.presence.get_presence(user_id).await else {
307		debug!(?user_id, "push not suppressed: presence unavailable");
308		return false;
309	};
310
311	if presence.content.presence != PresenceState::Online {
312		debug!(
313			?user_id,
314			presence = ?presence.content.presence,
315			"push not suppressed: presence not online"
316		);
317
318		return false;
319	}
320
321	let presence_age_ms = presence
322		.content
323		.last_active_ago
324		.map(u64::from)
325		.unwrap_or(u64::MAX);
326
327	if presence_age_ms >= ACTIVE_PRESENCE_AGE_MS {
328		debug!(?user_id, presence_age_ms, "push not suppressed: presence too old");
329		return false;
330	}
331
332	let sync_gap_ms = self
333		.services
334		.presence
335		.last_sync_gap_ms(user_id)
336		.await;
337
338	match sync_gap_ms {
339		| Some(gap) if gap < ACTIVE_SYNC_GAP_MS => {
340			debug!(
341				?user_id,
342				presence_age_ms,
343				sync_gap_ms = gap,
344				"suppressing push: active heuristic"
345			);
346
347			true
348		},
349		| Some(gap) => {
350			debug!(
351				?user_id,
352				presence_age_ms,
353				sync_gap_ms = gap,
354				"push not suppressed: sync gap too large"
355			);
356
357			false
358		},
359		| None => {
360			debug!(?user_id, presence_age_ms, "push not suppressed: no recent sync");
361
362			false
363		},
364	}
365}