matrix_sdk/encryption/dehydrated_devices.rs
1// Copyright 2026 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//! High-level interface for [Dehydrated Devices] ([MSC3814]).
16//!
17//! A dehydrated device is a virtual device the homeserver keeps on the user's
18//! behalf while no live device is online. Senders can encrypt to it using the
19//! same Olm session establishment as for any other device. When a new device
20//! comes online it rehydrates: it pulls the private parts of the virtual
21//! device back down, decrypts them with a pickle key, drains the queued
22//! to-device events, and imports the room keys they carry.
23//!
24//! # Lifecycle
25//!
26//! 1. `is_supported`: cheap probe of the homeserver.
27//! 2. `create`: build a fresh dehydrated device and upload it. The pickle key
28//! is supplied by the caller; storage and rotation of the pickle key are an
29//! application concern.
30//! 3. `rehydrate`: pull the existing dehydrated device, decrypt with the pickle
31//! key, absorb queued to-device events, and delete the device.
32//! 4. `delete`: remove the current dehydrated device without rehydrating.
33//!
34//! # Example
35//!
36//! ```no_run
37//! # use matrix_sdk::Client;
38//! # use matrix_sdk_base::crypto::store::types::DehydratedDeviceKey;
39//! # async fn example(client: Client, pickle_key: DehydratedDeviceKey)
40//! # -> anyhow::Result<()> {
41//! let dehydrated = client.encryption().dehydrated_devices();
42//!
43//! if !dehydrated.is_supported().await? {
44//! return Ok(());
45//! }
46//!
47//! // The pickle key comes from Secret Storage: only the key that encrypted
48//! // the existing device can rehydrate it. `start` manages that round trip.
49//! dehydrated.rehydrate(&pickle_key).await?;
50//! dehydrated.create(None, &pickle_key).await?;
51//! # Ok(())
52//! # }
53//! ```
54//!
55//! [Dehydrated Devices]: https://spec.matrix.org/unstable/client-server-api/#dehydrated-devices
56//! [MSC3814]: https://github.com/matrix-org/matrix-spec-proposals/pull/3814
57
58use std::{future::IntoFuture, time::Duration};
59
60use futures_core::Stream;
61use matrix_sdk_base::crypto::{
62 OlmError,
63 dehydrated_devices::{DehydrationError, RehydratedDevice},
64 store::types::DehydratedDeviceKey,
65 vodozemac::base64_decode,
66};
67use matrix_sdk_common::{
68 boxed_into_future, locks::Mutex as StdMutex, sleep::sleep, task_monitor::BackgroundTaskHandle,
69};
70use ruma::{
71 OwnedDeviceId,
72 api::{
73 client::dehydrated_device::{
74 DehydratedDeviceData, delete_dehydrated_device, get_dehydrated_device, get_events,
75 },
76 error::ErrorKind,
77 },
78 events::secret::request::SecretName,
79 serde::Raw,
80};
81use thiserror::Error;
82use tokio::sync::broadcast;
83use tokio_stream::wrappers::{BroadcastStream, errors::BroadcastStreamRecvError};
84use tracing::{Instrument, Span, debug, info, instrument, trace, warn};
85use zeroize::Zeroizing;
86
87use crate::{
88 Client, HttpError,
89 client::WeakClient,
90 encryption::{CryptoStoreError, secret_storage::SecretStore},
91};
92
93/// The default display name uploaded for a freshly created dehydrated device.
94const DEFAULT_DEVICE_DISPLAY_NAME: &str = "Dehydrated device";
95
96/// The name used to store the dehydrated-device pickle key in Secret Storage.
97///
98/// MSC3814 reserves `m.dehydrated_device` for the stable name; this is the
99/// unstable equivalent the implementation will publish until the MSC
100/// stabilizes.
101const PICKLE_KEY_SECRET_NAME: &str = "org.matrix.msc3814";
102
103/// How often the rotation task started by [`DehydratedDevices::start`]
104/// re-creates the dehydrated device. Fixed at one week.
105const DEHYDRATION_INTERVAL: Duration = Duration::from_secs(7 * 24 * 60 * 60);
106
107/// Crypto-store key under which the ID of the most recently uploaded dehydrated
108/// device is persisted, so the replay check in [`DehydratedDevices::rehydrate`]
109/// survives a client restart.
110const LAST_UPLOADED_DEVICE_ID_KEY: &str = "matrix-sdk-dehydrated-devices.last-uploaded-device-id";
111
112/// Defensive ceiling on the number of to-device events drained from a single
113/// dehydrated device during rehydration. MSC3814 does not specify a limit; this
114/// bounds the drain loop so a server that keeps returning events cannot make it
115/// run forever, and keeps the running count well clear of overflow.
116const MAX_TO_DEVICE_EVENTS: usize = 100_000;
117
118/// Errors that can occur while managing dehydrated devices.
119#[derive(Debug, Error)]
120pub enum DehydratedDeviceError {
121 /// The HTTP request to the homeserver failed.
122 #[error(transparent)]
123 Http(#[from] HttpError),
124
125 /// The cryptographic operation on the dehydrated device failed.
126 #[error(transparent)]
127 Crypto(#[from] DehydrationError),
128
129 /// Importing room keys from a rehydrated device's to-device events failed.
130 #[error(transparent)]
131 Olm(#[from] OlmError),
132
133 /// The to-device drain during rehydration stopped before the server's queue
134 /// was exhausted; the dehydrated device was left in place so a retry can
135 /// resume the drain.
136 #[error(
137 "the to-device drain stopped after {to_device_events} events with more still queued; the dehydrated device was kept so a retry can resume"
138 )]
139 DrainTruncated {
140 /// Number of to-device events processed before the drain stopped.
141 to_device_events: usize,
142 },
143
144 /// The crypto store could not be accessed.
145 #[error(transparent)]
146 Store(#[from] CryptoStoreError),
147
148 /// Reading or writing a secret to Secret Storage failed.
149 #[error(transparent)]
150 SecretStorage(#[from] crate::encryption::secret_storage::SecretStorageError),
151
152 /// The pickle key stored in Secret Storage was not valid base64.
153 #[error("the dehydrated-device pickle key in Secret Storage is not valid base64: {0}")]
154 PickleKeyDecode(#[from] vodozemac::Base64DecodeError),
155
156 /// The client is not logged in; the Olm machine is not available.
157 #[error("the client is not logged in")]
158 NotLoggedIn,
159}
160
161/// Return a [`SecretName`] for the dehydrated-device pickle key entry.
162fn pickle_key_secret_name() -> SecretName {
163 SecretName::from(PICKLE_KEY_SECRET_NAME)
164}
165
166/// Lifecycle events emitted by [`DehydratedDevices`].
167///
168/// Subscribe with [`DehydratedDevices::state_stream`] to observe creation,
169/// rehydration progress, and rotation outcomes. [`Self::RehydrationCompleted`]
170/// carries the final imported counts so a caller does not have to fold over the
171/// [`Self::RehydrationProgress`] events, and [`Self::RotationError`] surfaces
172/// background rotation failures the task would otherwise swallow.
173#[derive(Clone, Debug)]
174pub enum DehydratedDeviceEvent {
175 /// A fresh dehydrated device was constructed in the local crypto store,
176 /// before the upload PUT.
177 Created {
178 /// Device ID assigned to the new dehydrated device.
179 device_id: OwnedDeviceId,
180 },
181 /// The dehydrated device announced by the preceding [`Self::Created`] event
182 /// was accepted by the homeserver.
183 Uploaded {
184 /// Device ID of the dehydrated device now visible on the server.
185 device_id: OwnedDeviceId,
186 },
187 /// The dehydrated device currently on the server was deleted.
188 Deleted,
189 /// A pickle key was cached in the local crypto store.
190 KeyCached,
191 /// Rehydration of a dehydrated device began.
192 RehydrationStarted {
193 /// Device ID of the dehydrated device being rehydrated.
194 device_id: OwnedDeviceId,
195 },
196 /// A batch of to-device events has been imported during rehydration.
197 RehydrationProgress {
198 /// Cumulative number of room keys imported so far.
199 room_keys_imported: usize,
200 /// Cumulative number of to-device events processed so far.
201 to_device_events: usize,
202 },
203 /// Rehydration finished successfully.
204 RehydrationCompleted {
205 /// Device ID of the rehydrated device.
206 device_id: OwnedDeviceId,
207 /// Total number of room keys imported.
208 room_keys_imported: usize,
209 /// Total number of to-device events processed.
210 to_device_events: usize,
211 },
212 /// Rehydration failed before it could complete.
213 RehydrationError {
214 /// Human-readable description of the failure.
215 error: String,
216 },
217 /// A scheduled rotation tick failed; the rotation task remains scheduled
218 /// and will retry at the next tick.
219 RotationError {
220 /// Human-readable description of the failure.
221 error: String,
222 },
223}
224
225/// Process-wide state for the dehydrated-devices manager.
226///
227/// Held inside [`crate::encryption::EncryptionData`] so the event sender and
228/// any in-flight rotation task survive across
229/// `Client::encryption().dehydrated_devices()` calls.
230pub(crate) struct DehydratedDevicesState {
231 event_sender: broadcast::Sender<DehydratedDeviceEvent>,
232 rotation_task: StdMutex<Option<BackgroundTaskHandle>>,
233}
234
235impl Default for DehydratedDevicesState {
236 fn default() -> Self {
237 let (event_sender, _) = broadcast::channel(100);
238 Self { event_sender, rotation_task: StdMutex::new(None) }
239 }
240}
241
242/// High-level handle returned by
243/// [`Encryption::dehydrated_devices`](crate::encryption::Encryption::dehydrated_devices).
244#[derive(Debug, Clone)]
245pub struct DehydratedDevices {
246 pub(super) client: Client,
247}
248
249/// The dehydrated device the server currently holds on the user's behalf.
250struct DownloadedDevice {
251 device_id: OwnedDeviceId,
252 device_data: Raw<DehydratedDeviceData>,
253}
254
255/// Outcome of draining the dehydrated device's to-device queue.
256struct DrainOutcome {
257 room_keys_imported: usize,
258 to_device_events: usize,
259
260 /// Whether the drain stopped defensively (batch cap or a repeated cursor)
261 /// with events potentially still queued on the server.
262 truncated: bool,
263}
264
265impl DehydratedDevices {
266 /// Subscribe to the stream of [`DehydratedDeviceEvent`]s.
267 ///
268 /// Each call returns a fresh stream. If a subscriber is slow enough to fall
269 /// behind the channel's buffer, it receives a [`BroadcastStreamRecvError`]
270 /// reporting the number of skipped events and the stream continues from the
271 /// most recent event.
272 ///
273 /// # Example
274 ///
275 /// ```no_run
276 /// # use matrix_sdk::Client;
277 /// # use futures_util::StreamExt;
278 /// # async fn example(client: Client) -> anyhow::Result<()> {
279 /// let dehydrated = client.encryption().dehydrated_devices();
280 /// let mut stream = dehydrated.state_stream();
281 /// while let Some(Ok(event)) = stream.next().await {
282 /// println!("dehydrated devices: {event:?}");
283 /// }
284 /// # Ok(()) }
285 /// ```
286 pub fn state_stream(
287 &self,
288 ) -> impl Stream<Item = Result<DehydratedDeviceEvent, BroadcastStreamRecvError>> + use<> {
289 BroadcastStream::new(self.state().event_sender.subscribe())
290 }
291
292 fn state(&self) -> &DehydratedDevicesState {
293 &self.client.inner.e2ee.dehydrated_devices_state
294 }
295
296 fn emit(&self, event: DehydratedDeviceEvent) {
297 // A send failure means there are no subscribers; that is fine.
298 let _ = self.state().event_sender.send(event);
299 }
300
301 /// Return whether the homeserver advertises dehydrated-device support.
302 ///
303 /// Probes by issuing `GET /dehydrated_device` and inspecting the errcode of
304 /// the response:
305 ///
306 /// - `M_UNRECOGNIZED` means the server does not understand the endpoint.
307 /// - `M_NOT_FOUND` or a successful response means the server understands
308 /// the endpoint (whether or not the user currently has a dehydrated
309 /// device on file).
310 ///
311 /// Any other transport or API failure is propagated.
312 #[instrument(skip_all)]
313 pub async fn is_supported(&self) -> Result<bool, DehydratedDeviceError> {
314 let request = get_dehydrated_device::unstable::Request::new();
315 match self.client.send(request).await {
316 Ok(_) => Ok(true),
317 Err(e) => match e.client_api_error_kind() {
318 Some(ErrorKind::Unrecognized) => Ok(false),
319 Some(ErrorKind::NotFound) => Ok(true),
320 _ => Err(e.into()),
321 },
322 }
323 }
324
325 /// Create a fresh dehydrated device and upload it to the homeserver.
326 ///
327 /// The pickle key is used by [vodozemac] to encrypt the private parts of
328 /// the device. The application is responsible for safely storing the pickle
329 /// key (typically in Secret Storage so future sessions can rehydrate the
330 /// device).
331 ///
332 /// # Arguments
333 ///
334 /// - `display_name` - Optional human-readable name uploaded as the
335 /// dehydrated device's `initial_device_display_name`. Defaults to
336 /// `"Dehydrated device"`.
337 /// - `pickle_key` - 32-byte key used to encrypt the dehydrated device.
338 ///
339 /// # Example
340 ///
341 /// ```no_run
342 /// # use matrix_sdk::Client;
343 /// # use matrix_sdk_base::crypto::store::types::DehydratedDeviceKey;
344 /// # async fn example(client: Client) -> anyhow::Result<()> {
345 /// let pickle_key = DehydratedDeviceKey::new();
346 /// let device_id = client
347 /// .encryption()
348 /// .dehydrated_devices()
349 /// .create(Some("Offline catcher"), &pickle_key)
350 /// .await?;
351 /// println!("Uploaded dehydrated device {device_id}");
352 /// # Ok(()) }
353 /// ```
354 ///
355 /// [vodozemac]: https://docs.rs/vodozemac/
356 #[instrument(skip_all)]
357 pub async fn create(
358 &self,
359 display_name: Option<&str>,
360 pickle_key: &DehydratedDeviceKey,
361 ) -> Result<OwnedDeviceId, DehydratedDeviceError> {
362 let olm = self.client.olm_machine().await;
363 let machine = olm.as_ref().ok_or(DehydratedDeviceError::NotLoggedIn)?;
364
365 debug!("Creating a new dehydrated device in the crypto store");
366 let dehydrated_device = machine.dehydrated_devices().create().await?;
367
368 let display_name = display_name.unwrap_or(DEFAULT_DEVICE_DISPLAY_NAME);
369 let request =
370 dehydrated_device.keys_for_upload(display_name.to_owned(), pickle_key).await?;
371 let device_id = request.device_id.clone();
372 self.emit(DehydratedDeviceEvent::Created { device_id: device_id.clone() });
373
374 debug!(?device_id, "Uploading dehydrated device to the homeserver");
375 self.client.send(request).await?;
376 info!(?device_id, "Successfully uploaded dehydrated device");
377
378 machine.store().set_value(LAST_UPLOADED_DEVICE_ID_KEY, &device_id).await?;
379 self.emit(DehydratedDeviceEvent::Uploaded { device_id: device_id.clone() });
380 Ok(device_id)
381 }
382
383 /// Rehydrate the dehydrated device currently on the server, if any.
384 ///
385 /// Downloads the dehydrated device, decrypts it with `pickle_key`, drains
386 /// all queued to-device events to import their room keys, and finally
387 /// deletes the device from the server.
388 ///
389 /// Returns `Ok(false)` if the server reports no dehydrated device
390 /// (`M_NOT_FOUND`) or does not implement the endpoint (`M_UNRECOGNIZED`).
391 /// Returns `Ok(true)` once the rehydration cycle has completed end to end.
392 ///
393 /// If the drain stops defensively before the server's queue is exhausted,
394 /// the device is left on the server and
395 /// [`DehydratedDeviceError::DrainTruncated`] is returned. Calling this
396 /// method again restarts the drain from the beginning of the queue; room
397 /// keys imported by the earlier attempt import idempotently.
398 ///
399 /// # Example
400 ///
401 /// ```no_run
402 /// # use matrix_sdk::Client;
403 /// # use matrix_sdk_base::crypto::store::types::DehydratedDeviceKey;
404 /// # async fn example(client: Client, pickle_key: DehydratedDeviceKey)
405 /// # -> anyhow::Result<()> {
406 /// let rehydrated =
407 /// client.encryption().dehydrated_devices().rehydrate(&pickle_key).await?;
408 /// if rehydrated {
409 /// println!("Caught up on offline room keys");
410 /// }
411 /// # Ok(()) }
412 /// ```
413 #[instrument(skip_all)]
414 pub async fn rehydrate(
415 &self,
416 pickle_key: &DehydratedDeviceKey,
417 ) -> Result<bool, DehydratedDeviceError> {
418 let Some(downloaded) = self.download_device().await? else { return Ok(false) };
419 info!(device_id = ?downloaded.device_id, "Dehydrated device found");
420
421 // The server should serve back whichever id we last uploaded. A
422 // mismatch can indicate a stale or replayed payload; warn so the
423 // application can decide to drop or quarantine the imported keys. The
424 // last-uploaded id is persisted in the crypto store, so this check
425 // holds across restarts.
426 if let Some(expected) = self.last_uploaded_device_id().await?
427 && expected != downloaded.device_id
428 {
429 warn!(
430 ?expected,
431 got = ?downloaded.device_id,
432 "Server returned a different dehydrated-device id than the one we last uploaded; continuing but the payload may be stale"
433 );
434 }
435
436 self.emit(DehydratedDeviceEvent::RehydrationStarted {
437 device_id: downloaded.device_id.clone(),
438 });
439
440 let rehydrated = self.rehydrate_device(&downloaded, pickle_key).await?;
441 let drained = self.absorb_events(&downloaded.device_id, &rehydrated).await?;
442
443 // The server keeps undelivered events queued until the device is
444 // deleted, so a truncated drain must leave the device in place for a
445 // retry to resume; deleting now would discard the rest of the queue.
446 if drained.truncated {
447 return Err(DehydratedDeviceError::DrainTruncated {
448 to_device_events: drained.to_device_events,
449 });
450 }
451
452 self.emit(DehydratedDeviceEvent::RehydrationCompleted {
453 device_id: downloaded.device_id.clone(),
454 room_keys_imported: drained.room_keys_imported,
455 to_device_events: drained.to_device_events,
456 });
457
458 // Key import already succeeded; if the post-drain delete fails, log it
459 // but do not let the failure masquerade as a rehydration error. The
460 // next create() call will replace the device anyway.
461 if let Err(e) = self.delete_device().await {
462 warn!(device_id = ?downloaded.device_id, error = %e, "Post-rehydration delete failed; the next rotation will replace the device");
463 }
464
465 Ok(true)
466 }
467
468 /// Cache the pickle key in the local crypto store.
469 ///
470 /// Subsequent rehydration attempts can then resolve the key from the cache
471 /// without an account-data round-trip.
472 #[instrument(skip_all)]
473 pub(crate) async fn cache_key(
474 &self,
475 pickle_key: &DehydratedDeviceKey,
476 ) -> Result<(), DehydratedDeviceError> {
477 let olm = self.client.olm_machine().await;
478 let machine = olm.as_ref().ok_or(DehydratedDeviceError::NotLoggedIn)?;
479
480 machine.dehydrated_devices().save_dehydrated_device_pickle_key(pickle_key).await?;
481 self.emit(DehydratedDeviceEvent::KeyCached);
482 Ok(())
483 }
484
485 /// Return the pickle key currently cached in the local crypto store.
486 ///
487 /// `Ok(None)` if no key has been cached. The returned key matches the last
488 /// value persisted via [`cache_key`](Self::cache_key) or
489 /// [`reset_key`](Self::reset_key); it is not fetched from Secret Storage.
490 #[instrument(skip_all)]
491 pub(crate) async fn cached_key(
492 &self,
493 ) -> Result<Option<DehydratedDeviceKey>, DehydratedDeviceError> {
494 let olm = self.client.olm_machine().await;
495 let machine = olm.as_ref().ok_or(DehydratedDeviceError::NotLoggedIn)?;
496
497 Ok(machine.dehydrated_devices().get_dehydrated_device_pickle_key().await?)
498 }
499
500 /// Return the ID of the dehydrated device this client most recently
501 /// uploaded, as persisted in the crypto store, or `Ok(None)` if none has
502 /// been uploaded yet.
503 async fn last_uploaded_device_id(
504 &self,
505 ) -> Result<Option<OwnedDeviceId>, DehydratedDeviceError> {
506 let olm = self.client.olm_machine().await;
507 let machine = olm.as_ref().ok_or(DehydratedDeviceError::NotLoggedIn)?;
508
509 Ok(machine.store().get_value(LAST_UPLOADED_DEVICE_ID_KEY).await?)
510 }
511
512 /// Return whether the pickle key is stored in the given Secret Storage.
513 ///
514 /// The key is looked up by the account-data event type `org.matrix.msc3814`
515 /// (the unstable name reserved by MSC3814).
516 ///
517 /// # Example
518 ///
519 /// ```no_run
520 /// # use matrix_sdk::{Client, encryption::secret_storage::SecretStore};
521 /// # async fn example(client: Client, store: SecretStore)
522 /// # -> anyhow::Result<()> {
523 /// let stored =
524 /// client.encryption().dehydrated_devices().is_key_stored(&store).await?;
525 /// # Ok(()) }
526 /// ```
527 #[instrument(skip_all)]
528 pub async fn is_key_stored(
529 &self,
530 secret_store: &SecretStore,
531 ) -> Result<bool, DehydratedDeviceError> {
532 Ok(secret_store.get_secret(pickle_key_secret_name()).await?.is_some())
533 }
534
535 /// Generate a new random pickle key, persist it in Secret Storage, and
536 /// cache it in the local crypto store.
537 ///
538 /// The previous key (if any) is overwritten in both places. Any dehydrated
539 /// device that was encrypted with the previous key becomes unrehydratable
540 /// until rotated.
541 ///
542 /// # Example
543 ///
544 /// ```no_run
545 /// # use matrix_sdk::{Client, encryption::secret_storage::SecretStore};
546 /// # async fn example(client: Client, store: SecretStore)
547 /// # -> anyhow::Result<()> {
548 /// let fresh_key =
549 /// client.encryption().dehydrated_devices().reset_key(&store).await?;
550 /// # let _ = fresh_key;
551 /// # Ok(()) }
552 /// ```
553 #[instrument(skip_all)]
554 pub async fn reset_key(
555 &self,
556 secret_store: &SecretStore,
557 ) -> Result<DehydratedDeviceKey, DehydratedDeviceError> {
558 let key = DehydratedDeviceKey::new();
559 secret_store.put_secret(pickle_key_secret_name(), &key.to_base64()).await?;
560 self.cache_key(&key).await?;
561 Ok(key)
562 }
563
564 /// Resolve the pickle key.
565 ///
566 /// Looks for the key in this order:
567 ///
568 /// 1. The local crypto-store cache.
569 /// 2. The provided Secret Storage account-data entry. A successfully
570 /// fetched key is written back to the cache.
571 /// 3. If `create_if_missing`, a fresh random key is generated, stored, and
572 /// cached.
573 ///
574 /// Returns `Ok(None)` only when the key is absent from both sources and
575 /// `create_if_missing` is `false`.
576 #[instrument(skip_all)]
577 pub(crate) async fn load_key(
578 &self,
579 secret_store: &SecretStore,
580 create_if_missing: bool,
581 ) -> Result<Option<DehydratedDeviceKey>, DehydratedDeviceError> {
582 if let Some(cached) = self.cached_key().await? {
583 return Ok(Some(cached));
584 }
585
586 let Some(base64) = secret_store.get_secret(pickle_key_secret_name()).await? else {
587 return if create_if_missing {
588 Ok(Some(self.reset_key(secret_store).await?))
589 } else {
590 Ok(None)
591 };
592 };
593
594 let bytes = Zeroizing::new(base64_decode(&base64)?);
595 let key = DehydratedDeviceKey::from_slice(&bytes)?;
596 self.cache_key(&key).await?;
597 Ok(Some(key))
598 }
599
600 /// Start using dehydrated devices for this client.
601 ///
602 /// Returns a [`StartDehydration`] future. Await it directly for the default
603 /// behavior, or configure it first with
604 /// [`StartDehydration::create_new_key`],
605 /// [`StartDehydration::skip_rehydration`], and
606 /// [`StartDehydration::only_if_key_cached`].
607 ///
608 /// The caller is expected to have unlocked Secret Storage and bootstrapped
609 /// cross-signing before this call; otherwise the underlying account-data
610 /// reads and writes will fail mid-flight.
611 ///
612 /// Awaiting the future performs the following steps:
613 ///
614 /// 1. If [`StartDehydration::only_if_key_cached`] was set, return early
615 /// when no pickle key is cached locally.
616 /// 2. Stop any previously scheduled rotation.
617 /// 3. Unless [`StartDehydration::skip_rehydration`] was set, attempt to
618 /// rehydrate the existing dehydrated device. Failures are logged and
619 /// emitted as [`DehydratedDeviceEvent::RehydrationError`] but do not
620 /// abort the start.
621 /// 4. If [`StartDehydration::create_new_key`] was set and the rehydration
622 /// step succeeded (or was skipped), replace the pickle key in Secret
623 /// Storage with a fresh random one. A failed rehydration suppresses the
624 /// reset so the stored key can still recover the existing dehydrated
625 /// device on another client.
626 /// 5. Create a new dehydrated device now and schedule rotation once a week.
627 ///
628 /// The rotation task resolves the pickle key from the local crypto store on
629 /// each tick. The local cache is the only key source available to the task
630 /// because reopening Secret Storage requires the recovery key, which is not
631 /// retained. Per-tick failures emit
632 /// [`DehydratedDeviceEvent::RotationError`] without aborting the schedule.
633 ///
634 /// # Example
635 ///
636 /// ```no_run
637 /// # use matrix_sdk::{Client, encryption::secret_storage::SecretStore};
638 /// # async fn example(client: Client, store: SecretStore)
639 /// # -> anyhow::Result<()> {
640 /// client.encryption().dehydrated_devices().start(&store).await?;
641 /// # Ok(()) }
642 /// ```
643 pub fn start<'a>(&'a self, secret_store: &'a SecretStore) -> StartDehydration<'a> {
644 StartDehydration::new(self, secret_store)
645 }
646
647 /// Stop the scheduled dehydrated-device rotation, if any.
648 ///
649 /// Has no effect when no rotation is scheduled. Existing dehydrated devices
650 /// on the server are left in place; pair with [`Self::delete`] to clean
651 /// those up.
652 pub fn stop(&self) {
653 self.state().rotation_task.lock().take();
654 }
655
656 /// Create-and-upload the first dehydrated device now, then spawn the
657 /// rotation loop.
658 async fn schedule_dehydration(
659 &self,
660 secret_store: &SecretStore,
661 ) -> Result<(), DehydratedDeviceError> {
662 let key = self
663 .load_key(secret_store, true)
664 .await?
665 .expect("load_key(create_if_missing=true) always yields a key");
666 self.create(None, &key).await?;
667
668 let weak_client = WeakClient::from_client(&self.client);
669 let handle = self
670 .client
671 .task_monitor()
672 .spawn_infinite_task("dehydrated_devices::rotation", async move {
673 loop {
674 sleep(DEHYDRATION_INTERVAL).await;
675
676 let Some(client) = weak_client.get() else {
677 // The client is gone; this task is aborted when its
678 // handle drops, so just wait for that to happen.
679 continue;
680 };
681
682 client.encryption().dehydrated_devices().rotate_tick().await;
683 }
684 })
685 .abort_on_drop();
686
687 *self.state().rotation_task.lock() = Some(handle);
688 Ok(())
689 }
690
691 /// One iteration of the rotation timer: resolve the cached pickle key and
692 /// upload a fresh dehydrated device, emitting a
693 /// [`DehydratedDeviceEvent::RotationError`] on failure.
694 async fn rotate_tick(&self) {
695 let key = match self.cached_key().await {
696 Ok(Some(key)) => key,
697 Ok(None) => {
698 let msg = "no cached pickle key for dehydrated-device rotation".to_owned();
699 warn!("{msg}; skipping this rotation");
700 self.emit(DehydratedDeviceEvent::RotationError { error: msg });
701 return;
702 }
703 Err(e) => {
704 let msg = e.to_string();
705 warn!(error = %e, "Failed to load cached pickle key for rotation");
706 self.emit(DehydratedDeviceEvent::RotationError { error: msg });
707 return;
708 }
709 };
710
711 if let Err(e) = self.create(None, &key).await {
712 let msg = e.to_string();
713 warn!(error = msg, "Failed to rotate dehydrated device");
714 self.emit(DehydratedDeviceEvent::RotationError { error: msg });
715 }
716 }
717
718 /// Delete the current dehydrated device, if one exists.
719 ///
720 /// Also stops any scheduled rotation, so the next tick will not immediately
721 /// recreate the device the caller just asked to remove.
722 ///
723 /// Returns `Ok(())` silently if no dehydrated device is on the server or
724 /// the server does not implement the endpoint.
725 ///
726 /// # Example
727 ///
728 /// ```no_run
729 /// # use matrix_sdk::Client;
730 /// # async fn example(client: Client) -> anyhow::Result<()> {
731 /// client.encryption().dehydrated_devices().delete().await?;
732 /// # Ok(()) }
733 /// ```
734 #[instrument(skip_all)]
735 pub async fn delete(&self) -> Result<(), DehydratedDeviceError> {
736 self.stop();
737 self.delete_device().await
738 }
739
740 /// Issue the delete request for the current dehydrated device without
741 /// touching the rotation schedule.
742 async fn delete_device(&self) -> Result<(), DehydratedDeviceError> {
743 let request = delete_dehydrated_device::unstable::Request::new();
744 match self.client.send(request).await {
745 Ok(_) => {
746 self.emit(DehydratedDeviceEvent::Deleted);
747 Ok(())
748 }
749 Err(e) => match e.client_api_error_kind() {
750 Some(ErrorKind::Unrecognized) | Some(ErrorKind::NotFound) => Ok(()),
751 _ => Err(e.into()),
752 },
753 }
754 }
755
756 /// Fetch the dehydrated device payload from the server.
757 ///
758 /// Returns `Ok(None)` if the server reports `M_NOT_FOUND` or
759 /// `M_UNRECOGNIZED`.
760 async fn download_device(&self) -> Result<Option<DownloadedDevice>, DehydratedDeviceError> {
761 let request = get_dehydrated_device::unstable::Request::new();
762 match self.client.send(request).await {
763 Ok(response) => Ok(Some(DownloadedDevice {
764 device_id: response.device_id,
765 device_data: response.device_data,
766 })),
767 Err(e) => match e.client_api_error_kind() {
768 Some(ErrorKind::NotFound) | Some(ErrorKind::Unrecognized) => Ok(None),
769 _ => Err(e.into()),
770 },
771 }
772 }
773
774 /// Decrypt the downloaded device and stand up a [`RehydratedDevice`].
775 async fn rehydrate_device(
776 &self,
777 downloaded: &DownloadedDevice,
778 pickle_key: &DehydratedDeviceKey,
779 ) -> Result<RehydratedDevice, DehydratedDeviceError> {
780 let olm = self.client.olm_machine().await;
781 let machine = olm.as_ref().ok_or(DehydratedDeviceError::NotLoggedIn)?;
782
783 Ok(machine
784 .dehydrated_devices()
785 .rehydrate(pickle_key, &downloaded.device_id, downloaded.device_data.clone())
786 .await?)
787 }
788
789 /// Drain every queued to-device event from the dehydrated device's
790 /// server-side buffer, feeding each batch through the rehydrated machine so
791 /// the room keys are imported.
792 ///
793 /// An empty batch or an absent cursor ends the drain cleanly; the two
794 /// defensive stops (batch cap, repeated cursor) mark the returned
795 /// [`DrainOutcome`] as truncated instead.
796 async fn absorb_events(
797 &self,
798 device_id: &OwnedDeviceId,
799 rehydrated: &RehydratedDevice,
800 ) -> Result<DrainOutcome, DehydratedDeviceError> {
801 let settings = self.client.decryption_settings();
802
803 let mut next_batch: Option<String> = None;
804 let mut to_device_count: usize = 0;
805 let mut room_key_count: usize = 0;
806
807 let truncated = loop {
808 let mut request = get_events::unstable::Request::new(device_id.clone());
809 request.next_batch.clone_from(&next_batch);
810
811 let response = self.client.send(request).await?;
812 if response.events.is_empty() {
813 break false;
814 }
815
816 to_device_count = to_device_count.saturating_add(response.events.len());
817 let imported = rehydrated.receive_events(response.events, settings).await?;
818 room_key_count = room_key_count.saturating_add(imported.len());
819 trace!(to_device_count, room_key_count, "Absorbed a batch of to-device events");
820 self.emit(DehydratedDeviceEvent::RehydrationProgress {
821 room_keys_imported: room_key_count,
822 to_device_events: to_device_count,
823 });
824
825 if to_device_count >= MAX_TO_DEVICE_EVENTS {
826 warn!(
827 to_device_count,
828 "Reached the dehydrated-device drain limit; stopping to bound the loop"
829 );
830 break true;
831 }
832
833 // Guard against a server that repeats the same cursor and would
834 // otherwise keep us looping over an identical batch.
835 match &response.next_batch {
836 None => break false,
837 Some(token) if Some(token) == next_batch.as_ref() => {
838 warn!(
839 ?next_batch,
840 "Server returned the same next_batch twice; stopping to avoid a loop"
841 );
842 break true;
843 }
844 Some(token) => next_batch = Some(token.clone()),
845 }
846 };
847
848 info!(
849 to_device_count,
850 room_key_count, truncated, "Finished draining the dehydrated device's to-device queue"
851 );
852 Ok(DrainOutcome {
853 room_keys_imported: room_key_count,
854 to_device_events: to_device_count,
855 truncated,
856 })
857 }
858}
859
860/// Named future returned by [`DehydratedDevices::start`].
861///
862/// Await it to start dehydrated devices with the default behavior, or configure
863/// it first with [`Self::create_new_key`], [`Self::skip_rehydration`], and
864/// [`Self::only_if_key_cached`].
865#[derive(Debug)]
866pub struct StartDehydration<'a> {
867 devices: &'a DehydratedDevices,
868 secret_store: &'a SecretStore,
869 create_new_key: bool,
870 rehydrate: bool,
871 only_if_key_cached: bool,
872 tracing_span: Span,
873}
874
875impl<'a> StartDehydration<'a> {
876 fn new(devices: &'a DehydratedDevices, secret_store: &'a SecretStore) -> Self {
877 Self {
878 devices,
879 secret_store,
880 create_new_key: false,
881 rehydrate: true,
882 only_if_key_cached: false,
883 tracing_span: Span::current(),
884 }
885 }
886
887 /// Generate a fresh random pickle key on start, replacing any existing
888 /// entry in Secret Storage and the local cache.
889 ///
890 /// The reset is suppressed if a rehydration attempt in the same start
891 /// fails, so the stored key can still recover the existing dehydrated
892 /// device on another client.
893 pub fn create_new_key(mut self) -> Self {
894 self.create_new_key = true;
895 self
896 }
897
898 /// Skip the attempt to rehydrate the existing dehydrated device before
899 /// creating the next one.
900 ///
901 /// By default the existing device is rehydrated first.
902 pub fn skip_rehydration(mut self) -> Self {
903 self.rehydrate = false;
904 self
905 }
906
907 /// Do nothing unless a pickle key is already cached locally.
908 ///
909 /// Useful for an opportunistic start on a freshly opened client without
910 /// forcing a Secret Storage unlock.
911 pub fn only_if_key_cached(mut self) -> Self {
912 self.only_if_key_cached = true;
913 self
914 }
915}
916
917impl<'a> IntoFuture for StartDehydration<'a> {
918 type Output = Result<(), DehydratedDeviceError>;
919 boxed_into_future!(extra_bounds: 'a);
920
921 fn into_future(self) -> Self::IntoFuture {
922 let Self {
923 devices,
924 secret_store,
925 create_new_key,
926 rehydrate,
927 only_if_key_cached,
928 tracing_span,
929 } = self;
930
931 let future = async move {
932 if only_if_key_cached && devices.cached_key().await?.is_none() {
933 return Ok(());
934 }
935
936 devices.stop();
937
938 let mut rehydrate_failed = false;
939 if rehydrate
940 && let Some(key) = devices.load_key(secret_store, false).await?
941 && let Err(e) = devices.rehydrate(&key).await
942 {
943 let msg = e.to_string();
944 warn!(error = %e, "Rehydration failed during start; continuing");
945 devices.emit(DehydratedDeviceEvent::RehydrationError { error: msg });
946 rehydrate_failed = true;
947 }
948
949 // Refuse to clobber Secret Storage after a failed rehydration: the
950 // stored pickle key is the only one that can decrypt the existing
951 // dehydrated device, and overwriting it would discard that recovery
952 // chance for every other client signed in to this account.
953 if create_new_key {
954 if rehydrate_failed {
955 warn!(
956 "Skipping pickle-key reset after failed rehydration to preserve the chance of recovering the existing dehydrated device on another client"
957 );
958 } else {
959 devices.reset_key(secret_store).await?;
960 }
961 }
962
963 devices.schedule_dehydration(secret_store).await
964 };
965
966 Box::pin(future.instrument(tracing_span))
967 }
968}