Skip to main content

matrix_sdk/encryption/backups/
mod.rs

1// Copyright 2023 The Matrix.org Foundation C.I.C.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Room key backup support
16//!
17//! This module implements support for server-side key backups[[1]]. The module
18//! allows you to connect to an existing backup, create or delete backups from
19//! the homeserver, and download room keys from a backup.
20//!
21//! [1]: https://spec.matrix.org/unstable/client-server-api/#server-side-key-backups
22
23use 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/// The backups manager for the [`Client`].
67#[derive(Debug, Clone)]
68pub struct Backups {
69    pub(super) client: Client,
70}
71
72impl Backups {
73    /// Create a new backup version, encrypted with a new backup recovery key.
74    ///
75    /// The backup recovery key will be persisted locally and shared with
76    /// trusted devices as `m.secret.send` to-device messages.
77    ///
78    /// After the backup has been created, all room keys will be uploaded to the
79    /// homeserver.
80    ///
81    /// _Warning_: This will overwrite any existing backup.
82    ///
83    /// # Examples
84    ///
85    /// ```no_run
86    /// # use matrix_sdk::{Client, encryption::backups::BackupState};
87    /// # use url::Url;
88    /// # async {
89    /// # let homeserver = Url::parse("http://example.com")?;
90    /// # let client = Client::new(homeserver).await?;
91    /// let backups = client.encryption().backups();
92    /// backups.create().await?;
93    ///
94    /// assert_eq!(backups.state(), BackupState::Enabled);
95    /// # anyhow::Ok(()) };
96    /// ```
97    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        // Create a future so we can catch errors and go back to the `Unknown`
104        // state. This is a hack to get around the lack of `try` blocks in Rust.
105        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            // Create a new backup recovery key.
110            let decryption_key = BackupDecryptionKey::new();
111
112            // Get the info about the new backup key, this needs to be uploaded
113            // to the homeserver[1].
114            //
115            // We need to sign the `RoomKeyBackupInfo` so other clients which
116            // might want to start using the backup without having access to the
117            // `BackupDecryptionKey` can do so, as per [spec]:
118            //
119            // Clients must only store keys in backups after they have ensured
120            // that the `auth_data` has not been tampered with. This can be done
121            // either by:
122            //
123            // - checking that it is signed by the user's master cross-signing
124            //   key or by a verified device belonging to the same user, or
125            // - by deriving the public key from a private key that it obtained
126            //   from a trusted source. Trusted sources for the private key
127            //   include the user entering the key, retrieving the key stored in
128            //   secret storage, or obtaining the key via secret sharing from a
129            //   verified device belonging to the same user.
130            //
131            // [1]: https://spec.matrix.org/v1.8/client-server-api/#post_matrixclientv3room_keysversion
132            // [spec]: https://spec.matrix.org/v1.8/client-server-api/#server-side-key-backups
133            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            // Reset any state we might have had before the new backup was
145            // created. TODO: This should remove the old stored key and version.
146            olm_machine.backup_machine().disable_backup().await?;
147
148            let backup_key = decryption_key.megolm_v1_public_key();
149
150            // Save the newly created keys and the version we received from the
151            // server.
152            olm_machine
153                .backup_machine()
154                .save_decryption_key(Some(decryption_key), Some(version.to_owned()))
155                .await?;
156
157            // Enable the backup and start the upload of room keys.
158            self.enable(olm_machine, backup_key, version).await?;
159
160            #[cfg(feature = "experimental-push-secrets")]
161            {
162                // Push the backup key to our own verified devices.
163                // `push_secret_to_verified_devices` depends on having existing
164                // Olm sessions with the devices that the secret is being pushed
165                // to, so make sure we have Olm sessions with all our other
166                // devices.
167                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                // We can ignore errors here because the only way this function
176                // fails is if the secret is not found. But we saved the
177                // decryption key above, so this will never happen.
178                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    /// Disable and delete the currently active backup only if previously
194    /// enabled before, otherwise an error will be returned.
195    ///
196    /// For a more aggressive variant see [`Backups::disable_and_delete`] which
197    /// will delete the remote backup without checking the local state.
198    ///
199    /// # Examples
200    ///
201    /// ```no_run
202    /// # use matrix_sdk::{Client, encryption::backups::BackupState};
203    /// # use url::Url;
204    /// # async {
205    /// # let homeserver = Url::parse("http://example.com")?;
206    /// # let client = Client::new(homeserver).await?;
207    /// let backups = client.encryption().backups();
208    /// backups.disable().await?;
209    ///
210    /// assert_eq!(backups.state(), BackupState::Unknown);
211    /// # anyhow::Ok(()) };
212    /// ```
213    #[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        // Create a future so we can catch errors and go back to the `Unknown`
220        // state.
221        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    /// Completely disable and delete all backup versions, both locally and from
253    /// the server, no matter if they were previously setup locally or not.
254    ///
255    /// ⚠️ This method is mainly used when resetting the crypto identity and for
256    /// most other use cases its safer [`Backups::disable`] counterpart should
257    /// be used.
258    ///
259    /// It will fetch the current backup version from the server, delete it,
260    /// then repeat this until no backups remain, before proceeding to disabling
261    /// local backups as well
262    ///
263    /// # Examples
264    ///
265    /// ```no_run
266    /// # use matrix_sdk::{Client, encryption::backups::BackupState};
267    /// # use url::Url;
268    /// # async {
269    /// # let homeserver = Url::parse("http://example.com")?;
270    /// # let client = Client::new(homeserver).await?;
271    /// let backups = client.encryption().backups();
272    /// backups.disable_and_delete().await?;
273    ///
274    /// assert_eq!(backups.state(), BackupState::Unknown);
275    /// # anyhow::Ok(()) };
276    /// ```
277    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        // Create a future so we can catch errors and go back to the `Unknown`
283        // state.
284        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    /// Returns a future to wait for room keys to be uploaded.
305    ///
306    /// Awaiting the future will wake up a task to upload room keys which have
307    /// not yet been uploaded to the homeserver. It will then wait for the task
308    /// to finish uploading.
309    ///
310    /// # Examples
311    ///
312    /// ```no_run
313    /// # use matrix_sdk::{Client, encryption::backups::UploadState};
314    /// # use url::Url;
315    /// # async {
316    /// # let homeserver = Url::parse("http://example.com")?;
317    /// # let client = Client::new(homeserver).await?;
318    /// use futures_util::StreamExt;
319    ///
320    /// let backups = client.encryption().backups();
321    /// let wait_for_steady_state = backups.wait_for_steady_state();
322    ///
323    /// let mut progress_stream = wait_for_steady_state.subscribe_to_progress();
324    ///
325    /// tokio::spawn(async move {
326    ///     while let Some(update) = progress_stream.next().await {
327    ///         let Ok(update) = update else { break };
328    ///
329    ///         match update {
330    ///             UploadState::Uploading(counts) => {
331    ///                 println!(
332    ///                     "Uploaded {} out of {} room keys.",
333    ///                     counts.backed_up, counts.total
334    ///                 );
335    ///             }
336    ///             UploadState::Error => break,
337    ///             UploadState::Done => break,
338    ///             _ => (),
339    ///         }
340    ///     }
341    /// });
342    ///
343    /// wait_for_steady_state.await?;
344    ///
345    /// # anyhow::Ok(()) };
346    /// ```
347    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    /// Get a stream of updates to the [`BackupState`].
356    ///
357    /// This method will send out the current state as the first update.
358    ///
359    /// # Examples
360    ///
361    /// ```no_run
362    /// # use matrix_sdk::{Client, encryption::backups::BackupState};
363    /// # use url::Url;
364    /// # async {
365    /// # let homeserver = Url::parse("http://example.com")?;
366    /// # let client = Client::new(homeserver).await?;
367    /// use futures_util::StreamExt;
368    ///
369    /// let backups = client.encryption().backups();
370    ///
371    /// let mut state_stream = backups.state_stream();
372    ///
373    /// while let Some(update) = state_stream.next().await {
374    ///     let Ok(update) = update else { break };
375    ///
376    ///     match update {
377    ///         BackupState::Enabled => {
378    ///             println!("Backups have been enabled");
379    ///         }
380    ///         _ => (),
381    ///     }
382    /// }
383    /// # anyhow::Ok(()) };
384    /// ```
385    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    /// Get the current [`BackupState`] for this [`Client`].
392    pub fn state(&self) -> BackupState {
393        self.client.inner.e2ee.backup_state.global_state.get()
394    }
395
396    /// Are backups enabled for the current [`Client`]?
397    ///
398    /// This method will check if we locally have an active backup key and
399    /// backup version and are ready to upload room keys to a backup.
400    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    /// Does a backup exist on the server?
411    ///
412    /// This method will request info about the current backup from the
413    /// homeserver and if a backup exists return `true`, otherwise `false`.
414    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    /// Does a backup exist on the server?
421    ///
422    /// This method is identical to [`Self::fetch_exists_on_server`] except that
423    /// we cache the latest answer in memory and only empty the cache if the
424    /// local device adds or deletes a backup itself.
425    ///
426    /// Do not use this method if you need an accurate answer about whether a
427    /// backup exists - instead use [`Self::fetch_exists_on_server`]. This
428    /// method is useful when performance is more important than guaranteed
429    /// accuracy, such as when classifying UTDs.
430    pub async fn exists_on_server(&self) -> Result<bool, Error> {
431        // If we have an answer cached, return it immediately
432        if let Some(cached_value) = self.client.inner.e2ee.backup_state.backup_exists_on_server() {
433            return Ok(cached_value);
434        }
435
436        // Otherwise, delegate to fetch_exists_on_server. (It will update the
437        // cached value for us.)
438        self.fetch_exists_on_server().await
439    }
440
441    /// Subscribe to a stream that notifies when a room key for the specified
442    /// room is downloaded from the key backup.
443    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        // TODO: This is a bit crap to say the least. The type is
451        // non-descriptive and doesn't even contain all the important data. It
452        // should be a stream of `RoomKeyInfo` like the OlmMachine has... But on
453        // the other hand we should just be able to use the corresponding
454        // OlmMachine stream and remove this. Currently we can't do this because
455        // the OlmMachine gets destroyed and recreated all the time to be able
456        // to support the notifications-related multiprocessing on iOS.
457        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    /// Download all room keys for a certain room from the server-side key
470    /// backup.
471    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            // Transform response to standard format (map of room ID -> room
485            // key).
486            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    /// Download a single room key from the server-side key backup.
499    ///
500    /// Returns `true` if we managed to download a room key, `false` or an error
501    /// if we failed to download it. `false` indicates that there was no error,
502    /// we just don't have backups enabled so we can't download a room key.
503    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                // Transform response to standard format (map of room ID -> room
523                // key).
524                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    /// Set the state of the backup.
545    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    /// Set the backup state to the `Enabled` variant and insert the backup key
554    /// and version into the [`OlmMachine`].
555    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    /// Decrypt and forward a response containing backed up room keys to the
570    /// [`OlmMachine`].
571    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        // Since we can't use the usual room keys stream from the `OlmMachine`
619        // we're going to send things out in our own custom broadcaster.
620        let _ = self.client.inner.e2ee.backup_state.room_keys_broadcaster.send(result);
621
622        Ok(())
623    }
624
625    /// Download all room keys from the backup on the homeserver.
626    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    /// Get info about the currently active backup from the server.
649    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        // If the request succeeded, the backup is gone. If it failed, we are
681        // not really sure what the backup state is. Either way, clear the cache
682        // so we check next time we need to know.
683        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                            // TODO: If we're verified and there are other
736                            // devices besides us, request the new backup key
737                            // over `m.secret.send`.
738
739                            self.handle_deleted_backup_version(olm_machine).await?;
740                        }
741
742                        _ => (),
743                    }
744                }
745
746                Err(error.into())
747            }
748        }
749    }
750
751    /// Poll the [`OlmMachine`] for room keys which need to be backed up and
752    /// send out the request to the homeserver.
753    ///
754    /// This should only be called by the [`BackupUploadingTask`].
755    ///
756    /// [`BackupUploadingTask`]: crate::client::tasks::BackupUploadingTask
757    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    /// Set up a `m.secret.send` listener and re-enable backups if we have a
773    /// backup recovery key stored.
774    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    /// Try to enable backups with the given backup recovery key.
793    ///
794    /// This should be called if we receive a backup recovery, either:
795    ///
796    /// - As an `m.secret.send` to-device message from a trusted device.
797    /// - From 4S (i.e. from the `m.megolm_backup.v1` event global account
798    ///   data).
799    ///
800    /// In both cases the method will compare the currently active backup
801    /// version to the backup recovery key's version and, if there is a match,
802    /// activate backups on this device and start uploading room keys to the
803    /// backup.
804    ///
805    /// Returns true if backups were just enabled or were already enabled,
806    /// otherwise false.
807    #[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        // Create a future here which allows us to catch any failure that might
815        // happen so we can later on fall back to the correct `BackupState`.
816        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's try to see if there's a backup on the homeserver.
831            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", &current_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(&current_version.version)
844                && self.are_enabled().await
845            {
846                // If we already have a backup enabled which is using the
847                // currently active backup version, do nothing but tell the
848                // caller using the return value that backups are enabled.
849                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                // We're enabling a new backup, reset the `backed_up` flags on
857                // the room keys and remove any key/version we might have.
858                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                // Persist the new keys and enable the backup.
864                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 the user has set up the client to download any room keys,
873                // do so now. This is not really useful in a real scenario since
874                // the API to download room keys is not paginated.
875                //
876                // You need to download all room keys at once, parse a
877                // potentially huge JSON response and decrypt all the room keys
878                // found in the backup.
879                //
880                // This doesn't work for any sizeable account.
881                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                // Trigger the upload of any room keys we might need to upload.
894                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    /// Try to resume backups from a backup recovery key we have found in the
931    /// crypto store.
932    ///
933    /// Returns true if backups have been resumed, false otherwise.
934    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    /// Try to resume backups by iterating through the `m.secret.send` and
956    /// `io.element.msc4385.secret.push` to-device messages the [`OlmMachine`]
957    /// has received and stored in the secret inbox.
958    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                    // Ignore a bad backup decryption key here. We already
970                    // logged the details inside maybe_enable_backups().
971                }
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    /// Check and re-enable a backup if we have a backup recovery key locally.
982    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        // Let us first check if we have a stored backup recovery key and a
987        // backup version.
988        if !self.resume_backup_from_stored_backup_key(olm_machine).await? {
989            // We didn't manage to enable backups from a stored backup recovery
990            // key, let us check our secret inbox. Perhaps we can find a valid
991            // key there.
992            self.maybe_resume_from_secret_inbox(olm_machine).await?;
993        }
994
995        Ok(())
996    }
997
998    /// Listen for `m.secret.send` to-device messages and check the secret inbox
999    /// if we do receive one.
1000    #[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        // TODO: Because of our crude multi-process support, which reloads the
1005        // whole [`OlmMachine`] the `secrets_stream` might stop giving you
1006        // updates. Once that's fixed, stop listening to individual secret send
1007        // events and listen to the secrets stream.
1008        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    /// Listen for `io.element.msc4385.secret.push` to-device messages and check
1020    /// the pushed secret inbox if we do receive one.
1021    #[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        // TODO: As with `secret_send_event_handler`, because of our crude
1027        // multi-process support, which reloads the whole [`OlmMachine`] the
1028        // `secrets_stream` might stop giving you updates. Once that's fixed,
1029        // stop listening to individual secret push events and listen to the
1030        // secrets stream.
1031        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    /// Handle UTD events by triggering download from key backup.
1045    ///
1046    /// This function is registered as an event handler; it exists to deal with
1047    /// cases where [`Room::decrypt_event`] is not called and instead the event
1048    /// should be decrypted by the time this crate sees the event, such as for
1049    /// events received via `/sync` (as opposed to via `/messages`, `/context`,
1050    /// etc.)
1051    #[allow(clippy::unused_async)] // Because it's used as an event handler, which must be async.
1052    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    /// Send a notification to the task responsible for key backup downloads
1061    /// that it should attempt to download the keys for the given event.
1062    #[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    /// Send a notification to the task responsible for key backup downloads
1075    /// that it should attempt to download the keys for the given event.
1076    #[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    /// Send a notification to the task which is responsible for uploading room
1089    /// keys to the backup that it might have new room keys to back up.
1090    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    /// Disable our backups locally if we notice that the backup has been
1099    /// removed on the homeserver.
1100    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/// An error that happened while we were attempting to enable key backups.
1109#[derive(Debug, thiserror::Error)]
1110pub enum EnableBackupError {
1111    /// The private decryption key we found does not match the public key for
1112    /// the enabled backup.
1113    #[error("The backup decryption key does not match the latest backup version")]
1114    InconsistentBackupDecryptionKey,
1115
1116    /// A general error occurred while enabling key backup.
1117    #[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        // Given there is a backup decryption key waiting in the secrets inbox,
1231        // which is consistent with the public key
1232        queue_backup_decryption_key_secret(client, &backup_decryption_key.to_base64()).await;
1233
1234        // When we resume backups
1235        let res = backups.maybe_resume_backups().await;
1236
1237        // Then no error is returned
1238        assert_matches!(res, Ok(_));
1239
1240        // And the backup state is now enabled
1241        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        // Note: this was written when we added new error-surfacing logic to
1247        // matrix_sdk::encryption::backups::Backups::maybe_enable_backups, and
1248        // this test checks that the previous behaviour is preserved: we ignore
1249        // inconsistent backup keys when resuming backups. This may not turn out
1250        // to be the correct behaviour, so this test may need updating if we
1251        // decide this condition should be handled differently.
1252
1253        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        // Given there is a backup decryption key waiting in the secrets inbox,
1268        // but it doesn't match the public backup key
1269        queue_backup_decryption_key_secret(client, &backup_decryption_key.to_base64()).await;
1270
1271        // When we attempt to resume backups
1272        let res = backups.maybe_resume_backups().await;
1273
1274        // Then no error is returned ...
1275        assert_matches!(res, Ok(_));
1276
1277        // ... even though the backup setup failed because the decryption keys
1278        // were inconsistent
1279        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        // Given an invalid backup decryption key is waiting in the secrets
1289        // inbox
1290        queue_backup_decryption_key_secret(client, "not valid base64").await;
1291
1292        // When we attempt to resume backups
1293        let res = backups.maybe_resume_backups().await;
1294
1295        // Then an error is returned
1296        assert_matches!(res, Err(Error::SerdeJson(_)));
1297
1298        // And the backup state is unknown
1299        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        // Expect 2 requests to the server
1365        server.mock_room_keys_version().exists().expect(2).mount().await;
1366
1367        let backups = client.encryption().backups();
1368
1369        // Call fetch_exists_on_server twice
1370        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        // Create a mock stating that the request should only be made once
1481        server.mock_room_keys_version().exists().expect(1).mount().await;
1482
1483        let backups = client.encryption().backups();
1484
1485        // Call exists_on_server several times
1486        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        // We check expectations here, confirming that only one call was made
1498    }
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            // Call exists_on_server to fill the cache
1510            let exists = backups.exists_on_server().await.unwrap();
1511            assert!(!exists, "No backup exists at this point");
1512        }
1513
1514        // Create a new backup. Should invalidate the cache
1515        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            // Call exists_on_server to fill the cache
1537            let exists = backups.exists_on_server().await.unwrap();
1538            assert!(exists, "A backup exists at this point");
1539        }
1540
1541        // Delete the backup. Should invalidate the cache
1542        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    /// Given a private backup decryption key, return the matching public key
1619    /// for that backup
1620    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    /// Add a new secret to the secrets inbox of this client's OlmMachine
1633    /// containing the supplied backup decryption key.
1634    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        // Set up two devices for the test user
1683        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        // both devices are cross-signed
1695        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        // We create a new backup on one device
1706        client
1707            .encryption()
1708            .backups()
1709            .create()
1710            .await
1711            .expect("We should be able to create a new backup");
1712
1713        // which should result in pushing the key to our other device (i.e. a
1714        // to-device event should be sent)
1715        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}