diff --git a/crates/matrix-sdk/src/latest_events/mod.rs b/crates/matrix-sdk/src/latest_events/mod.rs index faff4e01f..3a4e071a4 100644 --- a/crates/matrix-sdk/src/latest_events/mod.rs +++ b/crates/matrix-sdk/src/latest_events/mod.rs @@ -141,6 +141,14 @@ impl LatestEvents { Ok(self.state.registered_rooms.for_room(room_id).await?.is_some()) } + /// Check whether the system listens to a particular room. + /// + /// Note: It's a test only method. + #[cfg(test)] + pub async fn is_listening_to_room(&self, room_id: &RoomId) -> bool { + self.state.registered_rooms.rooms.read().await.contains_key(room_id) + } + /// Start listening to updates (if not already) for a particular room, and /// return a [`Subscriber`] to get the current and future /// [`LatestEventValue`]s. diff --git a/crates/matrix-sdk/src/sliding_sync/client.rs b/crates/matrix-sdk/src/sliding_sync/client.rs index 8f624a7b9..01da84214 100644 --- a/crates/matrix-sdk/src/sliding_sync/client.rs +++ b/crates/matrix-sdk/src/sliding_sync/client.rs @@ -16,7 +16,7 @@ use ruma::{ use tracing::error; use super::{SlidingSync, SlidingSyncBuilder}; -use crate::{Client, Result}; +use crate::{Client, Result, sync::subscribe_to_room_latest_events}; /// A sliding sync version. #[derive(Clone, Debug)] @@ -197,6 +197,8 @@ impl SlidingSyncResponseProcessor { response: &http::Response, requested_required_states: &RequestedRequiredStates, ) -> Result<()> { + subscribe_to_room_latest_events(&self.client, response.rooms.keys()).await; + let previously_joined_rooms = self .client .joined_rooms() @@ -380,11 +382,11 @@ async fn handle_receipts_extension( #[cfg(all(test, not(target_family = "wasm")))] mod tests { - use std::collections::BTreeMap; + use std::{collections::BTreeMap, ops::Not}; use assert_matches::assert_matches; use matrix_sdk_base::{ - RequestedRequiredStates, RoomInfoNotableUpdate, RoomInfoNotableUpdateReasons, + RequestedRequiredStates, RoomInfoNotableUpdate, RoomInfoNotableUpdateReasons, RoomState, notification_settings::RoomNotificationMode, }; use matrix_sdk_test::async_test; @@ -637,6 +639,59 @@ mod tests { Ok(()) } + #[async_test] + async fn test_auto_listen_to_latest_events() -> Result<()> { + let client = MockClientBuilder::new(None).build().await; + let room_id = room_id!("!r0"); + + // Create the room beforehand. + client.base_client().get_or_create_room(room_id, RoomState::Joined); + + // Enable the event cache (required for the latest events). + client.event_cache().subscribe()?; + + // The latest event “listener” for this room has NOT been enabled. + assert!(client.latest_events().await.is_listening_to_room(room_id).await.not()); + + // Create the sliding sync client. + let sliding_sync = client + .sliding_sync("test")? + .add_list( + SlidingSyncList::builder("all") + .sync_mode(SlidingSyncMode::new_selective().add_range(0..=10)), + ) + .build() + .await?; + + // Receive a (mocked) response. + { + let server_response = assign!(http::Response::new("0".to_owned()), { + rooms: BTreeMap::from([( + room_id.to_owned(), + http::response::Room::default(), + )]), + }); + + let mut pos_guard = sliding_sync.inner.position.clone().lock_owned().await; + + sliding_sync + .handle_response( + server_response.clone(), + &mut pos_guard, + RequestedRequiredStates::default(), + ) + .await?; + } + + // The room still exists. + assert!(client.get_room(room_id).is_some()); + + // The latest event “listener” for this room has been enabled. + assert!(client.latest_events().await.is_listening_to_room(room_id).await); + + Ok(()) + } + #[async_test] async fn test_read_receipt_can_trigger_a_notable_update_reason() { use ruma::api::client::sync::sync_events::v5 as http; diff --git a/crates/matrix-sdk/src/sync.rs b/crates/matrix-sdk/src/sync.rs index 3eb69251d..5894e91ed 100644 --- a/crates/matrix-sdk/src/sync.rs +++ b/crates/matrix-sdk/src/sync.rs @@ -41,7 +41,7 @@ use ruma::{ serde::Raw, time::Instant, }; -use tracing::{debug, error, warn}; +use tracing::{debug, error, instrument, warn}; use crate::{Client, Result, Room, event_handler::HandlerKind}; @@ -151,6 +151,12 @@ impl Client { &self, response: sync_events::v3::Response, ) -> Result { + subscribe_to_room_latest_events( + self, + response.rooms.join.keys().chain(response.rooms.leave.keys()), + ) + .await; + let response = Box::pin(self.base_client().receive_sync_response(response)).await?; // Some new keys might have been received, so trigger a backup if needed. @@ -343,3 +349,21 @@ impl Client { *last_sync_time = Some(now); } } + +/// Call `LatestEvents::listen_to_room` for rooms in `response`. +/// +/// That way, the latest event is computed and updated for all rooms receiving +/// an update from the sync. +#[instrument(skip_all)] +pub(crate) async fn subscribe_to_room_latest_events<'a, R>(client: &'a Client, room_ids: R) +where + R: Iterator, +{ + let latest_events = client.latest_events().await; + + for room_id in room_ids { + if let Err(error) = latest_events.listen_to_room(room_id).await { + error!(?error, ?room_id, "Failed to listen to the latest event for this room"); + } + } +}