Skip to main content

tuwunel_api/client/sync/
v3.rs

1use std::{
2	collections::{BTreeMap, HashMap, HashSet},
3	time::Duration,
4};
5
6use axum::extract::State;
7use futures::{
8	FutureExt, StreamExt, TryFutureExt, TryStreamExt,
9	future::{join, join3, join4, join5, try_join},
10	pin_mut,
11};
12use ruma::{
13	DeviceId, EventId, OwnedEventId, OwnedRoomId, OwnedUserId, RoomId, UInt, UserId,
14	api::client::{
15		filter::FilterDefinition,
16		sync::sync_events::{
17			self, DeviceLists, UnreadNotificationsCount,
18			v3::{
19				Ephemeral, Filter, GlobalAccountData, InviteState, InvitedRoom, JoinedRoom,
20				KnockState, KnockedRoom, LeftRoom, Presence, RoomAccountData, RoomSummary, Rooms,
21				State as RoomState, StateEvents, Timeline, ToDevice,
22			},
23		},
24	},
25	events::{
26		AnyGlobalAccountDataEvent, AnyRawAccountDataEvent, AnyRoomAccountDataEvent,
27		AnySyncEphemeralRoomEvent, AnySyncStateEvent, StateEventType, SyncEphemeralRoomEvent,
28		TimelineEventType::*,
29		presence::{PresenceEvent, PresenceEventContent},
30		room::member::{MembershipState, RoomMemberEventContent},
31		typing::TypingEventContent,
32	},
33	serde::Raw,
34	uint,
35};
36use tokio::time;
37use tuwunel_core::{
38	Error, Result, at, debug,
39	debug::INFO_SPAN_LEVEL,
40	debug_error, err,
41	error::{inspect_debug_log, inspect_log},
42	extract_variant, is_equal_to, is_false, is_true,
43	matrix::{
44		Event,
45		event::{Matches, trim_event_fields},
46		pdu::{EventHash, PduCount, PduEvent},
47	},
48	pair_of, ref_at,
49	result::FlatOk,
50	smallvec::SmallVec,
51	trace,
52	utils::{
53		self, BoolExt, FutureBoolExt, IterStream, ReadyExt, TryFutureExtExt,
54		future::{OptionStream, ReadyBoolExt},
55		math::ruma_from_u64,
56		option::OptionExt,
57		result::MapExpect,
58		stream::{BroadbandExt, Tools, TryBroadbandExt, TryReadyExt, WidebandExt},
59	},
60	warn,
61};
62use tuwunel_service::{
63	Services,
64	presence::Ping,
65	rooms::{
66		lazy_loading,
67		lazy_loading::{Options, Witness},
68		read_receipt::PrivateReadEvents,
69		short::{ShortEventId, ShortStateHash, ShortStateKey},
70	},
71	users::InviteFilter,
72};
73
74use super::{
75	invite_permitted, load_timeline, profiles::collect as collect_profiles, share_encrypted_room,
76	strip_prev_state, timeline_prev_batch,
77};
78use crate::{
79	ClientIp, Ruma,
80	client::{ignored_filter, is_empty_account_data_event, with_membership},
81};
82
83struct SyncParams<'a> {
84	services: &'a Services,
85	sender_user: &'a UserId,
86	sender_device: Option<&'a DeviceId>,
87	since: Option<u64>,
88	next_batch: u64,
89	full_state: bool,
90	state_after: StateAfter,
91	filter: &'a FilterDefinition,
92}
93
94#[derive(Default)]
95struct StateChanges {
96	heroes: Option<Vec<OwnedUserId>>,
97	joined_member_count: Option<u64>,
98	invited_member_count: Option<u64>,
99	state_events: Vec<PduEvent>,
100
101	/// Witnessed members' current state, lazily loaded for display; never a
102	/// membership change, so the device-list scan leaves them out.
103	lazy_members: Vec<PduEvent>,
104}
105
106struct StateChangeParams<'a> {
107	full_state: bool,
108	state_after: StateAfter,
109	since_shortstatehash: Option<ShortStateHash>,
110	horizon_shortstatehash: Option<ShortStateHash>,
111	after_shortstatehash: Option<ShortStateHash>,
112	current_shortstatehash: ShortStateHash,
113	joined_since_last_sync: bool,
114	witness: Option<&'a Witness>,
115	include_heroes: bool,
116}
117
118struct RoomMetadata {
119	since_shortstatehash: Option<ShortStateHash>,
120	horizon_shortstatehash: Option<ShortStateHash>,
121	after_shortstatehash: Option<ShortStateHash>,
122	current_shortstatehash: Option<ShortStateHash>,
123	receipt_events: Vec<(OwnedUserId, Raw<AnySyncEphemeralRoomEvent>)>,
124	encrypted_room: Option<bool>,
125}
126
127struct UserMetadata {
128	witness: Option<Witness>,
129	last_notification_read: Option<u64>,
130	thread_last_reads: BTreeMap<OwnedEventId, u64>,
131	last_privateread_update: u64,
132	joined_since_last_sync: bool,
133}
134
135struct GateInputs<'a> {
136	last_notification_read: Option<u64>,
137	thread_last_reads: &'a BTreeMap<OwnedEventId, u64>,
138	quiet_round: bool,
139	since: u64,
140}
141
142struct NotificationGates<F> {
143	send_notification_counts: bool,
144	send_notification_count_filter: F,
145}
146
147struct BuildJoinedRoom {
148	receipt_events: Vec<(OwnedUserId, Raw<AnySyncEphemeralRoomEvent>)>,
149	typing_events: Vec<Raw<AnySyncEphemeralRoomEvent>>,
150	private_read_events: Option<PrivateReadEvents>,
151	state_events: Vec<Raw<AnySyncStateEvent>>,
152	account_data_events: Vec<Raw<AnyRoomAccountDataEvent>>,
153	room_events: Vec<PduEvent>,
154	heroes: Option<Vec<OwnedUserId>>,
155	joined_member_count: Option<u64>,
156	invited_member_count: Option<u64>,
157	unread_notifications: UnreadNotificationsCount,
158	unread_thread_notifications: BTreeMap<OwnedEventId, UnreadNotificationsCount>,
159	state_after: StateAfter,
160	limited: bool,
161	joined_since_last_sync: bool,
162	prev_batch: Option<PduCount>,
163}
164
165/// MSC4222 `state_after` opt-in: which room-state field the response carries.
166#[derive(Clone, Copy, Debug)]
167enum StateAfter {
168	Off,
169	Stable,
170	Unstable,
171}
172
173type PresenceUpdates = HashMap<OwnedUserId, PresenceEventContent>;
174type TimelineEventIds = SmallVec<[OwnedEventId; 1]>;
175
176impl StateAfter {
177	fn requested(self) -> bool { !matches!(self, Self::Off) }
178
179	fn wrap(self, events: StateEvents) -> RoomState {
180		match self {
181			| Self::Off => RoomState::Before(events),
182			| Self::Stable => RoomState::After(events),
183			| Self::Unstable => RoomState::AfterUnstable(events),
184		}
185	}
186}
187
188impl From<(bool, bool)> for StateAfter {
189	fn from((stable, unstable): (bool, bool)) -> Self {
190		// Unstable opt-in wins: such a client reads the unstable field name.
191		match (stable, unstable) {
192			| (_, true) => Self::Unstable,
193			| (true, _) => Self::Stable,
194			| _ => Self::Off,
195		}
196	}
197}
198
199/// # `GET /_matrix/client/r0/sync`
200///
201/// Synchronize the client's state with the latest state on the server.
202///
203/// - This endpoint takes a `since` parameter which should be the `next_batch`
204///   value from a previous request for incremental syncs.
205///
206/// Calling this endpoint without a `since` parameter returns:
207/// - Some of the most recent events of each timeline
208/// - Notification counts for each room
209/// - Joined and invited member counts, heroes
210/// - All state events
211///
212/// Calling this endpoint with a `since` parameter from a previous `next_batch`
213/// returns: For joined rooms:
214/// - Some of the most recent events of each timeline that happened after since
215/// - If user joined the room after since: All state events (unless lazy loading
216///   is activated) and all device list updates in that room
217/// - If the user was already in the room: A list of all events that are in the
218///   state now, but were not in the state at `since`
219/// - If the state we send contains a member event: Joined and invited member
220///   counts, heroes
221/// - Device list updates that happened after `since`
222/// - If there are events in the timeline we send or the user send updated his
223///   read mark: Notification counts
224/// - EDUs that are active now (read receipts, typing updates, presence)
225/// - TODO: Allow multiple sync streams to support Pantalaimon
226///
227/// For invited rooms:
228/// - If the user was invited after `since`: A subset of the state of the room
229///   at the point of the invite
230///
231/// For left rooms:
232/// - If the user left after `since`: `prev_batch` token, empty state (TODO:
233///   subset of the state at the point of the leave)
234#[tracing::instrument(
235	name = "sync",
236	level = "debug",
237	skip_all,
238	fields(
239		user_id = %body.sender_user(),
240		device_id = %body.sender_device.as_deref().map_or("<no device>", |x| x.as_str()),
241    )
242)]
243pub(crate) async fn sync_events_route(
244	State(services): State<crate::State>,
245	ClientIp(client): ClientIp,
246	body: Ruma<sync_events::v3::Request>,
247) -> Result<sync_events::v3::Response> {
248	let sender_user = body.sender_user();
249	let sender_device = body.sender_device.as_deref();
250
251	let filter = body
252		.body
253		.filter
254		.as_ref()
255		.map_async(async |filter| match filter {
256			| Filter::FilterDefinition(filter) => filter.clone(),
257			| Filter::FilterId(filter_id) => services
258				.users
259				.get_filter(sender_user, filter_id)
260				.await
261				.unwrap_or_default(),
262		});
263
264	let filter = filter.map(Option::unwrap_or_default);
265	let full_state = body.body.full_state;
266	let set_presence = &body.body.set_presence;
267	let state_after =
268		StateAfter::from((body.body.use_state_after, body.body.use_state_after_unstable));
269
270	let ping = Ping {
271		device_id: body.sender_device.as_deref(),
272		client_ip: Some(client),
273		new_state: Some(set_presence),
274		appservice: body.appservice_info.as_ref(),
275	};
276
277	let ping_presence = services
278		.presence
279		.maybe_ping_presence(sender_user, ping)
280		.inspect_err(inspect_log)
281		.ok();
282
283	// Record user as actively syncing for push suppression heuristic.
284	let note_sync = services
285		.presence
286		.note_sync(sender_user, body.appservice_info.as_ref());
287
288	let (filter, ..) = join3(filter, ping_presence, note_sync).await;
289
290	let mut since = body
291		.body
292		.since
293		.as_deref()
294		.map(|since| since.parse().unwrap_or(0));
295
296	let timeout = body
297		.body
298		.timeout
299		.as_ref()
300		.map(Duration::as_millis)
301		.map(TryInto::try_into)
302		.flat_ok()
303		.unwrap_or(services.config.client_sync_timeout_default)
304		.max(services.config.client_sync_timeout_min)
305		.min(services.config.client_sync_timeout_max);
306
307	let stop_at = time::Instant::now()
308		.checked_add(Duration::from_millis(timeout))
309		.expect("configuration must limit maximum timeout");
310
311	loop {
312		let watch_rooms = services
313			.state_cache
314			.rooms_joined(sender_user)
315			.chain(services.state_cache.rooms_invited(sender_user));
316
317		let watchers = services
318			.sync
319			.watch(sender_user, sender_device, watch_rooms)
320			.await;
321
322		let next_batch = services.globals.wait_pending().await?;
323		if since.is_some_and(|since| since > next_batch) {
324			debug_error!(?since, next_batch, "received since > next_batch, clamping");
325			since = Some(next_batch);
326		}
327
328		if since.is_none_or(|since| since < next_batch) || full_state {
329			let response = build_sync_events(SyncParams {
330				services: &services,
331				sender_user,
332				sender_device,
333				since,
334				next_batch,
335				full_state,
336				state_after,
337				filter: &filter,
338			})
339			.await?;
340
341			let empty = response.rooms.is_empty()
342				&& response.presence.is_empty()
343				&& response.account_data.is_empty()
344				&& response.device_lists.is_empty()
345				&& response.to_device.is_empty()
346				&& response.users.is_empty();
347
348			if !empty || full_state {
349				return Ok(response);
350			}
351		}
352
353		// Wait for activity
354		if time::timeout_at(stop_at, watchers).await.is_err() || services.server.is_stopping() {
355			let response =
356				build_empty_response(&services, sender_user, sender_device, next_batch).await;
357
358			trace!(?since, next_batch, "empty response");
359			return Ok(response);
360		}
361
362		trace!(
363			?since,
364			last_batch = ?next_batch,
365			count = ?services.globals.pending_count(),
366			stop_at = ?stop_at,
367			"notified by watcher"
368		);
369
370		since = Some(next_batch);
371	}
372}
373
374async fn build_empty_response(
375	services: &Services,
376	sender_user: &UserId,
377	sender_device: Option<&DeviceId>,
378	next_batch: u64,
379) -> sync_events::v3::Response {
380	let device_one_time_keys_count = sender_device.map_async(|sender_device| {
381		services
382			.users
383			.count_one_time_keys(sender_user, sender_device)
384	});
385
386	let device_unused_fallback_key_types = sender_device.map_async(|sender_device| {
387		services
388			.users
389			.unused_fallback_key_algorithms(sender_user, sender_device)
390			.collect::<Vec<_>>()
391	});
392
393	let (device_one_time_keys_count, device_unused_fallback_key_types) =
394		join(device_one_time_keys_count, device_unused_fallback_key_types).await;
395
396	sync_events::v3::Response {
397		device_one_time_keys_count: device_one_time_keys_count.unwrap_or_default(),
398		device_unused_fallback_key_types,
399		..sync_events::v3::Response::new(next_batch.to_string())
400	}
401}
402
403#[tracing::instrument(
404	name = "build",
405	level = INFO_SPAN_LEVEL,
406	skip_all,
407	fields(
408		?since,
409		%next_batch,
410		count = ?services.globals.pending_count(),
411    )
412)]
413async fn build_sync_events(
414	SyncParams {
415		services,
416		sender_user,
417		sender_device,
418		since,
419		next_batch,
420		full_state,
421		state_after,
422		filter,
423	}: SyncParams<'_>,
424) -> Result<sync_events::v3::Response> {
425	let profile_since = since;
426	let since = since.unwrap_or(0);
427
428	// MSC4155: a stored invite whose sender the recipient ignores or blocks is
429	// withheld here, and relaxing the configuration re-exposes it.
430	let invite_filter = services.users.invite_filter(sender_user).await;
431
432	let joined_rooms = collect_joined_rooms(
433		services,
434		sender_user,
435		sender_device,
436		since,
437		next_batch,
438		full_state,
439		state_after,
440		filter,
441	);
442
443	let left_rooms = collect_left_rooms(
444		services,
445		sender_user,
446		since,
447		next_batch,
448		full_state,
449		state_after,
450		filter,
451	);
452
453	let invited_rooms =
454		collect_invited_rooms(services, sender_user, since, next_batch, filter, &invite_filter);
455
456	let knocked_rooms = collect_knocked_rooms(services, sender_user, since, next_batch, filter);
457
458	let presence_updates = services
459		.config
460		.allow_local_presence
461		.then_async(|| {
462			process_presence_updates(services, since, next_batch, sender_user, filter)
463		});
464
465	let account_data = collect_global_account_data(services, sender_user, since, next_batch);
466
467	let keys_changed = services
468		.users
469		.keys_changed(sender_user, since, Some(next_batch))
470		.map(ToOwned::to_owned)
471		.collect::<HashSet<_>>();
472
473	let to_device_events = sender_device.map_async(|sender_device| {
474		services
475			.users
476			.get_to_device_events(sender_user, sender_device, Some(since), Some(next_batch))
477			.map(at!(1))
478			.collect::<Vec<_>>()
479	});
480
481	let device_one_time_keys_count = sender_device.map_async(|sender_device| {
482		services
483			.users
484			.count_one_time_keys(sender_user, sender_device)
485	});
486
487	let device_unused_fallback_key_types = sender_device.map_async(|sender_device| {
488		services
489			.users
490			.unused_fallback_key_algorithms(sender_user, sender_device)
491			.collect::<Vec<_>>()
492	});
493
494	// Remove all to-device events the device received *last time*
495	let remove_to_device_events = sender_device.map_async(|sender_device| {
496		services
497			.users
498			.remove_to_device_events(sender_user, sender_device, since)
499	});
500
501	let (
502		account_data,
503		keys_changed,
504		presence_updates,
505		(_, to_device_events, device_one_time_keys_count, device_unused_fallback_key_types),
506		(
507			(joined_rooms, mut device_list_updates, left_encrypted_users),
508			left_rooms,
509			invited_rooms,
510			knocked_rooms,
511		),
512	) = join5(
513		account_data,
514		keys_changed,
515		presence_updates,
516		join4(
517			remove_to_device_events,
518			to_device_events,
519			device_one_time_keys_count,
520			device_unused_fallback_key_types,
521		),
522		join4(joined_rooms, left_rooms, invited_rooms, knocked_rooms),
523	)
524	.boxed()
525	.await;
526
527	device_list_updates.extend(keys_changed);
528
529	let rooms = Rooms {
530		leave: left_rooms,
531		join: joined_rooms,
532		invite: invited_rooms,
533		knock: knocked_rooms,
534	};
535
536	// The profile reads and the left-device fan-out both pend on the pool.
537	let (device_list_left, users) = join(
538		collect_device_list_left(services, sender_user, left_encrypted_users),
539		collect_profiles(services, sender_user, profile_since, next_batch, filter, &rooms),
540	)
541	.await;
542
543	let presence_events = build_presence_events(presence_updates);
544	let device_lists = DeviceLists {
545		left: device_list_left,
546		changed: device_list_updates.into_iter().collect(),
547	};
548
549	let to_device = ToDevice {
550		events: to_device_events.unwrap_or_default(),
551	};
552
553	Ok(sync_events::v3::Response {
554		account_data: GlobalAccountData { events: account_data },
555		device_lists,
556		device_one_time_keys_count: device_one_time_keys_count.unwrap_or_default(),
557		device_unused_fallback_key_types,
558		next_batch: next_batch.to_string(),
559		presence: Presence { events: presence_events },
560		rooms,
561		to_device,
562		users: users?,
563	})
564}
565
566#[expect(clippy::too_many_arguments)]
567fn collect_joined_rooms<'a>(
568	services: &'a Services,
569	sender_user: &'a UserId,
570	sender_device: Option<&'a DeviceId>,
571	since: u64,
572	next_batch: u64,
573	full_state: bool,
574	state_after: StateAfter,
575	filter: &'a FilterDefinition,
576) -> impl Future<
577	Output = (BTreeMap<OwnedRoomId, JoinedRoom>, HashSet<OwnedUserId>, HashSet<OwnedUserId>),
578> + Send
579+ 'a {
580	services
581		.state_cache
582		.rooms_joined(sender_user)
583		.ready_filter(|&room_id| filter.room.matches(room_id))
584		.map(ToOwned::to_owned)
585		.broad_filter_map(move |room_id| {
586			load_joined_room(
587				services,
588				sender_user,
589				sender_device,
590				room_id.clone(),
591				since,
592				next_batch,
593				full_state,
594				state_after,
595				filter,
596			)
597			.map_ok(move |(joined_room, dlu, jeu)| (room_id, joined_room, dlu, jeu))
598			.ok()
599		})
600		.ready_fold(
601			(BTreeMap::new(), HashSet::new(), HashSet::new()),
602			|(mut joined_rooms, mut device_list_updates, mut left_encrypted_users),
603			 (room_id, joined_room, dlu, leu)| {
604				device_list_updates.extend(dlu);
605				left_encrypted_users.extend(leu);
606				if !joined_room.is_empty() {
607					joined_rooms.insert(room_id, joined_room);
608				}
609
610				(joined_rooms, device_list_updates, left_encrypted_users)
611			},
612		)
613}
614
615fn collect_left_rooms<'a>(
616	services: &'a Services,
617	sender_user: &'a UserId,
618	since: u64,
619	next_batch: u64,
620	full_state: bool,
621	state_after: StateAfter,
622	filter: &'a FilterDefinition,
623) -> impl Future<Output = BTreeMap<OwnedRoomId, LeftRoom>> + Send + 'a {
624	services
625		.state_cache
626		.rooms_left_state(sender_user)
627		.ready_filter(|(room_id, _)| filter.room.matches(room_id))
628		.broad_filter_map(move |(room_id, _)| {
629			handle_left_room(
630				services,
631				since,
632				room_id.clone(),
633				sender_user,
634				next_batch,
635				full_state,
636				state_after,
637				filter,
638			)
639			.map_ok(move |left_room| (room_id, left_room))
640			.ok()
641		})
642		.ready_filter_map(|(room_id, left_room)| left_room.map(|left_room| (room_id, left_room)))
643		.collect()
644}
645
646async fn collect_invited_rooms<'a>(
647	services: &'a Services,
648	sender_user: &'a UserId,
649	since: u64,
650	next_batch: u64,
651	filter: &'a FilterDefinition,
652	invite_filter: &'a InviteFilter,
653) -> BTreeMap<OwnedRoomId, InvitedRoom> {
654	services
655		.state_cache
656		.rooms_invited_state(sender_user)
657		.ready_filter(|(room_id, _)| filter.room.matches(room_id))
658		.ready_filter(|(room_id, invite_state)| {
659			invite_permitted(sender_user, room_id, invite_filter, invite_state)
660		})
661		.fold_default(async |mut invited_rooms: BTreeMap<_, _>, (room_id, invite_state)| {
662			let invite_count = services
663				.state_cache
664				.get_invite_count(&room_id, sender_user)
665				.await
666				.ok();
667
668			// Invited before last sync
669			if Some(since) >= invite_count || Some(next_batch) < invite_count {
670				return invited_rooms;
671			}
672
673			let invited_room = InvitedRoom {
674				invite_state: InviteState { events: invite_state },
675			};
676
677			invited_rooms.insert(room_id, invited_room);
678			invited_rooms
679		})
680		.await
681}
682
683async fn collect_knocked_rooms<'a>(
684	services: &'a Services,
685	sender_user: &'a UserId,
686	since: u64,
687	next_batch: u64,
688	filter: &'a FilterDefinition,
689) -> BTreeMap<OwnedRoomId, KnockedRoom> {
690	services
691		.state_cache
692		.rooms_knocked_state(sender_user)
693		.ready_filter(|(room_id, _)| filter.room.matches(room_id))
694		.fold_default(async |mut knocked_rooms: BTreeMap<_, _>, (room_id, knock_state)| {
695			let knock_count = services
696				.state_cache
697				.get_knock_count(&room_id, sender_user)
698				.await
699				.ok();
700
701			// Knocked before last sync; or after the cutoff for this sync
702			if Some(since) >= knock_count || Some(next_batch) < knock_count {
703				return knocked_rooms;
704			}
705
706			let knocked_room = KnockedRoom {
707				knock_state: KnockState { events: knock_state },
708			};
709
710			knocked_rooms.insert(room_id, knocked_room);
711			knocked_rooms
712		})
713		.await
714}
715
716fn collect_global_account_data<'a>(
717	services: &'a Services,
718	sender_user: &'a UserId,
719	since: u64,
720	next_batch: u64,
721) -> impl Future<Output = Vec<Raw<AnyGlobalAccountDataEvent>>> + Send + 'a {
722	services
723		.account_data
724		.changes_since(None, sender_user, since, Some(next_batch))
725		.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Global))
726		.ready_filter(move |e| since != 0 || !is_empty_account_data_event(e))
727		.collect()
728}
729
730fn collect_device_list_left<'a>(
731	services: &'a Services,
732	sender_user: &'a UserId,
733	left_encrypted_users: HashSet<OwnedUserId>,
734) -> impl Future<Output = Vec<OwnedUserId>> + Send + 'a {
735	left_encrypted_users
736		.into_iter()
737		.stream()
738		.broad_filter_map(async |user_id: OwnedUserId| {
739			share_encrypted_room(services, sender_user, &user_id, None)
740				.await
741				.eq(&false)
742				.then_some(user_id)
743		})
744		.collect()
745}
746
747fn build_presence_events(presence_updates: Option<PresenceUpdates>) -> Vec<Raw<PresenceEvent>> {
748	presence_updates
749		.into_iter()
750		.flat_map(IntoIterator::into_iter)
751		.map(|(sender, content)| PresenceEvent { content, sender })
752		.map(|ref event| Raw::new(event))
753		.filter_map(Result::ok)
754		.collect()
755}
756
757#[tracing::instrument(name = "presence", level = "debug", skip_all)]
758async fn process_presence_updates(
759	services: &Services,
760	since: u64,
761	next_batch: u64,
762	syncing_user: &UserId,
763	filter: &FilterDefinition,
764) -> PresenceUpdates {
765	services
766		.presence
767		.presence_since(since, Some(next_batch))
768		.ready_filter(|(user_id, ..)| filter.presence.matches(user_id))
769		.filter(|(user_id, ..)| {
770			services
771				.state_cache
772				.user_sees_user(syncing_user, user_id)
773		})
774		.filter_map(|(user_id, _, presence_bytes)| {
775			services
776				.presence
777				.from_json_bytes_to_event(presence_bytes, user_id)
778				.map_ok(move |event| (user_id, event))
779				.ok()
780		})
781		.map(|(user_id, event)| (user_id.to_owned(), event.content))
782		.collect()
783		.boxed()
784		.await
785}
786
787#[tracing::instrument(
788	name = "left",
789	level = "debug",
790	skip_all,
791	fields(
792		room_id = %room_id,
793		full = %full_state,
794	),
795)]
796#[expect(clippy::too_many_arguments)]
797async fn handle_left_room(
798	services: &Services,
799	since: u64,
800	ref room_id: OwnedRoomId,
801	sender_user: &UserId,
802	next_batch: u64,
803	full_state: bool,
804	state_after: StateAfter,
805	filter: &FilterDefinition,
806) -> Result<Option<LeftRoom>> {
807	let left_count = services
808		.state_cache
809		.get_left_count(room_id, sender_user)
810		.await
811		.unwrap_or(0);
812
813	if left_count == 0 || left_count > next_batch {
814		return Ok(None);
815	}
816
817	let include_leave = filter.room.include_leave;
818	if since == 0 && !include_leave {
819		return Ok(None);
820	}
821
822	// Cannot sync unless the event falls within the snapshot. The room is only
823	// sync'ed once to the client, after that it's too late.
824	if since != 0 && left_count <= since {
825		return Ok(None);
826	}
827
828	let is_not_found = services.metadata.exists(room_id).is_false();
829
830	let is_disabled = services.metadata.is_disabled(room_id);
831
832	let is_banned = services.metadata.is_banned(room_id);
833
834	pin_mut!(is_not_found, is_disabled, is_banned);
835	if is_not_found.or(is_disabled).or(is_banned).await {
836		// For rejected invites, deleted, missing, or broken room state this is the last
837		// resort to convey a the minimum of information to the client.
838		let event = PduEvent {
839			event_id: EventId::new_v1(services.globals.server_name()),
840			origin_server_ts: utils::millis_since_unix_epoch().try_into()?,
841			kind: RoomMember,
842			state_key: Some(sender_user.as_str().into()),
843			sender: sender_user.to_owned(),
844			content: serde_json::from_str(r#"{"membership":"leave"}"#)?,
845			// The following keys are dropped on conversion
846			room_id: room_id.clone(),
847			depth: uint!(1),
848			origin: None,
849			unsigned: None,
850			redacts: None,
851			hashes: EventHash::default(),
852			auth_events: Default::default(),
853			prev_events: Default::default(),
854		};
855
856		let state = state_after.wrap(StateEvents {
857			events: vec![trim_event_fields(event.into_format(), filter.event_fields.as_deref())],
858		});
859
860		return Ok(Some(LeftRoom {
861			account_data: RoomAccountData::default(),
862			state,
863			timeline: Timeline {
864				limited: false,
865				events: Default::default(),
866				prev_batch: Some(left_count.to_string()),
867			},
868		}));
869	}
870
871	load_left_room(
872		services,
873		sender_user,
874		room_id,
875		since,
876		left_count,
877		full_state,
878		state_after,
879		filter,
880	)
881	.await
882}
883
884#[tracing::instrument(name = "load", level = "debug", skip_all)]
885#[expect(clippy::too_many_arguments)]
886async fn load_left_room(
887	services: &Services,
888	sender_user: &UserId,
889	room_id: &RoomId,
890	since: u64,
891	left_count: u64,
892	full_state: bool,
893	state_after: StateAfter,
894	filter: &FilterDefinition,
895) -> Result<Option<LeftRoom>> {
896	let initial = since == 0;
897	let timeline_limit: usize = filter
898		.room
899		.timeline
900		.limit
901		.map(TryInto::try_into)
902		.map_expect("UInt to usize")
903		.unwrap_or(10)
904		.min(100);
905
906	let (timeline_pdus, limited, _) = load_timeline(
907		services,
908		sender_user,
909		room_id,
910		PduCount::Normal(since),
911		Some(PduCount::Normal(left_count)),
912		timeline_limit.max(1),
913	)
914	.await
915	.unwrap_or_default();
916
917	let since_shortstatehash = services
918		.timeline
919		.next_shortstatehash(room_id, PduCount::Normal(since))
920		.ok();
921
922	let horizon_shortstatehash = timeline_pdus
923		.first()
924		.map(at!(0))
925		.map_async(|count| {
926			services
927				.timeline
928				.get_shortstatehash(room_id, count)
929				.inspect_err(log_horizon_error)
930				.ok()
931		});
932
933	// MSC4222 `state_after` describes the state at the leave.
934	let after_shortstatehash = state_after.requested().then_async(|| {
935		services
936			.timeline
937			.shortstatehash_after(room_id, PduCount::Normal(left_count))
938			.inspect_err(inspect_debug_log)
939	});
940
941	let left_shortstatehash = services
942		.timeline
943		.get_shortstatehash(room_id, PduCount::Normal(left_count))
944		.inspect_err(inspect_debug_log)
945		.or_else(|_| services.state.get_room_shortstatehash(room_id))
946		.map_err(|_| err!(Database(error!("Room {room_id} has no state"))));
947
948	let (since_shortstatehash, horizon_shortstatehash, after_shortstatehash, left_shortstatehash) =
949		join4(
950			since_shortstatehash,
951			horizon_shortstatehash,
952			after_shortstatehash,
953			left_shortstatehash,
954		)
955		.boxed()
956		.await;
957
958	let StateChanges { state_events, .. } =
959		calculate_state_changes(services, sender_user, room_id, StateChangeParams {
960			full_state: full_state || initial,
961			state_after,
962			since_shortstatehash,
963			horizon_shortstatehash: horizon_shortstatehash.flatten(),
964			after_shortstatehash: after_shortstatehash.flat_ok(),
965			current_shortstatehash: left_shortstatehash?,
966			joined_since_last_sync: false,
967			witness: None,
968			include_heroes: true,
969		})
970		.boxed()
971		.await?;
972
973	let is_sender_membership = |event: &PduEvent| {
974		*event.kind() == RoomMember && event.state_key() == Some(sender_user.as_str())
975	};
976
977	let timeline_sender_member = timeline_limit
978		.eq(&0)
979		.then(|| timeline_pdus.last().map(ref_at!(1)).cloned())
980		.into_iter()
981		.flat_map(Option::into_iter);
982
983	let encrypted = services
984		.state_accessor
985		.is_encrypted_room(room_id)
986		.await;
987
988	let event_fields = filter.event_fields.as_deref();
989
990	let in_timeline = in_timeline(&timeline_pdus);
991
992	let state_events = state_events
993		.into_iter()
994		.filter(|pdu| filter.room.state.matches(pdu))
995		.filter(|pdu| timeline_limit > 0 || !is_sender_membership(pdu))
996		.chain(timeline_sender_member)
997		.stream()
998		.wide_then(|pdu| with_membership(services, pdu, sender_user, encrypted))
999		.map(|pdu| strip_prev_state(pdu, sender_user, &in_timeline))
1000		.map(|pdu| trim_event_fields(pdu.into_format(), event_fields))
1001		.collect();
1002
1003	let left_prev_batch = timeline_limit
1004		.eq(&0)
1005		.then_some(left_count)
1006		.map(PduCount::Normal);
1007
1008	let prev_batch = timeline_prev_batch(&timeline_pdus, PduCount::Normal(left_count))
1009		.filter(|_| timeline_limit > 0)
1010		.or(left_prev_batch)
1011		.as_ref()
1012		.map(ToString::to_string);
1013
1014	let timeline_events = timeline_pdus
1015		.into_iter()
1016		.stream()
1017		.wide_filter_map(|item| ignored_filter(services, item, sender_user))
1018		.map(at!(1))
1019		.ready_filter(|pdu| filter.room.timeline.matches(pdu))
1020		.take(timeline_limit)
1021		.wide_then(|pdu| with_membership(services, pdu, sender_user, encrypted))
1022		.wide_then(|pdu| {
1023			services
1024				.pdu_metadata
1025				.bundle_aggregations(sender_user, pdu)
1026		})
1027		.collect::<Vec<_>>();
1028
1029	let account_data_events = services
1030		.account_data
1031		.changes_since(Some(room_id), sender_user, since, None)
1032		.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Room))
1033		.ready_filter(move |e| since != 0 || !is_empty_account_data_event(e))
1034		.collect();
1035
1036	let (state_events, account_data_events, timeline_events) =
1037		join3(state_events, account_data_events, timeline_events)
1038			.boxed()
1039			.await;
1040
1041	let state = state_after.wrap(StateEvents { events: state_events });
1042
1043	Ok(Some(LeftRoom {
1044		account_data: RoomAccountData { events: account_data_events },
1045		state,
1046		timeline: Timeline {
1047			prev_batch,
1048			limited: limited || timeline_limit == 0,
1049			events: timeline_events
1050				.into_iter()
1051				.map(|pdu| trim_event_fields(pdu.into_format(), event_fields))
1052				.collect(),
1053		},
1054	}))
1055}
1056
1057/// Records a failed horizon read, quietly when the row is simply absent.
1058///
1059/// `append_to_state` writes no state hash for a room's `m.room.create` event,
1060/// there being no state before it, so the first window of every room misses
1061/// here by design. An absent row reads the same whatever left it absent, so
1062/// the loud path is what remains: a fault that is not a miss at all.
1063fn log_horizon_error(error: &Error) {
1064	match error.is_not_found() {
1065		| true => debug!(?error, "no shortstatehash at the window's first event"),
1066		| false => inspect_debug_log(error),
1067	}
1068}
1069
1070fn in_timeline(timeline_pdus: &[(PduCount, PduEvent)]) -> impl Fn(&PduEvent) -> bool + use<> {
1071	let timeline_ids: TimelineEventIds = timeline_pdus
1072		.iter()
1073		.map(ref_at!(1))
1074		.map(Event::event_id)
1075		.map(ToOwned::to_owned)
1076		.collect();
1077
1078	move |event: &PduEvent| {
1079		timeline_ids
1080			.iter()
1081			.any(is_equal_to!(event.event_id()))
1082	}
1083}
1084
1085#[tracing::instrument(
1086	name = "joined",
1087	level = "debug",
1088	skip_all,
1089	fields(
1090		room_id = ?room_id,
1091	),
1092)]
1093#[expect(clippy::too_many_arguments)]
1094async fn load_joined_room(
1095	services: &Services,
1096	sender_user: &UserId,
1097	sender_device: Option<&DeviceId>,
1098	ref room_id: OwnedRoomId,
1099	since: u64,
1100	next_batch: u64,
1101	full_state: bool,
1102	state_after: StateAfter,
1103	filter: &FilterDefinition,
1104) -> Result<(JoinedRoom, HashSet<OwnedUserId>, HashSet<OwnedUserId>)> {
1105	let initial = since == 0;
1106	let (timeline_pdus, limited, last_timeline_count) =
1107		load_join_timeline(services, sender_user, room_id, since, next_batch, filter).await?;
1108
1109	let timeline_changed = last_timeline_count.into_unsigned() > since;
1110	debug_assert!(
1111		timeline_pdus.is_empty() || timeline_changed,
1112		"if timeline events, last_timeline_count must be in the since window."
1113	);
1114
1115	let RoomMetadata {
1116		since_shortstatehash,
1117		horizon_shortstatehash,
1118		after_shortstatehash,
1119		current_shortstatehash,
1120		receipt_events,
1121		encrypted_room,
1122	} = gather_room_metadata(
1123		services,
1124		sender_user,
1125		room_id,
1126		since,
1127		next_batch,
1128		&timeline_pdus,
1129		last_timeline_count,
1130		timeline_changed || full_state || initial,
1131		state_after,
1132	)
1133	.boxed()
1134	.await?;
1135
1136	let UserMetadata {
1137		witness,
1138		last_notification_read,
1139		thread_last_reads,
1140		last_privateread_update,
1141		joined_since_last_sync,
1142	} = gather_user_metadata(
1143		services,
1144		sender_user,
1145		sender_device,
1146		room_id,
1147		filter,
1148		&timeline_pdus,
1149		&receipt_events,
1150		since,
1151		initial,
1152		timeline_changed,
1153		encrypted_room,
1154		since_shortstatehash,
1155	)
1156	.boxed()
1157	.await;
1158
1159	let (
1160		state_after,
1161		StateChanges {
1162			heroes,
1163			joined_member_count,
1164			invited_member_count,
1165			mut state_events,
1166			lazy_members,
1167		},
1168	) = compute_join_state_changes(
1169		services,
1170		sender_user,
1171		room_id,
1172		full_state || initial,
1173		state_after,
1174		since_shortstatehash,
1175		horizon_shortstatehash,
1176		after_shortstatehash,
1177		current_shortstatehash,
1178		joined_since_last_sync,
1179		witness.as_ref(),
1180	)
1181	.await?;
1182
1183	let quiet_round = timeline_pdus.is_empty();
1184
1185	let joined_sender_member = take_sender_membership_for_join(
1186		&mut state_events,
1187		sender_user,
1188		joined_since_last_sync,
1189		quiet_round,
1190		initial,
1191	);
1192
1193	let prev_batch =
1194		compute_join_prev_batch(&timeline_pdus, joined_sender_member.as_ref(), since, next_batch);
1195
1196	let in_window = in_window(since, next_batch);
1197
1198	let NotificationGates {
1199		send_notification_counts,
1200		send_notification_count_filter,
1201	} = compute_notification_gates(
1202		GateInputs {
1203			last_notification_read,
1204			thread_last_reads: &thread_last_reads,
1205			quiet_round,
1206			since,
1207		},
1208		in_window,
1209	);
1210
1211	// `encrypted_room` is `Some` whenever timeline or state events are emitted.
1212	let encrypted = encrypted_room.unwrap_or(false);
1213
1214	let aggregates = await_join_aggregates(
1215		services,
1216		sender_user,
1217		room_id,
1218		&state_events,
1219		timeline_pdus,
1220		joined_sender_member,
1221		encrypted,
1222		initial,
1223		since,
1224		next_batch,
1225		last_privateread_update,
1226		send_notification_counts,
1227		filter,
1228	)
1229	.await;
1230
1231	let state_events = state_events.into_iter().chain(lazy_members);
1232
1233	let (joined_room, device_list_updates, left_encrypted_users) = finalize_joined_room(
1234		services,
1235		sender_user,
1236		filter,
1237		state_events,
1238		aggregates,
1239		receipt_events,
1240		heroes,
1241		joined_member_count,
1242		invited_member_count,
1243		quiet_round.then_some(&thread_last_reads),
1244		send_notification_count_filter,
1245		FinalizeJoinFlags {
1246			encrypted,
1247			full_state,
1248			state_after,
1249			limited,
1250			joined_since_last_sync,
1251			initial,
1252		},
1253		in_window,
1254		prev_batch,
1255	)
1256	.await;
1257
1258	Ok((joined_room, device_list_updates, left_encrypted_users))
1259}
1260
1261#[expect(clippy::too_many_arguments)]
1262async fn compute_join_state_changes(
1263	services: &Services,
1264	sender_user: &UserId,
1265	room_id: &RoomId,
1266	full_state: bool,
1267	state_after: StateAfter,
1268	since_shortstatehash: Option<ShortStateHash>,
1269	horizon_shortstatehash: Option<ShortStateHash>,
1270	after_shortstatehash: Option<ShortStateHash>,
1271	current_shortstatehash: Option<ShortStateHash>,
1272	joined_since_last_sync: bool,
1273	witness: Option<&Witness>,
1274) -> Result<(StateAfter, StateChanges)> {
1275	let Some(current_shortstatehash) = current_shortstatehash else {
1276		return Ok((state_after, StateChanges::default()));
1277	};
1278
1279	let state_changes =
1280		calculate_state_changes(services, sender_user, room_id, StateChangeParams {
1281			full_state,
1282			state_after,
1283			since_shortstatehash,
1284			horizon_shortstatehash,
1285			after_shortstatehash,
1286			current_shortstatehash,
1287			joined_since_last_sync,
1288			witness,
1289			include_heroes: true,
1290		})
1291		.await;
1292
1293	let incremental = !full_state && !joined_since_last_sync && since_shortstatehash.is_some();
1294
1295	if !state_after.requested() || incremental {
1296		return state_changes.map(|state_changes| (state_after, state_changes));
1297	}
1298
1299	match state_changes {
1300		| Ok(state_changes) => Ok((state_after, state_changes)),
1301		| Err(after_error) => {
1302			let after_boundary = after_shortstatehash.unwrap_or(current_shortstatehash);
1303			let legacy_boundary = horizon_shortstatehash.unwrap_or(current_shortstatehash);
1304
1305			warn!(
1306				%room_id,
1307				%after_boundary,
1308				?after_error,
1309				"Failed to load requested state-after boundary; retrying legacy state."
1310			);
1311
1312			calculate_state_changes(services, sender_user, room_id, StateChangeParams {
1313				full_state,
1314				state_after: StateAfter::Off,
1315				since_shortstatehash,
1316				horizon_shortstatehash,
1317				after_shortstatehash,
1318				current_shortstatehash,
1319				joined_since_last_sync,
1320				witness,
1321				include_heroes: false,
1322			})
1323			.await
1324			.inspect_err(|legacy_error| {
1325				warn!(
1326					%room_id,
1327					%after_boundary,
1328					%legacy_boundary,
1329					?after_error,
1330					?legacy_error,
1331					"Failed to load state-after and legacy state boundaries."
1332				);
1333			})
1334			.map(|state_changes| (StateAfter::Off, state_changes))
1335		},
1336	}
1337}
1338
1339fn compute_join_prev_batch(
1340	timeline_pdus: &[(PduCount, PduEvent)],
1341	joined_sender_member: Option<&PduEvent>,
1342	since: u64,
1343	next_batch: u64,
1344) -> Option<PduCount> {
1345	timeline_prev_batch(timeline_pdus, PduCount::Normal(next_batch)).or_else(|| {
1346		joined_sender_member
1347			.is_some()
1348			.then_some(since)
1349			.map(Into::into)
1350	})
1351}
1352
1353#[expect(clippy::too_many_arguments)]
1354async fn assemble_join_state_events(
1355	services: &Services,
1356	state_events: impl Iterator<Item = PduEvent> + Send,
1357	sender_user: &UserId,
1358	encrypted: bool,
1359	room_events: &[PduEvent],
1360	filter: &FilterDefinition,
1361	full_state: bool,
1362	state_after: StateAfter,
1363) -> Vec<Raw<AnySyncStateEvent>> {
1364	let is_in_timeline = |event: &PduEvent| {
1365		room_events
1366			.iter()
1367			.map(Event::event_id)
1368			.any(is_equal_to!(event.event_id()))
1369	};
1370
1371	// MSC4222: when the client opts into `state_after`, state events that
1372	// took effect within the timeline appear in both the timeline and the
1373	// state section, so the in-timeline exclusion is bypassed.
1374	let include_in_state = |event: &PduEvent| {
1375		let filter = &filter.room.state;
1376		filter.matches(event) && (full_state || state_after.requested() || !is_in_timeline(event))
1377	};
1378
1379	assemble_state_events(
1380		services,
1381		state_events,
1382		sender_user,
1383		encrypted,
1384		include_in_state,
1385		&is_in_timeline,
1386		filter.event_fields.as_deref(),
1387	)
1388	.await
1389}
1390
1391async fn load_join_timeline(
1392	services: &Services,
1393	sender_user: &UserId,
1394	room_id: &RoomId,
1395	since: u64,
1396	next_batch: u64,
1397	filter: &FilterDefinition,
1398) -> Result<(Vec<(PduCount, PduEvent)>, bool, PduCount)> {
1399	let timeline_limit: usize = filter
1400		.room
1401		.timeline
1402		.limit
1403		.map(TryInto::try_into)
1404		.map_expect("UInt to usize")
1405		.unwrap_or(10)
1406		.min(100);
1407
1408	load_timeline(
1409		services,
1410		sender_user,
1411		room_id,
1412		PduCount::Normal(since),
1413		Some(PduCount::Normal(next_batch)),
1414		timeline_limit,
1415	)
1416	.await
1417}
1418
1419struct JoinAggregates {
1420	room_events: Vec<PduEvent>,
1421	account_data_events: Vec<Raw<AnyRoomAccountDataEvent>>,
1422	typing_events: Vec<Raw<AnySyncEphemeralRoomEvent>>,
1423	private_read_events: Option<PrivateReadEvents>,
1424	notification_count: Option<UInt>,
1425	highlight_count: Option<UInt>,
1426	thread_counts: Option<BTreeMap<OwnedEventId, (u64, u64)>>,
1427	device_list_updates: HashSet<OwnedUserId>,
1428	left_encrypted_users: HashSet<OwnedUserId>,
1429}
1430
1431#[expect(clippy::too_many_arguments)]
1432async fn await_join_aggregates(
1433	services: &Services,
1434	sender_user: &UserId,
1435	room_id: &RoomId,
1436	state_events: &[PduEvent],
1437	timeline_pdus: Vec<(PduCount, PduEvent)>,
1438	joined_sender_member: Option<PduEvent>,
1439	encrypted: bool,
1440	initial: bool,
1441	since: u64,
1442	next_batch: u64,
1443	last_privateread_update: u64,
1444	send_notification_counts: bool,
1445	filter: &FilterDefinition,
1446) -> JoinAggregates {
1447	let (notification_count, highlight_count, thread_counts) =
1448		notification_count_futures(services, sender_user, room_id, send_notification_counts);
1449
1450	let private_read_events = last_privateread_update.gt(&since).then_async(|| {
1451		services
1452			.read_receipt
1453			.private_read_get(room_id, sender_user)
1454			.unwrap_or_default()
1455	});
1456
1457	let typing_events = gather_typing_events(services, room_id, sender_user, since);
1458
1459	let device_list_updates = gather_device_list_updates(
1460		services,
1461		sender_user,
1462		room_id,
1463		timeline_membership_changes(&timeline_pdus, initial),
1464		state_events,
1465		initial,
1466		since,
1467		next_batch,
1468	);
1469
1470	let room_events = collect_room_events(
1471		services,
1472		sender_user,
1473		timeline_pdus,
1474		joined_sender_member,
1475		encrypted,
1476		filter,
1477	);
1478
1479	let account_data_events = collect_room_account_data(services, sender_user, room_id, since);
1480
1481	let (
1482		(room_events, account_data_events),
1483		(typing_events, private_read_events),
1484		(notification_count, highlight_count, thread_counts),
1485		(device_list_updates, left_encrypted_users),
1486	) = join4(
1487		join(room_events, account_data_events),
1488		join(typing_events, private_read_events),
1489		join3(notification_count, highlight_count, thread_counts),
1490		device_list_updates,
1491	)
1492	.boxed()
1493	.await;
1494
1495	JoinAggregates {
1496		room_events,
1497		account_data_events,
1498		typing_events,
1499		private_read_events,
1500		notification_count,
1501		highlight_count,
1502		thread_counts,
1503		device_list_updates,
1504		left_encrypted_users,
1505	}
1506}
1507
1508struct FinalizeJoinFlags {
1509	encrypted: bool,
1510	full_state: bool,
1511	state_after: StateAfter,
1512	limited: bool,
1513	joined_since_last_sync: bool,
1514	initial: bool,
1515}
1516
1517#[expect(clippy::too_many_arguments)]
1518async fn finalize_joined_room(
1519	services: &Services,
1520	sender_user: &UserId,
1521	filter: &FilterDefinition,
1522	state_events: impl Iterator<Item = PduEvent> + Send,
1523	aggregates: JoinAggregates,
1524	receipt_events: Vec<(OwnedUserId, Raw<AnySyncEphemeralRoomEvent>)>,
1525	heroes: Option<Vec<OwnedUserId>>,
1526	joined_member_count: Option<u64>,
1527	invited_member_count: Option<u64>,
1528	thread_last_reads: Option<&BTreeMap<OwnedEventId, u64>>,
1529	send_notification_count_filter: impl Fn(&UInt) -> bool,
1530	flags: FinalizeJoinFlags,
1531	in_window: impl Fn(u64) -> bool,
1532	prev_batch: Option<PduCount>,
1533) -> (JoinedRoom, HashSet<OwnedUserId>, HashSet<OwnedUserId>) {
1534	let JoinAggregates {
1535		room_events,
1536		account_data_events,
1537		typing_events,
1538		private_read_events,
1539		notification_count,
1540		highlight_count,
1541		thread_counts,
1542		device_list_updates,
1543		left_encrypted_users,
1544	} = aggregates;
1545
1546	let FinalizeJoinFlags {
1547		encrypted,
1548		full_state,
1549		state_after,
1550		limited,
1551		joined_since_last_sync,
1552		initial,
1553	} = flags;
1554
1555	let state_events = assemble_join_state_events(
1556		services,
1557		state_events,
1558		sender_user,
1559		encrypted,
1560		&room_events,
1561		filter,
1562		full_state,
1563		state_after,
1564	)
1565	.await;
1566
1567	let (unread_notifications, unread_thread_notifications) = assemble_unread_notifications(
1568		notification_count,
1569		highlight_count,
1570		thread_counts,
1571		thread_last_reads,
1572		send_notification_count_filter,
1573		filter.room.timeline.unread_thread_notifications,
1574		initial,
1575		in_window,
1576	);
1577
1578	let joined_room = build_joined_room(
1579		BuildJoinedRoom {
1580			receipt_events,
1581			typing_events,
1582			private_read_events,
1583			state_events,
1584			account_data_events,
1585			room_events,
1586			heroes,
1587			joined_member_count,
1588			invited_member_count,
1589			unread_notifications,
1590			unread_thread_notifications,
1591			state_after,
1592			limited,
1593			joined_since_last_sync,
1594			prev_batch,
1595		},
1596		filter.event_fields.as_deref(),
1597	);
1598
1599	(joined_room, device_list_updates, left_encrypted_users)
1600}
1601
1602fn build_joined_room(args: BuildJoinedRoom, event_fields: Option<&[String]>) -> JoinedRoom {
1603	let BuildJoinedRoom {
1604		receipt_events,
1605		typing_events,
1606		private_read_events,
1607		state_events,
1608		account_data_events,
1609		room_events,
1610		heroes,
1611		joined_member_count,
1612		invited_member_count,
1613		unread_notifications,
1614		unread_thread_notifications,
1615		state_after,
1616		limited,
1617		joined_since_last_sync,
1618		prev_batch,
1619	} = args;
1620
1621	let edus: Vec<Raw<AnySyncEphemeralRoomEvent>> = receipt_events
1622		.into_iter()
1623		.map(at!(1))
1624		.chain(typing_events)
1625		.chain(private_read_events.into_iter().flatten())
1626		.collect();
1627
1628	let state = state_after.wrap(StateEvents { events: state_events });
1629
1630	let heroes = heroes
1631		.into_iter()
1632		.flatten()
1633		.map(TryInto::try_into)
1634		.filter_map(Result::ok)
1635		.collect();
1636
1637	JoinedRoom {
1638		account_data: RoomAccountData { events: account_data_events },
1639		ephemeral: Ephemeral { events: edus },
1640		state,
1641		summary: RoomSummary {
1642			joined_member_count: joined_member_count.map(ruma_from_u64),
1643			invited_member_count: invited_member_count.map(ruma_from_u64),
1644			heroes,
1645		},
1646		timeline: Timeline {
1647			limited: limited || joined_since_last_sync,
1648			prev_batch: prev_batch.as_ref().map(ToString::to_string),
1649			events: room_events
1650				.into_iter()
1651				.map(|pdu| trim_event_fields(pdu.into_format(), event_fields))
1652				.collect(),
1653		},
1654		unread_notifications,
1655		unread_thread_notifications,
1656	}
1657}
1658
1659#[expect(clippy::too_many_arguments)]
1660async fn gather_room_metadata(
1661	services: &Services,
1662	sender_user: &UserId,
1663	room_id: &RoomId,
1664	since: u64,
1665	next_batch: u64,
1666	timeline_pdus: &[(PduCount, PduEvent)],
1667	last_timeline_count: PduCount,
1668	include_state: bool,
1669	state_after: StateAfter,
1670) -> Result<RoomMetadata> {
1671	let since_shortstatehash = include_state.then_async(|| {
1672		services
1673			.timeline
1674			.prev_shortstatehash(room_id, PduCount::Normal(since).saturating_add(1))
1675			.ok()
1676	});
1677
1678	let horizon_shortstatehash = timeline_pdus
1679		.first()
1680		.map(at!(0))
1681		.map_async(|count| {
1682			services
1683				.timeline
1684				.get_shortstatehash(room_id, count)
1685				.inspect_err(log_horizon_error)
1686		});
1687
1688	// MSC4222 `state_after` describes the state at the end of the timeline.
1689	let after_shortstatehash = state_after.requested().then_async(|| {
1690		services
1691			.timeline
1692			.shortstatehash_after(room_id, last_timeline_count)
1693			.inspect_err(inspect_debug_log)
1694	});
1695
1696	let current_shortstatehash = include_state.then_async(|| {
1697		// An empty timeline needs state after the last event, not before it.
1698		timeline_pdus
1699			.is_empty()
1700			.then(|| {
1701				services
1702					.timeline
1703					.next_shortstatehash(room_id, last_timeline_count)
1704					.left_future()
1705			})
1706			.unwrap_or_else(|| {
1707				services
1708					.timeline
1709					.get_shortstatehash(room_id, last_timeline_count)
1710					.inspect_err(inspect_debug_log)
1711					.right_future()
1712			})
1713			.or_else(|_| services.state.get_room_shortstatehash(room_id))
1714			.map_err(|_| err!(Database(error!("Room {room_id} has no state"))))
1715	});
1716
1717	let encrypted_room =
1718		include_state.then_async(|| services.state_accessor.is_encrypted_room(room_id));
1719
1720	let receipt_events = services
1721		.read_receipt
1722		.readreceipts_since(room_id, since, Some(next_batch))
1723		.filter_map(async |(read_user, _, edu)| {
1724			services
1725				.users
1726				.user_is_ignored(read_user, sender_user)
1727				.await
1728				.or_some((read_user.to_owned(), edu))
1729		})
1730		.collect::<Vec<(OwnedUserId, Raw<AnySyncEphemeralRoomEvent>)>>();
1731
1732	let (
1733		(
1734			since_shortstatehash,
1735			horizon_shortstatehash,
1736			after_shortstatehash,
1737			current_shortstatehash,
1738		),
1739		receipt_events,
1740		encrypted_room,
1741	) = join3(
1742		join4(
1743			since_shortstatehash,
1744			horizon_shortstatehash,
1745			after_shortstatehash,
1746			current_shortstatehash,
1747		),
1748		receipt_events,
1749		encrypted_room,
1750	)
1751	.boxed()
1752	.await;
1753
1754	Ok(RoomMetadata {
1755		since_shortstatehash: since_shortstatehash.flatten(),
1756		horizon_shortstatehash: horizon_shortstatehash.flat_ok(),
1757		after_shortstatehash: after_shortstatehash.flat_ok(),
1758		current_shortstatehash: current_shortstatehash.transpose()?,
1759		receipt_events,
1760		encrypted_room,
1761	})
1762}
1763
1764fn collect_room_events<'a>(
1765	services: &'a Services,
1766	sender_user: &'a UserId,
1767	timeline_pdus: Vec<(PduCount, PduEvent)>,
1768	joined_sender_member: Option<PduEvent>,
1769	encrypted: bool,
1770	filter: &'a FilterDefinition,
1771) -> impl Future<Output = Vec<PduEvent>> + Send + 'a {
1772	let include_in_timeline = |event: &PduEvent| filter.room.timeline.matches(event);
1773	timeline_pdus
1774		.into_iter()
1775		.stream()
1776		.wide_filter_map(|item| ignored_filter(services, item, sender_user))
1777		.map(at!(1))
1778		.chain(joined_sender_member.into_iter().stream())
1779		.ready_filter(include_in_timeline)
1780		.wide_then(move |pdu| with_membership(services, pdu, sender_user, encrypted))
1781		.wide_then(move |pdu| {
1782			services
1783				.pdu_metadata
1784				.bundle_aggregations(sender_user, pdu)
1785		})
1786		.collect::<Vec<_>>()
1787}
1788
1789fn collect_room_account_data<'a>(
1790	services: &'a Services,
1791	sender_user: &'a UserId,
1792	room_id: &'a RoomId,
1793	since: u64,
1794) -> impl Future<Output = Vec<Raw<AnyRoomAccountDataEvent>>> + Send + 'a {
1795	services
1796		.account_data
1797		.changes_since(Some(room_id), sender_user, since, None)
1798		.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Room))
1799		.ready_filter(move |e| since != 0 || !is_empty_account_data_event(e))
1800		.collect()
1801}
1802
1803#[expect(clippy::type_complexity)]
1804fn notification_count_futures<'a>(
1805	services: &'a Services,
1806	sender_user: &'a UserId,
1807	room_id: &'a RoomId,
1808	send: bool,
1809) -> (
1810	impl Future<Output = Option<UInt>> + Send + 'a,
1811	impl Future<Output = Option<UInt>> + Send + 'a,
1812	impl Future<Output = Option<BTreeMap<OwnedEventId, (u64, u64)>>> + Send + 'a,
1813) {
1814	let notification_count = send.then_async(move || {
1815		services
1816			.pusher
1817			.notification_count(sender_user, room_id)
1818			.map(TryInto::try_into)
1819			.unwrap_or(uint!(0))
1820	});
1821
1822	let highlight_count = send.then_async(move || {
1823		services
1824			.pusher
1825			.highlight_count(sender_user, room_id)
1826			.map(TryInto::try_into)
1827			.unwrap_or(uint!(0))
1828	});
1829
1830	// MSC3773: per-thread counts. Filtered downstream by per-thread last-read
1831	// so quiet threads are omitted on rounds where the main cursor advanced.
1832	let thread_counts = send.then_async(move || {
1833		services
1834			.pusher
1835			.thread_notification_counts(sender_user, room_id)
1836	});
1837
1838	(notification_count, highlight_count, thread_counts)
1839}
1840
1841fn take_sender_membership_for_join(
1842	state_events: &mut Vec<PduEvent>,
1843	sender_user: &UserId,
1844	joined_since_last_sync: bool,
1845	timeline_empty: bool,
1846	initial: bool,
1847) -> Option<PduEvent> {
1848	if !(joined_since_last_sync && timeline_empty && !initial) {
1849		return None;
1850	}
1851
1852	let is_sender_membership = |event: &PduEvent| {
1853		*event.event_type() == StateEventType::RoomMember.into()
1854			&& event
1855				.state_key()
1856				.is_some_and(is_equal_to!(sender_user.as_str()))
1857	};
1858
1859	state_events
1860		.iter()
1861		.position(is_sender_membership)
1862		.map(|pos| state_events.swap_remove(pos))
1863}
1864
1865#[expect(clippy::too_many_arguments)]
1866async fn gather_user_metadata(
1867	services: &Services,
1868	sender_user: &UserId,
1869	sender_device: Option<&DeviceId>,
1870	room_id: &RoomId,
1871	filter: &FilterDefinition,
1872	timeline_pdus: &[(PduCount, PduEvent)],
1873	receipt_events: &[(OwnedUserId, Raw<AnySyncEphemeralRoomEvent>)],
1874	since: u64,
1875	initial: bool,
1876	timeline_changed: bool,
1877	encrypted_room: Option<bool>,
1878	since_shortstatehash: Option<ShortStateHash>,
1879) -> UserMetadata {
1880	let lazy_load_options =
1881		[&filter.room.state.lazy_load_options, &filter.room.timeline.lazy_load_options];
1882
1883	let lazy_loading_enabled = encrypted_room.is_some_and(is_false!())
1884		&& lazy_load_options
1885			.iter()
1886			.any(|opts| opts.is_enabled());
1887
1888	let lazy_loading_context = &lazy_loading::Context {
1889		user_id: sender_user,
1890		device_id: sender_device,
1891		room_id,
1892		token: Some(since),
1893		options: Some(&filter.room.state.lazy_load_options),
1894		mode: lazy_loading::Mode::Update,
1895	};
1896
1897	// Reset lazy loading because this is an initial sync
1898	let lazy_load_reset =
1899		initial.then_async(|| services.lazy_loading.reset(lazy_loading_context));
1900
1901	lazy_load_reset.await;
1902	let witness = lazy_loading_enabled.then_async(|| {
1903		let witness: Witness = timeline_pdus
1904			.iter()
1905			.map(ref_at!(1))
1906			.map(Event::sender)
1907			.map(Into::into)
1908			.chain(receipt_events.iter().map(ref_at!(0)).cloned())
1909			.collect();
1910
1911		services
1912			.lazy_loading
1913			.witness_retain(witness, lazy_loading_context)
1914	});
1915
1916	let sender_joined_count = timeline_changed.then_async(|| {
1917		services
1918			.state_cache
1919			.get_joined_count(room_id, sender_user)
1920			.unwrap_or(0)
1921	});
1922
1923	let since_encryption = since_shortstatehash.map_async(|shortstatehash| {
1924		services
1925			.state_accessor
1926			.state_get(shortstatehash, &StateEventType::RoomEncryption, "")
1927	});
1928
1929	// Busy rounds need these cursors too, to authorize an explicit zero count.
1930	let last_notification_read = services
1931		.pusher
1932		.last_notification_read(sender_user, room_id)
1933		.ok();
1934
1935	let thread_last_reads = services
1936		.pusher
1937		.thread_last_notification_reads(sender_user, room_id);
1938
1939	let last_privateread_update = services
1940		.read_receipt
1941		.last_privateread_update(sender_user, room_id);
1942
1943	let (
1944		(last_privateread_update, last_notification_read, thread_last_reads),
1945		(sender_joined_count, since_encryption),
1946		witness,
1947	) = join3(
1948		join3(last_privateread_update, last_notification_read, thread_last_reads),
1949		join(sender_joined_count, since_encryption),
1950		witness,
1951	)
1952	.await;
1953
1954	let _encrypted_since_last_sync =
1955		!initial && encrypted_room.is_some_and(is_true!()) && since_encryption.is_none();
1956
1957	let joined_since_last_sync = sender_joined_count.unwrap_or(0) > since;
1958
1959	UserMetadata {
1960		witness,
1961		last_notification_read,
1962		thread_last_reads,
1963		last_privateread_update,
1964		joined_since_last_sync,
1965	}
1966}
1967
1968/// The window a sync round reports a cursor advance in.
1969///
1970/// The lower bound is the round's own token and the upper bound is the
1971/// response cutoff captured before the room is assembled, so a cursor stamped
1972/// after that cutoff falls outside it.
1973fn in_window(since: u64, next_batch: u64) -> impl Fn(u64) -> bool + Copy {
1974	move |count| count > since && count <= next_batch
1975}
1976
1977/// The gates deciding what notification counts a joined room reports.
1978///
1979/// A count reaches the client when the round carried timeline events or a read
1980/// cursor advanced within the window. An explicit zero additionally needs a
1981/// cursor newer than the request token, and because the response cutoff is
1982/// captured before that cursor is read, a cursor past the cutoff repeats the
1983/// zero until the token catches up.
1984fn compute_notification_gates<F: Fn(u64) -> bool>(
1985	GateInputs {
1986		last_notification_read,
1987		thread_last_reads,
1988		quiet_round,
1989		since,
1990	}: GateInputs<'_>,
1991	in_window: F,
1992) -> NotificationGates<impl Fn(&UInt) -> bool + use<F>> {
1993	// Thread-only resets leave the main cursor alone, so without the thread leg
1994	// a quiet round would never report them.
1995	let send_notification_counts = !quiet_round
1996		|| last_notification_read.is_none_or(&in_window)
1997		|| thread_last_reads
1998			.values()
1999			.copied()
2000			.any(&in_window);
2001
2002	let after_since = |count: u64| count > since;
2003
2004	// A cursor newer than since means the client may still cache the pre-reset
2005	// count, so allow an explicit zero to reconcile it.
2006	let send_notification_resets = last_notification_read.is_some_and(&after_since)
2007		|| thread_last_reads
2008			.values()
2009			.copied()
2010			.any(&after_since);
2011
2012	let send_notification_count_filter =
2013		move |count: &UInt| *count != uint!(0) || send_notification_resets;
2014
2015	NotificationGates {
2016		send_notification_counts,
2017		send_notification_count_filter,
2018	}
2019}
2020
2021async fn gather_typing_events(
2022	services: &Services,
2023	room_id: &RoomId,
2024	sender_user: &UserId,
2025	since: u64,
2026) -> Vec<Raw<AnySyncEphemeralRoomEvent>> {
2027	services
2028		.typing
2029		.last_typing_update(room_id)
2030		.and_then(async |count| {
2031			if count <= since {
2032				return Ok(Vec::<Raw<AnySyncEphemeralRoomEvent>>::new());
2033			}
2034
2035			let typings = typings_event_for_user(services, room_id, sender_user).await?;
2036
2037			Ok(vec![serde_json::from_str(&serde_json::to_string(&typings)?)?])
2038		})
2039		.unwrap_or(Vec::new())
2040		.await
2041}
2042
2043fn timeline_membership_changes(
2044	timeline_pdus: &[(PduCount, PduEvent)],
2045	initial: bool,
2046) -> Vec<(MembershipState, OwnedUserId)> {
2047	timeline_pdus
2048		.iter()
2049		.filter(|_| !initial)
2050		.map(ref_at!(1))
2051		.filter_map(extract_membership)
2052		.collect::<Vec<_>>()
2053}
2054
2055fn extract_membership(event: &PduEvent) -> Option<(MembershipState, OwnedUserId)> {
2056	let content: RoomMemberEventContent = event.get_content().ok()?;
2057	let user_id: OwnedUserId = event.state_key()?.parse().ok()?;
2058
2059	Some((content.membership, user_id))
2060}
2061
2062#[expect(clippy::too_many_arguments)]
2063async fn gather_device_list_updates(
2064	services: &Services,
2065	sender_user: &UserId,
2066	room_id: &RoomId,
2067	timeline_membership_changes: Vec<(MembershipState, OwnedUserId)>,
2068	state_events: &[PduEvent],
2069	initial: bool,
2070	since: u64,
2071	next_batch: u64,
2072) -> (HashSet<OwnedUserId>, HashSet<OwnedUserId>) {
2073	let keys_changed = services
2074		.users
2075		.room_keys_changed(room_id, since, Some(next_batch))
2076		.map(|(user_id, _)| user_id)
2077		.map(ToOwned::to_owned)
2078		.collect::<Vec<_>>();
2079
2080	let (mut dlu, leu) = state_events
2081		.iter()
2082		.stream()
2083		.ready_filter(|_| !initial)
2084		.ready_filter(|state_event| *state_event.event_type() == RoomMember)
2085		.ready_filter_map(extract_membership)
2086		.chain(timeline_membership_changes.into_iter().stream())
2087		.fold_default(async |(mut dlu, mut leu): pair_of!(HashSet<_>), (membership, user_id)| {
2088			use MembershipState::*;
2089
2090			let requires_update = async |user_id| {
2091				!share_encrypted_room(services, sender_user, user_id, Some(room_id)).await
2092			};
2093
2094			match membership {
2095				| Join if requires_update(&user_id).await => dlu.insert(user_id),
2096				| Leave => leu.insert(user_id),
2097				| _ => false,
2098			};
2099
2100			(dlu, leu)
2101		})
2102		.await;
2103
2104	dlu.extend(keys_changed.await);
2105	(dlu, leu)
2106}
2107
2108async fn assemble_state_events(
2109	services: &Services,
2110	state_events: impl Iterator<Item = PduEvent> + Send,
2111	sender_user: &UserId,
2112	encrypted: bool,
2113	include_in_state: impl Fn(&PduEvent) -> bool + Send + Sync,
2114	in_timeline: impl Fn(&PduEvent) -> bool + Send + Sync,
2115	event_fields: Option<&[String]>,
2116) -> Vec<Raw<AnySyncStateEvent>> {
2117	state_events
2118		.filter(include_in_state)
2119		.stream()
2120		.wide_then(|pdu| with_membership(services, pdu, sender_user, encrypted))
2121		.map(|pdu| strip_prev_state(pdu, sender_user, &in_timeline))
2122		.map(|pdu| trim_event_fields(pdu.into_format(), event_fields))
2123		.collect()
2124		.await
2125}
2126
2127#[expect(clippy::too_many_arguments)]
2128fn assemble_unread_notifications(
2129	notification_count: Option<UInt>,
2130	highlight_count: Option<UInt>,
2131	thread_counts: Option<BTreeMap<OwnedEventId, (u64, u64)>>,
2132	thread_last_reads: Option<&BTreeMap<OwnedEventId, u64>>,
2133	send_notification_count_filter: impl Fn(&UInt) -> bool,
2134	want_thread_unread: bool,
2135	initial: bool,
2136	in_window: impl Fn(u64) -> bool,
2137) -> (UnreadNotificationsCount, BTreeMap<OwnedEventId, UnreadNotificationsCount>) {
2138	let thread_counts = thread_counts.unwrap_or_default();
2139
2140	let (thread_total_notifications, thread_total_highlights) = thread_counts
2141		.values()
2142		.fold((0_u64, 0_u64), |(n, h), &(notifs, hl)| {
2143			(n.saturating_add(notifs), h.saturating_add(hl))
2144		});
2145
2146	// MSC3773: when the client opts in via the timeline filter, partition
2147	// notification counts per thread. Otherwise sum into the room total.
2148	let merge_total = |total: u64| {
2149		move |count: UInt| {
2150			want_thread_unread
2151				.is_false()
2152				.then(|| count.saturating_add(UInt::try_from(total).unwrap_or_default()))
2153				.unwrap_or(count)
2154		}
2155	};
2156
2157	let unread_notifications = UnreadNotificationsCount {
2158		highlight_count: highlight_count
2159			.map(merge_total(thread_total_highlights))
2160			.filter(&send_notification_count_filter),
2161		notification_count: notification_count
2162			.map(merge_total(thread_total_notifications))
2163			.filter(&send_notification_count_filter),
2164	};
2165
2166	// On quiet rounds (timeline empty) `thread_last_reads` is `Some`; emit
2167	// only threads whose read cursor advanced within the window. When the
2168	// timeline carried events `thread_last_reads` is `None`; emit all.
2169	// Initial sync (since == 0) is a full snapshot; bypass the gate so
2170	// clients with no prior cursor still see existing thread counts.
2171	let advanced_in_window = |root: &EventId| {
2172		initial
2173			|| thread_last_reads
2174				.is_none_or(|reads| reads.get(root).copied().is_some_and(&in_window))
2175	};
2176
2177	let unread_thread_notifications = thread_counts
2178		.into_iter()
2179		.filter(|_| want_thread_unread)
2180		.filter(|(root, _)| advanced_in_window(root))
2181		.map(|(root, (notifications, highlights))| {
2182			let counts = UnreadNotificationsCount {
2183				notification_count: UInt::try_from(notifications).ok(),
2184				highlight_count: UInt::try_from(highlights).ok(),
2185			};
2186
2187			(root, counts)
2188		})
2189		.collect();
2190
2191	(unread_notifications, unread_thread_notifications)
2192}
2193
2194#[tracing::instrument(
2195	name = "state",
2196	level = "trace",
2197	skip_all,
2198	fields(
2199	    full = %full_state,
2200	    after = ?state_after,
2201	    ss = ?since_shortstatehash,
2202	    hs = ?horizon_shortstatehash,
2203	    as = ?after_shortstatehash,
2204	    cs = %current_shortstatehash,
2205    )
2206)]
2207async fn calculate_state_changes<'a>(
2208	services: &Services,
2209	sender_user: &UserId,
2210	room_id: &RoomId,
2211	StateChangeParams {
2212		full_state,
2213		state_after,
2214		since_shortstatehash,
2215		horizon_shortstatehash,
2216		after_shortstatehash,
2217		current_shortstatehash,
2218		joined_since_last_sync,
2219		witness,
2220		include_heroes,
2221	}: StateChangeParams<'a>,
2222) -> Result<StateChanges> {
2223	let incremental = !full_state && !joined_since_last_sync && since_shortstatehash.is_some();
2224
2225	// MSC4222: `state_after` requests need state at the *end* of the
2226	// timeline; legacy `state` requests need state at the *start*. Pick
2227	// the right delta endpoint, falling back to the room's current
2228	// shortstatehash when the preferred lookup is unavailable.
2229	let horizon_shortstatehash = state_after
2230		.requested()
2231		.then_some(after_shortstatehash)
2232		.unwrap_or(horizon_shortstatehash)
2233		.unwrap_or(current_shortstatehash);
2234
2235	let since_shortstatehash = since_shortstatehash.unwrap_or(horizon_shortstatehash);
2236
2237	let state_get_shorteventid = |user_id: &'a UserId| {
2238		services
2239			.state_accessor
2240			.state_get_shortid(
2241				horizon_shortstatehash,
2242				&StateEventType::RoomMember,
2243				user_id.as_str(),
2244			)
2245			.ok()
2246	};
2247
2248	let get_pdu = |shorteventid: ShortEventId| {
2249		services
2250			.timeline
2251			.get_pdu_from_shorteventid(shorteventid)
2252			.ok()
2253	};
2254
2255	let lazy_state_ids = witness.map_async(|witness| {
2256		witness
2257			.iter()
2258			.stream()
2259			.ready_filter(|&user_id| user_id != sender_user)
2260			.broad_filter_map(|user_id| state_get_shorteventid(user_id))
2261			.into_future()
2262	});
2263
2264	let state_diff_ids = incremental.then_async(|| {
2265		services
2266			.state_accessor
2267			.state_added((since_shortstatehash, horizon_shortstatehash))
2268			.boxed()
2269			.into_future()
2270	});
2271
2272	let current_state_ids = (!incremental).then_async(|| {
2273		services
2274			.state_accessor
2275			.state_full_shortids(horizon_shortstatehash)
2276			.boxed()
2277			.into_future()
2278	});
2279
2280	// Full dump is strict; the delta relaxes the member filter under MSC4222.
2281	let after = state_after.requested();
2282	let state_events = current_state_ids
2283		.stream()
2284		.map_ok(|ids| (false, ids))
2285		.chain(
2286			state_diff_ids
2287				.stream()
2288				.map(move |ids| Ok((after, ids))),
2289		)
2290		.broad_and_then(async |(after, (shortstatekey, shorteventid))| {
2291			let event_id =
2292				lazy_filter(services, sender_user, witness, shortstatekey, shorteventid, after)
2293					.await;
2294
2295			Ok(event_id)
2296		})
2297		.ready_try_filter_map(Result::Ok)
2298		.broad_and_then(|shorteventid| get_pdu(shorteventid).map(Ok))
2299		.ready_try_filter_map(Result::Ok)
2300		.try_collect::<Vec<_>>();
2301
2302	let lazy_members = lazy_state_ids
2303		.stream()
2304		.broad_filter_map(get_pdu)
2305		.collect::<Vec<_>>()
2306		.map(Ok);
2307
2308	let (state_events, lazy_members) = try_join(state_events, lazy_members).await?;
2309
2310	let send_member_counts = state_events
2311		.iter()
2312		.chain(lazy_members.iter())
2313		.any(|event| *event.kind() == RoomMember);
2314
2315	let member_counts = send_member_counts
2316		.then_async(|| calculate_counts(services, room_id, sender_user, include_heroes));
2317
2318	let (joined_member_count, invited_member_count, heroes) =
2319		member_counts.await.unwrap_or((None, None, None));
2320
2321	Ok(StateChanges {
2322		heroes,
2323		joined_member_count,
2324		invited_member_count,
2325		state_events,
2326		lazy_members,
2327	})
2328}
2329
2330async fn lazy_filter(
2331	services: &Services,
2332	sender_user: &UserId,
2333	witness: Option<&Witness>,
2334	shortstatekey: ShortStateKey,
2335	shorteventid: ShortEventId,
2336	after: bool,
2337) -> Option<ShortEventId> {
2338	let Some(witness) = witness else {
2339		return Some(shorteventid);
2340	};
2341
2342	let (event_type, state_key) = services
2343		.short
2344		.get_statekey_from_short(shortstatekey)
2345		.await
2346		.ok()?;
2347
2348	// An MSC4222 delta also keeps changed members the witness will not re-add
2349	// (lazy_state_ids covers witnessed ones), avoiding both a miss and a duplicate.
2350	let keep = event_type != StateEventType::RoomMember
2351		|| state_key == sender_user.as_str()
2352		|| (after && <&UserId>::try_from(state_key.as_str()).is_ok_and(|u| !witness.contains(u)));
2353
2354	keep.then_some(shorteventid)
2355}
2356
2357async fn calculate_counts(
2358	services: &Services,
2359	room_id: &RoomId,
2360	sender_user: &UserId,
2361	include_heroes: bool,
2362) -> (Option<u64>, Option<u64>, Option<Vec<OwnedUserId>>) {
2363	let joined_member_count = services
2364		.state_cache
2365		.room_joined_count(room_id)
2366		.unwrap_or(0);
2367
2368	let invited_member_count = services
2369		.state_cache
2370		.room_invited_count(room_id)
2371		.unwrap_or(0);
2372
2373	let (joined_member_count, invited_member_count) =
2374		join(joined_member_count, invited_member_count).await;
2375
2376	let small_room = joined_member_count.saturating_add(invited_member_count) <= 5;
2377
2378	let heroes = services
2379		.config
2380		.calculate_heroes
2381		.and_is(include_heroes)
2382		.and_is(small_room)
2383		.then_async(|| calculate_heroes(services, room_id, sender_user));
2384
2385	(Some(joined_member_count), Some(invited_member_count), heroes.await)
2386}
2387
2388pub(crate) async fn calculate_heroes(
2389	services: &Services,
2390	room_id: &RoomId,
2391	sender_user: &UserId,
2392) -> Vec<OwnedUserId> {
2393	const LIMIT: usize = 5;
2394
2395	services
2396		.state_accessor
2397		.room_state_type_pdus(room_id, &StateEventType::RoomMember)
2398		.ready_filter_map(Result::ok)
2399		.filter_map(|pdu| filter_hero(services, room_id, sender_user, pdu))
2400		.take(LIMIT)
2401		.collect::<Vec<_>>()
2402		.await
2403}
2404
2405async fn filter_hero<Pdu: Event>(
2406	services: &Services,
2407	room_id: &RoomId,
2408	sender_user: &UserId,
2409	pdu: Pdu,
2410) -> Option<OwnedUserId> {
2411	let user_id = pdu.state_key().map(TryInto::try_into).flat_ok()?;
2412
2413	if user_id == sender_user {
2414		return None;
2415	}
2416
2417	let Ok(content): Result<RoomMemberEventContent, _> = pdu.get_content() else {
2418		return None;
2419	};
2420
2421	// The membership was and still is invite or join
2422	if !matches!(content.membership, MembershipState::Join | MembershipState::Invite) {
2423		return None;
2424	}
2425
2426	// The join read leads; a joined hero settles it before the invite read.
2427	let is_hero = services
2428		.state_cache
2429		.is_joined(user_id, room_id)
2430		.is_false()
2431		.and(services.state_cache.is_invited(user_id, room_id))
2432		.is_false();
2433
2434	is_hero.await.then(|| user_id.to_owned())
2435}
2436
2437async fn typings_event_for_user(
2438	services: &Services,
2439	room_id: &RoomId,
2440	sender_user: &UserId,
2441) -> Result<SyncEphemeralRoomEvent<TypingEventContent>> {
2442	Ok(SyncEphemeralRoomEvent {
2443		content: TypingEventContent {
2444			user_ids: services
2445				.typing
2446				.typing_users_for_user(room_id, sender_user)
2447				.await?,
2448		},
2449	})
2450}
2451
2452#[cfg(test)]
2453mod tests {
2454	use super::*;
2455
2456	const SINCE: u64 = 100;
2457	const NEXT_BATCH: u64 = 200;
2458	const INITIAL_SINCE: u64 = 0;
2459
2460	#[test]
2461	fn state_after_wraps_into_named_variant() {
2462		let events = StateEvents::default;
2463
2464		assert!(matches!(StateAfter::Off.wrap(events()), RoomState::Before(_)));
2465		assert!(matches!(StateAfter::Stable.wrap(events()), RoomState::After(_)));
2466		assert!(matches!(StateAfter::Unstable.wrap(events()), RoomState::AfterUnstable(_)));
2467
2468		assert!(!StateAfter::Off.requested());
2469		assert!(StateAfter::Stable.requested());
2470		assert!(StateAfter::Unstable.requested());
2471	}
2472
2473	#[test]
2474	fn state_after_selects_unstable_when_both_opted_in() {
2475		// (use_state_after, use_state_after_unstable)
2476		assert!(matches!(StateAfter::from((false, false)), StateAfter::Off));
2477		assert!(matches!(StateAfter::from((true, false)), StateAfter::Stable));
2478		assert!(matches!(StateAfter::from((false, true)), StateAfter::Unstable));
2479		assert!(matches!(StateAfter::from((true, true)), StateAfter::Unstable));
2480	}
2481
2482	#[test]
2483	fn busy_round_delivers_reset_zeros() {
2484		let gates = notification_gates(Some(150), &no_threads(), false, SINCE);
2485
2486		assert!(gates.send_notification_counts);
2487		assert!((gates.send_notification_count_filter)(&uint!(0)));
2488	}
2489
2490	#[test]
2491	fn busy_round_without_reset_suppresses_zeros() {
2492		// A cursor at 50 advanced before this window, so it is not a reset.
2493		let gates = notification_gates(Some(50), &no_threads(), false, SINCE);
2494
2495		assert!(gates.send_notification_counts);
2496		assert!(!(gates.send_notification_count_filter)(&uint!(0)));
2497		assert!((gates.send_notification_count_filter)(&uint!(5)));
2498	}
2499
2500	#[test]
2501	fn busy_round_without_cursor_sends_counts() {
2502		let gates = notification_gates(None, &no_threads(), false, SINCE);
2503
2504		assert!(gates.send_notification_counts);
2505		assert!(!(gates.send_notification_count_filter)(&uint!(0)));
2506	}
2507
2508	#[test]
2509	fn quiet_round_with_advanced_cursor_sends_zeros() {
2510		let gates = notification_gates(Some(150), &no_threads(), true, SINCE);
2511
2512		assert!(gates.send_notification_counts);
2513		assert!((gates.send_notification_count_filter)(&uint!(0)));
2514	}
2515
2516	#[test]
2517	fn quiet_round_with_stale_cursor_sends_nothing() {
2518		// A cursor at 50 advanced before this window, so nothing changed.
2519		let gates = notification_gates(Some(50), &no_threads(), true, SINCE);
2520
2521		assert!(!gates.send_notification_counts);
2522		assert!(!(gates.send_notification_count_filter)(&uint!(0)));
2523	}
2524
2525	#[test]
2526	fn quiet_round_without_cursor_suppresses_zeros() {
2527		let gates = notification_gates(None, &no_threads(), true, SINCE);
2528
2529		assert!(gates.send_notification_counts);
2530		assert!(!(gates.send_notification_count_filter)(&uint!(0)));
2531	}
2532
2533	#[test]
2534	fn initial_sync_delivers_reset_zeros() {
2535		// A client can start a fresh token while still caching a count, so an
2536		// initial round must be able to reconcile it.
2537		let gates = notification_gates(Some(150), &no_threads(), false, INITIAL_SINCE);
2538
2539		assert!(gates.send_notification_counts);
2540		assert!((gates.send_notification_count_filter)(&uint!(0)));
2541	}
2542
2543	#[test]
2544	fn cursor_past_the_cutoff_delivers_reset_zeros() {
2545		// The response cutoff is captured before the cursor is read, so a
2546		// cursor beyond it repeats the zero until the token catches up.
2547		let gates = notification_gates(Some(NEXT_BATCH + 1), &no_threads(), false, SINCE);
2548
2549		assert!(gates.send_notification_counts);
2550		assert!((gates.send_notification_count_filter)(&uint!(0)));
2551	}
2552
2553	#[test]
2554	fn thread_only_reset_delivers_zeros() {
2555		let root = EventId::parse("$thread:example.org").unwrap();
2556		let thread_reads = BTreeMap::from([(root, 150)]);
2557
2558		let gates = notification_gates(None, &thread_reads, false, SINCE);
2559
2560		assert!(gates.send_notification_counts);
2561		assert!((gates.send_notification_count_filter)(&uint!(0)));
2562	}
2563
2564	fn notification_gates(
2565		last_read: Option<u64>,
2566		thread_reads: &BTreeMap<OwnedEventId, u64>,
2567		quiet_round: bool,
2568		since: u64,
2569	) -> NotificationGates<impl Fn(&UInt) -> bool + use<>> {
2570		compute_notification_gates(
2571			GateInputs {
2572				last_notification_read: last_read,
2573				thread_last_reads: thread_reads,
2574				quiet_round,
2575				since,
2576			},
2577			in_window(since, NEXT_BATCH),
2578		)
2579	}
2580
2581	fn no_threads() -> BTreeMap<OwnedEventId, u64> { BTreeMap::new() }
2582}