Skip to main content

matrix_sdk_base/event_cache/store/
traits.rs

1// Copyright 2024 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 async_trait::async_trait;
18use matrix_sdk_common::{
19    AsyncTraitDeps,
20    cross_process_lock::CrossProcessLockGeneration,
21    linked_chunk::{
22        ChunkIdentifier, ChunkIdentifierGenerator, ChunkMetadata, LinkedChunkId, Position,
23        RawChunk, Update,
24    },
25};
26use ruma::{EventId, OwnedEventId, RoomId, events::relation::RelationType};
27
28use super::{
29    super::{Event, Gap, thread::ThreadInfo},
30    EventCacheStoreError,
31};
32
33/// A default capacity for linked chunks, when manipulating in conjunction with
34/// an `EventCacheStore` implementation.
35// TODO: move back?
36pub const DEFAULT_CHUNK_CAPACITY: usize = 128;
37
38/// An abstract trait that can be used to implement different store backends for
39/// the event cache of the SDK.
40#[cfg_attr(target_family = "wasm", async_trait(?Send))]
41#[cfg_attr(not(target_family = "wasm"), async_trait)]
42pub trait EventCacheStore: AsyncTraitDeps {
43    /// The error type used by this event cache store.
44    type Error: fmt::Debug + Into<EventCacheStoreError>;
45
46    /// Try to take a lock using the given store.
47    async fn try_take_leased_lock(
48        &self,
49        lease_duration_ms: u32,
50        key: &str,
51        holder: &str,
52    ) -> Result<Option<CrossProcessLockGeneration>, Self::Error>;
53
54    /// An [`Update`] reflects an operation that has happened inside a linked
55    /// chunk. The linked chunk is used by the event cache to store the events
56    /// in-memory. This method aims at forwarding this update inside this store.
57    async fn handle_linked_chunk_updates(
58        &self,
59        linked_chunk_id: LinkedChunkId<'_>,
60        updates: Vec<Update<Event, Gap>>,
61    ) -> Result<(), Self::Error>;
62
63    /// Return all the raw components of a linked chunk, so the caller may
64    /// reconstruct the linked chunk later.
65    #[doc(hidden)]
66    async fn load_all_chunks(
67        &self,
68        linked_chunk_id: LinkedChunkId<'_>,
69    ) -> Result<Vec<RawChunk<Event, Gap>>, Self::Error>;
70
71    /// Load all of the chunks' metadata for the given [`LinkedChunkId`].
72    ///
73    /// Chunks are unordered, and there's no guarantee that the chunks would
74    /// form a valid linked chunk after reconstruction.
75    async fn load_all_chunks_metadata(
76        &self,
77        linked_chunk_id: LinkedChunkId<'_>,
78    ) -> Result<Vec<ChunkMetadata>, Self::Error>;
79
80    /// Load the last chunk of the `LinkedChunk` holding all events of the room
81    /// identified by `room_id`.
82    ///
83    /// This is used to iteratively load events for the `EventCache`.
84    async fn load_last_chunk(
85        &self,
86        linked_chunk_id: LinkedChunkId<'_>,
87    ) -> Result<(Option<RawChunk<Event, Gap>>, ChunkIdentifierGenerator), Self::Error>;
88
89    /// Load the chunk before the chunk identified by `before_chunk_identifier`
90    /// of the `LinkedChunk` holding all events of the room identified by
91    /// `room_id`
92    ///
93    /// This is used to iteratively load events for the `EventCache`.
94    async fn load_previous_chunk(
95        &self,
96        linked_chunk_id: LinkedChunkId<'_>,
97        before_chunk_identifier: ChunkIdentifier,
98    ) -> Result<Option<RawChunk<Event, Gap>>, Self::Error>;
99
100    /// Load the [`ThreadInfo`] associated to `room_id` and `thread_id`.
101    ///
102    /// If `insert_default_if_missing` is `true`, this method **must create and
103    /// insert** the `ThreadInfo` if it doesn't exist. Consequently, in this
104    /// context, this method is also a way to remember a thread, and will
105    /// always return `Some(_)`.
106    ///
107    /// It does nothing regarding events or linked chunks. This is important if
108    /// one wants to list all threads, or remove specific events or linked
109    /// chunks.
110    async fn load_thread_info(
111        &self,
112        room_id: &RoomId,
113        thread_id: &EventId,
114        insert_default_if_missing: bool,
115    ) -> Result<Option<ThreadInfo>, Self::Error>;
116
117    /// Update the [`ThreadInfo`] associated to `room_id` and `thread_id`.
118    ///
119    /// If it does not exist, this method **must fail**! Normally, the
120    /// `ThreadInfo` must be created automatically with
121    /// [`Self::load_thread_info`], so it must always exist.
122    async fn update_thread_info(
123        &self,
124        room_id: &RoomId,
125        thread_id: &EventId,
126        thread_info: &ThreadInfo,
127    ) -> Result<(), Self::Error>;
128
129    /// Clear persisted events for all the rooms if `room_id` is `None`, or a
130    /// single room otherwise.
131    ///
132    /// This will empty and remove all the linked chunks stored previously,
133    /// using the above [`Self::handle_linked_chunk_updates`] methods. It _also_
134    /// deletes all the events' content.
135    ///
136    /// ⚠ This is meant only for super specific use cases, where there shouldn't
137    /// be any live in-memory linked chunks. In general, prefer using
138    /// `EventCache::clear_all_rooms()` from the common SDK crate.
139    async fn clear_all_events(&self, room_id: Option<&RoomId>) -> Result<(), Self::Error>;
140
141    /// Given a set of event IDs, return the duplicated events along with their
142    /// position if there are any.
143    async fn filter_duplicated_events(
144        &self,
145        linked_chunk_id: LinkedChunkId<'_>,
146        events: Vec<OwnedEventId>,
147    ) -> Result<Vec<(OwnedEventId, Position)>, Self::Error>;
148
149    /// Find an event by its ID in a room.
150    ///
151    /// This method must return events saved either in any linked chunks, _or_
152    /// events saved "out-of-band" with the [`Self::save_event`] method.
153    async fn find_event(
154        &self,
155        room_id: &RoomId,
156        event_id: &EventId,
157    ) -> Result<Option<Event>, Self::Error>;
158
159    /// Find all the events (alongside their position in the room's linked
160    /// chunk, if available) that relate to a given event.
161    ///
162    /// The only events which don't have a position are those which have been
163    /// saved out-of-band using [`Self::save_event`].
164    ///
165    /// Note: it doesn't process relations recursively: for instance, if
166    /// requesting only thread events, it will NOT return the aggregated events
167    /// affecting the returned events. It is the responsibility of the caller to
168    /// do so, if needed.
169    ///
170    /// An additional filter can be provided to only retrieve related events for
171    /// a certain relationship.
172    ///
173    /// This method must return events saved either in any linked chunks, _or_
174    /// events saved "out-of-band" with the [`Self::save_event`] method.
175    async fn find_event_relations(
176        &self,
177        room_id: &RoomId,
178        event_id: &EventId,
179        filter: Option<&[RelationType]>,
180    ) -> Result<Vec<(Event, Option<Position>)>, Self::Error>;
181
182    /// Get all events in this room.
183    ///
184    /// This method must return events saved either in any linked chunks, _or_
185    /// events saved "out-of-band" with the [`Self::save_event`] method.
186    async fn get_room_events(
187        &self,
188        room_id: &RoomId,
189        event_type: Option<&str>,
190        session_id: Option<&str>,
191    ) -> Result<Vec<Event>, Self::Error>;
192
193    /// Save an event, that might or might not be part of an existing linked
194    /// chunk.
195    ///
196    /// If the event has no event id, it will not be saved, and the function
197    /// must return an Ok result early.
198    ///
199    /// If the event was already stored with the same id, it must be replaced,
200    /// without causing an error.
201    async fn save_event(&self, room_id: &RoomId, event: Event) -> Result<(), Self::Error>;
202
203    /// Close the store, releasing all held resources (database connections,
204    /// file descriptors, file locks).
205    ///
206    /// In-flight operations complete before this method returns. After it
207    /// returns, operations will fail until [`Self::reopen()`] is called.
208    async fn close(&self) -> Result<(), Self::Error>;
209
210    /// Reopen the store after a [`Self::close()`], re-acquiring database
211    /// connections.
212    async fn reopen(&self) -> Result<(), Self::Error>;
213
214    /// Perform database optimizations if any are available, i.e. vacuuming in
215    /// SQLite.
216    ///
217    /// **Warning:** this was added to check if SQLite fragmentation was the
218    /// source of performance issues, **DO NOT use in production**.
219    #[doc(hidden)]
220    async fn optimize(&self) -> Result<(), Self::Error>;
221
222    /// Returns the size of the store in bytes, if known.
223    async fn get_size(&self) -> Result<Option<usize>, Self::Error>;
224}
225
226#[repr(transparent)]
227struct EraseEventCacheStoreError<T>(T);
228
229#[cfg(not(tarpaulin_include))]
230impl<T: fmt::Debug> fmt::Debug for EraseEventCacheStoreError<T> {
231    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
232        self.0.fmt(f)
233    }
234}
235
236#[cfg_attr(target_family = "wasm", async_trait(?Send))]
237#[cfg_attr(not(target_family = "wasm"), async_trait)]
238impl<T: EventCacheStore> EventCacheStore for EraseEventCacheStoreError<T> {
239    type Error = EventCacheStoreError;
240
241    async fn try_take_leased_lock(
242        &self,
243        lease_duration_ms: u32,
244        key: &str,
245        holder: &str,
246    ) -> Result<Option<CrossProcessLockGeneration>, Self::Error> {
247        self.0.try_take_leased_lock(lease_duration_ms, key, holder).await.map_err(Into::into)
248    }
249
250    async fn handle_linked_chunk_updates(
251        &self,
252        linked_chunk_id: LinkedChunkId<'_>,
253        updates: Vec<Update<Event, Gap>>,
254    ) -> Result<(), Self::Error> {
255        self.0.handle_linked_chunk_updates(linked_chunk_id, updates).await.map_err(Into::into)
256    }
257
258    async fn load_all_chunks(
259        &self,
260        linked_chunk_id: LinkedChunkId<'_>,
261    ) -> Result<Vec<RawChunk<Event, Gap>>, Self::Error> {
262        self.0.load_all_chunks(linked_chunk_id).await.map_err(Into::into)
263    }
264
265    async fn load_all_chunks_metadata(
266        &self,
267        linked_chunk_id: LinkedChunkId<'_>,
268    ) -> Result<Vec<ChunkMetadata>, Self::Error> {
269        self.0.load_all_chunks_metadata(linked_chunk_id).await.map_err(Into::into)
270    }
271
272    async fn load_last_chunk(
273        &self,
274        linked_chunk_id: LinkedChunkId<'_>,
275    ) -> Result<(Option<RawChunk<Event, Gap>>, ChunkIdentifierGenerator), Self::Error> {
276        self.0.load_last_chunk(linked_chunk_id).await.map_err(Into::into)
277    }
278
279    async fn load_previous_chunk(
280        &self,
281        linked_chunk_id: LinkedChunkId<'_>,
282        before_chunk_identifier: ChunkIdentifier,
283    ) -> Result<Option<RawChunk<Event, Gap>>, Self::Error> {
284        self.0
285            .load_previous_chunk(linked_chunk_id, before_chunk_identifier)
286            .await
287            .map_err(Into::into)
288    }
289
290    async fn load_thread_info(
291        &self,
292        room_id: &RoomId,
293        thread_id: &EventId,
294        insert_default_if_missing: bool,
295    ) -> Result<Option<ThreadInfo>, Self::Error> {
296        self.0
297            .load_thread_info(room_id, thread_id, insert_default_if_missing)
298            .await
299            .map_err(Into::into)
300    }
301
302    async fn update_thread_info(
303        &self,
304        room_id: &RoomId,
305        thread_id: &EventId,
306        thread_info: &ThreadInfo,
307    ) -> Result<(), Self::Error> {
308        self.0.update_thread_info(room_id, thread_id, thread_info).await.map_err(Into::into)
309    }
310
311    async fn clear_all_events(&self, room_id: Option<&RoomId>) -> Result<(), Self::Error> {
312        self.0.clear_all_events(room_id).await.map_err(Into::into)
313    }
314
315    async fn filter_duplicated_events(
316        &self,
317        linked_chunk_id: LinkedChunkId<'_>,
318        events: Vec<OwnedEventId>,
319    ) -> Result<Vec<(OwnedEventId, Position)>, Self::Error> {
320        self.0.filter_duplicated_events(linked_chunk_id, events).await.map_err(Into::into)
321    }
322
323    async fn find_event(
324        &self,
325        room_id: &RoomId,
326        event_id: &EventId,
327    ) -> Result<Option<Event>, Self::Error> {
328        self.0.find_event(room_id, event_id).await.map_err(Into::into)
329    }
330
331    async fn find_event_relations(
332        &self,
333        room_id: &RoomId,
334        event_id: &EventId,
335        filter: Option<&[RelationType]>,
336    ) -> Result<Vec<(Event, Option<Position>)>, Self::Error> {
337        self.0.find_event_relations(room_id, event_id, filter).await.map_err(Into::into)
338    }
339
340    async fn get_room_events(
341        &self,
342        room_id: &RoomId,
343        event_type: Option<&str>,
344        session_id: Option<&str>,
345    ) -> Result<Vec<Event>, Self::Error> {
346        self.0.get_room_events(room_id, event_type, session_id).await.map_err(Into::into)
347    }
348
349    async fn save_event(&self, room_id: &RoomId, event: Event) -> Result<(), Self::Error> {
350        self.0.save_event(room_id, event).await.map_err(Into::into)
351    }
352
353    async fn close(&self) -> Result<(), Self::Error> {
354        self.0.close().await.map_err(Into::into)
355    }
356
357    async fn reopen(&self) -> Result<(), Self::Error> {
358        self.0.reopen().await.map_err(Into::into)
359    }
360
361    async fn optimize(&self) -> Result<(), Self::Error> {
362        self.0.optimize().await.map_err(Into::into)?;
363        Ok(())
364    }
365
366    async fn get_size(&self) -> Result<Option<usize>, Self::Error> {
367        Ok(self.0.get_size().await.map_err(Into::into)?)
368    }
369}
370
371/// A type-erased [`EventCacheStore`].
372pub type DynEventCacheStore = dyn EventCacheStore<Error = EventCacheStoreError>;
373
374/// A type that can be type-erased into `Arc<dyn EventCacheStore>`.
375///
376/// This trait is not meant to be implemented directly outside
377/// `matrix-sdk-base`, but it is automatically implemented for everything that
378/// implements `EventCacheStore`.
379pub trait IntoEventCacheStore {
380    #[doc(hidden)]
381    fn into_event_cache_store(self) -> Arc<DynEventCacheStore>;
382}
383
384impl IntoEventCacheStore for Arc<DynEventCacheStore> {
385    fn into_event_cache_store(self) -> Arc<DynEventCacheStore> {
386        self
387    }
388}
389
390impl<T> IntoEventCacheStore for T
391where
392    T: EventCacheStore + Sized + 'static,
393{
394    fn into_event_cache_store(self) -> Arc<DynEventCacheStore> {
395        Arc::new(EraseEventCacheStoreError(self))
396    }
397}
398
399// Turns a given `Arc<T>` into `Arc<DynEventCacheStore>` by attaching the
400// `EventCacheStore` impl vtable of `EraseEventCacheStoreError<T>`.
401impl<T> IntoEventCacheStore for Arc<T>
402where
403    T: EventCacheStore + 'static,
404{
405    fn into_event_cache_store(self) -> Arc<DynEventCacheStore> {
406        let ptr: *const T = Arc::into_raw(self);
407        let ptr_erased = ptr as *const EraseEventCacheStoreError<T>;
408        // SAFETY: EraseEventCacheStoreError is repr(transparent) so T and
409        // EraseEventCacheStoreError<T> have the same layout and ABI
410        unsafe { Arc::from_raw(ptr_erased) }
411    }
412}