Skip to main content

matrix_sdk_base/store/
integration_tests.rs

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