Skip to main content

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}