Skip to main content

tuwunel_api/client/
threads.rs

1use std::collections::BTreeSet;
2
3use axum::extract::State;
4use futures::{StreamExt, TryStreamExt};
5use ruma::{
6	OwnedUserId,
7	api::client::threads::get_threads,
8	events::{
9		AnySyncMessageLikeEvent, GlobalAccountDataEventType,
10		ignored_user_list::IgnoredUserListEvent,
11	},
12	serde::Raw,
13};
14use tuwunel_core::{
15	Err, Result, at,
16	matrix::{
17		Event,
18		pdu::{PduCount, PduEvent},
19	},
20	result::{FlatOk, LogErr},
21	utils::stream::TryWidebandExt,
22};
23use tuwunel_service::rooms::pdu_metadata::IgnoredThreadView;
24
25use crate::Ruma;
26
27/// # `GET /_matrix/client/r0/rooms/{roomId}/threads`
28pub(crate) async fn get_threads_route(
29	State(services): State<crate::State>,
30	ref body: Ruma<get_threads::v1::Request>,
31) -> Result<get_threads::v1::Response> {
32	let sender_user = body.sender_user();
33	let room_id = &body.room_id;
34
35	if !services.metadata.exists(room_id).await {
36		return Err!(Request(Forbidden("Room does not exist to this server")));
37	}
38
39	if !services
40		.state_accessor
41		.user_can_see_room(sender_user, room_id)
42		.await
43	{
44		return Err!(Request(Forbidden("You don't have permission to view this room.")));
45	}
46
47	// Use limit or else 10, with maximum 100
48	let limit = body
49		.limit
50		.map(usize::try_from)
51		.flat_ok()
52		.unwrap_or(10)
53		.min(100);
54
55	let from: PduCount = body
56		.from
57		.as_deref()
58		.map(str::parse)
59		.transpose()?
60		.unwrap_or_else(PduCount::max);
61
62	// MSC3856: the requester's ignore list adjusts the served threads.
63	let ignored: BTreeSet<OwnedUserId> = services
64		.account_data
65		.get_global(sender_user, GlobalAccountDataEventType::IgnoredUserList)
66		.await
67		.map(|event: IgnoredUserListEvent| event.content.ignored_users.into_keys().collect())
68		.unwrap_or_default();
69
70	// One extra row probes whether the list continues past this page.
71	let mut threads: Vec<(PduCount, PduEvent)> = services
72		.threads
73		.threads_until(sender_user, room_id, from, &body.include)
74		.try_filter_map(async |(count, pdu)| {
75			Ok(services
76				.state_accessor
77				.user_can_see_event(sender_user, &pdu)
78				.await
79				.then_some((count, pdu)))
80		})
81		.and_then(async |(count, pdu)| {
82			let view = match ignored.is_empty() {
83				| true => IgnoredThreadView::Unchanged,
84				| false =>
85					services
86						.pdu_metadata
87						.ignored_thread_view(sender_user, &ignored, &pdu)
88						.await,
89			};
90
91			Ok((count, pdu, view))
92		})
93		.take(limit.saturating_add(1))
94		.wide_and_then(async |(count, pdu, view)| {
95			let pdu = services
96				.pdu_metadata
97				.bundle_aggregations(sender_user, pdu)
98				.await;
99
100			Ok((count, apply_ignored_view(pdu, view)))
101		})
102		.try_collect()
103		.await?;
104
105	let more = threads.len() > limit;
106
107	threads.truncate(limit);
108
109	Ok(get_threads::v1::Response {
110		next_batch: threads
111			.last()
112			.filter(|_| more)
113			.map(at!(0))
114			.as_ref()
115			.map(ToString::to_string),
116
117		chunk: threads
118			.into_iter()
119			.map(at!(1))
120			.map(Event::into_format)
121			.collect(),
122	})
123}
124
125/// MSC3856 ignored-user adjustments, applied after the bundle pass corrects
126/// the served `unsigned`: the redacted root replaces content only and keeps
127/// that `unsigned`, minus any `m.replace` bundle (a folded edit shares the
128/// root's sender, so it would re-serve the ignored content).
129fn apply_ignored_view(pdu: PduEvent, view: IgnoredThreadView) -> PduEvent {
130	match view {
131		| IgnoredThreadView::Unchanged => pdu,
132		| IgnoredThreadView::WithoutSummary { root } =>
133			without_thread_bundle(apply_redacted_root(pdu, root)),
134		| IgnoredThreadView::Adjusted { root, count, latest } =>
135			apply_redacted_root(adjust_thread_bundle(pdu, count, latest), root),
136	}
137}
138
139fn without_thread_bundle(mut pdu: PduEvent) -> PduEvent {
140	pdu.remove_thread_bundle().log_err().ok();
141	pdu
142}
143
144fn apply_redacted_root(pdu: PduEvent, root: Option<Box<PduEvent>>) -> PduEvent {
145	match root {
146		| None => pdu,
147		| Some(mut root) => {
148			root.unsigned = pdu.unsigned;
149			root.remove_replacement_bundle().log_err().ok();
150
151			*root
152		},
153	}
154}
155
156fn adjust_thread_bundle(
157	mut pdu: PduEvent,
158	count: Option<usize>,
159	latest: Option<Raw<AnySyncMessageLikeEvent>>,
160) -> PduEvent {
161	if let Some(count) = count {
162		pdu.set_thread_count(count).log_err().ok();
163	}
164
165	if let Some(latest) = latest {
166		pdu.set_thread_latest_event(&latest)
167			.log_err()
168			.ok();
169	}
170
171	pdu
172}