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 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#[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 match (stable, unstable) {
192 | (_, true) => Self::Unstable,
193 | (true, _) => Self::Stable,
194 | _ => Self::Off,
195 }
196 }
197}
198
199#[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 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 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 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 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 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 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 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 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 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 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 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
1057fn 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 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 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 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 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 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 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 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
1968fn in_window(since: u64, next_batch: u64) -> impl Fn(u64) -> bool + Copy {
1974 move |count| count > since && count <= next_batch
1975}
1976
1977fn 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 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 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 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 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 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 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 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 if !matches!(content.membership, MembershipState::Join | MembershipState::Invite) {
2423 return None;
2424 }
2425
2426 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 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 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 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 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 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}