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(())
724 }
725
726 async fn all_caches_for_room(
728 &self,
729 room_id: &RoomId,
730 ) -> Result<OwnedRwLockReadGuard<CachesByRoom, Caches>> {
731 match OwnedRwLockReadGuard::try_map(self.by_room.clone().read_owned().await, |by_room| {
734 by_room.get(room_id)
735 }) {
736 Ok(caches) => Ok(caches),
737
738 Err(by_room_guard) => {
739 drop(by_room_guard);
741 let by_room_guard = self.by_room.clone().write_owned().await;
742
743 let mut by_room_guard =
746 match OwnedRwLockWriteGuard::try_downgrade_map(by_room_guard, |by_room| {
747 by_room.get(room_id)
748 }) {
749 Ok(caches) => return Ok(caches),
750 Err(by_room_guard) => by_room_guard,
751 };
752
753 let caches = Caches::new(
754 &self.client,
755 room_id,
756 self.generic_update_sender.clone(),
757 self.linked_chunk_update_sender.clone(),
758 self.auto_shrink_sender.get().cloned().expect(
761 "we must have called `EventCache::subscribe()` before calling here.",
762 ),
763 &self.state,
764 self.automatic_pagination.get().cloned(),
765 )
766 .await?;
767
768 by_room_guard.insert(room_id.to_owned(), caches);
769
770 Ok(OwnedRwLockWriteGuard::try_downgrade_map(by_room_guard, |by_room| {
771 by_room.get(room_id)
772 })
773 .expect("`Caches` has just been inserted"))
774 }
775 }
776 }
777}
778
779#[derive(Debug, Clone)]
781pub enum EventsOrigin {
782 Sync,
784
785 Pagination,
787
788 Cache,
790}
791
792#[cfg(test)]
793mod tests {
794 use std::{ops::Not, sync::Arc, time::Duration};
795
796 use assert_matches::assert_matches;
797 use futures_util::FutureExt as _;
798 use matrix_sdk_base::{
799 RoomState,
800 linked_chunk::{ChunkIdentifier, LinkedChunkId, Position, Update},
801 sync::{JoinedRoomUpdate, RoomUpdates, Timeline},
802 };
803 use matrix_sdk_test::{
804 JoinedRoomBuilder, SyncResponseBuilder, async_test, event_factory::EventFactory,
805 };
806 use ruma::{event_id, room_id, user_id};
807 use tokio::time::sleep;
808
809 use super::{EventCacheError, RoomEventCacheGenericUpdate};
810 use crate::test_utils::{
811 assert_event_matches_msg, client::MockClientBuilder, logged_in_client,
812 };
813
814 #[async_test]
815 async fn test_must_explicitly_subscribe() {
816 let client = logged_in_client(None).await;
817
818 let event_cache = client.event_cache();
819
820 let room_id = room_id!("!omelette:fromage.fr");
823 let result = event_cache.room(room_id).await;
824
825 assert_matches!(result, Err(EventCacheError::NotSubscribedYet));
828 }
829
830 #[async_test]
831 async fn test_get_event_by_id() {
832 let client = logged_in_client(None).await;
833 let room_id1 = room_id!("!galette:saucisse.bzh");
834 let room_id2 = room_id!("!crepe:saucisse.bzh");
835
836 client.base_client().get_or_create_room(room_id1, RoomState::Joined);
837 client.base_client().get_or_create_room(room_id2, RoomState::Joined);
838
839 let event_cache = client.event_cache();
840 event_cache.subscribe().unwrap();
841
842 let f = EventFactory::new().room(room_id1).sender(user_id!("@ben:saucisse.bzh"));
844
845 let eid1 = event_id!("$1");
846 let eid2 = event_id!("$2");
847 let eid3 = event_id!("$3");
848
849 let joined_room_update1 = JoinedRoomUpdate {
850 timeline: Timeline {
851 events: vec![
852 f.text_msg("hey").event_id(eid1).into(),
853 f.text_msg("you").event_id(eid2).into(),
854 ],
855 ..Default::default()
856 },
857 ..Default::default()
858 };
859
860 let joined_room_update2 = JoinedRoomUpdate {
861 timeline: Timeline {
862 events: vec![f.text_msg("bjr").event_id(eid3).into()],
863 ..Default::default()
864 },
865 ..Default::default()
866 };
867
868 let mut updates = RoomUpdates::default();
869 updates.joined.insert(room_id1.to_owned(), joined_room_update1);
870 updates.joined.insert(room_id2.to_owned(), joined_room_update2);
871
872 event_cache.inner.handle_room_updates(updates).await.unwrap();
874
875 let room1 = client.get_room(room_id1).unwrap();
877
878 let (room_event_cache, _drop_handles) = room1.event_cache().await.unwrap();
879
880 let found1 = room_event_cache.find_event(eid1).await.unwrap().unwrap();
881 assert_event_matches_msg(&found1, "hey");
882
883 let found2 = room_event_cache.find_event(eid2).await.unwrap().unwrap();
884 assert_event_matches_msg(&found2, "you");
885
886 assert!(room_event_cache.find_event(eid3).await.unwrap().is_none());
889 }
890
891 #[async_test]
892 async fn test_generic_update_when_loading_rooms() {
893 let user = user_id!("@mnt_io:matrix.org");
895 let client = logged_in_client(None).await;
896 let room_id_0 = room_id!("!raclette:patate.ch");
897 let room_id_1 = room_id!("!fondue:patate.ch");
898
899 let event_factory = EventFactory::new().room(room_id_0).sender(user);
900
901 let event_cache = client.event_cache();
902 event_cache.subscribe().unwrap();
903
904 client.base_client().get_or_create_room(room_id_0, RoomState::Joined);
905 client.base_client().get_or_create_room(room_id_1, RoomState::Joined);
906
907 client
908 .event_cache_store()
909 .lock()
910 .await
911 .expect("Could not acquire the event cache lock")
912 .as_clean()
913 .expect("Could not acquire a clean event cache lock")
914 .handle_linked_chunk_updates(
915 LinkedChunkId::Room(room_id_0),
916 vec![
917 Update::NewItemsChunk {
919 previous: None,
920 new: ChunkIdentifier::new(0),
921 next: None,
922 },
923 Update::PushItems {
924 at: Position::new(ChunkIdentifier::new(0), 0),
925 items: vec![
926 event_factory
927 .text_msg("hello")
928 .sender(user)
929 .event_id(event_id!("$ev0"))
930 .into_event(),
931 ],
932 },
933 ],
934 )
935 .await
936 .unwrap();
937
938 let mut generic_stream = event_cache.subscribe_to_room_generic_updates();
939
940 {
942 let _room_event_cache = event_cache.room(room_id_0).await.unwrap();
943
944 assert_matches!(
945 generic_stream.recv().await,
946 Ok(RoomEventCacheGenericUpdate { room_id }) => {
947 assert_eq!(room_id, room_id_0);
948 }
949 );
950 }
951
952 {
954 let _room_event_cache = event_cache.room(room_id_1).await.unwrap();
955
956 assert!(generic_stream.recv().now_or_never().is_none());
957 }
958 }
959
960 #[async_test]
961 async fn test_generic_update_when_paginating_room() {
962 let user = user_id!("@mnt_io:matrix.org");
964 let client = logged_in_client(None).await;
965 let room_id = room_id!("!raclette:patate.ch");
966
967 let event_factory = EventFactory::new().room(room_id).sender(user);
968
969 let event_cache = client.event_cache();
970 event_cache.subscribe().unwrap();
971
972 client.base_client().get_or_create_room(room_id, RoomState::Joined);
973
974 client
975 .event_cache_store()
976 .lock()
977 .await
978 .expect("Could not acquire the event cache lock")
979 .as_clean()
980 .expect("Could not acquire a clean event cache lock")
981 .handle_linked_chunk_updates(
982 LinkedChunkId::Room(room_id),
983 vec![
984 Update::NewItemsChunk {
986 previous: None,
987 new: ChunkIdentifier::new(0),
988 next: None,
989 },
990 Update::NewItemsChunk {
992 previous: Some(ChunkIdentifier::new(0)),
993 new: ChunkIdentifier::new(1),
994 next: None,
995 },
996 Update::NewItemsChunk {
998 previous: Some(ChunkIdentifier::new(1)),
999 new: ChunkIdentifier::new(2),
1000 next: None,
1001 },
1002 Update::PushItems {
1003 at: Position::new(ChunkIdentifier::new(2), 0),
1004 items: vec![
1005 event_factory
1006 .text_msg("hello")
1007 .sender(user)
1008 .event_id(event_id!("$ev0"))
1009 .into_event(),
1010 ],
1011 },
1012 Update::NewItemsChunk {
1014 previous: Some(ChunkIdentifier::new(2)),
1015 new: ChunkIdentifier::new(3),
1016 next: None,
1017 },
1018 Update::PushItems {
1019 at: Position::new(ChunkIdentifier::new(3), 0),
1020 items: vec![
1021 event_factory
1022 .text_msg("world")
1023 .sender(user)
1024 .event_id(event_id!("$ev1"))
1025 .into_event(),
1026 ],
1027 },
1028 ],
1029 )
1030 .await
1031 .unwrap();
1032
1033 let mut generic_stream = event_cache.subscribe_to_room_generic_updates();
1034
1035 let (room_event_cache, _) = event_cache.room(room_id).await.unwrap();
1037
1038 assert_matches!(
1039 generic_stream.recv().await,
1040 Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) => {
1041 assert_eq!(room_id, expected_room_id);
1042 }
1043 );
1044
1045 let pagination = room_event_cache.pagination();
1046
1047 let pagination_outcome = pagination.run_backwards_once(1).await.unwrap();
1049
1050 assert_eq!(pagination_outcome.events.len(), 1);
1051 assert!(pagination_outcome.reached_start.not());
1052 assert_matches!(
1053 generic_stream.recv().await,
1054 Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) => {
1055 assert_eq!(room_id, expected_room_id);
1056 }
1057 );
1058
1059 let pagination_outcome = pagination.run_backwards_once(1).await.unwrap();
1061
1062 assert!(pagination_outcome.events.is_empty());
1063 assert!(pagination_outcome.reached_start.not());
1064 assert!(generic_stream.recv().now_or_never().is_none());
1065
1066 let pagination_outcome = pagination.run_backwards_once(1).await.unwrap();
1068
1069 assert!(pagination_outcome.reached_start);
1070 assert!(generic_stream.recv().now_or_never().is_none());
1071 }
1072
1073 #[async_test]
1074 async fn test_for_room_when_room_is_not_found() {
1075 let client = logged_in_client(None).await;
1076 let room_id = room_id!("!raclette:patate.ch");
1077
1078 let event_cache = client.event_cache();
1079 event_cache.subscribe().unwrap();
1080
1081 assert_matches!(
1083 event_cache.room(room_id).await,
1084 Err(EventCacheError::RoomNotFound { room_id: not_found_room_id }) => {
1085 assert_eq!(room_id, not_found_room_id);
1086 }
1087 );
1088
1089 client.base_client().get_or_create_room(room_id, RoomState::Joined);
1091
1092 assert!(event_cache.room(room_id).await.is_ok());
1094 }
1095
1096 #[cfg(not(target_family = "wasm"))]
1099 #[async_test]
1100 async fn test_no_refcycle_event_cache_tasks() {
1101 let client = MockClientBuilder::new(None).build().await;
1102
1103 sleep(Duration::from_secs(1)).await;
1105
1106 let event_cache_weak = Arc::downgrade(&client.event_cache().inner);
1107 assert_eq!(event_cache_weak.strong_count(), 1);
1108
1109 {
1110 let room_id = room_id!("!room:example.org");
1111
1112 let response = SyncResponseBuilder::default()
1114 .add_joined_room(JoinedRoomBuilder::new(room_id))
1115 .build_sync_response();
1116 client.inner.base_client.receive_sync_response(response).await.unwrap();
1117
1118 client.event_cache().subscribe().unwrap();
1119
1120 let (_room_event_cache, _drop_handles) =
1121 client.get_room(room_id).unwrap().event_cache().await.unwrap();
1122 }
1123
1124 drop(client);
1125
1126 sleep(Duration::from_secs(1)).await;
1128
1129 assert_eq!(
1131 event_cache_weak.strong_count(),
1132 0,
1133 "Too many strong references to the event cache {}",
1134 event_cache_weak.strong_count()
1135 );
1136 }
1137}