Skip to main content

matrix_sdk/event_cache/caches/pinned_events/
mod.rs

1// Copyright 2026 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
15mod updates;
16
17use std::{cmp::Ordering, collections::BTreeSet, fmt, sync::Arc};
18
19use eyeball_im::VectorDiff;
20use futures_util::{StreamExt as _, stream};
21use matrix_sdk_base::{
22    apply_redaction,
23    event_cache::{Event, Gap},
24    linked_chunk::{LinkedChunkId, OwnedLinkedChunkId, Position, Update},
25    serde_helpers::{extract_redaction_target, extract_relation},
26    sync::Timeline,
27    task_monitor::BackgroundTaskHandle,
28};
29use matrix_sdk_common::executor::spawn;
30use ruma::{
31    EventId, MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedRoomId, OwnedUserId, RoomId,
32    events::{relation::RelationType, room::redaction::SyncRoomRedactionEvent},
33    room_version_rules::RoomVersionRules,
34};
35use tokio::sync::broadcast::{Receiver, Sender};
36use tracing::{debug, instrument, trace, warn};
37
38pub(super) use self::updates::PinnedEventsCacheUpdateSender;
39#[cfg(feature = "e2e-encryption")]
40use super::super::redecryptor::MaybeResolvedEvent;
41use super::{
42    super::{
43        EventCacheError, EventsOrigin, Result,
44        deduplicator::{DeduplicationOutcome, filter_duplicate_events},
45        persistence::{find_event, send_updates_to_store},
46        states::{
47            CacheStateLock, ReloadPreprocessing, StateLock, StateLockWriteGuard,
48            selectors::PinnedEventsStateSelector,
49        },
50    },
51    EventLocation, TimelineVectorDiffs,
52    event_linked_chunk::{EventLinkedChunk, sort_positions_descending},
53    room::RoomEventCacheLinkedChunkUpdate,
54};
55use crate::{Room, client::WeakClient, config::RequestConfig, room::WeakRoom};
56
57pub struct PinnedEventsCacheState {
58    /// The ID of the room owning this list of pinned events.
59    room_id: OwnedRoomId,
60
61    /// The user's own user id.
62    own_user_id: OwnedUserId,
63
64    /// The rules for the version of this room.
65    room_version_rules: RoomVersionRules,
66
67    /// The linked chunk representing this room's pinned events.
68    ///
69    /// This linked chunk also contains related events. The events are sorted in
70    /// the chronological order (oldest to newest), since it would be otherwise
71    /// impossible to order them correctly, given that we fetch their relations
72    /// over time.
73    chunk: EventLinkedChunk,
74
75    /// Update sender for this pinned events cache.
76    pub update_sender: PinnedEventsCacheUpdateSender,
77
78    /// A sender for the globally observable linked chunk updates that happened
79    /// during a sync or a back-pagination.
80    ///
81    /// See also [`super::super::EventCacheInner::linked_chunk_update_sender`].
82    linked_chunk_update_sender: Sender<RoomEventCacheLinkedChunkUpdate>,
83}
84
85#[cfg(not(tarpaulin_include))]
86impl fmt::Debug for PinnedEventsCacheState {
87    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
88        f.debug_struct("PinnedEventsCacheState")
89            .field("room_id", &self.room_id)
90            .field("chunk", &self.chunk)
91            .finish_non_exhaustive()
92    }
93}
94
95impl<'a> StateLockWriteGuard<'a, PinnedEventsCacheState> {
96    /// Reload the pinned-events: only the last events will be reloaded,
97    /// shrinking the in-memory size of the cache.
98    ///
99    /// If `preprocessing` is set to [`ReloadPreprocessing::ForgetAll`], all
100    /// events will be erased before reloaded.
101    #[must_use = "Propagate `VectorDiff` updates via `TimelineVectorDiffs`"]
102    pub async fn reload(
103        &mut self,
104        preprocessing: ReloadPreprocessing,
105    ) -> Result<Vec<VectorDiff<Event>>> {
106        match preprocessing {
107            ReloadPreprocessing::ForgetAll => {
108                // Clear the `LinkedChunk` and broadcast the updates to the
109                // store.
110                self.state.chunk.reset();
111                self.propagate_changes().await?;
112            }
113
114            ReloadPreprocessing::None => {}
115        }
116
117        // The task will notice there is a desynchronisation and will reload
118        // from network.
119        self.reload_from_storage().await?;
120
121        Ok(self.state.chunk.updates_as_vector_diffs())
122    }
123
124    async fn handle_sync(&mut self, timeline: Timeline) -> Result<()> {
125        let DeduplicationOutcome {
126            all_events: events,
127            in_memory_duplicated_event_ids,
128            in_store_duplicated_event_ids,
129            non_empty_all_duplicates: all_duplicates,
130        } = filter_duplicate_events(
131            &self.state.own_user_id,
132            &self.store,
133            LinkedChunkId::PinnedEvents(&self.state.room_id),
134            &self.state.chunk,
135            timeline.events,
136        )
137        .await?;
138
139        if all_duplicates {
140            // If all events are duplicates, we don't need to do anything;
141            // ignore the new events.
142            return Ok(());
143        }
144
145        // Remove the old duplicated events.
146        //
147        // We don't have to worry about the removals can change the position of
148        // the existing events, because we are pushing all _new_ `events` at the
149        // back.
150        self.remove_events(in_memory_duplicated_event_ids, in_store_duplicated_event_ids).await?;
151
152        // We've found new relations; append them to the linked chunk.
153        self.state.chunk.push_live_events(None, &events);
154
155        self.propagate_changes().await?;
156        self.notify_subscribers(EventsOrigin::Sync);
157
158        // Do stuff for each event.
159        for event in &events {
160            // Handle redaction.
161            self.maybe_apply_new_redaction(event).await?;
162        }
163
164        Ok(())
165    }
166
167    /// Remove events by their position, in `EventLinkedChunk`.
168    ///
169    /// This method is purposely isolated because it must ensure that positions
170    /// are sorted appropriately or it can be disastrous.
171    #[instrument(skip_all)]
172    pub async fn remove_events(
173        &mut self,
174        in_memory_events: Vec<(OwnedEventId, Position)>,
175        in_store_events: Vec<(OwnedEventId, Position)>,
176    ) -> Result<()> {
177        // In-store events.
178        if !in_store_events.is_empty() {
179            let mut positions = in_store_events
180                .into_iter()
181                .map(|(_event_id, position)| position)
182                .collect::<Vec<_>>();
183
184            sort_positions_descending(&mut positions);
185
186            let updates =
187                positions.into_iter().map(|pos| Update::RemoveItem { at: pos }).collect::<Vec<_>>();
188
189            self.apply_store_only_updates(updates).await?;
190        }
191
192        // In-memory events.
193        if in_memory_events.is_empty() {
194            // Nothing else to do, return early.
195            return Ok(());
196        }
197
198        // `remove_events_by_position` is responsible of sorting positions.
199        self.state
200            .chunk
201            .remove_events_by_position(
202                in_memory_events.into_iter().map(|(_event_id, position)| position).collect(),
203            )
204            .expect("failed to remove an event");
205
206        self.propagate_changes().await
207    }
208
209    /// Apply some updates that are effective only on the store itself.
210    ///
211    /// This method should be used only for updates that happen _outside_ the
212    /// in-memory linked chunk. Such updates must be applied onto the persistent
213    /// storage.
214    async fn apply_store_only_updates(&mut self, updates: Vec<Update<Event, Gap>>) -> Result<()> {
215        self.send_updates_to_store(updates).await
216    }
217
218    /// If the given event is a redaction, try to retrieve the to-be-redacted
219    /// event in the chunk, and replace it by the redacted form.
220    #[instrument(skip_all)]
221    async fn maybe_apply_new_redaction(&mut self, event: &Event) -> Result<()> {
222        let Some(event_id) =
223            extract_redaction_target(event.raw(), &self.room_version_rules.redaction)
224        else {
225            return Ok(());
226        };
227
228        // Replace the redacted event by a redacted form, if we knew about it.
229        let Some((location, mut target_event)) = self.find_event(&event_id).await? else {
230            trace!("redacted event is missing from the linked chunk");
231            return Ok(());
232        };
233
234        let target_event_raw = target_event.raw();
235
236        // Don't redact already redacted events.
237        if let Ok(deserialized) = target_event_raw.deserialize()
238            && deserialized.is_redacted()
239        {
240            return Ok(());
241        }
242
243        if let Some(redacted_event) = apply_redaction(
244            target_event_raw,
245            event.raw().cast_ref_unchecked::<SyncRoomRedactionEvent>(),
246            &self.room_version_rules.redaction,
247        ) {
248            // It's safe to cast `redacted_event` here:
249            //
250            // - either the event was an `AnyTimelineEvent` cast to
251            //   `AnySyncTimelineEvent` when calling .raw(), so it's still one
252            //   under the hood.
253            // - or it wasn't, and it's a plain `AnySyncTimelineEvent` in this
254            //   case.
255            target_event.replace_raw(redacted_event.cast_unchecked());
256
257            self.replace_event_at(location, target_event.clone()).await?;
258        }
259
260        Ok(())
261    }
262
263    /// See documentation of [`find_event`].
264    pub(super) async fn find_event(
265        &self,
266        event_id: &EventId,
267    ) -> Result<Option<(EventLocation, Event)>> {
268        find_event(event_id, &self.room_id, &self.chunk, &self.store).await
269    }
270
271    /// Replaces a single event, be it saved in memory or in the store.
272    ///
273    /// If it was saved in memory, this will emit a notification to observers
274    /// that a single item has been replaced. Otherwise, such a notification is
275    /// not emitted, because observers are unlikely to observe the store updates
276    /// directly.
277    pub async fn replace_event_at(
278        &mut self,
279        location: EventLocation,
280        new_event: Event,
281    ) -> Result<()> {
282        match location {
283            EventLocation::Memory(position) => {
284                self.state
285                    .chunk
286                    .replace_event_at(position, new_event)
287                    .expect("should have been a valid position of an item");
288                // We just changed the in-memory representation; synchronize
289                // this with the store.
290                self.propagate_changes().await?;
291            }
292            EventLocation::Store => {
293                self.save_events([new_event]).await?;
294            }
295        }
296
297        Ok(())
298    }
299
300    /// Save events into the database, without notifying observers.
301    pub async fn save_events(&mut self, events: impl IntoIterator<Item = Event>) -> Result<()> {
302        let store = self.store.clone();
303        let room_id = self.state.room_id.clone();
304        let events = events.into_iter().collect::<Vec<_>>();
305
306        // Spawn a task so the save is uninterrupted by task cancellation.
307        spawn(async move {
308            for event in events {
309                store.save_event(&room_id, event).await?;
310            }
311
312            Result::Ok(())
313        })
314        .await
315        .expect("joining failed")?;
316
317        Ok(())
318    }
319
320    /// Reload all the pinned events from storage, replacing the current linked
321    /// chunk.
322    async fn reload_from_storage(&mut self) -> Result<()> {
323        let room_id = self.state.room_id.clone();
324        let linked_chunk_id = LinkedChunkId::PinnedEvents(&room_id);
325
326        let (last_chunk, chunk_id_gen) = self.store.load_last_chunk(linked_chunk_id).await?;
327
328        let Some(last_chunk) = last_chunk else {
329            // No pinned events stored, make sure the in-memory linked chunk is
330            // sync'd (i.e. empty), and return.
331            if self.state.chunk.events().next().is_some() {
332                self.state.chunk.reset();
333                self.notify_subscribers(EventsOrigin::Sync);
334            }
335
336            return Ok(());
337        };
338
339        {
340            let mut current_chunk_identifier = last_chunk.identifier;
341            self.state.chunk.shrink_to_last_reloaded_chunk(
342                Some(last_chunk),
343                chunk_id_gen,
344                // This cache doesn't use the `OrderTracker`.
345                None,
346            )?;
347
348            // Reload the entire chunk.
349            while let Some(previous_chunk) =
350                self.store.load_previous_chunk(linked_chunk_id, current_chunk_identifier).await?
351            {
352                current_chunk_identifier = previous_chunk.identifier;
353                self.state.chunk.insert_new_chunk_as_first(previous_chunk)?;
354            }
355        }
356
357        // Empty store updates, since we just reloaded from storage.
358        self.state.chunk.store_updates().take();
359
360        // Let observers know about it.
361        self.notify_subscribers(EventsOrigin::Cache);
362
363        Ok(())
364    }
365
366    async fn replace_all_events(&mut self, new_events: Vec<Event>) -> Result<()> {
367        trace!("resetting all pinned events in linked chunk");
368
369        let previous_pinned_event_ids = self.state.current_event_ids();
370
371        if new_events
372            .iter()
373            .filter_map(|e| e.event_id())
374            .map(ToOwned::to_owned)
375            .collect::<BTreeSet<_>>()
376            == previous_pinned_event_ids.into_iter().collect()
377        {
378            // No change in the list of pinned events.
379            return Ok(());
380        }
381
382        if self.state.chunk.events().next().is_some() {
383            self.state.chunk.reset();
384        }
385
386        self.state.chunk.push_live_events(None, &new_events);
387        self.propagate_changes().await?;
388        self.notify_subscribers(EventsOrigin::Sync);
389
390        Ok(())
391    }
392
393    /// Propagate the changes in this linked chunk to observers, and save the
394    /// changes on disk.
395    pub async fn propagate_changes(&mut self) -> Result<()> {
396        let updates = self.state.chunk.store_updates().take();
397
398        self.send_updates_to_store(updates).await
399    }
400
401    async fn send_updates_to_store(&mut self, updates: Vec<Update<Event, Gap>>) -> Result<()> {
402        let linked_chunk_id = OwnedLinkedChunkId::PinnedEvents(self.room_id.clone());
403
404        send_updates_to_store(
405            &self.store,
406            linked_chunk_id,
407            &self.state.linked_chunk_update_sender,
408            updates,
409        )
410        .await
411    }
412
413    /// Notify subscribers of timeline updates.
414    fn notify_subscribers(&mut self, origin: EventsOrigin) {
415        let diffs = self.state.chunk.updates_as_vector_diffs();
416
417        if !diffs.is_empty() {
418            self.update_sender.send(TimelineVectorDiffs { diffs, origin });
419        }
420    }
421}
422
423impl PinnedEventsCacheState {
424    /// Return a list of the current event IDs in this linked chunk.
425    pub(super) fn current_event_ids(&self) -> Vec<OwnedEventId> {
426        self.chunk
427            .events()
428            .filter_map(|(_position, event)| event.event_id().map(ToOwned::to_owned))
429            .collect()
430    }
431
432    /// Returns whether this contains exactly the given pinned events,
433    /// along with the events related to them (reactions, edits, redactions,
434    /// etc).
435    ///
436    /// Related events are ignored here when comparing the pinned event IDs
437    /// to all of the [`Self::current_event_ids`], as that would cause
438    /// differences if one pinned event has a related event, and then reload
439    /// them all over again.
440    fn has_exactly_pinned_events(&self, pinned_event_ids: &[OwnedEventId]) -> bool {
441        let pinned_event_ids: BTreeSet<&EventId> =
442            pinned_event_ids.iter().map(|event_id| &**event_id).collect();
443        let event_ids: BTreeSet<&EventId> =
444            self.chunk.events().filter_map(|(_position, event)| event.event_id()).collect();
445
446        if !pinned_event_ids.is_subset(&event_ids) {
447            return false;
448        }
449
450        // Every other event must relate to an event in this linked chunk,
451        // just like `aggregate_timeline_for_pinned_events` does for a sync.
452        // If not, then that event isn't pinned anymore.
453        self.chunk.events().all(|(_position, event)| {
454            event.event_id().is_some_and(|event_id| pinned_event_ids.contains(event_id))
455                || extract_relation(event.raw()).is_some_and(|(relation_type, related_event_id)| {
456                    relation_type != RelationType::Thread && event_ids.contains(&*related_event_id)
457                })
458                || extract_redaction_target(event.raw(), &self.room_version_rules.redaction)
459                    .is_some_and(|redacted_event_id| event_ids.contains(&*redacted_event_id))
460        })
461    }
462}
463
464/// All the information related to room's pinned events..
465///
466/// Cloning is shallow, and thus is cheap to do.
467#[derive(Clone)]
468pub struct PinnedEventsCache {
469    inner: Arc<PinnedEventsCacheInner>,
470
471    /// The task handling the refreshing of pinned events for this specific
472    /// room.
473    _task: Arc<BackgroundTaskHandle>,
474}
475
476/// The (non-cloneable) details of the `PinnedEventsCache`.
477struct PinnedEventsCacheInner {
478    /// The ID of the room owning this list of pinned events.
479    room_id: OwnedRoomId,
480
481    /// State of this `PinnedEventsCache`.
482    ///
483    /// It is behind an `Arc` because it is shared with the task.
484    state: CacheStateLock<PinnedEventsStateSelector>,
485}
486
487impl PinnedEventsCache {
488    /// Creates a new [`PinnedEventsCache`] for the given room.
489    pub(in super::super) async fn new(
490        weak_room: &WeakRoom,
491        own_user_id: OwnedUserId,
492        room_version_rules: RoomVersionRules,
493        linked_chunk_update_sender: Sender<RoomEventCacheLinkedChunkUpdate>,
494        state: &StateLock,
495    ) -> Result<Self> {
496        let room = weak_room.get().ok_or(EventCacheError::ClientDropped)?;
497        let room_id = room.room_id().to_owned();
498
499        let cache_state = state
500            .try_insert_once_with(
501                PinnedEventsStateSelector::new(room_id.clone()),
502                |_store_guard| async {
503                    Ok(PinnedEventsCacheState {
504                        room_id: room_id.clone(),
505                        own_user_id,
506                        room_version_rules,
507                        chunk: EventLinkedChunk::new(),
508                        update_sender: PinnedEventsCacheUpdateSender::new(),
509                        linked_chunk_update_sender,
510                    })
511                },
512            )
513            .await?;
514
515        let inner = Arc::new(PinnedEventsCacheInner { room_id, state: cache_state });
516
517        let task = room
518            .client()
519            .task_monitor()
520            .spawn_infinite_task(
521                "pinned_event_listener_task",
522                Self::pinned_event_listener_task(room, inner.clone()),
523            )
524            .abort_on_drop();
525
526        Ok(Self { inner, _task: Arc::new(task) })
527    }
528
529    /// Get the room ID of this cache.
530    pub fn room_id(&self) -> &RoomId {
531        &self.inner.room_id
532    }
533
534    /// Return a reference to the state.
535    pub(super) fn state(&self) -> &CacheStateLock<PinnedEventsStateSelector> {
536        &self.inner.state
537    }
538
539    /// Subscribe to live events from this room's pinned events cache.
540    pub async fn subscribe(&self) -> Result<(Vec<Event>, Receiver<TimelineVectorDiffs>)> {
541        let guard = self.inner.state.read().await?;
542        let events = guard.state.chunk.events().map(|(_position, item)| item.clone()).collect();
543
544        let recv = guard.state.update_sender.new_pinned_events_receiver();
545
546        Ok((events, recv))
547    }
548
549    /// Try to locate the events in the linked chunk corresponding to the given
550    /// list of resolved events, and replace them, while alerting observers
551    /// about the update.
552    #[cfg(feature = "e2e-encryption")]
553    pub(in super::super) async fn replace_in_memory_utds(
554        &self,
555        resolved_events: &[MaybeResolvedEvent],
556    ) -> Result<()> {
557        let mut state = self.inner.state.write().await?;
558
559        // Drain the updates to the store, events have already been updated
560        // before calling this method.
561        let _ = state.state.chunk.store_updates().take();
562
563        if state.state.chunk.replace_utds(resolved_events) {
564            state.propagate_changes().await?;
565            state.notify_subscribers(EventsOrigin::Cache);
566        }
567
568        Ok(())
569    }
570
571    /// Handle an update from a joined room.
572    #[instrument(skip_all, fields(room_id = %self.inner.room_id))]
573    pub(super) async fn handle_joined_room_update(&self, timeline: Timeline) -> Result<()> {
574        self.handle_timeline(timeline).await
575    }
576
577    /// Handle an update from a left room.
578    #[instrument(skip_all, fields(room_id = %self.inner.room_id))]
579    pub(super) async fn handle_left_room_update(&self, timeline: Timeline) -> Result<()> {
580        self.handle_timeline(timeline).await
581    }
582
583    /// Handle a [`Timeline`], i.e. new events received by a sync for this
584    /// thread.
585    async fn handle_timeline(&self, timeline: Timeline) -> Result<()> {
586        if timeline.events.is_empty() {
587            return Ok(());
588        }
589
590        trace!("adding new {} events", timeline.events.len());
591
592        self.inner.state.write().await?.handle_sync(timeline).await
593    }
594
595    #[instrument(fields(%room_id = room.room_id()), skip(room, inner))]
596    async fn pinned_event_listener_task(room: Room, inner: Arc<PinnedEventsCacheInner>) {
597        debug!("pinned events listener task started");
598
599        let reload_from_network = async |room: Room| {
600            let events = match Self::reload_pinned_events(room).await {
601                Ok(Some(events)) => events,
602                Ok(None) => Vec::new(),
603                Err(err) => {
604                    warn!("error when loading pinned events: {err}");
605                    return;
606                }
607            };
608
609            // Replace the whole linked chunk with those new events, and
610            // propagate updates to the observers.
611            match inner.state.write().await {
612                Ok(mut guard) => {
613                    guard.replace_all_events(events).await.unwrap_or_else(|err| {
614                        warn!("error when replacing pinned events: {err}");
615                    });
616                }
617
618                Err(err) => {
619                    warn!("error when acquiring write lock to replace pinned events: {err}");
620                }
621            }
622        };
623
624        // Reload the pinned events from the storage first.
625        match inner.state.write().await {
626            Ok(mut guard) => {
627                // On startup, reload the pinned events from storage.
628                guard.reload_from_storage().await.unwrap_or_else(|err| {
629                    warn!("error when reloading pinned events from storage, at start: {err}");
630                });
631
632                // Compare the initial list of pinned events to the one in the
633                // linked chunk.
634                let actual_pinned_events =
635                    pinned_event_ids_to_load(&room, room.pinned_event_ids().unwrap_or_default());
636
637                if !guard.state.has_exactly_pinned_events(&actual_pinned_events) {
638                    // Reload the list of pinned events from network.
639                    drop(guard);
640                    reload_from_network(room.clone()).await;
641                }
642            }
643
644            Err(err) => {
645                warn!("error when acquiring write lock to initialize pinned events: {err}");
646            }
647        }
648
649        let weak_room =
650            WeakRoom::new(WeakClient::from_client(&room.client()), room.room_id().to_owned());
651
652        let mut stream = room.pinned_event_ids_stream();
653
654        drop(room);
655
656        // Whenever the list of pinned events changes, reload it.
657        while let Some(new_list) = stream.next().await {
658            trace!("handling update");
659
660            let Some(room) = weak_room.get() else {
661                debug!("room has been dropped, ending pinned events listener task");
662                break;
663            };
664
665            let new_list = pinned_event_ids_to_load(&room, new_list);
666
667            let guard = match inner.state.read().await {
668                Ok(guard) => guard,
669                Err(err) => {
670                    warn!("error when acquiring read lock to handle pinned events update: {err}");
671                    break;
672                }
673            };
674
675            // Compare to the current linked chunk.
676            if guard.state.has_exactly_pinned_events(&new_list) {
677                // All the events in the pinned list are the same, don't reload.
678                continue;
679            }
680
681            drop(guard);
682
683            // Event IDs differ, so reload all the pinned events.
684            reload_from_network(room).await;
685        }
686
687        debug!("pinned events listener task ended");
688    }
689
690    /// Loads the pinned events in this room, using the cache first and then
691    /// requesting the event from the homeserver if it couldn't be found. This
692    /// method will perform as many concurrent requests for events as
693    /// `max_concurrent_requests` allows, to avoid overwhelming the server.
694    ///
695    /// Returns `None` if the list of pinned events hasn't changed since the
696    /// previous time we loaded them. May return an error if there was an issue
697    /// fetching the full events.
698    async fn reload_pinned_events(room: Room) -> Result<Option<Vec<Event>>> {
699        let max_concurrent_requests =
700            room.client().event_cache().config().max_pinned_events_concurrent_requests;
701
702        let pinned_event_ids =
703            pinned_event_ids_to_load(&room, room.pinned_event_ids().unwrap_or_default());
704
705        if pinned_event_ids.is_empty() {
706            return Ok(Some(Vec::new()));
707        }
708
709        let mut num_successful_loads = 0;
710
711        let mut loaded_events: Vec<Event> =
712            stream::iter(pinned_event_ids.clone().into_iter().map(|event_id| {
713                let room = room.clone();
714                let filter = vec![RelationType::Annotation, RelationType::Replacement];
715                let request_config = RequestConfig::default().retry_limit(3);
716
717                async move {
718                    let (target, mut relations) = room
719                        .load_or_fetch_event_with_relations(
720                            &event_id,
721                            Some(filter),
722                            Some(request_config),
723                        )
724                        .await?;
725
726                    relations.insert(0, target);
727                    Ok::<_, crate::Error>(relations)
728                }
729            }))
730            .buffer_unordered(max_concurrent_requests)
731            // Count successful queries.
732            .inspect(|result| {
733                if result.is_ok() {
734                    num_successful_loads += 1;
735                }
736            })
737            // Get rid of error results.
738            .flat_map(stream::iter)
739            // Flatten the list of `Vec<Event>` into a list of `Event`.
740            .flat_map(stream::iter)
741            .collect()
742            .await;
743
744        if num_successful_loads != pinned_event_ids.len() {
745            warn!(
746                "only successfully loaded {} out of {} pinned events",
747                num_successful_loads,
748                pinned_event_ids.len()
749            );
750        }
751
752        if loaded_events.is_empty() {
753            // If the list of loaded events is empty, we ran into an error to
754            // load _all_ the pinned events, which needs to be reported to the
755            // caller.
756            return Err(EventCacheError::UnableToLoadPinnedEvents);
757        }
758
759        // Since we have all the events and their related events, we can't
760        // nicely sort them, since we've lost all ordering information from
761        // using /event or /relations. Resort to sorting using chronological
762        // ordering (oldest -> newest).
763        loaded_events.sort_by(compare_pinned_items);
764
765        Ok(Some(loaded_events))
766    }
767}
768
769impl fmt::Debug for PinnedEventsCache {
770    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
771        f.debug_struct("PinnedEventsCache").finish_non_exhaustive()
772    }
773}
774
775/// Returns the IDs of the pinned events that this cache loads
776/// from the given list of all the room's pinned events.
777///
778/// This'll include only the most recently pinned ones, up to the limit of
779/// the chosen `max_pinned_events_to_load` cfg value.
780fn pinned_event_ids_to_load(room: &Room, pinned_event_ids: Vec<OwnedEventId>) -> Vec<OwnedEventId> {
781    let max_events_to_load = room.client().event_cache().config().max_pinned_events_to_load;
782    pinned_event_ids.into_iter().rev().take(max_events_to_load).rev().collect()
783}
784
785fn compare_pinned_items(a: &Event, b: &Event) -> Ordering {
786    let a_time: Option<MilliSecondsSinceUnixEpoch> = a.timestamp_raw();
787    let b_time: Option<MilliSecondsSinceUnixEpoch> = b.timestamp_raw();
788
789    compare_by_optional_timestamp(a_time, b_time)
790}
791
792fn compare_by_optional_timestamp(
793    a: Option<MilliSecondsSinceUnixEpoch>,
794    b: Option<MilliSecondsSinceUnixEpoch>,
795) -> Ordering {
796    match (a, b) {
797        (None, None) => Ordering::Equal,
798        (None, Some(_)) => Ordering::Greater,
799        (Some(_), None) => Ordering::Less,
800        (Some(a), Some(b)) => a.cmp(&b),
801    }
802}
803
804#[cfg(not(target_family = "wasm"))]
805#[cfg(test)]
806mod tests {
807    use proptest::prelude::*;
808    use ruma::UInt;
809
810    use super::*;
811
812    fn any_timestamp() -> impl Strategy<Value = Option<MilliSecondsSinceUnixEpoch>> {
813        prop::option::of(
814            any::<u32>().prop_map(|value| MilliSecondsSinceUnixEpoch(UInt::from(value))),
815        )
816    }
817
818    #[test]
819    fn sort_pinned_events_never_panics_only_nones() {
820        let mut vec = vec![None; 100_000];
821        vec.sort_by(|a, b| compare_by_optional_timestamp(*a, *b))
822    }
823
824    proptest! {
825    #[test]
826    fn sort_pinned_events_never_panics(mut v in prop::collection::vec(any_timestamp(), 0..1000)) {
827        v.sort_by(
828            |a, b| compare_by_optional_timestamp(*a, *b))
829    }
830
831    #[test]
832    fn compare_pinned_events_reflexive(a in any_timestamp()) {
833        prop_assert_eq!(compare_by_optional_timestamp(a, a), Ordering::Equal);
834    }
835
836    #[test]
837    fn compare_pinned_events_antisymmetric(a in any_timestamp(), b in any_timestamp()) {
838        let ab = compare_by_optional_timestamp(a, b);
839        let ba = compare_by_optional_timestamp(b, a);
840
841        prop_assert_eq!(ab, ba.reverse());
842    }
843
844    #[test]
845    fn compare_pinned_events_transitive(
846        a in any_timestamp(),
847        b in any_timestamp(),
848        c in any_timestamp()
849    ) {
850        let ab = compare_by_optional_timestamp(a, b);
851        let bc = compare_by_optional_timestamp(b, c);
852        let ac = compare_by_optional_timestamp(a, c);
853
854        if ab == Ordering::Less && bc == Ordering::Less {
855            prop_assert_eq!(ac, Ordering::Less);
856        }
857
858        if ab == Ordering::Equal && bc == Ordering::Equal {
859            prop_assert_eq!(ac, Ordering::Equal);
860        }
861
862        if ab == Ordering::Greater && bc == Ordering::Greater {
863            prop_assert_eq!(ac, Ordering::Greater);
864        }
865    }
866    }
867}