1use 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#[allow(async_fn_in_trait)]
68pub trait StateStoreIntegrationTests {
69 async fn populate(&self) -> TestResult;
71 async fn test_topic_redaction(&self) -> TestResult;
73 async fn test_populate_store(&self) -> TestResult;
75 async fn test_member_saving(&self) -> TestResult;
77 async fn test_filter_saving(&self) -> TestResult;
79 async fn test_user_avatar_url_saving(&self) -> TestResult;
81 async fn test_sync_token_saving(&self) -> TestResult;
83 async fn test_utd_hook_manager_data_saving(&self) -> TestResult;
85 async fn test_one_time_key_already_uploaded_data_saving(&self) -> TestResult;
87 async fn test_stripped_member_saving(&self) -> TestResult;
89 async fn test_power_level_saving(&self) -> TestResult;
91 async fn test_receipts_saving(&self) -> TestResult;
93 async fn test_custom_storage(&self) -> TestResult;
95 async fn test_stripped_non_stripped(&self) -> TestResult;
97 async fn test_room_removal(&self) -> TestResult;
99 async fn test_profile_removal(&self) -> TestResult;
101 async fn test_presence_saving(&self) -> TestResult;
103 async fn test_display_names_saving(&self) -> TestResult;
105 async fn test_send_queue(&self) -> TestResult;
107 async fn test_send_queue_priority(&self) -> TestResult;
109 async fn test_send_queue_dependents(&self) -> TestResult;
111 async fn test_update_send_queue_dependent(&self) -> TestResult;
113 async fn test_supported_versions_saving(&self) -> TestResult;
115 async fn test_well_known_saving(&self) -> TestResult;
117 async fn test_get_room_infos(&self) -> TestResult;
119 async fn test_thread_subscriptions(&self) -> TestResult;
121 async fn test_thread_subscriptions_bulk_upsert(&self) -> TestResult;
123 async fn test_global_profiles_saving(&self) -> TestResult;
125 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 profiles_to_delete: [(
1233 room_id.to_owned(),
1234 vec![user_id.to_owned(), invited_user_id.to_owned()],
1235 )]
1236 .into(),
1237
1238 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 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 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 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 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 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 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 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 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 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 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 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 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 let events = self.load_send_queue_requests(room_id).await?;
1391 assert!(events.is_empty());
1392
1393 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 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 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 let pending = self.load_send_queue_requests(room_id).await?;
1439
1440 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 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 let pending = self.load_send_queue_requests(room_id).await?;
1463
1464 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 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 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 self.remove_send_queue_request(room_id, &txn0).await?;
1510
1511 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 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 {
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 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 let events = self.load_send_queue_requests(room_id).await?;
1574 assert!(events.is_empty());
1575
1576 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 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 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 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 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 assert!(self.load_dependent_queued_requests(room_id).await?.is_empty());
1665
1666 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 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 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 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 let removed = self
1722 .remove_dependent_queued_request(room_id, &dependents[0].own_transaction_id)
1723 .await?;
1724 assert!(removed);
1725
1726 assert!(self.load_dependent_queued_requests(room_id).await?.is_empty());
1728
1729 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 let removed = self.remove_send_queue_request(room_id, &txn0).await?;
1769 assert!(removed);
1770
1771 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 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 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 self.update_dependent_queued_request(
1805 room_id,
1806 &child_txn,
1807 DependentQueuedRequestKind::ReactEvent { key: "👍".to_owned() },
1808 )
1809 .await?;
1810
1811 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 {
1834 assert_eq!(self.get_room_infos(&RoomLoadSettings::default()).await?.len(), 0);
1835 }
1836
1837 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 {
1845 let mut all_rooms = self.get_room_infos(&RoomLoadSettings::All).await?;
1846
1847 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 {
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 {
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 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 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 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 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 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 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 self.remove_thread_subscription(room_id(), second_thread).await?;
1961
1962 let maybe_sub = self.load_thread_subscription(room_id(), second_thread).await?;
1964 assert_eq!(maybe_sub, None);
1965
1966 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 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 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 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 StoredThreadSubscription {
2059 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2060 bump_stamp: None,
2061 },
2062 StoredThreadSubscription {
2064 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2065 bump_stamp: Some(14),
2066 },
2067 StoredThreadSubscription {
2069 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2070 bump_stamp: Some(1101),
2071 },
2072 StoredThreadSubscription {
2074 status: ThreadSubscriptionStatus::Subscribed { automatic: true },
2075 bump_stamp: Some(222),
2076 },
2077 StoredThreadSubscription {
2079 status: ThreadSubscriptionStatus::Unsubscribed,
2080 bump_stamp: Some(5),
2081 },
2082 StoredThreadSubscription {
2084 status: ThreadSubscriptionStatus::Unsubscribed,
2085 bump_stamp: Some(100),
2086 },
2087 ]);
2088
2089 self.upsert_thread_subscriptions(initial_subscriptions.clone()).await?;
2091
2092 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 self.upsert_thread_subscriptions(update_subscriptions).await?;
2100
2101 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 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 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#[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}