Skip to main content

matrix_sdk_ui/
encryption_sync_service.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 that specific language governing permissions and
13// limitations under the License.
14
15//! Encryption Sync API.
16//!
17//! The encryption sync API is a high-level helper that is designed to take care
18//! of handling the synchronization of encryption and to-device events (required
19//! for encryption), be they received within the app or within a dedicated
20//! extension process (e.g. the [NSE] process on iOS devices).
21//!
22//! Under the hood, this uses a sliding sync instance configured with no lists,
23//! but that enables the e2ee and to-device extensions, so that it can both
24//! handle encryption and manage encryption keys; that's sufficient to decrypt
25//! messages received in the notification processes.
26//!
27//! [NSE]: https://developer.apple.com/documentation/usernotifications/unnotificationserviceextension
28
29use std::{pin::Pin, time::Duration};
30
31use async_stream::stream;
32use futures_core::stream::Stream;
33use futures_util::{StreamExt, pin_mut};
34use matrix_sdk::{Client, LEASE_DURATION_MS, SlidingSync, sleep::sleep};
35use matrix_sdk_common::cross_process_lock::CrossProcessLockConfig;
36use ruma::{api::client::sync::sync_events::v5 as http, assign};
37use tokio::sync::OwnedMutexGuard;
38use tracing::{debug, instrument, trace};
39
40/// Unit type representing a permit to _use_ an [`EncryptionSyncService`].
41///
42/// This must be created once in the whole application's lifetime, wrapped in a
43/// mutex. Using an `EncryptionSyncService` must then lock that mutex in an
44/// owned way, so that there's at most a single `EncryptionSyncService` running
45/// at any time in the entire app.
46pub struct EncryptionSyncPermit(());
47
48impl EncryptionSyncPermit {
49    pub(crate) fn new() -> Self {
50        Self(())
51    }
52}
53
54impl EncryptionSyncPermit {
55    /// Test-only.
56    #[doc(hidden)]
57    pub fn new_for_testing() -> Self {
58        Self::new()
59    }
60}
61
62/// High-level helper for synchronizing encryption events using sliding sync.
63///
64/// See the module's documentation for more details.
65pub struct EncryptionSyncService {
66    client: Client,
67    sliding_sync: SlidingSync,
68}
69
70impl EncryptionSyncService {
71    /// Creates a new instance of a `EncryptionSyncService`.
72    ///
73    /// This will create and manage an instance of [`matrix_sdk::SlidingSync`].
74    pub async fn new(
75        client: Client,
76        poll_and_network_timeouts: Option<(Duration, Duration)>,
77    ) -> Result<Self, Error> {
78        // Make sure to use the same `conn_id` and caching store identifier,
79        // whichever process is running this sliding sync. There must be at most
80        // one sliding sync instance that enables the e2ee and to-device
81        // extensions.
82        let mut builder = client
83            .sliding_sync("encryption")
84            .map_err(Error::SlidingSync)?
85            //.share_pos() // TODO: This is racy, needs cross-process lock :')
86            .with_to_device_extension(
87                assign!(http::request::ToDevice::default(), { enabled: Some(true)}),
88            )
89            .with_e2ee_extension(assign!(http::request::E2EE::default(), { enabled: Some(true)}));
90
91        if let Some((poll_timeout, network_timeout)) = poll_and_network_timeouts {
92            builder = builder.poll_timeout(poll_timeout).network_timeout(network_timeout);
93        }
94
95        let sliding_sync = builder.build().await.map_err(Error::SlidingSync)?;
96
97        if let CrossProcessLockConfig::MultiProcess { holder_name } =
98            client.cross_process_lock_config()
99        {
100            // Gently try to enable the cross-process lock on behalf of the
101            // user.
102            match client.encryption().enable_cross_process_store_lock(holder_name.clone()).await {
103                Ok(()) | Err(matrix_sdk::Error::BadCryptoStoreState) => {
104                    // Ignore; we've already set the crypto store lock to
105                    // something, and that's sufficient as long as it uniquely
106                    // identifies the process.
107                }
108                Err(err) => {
109                    // Any other error is fatal
110                    return Err(Error::ClientError(err));
111                }
112            }
113        }
114
115        Ok(Self { client, sliding_sync })
116    }
117
118    /// Runs an `EncryptionSyncService` loop, yielding `Ok(())` after each
119    /// iteration so the caller can decide whether to continue or stop by
120    /// dropping the stream.
121    ///
122    /// Ends without yielding if the cross-process lock is configured but can't
123    /// be acquired (another process is expected to run the sync).
124    ///
125    /// Note: the [`EncryptionSyncPermit`] parameter ensures that there's at
126    /// most one encryption sync running at any time. See its documentation for
127    /// more details.
128    pub fn run_iterations(
129        self,
130        permit: OwnedMutexGuard<EncryptionSyncPermit>,
131    ) -> impl Stream<Item = Result<(), Error>> {
132        stream!({
133            // Move the permit into the stream, so that it's held for as long as
134            // the stream is alive.
135            let _permit = permit;
136
137            let _lock_guard = if let CrossProcessLockConfig::MultiProcess { .. } =
138                self.client.cross_process_lock_config()
139            {
140                let mut lock_guard = match self.client.encryption().try_lock_store_once().await {
141                    Ok(lock_guard) => lock_guard,
142                    Err(err) => {
143                        yield Err(Error::LockError(err));
144                        return;
145                    }
146                };
147
148                // Try to take the lock at the beginning; if it's busy, that
149                // means that another process already holds onto it, and as such
150                // we won't try to run the encryption sync loop at all (because
151                // we expect the other process to do so).
152
153                if lock_guard.is_none() {
154                    tracing::debug!(
155                        "Lock was already taken, and we're not the main loop; retrying in {}ms...",
156                        LEASE_DURATION_MS
157                    );
158
159                    sleep(Duration::from_millis(LEASE_DURATION_MS.into())).await;
160
161                    lock_guard = match self.client.encryption().try_lock_store_once().await {
162                        Ok(lock_guard) => lock_guard,
163                        Err(err) => {
164                            yield Err(Error::LockError(err));
165                            return;
166                        }
167                    };
168
169                    if lock_guard.is_none() {
170                        tracing::debug!(
171                            "Second attempt at locking outside the main app failed, aborting."
172                        );
173                        return;
174                    }
175                }
176
177                lock_guard
178            } else {
179                None
180            };
181
182            let sync = self.sliding_sync.sync();
183
184            pin_mut!(sync);
185
186            loop {
187                match sync.next().await {
188                    Some(Ok(update_summary)) => {
189                        // This API is only concerned with the e2ee and
190                        // to-device extensions. Warn if anything weird has been
191                        // received from the homeserver.
192                        if !update_summary.lists.is_empty() {
193                            debug!(?update_summary.lists, "unexpected non-empty list of lists in encryption sync API");
194                        }
195                        if !update_summary.rooms.is_empty() {
196                            debug!(?update_summary.rooms, "unexpected non-empty list of rooms in encryption sync API");
197                        }
198
199                        // Cool cool, let's do it again.
200                        trace!("Encryption sync received an update!");
201                        yield Ok(());
202                    }
203
204                    Some(Err(err)) => {
205                        trace!("Encryption sync stopped because of an error: {err:#}");
206                        yield Err(Error::SlidingSync(err));
207                        break;
208                    }
209
210                    None => {
211                        trace!("Encryption sync properly terminated.");
212                        break;
213                    }
214                }
215            }
216        })
217    }
218
219    /// Start synchronization.
220    ///
221    /// This should be regularly polled.
222    ///
223    /// Note: the [`EncryptionSyncPermit`] parameter ensures that there's at
224    /// most one encryption sync running at any time. See its documentation for
225    /// more details.
226    #[doc(hidden)] // Only public for testing purposes.
227    pub fn sync(
228        &self,
229        permit: OwnedMutexGuard<EncryptionSyncPermit>,
230    ) -> impl Stream<Item = Result<(), Error>> + '_ {
231        stream!({
232            // Move the permit into the stream, so that it's held for as long as
233            // the stream is alive.
234            let _permit = permit;
235
236            let sync = self.sliding_sync.sync();
237
238            pin_mut!(sync);
239
240            loop {
241                match self.next_sync_with_lock(&mut sync).await? {
242                    Some(Ok(update_summary)) => {
243                        // This API is only concerned with the e2ee and
244                        // to-device extensions. Warn if anything weird has been
245                        // received from the homeserver.
246                        if !update_summary.lists.is_empty() {
247                            debug!(?update_summary.lists, "unexpected non-empty list of lists in encryption sync API");
248                        }
249                        if !update_summary.rooms.is_empty() {
250                            debug!(?update_summary.rooms, "unexpected non-empty list of rooms in encryption sync API");
251                        }
252
253                        // Cool cool, let's do it again.
254                        trace!("Encryption sync received an update!");
255                        yield Ok(());
256                        continue;
257                    }
258
259                    Some(Err(err)) => {
260                        trace!("Encryption sync stopped because of an error: {err:#}");
261                        yield Err(Error::SlidingSync(err));
262                        break;
263                    }
264
265                    None => {
266                        trace!("Encryption sync properly terminated.");
267                        break;
268                    }
269                }
270            }
271        })
272    }
273
274    /// Helper function for `sync`. Take the cross-process store lock, and call
275    /// `sync.next()`
276    #[instrument(skip_all)]
277    async fn next_sync_with_lock<Item>(
278        &self,
279        sync: &mut Pin<&mut impl Stream<Item = Item>>,
280    ) -> Result<Option<Item>, Error> {
281        let _guard = if let CrossProcessLockConfig::MultiProcess { .. } =
282            self.client.cross_process_lock_config()
283        {
284            self.client.encryption().spin_lock_store(Some(60000)).await.map_err(Error::LockError)?
285        } else {
286            None
287        };
288
289        Ok(sync.next().await)
290    }
291
292    /// Requests that the underlying sliding sync be stopped.
293    ///
294    /// This will unlock the cross-process lock, if taken.
295    pub(crate) fn stop_sync(&self) -> Result<(), Error> {
296        // Stopping the sync loop will cause the next `next()` call to return
297        // `None`, so this will also release the cross-process lock
298        // automatically.
299        self.sliding_sync.stop_sync().map_err(Error::SlidingSync)?;
300
301        Ok(())
302    }
303
304    pub(crate) async fn expire_sync_session(&self) {
305        self.sliding_sync.expire_session().await;
306    }
307}
308
309/// Errors for the [`EncryptionSyncService`].
310#[derive(Debug, thiserror::Error)]
311pub enum Error {
312    #[error("Something wrong happened in sliding sync: {0:#}")]
313    SlidingSync(matrix_sdk::Error),
314
315    #[error("Locking failed: {0:#}")]
316    LockError(matrix_sdk::Error),
317
318    #[error(transparent)]
319    ClientError(matrix_sdk::Error),
320}