1#![forbid(missing_docs)]
29
30use std::{
31 collections::HashMap,
32 fmt,
33 ops::Deref,
34 sync::{Arc, OnceLock, RwLock as StdRwLock, RwLockReadGuard, RwLockWriteGuard},
35};
36
37use matrix_sdk_base::{
38 cross_process_lock::CrossProcessLockError,
39 event_cache::store::{EventCacheStoreError, EventCacheStoreLock},
40 linked_chunk::lazy_loader::LazyLoaderError,
41 sync::RoomUpdates,
42 task_monitor::BackgroundTaskHandle,
43};
44use ruma::{EventId, OwnedEventId, OwnedRoomId, RoomId};
45use tokio::sync::{
46 OwnedRwLockReadGuard, OwnedRwLockWriteGuard, RwLock,
47 broadcast::{Receiver, Sender, channel},
48 mpsc,
49};
50use tracing::{error, instrument, trace};
51
52use crate::{
53 Client,
54 client::{ClientInner, WeakClient},
55 paginators::PaginatorError,
56};
57
58mod automatic_pagination;
59mod caches;
60mod deduplicator;
61mod persistence;
62#[cfg(feature = "e2e-encryption")]
63mod redecryptor;
64mod states;
65mod tasks;
66
67#[cfg(feature = "e2e-encryption")]
68pub use redecryptor::{DecryptionRetryRequest, RedecryptorReport};
69
70pub use self::{
71 automatic_pagination::AutomaticPagination,
72 caches::{
73 TimelineVectorDiffs,
74 event_focused::{EventFocusThreadMode, EventFocusedCache, EventFocusedCacheKey},
75 pagination::{BackPaginationOutcome, PaginationStatus},
76 pinned_events::PinnedEventsCache,
77 room::{
78 RoomEventCache, RoomEventCacheGenericUpdate, RoomEventCacheUpdate,
79 pagination::RoomPagination,
80 },
81 subscriber::Subscriber,
82 thread::{ThreadEventCache, pagination::ThreadPagination},
83 },
84};
85use self::{
86 caches::{Caches, room::RoomEventCacheLinkedChunkUpdate, subscriber::AutoShrinkMessage},
87 states::StateLock,
88};
89
90#[derive(thiserror::Error, Clone, Debug)]
92pub enum EventCacheError {
93 #[error(
96 "The EventCache hasn't subscribed to sync responses yet, call `EventCache::subscribe()`"
97 )]
98 NotSubscribedYet,
99
100 #[error("Room cache `{room_id}` is not found.")]
102 RoomNotFound {
103 room_id: OwnedRoomId,
105 },
106
107 #[error("Thread cache `{thread_id}` of room `{room_id}` is not found.")]
109 ThreadNotFound {
110 room_id: OwnedRoomId,
112
113 thread_id: OwnedEventId,
115 },
116
117 #[error("Pinned-events cache for room `{room_id}` are not found.")]
119 PinnedEventsNotFound {
120 room_id: OwnedRoomId,
122 },
123
124 #[error("Event-focused cache `{event_focused_id:?}` of room `{room_id}` is not found.")]
126 EventFocusedNotFound {
127 room_id: OwnedRoomId,
129
130 event_focused_id: EventFocusedCacheKey,
132 },
133
134 #[error("The state of a cache is not found")]
137 CacheStateAlreadyExists,
138
139 #[error(transparent)]
141 PaginationError(Arc<crate::Error>),
142
143 #[error(transparent)]
145 InitialPaginationError(#[from] PaginatorError),
146
147 #[error(transparent)]
149 Storage(#[from] EventCacheStoreError),
150
151 #[error(transparent)]
153 LockingStorage(#[from] CrossProcessLockError),
154
155 #[error("The owning client of the event cache has been dropped.")]
159 ClientDropped,
160
161 #[error(transparent)]
166 LinkedChunkLoader(#[from] LazyLoaderError),
167
168 #[error("Unable to load any of the pinned events.")]
172 UnableToLoadPinnedEvents,
173
174 #[error("the linked chunk metadata is invalid: {details}")]
177 InvalidLinkedChunkMetadata {
178 details: String,
180 },
181}
182
183pub type Result<T> = std::result::Result<T, EventCacheError>;
185
186pub struct EventCacheDropHandles {
188 _listen_updates_task: BackgroundTaskHandle,
190
191 _ignore_user_list_update_task: BackgroundTaskHandle,
193
194 _auto_shrink_linked_chunk_task: BackgroundTaskHandle,
196
197 _thread_subscriber_task: BackgroundTaskHandle,
204
205 #[cfg(feature = "experimental-search")]
212 _search_indexing_task: BackgroundTaskHandle,
213
214 #[cfg(feature = "e2e-encryption")]
216 _redecryptor: redecryptor::Redecryptor,
217}
218
219impl fmt::Debug for EventCacheDropHandles {
220 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
221 f.debug_struct("EventCacheDropHandles").finish_non_exhaustive()
222 }
223}
224
225#[derive(Clone)]
231pub struct EventCache {
232 inner: Arc<EventCacheInner>,
234}
235
236impl fmt::Debug for EventCache {
237 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
238 f.debug_struct("EventCache").finish_non_exhaustive()
239 }
240}
241
242impl EventCache {
243 pub(crate) fn new(client: &Arc<ClientInner>, event_cache_store: EventCacheStoreLock) -> Self {
245 let (generic_update_sender, _) = channel(128);
246 let (linked_chunk_update_sender, _) = channel(128);
247
248 let weak_client = WeakClient::from_inner(client);
249
250 let (thread_subscriber_sender, _thread_subscriber_receiver) = channel(128);
251
252 #[cfg(feature = "e2e-encryption")]
253 let redecryption_channels = redecryptor::RedecryptorChannels::new();
254
255 Self {
256 inner: Arc::new(EventCacheInner {
257 client: weak_client,
258 config: StdRwLock::new(EventCacheConfig::default()),
259 state: StateLock::new(event_cache_store),
260 by_room: Default::default(),
261 drop_handles: Default::default(),
262 auto_shrink_sender: Default::default(),
263 generic_update_sender,
264 linked_chunk_update_sender,
265 #[cfg(feature = "e2e-encryption")]
266 redecryption_channels,
267 automatic_pagination: OnceLock::new(),
268 thread_subscriber_sender,
269 }),
270 }
271 }
272
273 pub fn config(&self) -> RwLockReadGuard<'_, EventCacheConfig> {
276 self.inner.config.read().unwrap()
277 }
278
279 pub fn config_mut(&self) -> RwLockWriteGuard<'_, EventCacheConfig> {
281 self.inner.config.write().unwrap()
282 }
283
284 #[cfg(feature = "testing")]
288 pub fn subscribe_thread_subscriber_updates(&self) -> Receiver<()> {
289 self.inner.thread_subscriber_sender.subscribe()
290 }
291
292 pub fn subscribe(&self) -> Result<()> {
298 let client = self.inner.client()?;
299
300 let _ = self.inner.drop_handles.get_or_init(|| {
302 let task_monitor = client.task_monitor();
303
304 let listen_updates_task = task_monitor.spawn_infinite_task("event_cache::room_updates_task", tasks::room_updates_task(
306 self.inner.clone(),
307 client.subscribe_to_all_room_updates(),
308 )).abort_on_drop();
309
310 let ignore_user_list_update_task = task_monitor.spawn_infinite_task("event_cache::ignore_user_list_update_task", tasks::ignore_user_list_update_task(
311 self.inner.clone(),
312 client.subscribe_to_ignore_user_list_changes(),
313 )).abort_on_drop();
314
315 let (auto_shrink_sender, auto_shrink_receiver) = mpsc::channel(32);
316
317 self.inner.auto_shrink_sender.get_or_init(|| auto_shrink_sender);
319
320 let auto_shrink_linked_chunk_task = task_monitor.spawn_infinite_task("event_cache::auto_shrink_linked_chunk_task", tasks::auto_shrink_linked_chunk_task(
321 Arc::downgrade(&self.inner),
322 auto_shrink_receiver,
323 )).abort_on_drop();
324
325 #[cfg(feature = "e2e-encryption")]
326 let redecryptor = {
327 let receiver = self
328 .inner
329 .redecryption_channels
330 .decryption_request_receiver
331 .lock()
332 .take()
333 .expect("We should have initialized the channel an subscribing should happen only once");
334
335 redecryptor::Redecryptor::new(&client, Arc::downgrade(&self.inner), receiver, &self.inner.linked_chunk_update_sender)
336 };
337
338 let thread_subscriber_task = client
339 .task_monitor()
340 .spawn_infinite_task(
341 "event_cache::thread_subscriber",
342 tasks::thread_subscriber_task(
343 self.inner.client.clone(),
344 self.inner.linked_chunk_update_sender.clone(),
345 self.inner.thread_subscriber_sender.clone(),
346 ),
347 )
348 .abort_on_drop();
349
350 #[cfg(feature = "experimental-search")]
351 let search_indexing_task = client
352 .task_monitor()
353 .spawn_infinite_task(
354 "event_cache::search_indexing",
355 tasks::search_indexing_task(
356 self.inner.client.clone(),
357 self.inner.linked_chunk_update_sender.clone(),
358 ),
359 )
360 .abort_on_drop();
361
362 if self.config().experimental_auto_backpagination {
363 trace!("spawning the automatic paginations API");
366 self.inner.automatic_pagination.get_or_init(|| AutomaticPagination::new(Arc::downgrade(&self.inner), task_monitor));
367 } else {
368 trace!("automatic paginations API is disabled");
369 }
370
371 Arc::new(EventCacheDropHandles {
372 _listen_updates_task: listen_updates_task,
373 _ignore_user_list_update_task: ignore_user_list_update_task,
374 _auto_shrink_linked_chunk_task: auto_shrink_linked_chunk_task,
375 #[cfg(feature = "e2e-encryption")]
376 _redecryptor: redecryptor,
377 _thread_subscriber_task: thread_subscriber_task,
378 #[cfg(feature = "experimental-search")]
379 _search_indexing_task: search_indexing_task,
380 })
381 });
382
383 Ok(())
384 }
385
386 #[doc(hidden)]
388 pub async fn handle_room_updates(&self, updates: RoomUpdates) -> Result<()> {
389 self.inner.handle_room_updates(updates).await
390 }
391
392 pub fn has_subscribed(&self) -> bool {
394 self.inner.drop_handles.get().is_some()
395 }
396
397 pub async fn room(
399 &self,
400 room_id: &RoomId,
401 ) -> Result<(RoomEventCache, Arc<EventCacheDropHandles>)> {
402 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
403 return Err(EventCacheError::NotSubscribedYet);
404 };
405
406 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
407
408 Ok((caches_for_room.room().clone(), drop_handles))
409 }
410
411 pub async fn thread(
413 &self,
414 room_id: &RoomId,
415 thread_id: &EventId,
416 ) -> Result<(ThreadEventCache, Arc<EventCacheDropHandles>)> {
417 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
418 return Err(EventCacheError::NotSubscribedYet);
419 };
420
421 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
422
423 Ok((caches_for_room.thread(thread_id.to_owned()).await?.deref().clone(), drop_handles))
424 }
425
426 pub async fn pinned_events(
428 &self,
429 room_id: &RoomId,
430 ) -> Result<(PinnedEventsCache, Arc<EventCacheDropHandles>)> {
431 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
432 return Err(EventCacheError::NotSubscribedYet);
433 };
434
435 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
436
437 Ok((caches_for_room.pinned_events().await?.clone(), drop_handles))
438 }
439
440 pub async fn event_focused(
442 &self,
443 room_id: &RoomId,
444 event_id: &EventId,
445 thread_mode: EventFocusThreadMode,
446 number_of_initial_events: u16,
447 ) -> Result<(EventFocusedCache, Arc<EventCacheDropHandles>)> {
448 let Some(drop_handles) = self.inner.drop_handles.get().cloned() else {
449 return Err(EventCacheError::NotSubscribedYet);
450 };
451
452 let caches_for_room = self.inner.all_caches_for_room(room_id).await?;
453
454 Ok((
455 caches_for_room
456 .event_focused(event_id.to_owned(), thread_mode, number_of_initial_events)
457 .await?
458 .deref()
459 .clone(),
460 drop_handles,
461 ))
462 }
463
464 pub async fn forget_room(&self, room_id: &RoomId) -> Result<()> {
468 self.inner.forget_room(room_id).await
469 }
470
471 pub async fn clear_all_rooms(&self) -> Result<()> {
475 self.inner.clear_all_rooms().await
476 }
477
478 pub fn subscribe_to_room_generic_updates(&self) -> Receiver<RoomEventCacheGenericUpdate> {
488 self.inner.generic_update_sender.subscribe()
489 }
490
491 pub fn automatic_pagination(&self) -> Option<AutomaticPagination> {
495 self.inner.automatic_pagination.get().cloned()
496 }
497}
498
499#[derive(Clone, Copy, Debug)]
501pub struct EventCacheConfig {
502 pub max_pinned_events_concurrent_requests: usize,
504
505 pub max_pinned_events_to_load: usize,
507
508 pub experimental_auto_backpagination: bool,
512
513 pub room_pagination_per_room_credit: usize,
522
523 pub room_pagination_batch_size: u16,
528}
529
530impl EventCacheConfig {
531 pub const DEFAULT_MAX_EVENTS_TO_LOAD: usize = 128;
533
534 pub const DEFAULT_MAX_CONCURRENT_REQUESTS: usize = 8;
537
538 pub const DEFAULT_ROOM_PAGINATION_CREDITS: usize = 20;
542
543 pub const DEFAULT_ROOM_PAGINATION_BATCH_SIZE: u16 = 30;
547}
548
549impl Default for EventCacheConfig {
550 fn default() -> Self {
551 Self {
552 max_pinned_events_concurrent_requests: Self::DEFAULT_MAX_CONCURRENT_REQUESTS,
553 max_pinned_events_to_load: Self::DEFAULT_MAX_EVENTS_TO_LOAD,
554 room_pagination_per_room_credit: Self::DEFAULT_ROOM_PAGINATION_CREDITS,
555 room_pagination_batch_size: Self::DEFAULT_ROOM_PAGINATION_BATCH_SIZE,
556 experimental_auto_backpagination: false,
557 }
558 }
559}
560
561type CachesByRoom = HashMap<OwnedRoomId, Caches>;
562
563struct EventCacheInner {
564 client: WeakClient,
567
568 config: StdRwLock<EventCacheConfig>,
570
571 state: StateLock,
574
575 by_room: Arc<RwLock<CachesByRoom>>,
579
580 drop_handles: OnceLock<Arc<EventCacheDropHandles>>,
582
583 auto_shrink_sender: OnceLock<mpsc::Sender<AutoShrinkMessage>>,
593
594 generic_update_sender: Sender<RoomEventCacheGenericUpdate>,
599
600 linked_chunk_update_sender: Sender<RoomEventCacheLinkedChunkUpdate>,
608
609 thread_subscriber_sender: Sender<()>,
615
616 #[cfg(feature = "e2e-encryption")]
617 redecryption_channels: redecryptor::RedecryptorChannels,
618
619 automatic_pagination: OnceLock<AutomaticPagination>,
624}
625
626impl EventCacheInner {
627 fn client(&self) -> Result<Client> {
628 self.client.get().ok_or(EventCacheError::ClientDropped)
629 }
630
631 async fn forget_room(&self, room_id: &RoomId) -> Result<()> {
633 let mut caches_for_all_rooms = self.by_room.write().await;
637 self.state.clear_and_reload(&caches_for_all_rooms, Some(room_id)).await?;
638
639 caches_for_all_rooms.remove(room_id);
641
642 Ok(())
643 }
644
645 async fn clear_all_rooms(&self) -> Result<()> {
647 let caches_for_all_rooms = self.by_room.write().await;
675
676 self.state.clear_and_reload(&caches_for_all_rooms, None).await?;
678
679 Ok(())
680 }
681
682 #[instrument(skip(self, updates))]
684 async fn handle_room_updates(&self, updates: RoomUpdates) -> Result<()> {
685 for (room_id, left_room_update) in updates.left {
692 let Ok(caches) = self.all_caches_for_room(&room_id).await else {
693 error!(?room_id, "Room must exist");
694 continue;
695 };
696
697 if let Err(err) = caches.handle_left_room_update(left_room_update).await {
698 error!("handling left room update: {err}");
700 }
701 }
702
703 for (room_id, joined_room_update) in updates.joined {
705 trace!(?room_id, "Handling a `JoinedRoomUpdate`");
706
707 let Ok(caches) = self.all_caches_for_room(&room_id).await else {
708 error!(?room_id, "Room must exist");
709 continue;
710 };
711
712 if let Err(err) = caches.handle_joined_room_update(joined_room_update).await {
713 error!(%room_id, "handling joined room update: {err}");
715 }
716 }
717
718 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.automatic_pagination.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}