Skip to content

Commit 32b5377

Browse files
committed
event_cache: test SpecificEventsCache
Loading a set with its relations in chronological order, sync updates reaching the cache (reactions and redactions) while unrelated events don't, set_event_ids being a no-op for the same set and reloading when it grows or shrinks, and the not-subscribed error.
1 parent 0dd2524 commit 32b5377

2 files changed

Lines changed: 305 additions & 0 deletions

File tree

‎crates/matrix-sdk/tests/integration/event_cache/mod.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@ use ruma::{
3838
use tokio::{spawn, sync::broadcast, task::yield_now, time::sleep};
3939

4040
mod read_receipts;
41+
mod specific_events;
4142
mod threads;
4243

4344
macro_rules! assert_event_id {
Lines changed: 304 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,304 @@
1+
use std::time::Duration;
2+
3+
use imbl::Vector;
4+
use matrix_sdk::{
5+
event_cache::{EventCacheError, TimelineVectorDiffs},
6+
test_utils::mocks::{MatrixMockServer, RoomRelationsResponseTemplate},
7+
timeout::timeout,
8+
};
9+
use matrix_sdk_base::event_cache::Event;
10+
use matrix_sdk_test::{ALICE, JoinedRoomBuilder, async_test, event_factory::EventFactory};
11+
use ruma::{
12+
EventId, OwnedEventId, event_id, events::room::message::RoomMessageEventContentWithoutRelation,
13+
room_id,
14+
};
15+
use tokio::sync::broadcast::Receiver;
16+
17+
fn room_id() -> &'static ruma::RoomId {
18+
room_id!("!galette:saucisse.bzh")
19+
}
20+
21+
/// Wait for the next update, apply it and whatever follows it closely
22+
/// (one sync can produce several) to `events`.
23+
async fn apply_next_updates(
24+
subscriber: &mut Receiver<TimelineVectorDiffs>,
25+
events: &mut Vector<Event>,
26+
) {
27+
let mut wait = Duration::from_secs(3);
28+
29+
while let Ok(update) = timeout(subscriber.recv(), wait).await {
30+
for diff in update.expect("the update channel should be open").diffs {
31+
diff.apply(events);
32+
}
33+
wait = Duration::from_millis(300);
34+
}
35+
36+
assert!(wait < Duration::from_secs(3), "an update should arrive");
37+
}
38+
39+
/// Assert that no update arrives for a while.
40+
async fn assert_no_update(subscriber: &mut Receiver<TimelineVectorDiffs>) {
41+
assert!(timeout(subscriber.recv(), Duration::from_millis(300)).await.is_err());
42+
}
43+
44+
fn event_ids(events: &Vector<Event>) -> Vec<OwnedEventId> {
45+
events.iter().filter_map(|event| event.event_id().map(ToOwned::to_owned)).collect()
46+
}
47+
48+
fn find<'a>(events: &'a Vector<Event>, event_id: &EventId) -> &'a Event {
49+
events.iter().find(|event| event.event_id() == Some(event_id)).expect("event should be present")
50+
}
51+
52+
async fn subscribed_client(server: &MatrixMockServer) -> matrix_sdk::Client {
53+
let client = server.client_builder().build().await;
54+
client.event_cache().subscribe().unwrap();
55+
server.sync_joined_room(&client, room_id()).await;
56+
client
57+
}
58+
59+
#[async_test]
60+
async fn test_specific_events_are_loaded_with_their_relations() {
61+
let f = EventFactory::new().room(room_id()).sender(*ALICE);
62+
63+
let first = f.text_msg("first").event_id(event_id!("$first")).server_ts(1).into_event();
64+
let second = f.text_msg("second").event_id(event_id!("$second")).server_ts(2).into_event();
65+
let reaction = f
66+
.reaction(event_id!("$first"), "👍")
67+
.event_id(event_id!("$reaction"))
68+
.server_ts(3)
69+
.into_event();
70+
let edit = f
71+
.text_msg("* second, edited")
72+
.edit(
73+
event_id!("$second"),
74+
RoomMessageEventContentWithoutRelation::text_plain("second, edited"),
75+
)
76+
.event_id(event_id!("$edit"))
77+
.server_ts(4)
78+
.into_event();
79+
80+
let server = MatrixMockServer::new().await;
81+
server.mock_room_event().match_event_id().ok(first).mount().await;
82+
server.mock_room_event().match_event_id().ok(second).mount().await;
83+
server
84+
.mock_room_relations()
85+
.match_target_event(event_id!("$first").to_owned())
86+
.ok(RoomRelationsResponseTemplate::default()
87+
.events(vec![reaction.raw().clone().cast_unchecked()]))
88+
.mount()
89+
.await;
90+
server
91+
.mock_room_relations()
92+
.match_target_event(event_id!("$second").to_owned())
93+
.ok(RoomRelationsResponseTemplate::default()
94+
.events(vec![edit.raw().clone().cast_unchecked()]))
95+
.mount()
96+
.await;
97+
98+
let client = subscribed_client(&server).await;
99+
100+
// The IDs are given in any order, with a duplicate; the cache loads each
101+
// once, with its relations, sorted chronologically.
102+
let (cache, _drop_handles) = client
103+
.event_cache()
104+
.specific_events(
105+
room_id(),
106+
vec![
107+
event_id!("$second").to_owned(),
108+
event_id!("$first").to_owned(),
109+
event_id!("$first").to_owned(),
110+
],
111+
)
112+
.await
113+
.unwrap();
114+
115+
let (events, _subscriber) = cache.subscribe().await.unwrap();
116+
assert_eq!(
117+
event_ids(&events.into()),
118+
[event_id!("$first"), event_id!("$second"), event_id!("$reaction"), event_id!("$edit")]
119+
);
120+
assert_eq!(cache.events().await.unwrap().len(), 4);
121+
}
122+
123+
#[async_test]
124+
async fn test_specific_events_skip_events_that_cannot_be_loaded() {
125+
let f = EventFactory::new().room(room_id()).sender(*ALICE);
126+
127+
let server = MatrixMockServer::new().await;
128+
server
129+
.mock_room_event()
130+
.match_event_id()
131+
.ok(f.text_msg("first").event_id(event_id!("$first")).into_event())
132+
.mount()
133+
.await;
134+
135+
let client = subscribed_client(&server).await;
136+
137+
// `$missing` isn't mocked: its request fails, the rest is loaded.
138+
let (cache, _drop_handles) = client
139+
.event_cache()
140+
.specific_events(
141+
room_id(),
142+
vec![event_id!("$first").to_owned(), event_id!("$missing").to_owned()],
143+
)
144+
.await
145+
.unwrap();
146+
assert_eq!(event_ids(&cache.events().await.unwrap().into()), [event_id!("$first")]);
147+
148+
// When nothing can be loaded, that's an error.
149+
let result = client
150+
.event_cache()
151+
.specific_events(room_id(), vec![event_id!("$missing").to_owned()])
152+
.await;
153+
assert!(matches!(result, Err(EventCacheError::UnableToLoadSpecificEvents)));
154+
}
155+
156+
#[async_test]
157+
async fn test_specific_events_follow_sync_updates() {
158+
let f = EventFactory::new().room(room_id()).sender(*ALICE);
159+
160+
let server = MatrixMockServer::new().await;
161+
server
162+
.mock_room_event()
163+
.match_event_id()
164+
.ok(f.text_msg("target").event_id(event_id!("$target")).server_ts(1).into_event())
165+
.mount()
166+
.await;
167+
168+
let client = subscribed_client(&server).await;
169+
170+
let (cache, _drop_handles) = client
171+
.event_cache()
172+
.specific_events(room_id(), vec![event_id!("$target").to_owned()])
173+
.await
174+
.unwrap();
175+
let (events, mut subscriber) = cache.subscribe().await.unwrap();
176+
let mut events: Vector<Event> = events.into();
177+
assert_eq!(event_ids(&events), [event_id!("$target")]);
178+
179+
// A reaction and an edit of the target arrive from sync.
180+
server
181+
.sync_room(
182+
&client,
183+
JoinedRoomBuilder::new(room_id())
184+
.add_timeline_event(
185+
f.reaction(event_id!("$target"), "👍").event_id(event_id!("$reaction")),
186+
)
187+
.add_timeline_event(
188+
f.text_msg("* edited")
189+
.edit(
190+
event_id!("$target"),
191+
RoomMessageEventContentWithoutRelation::text_plain("edited"),
192+
)
193+
.event_id(event_id!("$edit")),
194+
),
195+
)
196+
.await;
197+
apply_next_updates(&mut subscriber, &mut events).await;
198+
assert_eq!(
199+
event_ids(&events),
200+
[event_id!("$target"), event_id!("$reaction"), event_id!("$edit")]
201+
);
202+
203+
// An unrelated message doesn't reach the cache.
204+
server
205+
.sync_room(
206+
&client,
207+
JoinedRoomBuilder::new(room_id())
208+
.add_timeline_event(f.text_msg("noise").event_id(event_id!("$noise"))),
209+
)
210+
.await;
211+
assert_no_update(&mut subscriber).await;
212+
213+
// Redacting the reaction replaces it with its redacted form.
214+
server
215+
.sync_room(
216+
&client,
217+
JoinedRoomBuilder::new(room_id()).add_timeline_event(
218+
f.redaction(event_id!("$reaction")).event_id(event_id!("$redaction_of_reaction")),
219+
),
220+
)
221+
.await;
222+
apply_next_updates(&mut subscriber, &mut events).await;
223+
assert!(find(&events, event_id!("$reaction")).raw().deserialize().unwrap().is_redacted());
224+
assert!(!find(&events, event_id!("$target")).raw().deserialize().unwrap().is_redacted());
225+
226+
// So does redacting the target itself.
227+
server
228+
.sync_room(
229+
&client,
230+
JoinedRoomBuilder::new(room_id()).add_timeline_event(
231+
f.redaction(event_id!("$target")).event_id(event_id!("$redaction_of_target")),
232+
),
233+
)
234+
.await;
235+
apply_next_updates(&mut subscriber, &mut events).await;
236+
assert!(find(&events, event_id!("$target")).raw().deserialize().unwrap().is_redacted());
237+
238+
assert!(!event_ids(&events).contains(&event_id!("$noise").to_owned()));
239+
}
240+
241+
#[async_test]
242+
async fn test_specific_events_set_event_ids_reloads_when_the_set_changes() {
243+
let f = EventFactory::new().room(room_id()).sender(*ALICE);
244+
245+
let server = MatrixMockServer::new().await;
246+
server
247+
.mock_room_event()
248+
.match_event_id()
249+
.ok(f.text_msg("first").event_id(event_id!("$first")).server_ts(1).into_event())
250+
.mount()
251+
.await;
252+
server
253+
.mock_room_event()
254+
.match_event_id()
255+
.ok(f.text_msg("second").event_id(event_id!("$second")).server_ts(2).into_event())
256+
.mount()
257+
.await;
258+
259+
let client = subscribed_client(&server).await;
260+
261+
let (cache, _drop_handles) = client
262+
.event_cache()
263+
.specific_events(room_id(), vec![event_id!("$first").to_owned()])
264+
.await
265+
.unwrap();
266+
let (events, mut subscriber) = cache.subscribe().await.unwrap();
267+
let mut events: Vector<Event> = events.into();
268+
assert_eq!(event_ids(&events), [event_id!("$first")]);
269+
270+
// The same set, in any order, is a no-op.
271+
cache.set_event_ids(vec![event_id!("$first").to_owned()]).await.unwrap();
272+
assert_no_update(&mut subscriber).await;
273+
274+
// Growing the set loads the new event.
275+
cache
276+
.set_event_ids(vec![event_id!("$second").to_owned(), event_id!("$first").to_owned()])
277+
.await
278+
.unwrap();
279+
apply_next_updates(&mut subscriber, &mut events).await;
280+
assert_eq!(event_ids(&events), [event_id!("$first"), event_id!("$second")]);
281+
282+
// Shrinking it drops the removed event.
283+
cache.set_event_ids(vec![event_id!("$second").to_owned()]).await.unwrap();
284+
apply_next_updates(&mut subscriber, &mut events).await;
285+
assert_eq!(event_ids(&events), [event_id!("$second")]);
286+
287+
// A set that can't be loaded at all is an error, and the previous set is kept.
288+
let result = cache.set_event_ids(vec![event_id!("$missing").to_owned()]).await;
289+
assert!(matches!(result, Err(EventCacheError::UnableToLoadSpecificEvents)));
290+
assert_no_update(&mut subscriber).await;
291+
assert_eq!(event_ids(&cache.events().await.unwrap().into()), [event_id!("$second")]);
292+
cache.set_event_ids(vec![event_id!("$second").to_owned()]).await.unwrap();
293+
assert_no_update(&mut subscriber).await;
294+
}
295+
296+
#[async_test]
297+
async fn test_specific_events_require_the_event_cache_to_be_subscribed() {
298+
let server = MatrixMockServer::new().await;
299+
let client = server.client_builder().build().await;
300+
server.sync_joined_room(&client, room_id()).await;
301+
302+
let result = client.event_cache().specific_events(room_id(), vec![]).await;
303+
assert!(matches!(result, Err(EventCacheError::NotSubscribedYet)));
304+
}

0 commit comments

Comments
 (0)