diff --git a/crates/matrix-sdk/src/event_cache/redecryptor.rs b/crates/matrix-sdk/src/event_cache/redecryptor.rs index e3727c0e6..48c43dcf3 100644 --- a/crates/matrix-sdk/src/event_cache/redecryptor.rs +++ b/crates/matrix-sdk/src/event_cache/redecryptor.rs @@ -928,7 +928,7 @@ mod tests { use crate::{ Client, assert_let_timeout, encryption::EncryptionSettings, - event_cache::{DecryptionRetryRequest, RoomEventCacheUpdate}, + event_cache::{DecryptionRetryRequest, RoomEventCacheGenericUpdate, RoomEventCacheUpdate}, test_utils::mocks::MatrixMockServer, }; @@ -1236,13 +1236,14 @@ mod tests { // Let's now see what Bob's event cache does. - let (room_cache, _) = bob - .event_cache() + let event_cache = bob.event_cache(); + let (room_cache, _) = event_cache .for_room(room_id) .await .expect("We should be able to get to the event cache for a specific room"); let (_, mut subscriber) = room_cache.subscribe().await.unwrap(); + let mut generic_stream = event_cache.subscribe_to_room_generic_updates(); // We regenerate the Olm machine to check if the room key stream is recreated to // correctly. @@ -1272,6 +1273,12 @@ mod tests { assert_matches!(&diffs[0], VectorDiff::Append { values }); assert_matches!(&values[0].kind, TimelineEventKind::UnableToDecrypt { .. }); + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); + // Now we send the room key to Bob. matrix_mock_server .mock_sync() @@ -1295,6 +1302,12 @@ mod tests { assert_matches!(&diffs[0], VectorDiff::Set { index, value }); assert_eq!(*index, 0); assert_matches!(&value.kind, TimelineEventKind::Decrypted { .. }); + + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); } #[async_test] @@ -1311,14 +1324,15 @@ mod tests { // Let's now see what Bob's event cache does. - let (room_cache, _) = bob - .event_cache() + let event_cache = bob.event_cache(); + let (room_cache, _) = event_cache .for_room(room_id) .instrument(bob_span.clone()) .await .expect("We should be able to get to the event cache for a specific room"); let (_, mut subscriber) = room_cache.subscribe().await.unwrap(); + let mut generic_stream = event_cache.subscribe_to_room_generic_updates(); // Let us forward the event to Bob. matrix_mock_server @@ -1341,6 +1355,12 @@ mod tests { assert_matches!(&diffs[0], VectorDiff::Append { values }); assert_matches!(&values[0].kind, TimelineEventKind::UnableToDecrypt { .. }); + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); + // Now we send the room key to Bob. matrix_mock_server .mock_sync() @@ -1367,8 +1387,14 @@ mod tests { let encryption_info = value.encryption_info().unwrap(); assert_matches!(&encryption_info.verification_state, VerificationState::Unverified(_)); - let session_id = encryption_info.session_id().unwrap().to_owned(); + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); + + let session_id = encryption_info.session_id().unwrap().to_owned(); let alice_user_id = alice.user_id().unwrap(); // Alice now creates the identity. @@ -1410,6 +1436,12 @@ mod tests { VerificationState::Unverified(_), "The event should now know about the identity but still be unverified" ); + + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); } #[async_test] @@ -1425,14 +1457,16 @@ mod tests { let (event, room_key) = prepare_room(&matrix_mock_server, &event_factory, &alice, &bob, room_id).await; + let event_cache = bob.event_cache(); + // Let's now see what Bob's event cache does. - let (room_cache, _) = bob - .event_cache() + let (room_cache, _) = event_cache .for_room(room_id) .await .expect("We should be able to get to the event cache for a specific room"); let (_, mut subscriber) = room_cache.subscribe().await.unwrap(); + let mut generic_stream = event_cache.subscribe_to_room_generic_updates(); // Let us forward the event to Bob. matrix_mock_server @@ -1473,6 +1507,11 @@ mod tests { assert_matches!(&diffs[0], VectorDiff::Append { values }); assert_matches!(&values[0].kind, TimelineEventKind::UnableToDecrypt { .. }); + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + // Bob should receive a new update from the cache. assert_let_timeout!( Duration::from_secs(1), @@ -1484,5 +1523,11 @@ mod tests { assert_matches!(&diffs[0], VectorDiff::Set { index, value }); assert_eq!(*index, 0); assert_matches!(&value.kind, TimelineEventKind::Decrypted { .. }); + + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); } } diff --git a/crates/matrix-sdk/src/event_cache/room/mod.rs b/crates/matrix-sdk/src/event_cache/room/mod.rs index cc8c94e79..34bd4aa8e 100644 --- a/crates/matrix-sdk/src/event_cache/room/mod.rs +++ b/crates/matrix-sdk/src/event_cache/room/mod.rs @@ -2817,6 +2817,7 @@ mod timed_tests { assert_eq!(expected_room_id, room_id); } ); + assert!(generic_stream.is_empty()); // Check the storage. let linked_chunk = from_all_chunks::<3, _, _>( @@ -2896,6 +2897,7 @@ mod timed_tests { assert_eq!(expected_room_id, room_id); } ); + assert!(generic_stream.is_empty()); // The in-memory linked chunk keeps the bundled relation. { @@ -3047,6 +3049,13 @@ mod timed_tests { }); assert!(stream.is_empty()); + + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = + generic_stream.recv() + ); + assert_eq!(room_id, expected_room_id); + assert!(generic_stream.is_empty()); } // After clearing,… @@ -3064,6 +3073,7 @@ mod timed_tests { Ok(RoomEventCacheGenericUpdate { room_id: received_room_id }) = generic_stream.recv() ); assert_eq!(received_room_id, room_id); + assert!(generic_stream.is_empty()); // Events individually are not forgotten by the event cache, after clearing a // room. @@ -3174,6 +3184,7 @@ mod timed_tests { assert_eq!(room_id, expected_room_id); } ); + assert!(generic_stream.is_empty()); let (items, mut stream) = room_event_cache.subscribe().await.unwrap(); @@ -3207,6 +3218,7 @@ mod timed_tests { assert_eq!(expected_room_id, room_id); } ); + assert!(generic_stream.is_empty()); // A new update with one of these events leads to deduplication. let timeline = Timeline { limited: false, prev_batch: None, events: vec![ev2] }; @@ -3338,6 +3350,7 @@ mod timed_tests { assert_eq!(expected_room_id, room_id); } ); + assert!(generic_stream.is_empty()); { let mut state = room_event_cache.inner.state.write().await.unwrap(); @@ -3399,6 +3412,7 @@ mod timed_tests { assert_eq!(expected_room_id, room_id); } ); + assert!(generic_stream.is_empty()); { let state = room_event_cache.inner.state.read().await.unwrap(); @@ -3506,9 +3520,10 @@ mod timed_tests { // Same for the generic update. assert_let_timeout!( - Ok(RoomEventCacheGenericUpdate { room_id: received_room_id }) = generic_stream.recv() + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() ); - assert_eq!(received_room_id, room_id); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); // Shrink the linked chunk to the last chunk. let diffs = room_event_cache @@ -3768,6 +3783,8 @@ mod timed_tests { assert_eq!(events1[0].event_id().as_deref(), Some(evid2)); assert!(stream1.is_empty()); + let mut generic_stream = event_cache.subscribe_to_room_generic_updates(); + // Force loading the full linked chunk by back-paginating. let outcome = room_event_cache.pagination().run_backwards_once(20).await.unwrap(); assert_eq!(outcome.events.len(), 1); @@ -3786,6 +3803,12 @@ mod timed_tests { assert!(stream1.is_empty()); + assert_let_timeout!( + Ok(RoomEventCacheGenericUpdate { room_id: expected_room_id }) = generic_stream.recv() + ); + assert_eq!(expected_room_id, room_id); + assert!(generic_stream.is_empty()); + // Have another subscriber. // Since it's not the first one, and the previous one loaded some more events, // the second subscribers sees them all.