Skip to main content

matrix_sdk/encryption/backups/
futures.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//! Named futures for the backup support.
16
17use std::{future::IntoFuture, time::Duration};
18
19use futures_core::Stream;
20use futures_util::StreamExt;
21use matrix_sdk_common::boxed_into_future;
22use thiserror::Error;
23use tokio_stream::wrappers::errors::BroadcastStreamRecvError;
24use tracing::trace;
25
26use super::{Backups, UploadState};
27use crate::utils::ChannelObservable;
28
29/// Error describing the ways that waiting for the backup upload to settle down
30/// can fail.
31#[derive(Clone, Copy, Debug, Error)]
32pub enum SteadyStateError {
33    /// The currently active backup got either deleted or a new one was created.
34    ///
35    /// No further room keys will be uploaded to the currently active backup.
36    #[error("The backup got disabled while waiting for the room keys to be uploaded.")]
37    BackupDisabled,
38    /// Uploading the room keys to the homeserver failed due to a network error.
39    ///
40    /// Uploading will be retried again at a later point in time, or immediately
41    /// if you wait for the steady state again.
42    #[error("There was a network connection error.")]
43    Connection,
44    /// We missed some updates to the [`UploadState`] from the upload task.
45    ///
46    /// This error doesn't imply that there was an error with the uploading of
47    /// room keys, it just means that we didn't receive all the transitions in
48    /// the [`UploadState`]. You might want to retry waiting for the steady
49    /// state.
50    #[error("We couldn't read status updates from the upload task quickly enough.")]
51    Lagged,
52}
53
54/// Named future for the [`Backups::wait_for_steady_state()`] method.
55#[derive(Debug)]
56pub struct WaitForSteadyState<'a> {
57    pub(super) backups: &'a Backups,
58    pub(super) progress: ChannelObservable<UploadState>,
59    pub(super) timeout: Option<Duration>,
60}
61
62impl WaitForSteadyState<'_> {
63    /// Subscribe to the progress of the backup upload step while waiting for it
64    /// to settle down.
65    pub fn subscribe_to_progress(
66        &self,
67    ) -> impl Stream<Item = Result<UploadState, BroadcastStreamRecvError>> + use<> {
68        self.progress.subscribe()
69    }
70
71    /// Set the delay between each upload request.
72    ///
73    /// Uploading room keys might require multiple requests to be sent out. The
74    /// [`Client`] waits for a while before it sends the next request out.
75    ///
76    /// This method allows you to override how long the [`Client`] will wait.
77    /// The default value is 100 ms.
78    ///
79    /// [`Client`]: crate::Client
80    pub fn with_delay(mut self, delay: Duration) -> Self {
81        self.timeout = Some(delay);
82
83        self
84    }
85}
86
87impl<'a> IntoFuture for WaitForSteadyState<'a> {
88    type Output = Result<(), SteadyStateError>;
89    boxed_into_future!(extra_bounds: 'a);
90
91    fn into_future(self) -> Self::IntoFuture {
92        Box::pin(async move {
93            let Self { backups, timeout, progress } = self;
94
95            trace!("Creating a stream to wait for the steady state");
96
97            let mut progress_stream = progress.subscribe();
98
99            let old_delay = if let Some(delay) = timeout {
100                let mut lock = backups.client.inner.e2ee.backup_state.upload_delay.write().unwrap();
101                let old_delay = Some(lock.to_owned());
102
103                *lock = delay;
104
105                old_delay
106            } else {
107                None
108            };
109
110            trace!("Waiting for the upload steady state");
111
112            let ret = if backups.are_enabled().await {
113                backups.maybe_trigger_backup();
114
115                let mut ret = Ok(());
116
117                // TODO: Do we want to be smart here and remember the count when
118                // we started waiting and prevent the total from increasing, in
119                // case new room keys arrive after we started waiting.
120                while let Some(state) = progress_stream.next().await {
121                    trace!(?state, "Update state while waiting for the backup steady state");
122
123                    match state {
124                        Ok(UploadState::Done) => {
125                            ret = Ok(());
126                            break;
127                        }
128                        Ok(UploadState::Error) => {
129                            if backups.are_enabled().await {
130                                ret = Err(SteadyStateError::Connection);
131                            } else {
132                                ret = Err(SteadyStateError::BackupDisabled);
133                            }
134
135                            break;
136                        }
137                        Err(_) => {
138                            ret = Err(SteadyStateError::Lagged);
139                            break;
140                        }
141                        _ => (),
142                    }
143                }
144
145                ret
146            } else {
147                Err(SteadyStateError::BackupDisabled)
148            };
149
150            if let Some(old_delay) = old_delay {
151                let mut lock = backups.client.inner.e2ee.backup_state.upload_delay.write().unwrap();
152                *lock = old_delay;
153            }
154
155            ret
156        })
157    }
158}