1use std::collections::{BTreeMap, BTreeSet};
24
25use futures_core::Stream;
26use futures_util::StreamExt;
27#[cfg(feature = "experimental-encrypted-state-events")]
28use matrix_sdk_base::crypto::types::events::room::encrypted::EncryptedEvent;
29use matrix_sdk_base::crypto::{
30 OlmMachine, RoomKeyImportResult,
31 backups::MegolmV1BackupKey,
32 store::types::BackupDecryptionKey,
33 types::{RoomKeyBackupInfo, requests::KeysBackupRequest},
34};
35#[cfg(feature = "experimental-push-secrets")]
36use ruma::events::secret::push::ToDeviceSecretPushEvent;
37#[cfg(feature = "experimental-encrypted-state-events")]
38use ruma::serde::JsonCastable;
39use ruma::{
40 OwnedRoomId, RoomId, TransactionId,
41 api::{
42 client::backup::{
43 RoomKeyBackup, add_backup_keys, create_backup_version, get_backup_keys,
44 get_backup_keys_for_room, get_backup_keys_for_session, get_latest_backup_info,
45 },
46 error::ErrorKind,
47 },
48 events::{
49 room::encrypted::OriginalSyncRoomEncryptedEvent,
50 secret::{request::SecretName, send::ToDeviceSecretSendEvent},
51 },
52 serde::Raw,
53};
54use tokio_stream::wrappers::{BroadcastStream, errors::BroadcastStreamRecvError};
55use tracing::{Span, error, info, instrument, trace, warn};
56
57pub mod futures;
58pub(crate) mod types;
59
60use matrix_sdk_base::crypto::olm::ExportedRoomKey;
61pub use types::{BackupState, UploadState};
62
63use self::futures::WaitForSteadyState;
64use crate::{Client, Error, Room, encryption::BackupDownloadStrategy};
65
66#[derive(Debug, Clone)]
68pub struct Backups {
69 pub(super) client: Client,
70}
71
72impl Backups {
73 pub async fn create(&self) -> Result<(), Error> {
98 self.client.inner.e2ee.backup_state.clear_backup_exists_on_server();
99 let _guard = self.client.locks().backup_modify_lock.lock().await;
100
101 self.set_state(BackupState::Creating);
102
103 let future = async {
106 let olm_machine = self.client.olm_machine().await;
107 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
108
109 let decryption_key = BackupDecryptionKey::new();
111
112 let mut backup_info = decryption_key.to_backup_info();
134
135 if let Err(e) = olm_machine.backup_machine().sign_backup(&mut backup_info).await {
136 warn!("Unable to sign the newly created backup version: {e:?}");
137 }
138
139 let algorithm = Raw::new(&backup_info)?.cast();
140 let request = create_backup_version::v3::Request::new(algorithm);
141 let response = self.client.send(request).await?;
142 let version = response.version;
143
144 olm_machine.backup_machine().disable_backup().await?;
147
148 let backup_key = decryption_key.megolm_v1_public_key();
149
150 olm_machine
153 .backup_machine()
154 .save_decryption_key(Some(decryption_key), Some(version.to_owned()))
155 .await?;
156
157 self.enable(olm_machine, backup_key, version).await?;
159
160 #[cfg(feature = "experimental-push-secrets")]
161 {
162 if let Some((txn_id, keys_claim_request)) = olm_machine
168 .get_missing_sessions(vec![olm_machine.user_id()].into_iter())
169 .await?
170 {
171 let keys_claim_response = self.client.send(keys_claim_request).await?;
172 olm_machine.mark_request_as_sent(&txn_id, &keys_claim_response).await?;
173 }
174
175 let _ = olm_machine.push_secret_to_verified_devices(SecretName::RecoveryKey).await;
179 }
180
181 Ok(())
182 };
183
184 let result = future.await;
185
186 if result.is_err() {
187 self.set_state(BackupState::Unknown);
188 }
189
190 result
191 }
192
193 #[instrument(skip_all, fields(version))]
214 pub async fn disable(&self) -> Result<(), Error> {
215 let _guard = self.client.locks().backup_modify_lock.lock().await;
216
217 self.set_state(BackupState::Disabling);
218
219 let future = async {
222 let olm_machine = self.client.olm_machine().await;
223 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
224
225 let backup_keys = olm_machine.backup_machine().get_backup_keys().await?;
226
227 if let Some(version) = backup_keys.backup_version {
228 Span::current().record("version", &version);
229 info!("Deleting and disabling backup");
230
231 self.delete_backup_from_server(version).await?;
232 info!("Backup successfully deleted");
233
234 olm_machine.backup_machine().disable_backup().await?;
235
236 info!("Backup successfully disabled and deleted");
237
238 Ok(())
239 } else {
240 info!("Backup is not enabled, can't disable it");
241 Err(Error::BackupNotEnabled)
242 }
243 };
244
245 let result = future.await;
246
247 self.set_state(BackupState::Unknown);
248
249 result
250 }
251
252 pub async fn disable_and_delete(&self) -> Result<(), Error> {
278 let _guard = self.client.locks().backup_modify_lock.lock().await;
279
280 self.set_state(BackupState::Disabling);
281
282 let future = async {
285 while let Some(response) = self.get_current_version().await? {
286 self.delete_backup_from_server(response.version).await?;
287 }
288
289 let olm_machine = self.client.olm_machine().await;
290 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
291
292 olm_machine.backup_machine().disable_backup().await?;
293
294 Ok(())
295 };
296
297 let result = future.await;
298
299 self.set_state(BackupState::Unknown);
300
301 result
302 }
303
304 pub fn wait_for_steady_state(&self) -> WaitForSteadyState<'_> {
348 WaitForSteadyState {
349 backups: self,
350 progress: self.client.inner.e2ee.backup_state.upload_progress.clone(),
351 timeout: None,
352 }
353 }
354
355 pub fn state_stream(
386 &self,
387 ) -> impl Stream<Item = Result<BackupState, BroadcastStreamRecvError>> + use<> {
388 self.client.inner.e2ee.backup_state.global_state.subscribe()
389 }
390
391 pub fn state(&self) -> BackupState {
393 self.client.inner.e2ee.backup_state.global_state.get()
394 }
395
396 pub async fn are_enabled(&self) -> bool {
401 let olm_machine = self.client.olm_machine().await;
402
403 if let Some(machine) = olm_machine.as_ref() {
404 machine.backup_machine().enabled().await
405 } else {
406 false
407 }
408 }
409
410 pub async fn fetch_exists_on_server(&self) -> Result<bool, Error> {
415 let exists_on_server = self.get_current_version().await?.is_some();
416 self.client.inner.e2ee.backup_state.set_backup_exists_on_server(exists_on_server);
417 Ok(exists_on_server)
418 }
419
420 pub async fn exists_on_server(&self) -> Result<bool, Error> {
431 if let Some(cached_value) = self.client.inner.e2ee.backup_state.backup_exists_on_server() {
433 return Ok(cached_value);
434 }
435
436 self.fetch_exists_on_server().await
439 }
440
441 pub fn room_keys_for_room_stream(
444 &self,
445 room_id: &RoomId,
446 ) -> impl Stream<Item = Result<BTreeMap<String, BTreeSet<String>>, BroadcastStreamRecvError>> + use<>
447 {
448 let room_id = room_id.to_owned();
449
450 self.room_keys_stream().filter_map(move |import_result| {
458 let room_id = room_id.to_owned();
459
460 async move {
461 match import_result {
462 Ok(mut import_result) => import_result.keys.remove(&room_id).map(Ok),
463 Err(e) => Some(Err(e)),
464 }
465 }
466 })
467 }
468
469 pub async fn download_room_keys_for_room(&self, room_id: &RoomId) -> Result<(), Error> {
472 let olm_machine = self.client.olm_machine().await;
473 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
474
475 let backup_keys = olm_machine.store().load_backup_keys().await?;
476
477 if let Some(decryption_key) = backup_keys.decryption_key
478 && let Some(version) = backup_keys.backup_version
479 {
480 let request =
481 get_backup_keys_for_room::v3::Request::new(version.clone(), room_id.to_owned());
482 let response = self.client.send(request).await?;
483
484 let response = get_backup_keys::v3::Response::new(BTreeMap::from([(
487 room_id.to_owned(),
488 RoomKeyBackup::new(response.sessions),
489 )]));
490
491 self.handle_downloaded_room_keys(response, decryption_key, &version, olm_machine)
492 .await?;
493 }
494
495 Ok(())
496 }
497
498 pub async fn download_room_key(
504 &self,
505 room_id: &RoomId,
506 session_id: &str,
507 ) -> Result<bool, Error> {
508 let olm_machine = self.client.olm_machine().await;
509 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
510
511 let backup_keys = olm_machine.store().load_backup_keys().await?;
512
513 if let Some(decryption_key) = backup_keys.decryption_key {
514 if let Some(version) = backup_keys.backup_version {
515 let request = get_backup_keys_for_session::v3::Request::new(
516 version.clone(),
517 room_id.to_owned(),
518 session_id.to_owned(),
519 );
520 let response = self.client.send(request).await?;
521
522 let response = get_backup_keys::v3::Response::new(BTreeMap::from([(
525 room_id.to_owned(),
526 RoomKeyBackup::new(BTreeMap::from([(
527 session_id.to_owned(),
528 response.key_data,
529 )])),
530 )]));
531
532 self.handle_downloaded_room_keys(response, decryption_key, &version, olm_machine)
533 .await?;
534
535 Ok(true)
536 } else {
537 Ok(false)
538 }
539 } else {
540 Ok(false)
541 }
542 }
543
544 fn set_state(&self, new_state: BackupState) {
546 let old_state = self.client.inner.e2ee.backup_state.global_state.set(new_state);
547
548 if old_state != new_state {
549 info!("Backup state changed from {old_state:?} to {new_state:?}");
550 }
551 }
552
553 async fn enable(
556 &self,
557 olm_machine: &OlmMachine,
558 backup_key: MegolmV1BackupKey,
559 version: String,
560 ) -> Result<(), Error> {
561 backup_key.set_version(version);
562 olm_machine.backup_machine().enable_backup_v1(backup_key).await?;
563
564 self.set_state(BackupState::Enabled);
565
566 Ok(())
567 }
568
569 async fn handle_downloaded_room_keys(
572 &self,
573 backed_up_keys: get_backup_keys::v3::Response,
574 backup_decryption_key: BackupDecryptionKey,
575 backup_version: &str,
576 olm_machine: &OlmMachine,
577 ) -> Result<(), Error> {
578 let mut decrypted_room_keys: Vec<_> = Vec::new();
579
580 for (room_id, room_keys) in backed_up_keys.rooms {
581 for (session_id, room_key) in room_keys.sessions {
582 let room_key = match room_key.deserialize() {
583 Ok(k) => k,
584 Err(e) => {
585 warn!(
586 "Couldn't deserialize a room key we downloaded from backups, session \
587 ID: {session_id}, error: {e:?}"
588 );
589 continue;
590 }
591 };
592
593 let room_key =
594 match backup_decryption_key.decrypt_session_data(room_key.session_data) {
595 Ok(k) => k,
596 Err(e) => {
597 warn!(
598 "Couldn't decrypt a room key we downloaded from backups, session \
599 ID: {session_id}, error: {e:?}"
600 );
601 continue;
602 }
603 };
604
605 decrypted_room_keys.push(ExportedRoomKey::from_backed_up_room_key(
606 room_id.to_owned(),
607 session_id,
608 room_key,
609 ));
610 }
611 }
612
613 let result = olm_machine
614 .store()
615 .import_room_keys(decrypted_room_keys, Some(backup_version), |_, _| {})
616 .await?;
617
618 let _ = self.client.inner.e2ee.backup_state.room_keys_broadcaster.send(result);
621
622 Ok(())
623 }
624
625 async fn download_all_room_keys(
627 &self,
628 decryption_key: BackupDecryptionKey,
629 version: String,
630 ) -> Result<(), Error> {
631 let request = get_backup_keys::v3::Request::new(version.clone());
632 let response = self.client.send(request).await?;
633
634 let olm_machine = self.client.olm_machine().await;
635 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
636
637 self.handle_downloaded_room_keys(response, decryption_key, &version, olm_machine).await?;
638
639 Ok(())
640 }
641
642 fn room_keys_stream(
643 &self,
644 ) -> impl Stream<Item = Result<RoomKeyImportResult, BroadcastStreamRecvError>> + use<> {
645 BroadcastStream::new(self.client.inner.e2ee.backup_state.room_keys_broadcaster.subscribe())
646 }
647
648 async fn get_current_version(
650 &self,
651 ) -> Result<Option<get_latest_backup_info::v3::Response>, Error> {
652 let request = get_latest_backup_info::v3::Request::new();
653
654 match self.client.send(request).await {
655 Ok(r) => Ok(Some(r)),
656 Err(e) => {
657 if let Some(kind) = e.client_api_error_kind() {
658 if kind == &ErrorKind::NotFound { Ok(None) } else { Err(e.into()) }
659 } else {
660 Err(e.into())
661 }
662 }
663 }
664 }
665
666 async fn delete_backup_from_server(&self, version: String) -> Result<(), Error> {
667 let request = ruma::api::client::backup::delete_backup_version::v3::Request::new(version);
668
669 let ret = match self.client.send(request).await {
670 Ok(_) => Ok(()),
671 Err(e) => {
672 if let Some(kind) = e.client_api_error_kind() {
673 if kind == &ErrorKind::NotFound { Ok(()) } else { Err(e.into()) }
674 } else {
675 Err(e.into())
676 }
677 }
678 };
679
680 self.client.inner.e2ee.backup_state.clear_backup_exists_on_server();
684
685 ret
686 }
687
688 #[instrument(skip(self, olm_machine, request))]
689 async fn send_backup_request(
690 &self,
691 olm_machine: &OlmMachine,
692 request_id: &TransactionId,
693 request: KeysBackupRequest,
694 ) -> Result<(), Error> {
695 trace!("Uploading some room keys");
696
697 let add_backup_keys = add_backup_keys::v3::Request::new(request.version, request.rooms);
698
699 match self.client.send(add_backup_keys).await {
700 Ok(response) => {
701 olm_machine.mark_request_as_sent(request_id, &response).await?;
702
703 let new_counts = olm_machine.backup_machine().room_key_counts().await?;
704
705 self.client
706 .inner
707 .e2ee
708 .backup_state
709 .upload_progress
710 .set(UploadState::Uploading(new_counts));
711
712 let delay =
713 self.client.inner.e2ee.backup_state.upload_delay.read().unwrap().to_owned();
714 crate::sleep::sleep(delay).await;
715
716 Ok(())
717 }
718 Err(error) => {
719 if let Some(kind) = error.client_api_error_kind() {
720 match kind {
721 ErrorKind::NotFound => {
722 warn!(
723 "No backup found on the server, the backup likely got deleted, \
724 disabling backups."
725 );
726
727 self.handle_deleted_backup_version(olm_machine).await?;
728 }
729 ErrorKind::WrongRoomKeysVersion(wrong_version) => {
730 warn!(
731 new_version = wrong_version.current_version,
732 "A new backup version was found on the server, disabling backups."
733 );
734
735 self.handle_deleted_backup_version(olm_machine).await?;
740 }
741
742 _ => (),
743 }
744 }
745
746 Err(error.into())
747 }
748 }
749 }
750
751 pub(crate) async fn backup_room_keys(&self) -> Result<(), Error> {
758 let _guard = self.client.locks().backup_upload_lock.lock().await;
759
760 let olm_machine = self.client.olm_machine().await;
761 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
762
763 while let Some((request_id, request)) = olm_machine.backup_machine().backup().await? {
764 self.send_backup_request(olm_machine, &request_id, request).await?;
765 }
766
767 self.client.inner.e2ee.backup_state.upload_progress.set(UploadState::Done);
768
769 Ok(())
770 }
771
772 pub(crate) async fn setup_and_resume(&self) -> Result<(), Error> {
775 info!("Setting up secret listeners and trying to resume backups");
776
777 self.client.add_event_handler(Self::secret_send_event_handler);
778 #[cfg(feature = "experimental-push-secrets")]
779 self.client.add_event_handler(Self::secret_push_event_handler);
780
781 if self.client.inner.e2ee.encryption_settings.backup_download_strategy
782 == BackupDownloadStrategy::AfterDecryptionFailure
783 {
784 self.client.add_event_handler(Self::utd_event_handler);
785 }
786
787 self.maybe_resume_backups().await?;
788
789 Ok(())
790 }
791
792 #[instrument(skip_all)]
808 pub(crate) async fn maybe_enable_backups(
809 &self,
810 maybe_recovery_key: &str,
811 ) -> Result<bool, EnableBackupError> {
812 let _guard = self.client.locks().backup_modify_lock.lock().await;
813
814 let future = async {
817 self.set_state(BackupState::Enabling);
818
819 let olm_machine = self.client.olm_machine().await;
820 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
821 let backup_machine = olm_machine.backup_machine();
822
823 let decryption_key =
824 BackupDecryptionKey::from_base64(maybe_recovery_key).map_err(|e| {
825 <serde_json::Error as serde::de::Error>::custom(format!(
826 "Couldn't deserialize the backup recovery key: {e:?}"
827 ))
828 })?;
829
830 let current_version = self.get_current_version().await?;
832
833 let Some(current_version) = current_version else {
834 warn!("Tried to enable backups, but no backup version was found on the server.");
835 return Ok(false);
836 };
837
838 Span::current().record("backup_version", ¤t_version.version);
839
840 let backup_info: RoomKeyBackupInfo = current_version.algorithm.deserialize_as()?;
841 let stored_keys = backup_machine.get_backup_keys().await?;
842
843 if stored_keys.backup_version.as_ref() == Some(¤t_version.version)
844 && self.are_enabled().await
845 {
846 Ok(true)
850 } else if decryption_key.backup_key_matches(&backup_info) {
851 info!(
852 "We have found the correct backup recovery key. Storing the backup recovery \
853 key and enabling backups."
854 );
855
856 backup_machine.disable_backup().await?;
859
860 let backup_key = decryption_key.megolm_v1_public_key();
861 backup_key.set_version(current_version.version.to_owned());
862
863 backup_machine
865 .save_decryption_key(
866 Some(decryption_key.to_owned()),
867 Some(current_version.version.to_owned()),
868 )
869 .await?;
870 backup_machine.enable_backup_v1(backup_key).await?;
871
872 if self.client.inner.e2ee.encryption_settings.backup_download_strategy
882 == BackupDownloadStrategy::OneShot
883 {
884 self.set_state(BackupState::Downloading);
885
886 if let Err(e) =
887 self.download_all_room_keys(decryption_key, current_version.version).await
888 {
889 warn!("Couldn't automatically download all room keys from backup: {e:?}");
890 }
891 }
892
893 self.maybe_trigger_backup();
895
896 Ok(true)
897 } else {
898 let derived_key = decryption_key.megolm_v1_public_key();
899 let downloaded_key = current_version.algorithm;
900
901 warn!(
902 ?derived_key,
903 ?downloaded_key,
904 "Found an active backup but the recovery key we received isn't the one used for \
905 this backup version"
906 );
907
908 Err(EnableBackupError::InconsistentBackupDecryptionKey)
909 }
910 };
911
912 match future.await {
913 Ok(enabled) => {
914 if enabled {
915 self.set_state(BackupState::Enabled);
916 } else {
917 self.set_state(BackupState::Unknown);
918 }
919
920 Ok(enabled)
921 }
922 Err(e) => {
923 self.set_state(BackupState::Unknown);
924
925 Err(e)
926 }
927 }
928 }
929
930 async fn resume_backup_from_stored_backup_key(
935 &self,
936 olm_machine: &OlmMachine,
937 ) -> Result<bool, Error> {
938 let backup_keys = olm_machine.store().load_backup_keys().await?;
939
940 if let Some(decryption_key) = backup_keys.decryption_key {
941 if let Some(version) = backup_keys.backup_version {
942 let backup_key = decryption_key.megolm_v1_public_key();
943
944 self.enable(olm_machine, backup_key, version).await?;
945
946 Ok(true)
947 } else {
948 Ok(false)
949 }
950 } else {
951 Ok(false)
952 }
953 }
954
955 async fn maybe_resume_from_secret_inbox(&self, olm_machine: &OlmMachine) -> Result<(), Error> {
959 let secrets = olm_machine.store().get_secrets_from_inbox(&SecretName::RecoveryKey).await?;
960
961 for secret in secrets {
962 match self.maybe_enable_backups(&secret).await {
963 Ok(enabled) => {
964 if enabled {
965 break;
966 }
967 }
968 Err(EnableBackupError::InconsistentBackupDecryptionKey) => {
969 }
972 Err(EnableBackupError::Error(e)) => return Err(e),
973 }
974 }
975
976 olm_machine.store().delete_secrets_from_inbox(&SecretName::RecoveryKey).await?;
977
978 Ok(())
979 }
980
981 pub(super) async fn maybe_resume_backups(&self) -> Result<(), Error> {
983 let olm_machine = self.client.olm_machine().await;
984 let olm_machine = olm_machine.as_ref().ok_or(Error::NoOlmMachine)?;
985
986 if !self.resume_backup_from_stored_backup_key(olm_machine).await? {
989 self.maybe_resume_from_secret_inbox(olm_machine).await?;
993 }
994
995 Ok(())
996 }
997
998 #[instrument(skip_all)]
1001 pub(crate) async fn secret_send_event_handler(_: ToDeviceSecretSendEvent, client: Client) {
1002 let olm_machine = client.olm_machine().await;
1003
1004 if let Some(olm_machine) = olm_machine.as_ref() {
1009 if let Err(e) =
1010 client.encryption().backups().maybe_resume_from_secret_inbox(olm_machine).await
1011 {
1012 error!("Could not handle `m.secret.send` event: {e:?}");
1013 }
1014 } else {
1015 error!("Tried to handle a `m.secret.send` event but no OlmMachine was initialized");
1016 }
1017 }
1018
1019 #[cfg(feature = "experimental-push-secrets")]
1022 #[instrument(skip_all)]
1023 pub(crate) async fn secret_push_event_handler(_: ToDeviceSecretPushEvent, client: Client) {
1024 let olm_machine = client.olm_machine().await;
1025
1026 if let Some(olm_machine) = olm_machine.as_ref() {
1032 if let Err(e) =
1033 client.encryption().backups().maybe_resume_from_secret_inbox(olm_machine).await
1034 {
1035 error!("Could not handle `io.element.msc4385.secret.push` event: {e:?}");
1036 }
1037 } else {
1038 error!(
1039 "Tried to handle a `io.element.msc4385.secret.push` event but no OlmMachine was initialized"
1040 );
1041 }
1042 }
1043
1044 #[allow(clippy::unused_async)] pub(crate) async fn utd_event_handler(
1053 event: Raw<OriginalSyncRoomEncryptedEvent>,
1054 room: Room,
1055 client: Client,
1056 ) {
1057 client.encryption().backups().maybe_download_room_key(room.room_id().to_owned(), event);
1058 }
1059
1060 #[cfg(not(feature = "experimental-encrypted-state-events"))]
1063 pub(crate) fn maybe_download_room_key(
1064 &self,
1065 room_id: OwnedRoomId,
1066 event: Raw<OriginalSyncRoomEncryptedEvent>,
1067 ) {
1068 let tasks = self.client.inner.e2ee.tasks.lock();
1069 if let Some(task) = tasks.download_room_keys.as_ref() {
1070 task.trigger_download_for_utd_event(room_id, event);
1071 }
1072 }
1073
1074 #[cfg(feature = "experimental-encrypted-state-events")]
1077 pub(crate) fn maybe_download_room_key<T: JsonCastable<EncryptedEvent>>(
1078 &self,
1079 room_id: OwnedRoomId,
1080 event: Raw<T>,
1081 ) {
1082 let tasks = self.client.inner.e2ee.tasks.lock();
1083 if let Some(task) = tasks.download_room_keys.as_ref() {
1084 task.trigger_download_for_utd_event(room_id, event);
1085 }
1086 }
1087
1088 pub(crate) fn maybe_trigger_backup(&self) {
1091 let tasks = self.client.inner.e2ee.tasks.lock();
1092
1093 if let Some(tasks) = tasks.upload_room_keys.as_ref() {
1094 tasks.trigger_upload();
1095 }
1096 }
1097
1098 async fn handle_deleted_backup_version(&self, olm_machine: &OlmMachine) -> Result<(), Error> {
1101 olm_machine.backup_machine().disable_backup().await?;
1102 self.set_state(BackupState::Unknown);
1103
1104 Ok(())
1105 }
1106}
1107
1108#[derive(Debug, thiserror::Error)]
1110pub enum EnableBackupError {
1111 #[error("The backup decryption key does not match the latest backup version")]
1114 InconsistentBackupDecryptionKey,
1115
1116 #[error(transparent)]
1118 Error(Error),
1119}
1120
1121impl<T: Into<Error>> From<T> for EnableBackupError {
1122 fn from(value: T) -> Self {
1123 Self::Error(value.into())
1124 }
1125}
1126
1127#[cfg(all(test, not(target_family = "wasm")))]
1128mod test {
1129 use std::{assert_matches, time::Duration};
1130
1131 use matrix_sdk_base::crypto::{
1132 GossipRequest, GossippedSecret, SecretInfo,
1133 store::types::Changes,
1134 types::events::{
1135 olm_v1::{DecryptedSecretSendEvent, OlmV1Keys},
1136 secret_send::SecretSendContent,
1137 },
1138 };
1139 use matrix_sdk_test::async_test;
1140 #[cfg(feature = "experimental-push-secrets")]
1141 use ruma::{device_id, user_id};
1142 use serde_json::json;
1143 use vodozemac::Curve25519PublicKey;
1144 use wiremock::{
1145 Mock, MockServer, ResponseTemplate,
1146 matchers::{header, method, path},
1147 };
1148
1149 use super::*;
1150 use crate::test_utils::{logged_in_client, mocks::MatrixMockServer};
1151
1152 fn room_key() -> ExportedRoomKey {
1153 let json = json!({
1154 "algorithm": "m.megolm.v1.aes-sha2",
1155 "room_id": "!DovneieKSTkdHKpIXy:morpheus.localhost",
1156 "sender_key": "DeHIg4gwhClxzFYcmNntPNF9YtsdZbmMy8+3kzCMXHA",
1157 "session_id": "gM8i47Xhu0q52xLfgUXzanCMpLinoyVyH7R58cBuVBU",
1158 "session_key": "AQAAAABvWMNZjKFtebYIePKieQguozuoLgzeY6wKcyJjLJcJtQgy1dPqTBD12U+XrYLrRHn\
1159 lKmxoozlhFqJl456+9hlHCL+yq+6ScFuBHtJepnY1l2bdLb4T0JMDkNsNErkiLiLnD6yp3J\
1160 DSjIhkdHxmup/huygrmroq6/L5TaThEoqvW4DPIuO14btKudsS34FF82pwjKS4p6Mlch+0e\
1161 fHAblQV",
1162 "sender_claimed_keys":{},
1163 "forwarding_curve25519_key_chain":[]
1164 });
1165
1166 serde_json::from_value(json)
1167 .expect("We should be able to deserialize our exported room key")
1168 }
1169
1170 async fn backup_disabling_test_body(
1171 client: &Client,
1172 server: &MockServer,
1173 put_response: ResponseTemplate,
1174 ) {
1175 let _post_scope = Mock::given(method("POST"))
1176 .and(path("_matrix/client/unstable/room_keys/version"))
1177 .and(header("authorization", "Bearer 1234"))
1178 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
1179 "version": "1"
1180 })))
1181 .expect(1)
1182 .named("POST for the backup creation")
1183 .mount_as_scoped(server)
1184 .await;
1185
1186 let _put_scope = Mock::given(method("PUT"))
1187 .and(path("_matrix/client/unstable/room_keys/keys"))
1188 .and(header("authorization", "Bearer 1234"))
1189 .respond_with(put_response)
1190 .expect(1)
1191 .named("POST for the backup creation")
1192 .mount_as_scoped(server)
1193 .await;
1194
1195 client
1196 .encryption()
1197 .backups()
1198 .create()
1199 .await
1200 .expect("We should be able to create a new backup");
1201
1202 assert_eq!(client.encryption().backups().state(), BackupState::Enabled);
1203
1204 client
1205 .encryption()
1206 .backups()
1207 .backup_room_keys()
1208 .await
1209 .expect_err("Backups should be disabled");
1210
1211 assert_eq!(client.encryption().backups().state(), BackupState::Unknown);
1212 }
1213
1214 #[async_test]
1215 async fn test_resuming_backups_when_keys_are_consistent_makes_backups_enabled() {
1216 let server = MatrixMockServer::new().await;
1217 let client = server.client_builder().build().await;
1218 let backups = client.encryption().backups();
1219 let backup_decryption_key = BackupDecryptionKey::new();
1220
1221 let matching_public_key = derive_public_key_from(&backup_decryption_key);
1222
1223 server
1224 .mock_room_keys_version()
1225 .exists_with_key(&matching_public_key.to_base64())
1226 .expect(1)
1227 .mount()
1228 .await;
1229
1230 queue_backup_decryption_key_secret(client, &backup_decryption_key.to_base64()).await;
1233
1234 let res = backups.maybe_resume_backups().await;
1236
1237 assert_matches!(res, Ok(_));
1239
1240 assert_eq!(backups.state(), BackupState::Enabled);
1242 }
1243
1244 #[async_test]
1245 async fn test_resuming_backups_when_keys_are_inconsistent_has_no_effect() {
1246 let server = MatrixMockServer::new().await;
1254 let client = server.client_builder().build().await;
1255 let backups = client.encryption().backups();
1256 let backup_decryption_key = BackupDecryptionKey::new();
1257
1258 let non_matching_public_key = derive_public_key_from(&BackupDecryptionKey::new());
1259
1260 server
1261 .mock_room_keys_version()
1262 .exists_with_key(&non_matching_public_key.to_base64())
1263 .expect(1)
1264 .mount()
1265 .await;
1266
1267 queue_backup_decryption_key_secret(client, &backup_decryption_key.to_base64()).await;
1270
1271 let res = backups.maybe_resume_backups().await;
1273
1274 assert_matches!(res, Ok(_));
1276
1277 assert_eq!(backups.state(), BackupState::Unknown);
1280 }
1281
1282 #[async_test]
1283 async fn test_errors_when_resuming_backups_are_propagated() {
1284 let server = MatrixMockServer::new().await;
1285 let client = server.client_builder().build().await;
1286 let backups = client.encryption().backups();
1287
1288 queue_backup_decryption_key_secret(client, "not valid base64").await;
1291
1292 let res = backups.maybe_resume_backups().await;
1294
1295 assert_matches!(res, Err(Error::SerdeJson(_)));
1297
1298 assert_eq!(backups.state(), BackupState::Unknown);
1300 }
1301
1302 #[async_test]
1303 async fn test_backup_disabling_after_remote_deletion() {
1304 let server = MockServer::start().await;
1305 let client = logged_in_client(Some(server.uri())).await;
1306
1307 {
1308 let machine = client.olm_machine().await;
1309 machine
1310 .as_ref()
1311 .unwrap()
1312 .store()
1313 .import_exported_room_keys(vec![room_key()], |_, _| {})
1314 .await
1315 .expect("We should be able to import a room key");
1316 }
1317
1318 backup_disabling_test_body(
1319 &client,
1320 &server,
1321 ResponseTemplate::new(404).set_body_json(json!({
1322 "errcode": "M_NOT_FOUND",
1323 "error": "Unknown backup version"
1324 })),
1325 )
1326 .await;
1327
1328 backup_disabling_test_body(
1329 &client,
1330 &server,
1331 ResponseTemplate::new(403).set_body_json(json!({
1332 "current_version": "42",
1333 "errcode": "M_WRONG_ROOM_KEYS_VERSION",
1334 "error": "Wrong backup version."
1335 })),
1336 )
1337 .await;
1338
1339 server.verify().await;
1340 }
1341
1342 #[async_test]
1343 async fn test_when_a_backup_exists_then_fetch_exists_on_server_returns_true() {
1344 let server = MatrixMockServer::new().await;
1345 let client = server.client_builder().build().await;
1346
1347 server.mock_room_keys_version().exists().expect(1).mount().await;
1348
1349 let exists = client
1350 .encryption()
1351 .backups()
1352 .fetch_exists_on_server()
1353 .await
1354 .expect("We should be able to check if backups exist on the server");
1355
1356 assert!(exists, "We should deduce that a backup exists on the server");
1357 }
1358
1359 #[async_test]
1360 async fn test_repeated_calls_to_fetch_exists_on_server_makes_repeated_requests() {
1361 let server = MatrixMockServer::new().await;
1362 let client = server.client_builder().build().await;
1363
1364 server.mock_room_keys_version().exists().expect(2).mount().await;
1366
1367 let backups = client.encryption().backups();
1368
1369 backups.fetch_exists_on_server().await.unwrap();
1371 let exists = backups.fetch_exists_on_server().await.unwrap();
1372
1373 assert!(exists, "We should deduce that a backup exists on the server");
1374 }
1375
1376 #[async_test]
1377 async fn test_when_no_backup_exists_then_fetch_exists_on_server_returns_false() {
1378 let server = MatrixMockServer::new().await;
1379 let client = server.client_builder().build().await;
1380
1381 server.mock_room_keys_version().none().expect(1).mount().await;
1382
1383 let exists = client
1384 .encryption()
1385 .backups()
1386 .fetch_exists_on_server()
1387 .await
1388 .expect("We should be able to check if backups exist on the server");
1389
1390 assert!(!exists, "We should deduce that no backup exists on the server");
1391 }
1392
1393 #[async_test]
1394 async fn test_when_server_returns_an_error_then_fetch_exists_on_server_returns_an_error() {
1395 let server = MatrixMockServer::new().await;
1396 let client = server.client_builder().build().await;
1397
1398 {
1399 let _scope =
1400 server.mock_room_keys_version().error429().expect(1).mount_as_scoped().await;
1401
1402 client.encryption().backups().fetch_exists_on_server().await.expect_err(
1403 "If the /version endpoint returns a non 404 error we should throw an error",
1404 );
1405 }
1406
1407 {
1408 let _scope =
1409 server.mock_room_keys_version().error404().expect(1).mount_as_scoped().await;
1410
1411 client.encryption().backups().fetch_exists_on_server().await.expect_err(
1412 "If the /version endpoint returns a non-Matrix 404 error we should throw an error",
1413 );
1414 }
1415 }
1416
1417 #[async_test]
1418 async fn test_when_a_backup_exists_then_exists_on_server_returns_true() {
1419 let server = MatrixMockServer::new().await;
1420 let client = server.client_builder().build().await;
1421
1422 server.mock_room_keys_version().exists().expect(1).mount().await;
1423
1424 let exists = client
1425 .encryption()
1426 .backups()
1427 .exists_on_server()
1428 .await
1429 .expect("We should be able to check if backups exist on the server");
1430
1431 assert!(exists, "We should deduce that a backup exists on the server");
1432 }
1433
1434 #[async_test]
1435 async fn test_when_no_backup_exists_then_exists_on_server_returns_false() {
1436 let server = MatrixMockServer::new().await;
1437 let client = server.client_builder().build().await;
1438
1439 server.mock_room_keys_version().none().expect(1).mount().await;
1440
1441 let exists = client
1442 .encryption()
1443 .backups()
1444 .exists_on_server()
1445 .await
1446 .expect("We should be able to check if backups exist on the server");
1447
1448 assert!(!exists, "We should deduce that no backup exists on the server");
1449 }
1450
1451 #[async_test]
1452 async fn test_when_server_returns_an_error_then_exists_on_server_returns_an_error() {
1453 let server = MatrixMockServer::new().await;
1454 let client = server.client_builder().build().await;
1455
1456 {
1457 let _scope =
1458 server.mock_room_keys_version().error429().expect(1).mount_as_scoped().await;
1459
1460 client.encryption().backups().exists_on_server().await.expect_err(
1461 "If the /version endpoint returns a non 404 error we should throw an error",
1462 );
1463 }
1464
1465 {
1466 let _scope =
1467 server.mock_room_keys_version().error404().expect(1).mount_as_scoped().await;
1468
1469 client.encryption().backups().exists_on_server().await.expect_err(
1470 "If the /version endpoint returns a non-Matrix 404 error we should throw an error",
1471 );
1472 }
1473 }
1474
1475 #[async_test]
1476 async fn test_repeated_calls_to_exists_on_server_do_not_make_additional_requests() {
1477 let server = MatrixMockServer::new().await;
1478 let client = server.client_builder().build().await;
1479
1480 server.mock_room_keys_version().exists().expect(1).mount().await;
1482
1483 let backups = client.encryption().backups();
1484
1485 backups.exists_on_server().await.unwrap();
1487 backups.exists_on_server().await.unwrap();
1488 backups.exists_on_server().await.unwrap();
1489
1490 let exists = backups
1491 .exists_on_server()
1492 .await
1493 .expect("We should be able to check if backups exist on the server");
1494
1495 assert!(exists, "We should deduce that a backup exists on the server");
1496
1497 }
1499
1500 #[async_test]
1501 async fn test_adding_a_backup_invalidates_exists_on_server_cache() {
1502 let server = MatrixMockServer::new().await;
1503 let client = server.client_builder().build().await;
1504 let backups = client.encryption().backups();
1505
1506 {
1507 let _scope = server.mock_room_keys_version().none().expect(1).mount_as_scoped().await;
1508
1509 let exists = backups.exists_on_server().await.unwrap();
1511 assert!(!exists, "No backup exists at this point");
1512 }
1513
1514 server.mock_add_room_keys_version().ok().expect(1).mount().await;
1516 backups.create().await.expect("Failed to create a backup");
1517
1518 server.mock_room_keys_version().exists().expect(1).mount().await;
1519 let exists = backups
1520 .exists_on_server()
1521 .await
1522 .expect("We should be able to check if backups exist on the server");
1523
1524 assert!(exists, "But now a backup does exist");
1525 }
1526
1527 #[async_test]
1528 async fn test_removing_a_backup_invalidates_exists_on_server_cache() {
1529 let server = MatrixMockServer::new().await;
1530 let client = server.client_builder().build().await;
1531 let backups = client.encryption().backups();
1532
1533 {
1534 let _scope = server.mock_room_keys_version().exists().expect(1).mount_as_scoped().await;
1535
1536 let exists = backups.exists_on_server().await.unwrap();
1538 assert!(exists, "A backup exists at this point");
1539 }
1540
1541 server.mock_delete_room_keys_version().ok().expect(1).mount().await;
1543 backups.delete_backup_from_server("1".to_owned()).await.expect("Failed to delete a backup");
1544
1545 server.mock_room_keys_version().none().expect(1).mount().await;
1546 let exists = backups
1547 .exists_on_server()
1548 .await
1549 .expect("We should be able to check if backups exist on the server");
1550
1551 assert!(!exists, "But now there is no backup");
1552 }
1553
1554 #[async_test]
1555 async fn test_waiting_for_steady_state_resets_the_delay() {
1556 let server = MatrixMockServer::new().await;
1557 let client = server.client_builder().build().await;
1558
1559 server.mock_add_room_keys_version().ok().expect(1).mount().await;
1560
1561 client
1562 .encryption()
1563 .backups()
1564 .create()
1565 .await
1566 .expect("We should be able to create a new backup");
1567
1568 let backups = client.encryption().backups();
1569
1570 let old_duration =
1571 { client.inner.e2ee.backup_state.upload_delay.read().unwrap().to_owned() };
1572
1573 let wait_for_steady_state =
1574 backups.wait_for_steady_state().with_delay(Duration::from_nanos(100));
1575
1576 let mut progress_stream = wait_for_steady_state.subscribe_to_progress();
1577
1578 let task = matrix_sdk_common::executor::spawn({
1579 let client = client.to_owned();
1580 async move {
1581 while let Some(state) = progress_stream.next().await {
1582 let Ok(state) = state else {
1583 panic!("Error while waiting for the upload state")
1584 };
1585
1586 match state {
1587 UploadState::Idle => (),
1588 UploadState::Done => {
1589 let current_delay = {
1590 client
1591 .inner
1592 .e2ee
1593 .backup_state
1594 .upload_delay
1595 .read()
1596 .unwrap()
1597 .to_owned()
1598 };
1599
1600 assert_ne!(current_delay, old_duration);
1601 break;
1602 }
1603 _ => panic!("We should not have entered any other state"),
1604 }
1605 }
1606 }
1607 });
1608
1609 wait_for_steady_state.await.expect("We should be able to wait for the steady state");
1610 task.await.unwrap();
1611
1612 let current_duration =
1613 { client.inner.e2ee.backup_state.upload_delay.read().unwrap().to_owned() };
1614
1615 assert_eq!(old_duration, current_duration);
1616 }
1617
1618 fn derive_public_key_from(backup_decryption_key: &BackupDecryptionKey) -> Curve25519PublicKey {
1621 let backup_info = backup_decryption_key.to_backup_info();
1622 match backup_info {
1623 RoomKeyBackupInfo::MegolmBackupV1Curve25519AesSha2(megolm_v1_auth_data) => {
1624 megolm_v1_auth_data.public_key
1625 }
1626 RoomKeyBackupInfo::Other { .. } => {
1627 panic!("Unexpected backup info type")
1628 }
1629 }
1630 }
1631
1632 async fn queue_backup_decryption_key_secret(
1635 client: Client,
1636 secret_backup_decryption_key: &str,
1637 ) {
1638 let _guard = client.olm_machine().await;
1639 let machine = _guard.as_ref().unwrap();
1640 let transaction_id = TransactionId::new();
1641 let secret_info = SecretInfo::SecretRequest(SecretName::RecoveryKey);
1642 let user_id = machine.user_id().to_owned();
1643
1644 let gossip_request = GossipRequest {
1645 request_recipient: machine.user_id().to_owned(),
1646 request_id: transaction_id.clone(),
1647 info: secret_info.clone(),
1648 sent_out: true,
1649 };
1650
1651 let event = DecryptedSecretSendEvent {
1652 sender: user_id.clone(),
1653 recipient: user_id.clone(),
1654 keys: OlmV1Keys { ed25519: machine.identity_keys().ed25519 },
1655 recipient_keys: OlmV1Keys { ed25519: machine.identity_keys().ed25519 },
1656 sender_device_keys: None,
1657 content: SecretSendContent::new(
1658 transaction_id.to_owned(),
1659 secret_backup_decryption_key.to_owned(),
1660 ),
1661 };
1662
1663 let gossipped_secret =
1664 GossippedSecret { secret_name: SecretName::RecoveryKey, gossip_request, event };
1665
1666 let changes = Changes { secrets: vec![gossipped_secret.into()], ..Default::default() };
1667
1668 machine
1669 .store()
1670 .save_changes(changes)
1671 .await
1672 .expect("We should be able to import a room key");
1673 }
1674
1675 #[async_test]
1676 #[cfg(feature = "experimental-push-secrets")]
1677 async fn test_push_secret_on_create() {
1678 let server = MatrixMockServer::new().await;
1679 server.mock_add_room_keys_version().ok().mount().await;
1680 server.mock_crypto_endpoints_preset().await;
1681
1682 let client = server
1684 .client_builder_for_crypto_end_to_end(
1685 user_id!("@example:localhost"),
1686 device_id!("DEVICEID"),
1687 )
1688 .build()
1689 .await;
1690 let _other_client = server
1691 .set_up_new_device_for_encryption(&client, device_id!("OTHERDEVICEID"), vec![])
1692 .await;
1693
1694 client.encryption().bootstrap_cross_signing(None).await.unwrap();
1696 let other_device = client
1697 .encryption()
1698 .get_device(user_id!("@example:localhost"), device_id!("OTHERDEVICEID"))
1699 .await
1700 .unwrap()
1701 .unwrap();
1702 other_device.verify().await.unwrap();
1703 client.encryption().request_user_identity(user_id!("@example:localhost")).await.unwrap();
1704
1705 client
1707 .encryption()
1708 .backups()
1709 .create()
1710 .await
1711 .expect("We should be able to create a new backup");
1712
1713 let (_guard, to_device) =
1716 server.mock_capture_put_to_device(client.user_id().unwrap()).await;
1717 client.send_outgoing_requests().await.unwrap();
1718 to_device.await;
1719 }
1720}