Skip to main content

tuwunel_api/client/sync/
mod.rs

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	// Take the last events for the timeline.
89	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
128/// Returns the backward pagination token for a timeline slice.
129///
130/// A slice beginning at room creation has nothing before it, so its window's
131/// end lets a members query describe the room as received.
132/// Backward pagination from that token re-reads the slice once before ending.
133fn 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
165/// MSC4155: whether a stored invite may be served to the invitee.
166///
167/// The verdict is the recipient's own, judged against the sender of the
168/// stripped invite membership event. Unreadable invite state takes the same
169/// verdict as the sender-less case below, since neither can name a sender.
170async 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
190/// [`invite_permitted_room`] for a room whose stripped state is in hand.
191///
192/// Callers walking stored invites already hold the state and take this form,
193/// which spares them the load the room-keyed form pays per room. `room_id`
194/// names the affected row in the sender-less diagnostic and does not enter the
195/// verdict. That diagnostic stays at debug on a release build deliberately,
196/// since the condition repeats for every sync an affected user makes.
197fn invite_permitted(
198	user_id: &UserId,
199	room_id: &RoomId,
200	filter: &InviteFilter,
201	invite_state: &[Raw<AnyStrippedStateEvent>],
202) -> bool {
203	// Load-bearing: keeps a permissive user off the sender derivation below.
204	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
217/// The sender of the stripped membership event inviting `user_id`.
218///
219/// The last matching entry wins. An invite this server recorded after the
220/// federation route began sanitising stripped state holds one entry for this
221/// cell, our own copy of the signed membership PDU, whose sender the origin
222/// check authenticated. An invite stored before that still carries whatever
223/// the inviting server sent ahead of our copy, and the array has no ordering
224/// semantics in the spec, so reading the last entry is what keeps those
225/// answering with the authenticated sender too.
226fn 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
246/// State sections strip the stored `prev_content`/`prev_sender` pair
247/// (Synapse injects the pair on timeline fetches only). The requester's own
248/// membership and events duplicated from the returned timeline (MSC4222,
249/// full_state) keep it: clients read membership transitions from those
250/// copies.
251fn 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}