Skip to main content

matrix_sdk/event_cache/caches/thread/
pagination.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
15use std::{fmt, sync::Arc};
16
17use eyeball::SharedObservable;
18use eyeball_im::VectorDiff;
19use matrix_sdk_base::{
20    event_cache::{Event, Gap},
21    linked_chunk::{ChunkContent, LinkedChunkId, Update},
22};
23use ruma::api::Direction;
24use tracing::{error, trace};
25
26use super::{
27    super::{
28        super::{
29            EventCacheError, EventsOrigin, Result, TimelineVectorDiffs,
30            deduplicator::{DeduplicationOutcome, filter_duplicate_events},
31        },
32        pagination::{
33            BackPaginationOutcome, LoadMoreEventsBackwardsOutcome, PaginatedCache, Pagination,
34            SharedPaginationStatus,
35        },
36        room::RoomEventCacheGenericUpdate,
37    },
38    ThreadEventCacheInner,
39    updates::ThreadEventCacheUpdate,
40};
41use crate::room::{IncludeRelations, RelationsOptions};
42
43/// Intermediate type because the `ThreadEventCache` state doesn't provide all
44/// the feature for the moment.
45//
46// TODO: Remove this intermediate type.
47#[derive(Clone)]
48struct ThreadEventCacheWrapper {
49    cache: Arc<ThreadEventCacheInner>,
50
51    // Threads do not support pagination status for the moment but we need one, so let's use a
52    // dummy one for now.
53    dummy_pagination_status: SharedObservable<SharedPaginationStatus>,
54}
55
56/// An API object to run pagination queries on a `ThreadEventCache`.
57#[allow(missing_debug_implementations)]
58pub struct ThreadPagination(Pagination<ThreadEventCacheWrapper>);
59
60impl ThreadPagination {
61    /// Construct a new [`ThreadPagination`].
62    pub(super) fn new(cache: Arc<ThreadEventCacheInner>) -> Self {
63        Self(Pagination::new(ThreadEventCacheWrapper {
64            cache,
65            dummy_pagination_status: SharedObservable::new(SharedPaginationStatus::Idle {
66                hit_timeline_start: false,
67            }),
68        }))
69    }
70
71    /// Starts a back-pagination for the requested number of events.
72    ///
73    /// This automatically takes care of waiting for a pagination token from
74    /// sync, if we haven't done that before.
75    ///
76    /// It will run multiple back-paginations until one of these two conditions
77    /// is met:
78    /// - either we've reached the start of the timeline,
79    /// - or we've obtained enough events to fulfill the requested number of
80    ///   events.
81    pub async fn run_backwards_until(
82        &self,
83        num_requested_events: u16,
84    ) -> Result<BackPaginationOutcome> {
85        self.0.run_backwards_until(num_requested_events).await
86    }
87
88    /// Run a single back-pagination for the requested number of events.
89    ///
90    /// This automatically takes care of waiting for a pagination token from
91    /// sync, if we haven't done that before.
92    pub async fn run_backwards_once(&self, batch_size: u16) -> Result<BackPaginationOutcome> {
93        self.0.run_backwards_once(batch_size).await
94    }
95}
96
97impl PaginatedCache for ThreadEventCacheWrapper {
98    fn status(&self) -> &SharedObservable<SharedPaginationStatus> {
99        &self.dummy_pagination_status
100    }
101
102    async fn load_more_events_backwards(&self) -> Result<LoadMoreEventsBackwardsOutcome> {
103        let mut state = self.cache.state.write().await?;
104
105        // If any in-memory chunk is a gap, don't load more events, and let the caller
106        // resolve the gap.
107        if let Some(prev_token) = state.thread_linked_chunk().rgap().map(|gap| gap.token) {
108            trace!(%prev_token, "thread chunk has at least a gap");
109
110            return Ok(LoadMoreEventsBackwardsOutcome::Gap {
111                prev_token: Some(prev_token),
112                waited_for_initial_prev_token: state.waited_for_initial_prev_token(),
113            });
114        }
115
116        let prev_first_chunk = state.thread_linked_chunk().first_chunk();
117
118        // If we are here, it means all gaps have been resolved (see the `if` block
119        // above). So the first chunk is not a gap, we can load its previous chunk.
120        let linked_chunk_id = LinkedChunkId::Thread(&state.room_id, &state.thread_id);
121        let new_first_chunk = match state
122            .store
123            .load_previous_chunk(linked_chunk_id, prev_first_chunk.identifier())
124            .await
125        {
126            Ok(Some(new_first_chunk)) => {
127                // All good, let's continue with this chunk.
128                new_first_chunk
129            }
130
131            Ok(None) => {
132                // No previous chunk in the store.
133                //
134                // If the first in-memory event is the thread root, it's all good, we have
135                // effectively reached the start of the thread.
136                if let Some((_pos, first_event)) = state.thread_linked_chunk().events().next()
137                    && self.cache.thread_id
138                        == first_event.event_id().expect("Stored events all have an ID")
139                {
140                    trace!("thread chunk is fully loaded and non-empty: reached_start=true");
141
142                    return Ok(LoadMoreEventsBackwardsOutcome::StartOfTimeline);
143                }
144
145                // Otherwise, start back-pagination from the end of the thread.
146                return Ok(LoadMoreEventsBackwardsOutcome::Gap {
147                    prev_token: None,
148                    waited_for_initial_prev_token: state.waited_for_initial_prev_token(),
149                });
150            }
151
152            Err(err) => {
153                error!("error when loading the previous chunk of a linked chunk: {err}");
154
155                // Clear storage for this room.
156                state
157                    .store
158                    .handle_linked_chunk_updates(linked_chunk_id, vec![Update::Clear])
159                    .await?;
160
161                // Return the error.
162                return Err(err.into());
163            }
164        };
165
166        let chunk_content = new_first_chunk.content.clone();
167
168        // We've reached the start on disk, if and only if, there was no chunk prior to
169        // the one we just loaded.
170        //
171        // This value is correct, if and only if, it is used for a chunk content of kind
172        // `Items`.
173        let reached_start = new_first_chunk.previous.is_none();
174
175        if let Err(err) = state.thread_linked_chunk_mut().insert_new_chunk_as_first(new_first_chunk)
176        {
177            error!("error when inserting the previous chunk into its linked chunk: {err}");
178
179            // Clear storage for this thread.
180            state
181                .store
182                .handle_linked_chunk_updates(
183                    LinkedChunkId::Thread(&state.room_id, &state.thread_id),
184                    vec![Update::Clear],
185                )
186                .await?;
187
188            // Return the error.
189            return Err(err.into());
190        }
191
192        // ⚠️ Let's not propagate the updates to the store! We already have these data
193        // in the store! Let's drain them.
194        let _ = state.thread_linked_chunk_mut().store_updates().take();
195
196        // However, we want to get updates as `VectorDiff`s.
197        let timeline_event_diffs = state.thread_linked_chunk_mut().updates_as_vector_diffs();
198
199        Ok(match chunk_content {
200            ChunkContent::Gap(gap) => {
201                trace!("reloaded chunk from disk (gap)");
202
203                LoadMoreEventsBackwardsOutcome::Gap {
204                    prev_token: Some(gap.token),
205                    waited_for_initial_prev_token: state.waited_for_initial_prev_token(),
206                }
207            }
208
209            ChunkContent::Items(events) => {
210                trace!(?reached_start, "reloaded chunk from disk ({} items)", events.len());
211
212                LoadMoreEventsBackwardsOutcome::Events {
213                    events,
214                    timeline_event_diffs,
215                    reached_start,
216                }
217            }
218        })
219    }
220
221    async fn mark_has_waited_for_initial_prev_token(&self) -> Result<()> {
222        *self.cache.state.write().await?.waited_for_initial_prev_token_mut() = true;
223
224        Ok(())
225    }
226
227    async fn wait_for_prev_token(&self) {
228        self.cache.pagination_batch_token_notifier.notified().await
229    }
230
231    async fn paginate_backwards_with_network(
232        &self,
233        batch_size: u16,
234        prev_token: &Option<String>,
235    ) -> Result<Option<(Vec<Event>, Option<String>)>> {
236        let Some(room) = self.cache.weak_room.get() else {
237            // The client is shutting down.
238            return Ok(None);
239        };
240
241        let options = RelationsOptions {
242            from: prev_token.clone(),
243            dir: Direction::Backward,
244            limit: Some(batch_size.into()),
245            include_relations: IncludeRelations::AllRelations,
246            recurse: true,
247        };
248
249        let response = room
250            .relations(self.cache.thread_id.clone(), options)
251            .await
252            .map_err(|err| EventCacheError::PaginationError(Arc::new(err)))?;
253
254        Ok(Some((response.chunk, response.next_batch_token)))
255    }
256
257    async fn conclude_backwards_pagination_from_disk(
258        &self,
259        events: Vec<Event>,
260        timeline_event_diffs: Vec<VectorDiff<Event>>,
261        reached_start: bool,
262    ) -> BackPaginationOutcome {
263        if !timeline_event_diffs.is_empty() {
264            self.cache.update_sender.send(
265                ThreadEventCacheUpdate::UpdateTimelineEvents(TimelineVectorDiffs {
266                    diffs: timeline_event_diffs,
267                    origin: EventsOrigin::Cache,
268                }),
269                Some(RoomEventCacheGenericUpdate { room_id: self.cache.room_id.clone() }),
270            );
271        }
272
273        BackPaginationOutcome {
274            reached_start,
275            // This is a backwards pagination. `BackPaginationOutcome` expects events to
276            // be in “reverse order”.
277            events: events.into_iter().rev().collect(),
278        }
279    }
280
281    async fn conclude_backwards_pagination_from_network(
282        &self,
283        mut events: Vec<Event>,
284        prev_token: Option<String>,
285        mut new_token: Option<String>,
286    ) -> Result<Option<BackPaginationOutcome>> {
287        let Some(room) = self.cache.weak_room.get() else {
288            // The client is shutting down.
289            return Ok(None);
290        };
291
292        // The thread root event is **NOT** part of the `/relations` response.
293        // However, we want the thread root event to be part of the thread itself. It's
294        // easier in a lot of situations. Let's load it if necessary.
295        //
296        // It is necessary to load the thread root event when `new_token` is `None`,
297        // i.e. when we've reached the start of the thread usually.
298        //
299        // We must do this dance before acquiring the state lock because
300        // `Room::load_or_fetch_event` is hitting the state lock too.
301        if new_token.is_none() {
302            events.push(
303                room.load_or_fetch_event(&self.cache.thread_id, None)
304                    .await
305                    .map_err(|err| EventCacheError::PaginationError(Arc::new(err)))?,
306            );
307        }
308
309        let mut state = self.cache.state.write().await?;
310
311        // Check that the previous token still exists; otherwise it's a sign that the
312        // thread's timeline has been cleared.
313        let prev_gap_id = if let Some(token) = prev_token {
314            // Find the corresponding gap in the in-memory linked chunk.
315            let gap_chunk_id = state.thread_linked_chunk().chunk_identifier(|chunk| {
316                    matches!(chunk.content(), ChunkContent::Gap(Gap { token: prev_token }) if *prev_token == token)
317                });
318
319            if gap_chunk_id.is_none() {
320                // We got a previous-batch token from the linked chunk *before* running the
321                // request, but it is missing *after* completing the request.
322                //
323                // It may be a sign the linked chunk has been reset, but it's fine!
324                return Ok(None);
325            }
326
327            gap_chunk_id
328        } else {
329            None
330        };
331
332        let DeduplicationOutcome {
333            all_events: mut events,
334            in_memory_duplicated_event_ids,
335            in_store_duplicated_event_ids,
336            non_empty_all_duplicates: all_duplicates,
337        } = filter_duplicate_events(
338            &state.own_user_id,
339            &state.store,
340            LinkedChunkId::Thread(&state.room_id, &state.thread_id),
341            state.thread_linked_chunk(),
342            events,
343        )
344        .await?;
345
346        // If not all the events have been back-paginated, we need to remove the
347        // previous ones, otherwise we can end up with misordered events.
348        //
349        // Consider the following scenario:
350        // - sync returns [D, E, F]
351        // - then sync returns [] with a previous batch token PB1, so the internal
352        //   linked chunk state is [D, E, F, PB1].
353        // - back-paginating with PB1 may return [A, B, C, D, E, F].
354        //
355        // Only inserting the new events when replacing PB1 would result in a timeline
356        // ordering of [D, E, F, A, B, C], which is incorrect. So we do have to remove
357        // all the events, in case this happens (see also #4746).
358
359        if !all_duplicates {
360            // Let's forget all the previous events.
361            state
362                .remove_events(in_memory_duplicated_event_ids, in_store_duplicated_event_ids)
363                .await?;
364        } else {
365            // All new events are duplicated, they can all be ignored.
366            events.clear();
367            // The gap can be ditched too, as it won't be useful to backpaginate any
368            // further.
369            new_token = None;
370        }
371
372        // `/relations` has been called with `dir=b` (backwards), so the events are in
373        // the inverted order; reorder them.
374        let topo_ordered_events = events.iter().rev().cloned().collect::<Vec<_>>();
375
376        let new_gap = new_token.map(|prev_token| Gap { token: prev_token });
377        let reached_start = state.thread_linked_chunk_mut().push_backwards_pagination_events(
378            prev_gap_id,
379            new_gap,
380            &topo_ordered_events,
381        );
382
383        // Update the store.
384        state.state.propagate_changes(&state.store).await?;
385
386        // A back-pagination can't include new read receipt events, as those are
387        // ephemeral events not included in /relations responses, so we can
388        // safely set the receipt event to None here.
389        //
390        // Note: read receipts may be updated anyhow in the post-processing step, as the
391        // back-pagination may have revealed the event pointed to by the latest read
392        // receipt.
393        let receipt_event = None;
394
395        // Post-process newly inserted events.
396        state.post_process_upserted_events(topo_ordered_events.iter(), receipt_event).await?;
397
398        // Notify observers about the updates.
399        let timeline_event_diffs = state.thread_linked_chunk_mut().updates_as_vector_diffs();
400
401        if !timeline_event_diffs.is_empty() {
402            state.update_sender.send(
403                ThreadEventCacheUpdate::UpdateTimelineEvents(TimelineVectorDiffs {
404                    diffs: timeline_event_diffs,
405                    origin: EventsOrigin::Pagination,
406                }),
407                Some(RoomEventCacheGenericUpdate { room_id: state.room_id.clone() }),
408            );
409        }
410
411        Ok(Some(BackPaginationOutcome { reached_start, events }))
412    }
413}
414
415impl fmt::Debug for ThreadPagination {
416    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
417        formatter.debug_tuple("ThreadPagination").finish_non_exhaustive()
418    }
419}