1use std::{
4 collections::{BTreeMap, BTreeSet, HashMap},
5 str::FromStr,
6};
7
8use assert_matches::assert_matches;
9use growable_bloom_filter::GrowableBloomBuilder;
10use matrix_sdk_common::ttl::TtlValue;
11use matrix_sdk_test::{TestResult, event_factory::EventFactory};
12use ruma::{
13 EventId, MilliSecondsSinceUnixEpoch, OwnedUserId, RoomId, TransactionId, UserId,
14 api::{
15 FeatureFlag, MatrixVersion,
16 client::{discovery::discover_homeserver::HomeserverInfo, rtc::RtcTransport},
17 },
18 event_id,
19 events::{
20 AnyGlobalAccountDataEvent, AnyMessageLikeEventContent, AnyRoomAccountDataEvent,
21 AnyStrippedStateEvent, AnySyncStateEvent, GlobalAccountDataEventType,
22 RoomAccountDataEventType, StateEventType, SyncStateEvent,
23 presence::PresenceEvent,
24 receipt::{ReceiptThread, ReceiptType},
25 room::{
26 member::{
27 MembershipState, RoomMemberEventContent, StrippedRoomMemberEvent,
28 SyncRoomMemberEvent,
29 },
30 message::RoomMessageEventContent,
31 power_levels::RoomPowerLevelsEventContent,
32 redaction::SyncRoomRedactionEvent,
33 topic::RoomTopicEventContent,
34 },
35 tag::{TagInfo, TagName, Tags, UserTagName},
36 },
37 mxc_uri, owned_event_id, owned_mxc_uri,
38 presence::PresenceState,
39 profile::{ProfileFieldName, UserProfileChanges, UserProfileUpdate},
40 push::Ruleset,
41 room_id,
42 room_version_rules::AuthorizationRules,
43 serde::Raw,
44 uint, user_id,
45};
46use serde_json::json;
47use strass::assert_let;
48
49#[cfg(feature = "unstable-msc4354")]
50use super::QueuedRequestKind;
51use super::{
52 DependentQueuedRequestKind, DisplayName, DynStateStore, RoomLoadSettings,
53 SupportedVersionsResponse, WellKnownResponse, send_queue::SentRequestKey,
54};
55use crate::{
56 RoomInfo, RoomMemberships, RoomState, StateChanges, StateStoreDataKey, StateStoreDataValue,
57 deserialized_responses::MemberEvent,
58 store::{
59 ChildTransactionId, QueueWedgeError, SerializableEventContent, StateStoreExt,
60 StoredThreadSubscription, ThreadSubscriptionStatus,
61 },
62 utils::RawStateEventWithKeys,
63};
64
65#[allow(async_fn_in_trait)]
70pub trait StateStoreIntegrationTests {
71 async fn populate(&self) -> TestResult;
73 async fn test_topic_redaction(&self) -> TestResult;
75 async fn test_populate_store(&self) -> TestResult;
77 async fn test_member_saving(&self) -> TestResult;
79 async fn test_filter_saving(&self) -> TestResult;
81 async fn test_user_avatar_url_saving(&self) -> TestResult;
83 async fn test_sync_token_saving(&self) -> TestResult;
85 async fn test_utd_hook_manager_data_saving(&self) -> TestResult;
87 async fn test_one_time_key_already_uploaded_data_saving(&self) -> TestResult;
89 async fn test_stripped_member_saving(&self) -> TestResult;
91 async fn test_power_level_saving(&self) -> TestResult;
93 async fn test_receipts_saving(&self) -> TestResult;
95 async fn test_custom_storage(&self) -> TestResult;
97 async fn test_stripped_non_stripped(&self) -> TestResult;
99 async fn test_stripped_state_is_kept_across_saves(&self) -> TestResult;
101 async fn test_room_info_only_save_keeps_room_data(&self) -> TestResult;
103 async fn test_room_removal(&self) -> TestResult;
105 async fn test_profile_removal(&self) -> TestResult;
107 async fn test_presence_saving(&self) -> TestResult;
109 async fn test_display_names_saving(&self) -> TestResult;
111 async fn test_send_queue(&self) -> TestResult;
113 async fn test_send_queue_priority(&self) -> TestResult;
115 async fn test_send_queue_dependents(&self) -> TestResult;
117 async fn test_send_queue_sticky_events(&self) -> TestResult;
119 async fn test_update_send_queue_dependent(&self) -> TestResult;
121 async fn test_supported_versions_saving(&self) -> TestResult;
123 async fn test_well_known_saving(&self) -> TestResult;
125 async fn test_get_room_infos(&self) -> TestResult;
127 async fn test_thread_subscriptions(&self) -> TestResult;
129 async fn test_thread_subscriptions_bulk_upsert(&self) -> TestResult;
131 async fn test_global_profiles_saving(&self) -> TestResult;
133 async fn test_global_profiles_bulk_loading(&self) -> TestResult;
135}
136
137impl StateStoreIntegrationTests for DynStateStore {
138 async fn populate(&self) -> TestResult {
139 let mut changes = StateChanges::default();
140
141 let user_id = user_id();
142 let invited_user_id = invited_user_id();
143 let room_id = room_id();
144 let stripped_room_id = stripped_room_id();
145
146 changes.sync_token = Some("t392-516_47314_0_7_1_1_1_11444_1".to_owned());
147
148 let f = EventFactory::new().sender(user_id);
149 let presence_raw: Raw<PresenceEvent> = f
150 .presence(PresenceState::Online)
151 .avatar_url(mxc_uri!("mxc://localhost/wefuiwegh8742w"))
152 .currently_active(false)
153 .last_active_ago(1)
154 .status_msg("Making cupcakes")
155 .into();
156 let presence_event = presence_raw.deserialize()?;
157 changes.add_presence_event(presence_event, presence_raw);
158
159 let f = EventFactory::new().sender(user_id);
160 let pushrules_raw: Raw<AnyGlobalAccountDataEvent> =
161 f.push_rules(Ruleset::server_default(user_id)).into();
162 let pushrules_event = pushrules_raw.deserialize()?;
163 changes.account_data.insert(pushrules_event.event_type(), pushrules_raw);
164
165 let mut room = RoomInfo::new(room_id, RoomState::Joined);
166 room.mark_as_left();
167
168 let f = EventFactory::new().sender(user_id).room(room_id);
169 let mut tags = Tags::new();
170 tags.insert(TagName::Favorite, TagInfo::new());
171 tags.insert(TagName::User(UserTagName::from_str("u.work").unwrap()), TagInfo::new());
172 let tag_raw: Raw<AnyRoomAccountDataEvent> = f.tag(tags).into();
173 let tag_event = tag_raw.deserialize()?;
174 changes.add_room_account_data(room_id, tag_event, tag_raw);
175
176 let f = EventFactory::new().sender(user_id).room(room_id);
177 let name_raw: Raw<AnySyncStateEvent> = f.room_name("room name").into();
178 let name_event = name_raw.deserialize()?;
179 room.handle_state_event(
180 &mut RawStateEventWithKeys::try_from_raw_state_event(name_raw.clone())
181 .expect("generated state event should be valid"),
182 );
183 changes.add_state_event(room_id, name_event, name_raw);
184
185 let f = EventFactory::new().sender(user_id);
186 let receipt_content = f
187 .room(room_id)
188 .read_receipts()
189 .add(event_id!("$example"), user_id, ReceiptType::Read, ReceiptThread::Unthreaded)
190 .into_content();
191 changes.add_receipts(room_id, receipt_content);
192
193 let topic_event_id = topic_event_id();
194 let topic_raw: Raw<AnySyncStateEvent> = EventFactory::new()
195 .room(room_id)
196 .sender(user_id)
197 .room_topic("😀")
198 .event_id(topic_event_id)
199 .prev_content(RoomTopicEventContent::new("test".to_owned()))
200 .into_raw_sync_state();
201 let topic_event = topic_raw.deserialize()?;
202 room.handle_state_event(
203 &mut RawStateEventWithKeys::try_from_raw_state_event(topic_raw.clone())
204 .expect("generated state event should be valid"),
205 );
206 changes.add_state_event(room_id, topic_event, topic_raw);
207
208 let mut room_ambiguity_map = HashMap::new();
209 let mut room_profiles = BTreeMap::new();
210
211 let f = EventFactory::new().sender(user_id).room(room_id);
212 let member_raw: Raw<SyncRoomMemberEvent> =
213 f.member(user_id).display_name("example").previous(MembershipState::Invite).into();
214 let member_event: SyncRoomMemberEvent = member_raw.deserialize()?;
215 let displayname = DisplayName::new(
216 member_event.as_original().unwrap().content.displayname.as_ref().unwrap(),
217 );
218 room_ambiguity_map.insert(displayname.clone(), BTreeSet::from([user_id.to_owned()]));
219 room_profiles.insert(user_id.to_owned(), (&member_event).into());
220
221 let member_raw: Raw<AnySyncStateEvent> = member_raw.cast_unchecked();
222 let member_state_event = member_raw.deserialize()?;
223 changes.add_state_event(room_id, member_state_event, member_raw);
224
225 let f = EventFactory::new().sender(user_id).room(room_id);
226 let invited_member_raw: Raw<SyncRoomMemberEvent> = f
227 .member(invited_user_id)
228 .invited(invited_user_id)
229 .display_name("example")
230 .avatar_url(mxc_uri!("mxc://localhost/SEsfnsuifSDFSSEF"))
231 .reason("Looking for support")
232 .into();
233 let invited_member_event: SyncRoomMemberEvent = invited_member_raw.deserialize()?;
235 room_ambiguity_map.entry(displayname).or_default().insert(invited_user_id.to_owned());
236 room_profiles.insert(invited_user_id.to_owned(), (&invited_member_event).into());
237
238 let invited_member_raw: Raw<AnySyncStateEvent> = invited_member_raw.cast_unchecked();
239 let invited_member_state_event = invited_member_raw.deserialize()?;
240 changes.add_state_event(room_id, invited_member_state_event, invited_member_raw);
241
242 changes.ambiguity_maps.insert(room_id.to_owned(), room_ambiguity_map);
243 changes.profiles.insert(room_id.to_owned(), room_profiles);
244 changes.add_room(room);
245
246 let mut stripped_room = RoomInfo::new(stripped_room_id, RoomState::Invited);
247
248 let f = EventFactory::new().sender(user_id).room(stripped_room_id);
249 let stripped_name_raw: Raw<AnyStrippedStateEvent> = f.room_name("room name").into();
250 let mut stripped_name_event =
251 RawStateEventWithKeys::try_from_raw_state_event(stripped_name_raw).unwrap();
252 stripped_room.handle_stripped_state_event(&mut stripped_name_event);
253 changes.stripped_state.insert(
254 stripped_room_id.to_owned(),
255 BTreeMap::from([(
256 stripped_name_event.event_type,
257 BTreeMap::from([(stripped_name_event.state_key, stripped_name_event.raw)]),
258 )]),
259 );
260
261 changes.add_room(stripped_room);
262
263 let f = EventFactory::new().sender(user_id).room(stripped_room_id);
264 let stripped_member_raw: Raw<StrippedRoomMemberEvent> =
265 f.member(user_id).display_name("example").into();
266 changes.add_stripped_member(stripped_room_id, user_id, stripped_member_raw);
267
268 self.save_changes(&changes).await?;
269
270 Ok(())
271 }
272
273 async fn test_topic_redaction(&self) -> TestResult {
274 let room_id = room_id();
275 let f = EventFactory::new();
276 self.populate().await?;
277
278 assert!(self.get_kv_data(StateStoreDataKey::SyncToken).await?.is_some());
279 assert_eq!(
280 self.get_state_event_static::<RoomTopicEventContent>(room_id)
281 .await?
282 .expect("room topic found before redaction")
283 .deserialize()
284 .expect("can deserialize room topic before redaction")
285 .as_sync()
286 .expect("room topic is a sync state event")
287 .as_original()
288 .expect("room topic is not redacted yet")
289 .content
290 .topic,
291 "😀"
292 );
293
294 let mut changes = StateChanges::default();
295
296 let topic_event_id = topic_event_id();
297 let redaction_evt: Raw<SyncRoomRedactionEvent> =
298 f.room(room_id).sender(user_id()).redaction(topic_event_id).into_raw();
299
300 changes.add_redaction(room_id, topic_event_id, redaction_evt);
301 self.save_changes(&changes).await?;
302
303 let redacted_event = self
304 .get_state_event_static::<RoomTopicEventContent>(room_id)
305 .await?
306 .expect("room topic found after redaction")
307 .deserialize()
308 .expect("can deserialize room topic after redaction");
309
310 assert_matches!(redacted_event.as_sync(), Some(SyncStateEvent::Redacted(_)));
311
312 Ok(())
313 }
314
315 async fn test_populate_store(&self) -> TestResult {
316 let room_id = room_id();
317 let user_id = user_id();
318 let display_name = DisplayName::new("example");
319
320 self.populate().await?;
321
322 assert!(self.get_kv_data(StateStoreDataKey::SyncToken).await?.is_some());
323 assert!(self.get_presence_event(user_id).await?.is_some());
324 assert_eq!(
325 self.get_room_infos(&RoomLoadSettings::default()).await?.len(),
326 2,
327 "Expected to find 2 room infos"
328 );
329 assert!(
330 self.get_account_data_event(GlobalAccountDataEventType::PushRules).await?.is_some()
331 );
332
333 assert!(self.get_state_event(room_id, StateEventType::RoomName, "").await?.is_some());
334 assert_eq!(
335 self.get_state_events(room_id, StateEventType::RoomTopic).await?.len(),
336 1,
337 "Expected to find 1 room topic"
338 );
339 assert!(self.get_profile(room_id, user_id).await?.is_some());
340 assert!(self.get_member_event(room_id, user_id).await?.is_some());
341 assert_eq!(
342 self.get_user_ids(room_id, RoomMemberships::empty()).await?.len(),
343 2,
344 "Expected to find 2 members for room"
345 );
346 assert_eq!(
347 self.get_user_ids(room_id, RoomMemberships::INVITE).await?.len(),
348 1,
349 "Expected to find 1 invited user ids"
350 );
351 assert_eq!(
352 self.get_user_ids(room_id, RoomMemberships::JOIN).await?.len(),
353 1,
354 "Expected to find 1 joined user ids"
355 );
356 assert_eq!(
357 self.get_users_with_display_name(room_id, &display_name).await?.len(),
358 2,
359 "Expected to find 2 display names for room"
360 );
361 assert!(
362 self.get_room_account_data_event(room_id, RoomAccountDataEventType::Tag)
363 .await?
364 .is_some()
365 );
366 assert!(
367 self.get_user_room_receipt_event(
368 room_id,
369 ReceiptType::Read,
370 &ReceiptThread::Unthreaded,
371 user_id
372 )
373 .await?
374 .is_some()
375 );
376 assert_eq!(
377 self.get_event_room_receipt_events(
378 room_id,
379 ReceiptType::Read,
380 &ReceiptThread::Unthreaded,
381 first_receipt_event_id()
382 )
383 .await?
384 .len(),
385 1,
386 "Expected to find 1 read receipt"
387 );
388 Ok(())
389 }
390
391 async fn test_member_saving(&self) -> TestResult {
392 let room_id = room_id!("!test_member_saving:localhost");
393 let user_id = user_id();
394 let second_user_id = user_id!("@second:localhost");
395 let third_user_id = user_id!("@third:localhost");
396 let unknown_user_id = user_id!("@unknown:localhost");
397
398 let mut user_ids = vec![user_id.to_owned()];
400 assert!(self.get_member_event(room_id, user_id).await?.is_none());
401 let member_events = self
402 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(room_id, &user_ids)
403 .await;
404 assert!(member_events?.is_empty());
405 assert!(self.get_profile(room_id, user_id).await?.is_none());
406 let profiles = self.get_profiles(room_id, &user_ids).await;
407 assert!(profiles?.is_empty());
408
409 let mut changes = StateChanges::default();
411 let raw_member_event = membership_event();
412 let profile = raw_member_event.deserialize()?.into();
413 changes
414 .state
415 .entry(room_id.to_owned())
416 .or_default()
417 .entry(StateEventType::RoomMember)
418 .or_default()
419 .insert(user_id.into(), raw_member_event.cast());
420 changes.profiles.entry(room_id.to_owned()).or_default().insert(user_id.to_owned(), profile);
421 self.save_changes(&changes).await?;
422
423 assert!(self.get_member_event(room_id, user_id).await?.is_some());
424 let member_events = self
425 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(room_id, &user_ids)
426 .await;
427 assert_eq!(member_events?.len(), 1);
428 let members = self.get_user_ids(room_id, RoomMemberships::empty()).await?;
429 assert_eq!(members.len(), 1, "We expected to find members for the room");
430 assert!(self.get_profile(room_id, user_id).await?.is_some());
431 let profiles = self.get_profiles(room_id, &user_ids).await;
432 assert_eq!(profiles?.len(), 1);
433
434 let mut changes = StateChanges::default();
436 let changes_members = changes
437 .state
438 .entry(room_id.to_owned())
439 .or_default()
440 .entry(StateEventType::RoomMember)
441 .or_default();
442 let changes_profiles = changes.profiles.entry(room_id.to_owned()).or_default();
443 let raw_second_member_event =
444 custom_membership_event(second_user_id, event_id!("$second_member_event"));
445 let second_profile = raw_second_member_event.deserialize()?.into();
446 changes_members.insert(second_user_id.into(), raw_second_member_event.cast());
447 changes_profiles.insert(second_user_id.to_owned(), second_profile);
448 let raw_third_member_event =
449 custom_membership_event(third_user_id, event_id!("$third_member_event"));
450 let third_profile = raw_third_member_event.deserialize()?.into();
451 changes_members.insert(third_user_id.into(), raw_third_member_event.cast());
452 changes_profiles.insert(third_user_id.to_owned(), third_profile);
453 self.save_changes(&changes).await?;
454
455 user_ids.extend([second_user_id.to_owned(), third_user_id.to_owned()]);
456 assert!(self.get_member_event(room_id, second_user_id).await?.is_some());
457 assert!(self.get_member_event(room_id, third_user_id).await?.is_some());
458 let member_events = self
459 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(room_id, &user_ids)
460 .await;
461 assert_eq!(member_events?.len(), 3);
462 let members = self.get_user_ids(room_id, RoomMemberships::empty()).await?;
463 assert_eq!(members.len(), 3, "We expected to find members for the room");
464 assert!(self.get_profile(room_id, second_user_id).await?.is_some());
465 assert!(self.get_profile(room_id, third_user_id).await?.is_some());
466 let profiles = self.get_profiles(room_id, &user_ids).await;
467 assert_eq!(profiles?.len(), 3);
468
469 user_ids.push(unknown_user_id.to_owned());
471 let member_events = self
472 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(room_id, &user_ids)
473 .await;
474 assert_eq!(member_events?.len(), 3);
475 let profiles = self.get_profiles(room_id, &user_ids).await;
476 assert_eq!(profiles?.len(), 3);
477
478 let member_events = self
480 .get_state_events_for_keys_static::<RoomMemberEventContent, OwnedUserId, _>(
481 room_id,
482 &[],
483 )
484 .await;
485 assert!(member_events?.is_empty());
486 let profiles = self.get_profiles(room_id, &[]).await;
487 assert!(profiles?.is_empty());
488
489 Ok(())
490 }
491
492 async fn test_filter_saving(&self) -> TestResult {
493 let filter_name = "filter_name";
494 let filter_id = "filter_id_1234";
495
496 self.set_kv_data(
497 StateStoreDataKey::Filter(filter_name),
498 StateStoreDataValue::Filter(filter_id.to_owned()),
499 )
500 .await?;
501 assert_let!(
502 Ok(Some(StateStoreDataValue::Filter(stored_filter_id))) =
503 self.get_kv_data(StateStoreDataKey::Filter(filter_name)).await
504 );
505 assert_eq!(stored_filter_id, filter_id);
506
507 self.remove_kv_data(StateStoreDataKey::Filter(filter_name)).await?;
508 assert_matches!(self.get_kv_data(StateStoreDataKey::Filter(filter_name)).await, Ok(None));
509
510 Ok(())
511 }
512
513 async fn test_user_avatar_url_saving(&self) -> TestResult {
514 let user_id = user_id!("@alice:example.org");
515 let url = owned_mxc_uri!("mxc://example.org/poiuyt098");
516
517 self.set_kv_data(
518 StateStoreDataKey::UserAvatarUrl(user_id),
519 StateStoreDataValue::UserAvatarUrl(url.clone()),
520 )
521 .await?;
522
523 assert_let!(
524 Ok(Some(StateStoreDataValue::UserAvatarUrl(stored_url))) =
525 self.get_kv_data(StateStoreDataKey::UserAvatarUrl(user_id)).await
526 );
527 assert_eq!(stored_url, url);
528
529 self.remove_kv_data(StateStoreDataKey::UserAvatarUrl(user_id)).await?;
530 assert_matches!(
531 self.get_kv_data(StateStoreDataKey::UserAvatarUrl(user_id)).await,
532 Ok(None)
533 );
534
535 Ok(())
536 }
537
538 async fn test_supported_versions_saving(&self) -> TestResult {
539 let versions =
540 BTreeSet::from([MatrixVersion::V1_1, MatrixVersion::V1_2, MatrixVersion::V1_11]);
541 let supported_versions = SupportedVersionsResponse {
542 versions: versions.iter().map(|version| version.as_str().unwrap().to_owned()).collect(),
543 unstable_features: [("org.matrix.experimental".to_owned(), true)].into(),
544 };
545
546 self.set_kv_data(
547 StateStoreDataKey::SupportedVersions,
548 StateStoreDataValue::SupportedVersions(TtlValue::new(supported_versions.clone())),
549 )
550 .await?;
551
552 assert_let!(
553 Ok(Some(StateStoreDataValue::SupportedVersions(stored_supported_versions))) =
554 self.get_kv_data(StateStoreDataKey::SupportedVersions).await
555 );
556 let stored_supported_versions = stored_supported_versions.into_data();
557 assert_eq!(supported_versions, stored_supported_versions);
558
559 let stored_supported = stored_supported_versions.supported_versions();
560 assert_eq!(stored_supported.versions, versions);
561 assert_eq!(stored_supported.features.len(), 1);
562 assert!(stored_supported.features.contains(&FeatureFlag::from("org.matrix.experimental")));
563
564 self.remove_kv_data(StateStoreDataKey::SupportedVersions).await?;
565 assert_matches!(self.get_kv_data(StateStoreDataKey::SupportedVersions).await, Ok(None));
566
567 Ok(())
568 }
569
570 async fn test_well_known_saving(&self) -> TestResult {
571 let well_known = WellKnownResponse {
572 homeserver: HomeserverInfo::new("matrix.example.com".to_owned()),
573 identity_server: None,
574 tile_server: None,
575 rtc_foci: vec![RtcTransport::livekit("livekit.example.com".to_owned())],
576 };
577
578 self.set_kv_data(
579 StateStoreDataKey::WellKnown,
580 StateStoreDataValue::WellKnown(TtlValue::new(Some(well_known.clone()))),
581 )
582 .await?;
583
584 assert_let!(
585 Ok(Some(StateStoreDataValue::WellKnown(stored_well_known))) =
586 self.get_kv_data(StateStoreDataKey::WellKnown).await
587 );
588 let stored_well_known = stored_well_known.into_data();
589 assert_eq!(stored_well_known, Some(well_known));
590
591 self.remove_kv_data(StateStoreDataKey::WellKnown).await?;
592 assert_matches!(self.get_kv_data(StateStoreDataKey::WellKnown).await, Ok(None));
593
594 self.set_kv_data(
595 StateStoreDataKey::WellKnown,
596 StateStoreDataValue::WellKnown(TtlValue::new(None)),
597 )
598 .await?;
599
600 assert_let!(
601 Ok(Some(StateStoreDataValue::WellKnown(stored_well_known))) =
602 self.get_kv_data(StateStoreDataKey::WellKnown).await
603 );
604 let stored_well_known = stored_well_known.into_data();
605 assert_eq!(stored_well_known, None);
606
607 Ok(())
608 }
609
610 async fn test_sync_token_saving(&self) -> TestResult {
611 let sync_token_1 = "t392-516_47314_0_7_1";
612 let sync_token_2 = "t392-516_47314_0_7_2";
613
614 assert_matches!(self.get_kv_data(StateStoreDataKey::SyncToken).await, Ok(None));
615
616 let changes =
617 StateChanges { sync_token: Some(sync_token_1.to_owned()), ..Default::default() };
618 self.save_changes(&changes).await?;
619 assert_let!(
620 Ok(Some(StateStoreDataValue::SyncToken(stored_sync_token))) =
621 self.get_kv_data(StateStoreDataKey::SyncToken).await
622 );
623 assert_eq!(stored_sync_token, sync_token_1);
624
625 self.set_kv_data(
626 StateStoreDataKey::SyncToken,
627 StateStoreDataValue::SyncToken(sync_token_2.to_owned()),
628 )
629 .await?;
630 assert_let!(
631 Ok(Some(StateStoreDataValue::SyncToken(stored_sync_token))) =
632 self.get_kv_data(StateStoreDataKey::SyncToken).await
633 );
634 assert_eq!(stored_sync_token, sync_token_2);
635
636 self.remove_kv_data(StateStoreDataKey::SyncToken).await?;
637 assert_matches!(self.get_kv_data(StateStoreDataKey::SyncToken).await, Ok(None));
638
639 Ok(())
640 }
641
642 async fn test_utd_hook_manager_data_saving(&self) -> TestResult {
643 assert!(
645 self.get_kv_data(StateStoreDataKey::UtdHookManagerData)
646 .await
647 .expect("Could not read data")
648 .is_none(),
649 "Store was not empty at start"
650 );
651
652 let data = GrowableBloomBuilder::new().build();
654 self.set_kv_data(
655 StateStoreDataKey::UtdHookManagerData,
656 StateStoreDataValue::UtdHookManagerData(data.clone()),
657 )
658 .await
659 .expect("Could not save data");
660
661 let read_data = self
663 .get_kv_data(StateStoreDataKey::UtdHookManagerData)
664 .await
665 .expect("Could not read data")
666 .expect("no data found")
667 .into_utd_hook_manager_data()
668 .expect("not UtdHookManagerData");
669
670 assert_eq!(read_data, data);
671
672 Ok(())
673 }
674
675 async fn test_one_time_key_already_uploaded_data_saving(&self) -> TestResult {
676 assert!(
678 self.get_kv_data(StateStoreDataKey::OneTimeKeyAlreadyUploaded).await?.is_none(),
679 "Store was not empty at start"
680 );
681
682 self.set_kv_data(
683 StateStoreDataKey::OneTimeKeyAlreadyUploaded,
684 StateStoreDataValue::OneTimeKeyAlreadyUploaded,
685 )
686 .await?;
687
688 let data = self.get_kv_data(StateStoreDataKey::OneTimeKeyAlreadyUploaded).await?;
689 data.expect("The loaded data should be Some");
690
691 Ok(())
692 }
693
694 async fn test_stripped_member_saving(&self) -> TestResult {
695 let room_id = room_id!("!test_stripped_member_saving:localhost");
696 let user_id = user_id();
697 let second_user_id = user_id!("@second:localhost");
698 let third_user_id = user_id!("@third:localhost");
699 let unknown_user_id = user_id!("@unknown:localhost");
700
701 assert!(self.get_member_event(room_id, user_id).await?.is_none());
703 let member_events = self
704 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(
705 room_id,
706 &[user_id.to_owned()],
707 )
708 .await;
709 assert!(member_events?.is_empty());
710
711 let mut changes = StateChanges::default();
713 changes
714 .stripped_state
715 .entry(room_id.to_owned())
716 .or_default()
717 .entry(StateEventType::RoomMember)
718 .or_default()
719 .insert(user_id.into(), stripped_membership_event().cast());
720 self.save_changes(&changes).await?;
721
722 assert!(self.get_member_event(room_id, user_id).await?.is_some());
723 let member_events = self
724 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(
725 room_id,
726 &[user_id.to_owned()],
727 )
728 .await;
729 assert_eq!(member_events?.len(), 1);
730 let members = self.get_user_ids(room_id, RoomMemberships::empty()).await?;
731 assert_eq!(members.len(), 1, "We expected to find members for the room");
732
733 let mut changes = StateChanges::default();
735 let changes_members = changes
736 .stripped_state
737 .entry(room_id.to_owned())
738 .or_default()
739 .entry(StateEventType::RoomMember)
740 .or_default();
741 changes_members
742 .insert(second_user_id.into(), custom_stripped_membership_event(second_user_id).cast());
743 changes_members
744 .insert(third_user_id.into(), custom_stripped_membership_event(third_user_id).cast());
745 self.save_changes(&changes).await?;
746
747 assert!(self.get_member_event(room_id, second_user_id).await?.is_some());
748 assert!(self.get_member_event(room_id, third_user_id).await?.is_some());
749 let member_events = self
750 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(
751 room_id,
752 &[user_id.to_owned(), second_user_id.to_owned(), third_user_id.to_owned()],
753 )
754 .await;
755 assert_eq!(member_events?.len(), 3);
756 let members = self.get_user_ids(room_id, RoomMemberships::empty()).await?;
757 assert_eq!(members.len(), 3, "We expected to find members for the room");
758
759 let member_events = self
761 .get_state_events_for_keys_static::<RoomMemberEventContent, _, _>(
762 room_id,
763 &[
764 user_id.to_owned(),
765 second_user_id.to_owned(),
766 third_user_id.to_owned(),
767 unknown_user_id.to_owned(),
768 ],
769 )
770 .await;
771 assert_eq!(member_events?.len(), 3);
772
773 let member_events = self
775 .get_state_events_for_keys_static::<RoomMemberEventContent, OwnedUserId, _>(
776 room_id,
777 &[],
778 )
779 .await;
780 assert!(member_events?.is_empty());
781
782 Ok(())
783 }
784
785 async fn test_power_level_saving(&self) -> TestResult {
786 let room_id = room_id!("!test_power_level_saving:localhost");
787
788 let raw_event = power_level_event();
789 let event = raw_event.deserialize()?;
790
791 assert!(
792 self.get_state_event(room_id, StateEventType::RoomPowerLevels, "").await?.is_none()
793 );
794 let mut changes = StateChanges::default();
795 changes.add_state_event(room_id, event, raw_event);
796
797 self.save_changes(&changes).await?;
798 assert!(
799 self.get_state_event(room_id, StateEventType::RoomPowerLevels, "").await?.is_some()
800 );
801
802 Ok(())
803 }
804
805 async fn test_receipts_saving(&self) -> TestResult {
806 let room_id = room_id!("!test_receipts_saving:localhost");
807
808 let first_event_id = event_id!("$1435641916114394fHBLK:matrix.org");
809 let second_event_id = event_id!("$fHBLK1435641916114394:matrix.org");
810
811 let first_receipt_ts = uint!(1436451550);
812 let second_receipt_ts = uint!(1436451653);
813 let third_receipt_ts = uint!(1436474532);
814
815 let first_receipt_event = serde_json::from_value(json!({
816 first_event_id: {
817 "m.read": {
818 user_id(): {
819 "ts": first_receipt_ts,
820 }
821 }
822 }
823 }))?;
824
825 let second_receipt_event = serde_json::from_value(json!({
826 second_event_id: {
827 "m.read": {
828 user_id(): {
829 "ts": second_receipt_ts,
830 }
831 }
832 }
833 }))?;
834
835 let third_receipt_event = serde_json::from_value(json!({
836 second_event_id: {
837 "m.read": {
838 user_id(): {
839 "ts": third_receipt_ts,
840 "thread_id": "main",
841 }
842 }
843 }
844 }))?;
845
846 assert!(
847 self.get_user_room_receipt_event(
848 room_id,
849 ReceiptType::Read,
850 &ReceiptThread::Unthreaded,
851 user_id()
852 )
853 .await
854 .expect("failed to read unthreaded user room receipt")
855 .is_none()
856 );
857 assert!(
858 self.get_event_room_receipt_events(
859 room_id,
860 ReceiptType::Read,
861 &ReceiptThread::Unthreaded,
862 first_event_id
863 )
864 .await
865 .expect("failed to read unthreaded event room receipt for 1")
866 .is_empty()
867 );
868 assert!(
869 self.get_event_room_receipt_events(
870 room_id,
871 ReceiptType::Read,
872 &ReceiptThread::Unthreaded,
873 second_event_id
874 )
875 .await
876 .expect("failed to read unthreaded event room receipt for 2")
877 .is_empty()
878 );
879
880 let mut changes = StateChanges::default();
881 changes.add_receipts(room_id, first_receipt_event);
882
883 self.save_changes(&changes).await?;
884 let (unthreaded_user_receipt_event_id, unthreaded_user_receipt) = self
885 .get_user_room_receipt_event(
886 room_id,
887 ReceiptType::Read,
888 &ReceiptThread::Unthreaded,
889 user_id(),
890 )
891 .await
892 .expect("failed to read unthreaded user room receipt after save")
893 .unwrap();
894 assert_eq!(unthreaded_user_receipt_event_id, first_event_id);
895 assert_eq!(unthreaded_user_receipt.ts.unwrap().0, first_receipt_ts);
896 let first_event_unthreaded_receipts = self
897 .get_event_room_receipt_events(
898 room_id,
899 ReceiptType::Read,
900 &ReceiptThread::Unthreaded,
901 first_event_id,
902 )
903 .await
904 .expect("failed to read unthreaded event room receipt for 1 after save");
905 assert_eq!(
906 first_event_unthreaded_receipts.len(),
907 1,
908 "Found a wrong number of unthreaded receipts for 1 after save"
909 );
910 assert_eq!(first_event_unthreaded_receipts[0].0, user_id());
911 assert_eq!(first_event_unthreaded_receipts[0].1.ts.unwrap().0, first_receipt_ts);
912 assert!(
913 self.get_event_room_receipt_events(
914 room_id,
915 ReceiptType::Read,
916 &ReceiptThread::Unthreaded,
917 second_event_id
918 )
919 .await
920 .expect("failed to read unthreaded event room receipt for 2 after save")
921 .is_empty()
922 );
923
924 let mut changes = StateChanges::default();
925 changes.add_receipts(room_id, second_receipt_event);
926
927 self.save_changes(&changes).await.expect("Saving works");
928 let (unthreaded_user_receipt_event_id, unthreaded_user_receipt) = self
929 .get_user_room_receipt_event(
930 room_id,
931 ReceiptType::Read,
932 &ReceiptThread::Unthreaded,
933 user_id(),
934 )
935 .await
936 .expect("Getting unthreaded user room receipt after save failed")
937 .unwrap();
938 assert_eq!(unthreaded_user_receipt_event_id, second_event_id);
939 assert_eq!(unthreaded_user_receipt.ts.unwrap().0, second_receipt_ts);
940 assert!(
941 self.get_event_room_receipt_events(
942 room_id,
943 ReceiptType::Read,
944 &ReceiptThread::Unthreaded,
945 first_event_id
946 )
947 .await
948 .expect("Getting unthreaded event room receipt events for first event failed")
949 .is_empty()
950 );
951 let second_event_unthreaded_receipts = self
952 .get_event_room_receipt_events(
953 room_id,
954 ReceiptType::Read,
955 &ReceiptThread::Unthreaded,
956 second_event_id,
957 )
958 .await
959 .expect("Getting unthreaded event room receipt events for second event failed");
960 assert_eq!(
961 second_event_unthreaded_receipts.len(),
962 1,
963 "Found a wrong number of unthreaded receipts for second event after save"
964 );
965 assert_eq!(second_event_unthreaded_receipts[0].0, user_id());
966 assert_eq!(second_event_unthreaded_receipts[0].1.ts.unwrap().0, second_receipt_ts);
967
968 assert!(
969 self.get_user_room_receipt_event(
970 room_id,
971 ReceiptType::Read,
972 &ReceiptThread::Main,
973 user_id()
974 )
975 .await
976 .expect("failed to read threaded user room receipt")
977 .is_none()
978 );
979 assert!(
980 self.get_event_room_receipt_events(
981 room_id,
982 ReceiptType::Read,
983 &ReceiptThread::Main,
984 second_event_id
985 )
986 .await
987 .expect("Getting threaded event room receipts for 2 failed")
988 .is_empty()
989 );
990
991 let mut changes = StateChanges::default();
992 changes.add_receipts(room_id, third_receipt_event);
993
994 self.save_changes(&changes).await.expect("Saving works");
995 let (unthreaded_user_receipt_event_id, unthreaded_user_receipt) = self
997 .get_user_room_receipt_event(
998 room_id,
999 ReceiptType::Read,
1000 &ReceiptThread::Unthreaded,
1001 user_id(),
1002 )
1003 .await
1004 .expect("Getting unthreaded user room receipt after save failed")
1005 .unwrap();
1006 assert_eq!(unthreaded_user_receipt_event_id, second_event_id);
1007 assert_eq!(unthreaded_user_receipt.ts.unwrap().0, second_receipt_ts);
1008 let second_event_unthreaded_receipts = self
1009 .get_event_room_receipt_events(
1010 room_id,
1011 ReceiptType::Read,
1012 &ReceiptThread::Unthreaded,
1013 second_event_id,
1014 )
1015 .await
1016 .expect("Getting unthreaded event room receipt events for second event failed");
1017 assert_eq!(
1018 second_event_unthreaded_receipts.len(),
1019 1,
1020 "Found a wrong number of unthreaded receipts for second event after save"
1021 );
1022 assert_eq!(second_event_unthreaded_receipts[0].0, user_id());
1023 assert_eq!(second_event_unthreaded_receipts[0].1.ts.unwrap().0, second_receipt_ts);
1024 let (threaded_user_receipt_event_id, threaded_user_receipt) = self
1026 .get_user_room_receipt_event(
1027 room_id,
1028 ReceiptType::Read,
1029 &ReceiptThread::Main,
1030 user_id(),
1031 )
1032 .await
1033 .expect("Getting threaded user room receipt after save failed")
1034 .unwrap();
1035 assert_eq!(threaded_user_receipt_event_id, second_event_id);
1036 assert_eq!(threaded_user_receipt.ts.unwrap().0, third_receipt_ts);
1037 let second_event_threaded_receipts = self
1038 .get_event_room_receipt_events(
1039 room_id,
1040 ReceiptType::Read,
1041 &ReceiptThread::Main,
1042 second_event_id,
1043 )
1044 .await
1045 .expect("Getting threaded event room receipt events for second event failed");
1046 assert_eq!(
1047 second_event_threaded_receipts.len(),
1048 1,
1049 "Found a wrong number of threaded receipts for second event after save"
1050 );
1051 assert_eq!(second_event_threaded_receipts[0].0, user_id());
1052 assert_eq!(second_event_threaded_receipts[0].1.ts.unwrap().0, third_receipt_ts);
1053
1054 Ok(())
1055 }
1056
1057 async fn test_custom_storage(&self) -> TestResult {
1058 let key = "my_key";
1059 let value = &[0, 1, 2, 3];
1060
1061 self.set_custom_value(key.as_bytes(), value.to_vec()).await?;
1062
1063 let read = self.get_custom_value(key.as_bytes()).await?;
1064
1065 assert_eq!(Some(value.as_ref()), read.as_deref());
1066
1067 Ok(())
1068 }
1069
1070 async fn test_stripped_non_stripped(&self) -> TestResult {
1071 let room_id = room_id!("!test_stripped_non_stripped:localhost");
1072 let user_id = user_id();
1073
1074 assert!(self.get_member_event(room_id, user_id).await?.is_none());
1075 assert_eq!(self.get_room_infos(&RoomLoadSettings::default()).await?.len(), 0);
1076
1077 let mut changes = StateChanges::default();
1078 changes
1079 .state
1080 .entry(room_id.to_owned())
1081 .or_default()
1082 .entry(StateEventType::RoomMember)
1083 .or_default()
1084 .insert(user_id.into(), membership_event().cast());
1085 changes.add_room(RoomInfo::new(room_id, RoomState::Left));
1086 self.save_changes(&changes).await?;
1087
1088 let member_event = self.get_member_event(room_id, user_id).await?.unwrap().deserialize()?;
1089 assert!(matches!(member_event, MemberEvent::Sync(_)));
1090 assert_eq!(self.get_room_infos(&RoomLoadSettings::default()).await?.len(), 1);
1091
1092 let members = self.get_user_ids(room_id, RoomMemberships::empty()).await?;
1093 assert_eq!(members, vec![user_id.to_owned()]);
1094
1095 let mut changes = StateChanges::default();
1096 changes.add_stripped_member(room_id, user_id, custom_stripped_membership_event(user_id));
1097 changes.add_room(RoomInfo::new(room_id, RoomState::Invited));
1098 self.save_changes(&changes).await?;
1099
1100 let member_event = self.get_member_event(room_id, user_id).await?.unwrap().deserialize()?;
1101 assert!(matches!(member_event, MemberEvent::Stripped(_)));
1102 assert_eq!(self.get_room_infos(&RoomLoadSettings::default()).await?.len(), 1);
1103
1104 let members = self.get_user_ids(room_id, RoomMemberships::empty()).await?;
1105 assert_eq!(members, vec![user_id.to_owned()]);
1106
1107 Ok(())
1108 }
1109
1110 async fn test_stripped_state_is_kept_across_saves(&self) -> TestResult {
1111 for (room_id, room_state) in [
1114 (room_id!("!test_stripped_state_invited:localhost"), RoomState::Invited),
1115 (room_id!("!test_stripped_state_knocked:localhost"), RoomState::Knocked),
1116 ] {
1117 let user_id = user_id();
1118
1119 let mut changes = StateChanges::default();
1120 changes.add_stripped_member(
1121 room_id,
1122 user_id,
1123 custom_stripped_membership_event(user_id),
1124 );
1125
1126 let f = EventFactory::new().sender(user_id).room(room_id);
1127 let stripped_name_raw: Raw<AnyStrippedStateEvent> = f.room_name("room name").into();
1128 let stripped_name_event =
1129 RawStateEventWithKeys::try_from_raw_state_event(stripped_name_raw).unwrap();
1130 changes
1131 .stripped_state
1132 .entry(room_id.to_owned())
1133 .or_default()
1134 .entry(stripped_name_event.event_type)
1135 .or_default()
1136 .insert(stripped_name_event.state_key, stripped_name_event.raw);
1137
1138 changes.add_room(RoomInfo::new(room_id, room_state));
1139 self.save_changes(&changes).await?;
1140
1141 assert!(
1142 self.get_state_event(room_id, StateEventType::RoomName, "").await?.is_some(),
1143 "{room_state:?}: the stripped room name must be saved"
1144 );
1145 let member_event =
1146 self.get_member_event(room_id, user_id).await?.unwrap().deserialize()?;
1147 assert_matches!(member_event, MemberEvent::Stripped(_));
1148
1149 let mut changes = StateChanges::default();
1151 changes.add_room(RoomInfo::new(room_id, room_state));
1152 self.save_changes(&changes).await?;
1153
1154 assert!(
1155 self.get_state_event(room_id, StateEventType::RoomName, "").await?.is_some(),
1156 "{room_state:?}: the stripped room name must survive a room info save"
1157 );
1158 assert!(
1159 self.get_member_event(room_id, user_id).await?.is_some(),
1160 "{room_state:?}: the stripped member must survive a room info save"
1161 );
1162 assert_eq!(
1163 self.get_user_ids(room_id, RoomMemberships::empty()).await?,
1164 vec![user_id.to_owned()],
1165 "{room_state:?}: the stripped member must still be listed"
1166 );
1167 }
1168
1169 Ok(())
1170 }
1171
1172 async fn test_room_info_only_save_keeps_room_data(&self) -> TestResult {
1173 for (room_id, new_state) in [
1177 (room_id!("!test_room_info_only_save_knocked:localhost"), RoomState::Knocked),
1178 (room_id!("!test_room_info_only_save_invited:localhost"), RoomState::Invited),
1179 ] {
1180 let user_id = user_id();
1181
1182 let mut changes = StateChanges::default();
1183 changes
1184 .state
1185 .entry(room_id.to_owned())
1186 .or_default()
1187 .entry(StateEventType::RoomMember)
1188 .or_default()
1189 .insert(user_id.into(), membership_event().cast());
1190
1191 let f = EventFactory::new().sender(user_id).room(room_id);
1192 let name_raw: Raw<AnySyncStateEvent> = f.room_name("room name").into();
1193 let name_event = name_raw.deserialize()?;
1194 changes.add_state_event(room_id, name_event, name_raw);
1195
1196 changes.add_room(RoomInfo::new(room_id, RoomState::Left));
1197 self.save_changes(&changes).await?;
1198
1199 let member_event =
1200 self.get_member_event(room_id, user_id).await?.unwrap().deserialize()?;
1201 assert_matches!(member_event, MemberEvent::Sync(_));
1202 assert!(self.get_state_event(room_id, StateEventType::RoomName, "").await?.is_some());
1203
1204 let mut changes = StateChanges::default();
1206 changes.add_room(RoomInfo::new(room_id, new_state));
1207 self.save_changes(&changes).await?;
1208
1209 let member_event = self
1210 .get_member_event(room_id, user_id)
1211 .await?
1212 .unwrap_or_else(|| {
1213 panic!("{new_state:?}: the member event must survive a room info only save")
1214 })
1215 .deserialize()?;
1216 assert_matches!(member_event, MemberEvent::Sync(_));
1217 assert!(
1218 self.get_state_event(room_id, StateEventType::RoomName, "").await?.is_some(),
1219 "{new_state:?}: the room name must survive a room info only save"
1220 );
1221 assert_eq!(
1222 self.get_user_ids(room_id, RoomMemberships::empty()).await?,
1223 vec![user_id.to_owned()],
1224 "{new_state:?}: the member must still be listed"
1225 );
1226 }
1227
1228 Ok(())
1229 }
1230
1231 async fn test_room_removal(&self) -> TestResult {
1232 let room_id = room_id();
1233 let user_id = user_id();
1234 let display_name = DisplayName::new("example");
1235 let stripped_room_id = stripped_room_id();
1236
1237 self.populate().await?;
1238
1239 {
1240 let txn = TransactionId::new();
1242 let ev =
1243 SerializableEventContent::new(&RoomMessageEventContent::text_plain("sup").into())?;
1244 self.save_send_queue_request(
1245 room_id,
1246 txn.clone(),
1247 MilliSecondsSinceUnixEpoch::now(),
1248 ev.into(),
1249 0,
1250 )
1251 .await?;
1252
1253 self.save_dependent_queued_request(
1255 room_id,
1256 &txn,
1257 ChildTransactionId::new(),
1258 MilliSecondsSinceUnixEpoch::now(),
1259 DependentQueuedRequestKind::RedactEvent,
1260 )
1261 .await?;
1262 }
1263
1264 self.remove_room(room_id).await?;
1265
1266 assert_eq!(
1267 self.get_room_infos(&RoomLoadSettings::default()).await?.len(),
1268 1,
1269 "room is still there"
1270 );
1271
1272 assert!(self.get_state_event(room_id, StateEventType::RoomName, "").await?.is_none());
1273 assert!(
1274 self.get_state_events(room_id, StateEventType::RoomTopic).await?.is_empty(),
1275 "still state events found"
1276 );
1277 assert!(self.get_profile(room_id, user_id).await?.is_none());
1278 assert!(self.get_member_event(room_id, user_id).await?.is_none());
1279 assert!(
1280 self.get_user_ids(room_id, RoomMemberships::empty()).await?.is_empty(),
1281 "still user ids found"
1282 );
1283 assert!(
1284 self.get_user_ids(room_id, RoomMemberships::INVITE).await?.is_empty(),
1285 "still invited user ids found"
1286 );
1287 assert!(
1288 self.get_user_ids(room_id, RoomMemberships::JOIN).await?.is_empty(),
1289 "still joined users found"
1290 );
1291 assert!(
1292 self.get_users_with_display_name(room_id, &display_name).await?.is_empty(),
1293 "still display names found"
1294 );
1295 assert!(
1296 self.get_room_account_data_event(room_id, RoomAccountDataEventType::Tag)
1297 .await?
1298 .is_none()
1299 );
1300 assert!(
1301 self.get_user_room_receipt_event(
1302 room_id,
1303 ReceiptType::Read,
1304 &ReceiptThread::Unthreaded,
1305 user_id
1306 )
1307 .await?
1308 .is_none()
1309 );
1310 assert!(
1311 self.get_event_room_receipt_events(
1312 room_id,
1313 ReceiptType::Read,
1314 &ReceiptThread::Unthreaded,
1315 first_receipt_event_id()
1316 )
1317 .await?
1318 .is_empty(),
1319 "still event receipts in the store"
1320 );
1321 assert!(self.load_send_queue_requests(room_id).await?.is_empty());
1322 assert!(self.load_dependent_queued_requests(room_id).await?.is_empty());
1323
1324 self.remove_room(stripped_room_id).await?;
1325
1326 assert!(
1327 self.get_room_infos(&RoomLoadSettings::default()).await?.is_empty(),
1328 "still room info found"
1329 );
1330 Ok(())
1331 }
1332
1333 async fn test_profile_removal(&self) -> TestResult {
1334 let room_id = room_id();
1335
1336 let user_id = user_id();
1338 let invited_user_id = invited_user_id();
1339
1340 self.populate().await?;
1341
1342 let new_invite_member_json = json!({
1343 "content": {
1344 "avatar_url": "mxc://localhost/SEsfnsuifSDFSSEG",
1345 "displayname": "example after update",
1346 "membership": "invite",
1347 "reason": "Looking for support"
1348 },
1349 "event_id": "$143273582443PhrSm:localhost",
1350 "origin_server_ts": 1432735824,
1351 "room_id": room_id,
1352 "sender": user_id,
1353 "state_key": invited_user_id,
1354 "type": "m.room.member",
1355 });
1356 let new_invite_member_event: SyncRoomMemberEvent =
1357 serde_json::from_value(new_invite_member_json.clone())?;
1358
1359 let mut changes = StateChanges {
1360 profiles_to_delete: [(
1362 room_id.to_owned(),
1363 vec![user_id.to_owned(), invited_user_id.to_owned()],
1364 )]
1365 .into(),
1366
1367 profiles: {
1369 let mut map = BTreeMap::default();
1370 map.insert(
1371 room_id.to_owned(),
1372 [(invited_user_id.to_owned(), new_invite_member_event.into())]
1373 .into_iter()
1374 .collect(),
1375 );
1376 map
1377 },
1378
1379 ..StateChanges::default()
1380 };
1381
1382 let raw = serde_json::from_value::<Raw<AnySyncStateEvent>>(new_invite_member_json)
1383 .expect("can create sync-state-event for topic");
1384 let event = raw.deserialize()?;
1385 changes.add_state_event(room_id, event, raw);
1386
1387 self.save_changes(&changes).await?;
1388
1389 assert!(self.get_profile(room_id, user_id).await?.is_none());
1391 assert!(self.get_member_event(room_id, user_id).await?.is_some());
1392
1393 let invited_member_event = self.get_profile(room_id, invited_user_id).await?.unwrap();
1395 assert_eq!(
1396 invited_member_event.content.displayname.as_deref(),
1397 Some("example after update")
1398 );
1399 assert!(self.get_member_event(room_id, invited_user_id).await?.is_some());
1400
1401 Ok(())
1402 }
1403
1404 async fn test_presence_saving(&self) -> TestResult {
1405 let user_id = user_id();
1406 let second_user_id = user_id!("@second:localhost");
1407 let third_user_id = user_id!("@third:localhost");
1408 let unknown_user_id = user_id!("@unknown:localhost");
1409
1410 let mut user_ids = vec![user_id.to_owned()];
1412 let presence_event = self.get_presence_event(user_id).await;
1413 assert!(presence_event?.is_none());
1414 let presence_events = self.get_presence_events(&user_ids).await;
1415 assert!(presence_events?.is_empty());
1416
1417 let mut changes = StateChanges::default();
1419 changes.presence.insert(user_id.to_owned(), custom_presence_event(user_id));
1420 self.save_changes(&changes).await?;
1421
1422 let presence_event = self.get_presence_event(user_id).await;
1423 assert!(presence_event?.is_some());
1424 let presence_events = self.get_presence_events(&user_ids).await;
1425 assert_eq!(presence_events?.len(), 1);
1426
1427 let mut changes = StateChanges::default();
1429 changes.presence.insert(second_user_id.to_owned(), custom_presence_event(second_user_id));
1430 changes.presence.insert(third_user_id.to_owned(), custom_presence_event(third_user_id));
1431 self.save_changes(&changes).await?;
1432
1433 user_ids.extend([second_user_id.to_owned(), third_user_id.to_owned()]);
1434 let presence_event = self.get_presence_event(second_user_id).await;
1435 assert!(presence_event?.is_some());
1436 let presence_event = self.get_presence_event(third_user_id).await;
1437 assert!(presence_event?.is_some());
1438 let presence_events = self.get_presence_events(&user_ids).await;
1439 assert_eq!(presence_events?.len(), 3);
1440
1441 user_ids.push(unknown_user_id.to_owned());
1443 let member_events = self.get_presence_events(&user_ids).await;
1444 assert_eq!(member_events?.len(), 3);
1445
1446 let presence_events = self.get_presence_events(&[]).await;
1448 assert!(presence_events?.is_empty());
1449
1450 Ok(())
1451 }
1452
1453 async fn test_display_names_saving(&self) -> TestResult {
1454 let room_id = room_id!("!test_display_names_saving:localhost");
1455 let user_id = user_id();
1456 let user_display_name = DisplayName::new("User");
1457 let second_user_id = user_id!("@second:localhost");
1458 let third_user_id = user_id!("@third:localhost");
1459 let other_display_name = DisplayName::new("Raoul");
1460 let unknown_display_name = DisplayName::new("Unknown");
1461
1462 let mut display_names = vec![user_display_name.to_owned()];
1464 let users = self.get_users_with_display_name(room_id, &user_display_name).await?;
1465 assert!(users.is_empty());
1466 let names = self.get_users_with_display_names(room_id, &display_names).await?;
1467 assert!(names.is_empty());
1468
1469 let mut changes = StateChanges::default();
1471 changes
1472 .ambiguity_maps
1473 .entry(room_id.to_owned())
1474 .or_default()
1475 .insert(user_display_name.to_owned(), [user_id.to_owned()].into());
1476 self.save_changes(&changes).await?;
1477
1478 let users = self.get_users_with_display_name(room_id, &user_display_name).await?;
1479 assert_eq!(users.len(), 1);
1480 let names = self.get_users_with_display_names(room_id, &display_names).await?;
1481 assert_eq!(names.len(), 1);
1482 assert_eq!(names.get(&user_display_name).unwrap().len(), 1);
1483
1484 let mut changes = StateChanges::default();
1486 changes.ambiguity_maps.entry(room_id.to_owned()).or_default().insert(
1487 other_display_name.to_owned(),
1488 [second_user_id.to_owned(), third_user_id.to_owned()].into(),
1489 );
1490 self.save_changes(&changes).await?;
1491
1492 display_names.push(other_display_name.to_owned());
1493 let users = self.get_users_with_display_name(room_id, &user_display_name).await?;
1494 assert_eq!(users.len(), 1);
1495 let users = self.get_users_with_display_name(room_id, &other_display_name).await?;
1496 assert_eq!(users.len(), 2);
1497 let names = self.get_users_with_display_names(room_id, &display_names).await?;
1498 assert_eq!(names.len(), 2);
1499 assert_eq!(names.get(&user_display_name).unwrap().len(), 1);
1500 assert_eq!(names.get(&other_display_name).unwrap().len(), 2);
1501
1502 display_names.push(unknown_display_name.to_owned());
1504 let names = self.get_users_with_display_names(room_id, &display_names).await?;
1505 assert_eq!(names.len(), 2);
1506
1507 let names = self.get_users_with_display_names(room_id, &[]).await?;
1509 assert!(names.is_empty());
1510
1511 Ok(())
1512 }
1513
1514 #[allow(clippy::needless_range_loop)]
1515 async fn test_send_queue(&self) -> TestResult {
1516 let room_id = room_id!("!test_send_queue:localhost");
1517
1518 let events = self.load_send_queue_requests(room_id).await?;
1520 assert!(events.is_empty());
1521
1522 let txn0 = TransactionId::new();
1524 let event0 =
1525 SerializableEventContent::new(&RoomMessageEventContent::text_plain("msg0").into())?;
1526 self.save_send_queue_request(
1527 room_id,
1528 txn0.clone(),
1529 MilliSecondsSinceUnixEpoch::now(),
1530 event0.into(),
1531 0,
1532 )
1533 .await?;
1534
1535 let pending = self.load_send_queue_requests(room_id).await?;
1537
1538 assert_eq!(pending.len(), 1);
1539 {
1540 assert_eq!(pending[0].transaction_id, txn0);
1541
1542 let deserialized = pending[0].as_event().unwrap().deserialize()?;
1543 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1544 assert_eq!(content.body(), "msg0");
1545
1546 assert!(!pending[0].is_wedged());
1547 }
1548
1549 for i in 1..=3 {
1551 let txn = TransactionId::new();
1552 let event = SerializableEventContent::new(
1553 &RoomMessageEventContent::text_plain(format!("msg{i}")).into(),
1554 )?;
1555
1556 self.save_send_queue_request(
1557 room_id,
1558 txn,
1559 MilliSecondsSinceUnixEpoch::now(),
1560 event.into(),
1561 0,
1562 )
1563 .await?;
1564 }
1565
1566 let pending = self.load_send_queue_requests(room_id).await?;
1568
1569 assert_eq!(pending.len(), 4);
1571
1572 assert_eq!(pending[0].transaction_id, txn0);
1573
1574 for i in 0..4 {
1575 let deserialized = pending[i].as_event().unwrap().deserialize()?;
1576 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1577 assert_eq!(content.body(), format!("msg{i}"));
1578 assert!(!pending[i].is_wedged());
1579 }
1580
1581 let txn2 = &pending[2].transaction_id;
1583 self.update_send_queue_request_status(
1584 room_id,
1585 txn2,
1586 Some(QueueWedgeError::GenericApiError { msg: "Oops".to_owned() }),
1587 )
1588 .await?;
1589
1590 let pending = self.load_send_queue_requests(room_id).await?;
1592
1593 assert_eq!(pending.len(), 4);
1595 assert_eq!(pending[0].transaction_id, txn0);
1596 assert_eq!(pending[2].transaction_id, *txn2);
1597 assert!(pending[2].is_wedged());
1598 let error = pending[2].clone().error.unwrap();
1599 let generic_error = assert_matches!(error, QueueWedgeError::GenericApiError { msg } => msg);
1600 assert_eq!(generic_error, "Oops");
1601 for i in 0..4 {
1602 if i != 2 {
1603 assert!(!pending[i].is_wedged());
1604 }
1605 }
1606
1607 let event0 = SerializableEventContent::new(
1609 &RoomMessageEventContent::text_plain("wow that's a cool test").into(),
1610 )?;
1611 self.update_send_queue_request(room_id, txn2, event0.into()).await?;
1612
1613 let pending = self.load_send_queue_requests(room_id).await?;
1615
1616 assert_eq!(pending.len(), 4);
1617 {
1618 assert_eq!(pending[2].transaction_id, *txn2);
1619
1620 let deserialized = pending[2].as_event().unwrap().deserialize()?;
1621 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1622 assert_eq!(content.body(), "wow that's a cool test");
1623
1624 assert!(!pending[2].is_wedged());
1625
1626 for i in 0..4 {
1627 if i != 2 {
1628 let deserialized = pending[i].as_event().unwrap().deserialize()?;
1629 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1630 assert_eq!(content.body(), format!("msg{i}"));
1631
1632 assert!(!pending[i].is_wedged());
1633 }
1634 }
1635 }
1636
1637 self.remove_send_queue_request(room_id, &txn0).await?;
1639
1640 let pending = self.load_send_queue_requests(room_id).await?;
1642
1643 assert_eq!(pending.len(), 3);
1644 assert_eq!(pending[1].transaction_id, *txn2);
1645 for i in 0..3 {
1646 assert_ne!(pending[i].transaction_id, txn0);
1647 }
1648
1649 let room_id2 = room_id!("!test_send_queue_two:localhost");
1654 {
1655 let txn = TransactionId::new();
1656 let event = SerializableEventContent::new(
1657 &RoomMessageEventContent::text_plain("room2").into(),
1658 )?;
1659 self.save_send_queue_request(
1660 room_id2,
1661 txn.clone(),
1662 MilliSecondsSinceUnixEpoch::now(),
1663 event.into(),
1664 0,
1665 )
1666 .await?;
1667 }
1668
1669 {
1671 let room_id3 = room_id!("!test_send_queue_three:localhost");
1672 let txn = TransactionId::new();
1673 let event = SerializableEventContent::new(
1674 &RoomMessageEventContent::text_plain("room3").into(),
1675 )?;
1676 self.save_send_queue_request(
1677 room_id3,
1678 txn.clone(),
1679 MilliSecondsSinceUnixEpoch::now(),
1680 event.into(),
1681 0,
1682 )
1683 .await?;
1684
1685 self.remove_send_queue_request(room_id3, &txn).await?;
1686 }
1687
1688 let outstanding_rooms = self.load_rooms_with_unsent_requests().await?;
1691 assert_eq!(outstanding_rooms.len(), 2);
1692 assert!(outstanding_rooms.iter().any(|room| room == room_id));
1693 assert!(outstanding_rooms.iter().any(|room| room == room_id2));
1694
1695 Ok(())
1696 }
1697
1698 async fn test_send_queue_priority(&self) -> TestResult {
1699 let room_id = room_id!("!test_send_queue:localhost");
1700
1701 let events = self.load_send_queue_requests(room_id).await?;
1703 assert!(events.is_empty());
1704
1705 let low0_txn = TransactionId::new();
1707 let ev0 =
1708 SerializableEventContent::new(&RoomMessageEventContent::text_plain("low0").into())?;
1709 self.save_send_queue_request(
1710 room_id,
1711 low0_txn.clone(),
1712 MilliSecondsSinceUnixEpoch::now(),
1713 ev0.into(),
1714 2,
1715 )
1716 .await?;
1717
1718 let high_txn = TransactionId::new();
1720 let ev1 =
1721 SerializableEventContent::new(&RoomMessageEventContent::text_plain("high").into())?;
1722 self.save_send_queue_request(
1723 room_id,
1724 high_txn.clone(),
1725 MilliSecondsSinceUnixEpoch::now(),
1726 ev1.into(),
1727 10,
1728 )
1729 .await?;
1730
1731 let low1_txn = TransactionId::new();
1733 let ev2 =
1734 SerializableEventContent::new(&RoomMessageEventContent::text_plain("low1").into())?;
1735 self.save_send_queue_request(
1736 room_id,
1737 low1_txn.clone(),
1738 MilliSecondsSinceUnixEpoch::now(),
1739 ev2.into(),
1740 2,
1741 )
1742 .await?;
1743
1744 let pending = self.load_send_queue_requests(room_id).await?;
1747
1748 assert_eq!(pending.len(), 3);
1749 {
1750 assert_eq!(pending[0].transaction_id, high_txn);
1751
1752 let deserialized = pending[0].as_event().unwrap().deserialize()?;
1753 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1754 assert_eq!(content.body(), "high");
1755 }
1756
1757 {
1758 assert_eq!(pending[1].transaction_id, low0_txn);
1759
1760 let deserialized = pending[1].as_event().unwrap().deserialize()?;
1761 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1762 assert_eq!(content.body(), "low0");
1763 }
1764
1765 {
1766 assert_eq!(pending[2].transaction_id, low1_txn);
1767
1768 let deserialized = pending[2].as_event().unwrap().deserialize()?;
1769 assert_let!(AnyMessageLikeEventContent::RoomMessage(content) = deserialized);
1770 assert_eq!(content.body(), "low1");
1771 }
1772
1773 Ok(())
1774 }
1775
1776 #[cfg(not(feature = "unstable-msc4354"))]
1777 async fn test_send_queue_sticky_events(&self) -> TestResult {
1778 Ok(()) }
1780
1781 #[cfg(feature = "unstable-msc4354")]
1782 async fn test_send_queue_sticky_events(&self) -> TestResult {
1783 use ruma::events::sticky::StickyDurationMs;
1784
1785 let room_id = room_id!("!test_send_queue_sticky_events:localhost");
1786
1787 let content =
1789 SerializableEventContent::new(&RoomMessageEventContent::text_plain("regular").into())?;
1790 let regular_txn = TransactionId::new();
1791 self.save_send_queue_request(
1792 room_id,
1793 regular_txn.clone(),
1794 MilliSecondsSinceUnixEpoch::now(),
1795 content.into(),
1796 0,
1797 )
1798 .await?;
1799
1800 let content =
1802 SerializableEventContent::new(&RoomMessageEventContent::text_plain("sticky").into())?;
1803 let sticky_txn = TransactionId::new();
1804 self.save_send_queue_request(
1805 room_id,
1806 sticky_txn.clone(),
1807 MilliSecondsSinceUnixEpoch::now(),
1808 QueuedRequestKind::Event {
1809 content,
1810 sticky_duration: Some(StickyDurationMs::new_clamped(300_000u32)),
1811 },
1812 0,
1813 )
1814 .await?;
1815
1816 let pending = self.load_send_queue_requests(room_id).await?;
1817 assert_eq!(pending.len(), 2);
1818
1819 assert_eq!(pending[0].transaction_id, regular_txn);
1820 assert_matches!(&pending[0].kind, QueuedRequestKind::Event { sticky_duration: None, .. });
1821
1822 assert_eq!(pending[1].transaction_id, sticky_txn);
1823 assert_matches!(
1824 &pending[1].kind,
1825 QueuedRequestKind::Event { sticky_duration: Some(duration), .. } => {
1826 assert_eq!(duration.get(), 300_000);
1827 }
1828 );
1829
1830 Ok(())
1831 }
1832
1833 async fn test_send_queue_dependents(&self) -> TestResult {
1834 let room_id = room_id!("!test_send_queue_dependents:localhost");
1835
1836 let txn0 = TransactionId::new();
1838 let event0 =
1839 SerializableEventContent::new(&RoomMessageEventContent::text_plain("hey").into())?;
1840 self.save_send_queue_request(
1841 room_id,
1842 txn0.clone(),
1843 MilliSecondsSinceUnixEpoch::now(),
1844 event0.clone().into(),
1845 0,
1846 )
1847 .await?;
1848
1849 assert!(self.load_dependent_queued_requests(room_id).await?.is_empty());
1851
1852 let child_txn = ChildTransactionId::new();
1854 self.save_dependent_queued_request(
1855 room_id,
1856 &txn0,
1857 child_txn.clone(),
1858 MilliSecondsSinceUnixEpoch::now(),
1859 DependentQueuedRequestKind::RedactEvent,
1860 )
1861 .await?;
1862
1863 let dependents = self.load_dependent_queued_requests(room_id).await?;
1865 assert_eq!(dependents.len(), 1);
1866 assert_eq!(dependents[0].parent_transaction_id, txn0);
1867 assert_eq!(dependents[0].own_transaction_id, child_txn);
1868 assert!(dependents[0].parent_key.is_none());
1869 assert_matches!(dependents[0].kind, DependentQueuedRequestKind::RedactEvent);
1870
1871 let (event, event_type) = event0.raw();
1873 let event_id = owned_event_id!("$1");
1874 let num_updated = self
1875 .mark_dependent_queued_requests_as_ready(
1876 room_id,
1877 &txn0,
1878 SentRequestKey::Event {
1879 event_id: event_id.clone(),
1880 event: event.clone(),
1881 event_type: event_type.to_owned(),
1882 },
1883 )
1884 .await?;
1885 assert_eq!(num_updated, 1);
1886
1887 let dependents = self.load_dependent_queued_requests(room_id).await?;
1889 assert_eq!(dependents.len(), 1);
1890 assert_eq!(dependents[0].parent_transaction_id, txn0);
1891 assert_eq!(dependents[0].own_transaction_id, child_txn);
1892 assert_matches!(
1893 dependents[0].parent_key.as_ref(),
1894 Some(SentRequestKey::Event {
1895 event_id: received_event_id,
1896 event: received_event,
1897 event_type: received_event_type
1898 }) => {
1899 assert_eq!(received_event_id, &event_id);
1900 assert_eq!(received_event.json().to_string(), event.json().to_string());
1901 assert_eq!(received_event_type.as_str(), event_type);
1902 }
1903 );
1904 assert_matches!(dependents[0].kind, DependentQueuedRequestKind::RedactEvent);
1905
1906 let removed = self
1908 .remove_dependent_queued_request(room_id, &dependents[0].own_transaction_id)
1909 .await?;
1910 assert!(removed);
1911
1912 assert!(self.load_dependent_queued_requests(room_id).await?.is_empty());
1914
1915 let txn1 = TransactionId::new();
1918 let event1 =
1919 SerializableEventContent::new(&RoomMessageEventContent::text_plain("hey2").into())?;
1920 self.save_send_queue_request(
1921 room_id,
1922 txn1.clone(),
1923 MilliSecondsSinceUnixEpoch::now(),
1924 event1.into(),
1925 0,
1926 )
1927 .await?;
1928
1929 self.save_dependent_queued_request(
1930 room_id,
1931 &txn0,
1932 ChildTransactionId::new(),
1933 MilliSecondsSinceUnixEpoch::now(),
1934 DependentQueuedRequestKind::RedactEvent,
1935 )
1936 .await?;
1937 assert_eq!(self.load_dependent_queued_requests(room_id).await?.len(), 1);
1938
1939 self.save_dependent_queued_request(
1940 room_id,
1941 &txn1,
1942 ChildTransactionId::new(),
1943 MilliSecondsSinceUnixEpoch::now(),
1944 DependentQueuedRequestKind::EditEvent {
1945 new_content: SerializableEventContent::new(
1946 &RoomMessageEventContent::text_plain("edit").into(),
1947 )?,
1948 },
1949 )
1950 .await?;
1951 assert_eq!(self.load_dependent_queued_requests(room_id).await?.len(), 2);
1952
1953 let removed = self.remove_send_queue_request(room_id, &txn0).await?;
1955 assert!(removed);
1956
1957 let dependents = self.load_dependent_queued_requests(room_id).await?;
1959 assert_eq!(dependents.len(), 2);
1960
1961 Ok(())
1962 }
1963
1964 async fn test_update_send_queue_dependent(&self) -> TestResult {
1965 let room_id = room_id!("!test_send_queue_dependents:localhost");
1966
1967 let txn = TransactionId::new();
1968
1969 let child_txn = ChildTransactionId::new();
1971
1972 self.save_dependent_queued_request(
1973 room_id,
1974 &txn,
1975 child_txn.clone(),
1976 MilliSecondsSinceUnixEpoch::now(),
1977 DependentQueuedRequestKind::RedactEvent,
1978 )
1979 .await?;
1980
1981 let dependents = self.load_dependent_queued_requests(room_id).await?;
1983 assert_eq!(dependents.len(), 1);
1984 assert_eq!(dependents[0].parent_transaction_id, txn);
1985 assert_eq!(dependents[0].own_transaction_id, child_txn);
1986 assert!(dependents[0].parent_key.is_none());
1987 assert_matches!(dependents[0].kind, DependentQueuedRequestKind::RedactEvent);
1988
1989 self.update_dependent_queued_request(
1991 room_id,
1992 &child_txn,
1993 DependentQueuedRequestKind::ReactEvent { key: "👍".to_owned() },
1994 )
1995 .await?;
1996
1997 let dependents = self.load_dependent_queued_requests(room_id).await?;
1999 assert_eq!(dependents.len(), 1);
2000 assert_eq!(dependents[0].parent_transaction_id, txn);
2001 assert_eq!(dependents[0].own_transaction_id, child_txn);
2002 assert!(dependents[0].parent_key.is_none());
2003 assert_matches!(
2004 &dependents[0].kind,
2005 DependentQueuedRequestKind::ReactEvent { key } => {
2006 assert_eq!(key, "👍");
2007 }
2008 );
2009
2010 Ok(())
2011 }
2012
2013 async fn test_get_room_infos(&self) -> TestResult {
2014 let room_id_0 = room_id!("!r0");
2015 let room_id_1 = room_id!("!r1");
2016 let room_id_2 = room_id!("!r2");
2017
2018 {
2020 assert_eq!(self.get_room_infos(&RoomLoadSettings::default()).await?.len(), 0);
2021 }
2022
2023 let mut changes = StateChanges::default();
2025 changes.add_room(RoomInfo::new(room_id_0, RoomState::Joined));
2026 changes.add_room(RoomInfo::new(room_id_1, RoomState::Joined));
2027 self.save_changes(&changes).await?;
2028
2029 {
2031 let mut all_rooms = self.get_room_infos(&RoomLoadSettings::All).await?;
2032
2033 all_rooms.sort_by(|a, b| a.room_id.cmp(&b.room_id));
2036
2037 assert_eq!(all_rooms.len(), 2);
2038 assert_eq!(all_rooms[0].room_id, room_id_0);
2039 assert_eq!(all_rooms[1].room_id, room_id_1);
2040 }
2041
2042 {
2044 let all_rooms =
2045 self.get_room_infos(&RoomLoadSettings::One(room_id_1.to_owned())).await?;
2046
2047 assert_eq!(all_rooms.len(), 1);
2048 assert_eq!(all_rooms[0].room_id, room_id_1);
2049 }
2050
2051 {
2054 let all_rooms =
2055 self.get_room_infos(&RoomLoadSettings::One(room_id_2.to_owned())).await?;
2056
2057 assert_eq!(all_rooms.len(), 0);
2058 }
2059
2060 Ok(())
2061 }
2062
2063 async fn test_thread_subscriptions(&self) -> TestResult {
2064 let first_thread = event_id!("$t1");
2065 let second_thread = event_id!("$t2");
2066
2067 let maybe_sub = self.load_thread_subscription(room_id(), first_thread).await?;
2069 assert!(maybe_sub.is_none());
2070
2071 let maybe_sub = self.load_thread_subscription(room_id(), second_thread).await?;
2072 assert!(maybe_sub.is_none());
2073
2074 self.upsert_thread_subscriptions(vec![(
2076 room_id(),
2077 first_thread,
2078 StoredThreadSubscription {
2079 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2080 bump_stamp: None,
2081 },
2082 )])
2083 .await?;
2084
2085 self.upsert_thread_subscriptions(vec![(
2086 room_id(),
2087 second_thread,
2088 StoredThreadSubscription {
2089 status: ThreadSubscriptionStatus::Subscribed { automatic: false },
2090 bump_stamp: None,
2091 },
2092 )])
2093 .await?;
2094
2095 let maybe_sub = self.load_thread_subscription(room_id(), first_thread).await?;
2097 assert_eq!(
2098 maybe_sub,
2099 Some(StoredThreadSubscription {
2100 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2101 bump_stamp: None,
2102 })
2103 );
2104
2105 let maybe_sub = self.load_thread_subscription(room_id(), second_thread).await?;
2106 assert_eq!(
2107 maybe_sub,
2108 Some(StoredThreadSubscription {
2109 status: ThreadSubscriptionStatus::Subscribed { automatic: false },
2110 bump_stamp: None,
2111 })
2112 );
2113
2114 self.upsert_thread_subscriptions(vec![(
2116 room_id(),
2117 first_thread,
2118 StoredThreadSubscription {
2119 status: ThreadSubscriptionStatus::Unsubscribed,
2120 bump_stamp: None,
2121 },
2122 )])
2123 .await?;
2124
2125 let maybe_sub = self.load_thread_subscription(room_id(), first_thread).await?;
2127 assert_eq!(
2128 maybe_sub,
2129 Some(StoredThreadSubscription {
2130 status: ThreadSubscriptionStatus::Unsubscribed,
2131 bump_stamp: None,
2132 })
2133 );
2134
2135 let maybe_sub = self.load_thread_subscription(room_id(), second_thread).await?;
2137 assert_eq!(
2138 maybe_sub,
2139 Some(StoredThreadSubscription {
2140 status: ThreadSubscriptionStatus::Subscribed { automatic: false },
2141 bump_stamp: None,
2142 })
2143 );
2144
2145 self.remove_thread_subscription(room_id(), second_thread).await?;
2147
2148 let maybe_sub = self.load_thread_subscription(room_id(), second_thread).await?;
2150 assert_eq!(maybe_sub, None);
2151
2152 let maybe_sub = self.load_thread_subscription(room_id(), first_thread).await?;
2154 assert_eq!(
2155 maybe_sub,
2156 Some(StoredThreadSubscription {
2157 status: ThreadSubscriptionStatus::Unsubscribed,
2158 bump_stamp: None,
2159 })
2160 );
2161
2162 self.remove_thread_subscription(room_id(), second_thread).await?;
2164
2165 Ok(())
2166 }
2167
2168 async fn test_thread_subscriptions_bulk_upsert(&self) -> TestResult {
2169 let threads = [
2170 event_id!("$t1"),
2171 event_id!("$t2"),
2172 event_id!("$t3"),
2173 event_id!("$t4"),
2174 event_id!("$t5"),
2175 event_id!("$t6"),
2176 ];
2177 let build_subscription_updates = |subs: &[StoredThreadSubscription]| {
2181 threads
2182 .iter()
2183 .zip(subs)
2184 .map(|(&event_id, &sub)| (room_id(), event_id, sub))
2185 .collect::<Vec<_>>()
2186 };
2187
2188 let initial_subscriptions = build_subscription_updates(&[
2190 StoredThreadSubscription {
2191 status: ThreadSubscriptionStatus::Unsubscribed,
2192 bump_stamp: None,
2193 },
2194 StoredThreadSubscription {
2195 status: ThreadSubscriptionStatus::Unsubscribed,
2196 bump_stamp: Some(14),
2197 },
2198 StoredThreadSubscription {
2199 status: ThreadSubscriptionStatus::Unsubscribed,
2200 bump_stamp: None,
2201 },
2202 StoredThreadSubscription {
2203 status: ThreadSubscriptionStatus::Unsubscribed,
2204 bump_stamp: Some(210),
2205 },
2206 StoredThreadSubscription {
2207 status: ThreadSubscriptionStatus::Unsubscribed,
2208 bump_stamp: Some(5),
2209 },
2210 StoredThreadSubscription {
2211 status: ThreadSubscriptionStatus::Unsubscribed,
2212 bump_stamp: Some(100),
2213 },
2214 ]);
2215
2216 let update_subscriptions = build_subscription_updates(&[
2217 StoredThreadSubscription {
2218 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2219 bump_stamp: None,
2220 },
2221 StoredThreadSubscription {
2222 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2223 bump_stamp: None,
2224 },
2225 StoredThreadSubscription {
2226 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2227 bump_stamp: Some(1101),
2228 },
2229 StoredThreadSubscription {
2230 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2231 bump_stamp: Some(222),
2232 },
2233 StoredThreadSubscription {
2234 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2235 bump_stamp: Some(1),
2236 },
2237 StoredThreadSubscription {
2238 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2239 bump_stamp: Some(100),
2240 },
2241 ]);
2242
2243 let expected_subscriptions = build_subscription_updates(&[
2244 StoredThreadSubscription {
2246 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2247 bump_stamp: None,
2248 },
2249 StoredThreadSubscription {
2251 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2252 bump_stamp: Some(14),
2253 },
2254 StoredThreadSubscription {
2256 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2257 bump_stamp: Some(1101),
2258 },
2259 StoredThreadSubscription {
2261 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2262 bump_stamp: Some(222),
2263 },
2264 StoredThreadSubscription {
2266 status: ThreadSubscriptionStatus::Unsubscribed,
2267 bump_stamp: Some(5),
2268 },
2269 StoredThreadSubscription {
2271 status: ThreadSubscriptionStatus::Unsubscribed,
2272 bump_stamp: Some(100),
2273 },
2274 ]);
2275
2276 self.upsert_thread_subscriptions(initial_subscriptions.clone()).await?;
2278
2279 for (room_id, thread_id, expected_sub) in &initial_subscriptions {
2281 let stored_subscription = self.load_thread_subscription(room_id, thread_id).await?;
2282 assert_eq!(stored_subscription, Some(*expected_sub));
2283 }
2284
2285 self.upsert_thread_subscriptions(update_subscriptions).await?;
2287
2288 for (room_id, thread_id, expected_sub) in &expected_subscriptions {
2290 let stored_subscription = self.load_thread_subscription(room_id, thread_id).await?;
2291 assert_eq!(stored_subscription, Some(*expected_sub));
2292 }
2293
2294 for (room_id, thread_id, _) in &expected_subscriptions {
2296 self.remove_thread_subscription(room_id, thread_id).await?;
2297 }
2298
2299 let initial_subscriptions = build_subscription_updates(&[
2300 StoredThreadSubscription {
2301 status: ThreadSubscriptionStatus::Unsubscribed,
2302 bump_stamp: Some(1),
2303 },
2304 StoredThreadSubscription {
2305 status: ThreadSubscriptionStatus::Subscribed { automatic: false },
2306 bump_stamp: Some(1),
2307 },
2308 StoredThreadSubscription {
2309 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2310 bump_stamp: Some(1),
2311 },
2312 ]);
2313
2314 self.upsert_thread_subscriptions(initial_subscriptions.clone()).await?;
2315
2316 for (room_id, thread_id, expected_sub) in &initial_subscriptions {
2317 let stored_subscription = self.load_thread_subscription(room_id, thread_id).await?;
2318 assert_eq!(stored_subscription, Some(*expected_sub));
2319 }
2320
2321 let update_subscriptions = build_subscription_updates(&[
2322 StoredThreadSubscription {
2323 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2324 bump_stamp: Some(2),
2325 },
2326 StoredThreadSubscription {
2327 status: ThreadSubscriptionStatus::Unsubscribed,
2328 bump_stamp: Some(2),
2329 },
2330 StoredThreadSubscription {
2331 status: ThreadSubscriptionStatus::Subscribed { automatic: false },
2332 bump_stamp: Some(2),
2333 },
2334 ]);
2335
2336 self.upsert_thread_subscriptions(update_subscriptions.clone()).await?;
2337
2338 for (room_id, thread_id, expected_sub) in &update_subscriptions {
2339 let stored_subscription = self.load_thread_subscription(room_id, thread_id).await?;
2340 assert_eq!(stored_subscription, Some(*expected_sub));
2341 }
2342
2343 Ok(())
2344 }
2345
2346 async fn test_global_profiles_saving(&self) -> TestResult {
2347 let user_id = user_id();
2348
2349 assert_matches!(self.get_global_profile(user_id).await, Ok(None));
2350
2351 let mut changes = UserProfileChanges::new();
2352 changes.updated.insert(ProfileFieldName::DisplayName, json!("Alice"));
2353 changes.removed.push(ProfileFieldName::AvatarUrl);
2354 let update = UserProfileUpdate::Updated(changes);
2355
2356 let mut changes = StateChanges::default();
2357 changes.global_profiles.insert(user_id.to_owned(), update);
2358 self.save_changes(&changes).await?;
2359
2360 let loaded = self.get_global_profile(user_id).await?.expect("Profile should be saved");
2361 let loaded_map: BTreeMap<String, serde_json::Value> = loaded.into_iter().collect();
2362 assert_eq!(loaded_map.get("displayname"), Some(&json!("Alice")));
2363 assert!(!loaded_map.contains_key("avatar_url"));
2364
2365 let mut changes = UserProfileChanges::new();
2366 changes.removed.push(ProfileFieldName::DisplayName);
2367 changes.updated.insert(ProfileFieldName::AvatarUrl, json!("mxc://example.com/avatar"));
2368 let update2 = UserProfileUpdate::Updated(changes);
2369
2370 let mut changes = StateChanges::default();
2371 changes.global_profiles.insert(user_id.to_owned(), update2);
2372 self.save_changes(&changes).await?;
2373
2374 let loaded2 = self.get_global_profile(user_id).await?.expect("Profile should exist");
2375 let loaded_map2: BTreeMap<String, serde_json::Value> = loaded2.into_iter().collect();
2376 assert!(!loaded_map2.contains_key("displayname"));
2377 assert_eq!(loaded_map2.get("avatar_url"), Some(&json!("mxc://example.com/avatar")));
2378
2379 Ok(())
2380 }
2381
2382 async fn test_global_profiles_bulk_loading(&self) -> TestResult {
2383 let alice = user_id!("@alice:localhost");
2384 let bob = user_id!("@bob:localhost");
2385 let unknown = user_id!("@unknown:localhost");
2386
2387 let mut changes = StateChanges::default();
2388 changes.global_profiles.insert(alice.to_owned(), {
2389 let mut profile_changes = UserProfileChanges::new();
2390 profile_changes.updated.insert(ProfileFieldName::DisplayName, json!("Alice"));
2391 UserProfileUpdate::Updated(profile_changes)
2392 });
2393 changes.global_profiles.insert(bob.to_owned(), {
2394 let mut profile_changes = UserProfileChanges::new();
2395 profile_changes.updated.insert(ProfileFieldName::DisplayName, json!("Bob"));
2396 UserProfileUpdate::Updated(profile_changes)
2397 });
2398 self.save_changes(&changes).await?;
2399
2400 let requested = [alice.to_owned(), bob.to_owned(), unknown.to_owned()];
2402 let profiles = self.get_global_profiles(&requested).await?;
2403
2404 assert_eq!(profiles.len(), 2);
2405 assert!(!profiles.contains_key(unknown));
2406
2407 let alice_map: BTreeMap<String, serde_json::Value> = profiles
2408 .get(alice)
2409 .expect("Alice's profile should be loaded")
2410 .clone()
2411 .into_iter()
2412 .collect();
2413 assert_eq!(alice_map.get("displayname"), Some(&json!("Alice")));
2414
2415 let bob_map: BTreeMap<String, serde_json::Value> = profiles
2416 .get(bob)
2417 .expect("Bob's profile should be loaded")
2418 .clone()
2419 .into_iter()
2420 .collect();
2421 assert_eq!(bob_map.get("displayname"), Some(&json!("Bob")));
2422
2423 Ok(())
2424 }
2425}
2426
2427#[allow(unused_macros, unused_extern_crates)]
2454#[macro_export]
2455macro_rules! statestore_integration_tests {
2456 () => {
2457 mod statestore_integration_tests {
2458 use matrix_sdk_test::{TestResult, async_test};
2459 use $crate::store::{IntoStateStore, StateStoreIntegrationTests};
2460
2461 use super::get_store;
2462
2463 #[async_test]
2464 async fn test_topic_redaction() -> TestResult {
2465 let store = get_store().await?.into_state_store();
2466 store.test_topic_redaction().await
2467 }
2468
2469 #[async_test]
2470 async fn test_populate_store() -> TestResult {
2471 let store = get_store().await?.into_state_store();
2472 store.test_populate_store().await
2473 }
2474
2475 #[async_test]
2476 async fn test_member_saving() -> TestResult {
2477 let store = get_store().await?.into_state_store();
2478 store.test_member_saving().await
2479 }
2480
2481 #[async_test]
2482 async fn test_filter_saving() -> TestResult {
2483 let store = get_store().await?.into_state_store();
2484 store.test_filter_saving().await
2485 }
2486
2487 #[async_test]
2488 async fn test_user_avatar_url_saving() -> TestResult {
2489 let store = get_store().await?.into_state_store();
2490 store.test_user_avatar_url_saving().await
2491 }
2492
2493 #[async_test]
2494 async fn test_supported_versions_saving() -> TestResult {
2495 let store = get_store().await?.into_state_store();
2496 store.test_supported_versions_saving().await
2497 }
2498
2499 #[async_test]
2500 async fn test_well_known_saving() -> TestResult {
2501 let store = get_store().await?.into_state_store();
2502 store.test_well_known_saving().await
2503 }
2504
2505 #[async_test]
2506 async fn test_sync_token_saving() -> TestResult {
2507 let store = get_store().await?.into_state_store();
2508 store.test_sync_token_saving().await
2509 }
2510
2511 #[async_test]
2512 async fn test_utd_hook_manager_data_saving() -> TestResult {
2513 let store = get_store().await?.into_state_store();
2514 store.test_utd_hook_manager_data_saving().await
2515 }
2516
2517 #[async_test]
2518 async fn test_one_time_key_already_uploaded_data_saving() -> TestResult {
2519 let store = get_store().await?.into_state_store();
2520 store.test_one_time_key_already_uploaded_data_saving().await
2521 }
2522
2523 #[async_test]
2524 async fn test_stripped_member_saving() -> TestResult {
2525 let store = get_store().await?.into_state_store();
2526 store.test_stripped_member_saving().await
2527 }
2528
2529 #[async_test]
2530 async fn test_power_level_saving() -> TestResult {
2531 let store = get_store().await?.into_state_store();
2532 store.test_power_level_saving().await
2533 }
2534
2535 #[async_test]
2536 async fn test_receipts_saving() -> TestResult {
2537 let store = get_store().await?.into_state_store();
2538 store.test_receipts_saving().await
2539 }
2540
2541 #[async_test]
2542 async fn test_custom_storage() -> TestResult {
2543 let store = get_store().await?.into_state_store();
2544 store.test_custom_storage().await
2545 }
2546
2547 #[async_test]
2548 async fn test_stripped_non_stripped() -> TestResult {
2549 let store = get_store().await?.into_state_store();
2550 store.test_stripped_non_stripped().await
2551 }
2552
2553 #[async_test]
2554 async fn test_stripped_state_is_kept_across_saves() -> TestResult {
2555 let store = get_store().await?.into_state_store();
2556 store.test_stripped_state_is_kept_across_saves().await
2557 }
2558
2559 #[async_test]
2560 async fn test_room_info_only_save_keeps_room_data() -> TestResult {
2561 let store = get_store().await?.into_state_store();
2562 store.test_room_info_only_save_keeps_room_data().await
2563 }
2564
2565 #[async_test]
2566 async fn test_room_removal() -> TestResult {
2567 let store = get_store().await?.into_state_store();
2568 store.test_room_removal().await
2569 }
2570
2571 #[async_test]
2572 async fn test_profile_removal() -> TestResult {
2573 let store = get_store().await?.into_state_store();
2574 store.test_profile_removal().await
2575 }
2576
2577 #[async_test]
2578 async fn test_presence_saving() -> TestResult {
2579 let store = get_store().await?.into_state_store();
2580 store.test_presence_saving().await
2581 }
2582
2583 #[async_test]
2584 async fn test_display_names_saving() -> TestResult {
2585 let store = get_store().await?.into_state_store();
2586 store.test_display_names_saving().await
2587 }
2588
2589 #[async_test]
2590 async fn test_send_queue() -> TestResult {
2591 let store = get_store().await?.into_state_store();
2592 store.test_send_queue().await
2593 }
2594
2595 #[async_test]
2596 async fn test_send_queue_priority() -> TestResult {
2597 let store = get_store().await?.into_state_store();
2598 store.test_send_queue_priority().await
2599 }
2600
2601 #[async_test]
2602 async fn test_send_queue_dependents() -> TestResult {
2603 let store = get_store().await?.into_state_store();
2604 store.test_send_queue_dependents().await
2605 }
2606
2607 #[async_test]
2608 async fn test_send_queue_sticky_events() -> TestResult {
2609 let store = get_store().await?.into_state_store();
2610 store.test_send_queue_sticky_events().await
2611 }
2612
2613 #[async_test]
2614 async fn test_update_send_queue_dependent() -> TestResult {
2615 let store = get_store().await?.into_state_store();
2616 store.test_update_send_queue_dependent().await
2617 }
2618
2619 #[async_test]
2620 async fn test_get_room_infos() -> TestResult {
2621 let store = get_store().await?.into_state_store();
2622 store.test_get_room_infos().await
2623 }
2624
2625 #[async_test]
2626 async fn test_thread_subscriptions() -> TestResult {
2627 let store = get_store().await?.into_state_store();
2628 store.test_thread_subscriptions().await
2629 }
2630
2631 #[async_test]
2632 async fn test_thread_subscriptions_bulk_upsert() -> TestResult {
2633 let store = get_store().await?.into_state_store();
2634 store.test_thread_subscriptions_bulk_upsert().await
2635 }
2636
2637 #[async_test]
2638 async fn test_global_profiles_saving() -> TestResult {
2639 let store = get_store().await?.into_state_store();
2640 store.test_global_profiles_saving().await
2641 }
2642
2643 #[async_test]
2644 async fn test_global_profiles_bulk_loading() -> TestResult {
2645 let store = get_store().await?.into_state_store();
2646 store.test_global_profiles_bulk_loading().await
2647 }
2648 }
2649 };
2650}
2651
2652fn user_id() -> &'static UserId {
2653 user_id!("@example:localhost")
2654}
2655
2656fn invited_user_id() -> &'static UserId {
2657 user_id!("@invited:localhost")
2658}
2659
2660fn room_id() -> &'static RoomId {
2661 room_id!("!test:localhost")
2662}
2663
2664fn stripped_room_id() -> &'static RoomId {
2665 room_id!("!stripped:localhost")
2666}
2667
2668fn first_receipt_event_id() -> &'static EventId {
2669 event_id!("$example")
2670}
2671
2672fn topic_event_id() -> &'static EventId {
2673 event_id!("$topic_event")
2674}
2675
2676fn power_level_event() -> Raw<AnySyncStateEvent> {
2677 let content = RoomPowerLevelsEventContent::new(&AuthorizationRules::V1);
2678
2679 let event = json!({
2680 "event_id": "$h29iv0s8:example.com",
2681 "content": content,
2682 "sender": user_id(),
2683 "type": "m.room.power_levels",
2684 "origin_server_ts": 0u64,
2685 "state_key": "",
2686 });
2687
2688 serde_json::from_value(event).unwrap()
2689}
2690
2691fn stripped_membership_event() -> Raw<StrippedRoomMemberEvent> {
2692 custom_stripped_membership_event(user_id())
2693}
2694
2695fn custom_stripped_membership_event(user_id: &UserId) -> Raw<StrippedRoomMemberEvent> {
2696 let ev_json = json!({
2697 "type": "m.room.member",
2698 "content": RoomMemberEventContent::new(MembershipState::Join),
2699 "sender": user_id,
2700 "state_key": user_id,
2701 });
2702
2703 Raw::new(&ev_json).unwrap().cast_unchecked()
2704}
2705
2706fn membership_event() -> Raw<SyncRoomMemberEvent> {
2707 custom_membership_event(user_id(), event_id!("$h29iv0s8:example.com"))
2708}
2709
2710fn custom_membership_event(user_id: &UserId, event_id: &EventId) -> Raw<SyncRoomMemberEvent> {
2711 let ev_json = json!({
2712 "type": "m.room.member",
2713 "content": RoomMemberEventContent::new(MembershipState::Join),
2714 "event_id": event_id,
2715 "origin_server_ts": 198,
2716 "sender": user_id,
2717 "state_key": user_id,
2718 });
2719
2720 Raw::new(&ev_json).unwrap().cast_unchecked()
2721}
2722
2723fn custom_presence_event(user_id: &UserId) -> Raw<PresenceEvent> {
2724 let ev_json = json!({
2725 "content": {
2726 "presence": "online"
2727 },
2728 "sender": user_id,
2729 });
2730
2731 Raw::new(&ev_json).unwrap().cast_unchecked()
2732}