Skip to main content

matrix_sdk_ui/timeline/controller/
state_transaction.rs

1// Copyright 2025 The Matrix.org Foundation C.I.C.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// A vector transaction over the items themselves. Holds temporary state
48    /// until committed.
49    pub items: ObservableItemsTransaction<'a>,
50
51    /// Number of items when the transaction has been created/has started.
52    number_of_items_when_transaction_started: usize,
53
54    /// A clone of the previous meta, that we're operating on during the
55    /// transaction, and that will be committed to the previous meta location in
56    /// [`Self::commit`].
57    pub meta: TimelineMetadata,
58
59    /// Pointer to the previous meta, only used during [`Self::commit`].
60    previous_meta: &'a mut TimelineMetadata,
61
62    /// The kind of focus of this timeline.
63    pub focus: &'a TimelineFocusKind,
64
65    /// Phantom data for type parameter.
66    _phantom: std::marker::PhantomData<P>,
67}
68
69impl<'a, P: RoomDataProvider> TimelineStateTransaction<'a, P> {
70    /// Create a new [`TimelineStateTransaction`].
71    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    /// Handle updates on events as [`VectorDiff`]s.
91    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            // These are not used when handling an aggregation.
267            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            // This field is not used when handling an aggregation.
277            should_add_new_items: false,
278        };
279
280        // FIXME: Continuation of the hackjob to get UTDs for focused timelines
281        // working from `handle_remote_aggregations()`.
282        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            // Except when this is an UTD transitioning into a decrypted event.
291            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    /// Handle a set of live remote aggregations on events as [`VectorDiff`]s.
302    ///
303    /// This is like `handle_remote_events`, with two key differences:
304    /// - it only applies to aggregated events, not all the sync events.
305    /// - it will also not add the events to the `all_remote_events` array
306    ///   itself.
307    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                        // FIXME: This branch is a complete hackjob.
380                        //
381                        // The reason being is that this branch is here to handle UTD -> Decrypted
382                        // event remplacements for focused timelines. But this transition should
383                        // naturally happen the same way it happens for unfocused timelines.
384                        //
385                        // Why it doesn't work here? Because the event cache fires out a
386                        // VectorDiff::Set with an index that matches to the cache's view of the
387                        // timeline, which is unfiltered, while the focused timeline will only show
388                        // i.e. pinned events.
389                        //
390                        // The `test_pinned_events_are_decrypted_after_recovering` integration test
391                        // showcases this. The event cache fires out the `Set` with an index of 7,
392                        // but the timeline with the PinnedEvents focus has only 4 items.
393                        //
394                        // This hackjob continues in the `handle_remote_aggregation()` method as we
395                        // can't just handle any `TimelineAction::AddItem` due to:
396                        //  https://github.com/matrix-org/matrix-rust-sdk/pull/4645
397                        //
398                        // Doing so breaks the `test_new_pinned_events_are_not_added_on_sync` test.
399                        //
400                        // Relevant issue: https://github.com/matrix-org/matrix-rust-sdk/issues/5954.
401                        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                    // Do nothing. An aggregated redaction comes with a
418                    // redaction event, or as a redacted event in the first
419                    // place.
420                }
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    /// Check whether any read receipts are duplicated. If so, log an error
436    /// message.
437    ///
438    /// Returns true if some receipts are duplicated, and false if all is fine.
439    /// This return value is used for testing, and ignored in production (the
440    /// error log is sufficient).
441    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            // No duplicates - all good
458            false
459        } else {
460            // Some duplicates - log an error and return true (used in tests)
461            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    /// Whether the event should be added to the timeline as a new item.
493    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            // The user filtered out the event.
505            return false;
506        }
507
508        match &self.focus {
509            TimelineFocusKind::PinnedEvents { .. } => {
510                // The pinned events timeline only receives updates for, well, pinned events.
511                true
512            }
513
514            TimelineFocusKind::Event { .. } => {
515                // For event-focused timelines, thread filtering is now handled in the
516                // event cache layer. We accept all events from pagination.
517
518                // Retrieve the origin of the event.
519                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                    // Never add any item to a focused timeline when the item comes from sync.
533                    RemoteEventOrigin::Sync | RemoteEventOrigin::Unknown => false,
534                    RemoteEventOrigin::Cache | RemoteEventOrigin::Pagination => true,
535                }
536            }
537
538            TimelineFocusKind::Live { hide_threaded_events, .. } => {
539                // If the timeline's filtering out in-thread events, don't add items for
540                // threaded events.
541                thread_root.is_none() || !hide_threaded_events
542            }
543
544            TimelineFocusKind::Thread { root_event_id, .. } => {
545                // Add new items only for the thread root and the thread replies.
546                event.event_id() == root_event_id
547                    || thread_root.as_ref().is_some_and(|r| r == root_event_id)
548            }
549        }
550    }
551
552    /// Whether this event can show read receipts, or if they should be moved
553    /// to the previous event.
554    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    /// After a deserialization error, adds a failed-to-parse item to the
570    /// timeline if configured to do so, or logs the error (and optionally
571    /// save metadata) if not.
572    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        // A state event is an event that has a state key. Note that the two branches
592        // differ because the inferred return type for `get_field` is different
593        // in each case.
594        //
595        // If this was a state event but it didn't include a state_key, we'll assume it
596        // was a msg-like, because we can't do much more.
597        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            // If the event doesn't even have an event ID, we can't do anything with it.
609            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                // We have sufficient information to show an item in the timeline, and we've
625                // been requested to show it, let's do it.
626                #[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                // The event can be partially deserialized, and it is allowed to be added to
638                // the timeline.
639                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                // We either lack information for rendering an item, or we've been requested not
653                // to show it. Save it into the metadata and return.
654                warn!(
655                    ?event_type,
656                    ?event_id,
657                    "Failed to deserialize timeline event: {deserialization_error}"
658                );
659
660                // Remember the event before returning prematurely.
661                // See [`ObservableItems::all_remote_events`].
662                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    // Attempt to load a thread's latest reply as an embedded timeline item, either
677    // using the event cache or the storage.
678    #[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    /// Handle a remote event.
702    ///
703    /// Returns whether an item has been removed from the timeline.
704    #[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            // Classical path: the event is valid, can be deserialized, everything is alright.
757            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            // The event seems invalid…
799            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        // Remember the event.
811        // See [`ObservableItems::all_remote_events`].
812        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        // Handle the event to create or update a timeline item.
829        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            // A recycled timeline ID carries the TimelineUniqueId from a previously
869            // removed item (see VectorDiff::Remove), so that when the same event is
870            // re-added in the same diff batch the UI sees a stable identifier.
871            // It is only applicable when the event produces a single AddItem action;
872            // with multiple actions (e.g. beacon replace) there
873            // is no single item to associate it with, so it's safe to ignore.
874            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            // No item has been added to the timeline.
884            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                // If add was not called, that means the UTD event is one that
894                // wouldn't normally be visible. Remove it.
895                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    /// Remove one timeline item by its `event_index`.
905    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        // We need to be careful here.
913        //
914        // We must first remove the timeline item, which will update the mapping between
915        // remote events and timeline items. Removing the timeline item will “unlink”
916        // this mapping as the remote event will be updated to map to nothing. Only
917        // after that, we can remove the remote event. Doing this in the other order
918        // will update the mapping twice, and will result in a corrupted state.
919
920        let mut recycled_timeline_id = None;
921
922        // Remove the timeline item first.
923        if let Some(event_meta) = self.items.all_remote_events().get(event_index) {
924            // Fetch the `timeline_item_index` associated to the remote event.
925            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            // Now we can remove the remote event.
932            self.items.remove_remote_event(event_index);
933        }
934
935        recycled_timeline_id
936    }
937
938    pub(super) fn clear(&mut self) {
939        // By first checking if there are any local echoes first, we do a bit
940        // more work in case some are found, but it should be worth it because
941        // there will often not be any, and only emitting a single
942        // `VectorDiff::Clear` should be much more efficient to process for
943        // subscribers.
944        if self.items.has_local() {
945            // Remove all remote events and virtual items that aren't date dividers.
946            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            // Remove adjacent date dividers.
960            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                    // don't increment idx because all elements have shifted
967                } 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        // A similar event has been handled already. We can ignore it.
983        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        // Update the `subscriber_skip_count` value.
993        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        // Replace the pointer to the previous meta with the new one.
1007        *self.previous_meta = self.meta;
1008
1009        self.items.commit();
1010    }
1011
1012    /// Add or update a remote event in the
1013    /// [`ObservableItems::all_remote_events`] collection.
1014    ///
1015    /// This method also adjusts read receipt if needed.
1016    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                        // Since the event's visibility changed, we need to update the read
1049                        // receipts of the previous visible event.
1050                        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    /// This method replaces the `is_room_encrypted` value for all timeline
1075    /// items to its updated version and creates a `VectorDiff::Set` operation
1076    /// for each item which will be added to this transaction.
1077    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                // Replace the existing item with a new version with the right encryption flag
1090                let item = item.with_kind(cloned_event);
1091                self.items.replace(idx, item);
1092            }
1093        }
1094    }
1095}
1096
1097/// Retrieves the forwarder information for a given timeline event.
1098///
1099/// # Parameters
1100///
1101/// - `event`: The timeline event to extract forwarder information from.
1102/// - `room_data_provider`: A reference to the room data provider.
1103///
1104/// # Returns
1105///
1106/// A tuple containing:
1107/// - `Option<OwnedUserId>`: The user ID of the forwarder, if available.
1108/// - `Option<Profile>`: The profile of the forwarder, if available.
1109async 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        // Given a timeline with clashing receipts
1162        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        // When we check for duplicates
1168        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        // Then some duplicates were found
1180        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        // Given a timeline with receipts, but no clashes (users are different)
1195        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        // When we check for duplicates
1201        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        // Then no duplicates were found
1213        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        // Given a timeline with receipts, but no clashes (users are different)
1228        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        // When we check for duplicates
1237        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        // Then no duplicates were found
1249        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}