tuwunel_service/sending/sender/dispatch/push/
suppressed.rs1use 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
24struct 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
35const ACTIVE_PRESENCE_AGE_MS: u64 = 65_000;
37const ACTIVE_SYNC_GAP_MS: u64 = 32_000;
38
39#[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#[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#[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#[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}