test(sdk): Test RoomEventCacheGenericUpdates are broadcasted.
This commit is contained in:
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user