1#![cfg(feature = "e2e-encryption")]
17
18use std::collections::BTreeMap;
19
20use async_stream::stream;
21use futures_core::Stream;
22use futures_util::{StreamExt, stream_select};
23use matrix_sdk_base::crypto::{
24 IdentityState, IdentityStatusChange, RoomIdentityChange, RoomIdentityState,
25};
26use ruma::{OwnedUserId, UserId, events::room::member::SyncRoomMemberEvent};
27use tokio::sync::mpsc;
28use tokio_stream::wrappers::ReceiverStream;
29
30use super::Room;
31use crate::{
32 Client, Error, Result,
33 encryption::identities::{IdentityUpdates, UserIdentity},
34 event_handler::EventHandlerDropGuard,
35};
36
37#[derive(Debug)]
49pub struct IdentityStatusChanges {
50 room_identity_state: RoomIdentityState<Room>,
52
53 _drop_guard: EventHandlerDropGuard,
56}
57
58impl IdentityStatusChanges {
59 pub async fn create_stream(
90 room: Room,
91 ) -> Result<impl Stream<Item = Vec<IdentityStatusChange>>> {
92 let identity_updates = wrap_identity_updates(&room.client).await?;
93 let (drop_guard, room_member_events) = wrap_room_member_events(&room);
94 let mut unprocessed_stream = combine_streams(identity_updates, room_member_events);
95 let own_user_id = room.client.user_id().ok_or(Error::InsufficientData)?.to_owned();
96
97 let state = IdentityStatusChanges {
98 room_identity_state: RoomIdentityState::new(room).await,
99 _drop_guard: drop_guard,
100 };
101
102 Ok(stream!({
103 let mut state = state;
106
107 let mut current_state =
108 filter_for_initial_update(state.room_identity_state.current_state(), &own_user_id);
109
110 if !current_state.is_empty() {
111 current_state.sort();
112 yield current_state;
113 }
114
115 while let Some(item) = unprocessed_stream.next().await {
116 let mut update = filter_non_self(
117 state.room_identity_state.process_change(item).await,
118 &own_user_id,
119 );
120 if !update.is_empty() {
121 update.sort();
122 yield update;
123 }
124 }
125 }))
126 }
127}
128
129fn filter_for_initial_update(
130 mut input: Vec<IdentityStatusChange>,
131 own_user_id: &UserId,
132) -> Vec<IdentityStatusChange> {
133 input.retain(|change| {
138 change.user_id != own_user_id && change.changed_to != IdentityState::Verified
139 });
140
141 input
142}
143
144fn filter_non_self(
145 mut input: Vec<IdentityStatusChange>,
146 own_user_id: &UserId,
147) -> Vec<IdentityStatusChange> {
148 input.retain(|change| change.user_id != own_user_id);
150 input
151}
152
153fn combine_streams(
154 identity_updates: impl Stream<Item = RoomIdentityChange> + Unpin,
155 room_member_events: impl Stream<Item = RoomIdentityChange> + Unpin,
156) -> impl Stream<Item = RoomIdentityChange> {
157 stream_select!(identity_updates, room_member_events)
158}
159
160async fn wrap_identity_updates(
161 client: &Client,
162) -> Result<impl Stream<Item = RoomIdentityChange> + use<>> {
163 Ok(client
164 .encryption()
165 .user_identities_stream()
166 .await?
167 .map(|item| RoomIdentityChange::IdentityUpdates(to_base_updates(item))))
168}
169
170fn to_base_updates(
171 input: IdentityUpdates,
172) -> matrix_sdk_base::crypto::store::types::IdentityUpdates {
173 matrix_sdk_base::crypto::store::types::IdentityUpdates {
174 new: to_base_identities(input.new),
175 changed: to_base_identities(input.changed),
176 unchanged: Default::default(),
177 }
178}
179
180fn to_base_identities(
181 input: BTreeMap<OwnedUserId, UserIdentity>,
182) -> BTreeMap<OwnedUserId, matrix_sdk_base::crypto::UserIdentity> {
183 input.into_iter().map(|(k, v)| (k, v.underlying_identity())).collect()
184}
185
186fn wrap_room_member_events(
187 room: &Room,
188) -> (EventHandlerDropGuard, impl Stream<Item = RoomIdentityChange> + use<>) {
189 let own_user_id = room.own_user_id().to_owned();
190 let room_id = room.room_id();
191 let (sender, receiver) = mpsc::channel(16);
192 let handle =
193 room.client.add_room_event_handler(room_id, move |event: SyncRoomMemberEvent| async move {
194 if *event.state_key() == own_user_id {
195 return;
196 }
197 let _: Result<_, _> =
198 sender.send(RoomIdentityChange::SyncRoomMemberEvent(Box::new(event))).await;
199 });
200 let drop_guard = room.client.event_handler_drop_guard(handle);
201 (drop_guard, ReceiverStream::new(receiver))
202}
203
204#[cfg(all(test, not(target_family = "wasm")))]
205mod tests {
206 use std::time::Duration;
207
208 use futures_util::{FutureExt as _, StreamExt as _, pin_mut};
209 use matrix_sdk_base::crypto::IdentityState;
210 use matrix_sdk_test::{async_test, test_json::keys_query_sets::IdentityChangeDataSet};
211 use test_setup::TestSetup;
212
213 use crate::assert_next_with_timeout;
214
215 #[async_test]
216 async fn test_when_user_becomes_unpinned_we_report_it() {
217 let t = TestSetup::new_room_with_other_bob().await;
219
220 t.pin_bob().await;
222
223 let stream = t.subscribe_to_identity_status_changes().await;
225 pin_mut!(stream);
226
227 t.unpin_bob().await;
229
230 let change = assert_next_with_timeout!(stream);
232 assert_eq!(change[0].user_id, t.bob_user_id());
233 assert_eq!(change[0].changed_to, IdentityState::PinViolation);
234 assert_eq!(change.len(), 1);
235 }
236
237 #[async_test]
238 async fn test_when_user_becomes_verification_violation_we_report_it() {
239 let t = TestSetup::new_room_with_other_bob().await;
241
242 t.verify_bob().await;
244
245 let stream = t.subscribe_to_identity_status_changes().await;
247 pin_mut!(stream);
248
249 t.unpin_bob().await;
251
252 let change = assert_next_with_timeout!(stream);
254 assert_eq!(change[0].user_id, t.bob_user_id());
255 assert_eq!(change[0].changed_to, IdentityState::VerificationViolation);
256 assert_eq!(change.len(), 1);
257 }
258
259 #[async_test]
260 async fn test_when_user_becomes_pinned_we_report_it() {
261 let t = TestSetup::new_room_with_other_bob().await;
263
264 t.unpin_bob().await;
266
267 let stream = t.subscribe_to_identity_status_changes().await;
269 pin_mut!(stream);
270
271 t.pin_bob().await;
273
274 let change1 = assert_next_with_timeout!(stream);
276 assert_eq!(change1[0].user_id, t.bob_user_id());
277 assert_eq!(change1[0].changed_to, IdentityState::PinViolation);
278 assert_eq!(change1.len(), 1);
279
280 let change2 = assert_next_with_timeout!(stream);
282 assert_eq!(change2[0].user_id, t.bob_user_id());
283 assert_eq!(change2[0].changed_to, IdentityState::Pinned);
284 assert_eq!(change2.len(), 1);
285 }
286
287 #[async_test]
288 async fn test_when_user_becomes_verified_we_report_it() {
289 let t = TestSetup::new_room_with_other_bob().await;
291
292 let stream = t.subscribe_to_identity_status_changes().await;
294 pin_mut!(stream);
295
296 t.verify_bob().await;
298
299 let change = assert_next_with_timeout!(stream);
301 assert_eq!(change[0].user_id, t.bob_user_id());
302 assert_eq!(change[0].changed_to, IdentityState::Verified);
303 assert_eq!(change.len(), 1);
304
305 t.unpin_bob().await;
307
308 let change = assert_next_with_timeout!(stream);
310 assert_eq!(change[0].user_id, t.bob_user_id());
311 assert_eq!(change[0].changed_to, IdentityState::VerificationViolation);
312 assert_eq!(change.len(), 1);
313 }
314
315 #[async_test]
316 async fn test_when_an_unpinned_user_becomes_verified_we_report_it() {
317 let t = TestSetup::new_room_with_other_bob().await;
319
320 t.unpin_bob_with(IdentityChangeDataSet::key_query_with_identity_a()).await;
322
323 let stream = t.subscribe_to_identity_status_changes().await;
325 pin_mut!(stream);
326
327 t.verify_bob().await;
329
330 let change1 = assert_next_with_timeout!(stream);
332 assert_eq!(change1[0].user_id, t.bob_user_id());
333 assert_eq!(change1[0].changed_to, IdentityState::PinViolation);
334 assert_eq!(change1.len(), 1);
335
336 let change2 = assert_next_with_timeout!(stream);
338 assert_eq!(change2[0].user_id, t.bob_user_id());
339 assert_eq!(change2[0].changed_to, IdentityState::Verified);
340 assert_eq!(change2.len(), 1);
341 }
342
343 #[async_test]
344 async fn test_when_user_in_verification_violation_becomes_verified_we_report_it() {
345 let t = TestSetup::new_room_with_other_bob().await;
347
348 t.verify_bob_with(
350 IdentityChangeDataSet::key_query_with_identity_b(),
351 IdentityChangeDataSet::master_signing_keys_b(),
352 IdentityChangeDataSet::self_signing_keys_b(),
353 )
354 .await;
355 t.unpin_bob().await;
356
357 let stream = t.subscribe_to_identity_status_changes().await;
359 pin_mut!(stream);
360
361 t.verify_bob().await;
363
364 let change1 = assert_next_with_timeout!(stream);
366 assert_eq!(change1[0].user_id, t.bob_user_id());
367 assert_eq!(change1[0].changed_to, IdentityState::VerificationViolation);
368 assert_eq!(change1.len(), 1);
369
370 let change2 = assert_next_with_timeout!(stream);
372 assert_eq!(change2[0].user_id, t.bob_user_id());
373 assert_eq!(change2[0].changed_to, IdentityState::Verified);
374 assert_eq!(change2.len(), 1);
375 }
376
377 #[async_test]
378 async fn test_when_an_unpinned_user_joins_we_report_it() {
379 let mut t = TestSetup::new_just_me_room().await;
381
382 t.unpin_bob().await;
384
385 let stream = t.subscribe_to_identity_status_changes().await;
387 pin_mut!(stream);
388
389 t.bob_joins().await;
391
392 let change = assert_next_with_timeout!(stream);
394 assert_eq!(change[0].user_id, t.bob_user_id());
395 assert_eq!(change[0].changed_to, IdentityState::PinViolation);
396 assert_eq!(change.len(), 1);
397 }
398
399 #[async_test]
400 async fn test_when_an_verification_violating_user_joins_we_report_it() {
401 let mut t = TestSetup::new_just_me_room().await;
403
404 t.verify_bob().await;
406 t.unpin_bob().await;
407
408 let stream = t.subscribe_to_identity_status_changes().await;
410 pin_mut!(stream);
411
412 t.bob_joins().await;
414
415 let change = assert_next_with_timeout!(stream);
417 assert_eq!(change[0].user_id, t.bob_user_id());
418 assert_eq!(change[0].changed_to, IdentityState::VerificationViolation);
419 assert_eq!(change.len(), 1);
420 }
421
422 #[async_test]
423 async fn test_when_a_verified_user_joins_we_dont_report_it() {
424 let mut t = TestSetup::new_just_me_room().await;
426
427 t.verify_bob().await;
429
430 let stream = t.subscribe_to_identity_status_changes().await;
432 pin_mut!(stream);
433
434 t.bob_joins().await;
436
437 t.unpin_bob().await;
439
440 let change = assert_next_with_timeout!(stream);
442 assert_eq!(change[0].user_id, t.bob_user_id());
443 assert_eq!(change[0].changed_to, IdentityState::VerificationViolation);
444 assert_eq!(change.len(), 1);
445 }
446
447 #[async_test]
448 async fn test_when_a_pinned_user_joins_we_do_not_report() {
449 let mut t = TestSetup::new_just_me_room().await;
451
452 t.pin_bob().await;
454
455 let stream = t.subscribe_to_identity_status_changes().await;
457 pin_mut!(stream);
458
459 t.bob_joins().await;
461
462 tokio::time::sleep(Duration::from_millis(200)).await;
464 let change = stream.next().now_or_never();
465 assert!(change.is_none());
466 }
467
468 #[async_test]
469 async fn test_when_an_unpinned_user_leaves_we_report_it() {
470 let mut t = TestSetup::new_room_with_other_bob().await;
472
473 t.unpin_bob().await;
475
476 let stream = t.subscribe_to_identity_status_changes().await;
478 pin_mut!(stream);
479
480 t.bob_leaves().await;
482
483 let change1 = assert_next_with_timeout!(stream);
485 assert_eq!(change1[0].user_id, t.bob_user_id());
486 assert_eq!(change1[0].changed_to, IdentityState::PinViolation);
487 assert_eq!(change1.len(), 1);
488
489 let change2 = assert_next_with_timeout!(stream);
491 assert_eq!(change2[0].user_id, t.bob_user_id());
494 assert_eq!(change2[0].changed_to, IdentityState::Pinned);
495 assert_eq!(change2.len(), 1);
496 }
497
498 #[async_test]
499 async fn test_multiple_identity_changes_are_reported() {
500 let mut t = TestSetup::new_just_me_room().await;
502
503 t.unpin_bob().await;
505
506 let stream = t.subscribe_to_identity_status_changes().await;
508 pin_mut!(stream);
509
510 t.bob_joins().await;
520 let change1 = assert_next_with_timeout!(stream);
521
522 t.pin_bob().await;
524 let change2 = assert_next_with_timeout!(stream);
525
526 t.bob_leaves().await;
528 t.bob_joins().await;
529
530 t.unpin_bob().await;
532 let change3 = assert_next_with_timeout!(stream);
533
534 t.bob_leaves().await;
536 let change4 = assert_next_with_timeout!(stream);
537
538 assert_eq!(change1[0].user_id, t.bob_user_id());
539 assert_eq!(change2[0].user_id, t.bob_user_id());
540 assert_eq!(change3[0].user_id, t.bob_user_id());
541 assert_eq!(change4[0].user_id, t.bob_user_id());
542
543 assert_eq!(change1[0].changed_to, IdentityState::PinViolation);
544 assert_eq!(change2[0].changed_to, IdentityState::Pinned);
545 assert_eq!(change3[0].changed_to, IdentityState::PinViolation);
546 assert_eq!(change4[0].changed_to, IdentityState::Pinned);
547
548 assert_eq!(change1.len(), 1);
549 assert_eq!(change2.len(), 1);
550 assert_eq!(change3.len(), 1);
551 assert_eq!(change4.len(), 1);
552 }
553
554 #[async_test]
555 async fn test_when_an_unpinned_user_is_already_present_we_report_it_immediately() {
556 let t = TestSetup::new_room_with_other_bob().await;
558 t.unpin_bob().await;
559
560 let stream = t.subscribe_to_identity_status_changes().await;
562 pin_mut!(stream);
563
564 let change = assert_next_with_timeout!(stream);
566 assert_eq!(change[0].user_id, t.bob_user_id());
567 assert_eq!(change[0].changed_to, IdentityState::PinViolation);
568 assert_eq!(change.len(), 1);
569 }
570
571 #[async_test]
572 async fn test_when_a_verified_user_is_already_present_we_dont_report_it() {
573 let t = TestSetup::new_room_with_other_bob().await;
575 t.verify_bob().await;
576
577 let stream = t.subscribe_to_identity_status_changes().await;
579 pin_mut!(stream);
580
581 t.unpin_bob().await;
583
584 let next_change = assert_next_with_timeout!(stream);
586
587 assert_eq!(next_change[0].user_id, t.bob_user_id());
588 assert_eq!(next_change[0].changed_to, IdentityState::VerificationViolation);
589 assert_eq!(next_change.len(), 1);
590 }
591
592 mod test_setup {
597 use futures_core::Stream;
598 use matrix_sdk_base::{
599 RoomState,
600 crypto::{
601 IdentityStatusChange, OtherUserIdentity,
602 testing::simulate_key_query_response_for_verification,
603 },
604 };
605 use matrix_sdk_test::{
606 DEFAULT_TEST_ROOM_ID, JoinedRoomBuilder, SyncResponseBuilder,
607 event_factory::EventFactory, test_json,
608 test_json::keys_query_sets::IdentityChangeDataSet,
609 };
610 use ruma::{
611 OwnedUserId, TransactionId, UserId,
612 api::client::keys::{get_keys, get_keys::v3::Response as KeyQueryResponse},
613 events::room::member::MembershipState,
614 owned_user_id, user_id,
615 };
616 use serde_json::json;
617 use wiremock::{
618 Mock, MockServer, ResponseTemplate,
619 matchers::{header, method, path_regex},
620 };
621
622 use crate::{
623 Client, Room, encryption::identities::UserIdentity, test_utils::logged_in_client,
624 };
625
626 pub(super) struct TestSetup {
636 client: Client,
637 bob_user_id: OwnedUserId,
638 sync_response_builder: SyncResponseBuilder,
639 room: Room,
640 }
641
642 impl TestSetup {
643 pub(super) async fn new_just_me_room() -> Self {
644 let (client, user_id, mut sync_response_builder) = Self::init().await;
645 let room = create_just_me_room(&client, &mut sync_response_builder).await;
646 Self { client, bob_user_id: user_id, sync_response_builder, room }
647 }
648
649 pub(super) async fn new_room_with_other_bob() -> Self {
650 let (client, bob_user_id, mut sync_response_builder) = Self::init().await;
651 let room = create_room_with_other_member(
652 &mut sync_response_builder,
653 &client,
654 &bob_user_id,
655 )
656 .await;
657 Self { client, bob_user_id, sync_response_builder, room }
658 }
659
660 pub(super) fn bob_user_id(&self) -> &UserId {
661 &self.bob_user_id
662 }
663
664 pub(super) async fn pin_bob(&self) {
665 if self.bob_user_identity().await.is_some() {
666 assert!(
667 !self.bob_is_pinned().await,
668 "pin_bob() called when the identity is already pinned!"
669 );
670
671 self.bob_user_identity()
673 .await
674 .expect("User should exist")
675 .pin()
676 .await
677 .expect("Should not fail to pin");
678 } else {
679 self.change_bob_identity(IdentityChangeDataSet::key_query_with_identity_a())
682 .await;
683 }
684
685 assert!(self.bob_is_pinned().await);
687 }
688
689 pub(super) async fn unpin_bob(&self) {
690 self.unpin_bob_with(IdentityChangeDataSet::key_query_with_identity_b()).await;
691 }
692
693 pub(super) async fn unpin_bob_with(&self, requested: KeyQueryResponse) {
694 fn master_key_json(key_query_response: &KeyQueryResponse) -> String {
695 serde_json::to_string(
696 key_query_response
697 .master_keys
698 .first_key_value()
699 .expect("Master key should have a value")
700 .1,
701 )
702 .expect("Should be able to serialise master key")
703 }
704
705 let a = IdentityChangeDataSet::key_query_with_identity_a();
706 let b = IdentityChangeDataSet::key_query_with_identity_b();
707 let requested_master_key = master_key_json(&requested);
708 let a_master_key = master_key_json(&a);
709
710 if requested_master_key == a_master_key {
714 self.change_bob_identity(b).await;
715 if !self.bob_is_pinned().await {
716 self.pin_bob().await;
717 }
718 self.change_bob_identity(a).await;
719 } else {
720 self.change_bob_identity(a).await;
721 if !self.bob_is_pinned().await {
722 self.pin_bob().await;
723 }
724 self.change_bob_identity(b).await;
725 }
726
727 assert!(!self.bob_is_pinned().await);
729 }
730
731 pub(super) async fn verify_bob(&self) {
732 self.verify_bob_with(
733 IdentityChangeDataSet::key_query_with_identity_a(),
734 IdentityChangeDataSet::master_signing_keys_a(),
735 IdentityChangeDataSet::self_signing_keys_a(),
736 )
737 .await;
738 }
739
740 pub(super) async fn verify_bob_with(
741 &self,
742 key_query: KeyQueryResponse,
743 master_signing_key: serde_json::Value,
744 self_signing_key: serde_json::Value,
745 ) {
746 self.change_bob_identity(key_query).await;
748
749 let my_user_id = self.client.user_id().expect("I should have a user id");
750 let my_identity = self
751 .client
752 .encryption()
753 .get_user_identity(my_user_id)
754 .await
755 .expect("Should not fail to get own user identity")
756 .expect("Should have an own user identity")
757 .underlying_identity()
758 .own()
759 .expect("Our own identity should be of type Own");
760
761 let signature_upload_request = self
763 .bob_crypto_other_identity()
764 .await
765 .verify()
766 .await
767 .expect("Should be able to verify other identity");
768
769 let verification_response = simulate_key_query_response_for_verification(
770 signature_upload_request,
771 my_identity,
772 my_user_id,
773 self.bob_user_id(),
774 master_signing_key,
775 self_signing_key,
776 );
777
778 self.client
780 .mark_request_as_sent(&TransactionId::new(), &verification_response)
781 .await
782 .unwrap();
783
784 assert!(self.bob_is_verified().await);
786 }
787
788 pub(super) async fn bob_joins(&mut self) {
789 self.bob_membership_change(MembershipState::Join).await;
790 }
791
792 pub(super) async fn bob_leaves(&mut self) {
793 self.bob_membership_change(MembershipState::Leave).await;
794 }
795
796 pub(super) async fn subscribe_to_identity_status_changes(
797 &self,
798 ) -> impl Stream<Item = Vec<IdentityStatusChange>> + use<> {
799 self.room
800 .subscribe_to_identity_status_changes()
801 .await
802 .expect("Should be able to subscribe")
803 }
804
805 async fn init() -> (Client, OwnedUserId, SyncResponseBuilder) {
806 let (client, _server) = create_client_and_server().await;
807
808 client
810 .olm_machine()
811 .await
812 .as_ref()
813 .expect("We should have an Olm machine")
814 .bootstrap_cross_signing(true)
815 .await
816 .expect("Should be able to bootstrap cross-signing");
817
818 let bob_user_id = owned_user_id!("@bob:localhost");
821
822 let sync_response_builder = SyncResponseBuilder::default();
823
824 (client, bob_user_id, sync_response_builder)
825 }
826
827 async fn change_bob_identity(
828 &self,
829 key_query_response: get_keys::v3::Response,
830 ) -> OtherUserIdentity {
831 self.client
832 .mark_request_as_sent(&TransactionId::new(), &key_query_response)
833 .await
834 .expect("Should not fail to send identity changes");
835
836 self.bob_crypto_other_identity().await
837 }
838
839 async fn bob_membership_change(&mut self, new_state: MembershipState) {
840 let f = EventFactory::new().sender(user_id!("@example:localhost"));
841 let sync_response = self
842 .sync_response_builder
843 .add_joined_room(
844 JoinedRoomBuilder::new(&DEFAULT_TEST_ROOM_ID).add_state_event(
845 f.member(&self.bob_user_id).membership(new_state.clone()),
846 ),
847 )
848 .build_sync_response();
849 self.room.client.process_sync(sync_response).await.unwrap();
850
851 let m = self
853 .room
854 .get_member_no_sync(&self.bob_user_id)
855 .await
856 .expect("Should not fail to get member");
857
858 match (&new_state, m) {
859 (MembershipState::Leave, None) => {}
860 (_, None) => {
861 panic!("Member should exist")
862 }
863 (_, Some(m)) => {
864 assert_eq!(*m.membership(), new_state);
865 }
866 }
867 }
868
869 async fn bob_is_pinned(&self) -> bool {
870 !self.bob_crypto_other_identity().await.identity_needs_user_approval()
871 }
872
873 async fn bob_is_verified(&self) -> bool {
874 self.bob_crypto_other_identity().await.is_verified()
875 }
876
877 async fn bob_crypto_other_identity(&self) -> OtherUserIdentity {
878 self.bob_user_identity()
879 .await
880 .expect("User identity should exist")
881 .underlying_identity()
882 .other()
883 .expect("Identity should be Other, not Own")
884 }
885
886 async fn bob_user_identity(&self) -> Option<UserIdentity> {
887 self.client
888 .encryption()
889 .get_user_identity(&self.bob_user_id)
890 .await
891 .expect("Should not fail to get user identity")
892 }
893 }
894
895 async fn create_just_me_room(
896 client: &Client,
897 sync_response_builder: &mut SyncResponseBuilder,
898 ) -> Room {
899 let f = EventFactory::new().sender(user_id!("@example:localhost"));
900 let create_room_sync_response = sync_response_builder
901 .add_joined_room(JoinedRoomBuilder::new(&DEFAULT_TEST_ROOM_ID).add_state_event(
902 f.member(user_id!("@example:localhost")).display_name("example"),
903 ))
904 .build_sync_response();
905 client.process_sync(create_room_sync_response).await.unwrap();
906 let room = client.get_room(&DEFAULT_TEST_ROOM_ID).expect("Room should exist");
907 assert_eq!(room.state(), RoomState::Joined);
908 room
909 }
910
911 async fn create_room_with_other_member(
912 builder: &mut SyncResponseBuilder,
913 client: &Client,
914 other_user_id: &UserId,
915 ) -> Room {
916 let f = EventFactory::new().sender(user_id!("@example:localhost"));
917 let create_room_sync_response = builder
918 .add_joined_room(
919 JoinedRoomBuilder::new(&DEFAULT_TEST_ROOM_ID)
920 .add_state_event(
921 f.member(user_id!("@example:localhost")).display_name("example"),
922 )
923 .add_state_event(f.member(other_user_id).membership(MembershipState::Join)),
924 )
925 .build_sync_response();
926 client.process_sync(create_room_sync_response).await.unwrap();
927 let room = client.get_room(&DEFAULT_TEST_ROOM_ID).expect("Room should exist");
928 room.inner.mark_members_synced();
929
930 assert_eq!(room.state(), RoomState::Joined);
931 assert_eq!(
932 *room
933 .get_member_no_sync(other_user_id)
934 .await
935 .expect("Should not fail to get member")
936 .expect("Member should exist")
937 .membership(),
938 MembershipState::Join
939 );
940 room
941 }
942
943 async fn create_client_and_server() -> (Client, MockServer) {
944 let server = MockServer::start().await;
945 mock_members_request(&server).await;
946 mock_secret_storage_default_key(&server).await;
947 let client = logged_in_client(Some(server.uri())).await;
948 (client, server)
949 }
950
951 async fn mock_members_request(server: &MockServer) {
952 Mock::given(method("GET"))
953 .and(path_regex(r"^/_matrix/client/r0/rooms/.*/members"))
954 .and(header("authorization", "Bearer 1234"))
955 .respond_with(
956 ResponseTemplate::new(200).set_body_json(&*test_json::members::MEMBERS),
957 )
958 .mount(server)
959 .await;
960 }
961
962 async fn mock_secret_storage_default_key(server: &MockServer) {
963 Mock::given(method("GET"))
964 .and(path_regex(
965 r"^/_matrix/client/r0/user/.*/account_data/m.secret_storage.default_key",
966 ))
967 .and(header("authorization", "Bearer 1234"))
968 .respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
969 .mount(server)
970 .await;
971 }
972 }
973}