Skip to main content

tuwunel_api/client/
message.rs

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
43/// Shared inputs for [`get_messages`], the pagination core behind both the
44/// client-server `/messages` route and the admin room-messages endpoint.
45pub(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	/// Skip the room-visibility gate and the per-event visibility and ignore
56	/// filters, for admin callers that see all history.
57	pub bypass_visibility: bool,
58}
59
60/// list of safe and common non-state events to ignore if the user is ignored.
61/// MUST be sorted by `TimelineEventType::event_type_str()` for `binary_search`.
62const IGNORED_MESSAGE_TYPES: &[TimelineEventType] = &[
63	CallInvite,           // m.call.invite
64	KeyVerificationStart, // m.key.verification.start
65	Location,             // m.location
66	PollStart,            // m.poll.start
67	Reaction,             // m.reaction
68	RoomEncrypted,        // m.room.encrypted
69	RoomMessage,          // m.room.message
70	Sticker,              // m.sticker
71	Audio,                // org.matrix.msc1767.audio
72	Emote,                // org.matrix.msc1767.emote
73	File,                 // org.matrix.msc1767.file
74	Image,                // org.matrix.msc1767.image
75	Video,                // org.matrix.msc1767.video
76	Voice,                // org.matrix.msc3245.voice.v2
77	UnstablePollStart,    // org.matrix.msc3381.poll.start
78	Beacon,               // org.matrix.msc3672.beacon
79	CallNotify,           // org.matrix.msc4075.call.notify
80];
81
82/// MSC3440 `related_by_rel_types` entries, typed at the compare boundary.
83type RelTypes = SmallVec<[RelationType; 1]>;
84
85const LIMIT_MAX: usize = 1000;
86const LIMIT_DEFAULT: usize = 10;
87
88/// # `GET /_matrix/client/r0/rooms/{roomId}/messages`
89///
90/// Allows paginating through room history.
91///
92/// - Only works if the user is joined (TODO: always allow, but only show events
93///   where the user was joined, depending on `history_visibility`)
94pub(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
112/// Paginates a room's timeline, applying the request filter and (unless
113/// `bypass_visibility`) the per-user visibility and ignore filters. Powers the
114/// client-server `/messages` route and its admin bypass twin.
115pub(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	// `inspect` records the rejected boundary item, distinguishing a `to` stop
242	// from stream exhaustion.
243	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
338/// MSC3440 `related_by_*`: include an event only when another event relates
339/// to it matching the filter's reverse-relation criteria. A no-op stage when
340/// the filter carries neither field.
341pub(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	// exclude Synapse's dummy events from bloating up response bodies. clients
392	// don't need to see this.
393	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/// MSC4115: stamp `unsigned.membership` on a served PDU with the requesting
422/// user's membership at the time of the event. The MSC permits omitting the
423/// property when calculating it is expensive, so the project restricts it to
424/// encrypted rooms where membership-vs-event ordering matters for key share.
425#[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/// `annotate_membership` consume-and-return adapter for stream chains.
445#[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/// `with_membership` adapter for timeline-iterator items.
457#[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}