Skip to main content

matrix_sdk_ui/timeline/controller/
mod.rs

1// Copyright 2023 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::{
16    collections::BTreeSet,
17    fmt,
18    ops::Deref,
19    sync::{Arc, OnceLock},
20};
21
22use as_variant::as_variant;
23use eyeball_im::{VectorDiff, VectorSubscriberStream};
24use eyeball_im_util::vector::{FilterMap, VectorObserverExt};
25use futures_core::Stream;
26use futures_util::future::try_join_all;
27use imbl::{HashSet, Vector};
28use matrix_sdk::{
29    deserialized_responses::{ThreadSummary, TimelineEvent},
30    event_cache::{
31        DecryptionRetryRequest, EventCache, EventFocusedCache, PaginationStatus, PinnedEventsCache,
32        RoomEventCache, Subscriber as EventCacheSubscriber, ThreadEventCache,
33        ThreadEventCacheUpdate,
34    },
35    send_queue::{LocalEcho, LocalEchoContent, RoomSendQueueUpdate, SendHandle},
36    task_monitor::BackgroundTaskHandle,
37};
38use ruma::{
39    EventId, MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedTransactionId, OwnedUserId, RoomId,
40    TransactionId, UserId,
41    api::client::receipt::create_receipt::v3::ReceiptType as SendReceiptType,
42    events::{
43        AnyMessageLikeEventContent, AnySyncMessageLikeEvent, AnySyncTimelineEvent,
44        MessageLikeEventType,
45        poll::unstable_start::UnstablePollStartEventContent,
46        reaction::ReactionEventContent,
47        receipt::{Receipt, ReceiptEventContent, ReceiptThread, ReceiptType},
48        relation::{Annotation, RelationType},
49        room::message::{MessageType, Relation},
50    },
51    room_version_rules::RoomVersionRules,
52};
53use tokio::sync::{RwLock, RwLockWriteGuard};
54use tracing::{
55    Instrument as _, Span, debug, error, field::debug, info, info_span, instrument, trace, warn,
56};
57
58pub(super) use self::{
59    metadata::{RelativePosition, TimelineMetadata},
60    observable_items::{
61        AllRemoteEvents, ObservableItems, ObservableItemsEntry, ObservableItemsTransaction,
62        ObservableItemsTransactionEntry,
63    },
64    state::TimelineState,
65    state_transaction::TimelineStateTransaction,
66};
67use super::{
68    DateDividerMode, EmbeddedEvent, Error, EventSendState, EventTimelineItem, InReplyToDetails,
69    MediaUploadProgress, Profile, TimelineDetails, TimelineEventItemId, TimelineFocus,
70    TimelineItem, TimelineItemContent, TimelineItemKind, TimelineReadReceiptTracking,
71    VirtualTimelineItem,
72    algorithms::{rfind_event_by_id, rfind_event_item},
73    event_item::RemoteEventOrigin,
74    item::TimelineUniqueId,
75    subscriber::TimelineSubscriber,
76    traits::RoomDataProvider,
77};
78use crate::{
79    timeline::{
80        MsgLikeContent, MsgLikeKind, Room, SendTarget, TimelineEventFilterFn,
81        TimelineEventFocusThreadMode,
82        algorithms::rfind_event_by_item_id,
83        controller::decryption_retry_task::compute_redecryption_candidates,
84        date_dividers::DateDividerAdjuster,
85        event_item::TimelineItemHandle,
86        tasks::{event_focused_task, pinned_events_task, thread_updates_task},
87    },
88    unable_to_decrypt_hook::UtdHookManager,
89};
90
91pub(in crate::timeline) mod aggregations;
92mod decryption_retry_task;
93mod metadata;
94mod observable_items;
95mod read_receipts;
96mod state;
97mod state_transaction;
98
99pub(super) use aggregations::*;
100pub(super) use decryption_retry_task::{CryptoDropHandles, spawn_crypto_tasks};
101use matrix_sdk_base::{CallIntentConsensus, RoomInfo};
102
103/// The outcome of [`TimelineController::should_send_receipt`].
104pub(super) enum SendReceiptDecision {
105    /// No read receipt should be sent.
106    DoNotSend,
107
108    /// A read receipt should be sent, targeting this event.
109    ///
110    /// This may differ from the event the caller asked about, since a read
111    /// receipt should not point at one of the user's own events.
112    SendTo(OwnedEventId),
113}
114
115/// Data associated to the current timeline focus.
116///
117/// This is the private counterpart of [`TimelineFocus`], and it is an augmented
118/// version of it, including extra state that makes it useful over the lifetime
119/// of a timeline.
120#[derive(Debug)]
121pub(in crate::timeline) enum TimelineFocusKind {
122    /// The timeline receives live events from the sync.
123    Live {
124        /// Whether to hide in-thread events from the timeline.
125        hide_threaded_events: bool,
126
127        /// The cache holding all the events for this focus.
128        event_cache: RoomEventCache,
129    },
130
131    /// The timeline is focused on a single event, and it can expand in one
132    /// direction or another.
133    Event {
134        /// The focused event ID.
135        focused_event_id: OwnedEventId,
136
137        /// If the focused event is part or the root of a thread, what's the
138        /// thread root?
139        ///
140        /// This is determined once when initializing the event-focused cache,
141        /// and then it won't change for the duration of this timeline.
142        thread_root: OnceLock<OwnedEventId>,
143
144        /// The thread mode to use for this event-focused timeline, which is
145        /// part of the key for the memoized event-focused cache.
146        thread_mode: TimelineEventFocusThreadMode,
147
148        /// The cache holding all the events for this focus.
149        event_cache: EventFocusedCache,
150    },
151
152    /// A live timeline for a thread.
153    Thread {
154        /// The root event for the current thread.
155        thread_id: OwnedEventId,
156
157        /// The cache holding all the events for this focus.
158        event_cache: ThreadEventCache,
159    },
160
161    PinnedEvents {
162        /// The cache holding all the events for this focus.
163        event_cache: PinnedEventsCache,
164    },
165}
166
167impl TimelineFocusKind {
168    /// Get the room ID of this timeline.
169    pub(super) fn room_id(&self) -> &RoomId {
170        match self {
171            TimelineFocusKind::Live { event_cache, .. } => event_cache.room_id(),
172            TimelineFocusKind::Thread { event_cache, .. } => event_cache.room_id(),
173            TimelineFocusKind::Event { event_cache, .. } => event_cache.room_id(),
174            TimelineFocusKind::PinnedEvents { event_cache } => event_cache.room_id(),
175        }
176    }
177    /// Returns the [`ReceiptThread`] that should be used for the current
178    /// timeline focus.
179    ///
180    /// Live and event timelines will use the unthreaded read receipt type in
181    /// general, unless they hide in-thread events, in which case they will use
182    /// the main thread.
183    pub(super) fn receipt_thread(&self) -> ReceiptThread {
184        if let Some(thread_root) = self.thread_root() {
185            ReceiptThread::Thread(thread_root.to_owned())
186        } else if self.hide_threaded_events() {
187            ReceiptThread::Main
188        } else {
189            ReceiptThread::Unthreaded
190        }
191    }
192
193    /// Whether to hide in-thread events from the timeline.
194    fn hide_threaded_events(&self) -> bool {
195        match self {
196            TimelineFocusKind::Live { hide_threaded_events, .. } => *hide_threaded_events,
197            TimelineFocusKind::Event { thread_mode, .. } => {
198                matches!(
199                    thread_mode,
200                    TimelineEventFocusThreadMode::Automatic { hide_threaded_events: true }
201                )
202            }
203            TimelineFocusKind::Thread { .. } | TimelineFocusKind::PinnedEvents { .. } => false,
204        }
205    }
206
207    /// Whether the focus is on a thread (from a live thread or a thread
208    /// permalink).
209    fn is_thread(&self) -> bool {
210        self.thread_root().is_some()
211    }
212
213    /// If the focus is a thread or event-focused, returns its thread root event
214    /// ID if any.
215    fn thread_root(&self) -> Option<&EventId> {
216        match self {
217            TimelineFocusKind::Event { thread_root, .. } => thread_root.get().map(|v| &**v),
218            TimelineFocusKind::Live { .. } | TimelineFocusKind::PinnedEvents { .. } => None,
219            TimelineFocusKind::Thread { thread_id, .. } => Some(thread_id),
220        }
221    }
222}
223
224#[derive(Clone, Debug)]
225pub(super) struct TimelineController<P: RoomDataProvider = Room> {
226    /// Inner mutable state.
227    state: Arc<RwLock<TimelineState<P>>>,
228
229    /// Focus data.
230    focus: Arc<TimelineFocusKind>,
231
232    /// A [`RoomDataProvider`] implementation, providing data.
233    ///
234    /// The type is a `RoomDataProvider` to allow testing. In the real world,
235    /// this would normally be a [`Room`].
236    pub(crate) room_data_provider: P,
237
238    /// Settings applied to this timeline.
239    pub(super) settings: TimelineSettings,
240}
241
242#[derive(Clone)]
243pub(super) struct TimelineSettings {
244    /// Should the read receipts and read markers be handled and on which event
245    /// types?
246    pub(super) track_read_receipts: TimelineReadReceiptTracking,
247
248    /// Event filter that controls what's rendered as a timeline item (and thus
249    /// what can carry read receipts).
250    pub(super) event_filter: Arc<TimelineEventFilterFn>,
251
252    /// Are unparsable events added as timeline items of their own kind?
253    pub(super) add_failed_to_parse: bool,
254
255    /// Should the timeline items be grouped by day or month?
256    pub(super) date_divider_mode: DateDividerMode,
257}
258
259#[cfg(not(tarpaulin_include))]
260impl fmt::Debug for TimelineSettings {
261    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
262        f.debug_struct("TimelineSettings")
263            .field("track_read_receipts", &self.track_read_receipts)
264            .field("add_failed_to_parse", &self.add_failed_to_parse)
265            .finish_non_exhaustive()
266    }
267}
268
269impl Default for TimelineSettings {
270    fn default() -> Self {
271        Self {
272            track_read_receipts: TimelineReadReceiptTracking::Disabled,
273            event_filter: Arc::new(default_event_filter),
274            add_failed_to_parse: true,
275            date_divider_mode: DateDividerMode::Daily,
276        }
277    }
278}
279
280/// The default event filter for
281/// [`crate::timeline::TimelineBuilder::event_filter`].
282///
283/// It filters out events that are not rendered by the timeline, including but
284/// not limited to: reactions, edits, redactions on existing messages.
285///
286/// If you have a custom filter, it may be best to chain yours with this one if
287/// you do not want to run into situations where a read receipt is not visible
288/// because it's living on an event that doesn't have a matching timeline item.
289pub fn default_event_filter(event: &AnySyncTimelineEvent, rules: &RoomVersionRules) -> bool {
290    match event {
291        AnySyncTimelineEvent::MessageLike(AnySyncMessageLikeEvent::RoomRedaction(ev)) => {
292            if ev.redacts(&rules.redaction).is_some() {
293                // This is a redaction of an existing message, we'll only update
294                // the previous message and not render a new entry.
295                false
296            } else {
297                // This is a redacted entry, that we'll show only if the
298                // redacted entity wasn't a reaction.
299                ev.event_type() != MessageLikeEventType::Reaction
300            }
301        }
302
303        AnySyncTimelineEvent::MessageLike(msg) => {
304            match msg.original_content() {
305                None => {
306                    // This is a redacted entry, that we'll show only if the
307                    // redacted entity wasn't a reaction.
308                    msg.event_type() != MessageLikeEventType::Reaction
309                }
310
311                Some(original_content) => {
312                    match original_content {
313                        AnyMessageLikeEventContent::RoomMessage(content) => {
314                            if content
315                                .relates_to
316                                .as_ref()
317                                .is_some_and(|rel| matches!(rel, Relation::Replacement(_)))
318                            {
319                                // Edits aren't visible by default.
320                                return false;
321                            }
322
323                            match content.msgtype {
324                                MessageType::Audio(_)
325                                | MessageType::Emote(_)
326                                | MessageType::File(_)
327                                | MessageType::Image(_)
328                                | MessageType::Location(_)
329                                | MessageType::Notice(_)
330                                | MessageType::ServerNotice(_)
331                                | MessageType::Text(_)
332                                | MessageType::Video(_)
333                                | MessageType::VerificationRequest(_) => true,
334                                #[cfg(feature = "unstable-msc4274")]
335                                MessageType::Gallery(_) => true,
336                                _ => false,
337                            }
338                        }
339
340                        AnyMessageLikeEventContent::Sticker(_)
341                        | AnyMessageLikeEventContent::UnstablePollStart(
342                            UnstablePollStartEventContent::New(_),
343                        )
344                        | AnyMessageLikeEventContent::CallInvite(_)
345                        | AnyMessageLikeEventContent::RtcNotification(_)
346                        | AnyMessageLikeEventContent::RoomEncrypted(_) => true,
347
348                        // Beacon location-update events are aggregated onto
349                        // their parent `beacon_info` state event's timeline
350                        // item. They are never rendered as standalone items.
351                        AnyMessageLikeEventContent::Beacon(_) => false,
352                        // Ignore decline events, the matching RtcNotification
353                        // event will be updated to reflect the decline.
354                        AnyMessageLikeEventContent::RtcDecline(_) => false,
355
356                        _ => false,
357                    }
358                }
359            }
360        }
361
362        AnySyncTimelineEvent::State(_) => {
363            // All the state events may get displayed by default.
364            true
365        }
366    }
367}
368
369/// Result of calling [`TimelineController::init_focus`].
370pub(super) struct InitFocusResult {
371    /// Did the initialization result in having some events in the timeline?
372    pub has_events: bool,
373    /// If the timeline is a non-live timeline, an extra task that subscribes to
374    /// changes to the focus source.
375    pub focus_task: Option<BackgroundTaskHandle>,
376}
377
378/// Holds the various info about the current call
379#[derive(Clone, Debug, PartialEq)]
380pub struct ActiveCallInfo {
381    /// The list of users in the call
382    pub active_members: HashSet<OwnedUserId>,
383    /// The consensus intent of the call, audio/video
384    pub call_intent: CallIntentConsensus,
385    /// True if the user (with any device) is currently in the call, meaning
386    /// they have joined and haven't left yet.
387    pub is_joined: bool,
388    /// The timestamp of when the call started, in milliseconds since the unix
389    /// epoch. Currently, this is the origin_server_ts of the rtc.notification
390    /// event.
391    pub call_started_ts_millis: Option<MilliSecondsSinceUnixEpoch>,
392}
393
394impl ActiveCallInfo {
395    pub(crate) fn from_info(room_info: RoomInfo, owned_user_id: OwnedUserId) -> Option<Self> {
396        if room_info.has_active_room_call() {
397            Some(ActiveCallInfo {
398                active_members: HashSet::from(room_info.active_room_call_participants()),
399                call_intent: room_info.active_room_call_consensus_intent(),
400                is_joined: room_info.active_room_call_participants().contains(&owned_user_id),
401                call_started_ts_millis: None,
402            })
403        } else {
404            None
405        }
406    }
407
408    pub(crate) fn with_start_time(self, timestamp: Option<MilliSecondsSinceUnixEpoch>) -> Self {
409        Self { call_started_ts_millis: timestamp, ..self }
410    }
411}
412
413impl<P: RoomDataProvider> TimelineController<P> {
414    pub(super) async fn new(
415        room_data_provider: P,
416        focus: &TimelineFocus,
417        event_cache: &EventCache,
418        internal_id_prefix: Option<String>,
419        unable_to_decrypt_hook: Option<Arc<UtdHookManager>>,
420        is_room_encrypted: bool,
421        settings: TimelineSettings,
422    ) -> Result<Self, Error> {
423        let room_id = room_data_provider.room_id();
424
425        let focus = match focus {
426            TimelineFocus::Live { hide_threaded_events } => TimelineFocusKind::Live {
427                hide_threaded_events: *hide_threaded_events,
428                event_cache: event_cache.room(room_id).await?.0,
429            },
430
431            TimelineFocus::Event { target, thread_mode, num_context_events, .. } => {
432                TimelineFocusKind::Event {
433                    event_cache: event_cache
434                        .event_focused(room_id, target, (*thread_mode).into(), *num_context_events)
435                        .await?
436                        .0,
437                    focused_event_id: target.clone(),
438                    // This will be initialised in `Self::init_focus`.
439                    thread_root: OnceLock::new(),
440                    thread_mode: *thread_mode,
441                }
442            }
443
444            TimelineFocus::Thread { thread_id, .. } => TimelineFocusKind::Thread {
445                event_cache: event_cache.thread(room_id, thread_id).await?.0,
446                thread_id: thread_id.clone(),
447            },
448
449            TimelineFocus::PinnedEvents => TimelineFocusKind::PinnedEvents {
450                event_cache: event_cache.pinned_events(room_id).await?.0,
451            },
452        };
453
454        let focus = Arc::new(focus);
455        let state = Arc::new(RwLock::new(TimelineState::new(
456            event_cache.clone(),
457            focus.clone(),
458            room_data_provider.own_user_id().to_owned(),
459            room_data_provider.room_version_rules(),
460            internal_id_prefix,
461            unable_to_decrypt_hook,
462            is_room_encrypted,
463            None,
464        )));
465
466        Ok(Self { state, focus, room_data_provider, settings })
467    }
468
469    /// Listens to encryption state changes for the room in
470    /// [`matrix_sdk_base::RoomInfo`] and applies the new value to the existing
471    /// timeline items. This will then cause a refresh of those timeline items.
472    pub async fn handle_encryption_state_changes(&self) {
473        let mut room_info = self.room_data_provider.room_info();
474
475        // Small function helper to help mark as encrypted.
476        let mark_encrypted = || async {
477            let mut state = self.state.write().await;
478            state.meta.is_room_encrypted = true;
479            state.mark_all_events_as_encrypted();
480        };
481
482        if room_info.get().encryption_state().is_encrypted() {
483            // If the room was already encrypted, it won't toggle to
484            // unencrypted, so we can shut down this task early.
485            mark_encrypted().await;
486            return;
487        }
488
489        while let Some(info) = room_info.next().await {
490            if info.encryption_state().is_encrypted() {
491                mark_encrypted().await;
492                // Once the room is encrypted, it cannot switch back to
493                // unencrypted, so our work here is done.
494                break;
495            }
496        }
497    }
498
499    /// Run a lazy backwards pagination (in live mode).
500    ///
501    /// It adjusts the `count` value of the `Skip` higher-order stream so that
502    /// more items are pushed front in the timeline.
503    ///
504    /// If no more items are available (i.e. if the `count` is zero), this
505    /// method returns `Some(needs)` where `needs` is the number of events that
506    /// must be unlazily backwards paginated.
507    pub(super) async fn live_lazy_paginate_backwards(&self, num_events: u16) -> Option<usize> {
508        let state = self.state.read().await;
509
510        let (count, needs) = state
511            .meta
512            .subscriber_skip_count
513            .compute_next_when_paginating_backwards(num_events.into());
514
515        // This always happens on a live timeline.
516        let is_live_timeline = true;
517        state.meta.subscriber_skip_count.update(count, is_live_timeline);
518
519        needs
520    }
521
522    /// Is this timeline receiving events from sync (aka has a live focus)?
523    pub(super) fn is_live(&self) -> bool {
524        matches!(&*self.focus, TimelineFocusKind::Live { .. })
525    }
526
527    /// Is this timeline focused on a thread?
528    pub(super) fn is_threaded(&self) -> bool {
529        self.focus.is_thread()
530    }
531
532    /// The root of the current thread, for a live thread timeline or a
533    /// permalink to a thread message.
534    pub(super) fn thread_root(&self) -> Option<OwnedEventId> {
535        self.focus.thread_root().map(ToOwned::to_owned)
536    }
537
538    /// Get a copy of the current items in the list.
539    ///
540    /// Cheap because `im::Vector` is cheap to clone.
541    pub(super) async fn items(&self) -> Vector<Arc<TimelineItem>> {
542        self.state.read().await.items.clone_items()
543    }
544
545    #[cfg(test)]
546    pub(super) async fn subscribe_raw(
547        &self,
548    ) -> (Vector<Arc<TimelineItem>>, VectorSubscriberStream<Arc<TimelineItem>>) {
549        self.state.read().await.items.subscribe().into_values_and_stream()
550    }
551
552    pub(super) async fn subscribe(&self) -> (Vector<Arc<TimelineItem>>, TimelineSubscriber) {
553        let state = self.state.read().await;
554
555        TimelineSubscriber::new(&state.items, &state.meta.subscriber_skip_count)
556    }
557
558    pub(super) async fn subscribe_filter_map<U, F>(
559        &self,
560        f: F,
561    ) -> (Vector<U>, FilterMap<VectorSubscriberStream<Arc<TimelineItem>>, F>)
562    where
563        U: Clone,
564        F: Fn(Arc<TimelineItem>) -> Option<U>,
565    {
566        self.state.read().await.items.subscribe().filter_map(f)
567    }
568
569    /// Toggle a reaction locally.
570    ///
571    /// Returns true if the reaction was added, false if it was removed.
572    #[instrument(skip_all)]
573    pub(super) async fn toggle_reaction_local(
574        &self,
575        item_id: &TimelineEventItemId,
576        key: &str,
577        extra_content: Option<serde_json::Map<String, serde_json::Value>>,
578    ) -> Result<bool, Error> {
579        let mut state = self.state.write().await;
580
581        let Some((item_pos, item)) = rfind_event_by_item_id(&state.items, item_id) else {
582            warn!("Timeline item not found, can't add reaction");
583            return Err(Error::FailedToToggleReaction);
584        };
585
586        let user_id = self.room_data_provider.own_user_id();
587        let target = item.identifier();
588
589        // The item says whether we reacted; the registry has the handle and
590        // event id.
591        let has_reaction =
592            item.reactions().get(key).is_some_and(|by_user| by_user.contains_key(user_id));
593        let previous = has_reaction
594            .then(|| state.meta.aggregations.find_reaction(&target, key, user_id).cloned())
595            .flatten();
596
597        if has_reaction && previous.is_none() {
598            warn!("reaction is on the item but unknown to the aggregations");
599            return Ok(false);
600        }
601
602        let Some(previous) = previous else {
603            // Adding the new reaction.
604            match item.handle() {
605                TimelineItemHandle::Local(send_handle) => {
606                    if send_handle
607                        .react(key.to_owned())
608                        .await
609                        .map_err(|err| Error::SendQueueError(err.into()))?
610                        .is_some()
611                    {
612                        trace!("adding a reaction to a local echo");
613                        return Ok(true);
614                    }
615
616                    warn!("couldn't toggle reaction for local echo");
617                    return Ok(false);
618                }
619
620                TimelineItemHandle::Remote(event_id) => {
621                    // Add a reaction through the room data provider. No need to
622                    // reflect the effect locally, since the local echo handling
623                    // will take care of it.
624                    trace!("adding a reaction to a remote echo");
625                    let annotation = Annotation::new(event_id.to_owned(), key.to_owned());
626                    self.room_data_provider
627                        .send(ReactionEventContent::from(annotation).into(), extra_content)
628                        .await?;
629                    return Ok(true);
630                }
631            }
632        };
633
634        trace!("removing a previous reaction");
635
636        if previous.is_local() {
637            // Aborting the local echo is enough: its discard removes the
638            // reaction.
639            if let Some(handle) = &previous.send_handle {
640                if !handle.abort().await.map_err(|err| Error::SendQueueError(err.into()))? {
641                    // Impossible state: the reaction has moved from local to
642                    // echo under our feet, but the timeline was supposed to be
643                    // locked!
644                    warn!("unexpectedly unable to abort sending of local reaction");
645                }
646            } else {
647                warn!("no send handle (this should only happen in testing contexts)");
648            }
649            return Ok(false);
650        }
651
652        let TimelineEventItemId::EventId(event_id) = previous.own_id else {
653            warn!("sent reaction without an event id");
654            return Ok(false);
655        };
656
657        // Assume the redaction will work; we'll re-add the reaction if it
658        // didn't.
659        let Some(annotated_event_id) =
660            item.as_remote().map(|event_item| event_item.event_id.clone())
661        else {
662            warn!("remote reaction to remote event, but the associated item isn't remote");
663            return Ok(false);
664        };
665
666        let mut reactions = item.reactions().clone();
667        let reaction_info = reactions.remove_reaction(user_id, key);
668
669        if reaction_info.is_some() {
670            let new_item = item.with_reactions(reactions);
671            state.items.replace(item_pos, new_item);
672        } else {
673            warn!(
674                "reaction is missing on the item, not removing it locally, \
675                 but sending redaction."
676            );
677        }
678
679        // Release the lock before running the request.
680        drop(state);
681
682        trace!("sending redact for a previous reaction");
683        if let Err(err) = self.room_data_provider.redact(&event_id, None, None).await {
684            if let Some(reaction_info) = reaction_info {
685                debug!("sending redact failed, adding the reaction back to the list");
686
687                let mut state = self.state.write().await;
688                if let Some((item_pos, item)) = rfind_event_by_id(&state.items, &annotated_event_id)
689                {
690                    // Re-add the reaction to the mapping.
691                    let mut reactions = item.reactions().clone();
692                    reactions
693                        .entry(key.to_owned())
694                        .or_default()
695                        .insert(user_id.to_owned(), reaction_info);
696                    let new_item = item.with_reactions(reactions);
697                    state.items.replace(item_pos, new_item);
698                } else {
699                    warn!(
700                        "couldn't find item to re-add reaction anymore; \
701                         maybe it's been redacted?"
702                    );
703                }
704            }
705
706            return Err(err);
707        }
708
709        Ok(false)
710    }
711
712    /// The handle for a pending send on an item, see [`SendTarget`].
713    ///
714    /// `None` if there's nothing of that kind left to act on, which a caller
715    /// racing the remote echo can legitimately run into.
716    pub(super) async fn pending_send_handle(
717        &self,
718        item_id: &TimelineEventItemId,
719        target: SendTarget,
720    ) -> Result<Option<SendHandle>, Error> {
721        let state = self.state.read().await;
722
723        let Some((_, item)) = rfind_event_by_item_id(&state.items, item_id) else {
724            return Err(Error::EventNotInTimeline(item_id.clone()));
725        };
726
727        let own_user_id = self.room_data_provider.own_user_id();
728        let target_id = item.identifier();
729        let aggregations = &state.meta.aggregations;
730
731        let handle = match &target {
732            SendTarget::Event => item.local_echo_send_handle(),
733            SendTarget::Edit => aggregations
734                .pending_send_handle(&target_id, |kind| matches!(kind, AggregationKind::Edit(_))),
735            SendTarget::Redaction => aggregations
736                .pending_send_handle(&target_id, |kind| matches!(kind, AggregationKind::Redaction)),
737            SendTarget::Reaction { key } => {
738                aggregations.pending_send_handle(&target_id, |kind| {
739                    matches!(kind, AggregationKind::Reaction { key: k, sender, .. } if k == key && sender == own_user_id)
740                })
741            }
742        };
743
744        Ok(handle)
745    }
746
747    /// Handle updates on events as [`VectorDiff`]s.
748    pub(super) async fn handle_remote_events_with_diffs(
749        &self,
750        diffs: Vec<VectorDiff<TimelineEvent>>,
751        origin: RemoteEventOrigin,
752    ) {
753        if diffs.is_empty() {
754            return;
755        }
756
757        let mut state = self.state.write().await;
758        state
759            .handle_remote_events_with_diffs(
760                diffs,
761                origin,
762                &self.room_data_provider,
763                &self.settings,
764            )
765            .await
766    }
767
768    /// Only handle aggregations received as [`VectorDiff`]s.
769    pub(super) async fn handle_remote_aggregations(
770        &self,
771        diffs: Vec<VectorDiff<TimelineEvent>>,
772        origin: RemoteEventOrigin,
773    ) {
774        if diffs.is_empty() {
775            return;
776        }
777
778        let mut state = self.state.write().await;
779        state
780            .handle_remote_aggregations(diffs, origin, &self.room_data_provider, &self.settings)
781            .await
782    }
783
784    /// Handle an update of the thread summary of a single event that is a
785    /// thread root.
786    pub(super) async fn handle_thread_summary(
787        &self,
788        thread_root: OwnedEventId,
789        thread_summary: ThreadSummary,
790    ) {
791        let mut state = self.state.write().await;
792        state.handle_thread_summary(thread_root, thread_summary, &self.room_data_provider).await
793    }
794
795    pub(super) async fn clear(&self) {
796        self.state.write().await.clear();
797    }
798
799    /// Replaces the content of the current timeline with initial events.
800    ///
801    /// Also sets up read receipts and the read marker for a live timeline of a
802    /// room.
803    ///
804    /// This is all done with a single lock guard, since we don't want the state
805    /// to be modified between the clear and re-insertion of new events.
806    pub(super) async fn replace_with_initial_remote_events<Events>(
807        &self,
808        events: Events,
809        origin: RemoteEventOrigin,
810    ) where
811        Events: IntoIterator,
812        <Events as IntoIterator>::Item: Into<TimelineEvent>,
813    {
814        let mut state = self.state.write().await;
815
816        let track_read_markers = &self.settings.track_read_receipts;
817        if track_read_markers.is_enabled() {
818            state.populate_initial_user_receipt(&self.room_data_provider, ReceiptType::Read).await;
819            state
820                .populate_initial_user_receipt(&self.room_data_provider, ReceiptType::ReadPrivate)
821                .await;
822        }
823
824        // Replace the events if either the current event list or the new one
825        // aren't empty. Previously we just had to check the new one wasn't
826        // empty because we did a clear operation before so the current one
827        // would always be empty, but now we may want to replace a populated
828        // timeline with an empty one.
829        let mut events = events.into_iter().peekable();
830        if !state.items.is_empty() || events.peek().is_some() {
831            state
832                .replace_with_remote_events(
833                    events,
834                    origin,
835                    &self.room_data_provider,
836                    &self.settings,
837                )
838                .await;
839        }
840
841        if track_read_markers.is_enabled() {
842            if let Some(fully_read_event_id) =
843                self.room_data_provider.load_fully_read_marker().await
844            {
845                state.handle_fully_read_marker(fully_read_event_id);
846            } else if let Some(latest_receipt_event_id) = state
847                .latest_user_read_receipt_timeline_event_id(self.room_data_provider.own_user_id())
848            {
849                // Fall back to read receipt if no fully read marker exists.
850                debug!("no `m.fully_read` marker found, falling back to read receipt");
851                state.handle_fully_read_marker(latest_receipt_event_id);
852            }
853        }
854    }
855
856    pub(super) async fn handle_fully_read_marker(&self, fully_read_event_id: OwnedEventId) {
857        self.state.write().await.handle_fully_read_marker(fully_read_event_id);
858    }
859
860    pub(super) async fn handle_active_call_update(
861        &self,
862        maybe_active_call: Option<ActiveCallInfo>,
863    ) {
864        let mut state = self.state.write().await;
865        let mut txn = state.transaction();
866
867        // Store the current active call info in metadata for new
868        // RtcNotification items
869        txn.meta.active_call = maybe_active_call.clone();
870
871        if let Some(existing_event_id) = &txn.meta.active_rtc_notification_event_id {
872            // Clean up the notification event
873            let last_notification = rfind_event_by_id(&txn.items, existing_event_id);
874            if let Some((last_idx, last_notification)) = last_notification {
875                let updated_content = match last_notification.content() {
876                    TimelineItemContent::RtcNotification {
877                        call_intent,
878                        declined_by,
879                        active_call_info: _active_call_info,
880                    } => Some(TimelineItemContent::RtcNotification {
881                        call_intent: call_intent.to_owned(),
882                        declined_by: declined_by.clone(),
883                        active_call_info: maybe_active_call
884                            .clone()
885                            .map(|info| info.with_start_time(last_notification.timestamp.into())),
886                    }),
887                    _ => None,
888                };
889                if let Some(new_content) = updated_content {
890                    let new_event_item = last_notification.inner.with_content(new_content);
891                    let new_timeline_item =
892                        TimelineItem::new(new_event_item, last_notification.internal_id.clone());
893                    txn.items.replace(last_idx, new_timeline_item);
894                }
895
896                if maybe_active_call.is_none() {
897                    // There is no active rtc_notification anymore
898                    txn.meta.active_rtc_notification_event_id = None;
899                }
900            }
901        }
902
903        txn.commit();
904    }
905
906    pub(super) async fn handle_read_receipt_event(&self, event: ReceiptEventContent) {
907        // Don't even take the lock if there are no events to process.
908        if event.is_empty() {
909            return;
910        }
911
912        let mut state = self.state.write().await;
913        state.handle_read_receipt(event, &self.room_data_provider).await;
914    }
915
916    /// Creates the local echo for an event we're sending.
917    #[instrument(skip_all)]
918    pub(super) async fn handle_local_event(
919        &self,
920        txn_id: OwnedTransactionId,
921        content: AnyMessageLikeEventContent,
922        send_handle: Option<SendHandle>,
923    ) {
924        let sender = self.room_data_provider.own_user_id().to_owned();
925        let profile = self.room_data_provider.profile_from_user_id(&sender).await;
926
927        let date_divider_mode = self.settings.date_divider_mode.clone();
928
929        let mut state = self.state.write().await;
930        state
931            .handle_local_event(sender, profile, date_divider_mode, txn_id, send_handle, content)
932            .await;
933    }
934
935    /// Update the send state of a local event represented by a transaction ID.
936    ///
937    /// If the corresponding local timeline item is missing, a warning is
938    /// raised.
939    #[instrument(skip(self))]
940    pub(super) async fn update_event_send_state(
941        &self,
942        txn_id: &TransactionId,
943        send_state: EventSendState,
944    ) {
945        let mut state = self.state.write().await;
946        let mut txn = state.transaction();
947
948        let new_event_id: Option<&EventId> =
949            as_variant!(&send_state, EventSendState::Sent { event_id } => event_id);
950
951        // The local echoes are always at the end of the timeline, we must first
952        // make sure the remote echo hasn't showed up yet.
953        if rfind_event_item(&txn.items, |it| {
954            new_event_id.is_some() && it.event_id() == new_event_id && it.as_remote().is_some()
955        })
956        .is_some()
957        {
958            // Remote echo already received. This is very unlikely.
959            trace!("Remote echo received before send-event response");
960
961            let local_echo = rfind_event_item(&txn.items, |it| it.transaction_id() == Some(txn_id));
962
963            // If there's both the remote echo and a local echo, that means the
964            // remote echo was received before the response _and_ contained no
965            // transaction ID (and thus duplicated the local echo).
966            if let Some((idx, _)) = local_echo {
967                warn!("Message echo got duplicated, removing the local one");
968                txn.items.remove(idx);
969
970                // Adjust the date dividers, if needs be.
971                let mut adjuster =
972                    DateDividerAdjuster::new(self.settings.date_divider_mode.clone());
973                adjuster.run(&mut txn.items, &mut txn.meta);
974            }
975
976            txn.commit();
977            return;
978        }
979
980        // Look for the local event by the transaction ID or event ID.
981        let result = rfind_event_item(&txn.items, |it| {
982            it.transaction_id() == Some(txn_id)
983                || new_event_id.is_some()
984                    && it.event_id() == new_event_id
985                    && it.as_local().is_some()
986        });
987
988        let Some((idx, item)) = result else {
989            // Not a standalone item: maybe one of our aggregations.
990            if txn.meta.aggregations.update_send_state(
991                txn_id.to_owned(),
992                send_state,
993                &mut txn.items,
994                &txn.meta.room_version_rules,
995            ) {
996                trace!("Updated the send state of an aggregation");
997                txn.commit();
998                return;
999            }
1000
1001            warn!("Timeline item not found, can't update send state");
1002            return;
1003        };
1004
1005        let Some(local_item) = item.as_local() else {
1006            warn!("We looked for a local item, but it transitioned to remote.");
1007            return;
1008        };
1009
1010        // The event was already marked as sent, that's a broken state, let's
1011        // emit an error but also override to the given sent state.
1012        if let EventSendState::Sent { event_id: existing_event_id } = &local_item.send_state {
1013            error!(?existing_event_id, ?new_event_id, "Local echo already marked as sent");
1014        }
1015
1016        // If the event has just been marked as sent, update the aggregations
1017        // mapping to take that into account.
1018        if let Some(new_event_id) = new_event_id {
1019            txn.meta.aggregations.mark_target_as_sent(txn_id.to_owned(), new_event_id.to_owned());
1020        }
1021
1022        let new_item = item.with_inner_kind(local_item.with_send_state(send_state));
1023        txn.items.replace(idx, new_item);
1024
1025        txn.commit();
1026    }
1027
1028    pub(super) async fn discard_local_echo(&self, txn_id: &TransactionId) -> bool {
1029        let mut state = self.state.write().await;
1030
1031        if let Some((idx, _)) =
1032            rfind_event_item(&state.items, |it| it.transaction_id() == Some(txn_id))
1033        {
1034            let mut txn = state.transaction();
1035
1036            txn.items.remove(idx);
1037
1038            // A read marker or a date divider may have been inserted before the
1039            // local echo. Ensure both are up to date.
1040            let mut adjuster = DateDividerAdjuster::new(self.settings.date_divider_mode.clone());
1041            adjuster.run(&mut txn.items, &mut txn.meta);
1042
1043            txn.meta.update_read_marker(&mut txn.items);
1044
1045            txn.commit();
1046
1047            debug!("discarded local echo");
1048            return true;
1049        }
1050
1051        // Avoid multiple mutable and immutable borrows of the lock guard by
1052        // explicitly dereferencing it once.
1053        let mut txn = state.transaction();
1054
1055        // Look if this was a local aggregation.
1056        let found_aggregation = match txn.meta.aggregations.try_remove_aggregation(
1057            &TimelineEventItemId::TransactionId(txn_id.to_owned()),
1058            &mut txn.items,
1059        ) {
1060            Ok(val) => val,
1061            Err(err) => {
1062                warn!("error when discarding local echo for an aggregation: {err}");
1063                // The aggregation has been found, it's just that we couldn't
1064                // discard it.
1065                true
1066            }
1067        };
1068
1069        if found_aggregation {
1070            txn.commit();
1071        }
1072
1073        found_aggregation
1074    }
1075
1076    pub(super) async fn replace_local_echo(
1077        &self,
1078        txn_id: &TransactionId,
1079        content: AnyMessageLikeEventContent,
1080    ) -> bool {
1081        let AnyMessageLikeEventContent::RoomMessage(content) = content else {
1082            // Ideally, we'd support replacing local echoes for a reaction,
1083            // etc., but handling RoomMessage should be sufficient in most
1084            // cases. Worst case, the local echo will be sent Soonâ„¢ and we'll
1085            // get another chance at editing the event then.
1086            warn!("Replacing a local echo for a non-RoomMessage-like event NYI");
1087            return false;
1088        };
1089
1090        let mut state = self.state.write().await;
1091        let mut txn = state.transaction();
1092
1093        let Some((idx, prev_item)) =
1094            rfind_event_item(&txn.items, |it| it.transaction_id() == Some(txn_id))
1095        else {
1096            // Not a standalone item: maybe one of our pending edits.
1097            if let Some(Relation::Replacement(replacement)) = content.relates_to
1098                && txn.meta.aggregations.replace_local_edit(
1099                    txn_id,
1100                    replacement,
1101                    &mut txn.items,
1102                    &txn.meta.room_version_rules,
1103                )
1104            {
1105                debug!("Replaced local echo of an edit");
1106                txn.commit();
1107                return true;
1108            }
1109
1110            debug!("Can't find local echo to replace");
1111            return false;
1112        };
1113
1114        // Reuse the previous local echo's state, but reset the send state to
1115        // not sent (per API contract).
1116        let ti_kind = {
1117            let Some(prev_local_item) = prev_item.as_local() else {
1118                warn!("We looked for a local item, but it transitioned as remote??");
1119                return false;
1120            };
1121            // If the local echo had an upload progress, retain it.
1122            let progress = as_variant!(&prev_local_item.send_state,
1123                EventSendState::NotSentYet { progress } => progress.clone())
1124            .flatten();
1125            prev_local_item.with_send_state(EventSendState::NotSentYet { progress })
1126        };
1127
1128        // Replace the local-related state (kind) and the content state.
1129        let new_item = TimelineItem::new(
1130            prev_item.with_kind(ti_kind).with_content(TimelineItemContent::message(
1131                content.msgtype,
1132                content.mentions,
1133                prev_item.content().thread_root(),
1134                prev_item.content().in_reply_to(),
1135                prev_item.content().thread_summary(),
1136            )),
1137            prev_item.internal_id.to_owned(),
1138        );
1139
1140        txn.items.replace(idx, new_item);
1141
1142        // This doesn't change the original sending time, so there's no need to
1143        // adjust date dividers.
1144
1145        txn.commit();
1146
1147        debug!("Replaced local echo");
1148        true
1149    }
1150
1151    pub(super) async fn compute_redecryption_candidates(
1152        &self,
1153    ) -> (BTreeSet<String>, BTreeSet<String>) {
1154        let state = self.state.read().await;
1155        compute_redecryption_candidates(&state.items)
1156    }
1157
1158    pub(super) async fn set_sender_profiles_pending(&self) {
1159        self.set_non_ready_sender_profiles(TimelineDetails::Pending).await;
1160    }
1161
1162    pub(super) async fn set_sender_profiles_error(&self, error: Arc<matrix_sdk::Error>) {
1163        self.set_non_ready_sender_profiles(TimelineDetails::Error(error)).await;
1164    }
1165
1166    async fn set_non_ready_sender_profiles(&self, profile_state: TimelineDetails<Profile>) {
1167        self.state.write().await.items.for_each(|mut entry| {
1168            let Some(event_item) = entry.as_event() else { return };
1169            if !matches!(event_item.sender_profile(), TimelineDetails::Ready(_)) {
1170                let new_item = entry.with_kind(TimelineItemKind::Event(
1171                    event_item.with_sender_profile(profile_state.clone()),
1172                ));
1173                ObservableItemsEntry::replace(&mut entry, new_item);
1174            }
1175        });
1176    }
1177
1178    pub(super) async fn update_missing_sender_profiles(&self) {
1179        trace!("Updating missing sender profiles");
1180
1181        let mut state = self.state.write().await;
1182        let mut entries = state.items.entries();
1183        while let Some(mut entry) = entries.next() {
1184            let Some(event_item) = entry.as_event() else { continue };
1185            let event_id = event_item.event_id().map(debug);
1186            let transaction_id = event_item.transaction_id().map(debug);
1187
1188            if event_item.sender_profile().is_ready() {
1189                trace!(event_id, transaction_id, "Profile already set");
1190                continue;
1191            }
1192
1193            match self.room_data_provider.profile_from_user_id(event_item.sender()).await {
1194                Some(profile) => {
1195                    trace!(event_id, transaction_id, "Adding profile");
1196                    let updated_item =
1197                        event_item.with_sender_profile(TimelineDetails::Ready(profile));
1198                    let new_item = entry.with_kind(updated_item);
1199                    ObservableItemsEntry::replace(&mut entry, new_item);
1200                }
1201                None => {
1202                    if !event_item.sender_profile().is_unavailable() {
1203                        trace!(event_id, transaction_id, "Marking profile unavailable");
1204                        let updated_item =
1205                            event_item.with_sender_profile(TimelineDetails::Unavailable);
1206                        let new_item = entry.with_kind(updated_item);
1207                        ObservableItemsEntry::replace(&mut entry, new_item);
1208                    } else {
1209                        debug!(event_id, transaction_id, "Profile already marked unavailable");
1210                    }
1211                }
1212            }
1213        }
1214
1215        trace!("Done updating missing sender profiles");
1216    }
1217
1218    /// Update the profiles of the given senders, even if they are ready.
1219    pub(super) async fn force_update_sender_profiles(&self, sender_ids: &BTreeSet<&UserId>) {
1220        trace!("Forcing update of sender profiles: {sender_ids:?}");
1221
1222        let mut state = self.state.write().await;
1223        let mut entries = state.items.entries();
1224        while let Some(mut entry) = entries.next() {
1225            let Some(event_item) = entry.as_event() else { continue };
1226            if !sender_ids.contains(event_item.sender()) {
1227                continue;
1228            }
1229
1230            let event_id = event_item.event_id().map(debug);
1231            let transaction_id = event_item.transaction_id().map(debug);
1232
1233            match self.room_data_provider.profile_from_user_id(event_item.sender()).await {
1234                Some(profile) => {
1235                    if matches!(event_item.sender_profile(), TimelineDetails::Ready(old_profile) if *old_profile == profile)
1236                    {
1237                        debug!(event_id, transaction_id, "Profile already up-to-date");
1238                    } else {
1239                        trace!(event_id, transaction_id, "Updating profile");
1240                        let updated_item =
1241                            event_item.with_sender_profile(TimelineDetails::Ready(profile));
1242                        let new_item = entry.with_kind(updated_item);
1243                        ObservableItemsEntry::replace(&mut entry, new_item);
1244                    }
1245                }
1246                None => {
1247                    if !event_item.sender_profile().is_unavailable() {
1248                        trace!(event_id, transaction_id, "Marking profile unavailable");
1249                        let updated_item =
1250                            event_item.with_sender_profile(TimelineDetails::Unavailable);
1251                        let new_item = entry.with_kind(updated_item);
1252                        ObservableItemsEntry::replace(&mut entry, new_item);
1253                    } else {
1254                        debug!(event_id, transaction_id, "Profile already marked unavailable");
1255                    }
1256                }
1257            }
1258        }
1259
1260        trace!("Done forcing update of sender profiles");
1261    }
1262
1263    #[cfg(test)]
1264    pub(super) async fn handle_read_receipts(&self, receipt_event_content: ReceiptEventContent) {
1265        let own_user_id = self.room_data_provider.own_user_id();
1266        self.state.write().await.handle_read_receipts(receipt_event_content, own_user_id);
1267    }
1268
1269    /// Get the latest read receipt for the given user.
1270    ///
1271    /// Useful to get the latest read receipt, whether it's private or public.
1272    pub(super) async fn latest_user_read_receipt(
1273        &self,
1274        user_id: &UserId,
1275    ) -> Option<(OwnedEventId, Receipt)> {
1276        let receipt_thread = self.focus.receipt_thread();
1277
1278        self.state
1279            .read()
1280            .await
1281            .latest_user_read_receipt(
1282                user_id,
1283                receipt_thread,
1284                &self.room_data_provider,
1285                read_receipts::ImplicitReadReceipts::Include,
1286            )
1287            .await
1288    }
1289
1290    /// Get the ID of the timeline event with the latest read receipt for the
1291    /// given user.
1292    pub(super) async fn latest_user_read_receipt_timeline_event_id(
1293        &self,
1294        user_id: &UserId,
1295    ) -> Option<OwnedEventId> {
1296        self.state.read().await.latest_user_read_receipt_timeline_event_id(user_id)
1297    }
1298
1299    /// Subscribe to changes in the read receipts of our own user.
1300    pub async fn subscribe_own_user_read_receipts_changed(
1301        &self,
1302    ) -> impl Stream<Item = ()> + use<P> {
1303        self.state.read().await.meta.read_receipts.subscribe_own_user_read_receipts_changed()
1304    }
1305
1306    /// Handle a room send update that's a new local echo.
1307    pub(crate) async fn handle_local_echo(&self, echo: LocalEcho) {
1308        match echo.content {
1309            LocalEchoContent::Event { serialized_event, send_handle, send_error } => {
1310                let content = match serialized_event.deserialize() {
1311                    Ok(d) => d,
1312                    Err(err) => {
1313                        warn!("error deserializing local echo: {err}");
1314                        return;
1315                    }
1316                };
1317
1318                self.handle_local_event(echo.transaction_id.clone(), content, Some(send_handle))
1319                    .await;
1320
1321                if let Some(send_error) = send_error {
1322                    self.update_event_send_state(
1323                        &echo.transaction_id,
1324                        EventSendState::SendingFailed {
1325                            error: Arc::new(matrix_sdk::Error::SendQueueWedgeError(Box::new(
1326                                send_error,
1327                            ))),
1328                            is_recoverable: false,
1329                        },
1330                    )
1331                    .await;
1332                }
1333            }
1334
1335            LocalEchoContent::React { key, send_handle, applies_to } => {
1336                self.handle_local_reaction(key, send_handle, applies_to).await;
1337            }
1338
1339            LocalEchoContent::Redaction { redacts, send_handle, send_error, .. } => {
1340                self.handle_local_redaction(
1341                    echo.transaction_id.clone(),
1342                    redacts,
1343                    Some(send_handle),
1344                )
1345                .await;
1346
1347                if let Some(send_error) = send_error {
1348                    self.update_event_send_state(
1349                        &echo.transaction_id,
1350                        EventSendState::SendingFailed {
1351                            error: Arc::new(matrix_sdk::Error::SendQueueWedgeError(Box::new(
1352                                send_error,
1353                            ))),
1354                            is_recoverable: false,
1355                        },
1356                    )
1357                    .await;
1358                }
1359            }
1360        }
1361    }
1362
1363    /// Adds a reaction (local echo) to a local echo.
1364    #[instrument(skip(self, send_handle))]
1365    async fn handle_local_reaction(
1366        &self,
1367        reaction_key: String,
1368        send_handle: SendHandle,
1369        applies_to: OwnedTransactionId,
1370    ) {
1371        let mut state = self.state.write().await;
1372        let mut tr = state.transaction();
1373
1374        let target = TimelineEventItemId::TransactionId(applies_to);
1375
1376        let reaction_txn_id = send_handle.transaction_id().to_owned();
1377        let aggregation = Aggregation::new_local(
1378            TimelineEventItemId::TransactionId(reaction_txn_id),
1379            AggregationKind::Reaction {
1380                key: reaction_key.clone(),
1381                sender: self.room_data_provider.own_user_id().to_owned(),
1382                timestamp: MilliSecondsSinceUnixEpoch::now(),
1383            },
1384            Some(send_handle),
1385        );
1386
1387        tr.meta.aggregations.add(target.clone(), aggregation.clone());
1388        find_item_and_apply_aggregation(
1389            &tr.meta.aggregations,
1390            &mut tr.items,
1391            &target,
1392            aggregation,
1393            &tr.meta.room_version_rules,
1394        );
1395
1396        tr.commit();
1397    }
1398
1399    /// Applies a local echo of a redaction.
1400    pub(super) async fn handle_local_redaction(
1401        &self,
1402        txn_id: OwnedTransactionId,
1403        redacts: OwnedEventId,
1404        send_handle: Option<SendHandle>,
1405    ) {
1406        let mut state = self.state.write().await;
1407        let mut tr = state.transaction();
1408
1409        let target = TimelineEventItemId::EventId(redacts);
1410
1411        let aggregation = Aggregation::new_local(
1412            TimelineEventItemId::TransactionId(txn_id),
1413            AggregationKind::Redaction,
1414            send_handle,
1415        );
1416
1417        tr.meta.aggregations.add(target.clone(), aggregation.clone());
1418        find_item_and_apply_aggregation(
1419            &tr.meta.aggregations,
1420            &mut tr.items,
1421            &target,
1422            aggregation,
1423            &tr.meta.room_version_rules,
1424        );
1425
1426        tr.commit();
1427    }
1428
1429    /// Handle a single room send queue update.
1430    pub(crate) async fn handle_room_send_queue_update(&self, update: RoomSendQueueUpdate) {
1431        match update {
1432            RoomSendQueueUpdate::NewLocalEvent(echo) => {
1433                self.handle_local_echo(echo).await;
1434            }
1435
1436            RoomSendQueueUpdate::CancelledLocalEvent { transaction_id } => {
1437                if !self.discard_local_echo(&transaction_id).await {
1438                    warn!("couldn't find the local echo to discard");
1439                }
1440            }
1441
1442            RoomSendQueueUpdate::ReplacedLocalEvent { transaction_id, new_content } => {
1443                let content = match new_content.deserialize() {
1444                    Ok(d) => d,
1445                    Err(err) => {
1446                        warn!("error deserializing local echo (upon edit): {err}");
1447                        return;
1448                    }
1449                };
1450
1451                if !self.replace_local_echo(&transaction_id, content).await {
1452                    warn!("couldn't find the local echo to replace");
1453                }
1454            }
1455
1456            RoomSendQueueUpdate::SendError { transaction_id, error, is_recoverable } => {
1457                self.update_event_send_state(
1458                    &transaction_id,
1459                    EventSendState::SendingFailed { error, is_recoverable },
1460                )
1461                .await;
1462            }
1463
1464            RoomSendQueueUpdate::RetryEvent { transaction_id } => {
1465                self.update_event_send_state(
1466                    &transaction_id,
1467                    EventSendState::NotSentYet { progress: None },
1468                )
1469                .await;
1470            }
1471
1472            RoomSendQueueUpdate::SentEvent { transaction_id, event_id } => {
1473                self.update_event_send_state(&transaction_id, EventSendState::Sent { event_id })
1474                    .await;
1475            }
1476
1477            RoomSendQueueUpdate::MediaUpload { related_to, index, progress, .. } => {
1478                self.update_event_send_state(
1479                    &related_to,
1480                    EventSendState::NotSentYet {
1481                        progress: Some(MediaUploadProgress { index, progress }),
1482                    },
1483                )
1484                .await;
1485            }
1486        }
1487    }
1488
1489    /// Insert a timeline start item at the beginning of the room, if it's
1490    /// missing.
1491    pub async fn insert_timeline_start_if_missing(&self) {
1492        let mut state = self.state.write().await;
1493        let mut txn = state.transaction();
1494        txn.items.push_timeline_start_if_missing(
1495            txn.meta.new_timeline_item(VirtualTimelineItem::TimelineStart),
1496        );
1497        txn.commit();
1498    }
1499
1500    /// Create a [`EmbeddedEvent`] from an arbitrary event, be it in the
1501    /// timeline or not.
1502    ///
1503    /// Can be `None` if the event cannot be represented as a standalone item,
1504    /// because it's an aggregation.
1505    pub(super) async fn make_replied_to(
1506        &self,
1507        event: TimelineEvent,
1508    ) -> Result<Option<EmbeddedEvent>, Error> {
1509        let state = self.state.read().await;
1510        EmbeddedEvent::try_from_timeline_event(event, &self.room_data_provider, &state.meta).await
1511    }
1512}
1513
1514impl TimelineController {
1515    pub(super) fn room(&self) -> &Room {
1516        &self.room_data_provider
1517    }
1518
1519    /// Initializes the configured timeline focus with appropriate data.
1520    ///
1521    /// Should be called only once after creation of the [`TimelineController`],
1522    /// with all its fields set.
1523    pub(super) async fn init_focus(&self) -> Result<InitFocusResult, Error> {
1524        match self.focus.deref() {
1525            TimelineFocusKind::Live { event_cache, .. } => {
1526                // Retrieve the cached events, and add them to the timeline.
1527                let events = event_cache.events().await?;
1528
1529                let has_events = !events.is_empty();
1530
1531                self.replace_with_initial_remote_events(events, RemoteEventOrigin::Cache).await;
1532
1533                match event_cache.pagination().status().get() {
1534                    PaginationStatus::Idle { hit_timeline_start } => {
1535                        if hit_timeline_start {
1536                            // Eagerly insert the timeline start item, since
1537                            // pagination claims we've already hit the timeline
1538                            // start.
1539                            self.insert_timeline_start_if_missing().await;
1540                        }
1541                    }
1542                    PaginationStatus::Paginating => {}
1543                }
1544
1545                Ok(InitFocusResult { has_events, focus_task: None })
1546            }
1547
1548            TimelineFocusKind::Event {
1549                focused_event_id: event_id,
1550                thread_mode,
1551                thread_root: focus_thread_root,
1552                event_cache,
1553                ..
1554            } => {
1555                let (events, receiver) = event_cache.subscribe().await?;
1556
1557                let has_events = !events.is_empty();
1558
1559                // Ask the cache for the thread root, if it managed to extract
1560                // one or decided that the target event was the thread root.
1561                if let Some(thread_root) = event_cache.thread_root().await? {
1562                    focus_thread_root.get_or_init(|| thread_root);
1563                }
1564
1565                self.replace_with_initial_remote_events(events, RemoteEventOrigin::Pagination)
1566                    .await;
1567
1568                let task = self
1569                    .room_data_provider
1570                    .client()
1571                    .task_monitor()
1572                    .spawn_infinite_task(
1573                        "timeline::event_focused_cache_updates",
1574                        event_focused_task(
1575                            event_id.clone(),
1576                            (*thread_mode).into(),
1577                            event_cache.clone(),
1578                            self.clone(),
1579                            receiver,
1580                        ),
1581                    )
1582                    .abort_on_drop();
1583
1584                Ok(InitFocusResult { has_events, focus_task: Some(task) })
1585            }
1586
1587            TimelineFocusKind::Thread { event_cache, .. } => {
1588                let (has_events, subscriber) = self.init_with_thread_root(event_cache).await?;
1589
1590                let room = &self.room_data_provider;
1591                let span = info_span!(
1592                    parent: Span::none(),
1593                    "thread_live_update_handler",
1594                    room_id = ?room.room_id(),
1595                );
1596                span.follows_from(Span::current());
1597
1598                let task = room
1599                    .client()
1600                    .task_monitor()
1601                    .spawn_infinite_task(
1602                        "timeline::thread_event_cache_updates",
1603                        thread_updates_task(subscriber, event_cache.clone(), self.clone())
1604                            .instrument(span),
1605                    )
1606                    .abort_on_drop();
1607
1608                Ok(InitFocusResult { has_events, focus_task: Some(task) })
1609            }
1610
1611            TimelineFocusKind::PinnedEvents { event_cache } => {
1612                let (initial_events, pinned_events_recv) = event_cache.subscribe().await?;
1613
1614                let has_events = !initial_events.is_empty();
1615
1616                self.replace_with_initial_remote_events(
1617                    initial_events,
1618                    RemoteEventOrigin::Pagination,
1619                )
1620                .await;
1621
1622                let task = self
1623                    .room_data_provider
1624                    .client()
1625                    .task_monitor()
1626                    .spawn_infinite_task(
1627                        "timeline::pinned_events_cache_updates",
1628                        pinned_events_task(event_cache.clone(), self.clone(), pinned_events_recv),
1629                    )
1630                    .abort_on_drop();
1631
1632                Ok(InitFocusResult { has_events, focus_task: Some(task) })
1633            }
1634        }
1635    }
1636
1637    /// (Re-)initialise a timeline using [`TimelineFocus::Thread`] with cached
1638    /// threaded events and secondary relations.
1639    ///
1640    /// Returns whether there were any events added to the timeline, and a
1641    /// receiver to return updates after the initial events have been inserted
1642    /// in the timeline.
1643    pub(super) async fn init_with_thread_root(
1644        &self,
1645        event_cache: &ThreadEventCache,
1646    ) -> Result<(bool, EventCacheSubscriber<ThreadEventCacheUpdate>), Error> {
1647        let (events, subscriber) = event_cache.subscribe().await?;
1648        let has_events = !events.is_empty();
1649
1650        // For each event, we also need to find the related events, as they
1651        // don't include the thread relationship, they won't be included in the
1652        // initial list of events.
1653        //
1654        // The lookups are independent store queries, so run them together
1655        // rather than awaiting them one after the other. `try_join_all` keeps
1656        // the input order, so the related events are collected in the same
1657        // order as before.
1658        let lookups = events
1659            .iter()
1660            .filter_map(|event| event.event_id())
1661            .map(|event_id| event_cache.find_event_with_relations(event_id, None));
1662
1663        let mut related_events = Vector::new();
1664        for (_original, related) in try_join_all(lookups).await?.into_iter().flatten() {
1665            related_events.extend(related);
1666        }
1667
1668        self.replace_with_initial_remote_events(events, RemoteEventOrigin::Cache).await;
1669
1670        // Now that we've inserted the thread events, add the aggregations too.
1671        if !related_events.is_empty() {
1672            self.handle_remote_aggregations(
1673                vec![VectorDiff::Append { values: related_events }],
1674                RemoteEventOrigin::Cache,
1675            )
1676            .await;
1677        }
1678
1679        Ok((has_events, subscriber))
1680    }
1681
1682    /// Given an event identifier, will fetch the details for the event it's
1683    /// replying to, if applicable.
1684    #[instrument(skip(self))]
1685    pub(super) async fn fetch_in_reply_to_details(&self, event_id: &EventId) -> Result<(), Error> {
1686        let state_guard = self.state.write().await;
1687        let (index, item) = rfind_event_by_id(&state_guard.items, event_id)
1688            .ok_or(Error::EventNotInTimeline(TimelineEventItemId::EventId(event_id.to_owned())))?;
1689        let remote_item = item
1690            .as_remote()
1691            .ok_or(Error::EventNotInTimeline(TimelineEventItemId::EventId(event_id.to_owned())))?
1692            .clone();
1693
1694        let TimelineItemContent::MsgLike(msglike) = item.content().clone() else {
1695            debug!("Event is not a message");
1696            return Ok(());
1697        };
1698        let Some(in_reply_to) = msglike.in_reply_to.clone() else {
1699            debug!("Event is not a reply");
1700            return Ok(());
1701        };
1702        if let TimelineDetails::Pending = &in_reply_to.event {
1703            debug!("Replied-to event is already being fetched");
1704            return Ok(());
1705        }
1706        if let TimelineDetails::Ready(_) = &in_reply_to.event {
1707            debug!("Replied-to event has already been fetched");
1708            return Ok(());
1709        }
1710
1711        let internal_id = item.internal_id.to_owned();
1712        let item = item.clone();
1713        let event = fetch_replied_to_event(
1714            state_guard,
1715            &self.state,
1716            index,
1717            &item,
1718            internal_id,
1719            &msglike,
1720            &in_reply_to.event_id,
1721            self.room(),
1722        )
1723        .await?;
1724
1725        // We need to be sure to have the latest position of the event as it
1726        // might have changed while waiting for the request.
1727        let mut state = self.state.write().await;
1728        let (index, item) = rfind_event_by_id(&state.items, &remote_item.event_id)
1729            .ok_or(Error::EventNotInTimeline(TimelineEventItemId::EventId(event_id.to_owned())))?;
1730
1731        // Check the state of the event again, it might have been redacted while
1732        // the request was in-flight.
1733        let TimelineItemContent::MsgLike(MsgLikeContent {
1734            kind: MsgLikeKind::Message(message),
1735            thread_root,
1736            in_reply_to,
1737            thread_summary,
1738        }) = item.content().clone()
1739        else {
1740            info!("Event is no longer a message (redacted?)");
1741            return Ok(());
1742        };
1743        let Some(in_reply_to) = in_reply_to else {
1744            warn!("Event no longer has a reply (bug?)");
1745            return Ok(());
1746        };
1747
1748        // Now that we've received the content of the replied-to event, replace
1749        // the replied-to content in the item with it.
1750        trace!("Updating in-reply-to details");
1751        let internal_id = item.internal_id.to_owned();
1752        let mut item = item.clone();
1753        item.set_content(TimelineItemContent::MsgLike(MsgLikeContent {
1754            kind: MsgLikeKind::Message(message),
1755            thread_root,
1756            in_reply_to: Some(InReplyToDetails { event_id: in_reply_to.event_id, event }),
1757            thread_summary,
1758        }));
1759        state.items.replace(index, TimelineItem::new(item, internal_id));
1760
1761        Ok(())
1762    }
1763
1764    /// Returns the thread that should be used for a read receipt based on the
1765    /// current focus of the timeline and the receipt type.
1766    ///
1767    /// A `SendReceiptType::FullyRead` will always use
1768    /// `ReceiptThread::Unthreaded`
1769    pub(super) fn infer_thread_for_read_receipt(
1770        &self,
1771        receipt_type: &SendReceiptType,
1772    ) -> ReceiptThread {
1773        if matches!(receipt_type, SendReceiptType::FullyRead) {
1774            ReceiptThread::Unthreaded
1775        } else {
1776            self.focus.receipt_thread()
1777        }
1778    }
1779
1780    /// Decide whether a read receipt should be sent, and which event it should
1781    /// target.
1782    ///
1783    /// The returned event may differ from `event_id`: a read receipt should not
1784    /// point at one of the user's own events (see the Matrix spec's
1785    /// [Receipts module]), so if `event_id` is one of theirs, the latest unread
1786    /// event before it that is targeted instead.
1787    ///
1788    /// - When there's no such earlier event and `is_marking_room_as_read` is
1789    ///   `false`, [`SendReceiptDecision::DoNotSend`] is returned.
1790    /// - When `is_marking_room_as_read` is `true`, the receipt falls back to
1791    ///   `event_id` itself if it is explicitly unread, so that the homeserver
1792    ///   still recomputes the push/badge count.
1793    ///
1794    /// [Receipts module]: https://spec.matrix.org/latest/client-server-api/#receipts
1795    pub(super) async fn should_send_receipt(
1796        &self,
1797        receipt_type: &SendReceiptType,
1798        receipt_thread: &ReceiptThread,
1799        event_id: &EventId,
1800        is_marking_room_as_read: bool,
1801    ) -> SendReceiptDecision {
1802        let own_user_id = self.room().own_user_id();
1803        let state = self.state.read().await;
1804        let room = self.room();
1805        let all_remote_events = state.items.all_remote_events();
1806
1807        // Resolve the event the receipt should target, redirecting away from
1808        // the user's own events for read receipts.
1809        let target_event_id = match receipt_type {
1810            SendReceiptType::Read | SendReceiptType::ReadPrivate => {
1811                let is_own_event = all_remote_events
1812                    .get_by_event_id(event_id)
1813                    .and_then(|event_meta| event_meta.sender.as_deref())
1814                    == Some(own_user_id);
1815
1816                if is_own_event {
1817                    let filter_out_thread_events = match self.focus() {
1818                        TimelineFocusKind::Thread { .. } | TimelineFocusKind::Event { .. } => false,
1819                        TimelineFocusKind::Live { hide_threaded_events, .. } => {
1820                            *hide_threaded_events
1821                        }
1822                        TimelineFocusKind::PinnedEvents { .. } => true,
1823                    };
1824
1825                    let previous_event = all_remote_events
1826                        .iter()
1827                        .rev()
1828                        // Only consider the events that precede the requested one.
1829                        .skip_while(|event_meta| event_meta.event_id != *event_id)
1830                        .skip(1)
1831                        // Never point a read receipt at one of the user's own events.
1832                        .filter(|event_meta| event_meta.sender.as_deref() != Some(own_user_id))
1833                        .find_map(|event_meta| {
1834                            if !filter_out_thread_events {
1835                                Some(event_meta.event_id.clone())
1836                            } else if event_meta.thread_root_id.is_none() {
1837                                if let Some(TimelineEventItemId::EventId(aggregated_event_id)) =
1838                                    state.meta.aggregations.is_aggregation_of(
1839                                        &TimelineEventItemId::EventId(event_meta.event_id.clone()),
1840                                    )
1841                                    && let Some(target_meta) =
1842                                        all_remote_events.get_by_event_id(aggregated_event_id)
1843                                    && target_meta.thread_root_id.is_some()
1844                                {
1845                                    None
1846                                } else {
1847                                    Some(event_meta.event_id.clone())
1848                                }
1849                            } else {
1850                                None
1851                            }
1852                        });
1853
1854                    match previous_event {
1855                        Some(event_id) => event_id,
1856                        // Nothing from another user to point at. When marking
1857                        // the room as read, fall back to the user's own event
1858                        // so the homeserver still recomputes its push/badge
1859                        // count; otherwise there's nothing to send.
1860                        None if is_marking_room_as_read => event_id.to_owned(),
1861                        None => return SendReceiptDecision::DoNotSend,
1862                    }
1863                } else {
1864                    event_id.to_owned()
1865                }
1866            }
1867
1868            _ => event_id.to_owned(),
1869        };
1870
1871        // Find the real receipt the homeserver already knows about.
1872        let previous_event_id = match receipt_type {
1873            SendReceiptType::Read => state
1874                .meta
1875                .user_receipt(
1876                    own_user_id,
1877                    ReceiptType::Read,
1878                    receipt_thread.clone(),
1879                    room,
1880                    all_remote_events,
1881                    read_receipts::ImplicitReadReceipts::Exclude,
1882                )
1883                .await
1884                .map(|(event_id, _)| event_id),
1885
1886            // Implicit read receipts are saved as public read receipts, so get
1887            // the latest. It also doesn't make sense to have a private read
1888            // receipt behind a public one.
1889            SendReceiptType::ReadPrivate => state
1890                .latest_user_read_receipt(
1891                    own_user_id,
1892                    receipt_thread.clone(),
1893                    room,
1894                    read_receipts::ImplicitReadReceipts::Exclude,
1895                )
1896                .await
1897                .map(|(event_id, _)| event_id),
1898
1899            SendReceiptType::FullyRead => self.room_data_provider.load_fully_read_marker().await,
1900
1901            _ => None,
1902        };
1903
1904        // Don't send anything if the resolved event isn't more recent than
1905        // that.
1906        if let Some(previous_event_id) = previous_event_id {
1907            trace!(%previous_event_id, "found a previous receipt");
1908            if let Some(relative_pos) = TimelineMetadata::compare_events_positions(
1909                &previous_event_id,
1910                &target_event_id,
1911                all_remote_events,
1912            ) && relative_pos != RelativePosition::After
1913            {
1914                return SendReceiptDecision::DoNotSend;
1915            }
1916        }
1917
1918        // No previous receipt was found (or it's an unknown one): let the
1919        // server handle it.
1920        SendReceiptDecision::SendTo(target_event_id)
1921    }
1922
1923    /// Returns the latest event identifier, even if it's not visible, or if
1924    /// it's folded into another timeline item.
1925    pub(crate) async fn latest_event_id(&self) -> Option<OwnedEventId> {
1926        let state = self.state.read().await;
1927        let filter_out_thread_events = match self.focus() {
1928            TimelineFocusKind::Thread { .. } => false,
1929            TimelineFocusKind::Live { hide_threaded_events, .. } => *hide_threaded_events,
1930            TimelineFocusKind::Event { .. } => {
1931                // For event-focused timelines, filtering is handled in the
1932                // event cache layer.
1933                false
1934            }
1935            TimelineFocusKind::PinnedEvents { .. } => true,
1936        };
1937
1938        state
1939            .items
1940            .all_remote_events()
1941            .iter()
1942            .rev()
1943            .filter_map(|event_meta| {
1944                if !filter_out_thread_events {
1945                    // For an unthreaded timeline, the last event is always the
1946                    // latest event.
1947                    Some(event_meta.event_id.clone())
1948                } else if event_meta.thread_root_id.is_none() {
1949                    // For the main-thread timeline, only non-threaded events
1950                    // are valid candidates for the latest event.
1951                    //
1952                    // But! An event could be an aggregation that relate to an
1953                    // in-thread event. In this case, it's not a valid latest
1954                    // event.
1955                    if let Some(TimelineEventItemId::EventId(target_event_id)) =
1956                        state.meta.aggregations.is_aggregation_of(&TimelineEventItemId::EventId(
1957                            event_meta.event_id.clone(),
1958                        ))
1959                        && let Some(target_meta) =
1960                            state.items.all_remote_events().get_by_event_id(target_event_id)
1961                        && target_meta.thread_root_id.is_some()
1962                    {
1963                        // This event is an aggregation of an in-thread event,
1964                        // so skip it.
1965                        None
1966                    } else {
1967                        // Not in a thread, and not the aggregation of an
1968                        // in-thread event, so it's a valid candidate for the
1969                        // latest event.
1970                        Some(event_meta.event_id.clone())
1971                    }
1972                } else {
1973                    // An in-thread event, when we're filtering out threaded
1974                    // events, is never a valid candidate for the latest event.
1975                    None
1976                }
1977            })
1978            .next()
1979    }
1980
1981    #[instrument(skip(self), fields(room_id = ?self.room().room_id()))]
1982    pub(super) async fn retry_event_decryption(&self, session_ids: Option<BTreeSet<String>>) {
1983        let (utds, decrypted) = self.compute_redecryption_candidates().await;
1984
1985        let request = DecryptionRetryRequest {
1986            room_id: self.room().room_id().to_owned(),
1987            utd_session_ids: utds,
1988            refresh_info_session_ids: decrypted,
1989        };
1990
1991        self.room().client().event_cache().request_decryption(request);
1992    }
1993
1994    /// Combine the global (event cache) pagination status with the local state
1995    /// of the timeline.
1996    ///
1997    /// This only changes the global pagination status of this room, in one
1998    /// case: if the timeline has a skip count greater than 0, it will ensure
1999    /// that the pagination status says that we haven't reached the timeline
2000    /// start yet.
2001    pub(super) async fn map_pagination_status(&self, status: PaginationStatus) -> PaginationStatus {
2002        match status {
2003            PaginationStatus::Idle { hit_timeline_start } => {
2004                if hit_timeline_start {
2005                    let state = self.state.read().await;
2006                    // If the skip count is greater than 0, it means that a
2007                    // subsequent pagination could return more items, so pretend
2008                    // we didn't get the information that the timeline start was
2009                    // hit.
2010                    if state.meta.subscriber_skip_count.get() > 0 {
2011                        return PaginationStatus::Idle { hit_timeline_start: false };
2012                    }
2013                }
2014            }
2015            PaginationStatus::Paginating => {}
2016        }
2017
2018        // You're perfect, just the way you are.
2019        status
2020    }
2021}
2022
2023impl<P: RoomDataProvider> TimelineController<P> {
2024    /// Returns the timeline focus of the [`TimelineController`].
2025    pub(super) fn focus(&self) -> &TimelineFocusKind {
2026        &self.focus
2027    }
2028
2029    /// Find an event by ID in this timeline, along with its related events.
2030    ///
2031    /// The related events can be filtered by relation type.
2032    pub(in crate::timeline) async fn find_event_with_relations(
2033        &self,
2034        event_id: &EventId,
2035        filter: Option<Vec<RelationType>>,
2036    ) -> Result<(TimelineEvent, Vec<TimelineEvent>), Error> {
2037        self.room_data_provider
2038            .load_or_fetch_event_with_relations(event_id, filter)
2039            .await
2040            .map_err(Into::into)
2041    }
2042}
2043
2044#[allow(clippy::too_many_arguments)]
2045async fn fetch_replied_to_event<P: RoomDataProvider>(
2046    mut state_guard: RwLockWriteGuard<'_, TimelineState<P>>,
2047    state_lock: &RwLock<TimelineState<P>>,
2048    index: usize,
2049    item: &EventTimelineItem,
2050    internal_id: TimelineUniqueId,
2051    msglike: &MsgLikeContent,
2052    in_reply_to: &EventId,
2053    room: &Room,
2054) -> Result<TimelineDetails<Box<EmbeddedEvent>>, Error> {
2055    if let Some((_, item)) = rfind_event_by_id(&state_guard.items, in_reply_to) {
2056        let details = TimelineDetails::Ready(Box::new(EmbeddedEvent::from_timeline_item(&item)));
2057        trace!("Found replied-to event locally");
2058        return Ok(details);
2059    }
2060
2061    // Replace the item with a new timeline item that has the fetching status of
2062    // the replied-to event to pending.
2063    trace!("Setting in-reply-to details to pending");
2064    let in_reply_to_details =
2065        InReplyToDetails { event_id: in_reply_to.to_owned(), event: TimelineDetails::Pending };
2066
2067    let event_item = item
2068        .with_content(TimelineItemContent::MsgLike(msglike.with_in_reply_to(in_reply_to_details)));
2069
2070    let new_timeline_item = TimelineItem::new(event_item, internal_id);
2071    state_guard.items.replace(index, new_timeline_item);
2072
2073    // Don't hold the state lock while the network request is made.
2074    drop(state_guard);
2075
2076    trace!("Fetching replied-to event");
2077    let res = match room.load_or_fetch_event(in_reply_to, None).await {
2078        Ok(timeline_event) => {
2079            let state = state_lock.read().await;
2080
2081            let replied_to_item =
2082                EmbeddedEvent::try_from_timeline_event(timeline_event, room, &state.meta).await?;
2083
2084            if let Some(item) = replied_to_item {
2085                TimelineDetails::Ready(Box::new(item))
2086            } else {
2087                // The replied-to item is an aggregation, not a standalone item.
2088                return Err(Error::UnsupportedEvent);
2089            }
2090        }
2091
2092        Err(e) => TimelineDetails::Error(Arc::new(e)),
2093    };
2094
2095    Ok(res)
2096}