tuwunel_api/client/
threads.rs1use 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
27pub(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 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 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 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
125fn 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}