Skip to main content

tuwunel_service/pusher/
suppressed.rs

1//! Deferred push suppression queues.
2//!
3//! Stores suppressed push events in memory until they can be flushed. This is
4//! intentionally in-memory only: suppressed events are discarded on restart.
5
6use std::{
7	collections::{HashMap, VecDeque},
8	sync::Mutex,
9};
10
11use ruma::{OwnedRoomId, OwnedUserId, RoomId, UserId};
12use tuwunel_core::{debug, implement, trace, utils};
13
14use crate::rooms::timeline::RawPduId;
15
16const SUPPRESSED_MAX_EVENTS_PER_ROOM: usize = 512;
17const SUPPRESSED_MAX_EVENTS_PER_PUSHKEY: usize = 4096;
18const SUPPRESSED_MAX_ROOMS_PER_PUSHKEY: usize = 256;
19
20/// The PDUs suppressed for one pusher, grouped by room.
21///
22/// The PDUs of one room are listed in the order they were withheld in and the
23/// rooms themselves are unordered; the flush admits a room's PDUs in that
24/// order but sends a few at a time, so delivery order is not promised.
25pub type SuppressedRooms = Vec<(OwnedRoomId, Vec<RawPduId>)>;
26
27/// The suppressed PDUs of every pusher a user owns, keyed by pushkey.
28///
29/// This is the shape a whole-user flush drains.
30pub type SuppressedPushes = Vec<(String, SuppressedRooms)>;
31
32#[derive(Default)]
33pub(super) struct SuppressedQueue {
34	inner: Mutex<HashMap<OwnedUserId, HashMap<String, PushkeyQueue>>>,
35}
36
37#[derive(Default)]
38struct PushkeyQueue {
39	rooms: HashMap<OwnedRoomId, VecDeque<SuppressedEvent>>,
40	total_events: usize,
41}
42
43#[derive(Clone, Debug)]
44struct SuppressedEvent {
45	pdu_id: RawPduId,
46	_inserted_at_ms: u64,
47}
48
49impl SuppressedQueue {
50	fn lock(
51		&self,
52	) -> std::sync::MutexGuard<'_, HashMap<OwnedUserId, HashMap<String, PushkeyQueue>>> {
53		self.inner
54			.lock()
55			.unwrap_or_else(std::sync::PoisonError::into_inner)
56	}
57
58	fn drain_room(queue: VecDeque<SuppressedEvent>) -> Vec<RawPduId> {
59		queue
60			.into_iter()
61			.map(|event| event.pdu_id)
62			.collect()
63	}
64
65	fn drop_one_front(queue: &mut VecDeque<SuppressedEvent>, total_events: &mut usize) -> bool {
66		if queue.pop_front().is_some() {
67			*total_events = total_events.saturating_sub(1);
68			return true;
69		}
70
71		false
72	}
73}
74
75/// Enqueue a PDU for later push delivery when suppression is active.
76#[implement(super::Service)]
77pub fn queue_suppressed_push(
78	&self,
79	user_id: &UserId,
80	pushkey: &str,
81	room_id: &RoomId,
82	pdu_id: RawPduId,
83) -> bool {
84	let mut inner = self.suppressed.lock();
85	let user_entry = inner.entry(user_id.to_owned()).or_default();
86	let push_entry = user_entry.entry(pushkey.to_owned()).or_default();
87
88	if !push_entry.rooms.contains_key(room_id)
89		&& push_entry.rooms.len() >= SUPPRESSED_MAX_ROOMS_PER_PUSHKEY
90	{
91		debug!(
92			?user_id,
93			?room_id,
94			pushkey,
95			max_rooms = SUPPRESSED_MAX_ROOMS_PER_PUSHKEY,
96			"Suppressed push queue full (rooms); dropping event"
97		);
98		return false;
99	}
100
101	let queue = push_entry
102		.rooms
103		.entry(room_id.to_owned())
104		.or_default();
105
106	if queue
107		.back()
108		.is_some_and(|event| event.pdu_id == pdu_id)
109	{
110		trace!(?user_id, ?room_id, pushkey, "Suppressed push event is duplicate; skipping");
111		return false;
112	}
113
114	if push_entry.total_events >= SUPPRESSED_MAX_EVENTS_PER_PUSHKEY && queue.is_empty() {
115		debug!(
116			?user_id,
117			?room_id,
118			pushkey,
119			max_events = SUPPRESSED_MAX_EVENTS_PER_PUSHKEY,
120			"Suppressed push queue full (total); dropping event"
121		);
122		return false;
123	}
124
125	while queue.len() >= SUPPRESSED_MAX_EVENTS_PER_ROOM
126		|| push_entry.total_events >= SUPPRESSED_MAX_EVENTS_PER_PUSHKEY
127	{
128		if !SuppressedQueue::drop_one_front(queue, &mut push_entry.total_events) {
129			break;
130		}
131	}
132
133	queue.push_back(SuppressedEvent {
134		pdu_id,
135		_inserted_at_ms: utils::millis_since_unix_epoch(),
136	});
137	push_entry.total_events = push_entry.total_events.saturating_add(1);
138
139	true
140}
141
142/// Take and remove all suppressed PDUs for a given user + pushkey.
143#[implement(super::Service)]
144pub fn take_suppressed_for_pushkey(&self, user_id: &UserId, pushkey: &str) -> SuppressedRooms {
145	let mut inner = self.suppressed.lock();
146	let Some(user_entry) = inner.get_mut(user_id) else {
147		return Vec::new();
148	};
149
150	let Some(push_entry) = user_entry.remove(pushkey) else {
151		return Vec::new();
152	};
153
154	if user_entry.is_empty() {
155		inner.remove(user_id);
156	}
157
158	push_entry
159		.rooms
160		.into_iter()
161		.map(|(room_id, queue)| (room_id, SuppressedQueue::drain_room(queue)))
162		.collect()
163}
164
165/// Take and remove all suppressed PDUs for a given user across all pushkeys.
166#[implement(super::Service)]
167pub fn take_suppressed_for_user(&self, user_id: &UserId) -> SuppressedPushes {
168	let mut inner = self.suppressed.lock();
169	let Some(user_entry) = inner.remove(user_id) else {
170		return Vec::new();
171	};
172
173	user_entry
174		.into_iter()
175		.map(|(pushkey, queue)| {
176			let rooms = queue
177				.rooms
178				.into_iter()
179				.map(|(room_id, q)| (room_id, SuppressedQueue::drain_room(q)))
180				.collect();
181			(pushkey, rooms)
182		})
183		.collect()
184}
185
186/// Clear suppressed PDUs for a specific room (across all pushkeys).
187#[implement(super::Service)]
188pub fn clear_suppressed_room(&self, user_id: &UserId, room_id: &RoomId) -> usize {
189	let mut inner = self.suppressed.lock();
190	let Some(user_entry) = inner.get_mut(user_id) else {
191		return 0;
192	};
193
194	let mut removed: usize = 0;
195	user_entry.retain(|_, push_entry| {
196		if let Some(queue) = push_entry.rooms.remove(room_id) {
197			removed = removed.saturating_add(queue.len());
198			push_entry.total_events = push_entry
199				.total_events
200				.saturating_sub(queue.len());
201		}
202
203		!push_entry.rooms.is_empty()
204	});
205
206	if user_entry.is_empty() {
207		inner.remove(user_id);
208	}
209
210	removed
211}
212
213/// Clear suppressed PDUs for a specific pushkey.
214#[implement(super::Service)]
215pub fn clear_suppressed_pushkey(&self, user_id: &UserId, pushkey: &str) -> usize {
216	let mut inner = self.suppressed.lock();
217	let Some(user_entry) = inner.get_mut(user_id) else {
218		return 0;
219	};
220
221	let removed = user_entry
222		.remove(pushkey)
223		.map(|queue| queue.total_events)
224		.unwrap_or(0);
225
226	if user_entry.is_empty() {
227		inner.remove(user_id);
228	}
229
230	removed
231}