1use 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
33pub const DEFAULT_CHUNK_CAPACITY: usize = 128;
37
38#[cfg_attr(target_family = "wasm", async_trait(?Send))]
41#[cfg_attr(not(target_family = "wasm"), async_trait)]
42pub trait EventCacheStore: AsyncTraitDeps {
43 type Error: fmt::Debug + Into<EventCacheStoreError>;
45
46 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 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 #[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 async fn load_all_chunks_metadata(
76 &self,
77 linked_chunk_id: LinkedChunkId<'_>,
78 ) -> Result<Vec<ChunkMetadata>, Self::Error>;
79
80 async fn load_last_chunk(
85 &self,
86 linked_chunk_id: LinkedChunkId<'_>,
87 ) -> Result<(Option<RawChunk<Event, Gap>>, ChunkIdentifierGenerator), Self::Error>;
88
89 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 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 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 async fn clear_all_events(&self, room_id: Option<&RoomId>) -> Result<(), Self::Error>;
140
141 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 async fn find_event(
154 &self,
155 room_id: &RoomId,
156 event_id: &EventId,
157 ) -> Result<Option<Event>, Self::Error>;
158
159 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 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 async fn save_event(&self, room_id: &RoomId, event: Event) -> Result<(), Self::Error>;
202
203 async fn close(&self) -> Result<(), Self::Error>;
209
210 async fn reopen(&self) -> Result<(), Self::Error>;
213
214 #[doc(hidden)]
220 async fn optimize(&self) -> Result<(), Self::Error>;
221
222 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
371pub type DynEventCacheStore = dyn EventCacheStore<Error = EventCacheStoreError>;
373
374pub 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
399impl<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 unsafe { Arc::from_raw(ptr_erased) }
411 }
412}