1use std::collections::{HashMap, HashSet};
16
17use eyeball_im::VectorDiff;
18use itertools::Itertools as _;
19use matrix_sdk::deserialized_responses::{
20 ThreadSummaryStatus, TimelineEvent, TimelineEventKind, UnsignedEventLocation,
21};
22use ruma::{
23 EventId, MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedTransactionId, OwnedUserId, UserId,
24 events::AnySyncTimelineEvent, push::Action, serde::Raw,
25};
26use tracing::{debug, instrument, trace, warn};
27
28use super::{
29 super::{
30 controller::ObservableItemsTransactionEntry,
31 date_dividers::DateDividerAdjuster,
32 event_handler::{Flow, TimelineEventContext, TimelineEventHandler, TimelineItemPosition},
33 event_item::RemoteEventOrigin,
34 traits::RoomDataProvider,
35 },
36 ObservableItems, ObservableItemsTransaction, TimelineMetadata, TimelineReadReceiptTracking,
37 TimelineSettings,
38 metadata::EventMeta,
39};
40use crate::timeline::{
41 EmbeddedEvent, Profile, ThreadSummary, TimelineDetails, TimelineUniqueId, VirtualTimelineItem,
42 controller::TimelineFocusKind,
43 event_handler::{FailedToParseEvent, RemovedItem, TimelineAction},
44};
45
46pub(in crate::timeline) struct TimelineStateTransaction<'a, P: RoomDataProvider> {
47 pub items: ObservableItemsTransaction<'a>,
50
51 number_of_items_when_transaction_started: usize,
53
54 pub meta: TimelineMetadata,
58
59 previous_meta: &'a mut TimelineMetadata,
61
62 pub focus: &'a TimelineFocusKind,
64
65 _phantom: std::marker::PhantomData<P>,
67}
68
69impl<'a, P: RoomDataProvider> TimelineStateTransaction<'a, P> {
70 pub(super) fn new(
72 items: &'a mut ObservableItems,
73 meta: &'a mut TimelineMetadata,
74 focus: &'a TimelineFocusKind,
75 ) -> Self {
76 let previous_meta = meta;
77 let meta = previous_meta.clone();
78 let items = items.transaction();
79
80 Self {
81 number_of_items_when_transaction_started: items.len(),
82 items,
83 previous_meta,
84 meta,
85 focus,
86 _phantom: std::marker::PhantomData,
87 }
88 }
89
90 pub(super) async fn handle_remote_events_with_diffs(
92 &mut self,
93 diffs: Vec<VectorDiff<TimelineEvent>>,
94 origin: RemoteEventOrigin,
95 room_data_provider: &P,
96 settings: &TimelineSettings,
97 ) {
98 let mut date_divider_adjuster =
99 DateDividerAdjuster::new(settings.date_divider_mode.clone());
100
101 let mut cached_profiles: HashMap<OwnedUserId, Option<Profile>> = HashMap::new();
102
103 let mut recycled_timeline_ids = HashMap::new();
104
105 for diff in diffs {
106 match diff {
107 VectorDiff::Append { values: events } => {
108 for event in events {
109 let recycled_timeline_id = event
110 .event_id()
111 .and_then(|event_id| recycled_timeline_ids.remove(event_id));
112 self.handle_remote_event(
113 event,
114 TimelineItemPosition::End { origin },
115 room_data_provider,
116 settings,
117 &mut date_divider_adjuster,
118 &mut cached_profiles,
119 recycled_timeline_id,
120 )
121 .await;
122 }
123 }
124
125 VectorDiff::PushFront { value: event } => {
126 let recycled_timeline_id = event
127 .event_id()
128 .and_then(|event_id| recycled_timeline_ids.remove(event_id));
129 self.handle_remote_event(
130 event,
131 TimelineItemPosition::Start { origin },
132 room_data_provider,
133 settings,
134 &mut date_divider_adjuster,
135 &mut cached_profiles,
136 recycled_timeline_id,
137 )
138 .await;
139 }
140
141 VectorDiff::PushBack { value: event } => {
142 let recycled_timeline_id = event
143 .event_id()
144 .and_then(|event_id| recycled_timeline_ids.remove(event_id));
145 self.handle_remote_event(
146 event,
147 TimelineItemPosition::End { origin },
148 room_data_provider,
149 settings,
150 &mut date_divider_adjuster,
151 &mut cached_profiles,
152 recycled_timeline_id,
153 )
154 .await;
155 }
156
157 VectorDiff::Insert { index: event_index, value: event } => {
158 let recycled_timeline_id = event
159 .event_id()
160 .and_then(|event_id| recycled_timeline_ids.remove(event_id));
161 self.handle_remote_event(
162 event,
163 TimelineItemPosition::At { event_index, origin },
164 room_data_provider,
165 settings,
166 &mut date_divider_adjuster,
167 &mut cached_profiles,
168 recycled_timeline_id,
169 )
170 .await;
171 }
172
173 VectorDiff::Set { index: event_index, value: event } => {
174 if let Some(timeline_item_index) = self
175 .items
176 .all_remote_events()
177 .get(event_index)
178 .and_then(|meta| meta.timeline_item_index)
179 {
180 self.handle_remote_event(
181 event,
182 TimelineItemPosition::UpdateAt { timeline_item_index },
183 room_data_provider,
184 settings,
185 &mut date_divider_adjuster,
186 &mut cached_profiles,
187 None,
188 )
189 .await;
190 } else {
191 warn!(
192 event_index,
193 "Set update dropped because there wasn't any attached timeline item index."
194 );
195 }
196 }
197
198 VectorDiff::Remove { index: event_index } => {
199 if let Some((timeline_id, event_id)) =
200 self.remove_timeline_item(event_index, &mut date_divider_adjuster)
201 {
202 recycled_timeline_ids.insert(event_id, timeline_id);
203 }
204 }
205
206 VectorDiff::Clear => {
207 self.clear();
208 }
209
210 v => unimplemented!("{v:?}"),
211 }
212 }
213
214 self.adjust_date_dividers(date_divider_adjuster);
215 self.check_invariants();
216 }
217
218 async fn handle_remote_aggregation(
219 &mut self,
220 event: TimelineEvent,
221 position: TimelineItemPosition,
222 room_data_provider: &P,
223 date_divider_adjuster: &mut DateDividerAdjuster,
224 ) {
225 let raw_event = event.raw();
226
227 let deserialized = match raw_event.deserialize() {
228 Ok(deserialized) => deserialized,
229 Err(err) => {
230 warn!("Failed to deserialize timeline event: {err}");
231 return;
232 }
233 };
234
235 let sender = deserialized.sender().to_owned();
236 let timestamp = deserialized.origin_server_ts();
237 let event_id = deserialized.event_id().to_owned();
238 let txn_id = deserialized.transaction_id().map(ToOwned::to_owned);
239
240 let timeline_actions = TimelineAction::from_event(
241 deserialized,
242 raw_event,
243 room_data_provider,
244 None,
245 None,
246 None,
247 None,
248 )
249 .await;
250
251 if timeline_actions.is_empty() {
252 return;
253 }
254
255 let encryption_info = event.kind.encryption_info().cloned();
256 let sender_profile = room_data_provider.profile_from_user_id(&sender).await;
257
258 let (forwarder, forwarder_profile) = get_forwarder_info(&event, room_data_provider).await;
259
260 let mut ctx = TimelineEventContext {
261 sender,
262 sender_profile,
263 forwarder,
264 forwarder_profile,
265 timestamp,
266 read_receipts: Default::default(),
268 is_highlighted: false,
269 flow: Flow::Remote {
270 event_id: event_id.clone(),
271 raw_event: event.raw().clone(),
272 encryption_info,
273 txn_id,
274 position,
275 },
276 should_add_new_items: false,
278 };
279
280 if let [TimelineAction::AddItem { .. }] = timeline_actions.as_slice()
283 && let TimelineItemPosition::UpdateAt { timeline_item_index } = position
284 && let Some(item) = self.items.get(timeline_item_index)
285 && item
286 .as_event()
287 .map(|e| e.content().is_unable_to_decrypt() && e.event_id() == Some(&event_id))
288 .unwrap_or_default()
289 {
290 ctx.should_add_new_items = true;
292 }
293
294 for action in timeline_actions {
295 TimelineEventHandler::new(self, &ctx)
296 .handle_event(date_divider_adjuster, action, None)
297 .await;
298 }
299 }
300
301 pub(super) async fn handle_remote_aggregations(
308 &mut self,
309 diffs: Vec<VectorDiff<TimelineEvent>>,
310 origin: RemoteEventOrigin,
311 room_data_provider: &P,
312 settings: &TimelineSettings,
313 ) {
314 let mut date_divider_adjuster =
315 DateDividerAdjuster::new(settings.date_divider_mode.clone());
316
317 for diff in diffs {
318 match diff {
319 VectorDiff::Append { values: events } => {
320 for event in events {
321 self.handle_remote_aggregation(
322 event,
323 TimelineItemPosition::End { origin },
324 room_data_provider,
325 &mut date_divider_adjuster,
326 )
327 .await;
328 }
329 }
330
331 VectorDiff::PushFront { value: event } => {
332 self.handle_remote_aggregation(
333 event,
334 TimelineItemPosition::Start { origin },
335 room_data_provider,
336 &mut date_divider_adjuster,
337 )
338 .await;
339 }
340
341 VectorDiff::PushBack { value: event } => {
342 self.handle_remote_aggregation(
343 event,
344 TimelineItemPosition::End { origin },
345 room_data_provider,
346 &mut date_divider_adjuster,
347 )
348 .await;
349 }
350
351 VectorDiff::Insert { index: event_index, value: event } => {
352 self.handle_remote_aggregation(
353 event,
354 TimelineItemPosition::At { event_index, origin },
355 room_data_provider,
356 &mut date_divider_adjuster,
357 )
358 .await;
359 }
360
361 VectorDiff::Set { index: event_index, value: event } => {
362 if let Some(timeline_item_index) = self
363 .items
364 .all_remote_events()
365 .get(event_index)
366 .and_then(|meta| meta.timeline_item_index)
367 {
368 self.handle_remote_aggregation(
369 event,
370 TimelineItemPosition::UpdateAt { timeline_item_index },
371 room_data_provider,
372 &mut date_divider_adjuster,
373 )
374 .await;
375 } else if let Some(event_id) = event.event_id()
376 && let Some(meta) = self.items.all_remote_events().get_by_event_id(event_id)
377 && let Some(timeline_item_index) = meta.timeline_item_index
378 {
379 self.handle_remote_aggregation(
402 event,
403 TimelineItemPosition::UpdateAt { timeline_item_index },
404 room_data_provider,
405 &mut date_divider_adjuster,
406 )
407 .await;
408 } else {
409 warn!(
410 event_index,
411 "Set update dropped because there wasn't any attached timeline item index."
412 );
413 }
414 }
415
416 VectorDiff::Remove { .. } | VectorDiff::Clear => {
417 }
421
422 v => unimplemented!("{v:?}"),
423 }
424 }
425
426 self.adjust_date_dividers(date_divider_adjuster);
427 self.check_invariants();
428 }
429
430 fn check_invariants(&self) {
431 self.check_no_duplicate_read_receipts();
432 self.check_no_unused_unique_ids();
433 }
434
435 fn check_no_duplicate_read_receipts(&self) -> bool {
442 let mut by_user_and_thread = HashMap::new();
443 let mut duplicates = HashSet::new();
444
445 for item in self.items.iter_remotes_region().filter_map(|(_, item)| item.as_event()) {
446 if let Some(event_id) = item.event_id() {
447 for (user_id, read_receipt) in item.read_receipts() {
448 let key = (user_id, read_receipt.thread.as_str());
449 if let Some(prev_event_id) = by_user_and_thread.insert(key, event_id) {
450 duplicates.insert((key.0, key.1, prev_event_id, event_id));
451 }
452 }
453 }
454 }
455
456 if duplicates.is_empty() {
457 false
459 } else {
460 tracing::error!(
462 ?duplicates,
463 items = ?self.items,
464 "duplicate read receipts in this timeline",
465 );
466
467 true
468 }
469 }
470
471 fn check_no_unused_unique_ids(&self) {
472 let duplicates = self
473 .items
474 .iter_all_regions()
475 .duplicates_by(|(_nth, item)| item.unique_id())
476 .map(|(_nth, item)| item.unique_id())
477 .collect::<Vec<_>>();
478
479 if !duplicates.is_empty() {
480 #[cfg(any(debug_assertions, test))]
481 panic!("duplicate unique ids in this timeline: {duplicates:?}\n{:?}", self.items);
482
483 #[cfg(not(any(debug_assertions, test)))]
484 tracing::error!(
485 ?duplicates,
486 items = ?self.items,
487 "duplicate unique ids in this timeline",
488 );
489 }
490 }
491
492 fn should_add_event_item(
494 &self,
495 room_data_provider: &P,
496 settings: &TimelineSettings,
497 event: &AnySyncTimelineEvent,
498 thread_root: Option<&EventId>,
499 position: TimelineItemPosition,
500 ) -> bool {
501 let rules = room_data_provider.room_version_rules();
502
503 if !(settings.event_filter)(event, &rules) {
504 return false;
506 }
507
508 match &self.focus {
509 TimelineFocusKind::PinnedEvents { .. } => {
510 true
512 }
513
514 TimelineFocusKind::Event { .. } => {
515 let origin = match position {
520 TimelineItemPosition::End { origin }
521 | TimelineItemPosition::Start { origin }
522 | TimelineItemPosition::At { origin, .. } => origin,
523
524 TimelineItemPosition::UpdateAt { timeline_item_index: idx } => self
525 .items
526 .get(idx)
527 .and_then(|item| item.as_event()?.as_remote())
528 .map_or(RemoteEventOrigin::Unknown, |item| item.origin),
529 };
530
531 match origin {
532 RemoteEventOrigin::Sync | RemoteEventOrigin::Unknown => false,
534 RemoteEventOrigin::Cache | RemoteEventOrigin::Pagination => true,
535 }
536 }
537
538 TimelineFocusKind::Live { hide_threaded_events, .. } => {
539 thread_root.is_none() || !hide_threaded_events
542 }
543
544 TimelineFocusKind::Thread { root_event_id, .. } => {
545 event.event_id() == root_event_id
547 || thread_root.as_ref().is_some_and(|r| r == root_event_id)
548 }
549 }
550 }
551
552 fn can_show_read_receipts(
555 &self,
556 settings: &TimelineSettings,
557 event: &AnySyncTimelineEvent,
558 ) -> bool {
559 match event {
560 AnySyncTimelineEvent::State(_) => {
561 matches!(settings.track_read_receipts, TimelineReadReceiptTracking::AllEvents)
562 }
563 AnySyncTimelineEvent::MessageLike(_) => {
564 !matches!(settings.track_read_receipts, TimelineReadReceiptTracking::Disabled)
565 }
566 }
567 }
568
569 async fn maybe_add_error_item(
573 &mut self,
574 position: TimelineItemPosition,
575 room_data_provider: &P,
576 raw: &Raw<AnySyncTimelineEvent>,
577 deserialization_error: serde_json::Error,
578 settings: &TimelineSettings,
579 ) -> Option<(
580 OwnedEventId,
581 OwnedUserId,
582 MilliSecondsSinceUnixEpoch,
583 Option<OwnedTransactionId>,
584 Vec<TimelineAction>,
585 Option<OwnedEventId>,
586 bool,
587 bool,
588 )> {
589 let state_key: Option<String> = raw.get_field("state_key").ok().flatten();
590
591 let event_type = if let Some(state_key) = state_key {
598 raw.get_field("type")
599 .ok()
600 .flatten()
601 .map(|event_type| FailedToParseEvent::State { event_type, state_key })
602 } else {
603 raw.get_field("type").ok().flatten().map(FailedToParseEvent::MsgLike)
604 };
605
606 let event_id: Option<OwnedEventId> = raw.get_field("event_id").ok().flatten();
607 let Some(event_id) = event_id else {
608 warn!(
610 ?event_type,
611 "Failed to deserialize timeline event (with no ID): {deserialization_error}"
612 );
613 return None;
614 };
615
616 let sender: Option<OwnedUserId> = raw.get_field("sender").ok().flatten();
617 let origin_server_ts: Option<MilliSecondsSinceUnixEpoch> =
618 raw.get_field("origin_server_ts").ok().flatten();
619
620 match (sender, origin_server_ts, event_type) {
621 (Some(sender), Some(origin_server_ts), Some(event_type))
622 if settings.add_failed_to_parse =>
623 {
624 #[derive(serde::Deserialize)]
627 struct Unsigned {
628 transaction_id: Option<OwnedTransactionId>,
629 }
630
631 let transaction_id: Option<OwnedTransactionId> = raw
632 .get_field::<Unsigned>("unsigned")
633 .ok()
634 .flatten()
635 .and_then(|unsigned| unsigned.transaction_id);
636
637 Some((
640 event_id,
641 sender,
642 origin_server_ts,
643 transaction_id,
644 vec![TimelineAction::failed_to_parse(event_type, deserialization_error)],
645 None,
646 true,
647 true,
648 ))
649 }
650
651 (sender, origin_server_ts, event_type) => {
652 warn!(
655 ?event_type,
656 ?event_id,
657 "Failed to deserialize timeline event: {deserialization_error}"
658 );
659
660 self.add_or_update_remote_event(
663 EventMeta::new(event_id, sender.as_deref(), false, false, None),
664 sender.as_deref(),
665 origin_server_ts,
666 position,
667 room_data_provider,
668 settings,
669 )
670 .await;
671 None
672 }
673 }
674 }
675
676 #[instrument(skip(self, room_data_provider))]
679 async fn fetch_latest_thread_reply(
680 &mut self,
681 event_id: &EventId,
682 room_data_provider: &P,
683 ) -> Option<Box<EmbeddedEvent>> {
684 let event = RoomDataProvider::load_event(room_data_provider, event_id)
685 .await
686 .inspect_err(|err| {
687 warn!("Failed to load thread latest event: {err}");
688 })
689 .ok()?;
690
691 EmbeddedEvent::try_from_timeline_event(event, room_data_provider, &self.meta)
692 .await
693 .inspect_err(|err| {
694 warn!("Failed to extract thread latest event into a timeline item content: {err}");
695 })
696 .ok()
697 .flatten()
698 .map(Box::new)
699 }
700
701 #[allow(clippy::too_many_arguments)]
705 pub(super) async fn handle_remote_event(
706 &mut self,
707 event: TimelineEvent,
708 position: TimelineItemPosition,
709 room_data_provider: &P,
710 settings: &TimelineSettings,
711 date_divider_adjuster: &mut DateDividerAdjuster,
712 profiles: &mut HashMap<OwnedUserId, Option<Profile>>,
713 recycled_timeline_id: Option<TimelineUniqueId>,
714 ) -> RemovedItem {
715 let is_highlighted =
716 event.push_actions().is_some_and(|actions| actions.iter().any(Action::is_highlight));
717
718 let thread_summary = if let ThreadSummaryStatus::Some(ref summary) = event.thread_summary {
719 let latest_reply_item = if let Some(ref latest_reply) = summary.latest_reply {
720 self.fetch_latest_thread_reply(latest_reply, room_data_provider).await
721 } else {
722 None
723 };
724
725 Some(ThreadSummary {
726 latest_event: TimelineDetails::from_initial_value(latest_reply_item),
727 num_replies: summary.num_replies,
728 })
729 } else {
730 None
731 };
732
733 let encryption_info = event.kind.encryption_info().cloned();
734
735 let bundled_edit_encryption_info = event.kind.unsigned_encryption_map().and_then(|map| {
736 map.get(&UnsignedEventLocation::RelationsReplace)?.encryption_info().cloned()
737 });
738
739 let (forwarder, forwarder_profile) = get_forwarder_info(&event, room_data_provider).await;
740
741 let (raw, utd_info) = match event.kind {
742 TimelineEventKind::UnableToDecrypt { utd_info, event } => (event, Some(utd_info)),
743 _ => (event.kind.into_raw(), None),
744 };
745
746 let (
747 event_id,
748 sender,
749 timestamp,
750 txn_id,
751 timeline_actions,
752 thread_root,
753 should_add,
754 can_show_read_receipts,
755 ) = match raw.deserialize() {
756 Ok(event) => {
758 let (in_reply_to, thread_root) = self.meta.process_event_relations(
759 &event,
760 &raw,
761 bundled_edit_encryption_info,
762 &self.items,
763 self.focus.is_thread(),
764 );
765
766 let should_add = self.should_add_event_item(
767 room_data_provider,
768 settings,
769 &event,
770 thread_root.as_deref(),
771 position,
772 );
773
774 let can_show_read_receipts = self.can_show_read_receipts(settings, &event);
775
776 (
777 event.event_id().to_owned(),
778 event.sender().to_owned(),
779 event.origin_server_ts(),
780 event.transaction_id().map(ToOwned::to_owned),
781 TimelineAction::from_event(
782 event,
783 &raw,
784 room_data_provider,
785 utd_info
786 .map(|utd_info| (utd_info, self.meta.unable_to_decrypt_hook.as_ref())),
787 in_reply_to,
788 thread_root.clone(),
789 thread_summary,
790 )
791 .await,
792 thread_root,
793 should_add,
794 can_show_read_receipts,
795 )
796 }
797
798 Err(e) => {
800 if let Some(tuple) =
801 self.maybe_add_error_item(position, room_data_provider, &raw, e, settings).await
802 {
803 tuple
804 } else {
805 return false;
806 }
807 }
808 };
809
810 self.add_or_update_remote_event(
813 EventMeta::new(
814 event_id.clone(),
815 Some(&sender),
816 should_add,
817 can_show_read_receipts,
818 thread_root,
819 ),
820 Some(&sender),
821 Some(timestamp),
822 position,
823 room_data_provider,
824 settings,
825 )
826 .await;
827
828 let item_added = if !timeline_actions.is_empty() {
830 let sender_profile = if let Some(profile) = profiles.get(&sender) {
831 profile.clone()
832 } else {
833 let profile = room_data_provider.profile_from_user_id(&sender).await;
834 profiles.insert(sender.clone(), profile.clone());
835 profile
836 };
837
838 let mut item_added = false;
839 let ctx = TimelineEventContext {
840 sender,
841 sender_profile,
842 forwarder,
843 forwarder_profile,
844 timestamp,
845 read_receipts: if settings.track_read_receipts.is_enabled()
846 && should_add
847 && can_show_read_receipts
848 {
849 self.meta.read_receipts.compute_event_receipts(
850 &event_id,
851 &mut self.items,
852 matches!(position, TimelineItemPosition::End { .. }),
853 )
854 } else {
855 Default::default()
856 },
857
858 is_highlighted,
859 flow: Flow::Remote {
860 event_id: event_id.clone(),
861 raw_event: raw,
862 encryption_info,
863 txn_id,
864 position,
865 },
866 should_add_new_items: should_add,
867 };
868 let recycled_timeline_id = recycled_timeline_id.filter(|_| timeline_actions.len() == 1);
875
876 for action in timeline_actions {
877 item_added |= TimelineEventHandler::new(self, &ctx)
878 .handle_event(date_divider_adjuster, action, recycled_timeline_id.clone())
879 .await;
880 }
881 item_added
882 } else {
883 false
885 };
886
887 let mut item_removed = false;
888
889 if !item_added {
890 trace!("No new item added");
891
892 if let TimelineItemPosition::UpdateAt { timeline_item_index } = position {
893 trace!("Removing UTD that was successfully retried");
896 self.items.remove(timeline_item_index);
897 item_removed = true;
898 }
899 }
900
901 item_removed
902 }
903
904 fn remove_timeline_item(
906 &mut self,
907 event_index: usize,
908 day_divider_adjuster: &mut DateDividerAdjuster,
909 ) -> Option<(TimelineUniqueId, OwnedEventId)> {
910 day_divider_adjuster.mark_used();
911
912 let mut recycled_timeline_id = None;
921
922 if let Some(event_meta) = self.items.all_remote_events().get(event_index) {
924 if let Some(timeline_item_index) = event_meta.timeline_item_index {
926 let event_id = event_meta.event_id.clone();
927 let timeline_item = self.items.remove(timeline_item_index);
928 recycled_timeline_id = Some((timeline_item.unique_id().clone(), event_id));
929 }
930
931 self.items.remove_remote_event(event_index);
933 }
934
935 recycled_timeline_id
936 }
937
938 pub(super) fn clear(&mut self) {
939 if self.items.has_local() {
945 self.items.for_each(|entry| {
947 if entry.is_remote_event()
948 || entry.as_virtual().is_some_and(|vitem| match vitem {
949 VirtualTimelineItem::DateDivider(_) => false,
950 VirtualTimelineItem::ReadMarker | VirtualTimelineItem::TimelineStart => {
951 true
952 }
953 })
954 {
955 ObservableItemsTransactionEntry::remove_timeline_index_and_remote_event(entry);
956 }
957 });
958
959 let mut idx = 0;
961 while idx < self.items.len() {
962 if self.items[idx].is_date_divider()
963 && self.items.get(idx + 1).is_none_or(|item| item.is_date_divider())
964 {
965 self.items.remove(idx);
966 } else {
968 idx += 1;
969 }
970 }
971 } else {
972 self.items.clear();
973 }
974
975 self.meta.clear();
976
977 debug!(remaining_items = self.items.len(), "Timeline cleared");
978 }
979
980 #[instrument(skip_all)]
981 pub(super) fn set_fully_read_event(&mut self, fully_read_event_id: OwnedEventId) {
982 if self.meta.fully_read_event.as_ref().is_some_and(|id| *id == fully_read_event_id) {
984 return;
985 }
986
987 self.meta.fully_read_event = Some(fully_read_event_id);
988 self.meta.update_read_marker(&mut self.items);
989 }
990
991 pub(super) fn commit(self) {
992 let previous_number_of_items = self.number_of_items_when_transaction_started;
994 let next_number_of_items = self.items.len();
995
996 if previous_number_of_items != next_number_of_items {
997 let count = self
998 .meta
999 .subscriber_skip_count
1000 .compute_next(previous_number_of_items, next_number_of_items);
1001 self.meta
1002 .subscriber_skip_count
1003 .update(count, matches!(self.focus, TimelineFocusKind::Live { .. }));
1004 }
1005
1006 *self.previous_meta = self.meta;
1008
1009 self.items.commit();
1010 }
1011
1012 async fn add_or_update_remote_event(
1017 &mut self,
1018 event_meta: EventMeta,
1019 sender: Option<&UserId>,
1020 timestamp: Option<MilliSecondsSinceUnixEpoch>,
1021 position: TimelineItemPosition,
1022 room_data_provider: &P,
1023 settings: &TimelineSettings,
1024 ) {
1025 let event_id = event_meta.event_id.clone();
1026
1027 match position {
1028 TimelineItemPosition::Start { .. } => self.items.push_front_remote_event(event_meta),
1029
1030 TimelineItemPosition::End { .. } => {
1031 self.items.push_back_remote_event(event_meta);
1032 }
1033
1034 TimelineItemPosition::At { event_index, .. } => {
1035 self.items.insert_remote_event(event_index, event_meta);
1036 }
1037
1038 TimelineItemPosition::UpdateAt { .. } => {
1039 if let Some(event) =
1040 self.items.get_remote_event_by_event_id_mut(&event_meta.event_id)
1041 && (event.visible != event_meta.visible
1042 || event.can_show_read_receipts != event_meta.can_show_read_receipts)
1043 {
1044 event.visible = event_meta.visible;
1045 event.can_show_read_receipts = event_meta.can_show_read_receipts;
1046
1047 if settings.track_read_receipts.is_enabled() {
1048 self.maybe_update_read_receipts_of_prev_event(&event_meta.event_id);
1051 }
1052 }
1053 }
1054 }
1055
1056 if settings.track_read_receipts.is_enabled()
1057 && matches!(
1058 position,
1059 TimelineItemPosition::Start { .. }
1060 | TimelineItemPosition::End { .. }
1061 | TimelineItemPosition::At { .. }
1062 )
1063 {
1064 self.load_read_receipts_for_event(&event_id, room_data_provider).await;
1065
1066 self.maybe_add_implicit_read_receipt(&event_id, sender, timestamp);
1067 }
1068 }
1069
1070 pub(super) fn adjust_date_dividers(&mut self, mut adjuster: DateDividerAdjuster) {
1071 adjuster.run(&mut self.items, &mut self.meta);
1072 }
1073
1074 pub(super) fn mark_all_events_as_encrypted(&mut self) {
1078 for idx in 0..self.items.len() {
1079 let item = &self.items[idx];
1080
1081 if let Some(event) = item.as_event() {
1082 if event.is_room_encrypted {
1083 continue;
1084 }
1085
1086 let mut cloned_event = event.clone();
1087 cloned_event.is_room_encrypted = true;
1088
1089 let item = item.with_kind(cloned_event);
1091 self.items.replace(idx, item);
1092 }
1093 }
1094 }
1095}
1096
1097async fn get_forwarder_info<P: RoomDataProvider>(
1110 event: &TimelineEvent,
1111 room_data_provider: &P,
1112) -> (Option<OwnedUserId>, Option<Profile>) {
1113 let forwarder = event
1114 .kind
1115 .encryption_info()
1116 .and_then(|info| info.forwarder.as_ref())
1117 .map(|info| info.user_id.clone());
1118
1119 let forwarder_profile = if let Some(ref forwarder_id) = forwarder {
1120 Some(room_data_provider.profile_from_user_id(forwarder_id).await)
1121 } else {
1122 None
1123 };
1124
1125 (forwarder, forwarder_profile.flatten())
1126}
1127
1128#[cfg(test)]
1129mod tests {
1130 use std::sync::Arc;
1131
1132 use matrix_sdk::{Room, test_utils::mocks::MatrixMockServer};
1133 use matrix_sdk_test::async_test;
1134 use ruma::{
1135 MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedUserId,
1136 events::receipt::{Receipt, ReceiptThread},
1137 owned_event_id, owned_user_id, room_id,
1138 room_version_rules::RoomVersionRules,
1139 };
1140
1141 use crate::timeline::{
1142 EventTimelineItem, MsgLikeContent, TimelineDetails, TimelineItem, TimelineItemContent,
1143 TimelineItemKind, TimelineUniqueId,
1144 controller::{
1145 ObservableItems, TimelineFocusKind, TimelineMetadata, TimelineStateTransaction,
1146 },
1147 event_item::{EventTimelineItemKind, RemoteEventOrigin, RemoteEventTimelineItem},
1148 };
1149
1150 #[async_test]
1151 async fn test_detects_duplicate_read_receipts_in_same_receipt_thread() {
1152 let server = MatrixMockServer::new().await;
1153 let client = server.client_builder().build().await;
1154 let room_id = room_id!("!r0");
1155
1156 let _ = server.sync_joined_room(&client, room_id).await;
1157
1158 let event_cache = client.event_cache();
1159 event_cache.subscribe().unwrap();
1160
1161 let mut items = create_items_with_receipts(vec![
1163 (owned_user_id!("@user1:s.co"), create_receipt(ReceiptThread::Unthreaded)),
1164 (owned_user_id!("@user1:s.co"), create_receipt(ReceiptThread::Unthreaded)),
1165 ]);
1166
1167 let user_id = owned_user_id!("@foo:s.co");
1169 let mut meta = TimelineMetadata::new(user_id, RoomVersionRules::V12, None, None, true);
1170 let focus = TimelineFocusKind::Live {
1171 hide_threaded_events: false,
1172 event_cache: event_cache.room(room_id).await.unwrap().0,
1173 };
1174 let transaction: TimelineStateTransaction<'_, Room> =
1175 TimelineStateTransaction::new(&mut items, &mut meta, &focus);
1176
1177 let dups = transaction.check_no_duplicate_read_receipts();
1178
1179 assert!(dups);
1181 }
1182
1183 #[async_test]
1184 async fn test_if_there_are_no_duplicate_receipts_we_report_no_duplicates() {
1185 let server = MatrixMockServer::new().await;
1186 let client = server.client_builder().build().await;
1187 let room_id = room_id!("!r0");
1188
1189 let _ = server.sync_joined_room(&client, room_id).await;
1190
1191 let event_cache = client.event_cache();
1192 event_cache.subscribe().unwrap();
1193
1194 let mut items = create_items_with_receipts(vec![
1196 (owned_user_id!("@user1:s.co"), create_receipt(ReceiptThread::Unthreaded)),
1197 (owned_user_id!("@user2:s.co"), create_receipt(ReceiptThread::Unthreaded)),
1198 ]);
1199
1200 let user_id = owned_user_id!("@foo:s.co");
1202 let mut meta = TimelineMetadata::new(user_id, RoomVersionRules::V12, None, None, true);
1203 let focus = TimelineFocusKind::Live {
1204 hide_threaded_events: false,
1205 event_cache: event_cache.room(room_id).await.unwrap().0,
1206 };
1207 let transaction: TimelineStateTransaction<'_, Room> =
1208 TimelineStateTransaction::new(&mut items, &mut meta, &focus);
1209
1210 let dups = transaction.check_no_duplicate_read_receipts();
1211
1212 assert!(!dups);
1214 }
1215
1216 #[async_test]
1217 async fn test_if_there_are_receipts_for_different_receipt_threads_we_report_no_duplicates() {
1218 let server = MatrixMockServer::new().await;
1219 let client = server.client_builder().build().await;
1220 let room_id = room_id!("!r0");
1221
1222 let _ = server.sync_joined_room(&client, room_id).await;
1223
1224 let event_cache = client.event_cache();
1225 event_cache.subscribe().unwrap();
1226
1227 let mut items = create_items_with_receipts(vec![
1229 (owned_user_id!("@user1:s.co"), create_receipt(ReceiptThread::Unthreaded)),
1230 (
1231 owned_user_id!("@user1:s.co"),
1232 create_receipt(ReceiptThread::Thread(owned_event_id!("$thread_root"))),
1233 ),
1234 ]);
1235
1236 let user_id = owned_user_id!("@foo:s.co");
1238 let mut meta = TimelineMetadata::new(user_id, RoomVersionRules::V12, None, None, true);
1239 let focus = TimelineFocusKind::Live {
1240 hide_threaded_events: false,
1241 event_cache: event_cache.room(room_id).await.unwrap().0,
1242 };
1243 let transaction: TimelineStateTransaction<'_, Room> =
1244 TimelineStateTransaction::new(&mut items, &mut meta, &focus);
1245
1246 let dups = transaction.check_no_duplicate_read_receipts();
1247
1248 assert!(!dups);
1250 }
1251
1252 fn create_receipt(receipt_thread: ReceiptThread) -> Receipt {
1253 let mut receipt = Receipt::new(MilliSecondsSinceUnixEpoch::now());
1254 receipt.thread = receipt_thread;
1255 receipt
1256 }
1257
1258 fn create_items_with_receipts(receipts: Vec<(OwnedUserId, Receipt)>) -> ObservableItems {
1259 let mut items = ObservableItems::new();
1260 {
1261 let mut t = items.transaction();
1262
1263 for (num, receipt) in receipts.into_iter().enumerate() {
1264 let event_id = OwnedEventId::try_from(format!("$event-{num}")).unwrap();
1265 let timeline_id = format!("timeline-{num}");
1266
1267 t.push_front(event_with_receipt_item(event_id, timeline_id, receipt), None);
1268 }
1269 t.commit();
1270 }
1271 items
1272 }
1273
1274 fn event_with_receipt_item(
1275 event_id: OwnedEventId,
1276 timeline_id: String,
1277 receipt: (OwnedUserId, Receipt),
1278 ) -> Arc<TimelineItem> {
1279 let event_kind = EventTimelineItemKind::Remote(RemoteEventTimelineItem {
1280 event_id,
1281 transaction_id: None,
1282 read_receipts: [receipt].into(),
1283 is_own: false,
1284 is_highlighted: false,
1285 encryption_info: None,
1286 original_json: None,
1287 latest_edit_json: None,
1288 origin: RemoteEventOrigin::Sync,
1289 });
1290
1291 let event_timeline_item = EventTimelineItem::new(
1292 owned_user_id!("@u:s.co"),
1293 TimelineDetails::Pending,
1294 None,
1295 None,
1296 MilliSecondsSinceUnixEpoch::now(),
1297 TimelineItemContent::MsgLike(MsgLikeContent::redacted()),
1298 event_kind,
1299 false,
1300 );
1301
1302 TimelineItem::new(
1303 TimelineItemKind::Event(event_timeline_item),
1304 TimelineUniqueId(timeline_id),
1305 )
1306 }
1307}