1use axum::extract::State;
2use futures::{FutureExt, StreamExt, TryFutureExt, pin_mut};
3use ruma::{
4 DeviceId, RoomId, UInt, UserId,
5 api::{
6 Direction,
7 client::{filter::RoomEventFilter, message::get_message_events},
8 },
9 events::{
10 AnyStateEvent, StateEventType, TimelineEventType, TimelineEventType::*,
11 relation::RelationType,
12 },
13 serde::Raw,
14};
15use tuwunel_core::{
16 Err, PduId, Result, at, err,
17 matrix::{
18 event::{Event, Matches},
19 pdu::{PduCount, PduEvent},
20 },
21 ref_at,
22 smallvec::SmallVec,
23 utils::{
24 BoolExt, IterStream, ReadyExt,
25 math::usize_from_ruma_bounded,
26 result::LogErr,
27 stream::{BroadbandExt, TryIgnore, WidebandExt},
28 },
29};
30use tuwunel_service::{
31 Services,
32 rooms::{
33 lazy_loading,
34 lazy_loading::{Options, Witness},
35 short::ShortRoomId,
36 timeline::PdusIterItem,
37 },
38};
39
40use super::visibility_filter;
41use crate::Ruma;
42
43pub(crate) struct MessagesArgs<'a> {
46 pub room_id: &'a RoomId,
47 pub sender_user: &'a UserId,
48 pub sender_device: Option<&'a DeviceId>,
49 pub from: Option<&'a str>,
50 pub to: Option<&'a str>,
51 pub dir: Direction,
52 pub limit: Option<UInt>,
53 pub filter: &'a RoomEventFilter,
54
55 pub bypass_visibility: bool,
58}
59
60const IGNORED_MESSAGE_TYPES: &[TimelineEventType] = &[
63 CallInvite, KeyVerificationStart, Location, PollStart, Reaction, RoomEncrypted, RoomMessage, Sticker, Audio, Emote, File, Image, Video, Voice, UnstablePollStart, Beacon, CallNotify, ];
81
82type RelTypes = SmallVec<[RelationType; 1]>;
84
85const LIMIT_MAX: usize = 1000;
86const LIMIT_DEFAULT: usize = 10;
87
88pub(crate) async fn get_message_events_route(
95 State(services): State<crate::State>,
96 body: Ruma<get_message_events::v3::Request>,
97) -> Result<get_message_events::v3::Response> {
98 get_messages(&services, MessagesArgs {
99 room_id: &body.room_id,
100 sender_user: body.sender_user(),
101 sender_device: body.sender_device.as_deref(),
102 from: body.from.as_deref(),
103 to: body.to.as_deref(),
104 dir: body.dir,
105 limit: Some(body.limit),
106 filter: &body.filter,
107 bypass_visibility: false,
108 })
109 .await
110}
111
112pub(crate) async fn get_messages(
116 services: &Services,
117 args: MessagesArgs<'_>,
118) -> Result<get_message_events::v3::Response> {
119 let MessagesArgs {
120 room_id,
121 sender_user,
122 sender_device,
123 from,
124 to,
125 dir,
126 limit,
127 filter,
128 bypass_visibility,
129 } = args;
130
131 if !services.metadata.exists(room_id).await {
132 return Err!(Request(Forbidden("Room does not exist to this server")));
133 }
134
135 if !bypass_visibility
136 && !services
137 .state_accessor
138 .user_can_see_room(sender_user, room_id)
139 .await
140 {
141 return Err!(Request(Forbidden("You don't have permission to view this room.")));
142 }
143
144 let from: PduCount = from
145 .map(str::parse)
146 .transpose()
147 .map_err(|_| err!(Request(InvalidParam("Invalid `from` token."))))?
148 .unwrap_or_else(|| match dir {
149 | Direction::Forward => PduCount::min(),
150 | Direction::Backward => PduCount::max(),
151 });
152
153 let to: Option<PduCount> = to
154 .map(str::parse)
155 .transpose()
156 .map_err(|_| err!(Request(InvalidParam("Invalid `to` token."))))?;
157
158 let limit = limit
159 .map_or(LIMIT_DEFAULT, |limit| usize_from_ruma_bounded(limit, LIMIT_DEFAULT, LIMIT_MAX));
160
161 if matches!(dir, Direction::Backward) {
162 services
163 .timeline
164 .backfill_if_required(room_id, from)
165 .await
166 .log_err()
167 .ok();
168 }
169
170 let it = match dir {
171 | Direction::Forward => services
172 .timeline
173 .pdus(Some(sender_user), room_id, Some(from))
174 .ignore_err()
175 .left_stream(),
176
177 | Direction::Backward => services
178 .timeline
179 .pdus_rev(Some(sender_user), room_id, Some(from))
180 .ignore_err()
181 .right_stream(),
182 };
183
184 let encrypted = services
185 .state_accessor
186 .is_encrypted_room(room_id)
187 .await;
188
189 let shortroomid = services.short.get_shortroomid(room_id).await?;
190 let mut scanned = None;
191 let reached_to = |count: PduCount| {
192 to.is_some_and(|to| match dir {
193 | Direction::Forward => count >= to,
194 | Direction::Backward => count <= to,
195 })
196 };
197
198 let events: Vec<_> = it
199 .inspect(|(count, _)| scanned = Some(*count))
200 .ready_take_while(|(count, _)| !reached_to(*count))
201 .ready_filter_map(|item| event_filter(item, filter))
202 .wide_filter_map(|item| related_by_filter(services, shortroomid, filter, item))
203 .wide_filter_map(|item| event_filters(services, sender_user, item, bypass_visibility))
204 .take(limit)
205 .wide_then(|item| add_membership_unsigned(services, item, sender_user, encrypted))
206 .wide_then(async |(count, pdu)| {
207 let pdu = services
208 .pdu_metadata
209 .bundle_aggregations(sender_user, pdu)
210 .await;
211
212 (count, pdu)
213 })
214 .collect()
215 .await;
216
217 let lazy_loading_context = lazy_loading::Context {
218 user_id: sender_user,
219 device_id: sender_device,
220 room_id,
221 token: Some(from.into_unsigned()),
222 options: Some(&filter.lazy_load_options),
223 mode: lazy_loading::Mode::Update,
224 };
225
226 let witness = filter
227 .lazy_load_options
228 .is_enabled()
229 .then_async(|| lazy_loading_witness(services, &lazy_loading_context, events.iter()));
230
231 let state = witness
232 .map(Option::into_iter)
233 .map(|option| option.flat_map(Witness::into_iter))
234 .map(IterStream::stream)
235 .into_stream()
236 .flatten()
237 .broad_filter_map(async |user_id| get_member_event(services, room_id, &user_id).await)
238 .collect()
239 .await;
240
241 let stopped_at_to = scanned.is_some_and(reached_to);
244 let exhausted = matches!(dir, Direction::Backward) && events.len() < limit && !stopped_at_to;
245 let next_token = if exhausted { scanned } else { events.last().map(at!(0)) };
246
247 let chunk = events
248 .into_iter()
249 .map(at!(1))
250 .map(Event::into_format)
251 .collect();
252
253 Ok(get_message_events::v3::Response {
254 start: from.to_string(),
255 end: next_token.as_ref().map(ToString::to_string),
256 chunk,
257 state,
258 })
259}
260
261pub(crate) async fn lazy_loading_witness<'a, I>(
262 services: &Services,
263 lazy_loading_context: &lazy_loading::Context<'_>,
264 events: I,
265) -> Witness
266where
267 I: Iterator<Item = &'a PdusIterItem> + Clone + Send,
268{
269 let oldest = events
270 .clone()
271 .map(|(count, _)| count)
272 .copied()
273 .min()
274 .unwrap_or_else(PduCount::max);
275
276 let newest = events
277 .clone()
278 .map(|(count, _)| count)
279 .copied()
280 .max()
281 .unwrap_or_else(PduCount::max);
282
283 let receipts = services.read_receipt.readreceipts_since(
284 lazy_loading_context.room_id,
285 oldest.into_unsigned(),
286 Some(newest.into_unsigned()),
287 );
288
289 pin_mut!(receipts);
290 let witness: Witness = events
291 .stream()
292 .map(ref_at!(1))
293 .map(Event::sender)
294 .map(ToOwned::to_owned)
295 .chain(
296 receipts
297 .ready_take_while(|(_, c, _)| *c <= newest.into_unsigned())
298 .map(|(user_id, ..)| user_id.to_owned()),
299 )
300 .collect()
301 .await;
302
303 services
304 .lazy_loading
305 .witness_retain(witness, lazy_loading_context)
306 .await
307}
308
309async fn get_member_event(
310 services: &Services,
311 room_id: &RoomId,
312 user_id: &UserId,
313) -> Option<Raw<AnyStateEvent>> {
314 services
315 .state_accessor
316 .room_state_get(room_id, &StateEventType::RoomMember, user_id.as_str())
317 .map_ok(Event::into_format)
318 .await
319 .ok()
320}
321
322pub(crate) async fn event_filters(
323 services: &Services,
324 user_id: &UserId,
325 item: PdusIterItem,
326 bypass_visibility: bool,
327) -> Option<PdusIterItem> {
328 if bypass_visibility {
329 return Some(item);
330 }
331
332 let item = ignored_filter(services, item, user_id).await?;
333 let item = visibility_filter(services, item, user_id).await?;
334
335 Some(item)
336}
337
338pub(crate) async fn related_by_filter(
342 services: &Services,
343 shortroomid: ShortRoomId,
344 filter: &RoomEventFilter,
345 item: PdusIterItem,
346) -> Option<PdusIterItem> {
347 if filter.related_by_senders.is_empty() && filter.related_by_rel_types.is_empty() {
348 return Some(item);
349 }
350
351 let rel_types: RelTypes = filter
352 .related_by_rel_types
353 .iter()
354 .map(String::as_str)
355 .map(RelationType::from)
356 .collect();
357
358 let (count, _) = &item;
359 let target = PduId { shortroomid, count: *count };
360
361 services
362 .pdu_metadata
363 .has_incoming_relation(target, &filter.related_by_senders, &rel_types)
364 .await
365 .then_some(item)
366}
367
368#[inline]
369pub(crate) async fn ignored_filter(
370 services: &Services,
371 item: PdusIterItem,
372 user_id: &UserId,
373) -> Option<PdusIterItem> {
374 let (_, ref pdu) = item;
375
376 is_ignored_pdu(services, pdu, user_id)
377 .await
378 .is_false()
379 .then_some(item)
380}
381
382#[inline]
383pub(crate) async fn is_ignored_pdu<Pdu>(
384 services: &Services,
385 event: &Pdu,
386 user_id: &UserId,
387) -> bool
388where
389 Pdu: Event,
390{
391 if event.kind().to_cow_str() == "org.matrix.dummy_event" {
394 return true;
395 }
396
397 if IGNORED_MESSAGE_TYPES
398 .binary_search(event.kind())
399 .is_err()
400 {
401 return false;
402 }
403
404 let ignored_server = services
405 .config
406 .is_forbidden_remote_server_name(event.sender().server_name());
407
408 ignored_server
409 || services
410 .users
411 .user_is_ignored(event.sender(), user_id)
412 .await
413}
414
415#[inline]
416pub(crate) fn event_filter(item: PdusIterItem, filter: &RoomEventFilter) -> Option<PdusIterItem> {
417 let (_, pdu) = &item;
418 filter.matches(pdu).then_some(item)
419}
420
421#[inline]
426pub(crate) async fn annotate_membership(
427 services: &Services,
428 pdu: &mut PduEvent,
429 user_id: &UserId,
430 encrypted: bool,
431) {
432 if !encrypted {
433 return;
434 }
435
436 let membership = services
437 .state_accessor
438 .user_membership_at_pdu(user_id, pdu)
439 .await;
440
441 pdu.add_membership(&membership).log_err().ok();
442}
443
444#[inline]
446pub(crate) async fn with_membership(
447 services: &Services,
448 mut pdu: PduEvent,
449 user_id: &UserId,
450 encrypted: bool,
451) -> PduEvent {
452 annotate_membership(services, &mut pdu, user_id, encrypted).await;
453 pdu
454}
455
456#[inline]
458pub(crate) async fn add_membership_unsigned(
459 services: &Services,
460 (count, pdu): PdusIterItem,
461 user_id: &UserId,
462 encrypted: bool,
463) -> PdusIterItem {
464 (count, with_membership(services, pdu, user_id, encrypted).await)
465}
466
467#[cfg_attr(debug_assertions, tuwunel_core::ctor(unsafe))]
468fn _is_sorted() {
469 debug_assert!(
470 IGNORED_MESSAGE_TYPES.is_sorted(),
471 "IGNORED_MESSAGE_TYPES must be sorted by the developer"
472 );
473}