1mod profiles;
2#[cfg(test)]
3mod tests;
4mod v3;
5mod v5;
6
7use futures::{StreamExt, pin_mut};
8use ruma::{
9 OwnedUserId, RoomId, UserId,
10 events::{
11 AnyStrippedStateEvent,
12 TimelineEventType::{RoomCreate, RoomMember},
13 },
14 serde::Raw,
15};
16use tuwunel_core::{
17 Error, PduCount, Result, debug_warn, is_equal_to,
18 matrix::{Event, pdu::PduEvent},
19 utils::{ReadyExt, result::LogErr, stream::BroadbandExt},
20};
21use tuwunel_service::{Services, users::InviteFilter};
22
23pub(crate) use self::{
24 v3::{calculate_heroes, sync_events_route},
25 v5::sync_events_v5_route,
26};
27
28#[derive(Clone, Copy)]
29enum TimelineErrors {
30 Ignore,
31 Propagate,
32}
33
34async fn load_timeline(
35 services: &Services,
36 sender_user: &UserId,
37 room_id: &RoomId,
38 roomsincecount: PduCount,
39 next_batch: Option<PduCount>,
40 limit: usize,
41) -> Result<(Vec<(PduCount, PduEvent)>, bool, PduCount), Error> {
42 load_timeline_with_errors(
43 services,
44 sender_user,
45 room_id,
46 roomsincecount,
47 next_batch,
48 limit,
49 TimelineErrors::Ignore,
50 )
51 .await
52}
53
54async fn load_timeline_fallible(
55 services: &Services,
56 sender_user: &UserId,
57 room_id: &RoomId,
58 roomsincecount: PduCount,
59 next_batch: Option<PduCount>,
60 limit: usize,
61) -> Result<(Vec<(PduCount, PduEvent)>, bool, PduCount), Error> {
62 load_timeline_with_errors(
63 services,
64 sender_user,
65 room_id,
66 roomsincecount,
67 next_batch,
68 limit,
69 TimelineErrors::Propagate,
70 )
71 .await
72}
73
74async fn load_timeline_with_errors(
75 services: &Services,
76 sender_user: &UserId,
77 room_id: &RoomId,
78 roomsincecount: PduCount,
79 next_batch: Option<PduCount>,
80 limit: usize,
81 errors: TimelineErrors,
82) -> Result<(Vec<(PduCount, PduEvent)>, bool, PduCount), Error> {
83 let until = next_batch.map(|count| count.saturating_add(1));
84 let pdus = services
85 .timeline
86 .pdus_rev(Some(sender_user), room_id, until);
87
88 pin_mut!(pdus);
90 let mut timeline_pdus = Vec::new();
91 let mut last_timeline_count = PduCount::max();
92 let mut first = true;
93 let mut limited = false;
94
95 while let Some(pdu) = pdus.next().await {
96 let (pducount, pdu) = match pdu {
97 | Ok(pdu) => pdu,
98 | Err(error) if first || matches!(errors, TimelineErrors::Propagate) => {
99 return Err(error);
100 },
101 | Err(_) => continue,
102 };
103
104 if first {
105 first = false;
106 last_timeline_count = matches!(pducount, PduCount::Normal(_))
107 .then_some(pducount)
108 .unwrap_or_else(PduCount::max);
109 }
110
111 if pducount <= roomsincecount {
112 break;
113 }
114
115 if timeline_pdus.len() == limit {
116 limited = true;
117 break;
118 }
119
120 timeline_pdus.push((pducount, pdu));
121 }
122
123 timeline_pdus.reverse();
124
125 Ok((timeline_pdus, limited, last_timeline_count))
126}
127
128fn timeline_prev_batch(
134 timeline_pdus: &[(PduCount, PduEvent)],
135 window_end: PduCount,
136) -> Option<PduCount> {
137 timeline_pdus
138 .first()
139 .map(|(count, pdu)| match pdu.kind() {
140 | RoomCreate => window_end,
141 | _ => *count,
142 })
143}
144
145async fn share_encrypted_room(
146 services: &Services,
147 sender_user: &UserId,
148 user_id: &UserId,
149 ignore_room: Option<&RoomId>,
150) -> bool {
151 services
152 .state_cache
153 .get_shared_rooms(sender_user, user_id)
154 .ready_filter(|&room_id| Some(room_id) != ignore_room)
155 .map(ToOwned::to_owned)
156 .broad_any(async |other_room_id| {
157 services
158 .state_accessor
159 .is_encrypted_room(&other_room_id)
160 .await
161 })
162 .await
163}
164
165async fn invite_permitted_room(
171 services: &Services,
172 user_id: &UserId,
173 room_id: &RoomId,
174 filter: &InviteFilter,
175) -> bool {
176 filter.is_permissive()
177 || services
178 .state_cache
179 .invite_state(user_id, room_id)
180 .await
181 .map_or_else(
182 |error| {
183 debug_warn!(%user_id, %room_id, ?error, "invite state is unreadable; skipping the sender rules");
184 filter.permits(None)
185 },
186 |invite_state| invite_permitted(user_id, room_id, filter, &invite_state),
187 )
188}
189
190fn invite_permitted(
198 user_id: &UserId,
199 room_id: &RoomId,
200 filter: &InviteFilter,
201 invite_state: &[Raw<AnyStrippedStateEvent>],
202) -> bool {
203 if filter.is_permissive() {
205 return true;
206 }
207
208 let sender = invite_sender(user_id, invite_state);
209
210 if sender.is_none() {
211 debug_warn!(%user_id, %room_id, "invite state names no sender; skipping the sender rules");
212 }
213
214 filter.permits(sender.as_deref())
215}
216
217fn invite_sender(
227 user_id: &UserId,
228 invite_state: &[Raw<AnyStrippedStateEvent>],
229) -> Option<OwnedUserId> {
230 invite_state
231 .iter()
232 .rev()
233 .filter(|event| {
234 event
235 .get_field::<&str>("state_key")
236 .is_ok_and(|state_key| state_key.is_some_and(is_equal_to!(user_id.as_str())))
237 })
238 .filter_map(|event| event.deserialize().ok())
239 .find_map(|event| match event {
240 | AnyStrippedStateEvent::RoomMember(member) if member.state_key == user_id =>
241 Some(member.sender),
242 | _ => None,
243 })
244}
245
246fn strip_prev_state(
252 mut pdu: PduEvent,
253 sender_user: &UserId,
254 in_timeline: impl Fn(&PduEvent) -> bool,
255) -> PduEvent {
256 let own_membership =
257 *pdu.kind() == RoomMember && pdu.state_key() == Some(sender_user.as_str());
258
259 if !own_membership && !in_timeline(&pdu) {
260 pdu.remove_prev_state().log_err().ok();
261 }
262
263 pdu
264}