tuwunel_service/pusher/
suppressed.rs1use 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
20pub type SuppressedRooms = Vec<(OwnedRoomId, Vec<RawPduId>)>;
26
27pub 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#[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#[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#[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#[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#[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}