1#![forbid(missing_docs)]
29
30use std::{
31 collections::HashMap,
32 fmt,
33 num::NonZeroUsize,
34 ops::Deref,
35 sync::{Arc, OnceLock, RwLock as StdRwLock, RwLockReadGuard, RwLockWriteGuard},
36};
37
38use matrix_sdk_base::{
39 cross_process_lock::CrossProcessLockError,
40 event_cache::store::{EventCacheStoreError, EventCacheStoreLock},
41 linked_chunk::lazy_loader::LazyLoaderError,
42 sync::RoomUpdates,
43 task_monitor::BackgroundTaskHandle,
44};
45use ruma::{EventId, OwnedEventId, OwnedRoomId, RoomId};
46use tokio::sync::{
47 OwnedRwLockReadGuard, OwnedRwLockWriteGuard, RwLock,
48 broadcast::{Receiver, Sender, channel},
49 mpsc,
50};
51use tracing::{error, instrument, trace};
52
53use crate::{
54 Client,
55 client::{ClientInner, WeakClient},
56 paginators::PaginatorError,
57};
58
59pub(crate) mod back_pagination_queue;
60mod caches;
61mod deduplicator;
62mod persistence;
63#[cfg(feature = "e2e-encryption")]
64mod redecryptor;
65mod states;
66mod tasks;
67
68#[cfg(feature = "e2e-encryption")]
69pub use redecryptor::{DecryptionRetryRequest, RedecryptorReport};
70
71pub use self::{
72 back_pagination_queue::BackPaginationQueue,
73 caches::{
74 TimelineVectorDiffs,
75 event_focused::{EventFocusThreadMode, EventFocusedCache, EventFocusedCacheKey},
76 pagination::{BackPaginationOutcome, PaginationStatus},
77 pinned_events::PinnedEventsCache,
78 room::{
79 RoomEventCache, RoomEventCacheGenericUpdate, RoomEventCacheUpdate,
80 pagination::RoomPagination,
81 },
82 subscriber::Subscriber,
83 thread::{ThreadEventCache, pagination::ThreadPagination},
84 },
85};
86use self::{
87 caches::{Caches, room::RoomEventCacheLinkedChunkUpdate, subscriber::AutoShrinkMessage},
88 states::StateLock,
89};
90
91#[derive(thiserror::Error, Clone, Debug)]
93pub enum EventCacheError {
94 #[error(
97 "The EventCache hasn't subscribed to sync responses yet, call `EventCache::subscribe()`"
98 )]
99 NotSubscribedYet,
100
101 #[error("Room cache `{room_id}` is not found.")]
103 RoomNotFound {
104 room_id: OwnedRoomId,
106 },
107
108 #[error("Thread cache `{thread_id}` of room `{room_id}` is not found.")]
110 ThreadNotFound {
111 room_id: OwnedRoomId,
113
114 thread_id: OwnedEventId,
116 },
117
118 #[error("Pinned-events cache for room `{room_id}` are not found.")]
120 PinnedEventsNotFound {
121 room_id: OwnedRoomId,
123 },
124
125 #[error("Event-focused cache `{event_focused_id:?}` of room `{room_id}` is not found.")]
127 EventFocusedNotFound {
128 room_id: OwnedRoomId,
130
131 event_focused_id: EventFocusedCacheKey,
133 },
134
135 #[error("The state of a cache is not found")]
138 CacheStateAlreadyExists,
139
140 #[error(transparent)]
142 PaginationError(Arc<crate::Error>),
143
144 #[error(transparent)]
146 InitialPaginationError(#[from] PaginatorError),
147
148 #[error(transparent)]
150 Storage(#[from] EventCacheStoreError),
151
152 #[error(transparent)]
154 LockingStorage(#[from] CrossProcessLockError),
155
156 #[error("The owning client of the event cache has been dropped.")]
160 ClientDropped,
161
162 #[error(transparent)]
167 LinkedChunkLoader(#[from] LazyLoaderError),
168
169 #[error("Unable to load any of the pinned events.")]
173 UnableToLoadPinnedEvents,
174
175 #[error("the linked chunk metadata is invalid: {details}")]
178 InvalidLinkedChunkMetadata {
179 details: String,
181 },
182}
183
184pub type Result<T> = std::result::Result<T, EventCacheError>;
186
187pub struct EventCacheDropHandles {
189 _listen_updates_task: BackgroundTaskHandle,
191
192 _ignore_user_list_update_task: BackgroundTaskHandle,
194
195 _auto_shrink_linked_chunk_task: BackgroundTaskHandle,
197
198 _thread_subscriber_task: BackgroundTaskHandle,
205
206 #[cfg(feature = "experimental-search")]
213 _search_indexing_task: BackgroundTaskHandle,
214
215 #[cfg(feature = "e2e-encryption")]
217 _redecryptor: redecryptor::Redecryptor,
218}
219
220impl fmt::Debug for EventCacheDropHandles {
221 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
222 f.debug_struct("EventCacheDropHandles").finish_non_exhaustive()
223 }
224}
225
226#[derive(Clone)]
232pub struct EventCache {
233 inner: Arc<EventCacheInner>,
235}
236
237impl fmt::Debug for EventCache {
238 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
239 f.debug_struct("EventCache").finish_non_exhaustive()
240 }
241}
242
243impl EventCache {
244 pub(crate) fn new(
246 client: &Arc<ClientInner>,
247 event_cache_store: EventCacheStoreLock,
248 enable_automatic_back_pagination: bool,
249 ) -> Self {
250 let (generic_update_sender, _) = channel(128);
251 let (linked_chunk_update_sender, _) = channel(128);
252
253 let weak_client = WeakClient::from_inner(client);
254
255 let (thread_subscriber_sender, _thread_subscriber_receiver) = channel(128);
256
257 #[cfg(feature = "e2e-encryption")]
258 let redecryption_channels = redecryptor::RedecryptorChannels::new();
259
260 Self {
261 inner: Arc::new(EventCacheInner {
262 client: weak_client,
263 config: StdRwLock::new(EventCacheConfig::default()),
264 state: StateLock::new(event_cache_store),
265 by_room: Default::default(),
266 drop_handles: Default::default(),
267 auto_shrink_sender: Default::default(),
268 generic_update_sender,
269 linked_chunk_update_sender,
270 #[cfg(feature = "e2e-encryption")]
271 redecryption_channels,
272 enable_automatic_back_pagination,
273 back_pagination_queue: OnceLock::new(),
274 thread_subscriber_sender,
275 }),
276 }
277 }
278
279 pub fn config(&self) -> RwLockReadGuard<'_, EventCacheConfig> {
282 self.inner.config.read().unwrap()
283 }
284
285 pub fn config_mut(&self) -> RwLockWriteGuard<'_, EventCacheConfig> {
287 self.inner.config.write().unwrap()
288 }
289
290 #[cfg(feature = "testing")]
294 pub fn subscribe_thread_subscriber_updates(&self) -> Receiver<()> {
295 self.inner.thread_subscriber_sender.subscribe()
296 }
297
298 pub fn subscribe(&self) -> Result<()> {
304 let client = self.inner.client()?;
305
306 let _ = self.inner.drop_handles.get_or_init(|| {
308 let task_monitor = client.task_monitor();
309
310 let listen_updates_task = task_monitor.spawn_infinite_task("event_cache::room_updates_task", tasks::room_updates_task(
312 self.inner.clone(),
313 client.subscribe_to_all_room_updates(),
314 )).abort_on_drop();
315
316 let ignore_user_list_update_task = task_monitor.spawn_infinite_task("event_cache::ignore_user_list_update_task", tasks::ignore_user_list_update_task(
317 self.inner.clone(),
318 client.subscribe_to_ignore_user_list_changes(),
319 )).abort_on_drop();
320
321 let (auto_shrink_sender, auto_shrink_receiver) = mpsc::channel(32);
322
323 self.inner.auto_shrink_sender.get_or_init(|| auto_shrink_sender);
325
326 let auto_shrink_linked_chunk_task = task_monitor.spawn_infinite_task("event_cache::auto_shrink_linked_chunk_task", tasks::auto_shrink_linked_chunk_task(
327 Arc::downgrade(&self.inner),
328 auto_shrink_receiver,
329 )).abort_on_drop();
330
331 #[cfg(feature = "e2e-encryption")]
332 let redecryptor = {
333 let receiver = self
334 .inner
335 .redecryption_channels
336 .decryption_request_receiver
337 .lock()
338 .take()
339 .expect("We should have initialized the channel an subscribing should happen only once");
340
341 redecryptor::Redecryptor::new(&client, Arc::downgrade(&self.inner), receiver, &self.inner.linked_chunk_update_sender)
342 };
343
344 let thread_subscriber_task = client
345 .task_monitor()
346 .spawn_infinite_task(
347 "event_cache::thread_subscriber",
348 tasks::thread_subscriber_task(
349 self.inner.client.clone(),
350 self.inner.linked_chunk_update_sender.clone(),
351 self.inner.thread_subscriber_sender.clone(),
352 ),
353 )
354 .abort_on_drop();
355
356 #[cfg(feature = "experimental-search")]
357 let search_indexing_task = client
358 .task_monitor()
359 .spawn_infinite_task(
360 "event_cache::search_indexing",
361 tasks::search_indexing_task(
362 self.inner.client.clone(),
363 self.inner.linked_chunk_update_sender.clone(),
364 ),
365 )
366 .abort_on_drop();
367
368 if self.inner.enable_automatic_back_pagination {
369 trace!("spawning the back-pagination queue");
371 let max_concurrent = self.config().max_concurrent_back_paginations;
372 self.inner.back_pagination_queue.get_or_init(|| {
373 BackPaginationQueue::new(
374 Arc::downgrade(&self.inner),
375 max_concurrent,
376 task_monitor,
377 )
378 });
379 } else {
380 trace!("back-pagination queue is disabled");
381 }
382
383 Arc::new(EventCacheDropHandles {
384 _listen_updates_task: listen_updates_task,
385 _ignore_user_list_update_task: ignore_user_list_update_task,
386 _auto_shrink_linked_chunk_task: auto_shrink_linked_chunk_task,
387 #[cfg(feature = "e2e-encryption")]
388 _redecryptor: redecryptor,
389 _thread_subscriber_task: thread_subscriber_task,
390 #[cfg(feature = "experimental-search")]
391 _search_indexing_task: search_indexing_task,
392 })
393 });
394
395 Ok(())
396 }
397
398 #[doc(hidden)]
400 pub async fn handle_room_updates(&self, updates: RoomUpdates) -> Result<()> {
401 self.inner.handle_room_updates(updates).await
402 }
403
404 pub fn has_subscribed(&self) -> bool {
406 self.inner.drop_handles.get().is_some()
407 }
408
409 pub async fn room(
411 &self,
412 room_id: &RoomId,
413 ) -> Result<(RoomEventCache, Arc<EventCacheDropHandles>)> {
414 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
415 return Err(EventCacheError::NotSubscribedYet);
416 };
417
418 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
419
420 Ok((caches_for_room.room().clone(), drop_handles))
421 }
422
423 pub async fn thread(
425 &self,
426 room_id: &RoomId,
427 thread_id: &EventId,
428 ) -> Result<(ThreadEventCache, Arc<EventCacheDropHandles>)> {
429 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
430 return Err(EventCacheError::NotSubscribedYet);
431 };
432
433 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
434
435 Ok((caches_for_room.thread(thread_id.to_owned()).await?.deref().clone(), drop_handles))
436 }
437
438 pub async fn pinned_events(
440 &self,
441 room_id: &RoomId,
442 ) -> Result<(PinnedEventsCache, Arc<EventCacheDropHandles>)> {
443 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
444 return Err(EventCacheError::NotSubscribedYet);
445 };
446
447 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
448
449 Ok((caches_for_room.pinned_events().await?.clone(), drop_handles))
450 }
451
452 pub async fn event_focused(
454 &self,
455 room_id: &RoomId,
456 event_id: &EventId,
457 thread_mode: EventFocusThreadMode,
458 number_of_initial_events: u16,
459 ) -> Result<(EventFocusedCache, Arc<EventCacheDropHandles>)> {
460 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
461 return Err(EventCacheError::NotSubscribedYet);
462 };
463
464 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
465
466 Ok((
467 caches_for_room
468 .event_focused(event_id.to_owned(), thread_mode, number_of_initial_events)
469 .await?
470 .deref()
471 .clone(),
472 drop_handles,
473 ))
474 }
475
476 pub async fn forget_room(&self, room_id: &RoomId) -> Result<()> {
480 self.inner.forget_room(room_id).await
481 }
482
483 pub async fn clear_all_rooms(&self) -> Result<()> {
487 self.inner.clear_all_rooms().await
488 }
489
490 pub fn subscribe_to_room_generic_updates(&self) -> Receiver<RoomEventCacheGenericUpdate> {
500 self.inner.generic_update_sender.subscribe()
501 }
502
503 pub fn back_pagination_queue(&self) -> Option<BackPaginationQueue> {
508 self.inner.back_pagination_queue.get().cloned()
509 }
510}
511
512#[derive(Clone, Copy, Debug)]
514pub struct EventCacheConfig {
515 pub max_pinned_events_concurrent_requests: usize,
517
518 pub max_pinned_events_to_load: usize,
520
521 pub max_concurrent_back_paginations: NonZeroUsize,
527}
528
529impl EventCacheConfig {
530 pub const DEFAULT_MAX_EVENTS_TO_LOAD: usize = 128;
532
533 pub const DEFAULT_MAX_CONCURRENT_REQUESTS: usize = 8;
536
537 pub const DEFAULT_MAX_CONCURRENT_BACK_PAGINATIONS: NonZeroUsize = NonZeroUsize::new(3).unwrap();
540}
541
542impl Default for EventCacheConfig {
543 fn default() -> Self {
544 Self {
545 max_pinned_events_concurrent_requests: Self::DEFAULT_MAX_CONCURRENT_REQUESTS,
546 max_pinned_events_to_load: Self::DEFAULT_MAX_EVENTS_TO_LOAD,
547 max_concurrent_back_paginations: Self::DEFAULT_MAX_CONCURRENT_BACK_PAGINATIONS,
548 }
549 }
550}
551
552type CachesByRoom = HashMap<OwnedRoomId, Caches>;
553
554struct EventCacheInner {
555 client: WeakClient,
558
559 config: StdRwLock<EventCacheConfig>,
561
562 state: StateLock,
565
566 by_room: Arc<RwLock<CachesByRoom>>,
570
571 drop_handles: OnceLock<Arc<EventCacheDropHandles>>,
573
574 auto_shrink_sender: OnceLock<mpsc::Sender<AutoShrinkMessage>>,
584
585 generic_update_sender: Sender<RoomEventCacheGenericUpdate>,
590
591 linked_chunk_update_sender: Sender<RoomEventCacheLinkedChunkUpdate>,
599
600 thread_subscriber_sender: Sender<()>,
606
607 #[cfg(feature = "e2e-encryption")]
608 redecryption_channels: redecryptor::RedecryptorChannels,
609
610 enable_automatic_back_pagination: bool,
616
617 back_pagination_queue: OnceLock<BackPaginationQueue>,
622}
623
624impl EventCacheInner {
625 fn client(&self) -> Result<Client> {
626 self.client.get().ok_or(EventCacheError::ClientDropped)
627 }
628
629 async fn forget_room(&self, room_id: &RoomId) -> Result<()> {
631 let mut caches_for_all_rooms = self.by_room.write().await;
635 self.state.clear_and_reload(&caches_for_all_rooms, Some(room_id)).await?;
636
637 caches_for_all_rooms.remove(room_id);
639
640 Ok(())
641 }
642
643 async fn clear_all_rooms(&self) -> Result<()> {
645 let caches_for_all_rooms = self.by_room.write().await;
673
674 self.state.clear_and_reload(&caches_for_all_rooms, None).await?;
676
677 Ok(())
678 }
679
680 #[instrument(skip(self, updates))]
682 async fn handle_room_updates(&self, updates: RoomUpdates) -> Result<()> {
683 for (room_id, left_room_update) in updates.left {
690 let Ok(caches) = self.all_caches_for_room(&room_id).await else {
691 error!(?room_id, "Room must exist");
692 continue;
693 };
694
695 if let Err(err) = caches.handle_left_room_update(left_room_update).await {
696 error!("handling left room update: {err}");
698 }
699 }
700
701 for (room_id, joined_room_update) in updates.joined {
703 trace!(?room_id, "Handling a `JoinedRoomUpdate`");
704
705 let Ok(caches) = self.all_caches_for_room(&room_id).await else {
706 error!(?room_id, "Room must exist");
707 continue;
708 };
709
710 if let Err(err) = caches.handle_joined_room_update(joined_room_update).await {
711 error!(%room_id, "handling joined room update: {err}");
713 }
714 }
715
716 Ok(())
722 }
723
724 async fn all_caches_for_room(
726 &self,
727 room_id: &RoomId,
728 ) -> Result<OwnedRwLockReadGuard<CachesByRoom, Caches>> {
729 match OwnedRwLockReadGuard::try_map(self.by_room.clone().read_owned().await, |by_room| {
732 by_room.get(room_id)
733 }) {
734 Ok(caches) => Ok(caches),
735
736 Err(by_room_guard) => {
737 drop(by_room_guard);
739 let by_room_guard = self.by_room.clone().write_owned().await;
740
741 let mut by_room_guard =
744 match OwnedRwLockWriteGuard::try_downgrade_map(by_room_guard, |by_room| {
745 by_room.get(room_id)
746 }) {
747 Ok(caches) => return Ok(caches),
748 Err(by_room_guard) => by_room_guard,
749 };
750
751 let caches = Caches::new(
752 &self.client,
753 room_id,
754 self.generic_update_sender.clone(),
755 self.linked_chunk_update_sender.clone(),
756 self.auto_shrink_sender.get().cloned().expect(
759 "we must have called `EventCache::subscribe()` before calling here.",
760 ),
761 &self.state,
762 self.back_pagination_queue.get().cloned(),
763 )
764 .await?;
765
766 by_room_guard.insert(room_id.to_owned(), caches);
767
768 Ok(OwnedRwLockWriteGuard::try_downgrade_map(by_room_guard, |by_room| {
769 by_room.get(room_id)
770 })
771 .expect("`Caches` has just been inserted"))
772 }
773 }
774 }
775}
776
777#[derive(Debug, Clone)]
779pub enum EventsOrigin {
780 Sync,
782
783 Pagination,
785
786 Cache,
788}
789
790#[cfg(test)]
791mod tests {
792 use std::{ops::Not, sync::Arc, time::Duration};
793
794 use assert_matches::assert_matches;
795 use futures_util::FutureExt as _;
796 use matrix_sdk_base::{
797 RoomState,
798 linked_chunk::{ChunkIdentifier, LinkedChunkId, Position, Update},
799 sync::{JoinedRoomUpdate, RoomUpdates, Timeline},
800 };
801 use matrix_sdk_test::{
802 JoinedRoomBuilder, SyncResponseBuilder, async_test, event_factory::EventFactory,
803 };
804 use ruma::{event_id, room_id, user_id};
805 use tokio::time::sleep;
806
807 use super::{EventCacheError, RoomEventCacheGenericUpdate};
808 use crate::test_utils::{
809 assert_event_matches_msg, client::MockClientBuilder, logged_in_client,
810 };
811
812 #[async_test]
813 async fn test_must_explicitly_subscribe() {
814 let client = logged_in_client(None).await;
815
816 let event_cache = client.event_cache();
817
818 let room_id = room_id!("!omelette:fromage.fr");
821 let result = event_cache.room(room_id).await;
822
823 assert_matches!(result, Err(EventCacheError::NotSubscribedYet));
826 }
827
828 #[async_test]
829 async fn test_get_event_by_id() {
830 let client = logged_in_client(None).await;
831 let room_id1 = room_id!("!galette:saucisse.bzh");
832 let room_id2 = room_id!("!crepe:saucisse.bzh");
833
834 client.base_client().get_or_create_room(room_id1, RoomState::Joined);
835 client.base_client().get_or_create_room(room_id2, RoomState::Joined);
836
837 let event_cache = client.event_cache();
838 event_cache.subscribe().unwrap();
839
840 let f = EventFactory::new().room(room_id1).sender(user_id!("@ben:saucisse.bzh"));
842
843 let eid1 = event_id!("$1");
844 let eid2 = event_id!("$2");
845 let eid3 = event_id!("$3");
846
847 let joined_room_update1 = JoinedRoomUpdate {
848 timeline: Timeline {
849 events: vec![
850 f.text_msg("hey").event_id(eid1).into(),
851 f.text_msg("you").event_id(eid2).into(),
852 ],
853 ..Default::default()
854 },
855 ..Default::default()
856 };
857
858 let joined_room_update2 = JoinedRoomUpdate {
859 timeline: Timeline {
860 events: vec![f.text_msg("bjr").event_id(eid3).into()],
861 ..Default::default()
862 },
863 ..Default::default()
864 };
865
866 let mut updates = RoomUpdates::default();
867 updates.joined.insert(room_id1.to_owned(), joined_room_update1);
868 updates.joined.insert(room_id2.to_owned(), joined_room_update2);
869
870 event_cache.inner.handle_room_updates(updates).await.unwrap();
872
873 let room1 = client.get_room(room_id1).unwrap();
875
876 let (room_event_cache, _drop_handles) = room1.event_cache().await.unwrap();
877
878 let found1 = room_event_cache.find_event(eid1).await.unwrap().unwrap();
879 assert_event_matches_msg(&found1, "hey");
880
881 let found2 = room_event_cache.find_event(eid2).await.unwrap().unwrap();
882 assert_event_matches_msg(&found2, "you");
883
884 assert!(room_event_cache.find_event(eid3).await.unwrap().is_none());
887 }
888
889 #[async_test]
890 async fn test_generic_update_when_loading_rooms() {
891 let user = user_id!("@mnt_io:matrix.org");
893 let client = logged_in_client(None).await;
894 let room_id_0 = room_id!("!raclette:patate.ch");
895 let room_id_1 = room_id!("!fondue:patate.ch");
896
897 let event_factory = EventFactory::new().room(room_id_0).sender(user);
898
899 let event_cache = client.event_cache();
900 event_cache.subscribe().unwrap();
901
902 client.base_client().get_or_create_room(room_id_0, RoomState::Joined);
903 client.base_client().get_or_create_room(room_id_1, RoomState::Joined);
904
905 client
906 .event_cache_store()
907 .lock()
908 .await
909 .expect("Could not acquire the event cache lock")
910 .as_clean()
911 .expect("Could not acquire a clean event cache lock")
912 .handle_linked_chunk_updates(
913 LinkedChunkId::Room(room_id_0),
914 vec![
915 Update::NewItemsChunk {
917 previous: None,
918 new: ChunkIdentifier::new(0),
919 next: None,
920 },
921 Update::PushItems {
922 at: Position::new(ChunkIdentifier::new(0), 0),
923 items: vec![
924 event_factory
925 .text_msg("hello")
926 .sender(user)
927 .event_id(event_id!("$ev0"))
928 .into_event(),
929 ],
930 },
931 ],
932 )
933 .await
934 .unwrap();
935
936 let mut generic_stream = event_cache.subscribe_to_room_generic_updates();
937
938 {
940 let _room_event_cache = event_cache.room(room_id_0).await.unwrap();
941
942 assert_matches!(
943 generic_stream.recv().await,
944 Ok(RoomEventCacheGenericUpdate { room_id }) => {
945 assert_eq!(room_id, room_id_0);
946 }
947 );
948 }
949
950 {
952 let _room_event_cache = event_cache.room(room_id_1).await.unwrap();
953
954 assert!(generic_stream.recv().now_or_never().is_none());
955 }
956 }
957
958 #[async_test]
959 async fn test_generic_update_when_paginating_room() {
960 let user = user_id!("@mnt_io:matrix.org");
962 let client = logged_in_client(None).await;
963 let room_id = room_id!("!raclette:patate.ch");
964
965 let event_factory = EventFactory::new().room(room_id).sender(user);
966
967 let event_cache = client.event_cache();
968 event_cache.subscribe().unwrap();
969
970 client.base_client().get_or_create_room(room_id, RoomState::Joined);
971
972 client
973 .event_cache_store()
974 .lock()
975 .await
976 .expect("Could not acquire the event cache lock")
977 .as_clean()
978 .expect("Could not acquire a clean event cache lock")
979 .handle_linked_chunk_updates(
980 LinkedChunkId::Room(room_id),
981 vec![
982 Update::NewItemsChunk {
984 previous: None,
985 new: ChunkIdentifier::new(0),
986 next: None,
987 },
988 Update::NewItemsChunk {
990 previous: Some(ChunkIdentifier::new(0)),
991 new: ChunkIdentifier::new(1),
992 next: None,
993 },
994 Update::NewItemsChunk {
996 previous: Some(ChunkIdentifier::new(1)),
997 new: ChunkIdentifier::new(2),
998 next: None,
999 },
1000 Update::PushItems {
1001 at: Position::new(ChunkIdentifier::new(2), 0),
1002 items: vec![
1003 event_factory
1004 .text_msg("hello")
1005 .sender(user)
1006 .event_id(event_id!("$ev0"))
1007 .into_event(),
1008 ],
1009 },
1010 Update::NewItemsChunk {
1012 previous: Some(ChunkIdentifier::new(2)),
1013 new: ChunkIdentifier::new(3),
1014 next: None,
1015 },
1016 Update::PushItems {
1017 at: Position::new(ChunkIdentifier::new(3), 0),
1018 items: vec![
1019 event_factory
1020 .text_msg("world")
1021 .sender(user)
1022 .event_id(event_id!("$ev1"))
1023 .into_event(),
1024 ],
1025 },
1026 ],
1027 )
1028 .await
1029 .unwrap();
1030
1031 let mut generic_stream = event_cache.subscribe_to_room_generic_updates();
1032
1033 let (room_event_cache, _) = event_cache.room(room_id).await.unwrap();
1035
1036 assert_matches!(
1037 generic_stream.recv().await,
1038 Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) => {
1039 assert_eq!(room_id, expected_room_id);
1040 }
1041 );
1042
1043 let pagination = room_event_cache.pagination();
1044
1045 let pagination_outcome = pagination.run_backwards_once(1).await.unwrap();
1047
1048 assert_eq!(pagination_outcome.events.len(), 1);
1049 assert!(pagination_outcome.reached_start.not());
1050 assert_matches!(
1051 generic_stream.recv().await,
1052 Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) => {
1053 assert_eq!(room_id, expected_room_id);
1054 }
1055 );
1056
1057 let pagination_outcome = pagination.run_backwards_once(1).await.unwrap();
1059
1060 assert!(pagination_outcome.events.is_empty());
1061 assert!(pagination_outcome.reached_start.not());
1062 assert!(generic_stream.recv().now_or_never().is_none());
1063
1064 let pagination_outcome = pagination.run_backwards_once(1).await.unwrap();
1066
1067 assert!(pagination_outcome.reached_start);
1068 assert!(generic_stream.recv().now_or_never().is_none());
1069 }
1070
1071 #[async_test]
1072 async fn test_for_room_when_room_is_not_found() {
1073 let client = logged_in_client(None).await;
1074 let room_id = room_id!("!raclette:patate.ch");
1075
1076 let event_cache = client.event_cache();
1077 event_cache.subscribe().unwrap();
1078
1079 assert_matches!(
1081 event_cache.room(room_id).await,
1082 Err(EventCacheError::RoomNotFound { room_id: not_found_room_id }) => {
1083 assert_eq!(room_id, not_found_room_id);
1084 }
1085 );
1086
1087 client.base_client().get_or_create_room(room_id, RoomState::Joined);
1089
1090 assert!(event_cache.room(room_id).await.is_ok());
1092 }
1093
1094 #[cfg(not(target_family = "wasm"))]
1097 #[async_test]
1098 async fn test_no_refcycle_event_cache_tasks() {
1099 let client = MockClientBuilder::new(None).build().await;
1100
1101 sleep(Duration::from_secs(1)).await;
1103
1104 let event_cache_weak = Arc::downgrade(&client.event_cache().inner);
1105 assert_eq!(event_cache_weak.strong_count(), 1);
1106
1107 {
1108 let room_id = room_id!("!room:example.org");
1109
1110 let response = SyncResponseBuilder::default()
1112 .add_joined_room(JoinedRoomBuilder::new(room_id))
1113 .build_sync_response();
1114 client.inner.base_client.receive_sync_response(response).await.unwrap();
1115
1116 client.event_cache().subscribe().unwrap();
1117
1118 let (_room_event_cache, _drop_handles) =
1119 client.get_room(room_id).unwrap().event_cache().await.unwrap();
1120 }
1121
1122 drop(client);
1123
1124 sleep(Duration::from_secs(1)).await;
1126
1127 assert_eq!(
1129 event_cache_weak.strong_count(),
1130 0,
1131 "Too many strong references to the event cache {}",
1132 event_cache_weak.strong_count()
1133 );
1134 }
1135}