1use 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
103pub(super) enum SendReceiptDecision {
105 DoNotSend,
107
108 SendTo(OwnedEventId),
113}
114
115#[derive(Debug)]
121pub(in crate::timeline) enum TimelineFocusKind {
122 Live {
124 hide_threaded_events: bool,
126
127 event_cache: RoomEventCache,
129 },
130
131 Event {
134 focused_event_id: OwnedEventId,
136
137 thread_root: OnceLock<OwnedEventId>,
143
144 thread_mode: TimelineEventFocusThreadMode,
147
148 event_cache: EventFocusedCache,
150 },
151
152 Thread {
154 thread_id: OwnedEventId,
156
157 event_cache: ThreadEventCache,
159 },
160
161 PinnedEvents {
162 event_cache: PinnedEventsCache,
164 },
165}
166
167impl TimelineFocusKind {
168 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 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 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 fn is_thread(&self) -> bool {
210 self.thread_root().is_some()
211 }
212
213 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 state: Arc<RwLock<TimelineState<P>>>,
228
229 focus: Arc<TimelineFocusKind>,
231
232 pub(crate) room_data_provider: P,
237
238 pub(super) settings: TimelineSettings,
240}
241
242#[derive(Clone)]
243pub(super) struct TimelineSettings {
244 pub(super) track_read_receipts: TimelineReadReceiptTracking,
247
248 pub(super) event_filter: Arc<TimelineEventFilterFn>,
251
252 pub(super) add_failed_to_parse: bool,
254
255 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
280pub 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 false
296 } else {
297 ev.event_type() != MessageLikeEventType::Reaction
300 }
301 }
302
303 AnySyncTimelineEvent::MessageLike(msg) => {
304 match msg.original_content() {
305 None => {
306 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 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 AnyMessageLikeEventContent::Beacon(_) => false,
352 AnyMessageLikeEventContent::RtcDecline(_) => false,
355
356 _ => false,
357 }
358 }
359 }
360 }
361
362 AnySyncTimelineEvent::State(_) => {
363 true
365 }
366 }
367}
368
369pub(super) struct InitFocusResult {
371 pub has_events: bool,
373 pub focus_task: Option<BackgroundTaskHandle>,
376}
377
378#[derive(Clone, Debug, PartialEq)]
380pub struct ActiveCallInfo {
381 pub active_members: HashSet<OwnedUserId>,
383 pub call_intent: CallIntentConsensus,
385 pub is_joined: bool,
388 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 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 pub async fn handle_encryption_state_changes(&self) {
473 let mut room_info = self.room_data_provider.room_info();
474
475 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 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 break;
495 }
496 }
497 }
498
499 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 let is_live_timeline = true;
517 state.meta.subscriber_skip_count.update(count, is_live_timeline);
518
519 needs
520 }
521
522 pub(super) fn is_live(&self) -> bool {
524 matches!(&*self.focus, TimelineFocusKind::Live { .. })
525 }
526
527 pub(super) fn is_threaded(&self) -> bool {
529 self.focus.is_thread()
530 }
531
532 pub(super) fn thread_root(&self) -> Option<OwnedEventId> {
535 self.focus.thread_root().map(ToOwned::to_owned)
536 }
537
538 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 #[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 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 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 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 if let Some(handle) = &previous.send_handle {
640 if !handle.abort().await.map_err(|err| Error::SendQueueError(err.into()))? {
641 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 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 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 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 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 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 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 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 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 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 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 txn.meta.active_call = maybe_active_call.clone();
870
871 if let Some(existing_event_id) = &txn.meta.active_rtc_notification_event_id {
872 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 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 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 #[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 #[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 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 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 let Some((idx, _)) = local_echo {
967 warn!("Message echo got duplicated, removing the local one");
968 txn.items.remove(idx);
969
970 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 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 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 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 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 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 let mut txn = state.transaction();
1054
1055 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 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 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 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 pub(super) async fn init_focus(&self) -> Result<InitFocusResult, Error> {
1524 match self.focus.deref() {
1525 TimelineFocusKind::Live { event_cache, .. } => {
1526 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 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 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 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 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 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 #[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 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 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 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 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 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 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 .skip_while(|event_meta| event_meta.event_id != *event_id)
1830 .skip(1)
1831 .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 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 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 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 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 SendReceiptDecision::SendTo(target_event_id)
1921 }
1922
1923 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 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 Some(event_meta.event_id.clone())
1948 } else if event_meta.thread_root_id.is_none() {
1949 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 None
1966 } else {
1967 Some(event_meta.event_id.clone())
1971 }
1972 } else {
1973 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 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 state.meta.subscriber_skip_count.get() > 0 {
2011 return PaginationStatus::Idle { hit_timeline_start: false };
2012 }
2013 }
2014 }
2015 PaginationStatus::Paginating => {}
2016 }
2017
2018 status
2020 }
2021}
2022
2023impl<P: RoomDataProvider> TimelineController<P> {
2024 pub(super) fn focus(&self) -> &TimelineFocusKind {
2026 &self.focus
2027 }
2028
2029 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 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 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 return Err(Error::UnsupportedEvent);
2089 }
2090 }
2091
2092 Err(e) => TimelineDetails::Error(Arc::new(e)),
2093 };
2094
2095 Ok(res)
2096}