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}