Skip to main content

matrix_sdk_indexeddb/state_store/
mod.rs

1// Copyright 2021 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
15use std::{
16    cmp::Reverse,
17    collections::{BTreeMap, BTreeSet, HashMap, HashSet},
18    str::FromStr as _,
19    sync::Arc,
20};
21
22use async_trait::async_trait;
23use gloo_utils::format::JsValueSerdeExt;
24use growable_bloom_filter::GrowableBloom;
25use indexed_db_futures::{
26    KeyRange, cursor::CursorDirection, database::Database, error::OpenDbError, prelude::*,
27    transaction::TransactionMode,
28};
29use matrix_sdk_base::{
30    MinimalRoomMemberEvent, ROOM_VERSION_FALLBACK, ROOM_VERSION_RULES_FALLBACK, RoomInfo,
31    RoomMemberships, StateStoreDataKey, StateStoreDataValue, ThreadSubscriptionCatchupToken,
32    deserialized_responses::{DisplayName, RawAnySyncOrStrippedState},
33    store::{
34        ChildTransactionId, ComposerDraft, DependentQueuedRequest, DependentQueuedRequestKind,
35        QueuedRequest, QueuedRequestKind, RoomLoadSettings, SentRequestKey,
36        SerializableEventContent, StateChanges, StateStore, StoreError, StoredThreadSubscription,
37        SupportedVersionsResponse, ThreadSubscriptionStatus, WellKnownResponse,
38        compare_thread_subscription_bump_stamps,
39    },
40    ttl::TtlValue,
41};
42use matrix_sdk_store_encryption::{Error as EncryptionError, StoreCipher};
43use ruma::{
44    CanonicalJsonObject, EventId, MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedMxcUri,
45    OwnedRoomId, OwnedTransactionId, OwnedUserId, RoomId, TransactionId, UserId,
46    api::client::discovery::get_capabilities::v3::Capabilities,
47    canonical_json::{RedactedBecause, redact},
48    events::{
49        AnyGlobalAccountDataEvent, AnyRoomAccountDataEvent, AnySyncStateEvent,
50        GlobalAccountDataEventType, RoomAccountDataEventType, StateEventType, SyncStateEvent,
51        presence::PresenceEvent,
52        receipt::{Receipt, ReceiptThread, ReceiptType},
53        room::member::{
54            MembershipState, RoomMemberEventContent, StrippedRoomMemberEvent, SyncRoomMemberEvent,
55        },
56    },
57    profile::{UserProfile, UserProfileUpdate},
58    serde::Raw,
59};
60use serde::{Deserialize, Serialize, de::DeserializeOwned, ser::Error};
61use tracing::{debug, warn};
62use wasm_bindgen::JsValue;
63
64mod migrations;
65
66pub use self::migrations::MigrationConflictStrategy;
67use self::migrations::{upgrade_inner_db, upgrade_meta_db};
68use crate::{error::GenericError, serializer::safe_encode::traits::SafeEncode};
69
70#[derive(Debug, thiserror::Error)]
71pub enum IndexeddbStateStoreError {
72    #[error(transparent)]
73    Json(#[from] serde_json::Error),
74    #[error(transparent)]
75    Encryption(#[from] EncryptionError),
76    #[error("DomException {name} ({code}): {message}")]
77    DomException { name: String, message: String, code: u16 },
78    #[error(transparent)]
79    StoreError(#[from] StoreError),
80    #[error(
81        "Can't migrate {name} from {old_version} to {new_version} without deleting data. \
82         See MigrationConflictStrategy for ways to configure."
83    )]
84    MigrationConflict { name: String, old_version: u32, new_version: u32 },
85}
86
87impl From<GenericError> for IndexeddbStateStoreError {
88    fn from(value: GenericError) -> Self {
89        Self::StoreError(value.into())
90    }
91}
92
93impl From<web_sys::DomException> for IndexeddbStateStoreError {
94    fn from(frm: web_sys::DomException) -> IndexeddbStateStoreError {
95        IndexeddbStateStoreError::DomException {
96            name: frm.name(),
97            message: frm.message(),
98            code: frm.code(),
99        }
100    }
101}
102
103impl From<IndexeddbStateStoreError> for StoreError {
104    fn from(e: IndexeddbStateStoreError) -> Self {
105        match e {
106            IndexeddbStateStoreError::Json(e) => StoreError::Json(e),
107            IndexeddbStateStoreError::StoreError(e) => e,
108            IndexeddbStateStoreError::Encryption(e) => StoreError::Encryption(e),
109            _ => StoreError::backend(e),
110        }
111    }
112}
113
114impl From<indexed_db_futures::error::DomException> for IndexeddbStateStoreError {
115    fn from(value: indexed_db_futures::error::DomException) -> Self {
116        web_sys::DomException::from(value).into()
117    }
118}
119
120impl From<indexed_db_futures::error::SerialisationError> for IndexeddbStateStoreError {
121    fn from(value: indexed_db_futures::error::SerialisationError) -> Self {
122        Self::Json(serde_json::Error::custom(value.to_string()))
123    }
124}
125
126impl From<indexed_db_futures::error::UnexpectedDataError> for IndexeddbStateStoreError {
127    fn from(value: indexed_db_futures::error::UnexpectedDataError) -> Self {
128        IndexeddbStateStoreError::StoreError(StoreError::backend(value))
129    }
130}
131
132impl From<indexed_db_futures::error::JSError> for IndexeddbStateStoreError {
133    fn from(value: indexed_db_futures::error::JSError) -> Self {
134        GenericError::from(value.to_string()).into()
135    }
136}
137
138impl From<indexed_db_futures::error::Error> for IndexeddbStateStoreError {
139    fn from(value: indexed_db_futures::error::Error) -> Self {
140        use indexed_db_futures::error::Error;
141        match value {
142            Error::DomException(e) => e.into(),
143            Error::Serialisation(e) => e.into(),
144            Error::MissingData(e) => e.into(),
145            Error::Unknown(e) => e.into(),
146        }
147    }
148}
149
150impl From<OpenDbError> for IndexeddbStateStoreError {
151    fn from(value: OpenDbError) -> Self {
152        match value {
153            OpenDbError::Base(error) => error.into(),
154            _ => GenericError::from(value.to_string()).into(),
155        }
156    }
157}
158
159mod keys {
160    pub const INTERNAL_STATE: &str = "matrix-sdk-state";
161    pub const BACKUPS_META: &str = "backups";
162
163    pub const ACCOUNT_DATA: &str = "account_data";
164
165    /// Room profiles.
166    pub const PROFILES: &str = "profiles";
167    pub const DISPLAY_NAMES: &str = "display_names";
168    pub const USER_IDS: &str = "user_ids";
169
170    pub const ROOM_STATE: &str = "room_state";
171    pub const ROOM_INFOS: &str = "room_infos";
172    pub const PRESENCE: &str = "presence";
173    pub const ROOM_ACCOUNT_DATA: &str = "room_account_data";
174    /// Table used to save send queue events.
175    pub const ROOM_SEND_QUEUE: &str = "room_send_queue";
176    /// Table used to save dependent send queue events.
177    pub const DEPENDENT_SEND_QUEUE: &str = "room_dependent_send_queue";
178    pub const THREAD_SUBSCRIPTIONS: &str = "room_thread_subscriptions";
179
180    pub const STRIPPED_ROOM_STATE: &str = "stripped_room_state";
181    pub const STRIPPED_USER_IDS: &str = "stripped_user_ids";
182
183    pub const ROOM_USER_RECEIPTS: &str = "room_user_receipts";
184    pub const ROOM_EVENT_RECEIPTS: &str = "room_event_receipts";
185
186    pub const GLOBAL_PROFILES: &str = "global_profiles";
187
188    pub const CUSTOM: &str = "custom";
189    pub const KV: &str = "kv";
190
191    /// All names of the current state stores for convenience.
192    pub const ALL_STORES: &[&str] = &[
193        ACCOUNT_DATA,
194        PROFILES,
195        DISPLAY_NAMES,
196        USER_IDS,
197        ROOM_STATE,
198        ROOM_INFOS,
199        PRESENCE,
200        ROOM_ACCOUNT_DATA,
201        STRIPPED_ROOM_STATE,
202        STRIPPED_USER_IDS,
203        ROOM_USER_RECEIPTS,
204        ROOM_EVENT_RECEIPTS,
205        ROOM_SEND_QUEUE,
206        THREAD_SUBSCRIPTIONS,
207        DEPENDENT_SEND_QUEUE,
208        GLOBAL_PROFILES,
209        CUSTOM,
210        KV,
211    ];
212
213    // static keys
214
215    pub const STORE_KEY: &str = "store_key";
216}
217
218pub use keys::ALL_STORES;
219use matrix_sdk_base::store::QueueWedgeError;
220
221/// Encrypt (if needs be) then JSON-serialize a value.
222fn serialize_value(store_cipher: Option<&StoreCipher>, event: &impl Serialize) -> Result<JsValue> {
223    Ok(match store_cipher {
224        Some(cipher) => {
225            let data = serde_json::to_vec(event)?;
226            JsValue::from_serde(&cipher.encrypt_value_data(data)?)?
227        }
228        None => JsValue::from_serde(event)?,
229    })
230}
231
232/// Deserialize a JSON value and then decrypt it (if needs be).
233fn deserialize_value<T: DeserializeOwned>(
234    store_cipher: Option<&StoreCipher>,
235    event: &JsValue,
236) -> Result<T> {
237    match store_cipher {
238        Some(cipher) => {
239            use zeroize::Zeroize;
240            let mut plaintext = cipher.decrypt_value_data(event.into_serde()?)?;
241            let ret = serde_json::from_slice(&plaintext);
242            plaintext.zeroize();
243            Ok(ret?)
244        }
245        None => Ok(event.into_serde()?),
246    }
247}
248
249fn encode_key<T>(store_cipher: Option<&StoreCipher>, table_name: &str, key: T) -> JsValue
250where
251    T: SafeEncode,
252{
253    match store_cipher {
254        Some(cipher) => key.as_secure_string(table_name, cipher),
255        None => key.as_encoded_string(),
256    }
257    .into()
258}
259
260fn encode_to_range<T>(
261    store_cipher: Option<&StoreCipher>,
262    table_name: &str,
263    key: T,
264) -> KeyRange<JsValue>
265where
266    T: SafeEncode,
267{
268    match store_cipher {
269        Some(cipher) => key.encode_to_range_secure(table_name, cipher),
270        None => key.encode_to_range(),
271    }
272}
273
274/// Builder for [`IndexeddbStateStore`].
275#[derive(Debug)]
276pub struct IndexeddbStateStoreBuilder {
277    name: Option<String>,
278    passphrase: Option<String>,
279    migration_conflict_strategy: MigrationConflictStrategy,
280}
281
282impl IndexeddbStateStoreBuilder {
283    fn new() -> Self {
284        Self {
285            name: None,
286            passphrase: None,
287            migration_conflict_strategy: MigrationConflictStrategy::BackupAndDrop,
288        }
289    }
290
291    /// Set the name for the indexeddb store to use, `state` is none given.
292    pub fn name(mut self, value: String) -> Self {
293        self.name = Some(value);
294        self
295    }
296
297    /// Set the password the indexeddb should be encrypted with.
298    ///
299    /// If not given, the DB is not encrypted.
300    pub fn passphrase(mut self, value: String) -> Self {
301        self.passphrase = Some(value);
302        self
303    }
304
305    /// The strategy to use when a merge conflict is found.
306    ///
307    /// See [`MigrationConflictStrategy`] for details.
308    pub fn migration_conflict_strategy(mut self, value: MigrationConflictStrategy) -> Self {
309        self.migration_conflict_strategy = value;
310        self
311    }
312
313    pub async fn build(self) -> Result<IndexeddbStateStore> {
314        let migration_strategy = self.migration_conflict_strategy.clone();
315        let name = self.name.unwrap_or_else(|| "state".to_owned());
316
317        let meta_name = format!("{name}::{}", keys::INTERNAL_STATE);
318
319        let (meta, store_cipher) = upgrade_meta_db(&meta_name, self.passphrase.as_deref()).await?;
320        let inner =
321            upgrade_inner_db(&name, store_cipher.as_deref(), migration_strategy, &meta).await?;
322
323        Ok(IndexeddbStateStore { name, inner, meta, store_cipher })
324    }
325}
326
327pub struct IndexeddbStateStore {
328    name: String,
329    pub(crate) inner: Database,
330    pub(crate) meta: Database,
331    pub(crate) store_cipher: Option<Arc<StoreCipher>>,
332}
333
334#[cfg(not(tarpaulin_include))]
335impl std::fmt::Debug for IndexeddbStateStore {
336    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
337        f.debug_struct("IndexeddbStateStore").field("name", &self.name).finish()
338    }
339}
340
341type Result<A, E = IndexeddbStateStoreError> = std::result::Result<A, E>;
342
343impl IndexeddbStateStore {
344    /// Generate a IndexeddbStateStoreBuilder with default parameters
345    pub fn builder() -> IndexeddbStateStoreBuilder {
346        IndexeddbStateStoreBuilder::new()
347    }
348
349    /// The version of the database containing the data.
350    pub fn version(&self) -> u32 {
351        self.inner.version() as u32
352    }
353
354    /// The version of the database containing the metadata.
355    pub fn meta_version(&self) -> u32 {
356        self.meta.version() as u32
357    }
358
359    /// Whether this database has any migration backups
360    pub async fn has_backups(&self) -> Result<bool> {
361        Ok(self
362            .meta
363            .transaction(keys::BACKUPS_META)
364            .with_mode(TransactionMode::Readonly)
365            .build()?
366            .object_store(keys::BACKUPS_META)?
367            .count()
368            .await?
369            > 0)
370    }
371
372    /// What's the database name of the latest backup<
373    pub async fn latest_backup(&self) -> Result<Option<String>> {
374        if let Some(mut cursor) = self
375            .meta
376            .transaction(keys::BACKUPS_META)
377            .with_mode(TransactionMode::Readonly)
378            .build()?
379            .object_store(keys::BACKUPS_META)?
380            .open_cursor()
381            .with_direction(CursorDirection::Prev)
382            .await?
383            && let Some(record) = cursor.next_record::<JsValue>().await?
384        {
385            Ok(record.as_string())
386        } else {
387            Ok(None)
388        }
389    }
390
391    /// Encrypt (if needs be) then JSON-serialize a value.
392    fn serialize_value(&self, event: &impl Serialize) -> Result<JsValue> {
393        serialize_value(self.store_cipher.as_deref(), event)
394    }
395
396    /// Deserialize a JSON value and then decrypt it (if needs be).
397    fn deserialize_value<T: DeserializeOwned>(&self, event: &JsValue) -> Result<T> {
398        deserialize_value(self.store_cipher.as_deref(), event)
399    }
400
401    fn encode_key<T>(&self, table_name: &str, key: T) -> JsValue
402    where
403        T: SafeEncode,
404    {
405        encode_key(self.store_cipher.as_deref(), table_name, key)
406    }
407
408    fn encode_to_range<T>(&self, table_name: &str, key: T) -> KeyRange<JsValue>
409    where
410        T: SafeEncode,
411    {
412        encode_to_range(self.store_cipher.as_deref(), table_name, key)
413    }
414
415    /// Get user IDs for the given room with the given memberships and stripped
416    /// state.
417    pub async fn get_user_ids_inner(
418        &self,
419        room_id: &RoomId,
420        memberships: RoomMemberships,
421        stripped: bool,
422    ) -> Result<Vec<OwnedUserId>> {
423        let store_name = if stripped { keys::STRIPPED_USER_IDS } else { keys::USER_IDS };
424
425        let tx = self.inner.transaction(store_name).with_mode(TransactionMode::Readonly).build()?;
426        let store = tx.object_store(store_name)?;
427        let range = self.encode_to_range(store_name, room_id);
428
429        let user_ids = if memberships.is_empty() {
430            // It should be faster to just get all user IDs in this case.
431            store
432                .get_all()
433                .with_query(&range)
434                .await?
435                .filter_map(Result::ok)
436                .filter_map(|f| self.deserialize_value::<RoomMember>(&f).ok().map(|m| m.user_id))
437                .collect::<Vec<_>>()
438        } else {
439            let mut user_ids = Vec::new();
440            let cursor = store.open_cursor().with_query(&range).await?;
441
442            if let Some(mut cursor) = cursor {
443                while let Some(value) = cursor.next_record().await? {
444                    let member = self.deserialize_value::<RoomMember>(&value)?;
445
446                    if memberships.matches(&member.membership) {
447                        user_ids.push(member.user_id);
448                    }
449                }
450            }
451
452            user_ids
453        };
454
455        Ok(user_ids)
456    }
457
458    async fn get_custom_value_for_js(&self, jskey: &JsValue) -> Result<Option<Vec<u8>>> {
459        self.inner
460            .transaction(keys::CUSTOM)
461            .with_mode(TransactionMode::Readonly)
462            .build()?
463            .object_store(keys::CUSTOM)?
464            .get(jskey)
465            .await?
466            .map(|f| self.deserialize_value(&f))
467            .transpose()
468    }
469
470    fn encode_kv_data_key(&self, key: StateStoreDataKey<'_>) -> JsValue {
471        // Use the key (prefix) for the table name as well, to keep encoded
472        // keys compatible for the sync token and filters, which were in
473        // separate tables initially.
474        match key {
475            StateStoreDataKey::SyncToken => {
476                self.encode_key(StateStoreDataKey::SYNC_TOKEN, StateStoreDataKey::SYNC_TOKEN)
477            }
478            StateStoreDataKey::SupportedVersions => self.encode_key(
479                StateStoreDataKey::SUPPORTED_VERSIONS,
480                StateStoreDataKey::SUPPORTED_VERSIONS,
481            ),
482            StateStoreDataKey::WellKnown => {
483                self.encode_key(keys::KV, StateStoreDataKey::WELL_KNOWN)
484            }
485            StateStoreDataKey::Filter(filter_name) => {
486                self.encode_key(StateStoreDataKey::FILTER, (StateStoreDataKey::FILTER, filter_name))
487            }
488            StateStoreDataKey::UserAvatarUrl(user_id) => {
489                self.encode_key(keys::KV, (StateStoreDataKey::USER_AVATAR_URL, user_id))
490            }
491            StateStoreDataKey::RecentlyVisitedRooms(user_id) => {
492                self.encode_key(keys::KV, (StateStoreDataKey::RECENTLY_VISITED_ROOMS, user_id))
493            }
494            StateStoreDataKey::UtdHookManagerData => {
495                self.encode_key(keys::KV, StateStoreDataKey::UTD_HOOK_MANAGER_DATA)
496            }
497            StateStoreDataKey::OneTimeKeyAlreadyUploaded => {
498                self.encode_key(keys::KV, StateStoreDataKey::ONE_TIME_KEY_ALREADY_UPLOADED)
499            }
500            StateStoreDataKey::ComposerDraft(room_id, thread_root) => {
501                if let Some(thread_root) = thread_root {
502                    self.encode_key(
503                        keys::KV,
504                        (StateStoreDataKey::COMPOSER_DRAFT, (room_id, thread_root)),
505                    )
506                } else {
507                    self.encode_key(keys::KV, (StateStoreDataKey::COMPOSER_DRAFT, room_id))
508                }
509            }
510            StateStoreDataKey::SeenKnockRequests(room_id) => {
511                self.encode_key(keys::KV, (StateStoreDataKey::SEEN_KNOCK_REQUESTS, room_id))
512            }
513            StateStoreDataKey::ThreadSubscriptionsCatchupTokens => {
514                self.encode_key(keys::KV, StateStoreDataKey::THREAD_SUBSCRIPTIONS_CATCHUP_TOKENS)
515            }
516            StateStoreDataKey::HomeserverCapabilities => {
517                self.encode_key(keys::KV, StateStoreDataKey::HOMESERVER_CAPABILITIES)
518            }
519        }
520    }
521}
522
523/// A superset of [`QueuedRequest`] that also contains the room id, since we
524/// want to return them.
525#[derive(Serialize, Deserialize)]
526struct PersistedQueuedRequest {
527    /// In which room is this event going to be sent.
528    pub room_id: OwnedRoomId,
529
530    // All these fields are the same as in [`QueuedRequest`].
531    /// Kind. Optional because it might be missing from previous formats.
532    kind: Option<QueuedRequestKind>,
533    transaction_id: OwnedTransactionId,
534
535    pub error: Option<QueueWedgeError>,
536
537    priority: Option<usize>,
538
539    /// The time the original message was first attempted to be sent at.
540    #[serde(default = "created_now")]
541    created_at: MilliSecondsSinceUnixEpoch,
542
543    // Migrated fields: keep these private, they're not used anymore elsewhere in the code base.
544    /// Deprecated (from old format), now replaced with error field.
545    is_wedged: Option<bool>,
546
547    event: Option<SerializableEventContent>,
548}
549
550fn created_now() -> MilliSecondsSinceUnixEpoch {
551    MilliSecondsSinceUnixEpoch::now()
552}
553
554impl PersistedQueuedRequest {
555    fn into_queued_request(self) -> Option<QueuedRequest> {
556        let kind =
557            self.kind.or_else(|| self.event.map(|content| QueuedRequestKind::Event { content }))?;
558
559        let error = match self.is_wedged {
560            Some(true) => {
561                // Migrate to a generic error.
562                Some(QueueWedgeError::GenericApiError {
563                    msg: "local echo failed to send in a previous session".into(),
564                })
565            }
566            _ => self.error,
567        };
568
569        // By default, events without a priority have a priority of 0.
570        let priority = self.priority.unwrap_or(0);
571
572        Some(QueuedRequest {
573            kind,
574            transaction_id: self.transaction_id,
575            error,
576            priority,
577            created_at: self.created_at,
578        })
579    }
580}
581
582#[derive(Serialize, Deserialize, PartialEq)]
583struct PersistedThreadSubscription {
584    status: String,
585    bump_stamp: Option<u64>,
586}
587
588impl From<StoredThreadSubscription> for PersistedThreadSubscription {
589    fn from(value: StoredThreadSubscription) -> Self {
590        Self { status: value.status.as_str().to_owned(), bump_stamp: value.bump_stamp }
591    }
592}
593
594// Small hack to have the following macro invocation act as the appropriate
595// trait impl block on wasm, but still be compiled on non-wasm as a regular
596// impl block otherwise.
597//
598// The trait impl doesn't compile on non-wasm due to unfulfilled trait bounds,
599// this hack allows us to still have most of rust-analyzer's IDE functionality
600// within the impl block without having to set it up to check things against
601// the wasm target (which would disable many other parts of the codebase).
602#[cfg(target_family = "wasm")]
603macro_rules! impl_state_store {
604    ({ $($body:tt)* }) => {
605        #[async_trait(?Send)]
606        impl StateStore for IndexeddbStateStore {
607            type Error = IndexeddbStateStoreError;
608
609            $($body)*
610        }
611    };
612}
613
614#[cfg(not(target_family = "wasm"))]
615macro_rules! impl_state_store {
616    ({ $($body:tt)* }) => {
617        impl IndexeddbStateStore {
618            $($body)*
619        }
620    };
621}
622
623impl_state_store!({
624    async fn get_kv_data(&self, key: StateStoreDataKey<'_>) -> Result<Option<StateStoreDataValue>> {
625        let encoded_key = self.encode_kv_data_key(key);
626
627        let value = self
628            .inner
629            .transaction(keys::KV)
630            .with_mode(TransactionMode::Readonly)
631            .build()?
632            .object_store(keys::KV)?
633            .get(&encoded_key)
634            .await?;
635
636        let value = match key {
637            StateStoreDataKey::SyncToken => value
638                .map(|f| self.deserialize_value::<String>(&f))
639                .transpose()?
640                .map(StateStoreDataValue::SyncToken),
641            StateStoreDataKey::SupportedVersions => value
642                .map(|f| self.deserialize_value::<TtlValue<SupportedVersionsResponse>>(&f))
643                .transpose()?
644                .map(StateStoreDataValue::SupportedVersions),
645            StateStoreDataKey::WellKnown => value
646                .map(|f| self.deserialize_value::<TtlValue<Option<WellKnownResponse>>>(&f))
647                .transpose()?
648                .map(StateStoreDataValue::WellKnown),
649            StateStoreDataKey::Filter(_) => value
650                .map(|f| self.deserialize_value::<String>(&f))
651                .transpose()?
652                .map(StateStoreDataValue::Filter),
653            StateStoreDataKey::UserAvatarUrl(_) => value
654                .map(|f| self.deserialize_value::<OwnedMxcUri>(&f))
655                .transpose()?
656                .map(StateStoreDataValue::UserAvatarUrl),
657            StateStoreDataKey::RecentlyVisitedRooms(_) => value
658                .map(|f| self.deserialize_value::<Vec<OwnedRoomId>>(&f))
659                .transpose()?
660                .map(StateStoreDataValue::RecentlyVisitedRooms),
661            StateStoreDataKey::UtdHookManagerData => value
662                .map(|f| self.deserialize_value::<GrowableBloom>(&f))
663                .transpose()?
664                .map(StateStoreDataValue::UtdHookManagerData),
665            StateStoreDataKey::OneTimeKeyAlreadyUploaded => value
666                .map(|f| self.deserialize_value::<bool>(&f))
667                .transpose()?
668                .map(|_| StateStoreDataValue::OneTimeKeyAlreadyUploaded),
669            StateStoreDataKey::ComposerDraft(_, _) => value
670                .map(|f| self.deserialize_value::<ComposerDraft>(&f))
671                .transpose()?
672                .map(StateStoreDataValue::ComposerDraft),
673            StateStoreDataKey::SeenKnockRequests(_) => value
674                .map(|f| self.deserialize_value::<BTreeMap<OwnedEventId, OwnedUserId>>(&f))
675                .transpose()?
676                .map(StateStoreDataValue::SeenKnockRequests),
677            StateStoreDataKey::ThreadSubscriptionsCatchupTokens => value
678                .map(|f| self.deserialize_value::<Vec<ThreadSubscriptionCatchupToken>>(&f))
679                .transpose()?
680                .map(StateStoreDataValue::ThreadSubscriptionsCatchupTokens),
681            StateStoreDataKey::HomeserverCapabilities => value
682                .map(|f| self.deserialize_value::<TtlValue<Capabilities>>(&f))
683                .transpose()?
684                .map(StateStoreDataValue::HomeserverCapabilities),
685        };
686
687        Ok(value)
688    }
689
690    async fn set_kv_data(
691        &self,
692        key: StateStoreDataKey<'_>,
693        value: StateStoreDataValue,
694    ) -> Result<()> {
695        let encoded_key = self.encode_kv_data_key(key);
696
697        let serialized_value = match key {
698            StateStoreDataKey::SyncToken => self
699                .serialize_value(&value.into_sync_token().expect("Session data not a sync token")),
700            StateStoreDataKey::SupportedVersions => self.serialize_value(
701                &value
702                    .into_supported_versions()
703                    .expect("Session data not containing supported versions"),
704            ),
705            StateStoreDataKey::WellKnown => self.serialize_value(
706                &value.into_well_known().expect("Session data not containing well-known"),
707            ),
708            StateStoreDataKey::Filter(_) => {
709                self.serialize_value(&value.into_filter().expect("Session data not a filter"))
710            }
711            StateStoreDataKey::UserAvatarUrl(_) => self.serialize_value(
712                &value.into_user_avatar_url().expect("Session data not an user avatar url"),
713            ),
714            StateStoreDataKey::RecentlyVisitedRooms(_) => self.serialize_value(
715                &value
716                    .into_recently_visited_rooms()
717                    .expect("Session data not a recently visited room list"),
718            ),
719            StateStoreDataKey::UtdHookManagerData => self.serialize_value(
720                &value.into_utd_hook_manager_data().expect("Session data not UtdHookManagerData"),
721            ),
722            StateStoreDataKey::OneTimeKeyAlreadyUploaded => self.serialize_value(&true),
723            StateStoreDataKey::ComposerDraft(_, _) => self.serialize_value(
724                &value.into_composer_draft().expect("Session data not a composer draft"),
725            ),
726            StateStoreDataKey::SeenKnockRequests(_) => self.serialize_value(
727                &value
728                    .into_seen_knock_requests()
729                    .expect("Session data is not a set of seen knock request ids"),
730            ),
731            StateStoreDataKey::ThreadSubscriptionsCatchupTokens => self.serialize_value(
732                &value
733                    .into_thread_subscriptions_catchup_tokens()
734                    .expect("Session data is not a list of thread subscription catchup tokens"),
735            ),
736            StateStoreDataKey::HomeserverCapabilities => self.serialize_value(
737                &value
738                    .into_homeserver_capabilities()
739                    .expect("Session data is not a homeserver capabilities"),
740            ),
741        };
742
743        let tx = self.inner.transaction(keys::KV).with_mode(TransactionMode::Readwrite).build()?;
744
745        let obj = tx.object_store(keys::KV)?;
746
747        obj.put(&serialized_value?).with_key(encoded_key).build()?;
748
749        tx.commit().await?;
750
751        Ok(())
752    }
753
754    async fn remove_kv_data(&self, key: StateStoreDataKey<'_>) -> Result<()> {
755        let encoded_key = self.encode_kv_data_key(key);
756
757        let tx = self.inner.transaction(keys::KV).with_mode(TransactionMode::Readwrite).build()?;
758        let obj = tx.object_store(keys::KV)?;
759
760        obj.delete(&encoded_key).build()?;
761
762        tx.commit().await?;
763
764        Ok(())
765    }
766
767    async fn save_changes(&self, changes: &StateChanges) -> Result<()> {
768        let mut stores: HashSet<&'static str> = [
769            (changes.sync_token.is_some(), keys::KV),
770            (!changes.ambiguity_maps.is_empty(), keys::DISPLAY_NAMES),
771            (!changes.account_data.is_empty(), keys::ACCOUNT_DATA),
772            (!changes.presence.is_empty(), keys::PRESENCE),
773            (
774                !changes.profiles.is_empty() || !changes.profiles_to_delete.is_empty(),
775                keys::PROFILES,
776            ),
777            (!changes.room_account_data.is_empty(), keys::ROOM_ACCOUNT_DATA),
778            (!changes.receipts.is_empty(), keys::ROOM_EVENT_RECEIPTS),
779            (!changes.global_profiles.is_empty(), keys::GLOBAL_PROFILES),
780        ]
781        .iter()
782        .filter_map(|(id, key)| if *id { Some(*key) } else { None })
783        .collect();
784
785        if !changes.state.is_empty() {
786            stores.extend([
787                keys::ROOM_STATE,
788                keys::USER_IDS,
789                keys::STRIPPED_USER_IDS,
790                keys::STRIPPED_ROOM_STATE,
791                keys::PROFILES,
792            ]);
793        }
794
795        if !changes.redactions.is_empty() {
796            stores.extend([keys::ROOM_STATE, keys::ROOM_INFOS]);
797        }
798
799        if !changes.room_infos.is_empty() {
800            stores.insert(keys::ROOM_INFOS);
801        }
802
803        if !changes.stripped_state.is_empty() {
804            stores.extend([keys::STRIPPED_ROOM_STATE, keys::STRIPPED_USER_IDS]);
805        }
806
807        if !changes.receipts.is_empty() {
808            stores.extend([keys::ROOM_EVENT_RECEIPTS, keys::ROOM_USER_RECEIPTS])
809        }
810
811        if stores.is_empty() {
812            // nothing to do, quit early
813            return Ok(());
814        }
815
816        let stores: Vec<&'static str> = stores.into_iter().collect();
817        let tx = self.inner.transaction(stores).with_mode(TransactionMode::Readwrite).build()?;
818
819        if let Some(s) = &changes.sync_token {
820            tx.object_store(keys::KV)?
821                .put(&self.serialize_value(s)?)
822                .with_key(self.encode_kv_data_key(StateStoreDataKey::SyncToken))
823                .build()?;
824        }
825
826        if !changes.ambiguity_maps.is_empty() {
827            let store = tx.object_store(keys::DISPLAY_NAMES)?;
828            for (room_id, ambiguity_maps) in &changes.ambiguity_maps {
829                for (display_name, map) in ambiguity_maps {
830                    let key = self.encode_key(
831                        keys::DISPLAY_NAMES,
832                        (
833                            room_id,
834                            display_name
835                                .as_normalized_str()
836                                .unwrap_or_else(|| display_name.as_raw_str()),
837                        ),
838                    );
839
840                    store.put(&self.serialize_value(&map)?).with_key(key).build()?;
841                }
842            }
843        }
844
845        if !changes.account_data.is_empty() {
846            let store = tx.object_store(keys::ACCOUNT_DATA)?;
847            for (event_type, event) in &changes.account_data {
848                store
849                    .put(&self.serialize_value(&event)?)
850                    .with_key(self.encode_key(keys::ACCOUNT_DATA, event_type))
851                    .build()?;
852            }
853        }
854
855        if !changes.room_account_data.is_empty() {
856            let store = tx.object_store(keys::ROOM_ACCOUNT_DATA)?;
857            for (room, events) in &changes.room_account_data {
858                for (event_type, event) in events {
859                    let key = self.encode_key(keys::ROOM_ACCOUNT_DATA, (room, event_type));
860                    store.put(&self.serialize_value(&event)?).with_key(key).build()?;
861                }
862            }
863        }
864
865        if !changes.state.is_empty() {
866            let state = tx.object_store(keys::ROOM_STATE)?;
867            let profiles = tx.object_store(keys::PROFILES)?;
868            let user_ids = tx.object_store(keys::USER_IDS)?;
869            let stripped_state = tx.object_store(keys::STRIPPED_ROOM_STATE)?;
870            let stripped_user_ids = tx.object_store(keys::STRIPPED_USER_IDS)?;
871
872            for (room, user_ids) in &changes.profiles_to_delete {
873                for user_id in user_ids {
874                    let key = self.encode_key(keys::PROFILES, (room, user_id));
875                    profiles.delete(&key).build()?;
876                }
877            }
878
879            for (room, event_types) in &changes.state {
880                let profile_changes = changes.profiles.get(room);
881
882                for (event_type, events) in event_types {
883                    for (state_key, raw_event) in events {
884                        let key = self.encode_key(keys::ROOM_STATE, (room, event_type, state_key));
885                        state
886                            .put(&self.serialize_value(&raw_event)?)
887                            .with_key(key.clone())
888                            .build()?;
889                        stripped_state.delete(&key).build()?;
890
891                        if *event_type == StateEventType::RoomMember {
892                            let event =
893                                match raw_event.deserialize_as_unchecked::<SyncRoomMemberEvent>() {
894                                    Ok(ev) => ev,
895                                    Err(e) => {
896                                        let event_id: Option<String> =
897                                            raw_event.get_field("event_id").ok().flatten();
898                                        debug!(event_id, "Failed to deserialize member event: {e}");
899                                        continue;
900                                    }
901                                };
902
903                            let key = (room, state_key);
904
905                            stripped_user_ids
906                                .delete(&self.encode_key(keys::STRIPPED_USER_IDS, key))
907                                .build()?;
908
909                            user_ids
910                                .put(&self.serialize_value(&RoomMember::from(&event))?)
911                                .with_key(self.encode_key(keys::USER_IDS, key))
912                                .build()?;
913
914                            if let Some(profile) =
915                                profile_changes.and_then(|p| p.get(event.state_key()))
916                            {
917                                profiles
918                                    .put(&self.serialize_value(&profile)?)
919                                    .with_key(self.encode_key(keys::PROFILES, key))
920                                    .build()?;
921                            }
922                        }
923                    }
924                }
925            }
926        }
927
928        if !changes.room_infos.is_empty() {
929            let room_infos = tx.object_store(keys::ROOM_INFOS)?;
930            for (room_id, room_info) in &changes.room_infos {
931                room_infos
932                    .put(&self.serialize_value(&room_info)?)
933                    .with_key(self.encode_key(keys::ROOM_INFOS, room_id))
934                    .build()?;
935            }
936        }
937
938        if !changes.presence.is_empty() {
939            let store = tx.object_store(keys::PRESENCE)?;
940            for (sender, event) in &changes.presence {
941                store
942                    .put(&self.serialize_value(&event)?)
943                    .with_key(self.encode_key(keys::PRESENCE, sender))
944                    .build()?;
945            }
946        }
947
948        if !changes.stripped_state.is_empty() {
949            let store = tx.object_store(keys::STRIPPED_ROOM_STATE)?;
950            let user_ids = tx.object_store(keys::STRIPPED_USER_IDS)?;
951
952            for (room, event_types) in &changes.stripped_state {
953                for (event_type, events) in event_types {
954                    for (state_key, raw_event) in events {
955                        let key = self
956                            .encode_key(keys::STRIPPED_ROOM_STATE, (room, event_type, state_key));
957                        store.put(&self.serialize_value(&raw_event)?).with_key(key).build()?;
958
959                        if *event_type == StateEventType::RoomMember {
960                            let event = match raw_event
961                                .deserialize_as_unchecked::<StrippedRoomMemberEvent>()
962                            {
963                                Ok(ev) => ev,
964                                Err(e) => {
965                                    let event_id: Option<String> =
966                                        raw_event.get_field("event_id").ok().flatten();
967                                    debug!(
968                                        event_id,
969                                        "Failed to deserialize stripped member event: {e}"
970                                    );
971                                    continue;
972                                }
973                            };
974
975                            let key = (room, state_key);
976
977                            user_ids
978                                .put(&self.serialize_value(&RoomMember::from(&event))?)
979                                .with_key(self.encode_key(keys::STRIPPED_USER_IDS, key))
980                                .build()?;
981                        }
982                    }
983                }
984            }
985        }
986
987        if !changes.receipts.is_empty() {
988            let room_user_receipts = tx.object_store(keys::ROOM_USER_RECEIPTS)?;
989            let room_event_receipts = tx.object_store(keys::ROOM_EVENT_RECEIPTS)?;
990
991            for (room, content) in &changes.receipts {
992                for (event_id, receipts) in &content.0 {
993                    for (receipt_type, receipts) in receipts {
994                        for (user_id, receipt) in receipts {
995                            let key = match receipt.thread.as_str() {
996                                Some(thread_id) => self.encode_key(
997                                    keys::ROOM_USER_RECEIPTS,
998                                    (room, receipt_type, thread_id, user_id),
999                                ),
1000                                None => self.encode_key(
1001                                    keys::ROOM_USER_RECEIPTS,
1002                                    (room, receipt_type, user_id),
1003                                ),
1004                            };
1005
1006                            if let Some((old_event, _)) =
1007                                room_user_receipts.get(&key).await?.and_then(|f| {
1008                                    self.deserialize_value::<(OwnedEventId, Receipt)>(&f).ok()
1009                                })
1010                            {
1011                                let key = match receipt.thread.as_str() {
1012                                    Some(thread_id) => self.encode_key(
1013                                        keys::ROOM_EVENT_RECEIPTS,
1014                                        (room, receipt_type, thread_id, old_event, user_id),
1015                                    ),
1016                                    None => self.encode_key(
1017                                        keys::ROOM_EVENT_RECEIPTS,
1018                                        (room, receipt_type, old_event, user_id),
1019                                    ),
1020                                };
1021                                room_event_receipts.delete(&key).build()?;
1022                            }
1023
1024                            room_user_receipts
1025                                .put(&self.serialize_value(&(event_id, receipt))?)
1026                                .with_key(key)
1027                                .build()?;
1028
1029                            // Add the receipt to the room event receipts
1030                            let key = match receipt.thread.as_str() {
1031                                Some(thread_id) => self.encode_key(
1032                                    keys::ROOM_EVENT_RECEIPTS,
1033                                    (room, receipt_type, thread_id, event_id, user_id),
1034                                ),
1035                                None => self.encode_key(
1036                                    keys::ROOM_EVENT_RECEIPTS,
1037                                    (room, receipt_type, event_id, user_id),
1038                                ),
1039                            };
1040                            room_event_receipts
1041                                .put(&self.serialize_value(&(user_id, receipt))?)
1042                                .with_key(key)
1043                                .build()?;
1044                        }
1045                    }
1046                }
1047            }
1048        }
1049
1050        if !changes.redactions.is_empty() {
1051            let state = tx.object_store(keys::ROOM_STATE)?;
1052            let room_info = tx.object_store(keys::ROOM_INFOS)?;
1053
1054            for (room_id, redactions) in &changes.redactions {
1055                let range = self.encode_to_range(keys::ROOM_STATE, room_id);
1056                let Some(mut cursor) = state.open_cursor().with_query(&range).await? else {
1057                    continue;
1058                };
1059
1060                let mut redaction_rules = None;
1061
1062                while let Some(value) = cursor.next_record().await? {
1063                    let Some(key) = cursor.key::<JsValue>()? else {
1064                        break;
1065                    };
1066
1067                    let raw_evt = self.deserialize_value::<Raw<AnySyncStateEvent>>(&value)?;
1068                    if let Ok(Some(event_id)) = raw_evt.get_field::<OwnedEventId>("event_id")
1069                        && let Some(redaction) = redactions.get(&event_id)
1070                    {
1071                        let redaction_rules = match &redaction_rules {
1072                            Some(r) => r,
1073                            None => {
1074                                let value = room_info
1075                                    .get(&self.encode_key(keys::ROOM_INFOS, room_id))
1076                                    .await?
1077                                    .and_then(|f| self.deserialize_value::<RoomInfo>(&f).ok())
1078                                    .map(|info| info.room_version_rules_or_default())
1079                                    .unwrap_or_else(|| {
1080                                        warn!(
1081                                            ?room_id,
1082                                            "Unable to get the room version rules, \
1083                                             defaulting to rules for room version \
1084                                             {ROOM_VERSION_FALLBACK}"
1085                                        );
1086                                        ROOM_VERSION_RULES_FALLBACK
1087                                    })
1088                                    .redaction;
1089                                redaction_rules.get_or_insert(value)
1090                            }
1091                        };
1092
1093                        let redacted = redact(
1094                            raw_evt.deserialize_as::<CanonicalJsonObject>()?,
1095                            redaction_rules,
1096                            Some(RedactedBecause::from_raw_event(redaction)?),
1097                        )
1098                        .map_err(StoreError::Redaction)?;
1099                        state.put(&self.serialize_value(&redacted)?).with_key(key).build()?;
1100                    }
1101                }
1102            }
1103        }
1104
1105        if !changes.global_profiles.is_empty() {
1106            let store = tx.object_store(keys::GLOBAL_PROFILES)?;
1107            for (user_id, profile_update) in &changes.global_profiles {
1108                let key = self.encode_key(keys::GLOBAL_PROFILES, user_id);
1109                match profile_update {
1110                    UserProfileUpdate::Updated(profile_changes) => {
1111                        let existing: Option<UserProfile> = store
1112                            .get(&key)
1113                            .await?
1114                            .map(|f| self.deserialize_value(&f))
1115                            .transpose()?;
1116
1117                        let mut profile = existing.unwrap_or_default();
1118                        profile.apply(profile_changes.clone());
1119
1120                        store.put(&self.serialize_value(&profile)?).with_key(key).build()?;
1121                    }
1122                    UserProfileUpdate::Dropped => {
1123                        store.delete(&key).build()?;
1124                    }
1125                    _ => {
1126                        warn!(%user_id, "Unhandled UserProfileUpdate variant; ignoring");
1127                    }
1128                }
1129            }
1130        }
1131
1132        tx.commit().await.map_err(|e| e.into())
1133    }
1134
1135    async fn get_presence_event(&self, user_id: &UserId) -> Result<Option<Raw<PresenceEvent>>> {
1136        self.inner
1137            .transaction(keys::PRESENCE)
1138            .with_mode(TransactionMode::Readonly)
1139            .build()?
1140            .object_store(keys::PRESENCE)?
1141            .get(&self.encode_key(keys::PRESENCE, user_id))
1142            .await?
1143            .map(|f| self.deserialize_value(&f))
1144            .transpose()
1145    }
1146
1147    async fn get_presence_events(
1148        &self,
1149        user_ids: &[OwnedUserId],
1150    ) -> Result<Vec<Raw<PresenceEvent>>> {
1151        if user_ids.is_empty() {
1152            return Ok(Vec::new());
1153        }
1154
1155        let txn =
1156            self.inner.transaction(keys::PRESENCE).with_mode(TransactionMode::Readonly).build()?;
1157        let store = txn.object_store(keys::PRESENCE)?;
1158
1159        let mut events = Vec::with_capacity(user_ids.len());
1160
1161        for user_id in user_ids {
1162            if let Some(event) = store
1163                .get(&self.encode_key(keys::PRESENCE, user_id))
1164                .await?
1165                .map(|f| self.deserialize_value(&f))
1166                .transpose()?
1167            {
1168                events.push(event)
1169            }
1170        }
1171
1172        Ok(events)
1173    }
1174
1175    async fn get_state_event(
1176        &self,
1177        room_id: &RoomId,
1178        event_type: StateEventType,
1179        state_key: &str,
1180    ) -> Result<Option<RawAnySyncOrStrippedState>> {
1181        Ok(self
1182            .get_state_events_for_keys(room_id, event_type, &[state_key])
1183            .await?
1184            .into_iter()
1185            .next())
1186    }
1187
1188    async fn get_state_events(
1189        &self,
1190        room_id: &RoomId,
1191        event_type: StateEventType,
1192    ) -> Result<Vec<RawAnySyncOrStrippedState>> {
1193        let stripped_range =
1194            self.encode_to_range(keys::STRIPPED_ROOM_STATE, (room_id, &event_type));
1195        let stripped_events = self
1196            .inner
1197            .transaction(keys::STRIPPED_ROOM_STATE)
1198            .with_mode(TransactionMode::Readonly)
1199            .build()?
1200            .object_store(keys::STRIPPED_ROOM_STATE)?
1201            .get_all()
1202            .with_query(&stripped_range)
1203            .await?
1204            .filter_map(Result::ok)
1205            .filter_map(|f| {
1206                self.deserialize_value(&f).ok().map(RawAnySyncOrStrippedState::Stripped)
1207            })
1208            .collect::<Vec<_>>();
1209
1210        if !stripped_events.is_empty() {
1211            return Ok(stripped_events);
1212        }
1213
1214        let range = self.encode_to_range(keys::ROOM_STATE, (room_id, event_type));
1215        Ok(self
1216            .inner
1217            .transaction(keys::ROOM_STATE)
1218            .with_mode(TransactionMode::Readonly)
1219            .build()?
1220            .object_store(keys::ROOM_STATE)?
1221            .get_all()
1222            .with_query(&range)
1223            .await?
1224            .filter_map(Result::ok)
1225            .filter_map(|f| self.deserialize_value(&f).ok().map(RawAnySyncOrStrippedState::Sync))
1226            .collect::<Vec<_>>())
1227    }
1228
1229    async fn get_state_events_for_keys(
1230        &self,
1231        room_id: &RoomId,
1232        event_type: StateEventType,
1233        state_keys: &[&str],
1234    ) -> Result<Vec<RawAnySyncOrStrippedState>> {
1235        if state_keys.is_empty() {
1236            return Ok(Vec::new());
1237        }
1238
1239        let mut events = Vec::with_capacity(state_keys.len());
1240
1241        {
1242            let txn = self
1243                .inner
1244                .transaction(keys::STRIPPED_ROOM_STATE)
1245                .with_mode(TransactionMode::Readonly)
1246                .build()?;
1247            let store = txn.object_store(keys::STRIPPED_ROOM_STATE)?;
1248
1249            for state_key in state_keys {
1250                if let Some(event) =
1251                    store
1252                        .get(&self.encode_key(
1253                            keys::STRIPPED_ROOM_STATE,
1254                            (room_id, &event_type, state_key),
1255                        ))
1256                        .await?
1257                        .map(|f| self.deserialize_value(&f))
1258                        .transpose()?
1259                {
1260                    events.push(RawAnySyncOrStrippedState::Stripped(event));
1261                }
1262            }
1263
1264            if !events.is_empty() {
1265                return Ok(events);
1266            }
1267        }
1268
1269        let txn = self
1270            .inner
1271            .transaction(keys::ROOM_STATE)
1272            .with_mode(TransactionMode::Readonly)
1273            .build()?;
1274        let store = txn.object_store(keys::ROOM_STATE)?;
1275
1276        for state_key in state_keys {
1277            if let Some(event) = store
1278                .get(&self.encode_key(keys::ROOM_STATE, (room_id, &event_type, state_key)))
1279                .await?
1280                .map(|f| self.deserialize_value(&f))
1281                .transpose()?
1282            {
1283                events.push(RawAnySyncOrStrippedState::Sync(event));
1284            }
1285        }
1286
1287        Ok(events)
1288    }
1289
1290    async fn get_profile(
1291        &self,
1292        room_id: &RoomId,
1293        user_id: &UserId,
1294    ) -> Result<Option<MinimalRoomMemberEvent>> {
1295        self.inner
1296            .transaction(keys::PROFILES)
1297            .with_mode(TransactionMode::Readonly)
1298            .build()?
1299            .object_store(keys::PROFILES)?
1300            .get(&self.encode_key(keys::PROFILES, (room_id, user_id)))
1301            .await?
1302            .map(|f| self.deserialize_value(&f))
1303            .transpose()
1304    }
1305
1306    async fn get_profiles<'a>(
1307        &self,
1308        room_id: &RoomId,
1309        user_ids: &'a [OwnedUserId],
1310    ) -> Result<BTreeMap<&'a UserId, MinimalRoomMemberEvent>> {
1311        if user_ids.is_empty() {
1312            return Ok(BTreeMap::new());
1313        }
1314
1315        let txn =
1316            self.inner.transaction(keys::PROFILES).with_mode(TransactionMode::Readonly).build()?;
1317        let store = txn.object_store(keys::PROFILES)?;
1318
1319        let mut profiles = BTreeMap::new();
1320        for user_id in user_ids {
1321            if let Some(profile) = store
1322                .get(&self.encode_key(keys::PROFILES, (room_id, user_id)))
1323                .await?
1324                .map(|f| self.deserialize_value(&f))
1325                .transpose()?
1326            {
1327                profiles.insert(user_id.as_ref(), profile);
1328            }
1329        }
1330
1331        Ok(profiles)
1332    }
1333
1334    async fn get_room_infos(&self, room_load_settings: &RoomLoadSettings) -> Result<Vec<RoomInfo>> {
1335        let transaction = self
1336            .inner
1337            .transaction(keys::ROOM_INFOS)
1338            .with_mode(TransactionMode::Readonly)
1339            .build()?;
1340
1341        let object_store = transaction.object_store(keys::ROOM_INFOS)?;
1342
1343        Ok(match room_load_settings {
1344            RoomLoadSettings::All => object_store
1345                .get_all()
1346                .await?
1347                .map(|room_info| self.deserialize_value::<RoomInfo>(&room_info?))
1348                .collect::<Result<_>>()?,
1349
1350            RoomLoadSettings::One(room_id) => {
1351                match object_store.get(&self.encode_key(keys::ROOM_INFOS, room_id)).await? {
1352                    Some(room_info) => vec![self.deserialize_value::<RoomInfo>(&room_info)?],
1353                    None => vec![],
1354                }
1355            }
1356        })
1357    }
1358
1359    async fn get_users_with_display_name(
1360        &self,
1361        room_id: &RoomId,
1362        display_name: &DisplayName,
1363    ) -> Result<BTreeSet<OwnedUserId>> {
1364        self.inner
1365            .transaction(keys::DISPLAY_NAMES)
1366            .with_mode(TransactionMode::Readonly)
1367            .build()?
1368            .object_store(keys::DISPLAY_NAMES)?
1369            .get(&self.encode_key(
1370                keys::DISPLAY_NAMES,
1371                (
1372                    room_id,
1373                    display_name.as_normalized_str().unwrap_or_else(|| display_name.as_raw_str()),
1374                ),
1375            ))
1376            .await?
1377            .map(|f| self.deserialize_value::<BTreeSet<OwnedUserId>>(&f))
1378            .unwrap_or_else(|| Ok(Default::default()))
1379    }
1380
1381    async fn get_users_with_display_names<'a>(
1382        &self,
1383        room_id: &RoomId,
1384        display_names: &'a [DisplayName],
1385    ) -> Result<HashMap<&'a DisplayName, BTreeSet<OwnedUserId>>> {
1386        let mut map = HashMap::new();
1387
1388        if display_names.is_empty() {
1389            return Ok(map);
1390        }
1391
1392        let txn = self
1393            .inner
1394            .transaction(keys::DISPLAY_NAMES)
1395            .with_mode(TransactionMode::Readonly)
1396            .build()?;
1397        let store = txn.object_store(keys::DISPLAY_NAMES)?;
1398
1399        for display_name in display_names {
1400            if let Some(user_ids) = store
1401                .get(
1402                    &self.encode_key(
1403                        keys::DISPLAY_NAMES,
1404                        (
1405                            room_id,
1406                            display_name
1407                                .as_normalized_str()
1408                                .unwrap_or_else(|| display_name.as_raw_str()),
1409                        ),
1410                    ),
1411                )
1412                .await?
1413                .map(|f| self.deserialize_value::<BTreeSet<OwnedUserId>>(&f))
1414                .transpose()?
1415            {
1416                map.insert(display_name, user_ids);
1417            }
1418        }
1419
1420        Ok(map)
1421    }
1422
1423    async fn get_account_data_event(
1424        &self,
1425        event_type: GlobalAccountDataEventType,
1426    ) -> Result<Option<Raw<AnyGlobalAccountDataEvent>>> {
1427        self.inner
1428            .transaction(keys::ACCOUNT_DATA)
1429            .with_mode(TransactionMode::Readonly)
1430            .build()?
1431            .object_store(keys::ACCOUNT_DATA)?
1432            .get(&self.encode_key(keys::ACCOUNT_DATA, event_type))
1433            .await?
1434            .map(|f| self.deserialize_value(&f))
1435            .transpose()
1436    }
1437
1438    async fn get_room_account_data_event(
1439        &self,
1440        room_id: &RoomId,
1441        event_type: RoomAccountDataEventType,
1442    ) -> Result<Option<Raw<AnyRoomAccountDataEvent>>> {
1443        self.inner
1444            .transaction(keys::ROOM_ACCOUNT_DATA)
1445            .with_mode(TransactionMode::Readonly)
1446            .build()?
1447            .object_store(keys::ROOM_ACCOUNT_DATA)?
1448            .get(&self.encode_key(keys::ROOM_ACCOUNT_DATA, (room_id, event_type)))
1449            .await?
1450            .map(|f| self.deserialize_value(&f))
1451            .transpose()
1452    }
1453
1454    async fn get_user_room_receipt_event(
1455        &self,
1456        room_id: &RoomId,
1457        receipt_type: ReceiptType,
1458        thread: ReceiptThread,
1459        user_id: &UserId,
1460    ) -> Result<Option<(OwnedEventId, Receipt)>> {
1461        let key = match thread.as_str() {
1462            Some(thread_id) => self
1463                .encode_key(keys::ROOM_USER_RECEIPTS, (room_id, receipt_type, thread_id, user_id)),
1464            None => self.encode_key(keys::ROOM_USER_RECEIPTS, (room_id, receipt_type, user_id)),
1465        };
1466        self.inner
1467            .transaction(keys::ROOM_USER_RECEIPTS)
1468            .with_mode(TransactionMode::Readonly)
1469            .build()?
1470            .object_store(keys::ROOM_USER_RECEIPTS)?
1471            .get(&key)
1472            .await?
1473            .map(|f| self.deserialize_value(&f))
1474            .transpose()
1475    }
1476
1477    async fn get_event_room_receipt_events(
1478        &self,
1479        room_id: &RoomId,
1480        receipt_type: ReceiptType,
1481        thread: ReceiptThread,
1482        event_id: &EventId,
1483    ) -> Result<Vec<(OwnedUserId, Receipt)>> {
1484        let range = match thread.as_str() {
1485            Some(thread_id) => self.encode_to_range(
1486                keys::ROOM_EVENT_RECEIPTS,
1487                (room_id, receipt_type, thread_id, event_id),
1488            ),
1489            None => {
1490                self.encode_to_range(keys::ROOM_EVENT_RECEIPTS, (room_id, receipt_type, event_id))
1491            }
1492        };
1493        let tx = self
1494            .inner
1495            .transaction(keys::ROOM_EVENT_RECEIPTS)
1496            .with_mode(TransactionMode::Readonly)
1497            .build()?;
1498        let store = tx.object_store(keys::ROOM_EVENT_RECEIPTS)?;
1499
1500        Ok(store
1501            .get_all()
1502            .with_query(&range)
1503            .await?
1504            .filter_map(Result::ok)
1505            .filter_map(|f| self.deserialize_value(&f).ok())
1506            .collect::<Vec<_>>())
1507    }
1508
1509    async fn get_custom_value(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
1510        let jskey = &JsValue::from_str(core::str::from_utf8(key).map_err(StoreError::Codec)?);
1511        self.get_custom_value_for_js(jskey).await
1512    }
1513
1514    async fn set_custom_value(&self, key: &[u8], value: Vec<u8>) -> Result<Option<Vec<u8>>> {
1515        let jskey = JsValue::from_str(core::str::from_utf8(key).map_err(StoreError::Codec)?);
1516
1517        let prev = self.get_custom_value_for_js(&jskey).await?;
1518
1519        let tx =
1520            self.inner.transaction(keys::CUSTOM).with_mode(TransactionMode::Readwrite).build()?;
1521
1522        tx.object_store(keys::CUSTOM)?
1523            .put(&self.serialize_value(&value)?)
1524            .with_key(jskey)
1525            .build()?;
1526
1527        tx.commit().await.map_err(IndexeddbStateStoreError::from)?;
1528        Ok(prev)
1529    }
1530
1531    async fn remove_custom_value(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
1532        let jskey = JsValue::from_str(core::str::from_utf8(key).map_err(StoreError::Codec)?);
1533
1534        let prev = self.get_custom_value_for_js(&jskey).await?;
1535
1536        let tx =
1537            self.inner.transaction(keys::CUSTOM).with_mode(TransactionMode::Readwrite).build()?;
1538
1539        tx.object_store(keys::CUSTOM)?.delete(&jskey).build()?;
1540
1541        tx.commit().await.map_err(IndexeddbStateStoreError::from)?;
1542        Ok(prev)
1543    }
1544
1545    async fn remove_room(&self, room_id: &RoomId) -> Result<()> {
1546        // All the stores which use a RoomId as their key (and nothing additional).
1547        let direct_stores = [keys::ROOM_INFOS, keys::ROOM_SEND_QUEUE, keys::DEPENDENT_SEND_QUEUE];
1548
1549        // All the stores which use a RoomId as the first part of their key, but may
1550        // have some additional data in the key.
1551        let prefixed_stores = [
1552            keys::PROFILES,
1553            keys::DISPLAY_NAMES,
1554            keys::USER_IDS,
1555            keys::ROOM_STATE,
1556            keys::ROOM_ACCOUNT_DATA,
1557            keys::ROOM_EVENT_RECEIPTS,
1558            keys::ROOM_USER_RECEIPTS,
1559            keys::STRIPPED_ROOM_STATE,
1560            keys::STRIPPED_USER_IDS,
1561            keys::THREAD_SUBSCRIPTIONS,
1562        ];
1563
1564        let all_stores = {
1565            let mut v = Vec::new();
1566            v.extend(prefixed_stores);
1567            v.extend(direct_stores);
1568            v
1569        };
1570
1571        let tx =
1572            self.inner.transaction(all_stores).with_mode(TransactionMode::Readwrite).build()?;
1573
1574        for store_name in direct_stores {
1575            tx.object_store(store_name)?.delete(&self.encode_key(store_name, room_id)).build()?;
1576        }
1577
1578        for store_name in prefixed_stores {
1579            let store = tx.object_store(store_name)?;
1580            let range = self.encode_to_range(store_name, room_id);
1581            for key in store.get_all_keys::<JsValue>().with_query(&range).await? {
1582                store.delete(&key?).build()?;
1583            }
1584        }
1585
1586        tx.commit().await.map_err(|e| e.into())
1587    }
1588
1589    async fn get_user_ids(
1590        &self,
1591        room_id: &RoomId,
1592        memberships: RoomMemberships,
1593    ) -> Result<Vec<OwnedUserId>> {
1594        let ids = self.get_user_ids_inner(room_id, memberships, true).await?;
1595        if !ids.is_empty() {
1596            return Ok(ids);
1597        }
1598        self.get_user_ids_inner(room_id, memberships, false).await
1599    }
1600
1601    async fn save_send_queue_request(
1602        &self,
1603        room_id: &RoomId,
1604        transaction_id: OwnedTransactionId,
1605        created_at: MilliSecondsSinceUnixEpoch,
1606        kind: QueuedRequestKind,
1607        priority: usize,
1608    ) -> Result<()> {
1609        let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1610
1611        let tx = self
1612            .inner
1613            .transaction(keys::ROOM_SEND_QUEUE)
1614            .with_mode(TransactionMode::Readwrite)
1615            .build()?;
1616
1617        let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1618
1619        // We store an encoded vector of the queued requests, with their transaction
1620        // ids.
1621
1622        // Reload the previous vector for this room, or create an empty one.
1623        let prev = obj.get(&encoded_key).await?;
1624
1625        let mut prev = prev.map_or_else(
1626            || Ok(Vec::new()),
1627            |val| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val),
1628        )?;
1629
1630        // Push the new request.
1631        prev.push(PersistedQueuedRequest {
1632            room_id: room_id.to_owned(),
1633            kind: Some(kind),
1634            transaction_id,
1635            error: None,
1636            is_wedged: None,
1637            event: None,
1638            priority: Some(priority),
1639            created_at,
1640        });
1641
1642        // Save the new vector into db.
1643        obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1644
1645        tx.commit().await?;
1646
1647        Ok(())
1648    }
1649
1650    async fn update_send_queue_request(
1651        &self,
1652        room_id: &RoomId,
1653        transaction_id: &TransactionId,
1654        kind: QueuedRequestKind,
1655    ) -> Result<bool> {
1656        let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1657
1658        let tx = self
1659            .inner
1660            .transaction(keys::ROOM_SEND_QUEUE)
1661            .with_mode(TransactionMode::Readwrite)
1662            .build()?;
1663
1664        let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1665
1666        // We store an encoded vector of the queued requests, with their transaction
1667        // ids.
1668
1669        // Reload the previous vector for this room, or create an empty one.
1670        let prev = obj.get(&encoded_key).await?;
1671
1672        let mut prev = prev.map_or_else(
1673            || Ok(Vec::new()),
1674            |val| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val),
1675        )?;
1676
1677        // Modify the one request.
1678        if let Some(entry) = prev.iter_mut().find(|entry| entry.transaction_id == transaction_id) {
1679            entry.kind = Some(kind);
1680            // Reset the error state.
1681            entry.error = None;
1682            // Remove migrated fields.
1683            entry.is_wedged = None;
1684            entry.event = None;
1685
1686            // Save the new vector into db.
1687            obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1688            tx.commit().await?;
1689
1690            Ok(true)
1691        } else {
1692            Ok(false)
1693        }
1694    }
1695
1696    async fn remove_send_queue_request(
1697        &self,
1698        room_id: &RoomId,
1699        transaction_id: &TransactionId,
1700    ) -> Result<bool> {
1701        let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1702
1703        let tx = self
1704            .inner
1705            .transaction([keys::ROOM_SEND_QUEUE, keys::DEPENDENT_SEND_QUEUE])
1706            .with_mode(TransactionMode::Readwrite)
1707            .build()?;
1708
1709        let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1710
1711        // We store an encoded vector of the queued requests, with their transaction
1712        // ids.
1713
1714        // Reload the previous vector for this room.
1715        if let Some(val) = obj.get(&encoded_key).await? {
1716            let mut prev = self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val)?;
1717            if let Some(pos) = prev.iter().position(|item| item.transaction_id == transaction_id) {
1718                prev.remove(pos);
1719
1720                if prev.is_empty() {
1721                    obj.delete(&encoded_key).build()?;
1722                } else {
1723                    obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1724                }
1725
1726                tx.commit().await?;
1727                return Ok(true);
1728            }
1729        }
1730
1731        Ok(false)
1732    }
1733
1734    async fn load_send_queue_requests(&self, room_id: &RoomId) -> Result<Vec<QueuedRequest>> {
1735        let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1736
1737        // We store an encoded vector of the queued requests, with their transaction
1738        // ids.
1739        let prev = self
1740            .inner
1741            .transaction(keys::ROOM_SEND_QUEUE)
1742            .with_mode(TransactionMode::Readwrite)
1743            .build()?
1744            .object_store(keys::ROOM_SEND_QUEUE)?
1745            .get(&encoded_key)
1746            .await?;
1747
1748        let mut prev = prev.map_or_else(
1749            || Ok(Vec::new()),
1750            |val| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val),
1751        )?;
1752
1753        // Inverted stable ordering on priority.
1754        prev.sort_by_key(|item| Reverse(item.priority.unwrap_or(0)));
1755
1756        Ok(prev.into_iter().filter_map(PersistedQueuedRequest::into_queued_request).collect())
1757    }
1758
1759    async fn update_send_queue_request_status(
1760        &self,
1761        room_id: &RoomId,
1762        transaction_id: &TransactionId,
1763        error: Option<QueueWedgeError>,
1764    ) -> Result<()> {
1765        let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1766
1767        let tx = self
1768            .inner
1769            .transaction(keys::ROOM_SEND_QUEUE)
1770            .with_mode(TransactionMode::Readwrite)
1771            .build()?;
1772
1773        let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1774
1775        if let Some(val) = obj.get(&encoded_key).await? {
1776            let mut prev = self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val)?;
1777            if let Some(request) =
1778                prev.iter_mut().find(|item| item.transaction_id == transaction_id)
1779            {
1780                request.is_wedged = None;
1781                request.error = error;
1782                obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1783            }
1784        }
1785
1786        tx.commit().await?;
1787
1788        Ok(())
1789    }
1790
1791    async fn load_rooms_with_unsent_requests(&self) -> Result<Vec<OwnedRoomId>> {
1792        let tx = self
1793            .inner
1794            .transaction(keys::ROOM_SEND_QUEUE)
1795            .with_mode(TransactionMode::Readwrite)
1796            .build()?;
1797
1798        let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1799
1800        let all_entries = obj
1801            .get_all()
1802            .await?
1803            .map(|item| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&item?))
1804            .collect::<Result<Vec<Vec<PersistedQueuedRequest>>, _>>()?
1805            .into_iter()
1806            .flat_map(|vec| vec.into_iter().map(|item| item.room_id))
1807            .collect::<BTreeSet<_>>();
1808
1809        Ok(all_entries.into_iter().collect())
1810    }
1811
1812    async fn save_dependent_queued_request(
1813        &self,
1814        room_id: &RoomId,
1815        parent_txn_id: &TransactionId,
1816        own_txn_id: ChildTransactionId,
1817        created_at: MilliSecondsSinceUnixEpoch,
1818        content: DependentQueuedRequestKind,
1819    ) -> Result<()> {
1820        let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1821
1822        let tx = self
1823            .inner
1824            .transaction(keys::DEPENDENT_SEND_QUEUE)
1825            .with_mode(TransactionMode::Readwrite)
1826            .build()?;
1827
1828        let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1829
1830        // We store an encoded vector of the dependent requests.
1831        // Reload the previous vector for this room, or create an empty one.
1832        let prev = obj.get(&encoded_key).await?;
1833
1834        let mut prev = prev.map_or_else(
1835            || Ok(Vec::new()),
1836            |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1837        )?;
1838
1839        // Push the new request.
1840        prev.push(DependentQueuedRequest {
1841            kind: content,
1842            parent_transaction_id: parent_txn_id.to_owned(),
1843            own_transaction_id: own_txn_id,
1844            parent_key: None,
1845            created_at,
1846        });
1847
1848        // Save the new vector into db.
1849        obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1850
1851        tx.commit().await?;
1852
1853        Ok(())
1854    }
1855
1856    async fn update_dependent_queued_request(
1857        &self,
1858        room_id: &RoomId,
1859        own_transaction_id: &ChildTransactionId,
1860        new_content: DependentQueuedRequestKind,
1861    ) -> Result<bool> {
1862        let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1863
1864        let tx = self
1865            .inner
1866            .transaction(keys::DEPENDENT_SEND_QUEUE)
1867            .with_mode(TransactionMode::Readwrite)
1868            .build()?;
1869
1870        let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1871
1872        // We store an encoded vector of the dependent requests.
1873        // Reload the previous vector for this room, or create an empty one.
1874        let prev = obj.get(&encoded_key).await?;
1875
1876        let mut prev = prev.map_or_else(
1877            || Ok(Vec::new()),
1878            |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1879        )?;
1880
1881        // Modify the dependent request, if found.
1882        let mut found = false;
1883        for entry in prev.iter_mut() {
1884            if entry.own_transaction_id == *own_transaction_id {
1885                found = true;
1886                entry.kind = new_content;
1887                break;
1888            }
1889        }
1890
1891        if found {
1892            obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1893            tx.commit().await?;
1894        }
1895
1896        Ok(found)
1897    }
1898
1899    async fn mark_dependent_queued_requests_as_ready(
1900        &self,
1901        room_id: &RoomId,
1902        parent_txn_id: &TransactionId,
1903        parent_key: SentRequestKey,
1904    ) -> Result<usize> {
1905        let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1906
1907        let tx = self
1908            .inner
1909            .transaction(keys::DEPENDENT_SEND_QUEUE)
1910            .with_mode(TransactionMode::Readwrite)
1911            .build()?;
1912
1913        let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1914
1915        // We store an encoded vector of the dependent requests.
1916        // Reload the previous vector for this room, or create an empty one.
1917        let prev = obj.get(&encoded_key).await?;
1918
1919        let mut prev = prev.map_or_else(
1920            || Ok(Vec::new()),
1921            |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1922        )?;
1923
1924        // Modify all requests that match.
1925        let mut num_updated = 0;
1926        for entry in prev.iter_mut().filter(|entry| entry.parent_transaction_id == parent_txn_id) {
1927            entry.parent_key = Some(parent_key.clone());
1928            num_updated += 1;
1929        }
1930
1931        if num_updated > 0 {
1932            obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1933            tx.commit().await?;
1934        }
1935
1936        Ok(num_updated)
1937    }
1938
1939    async fn remove_dependent_queued_request(
1940        &self,
1941        room_id: &RoomId,
1942        txn_id: &ChildTransactionId,
1943    ) -> Result<bool> {
1944        let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1945
1946        let tx = self
1947            .inner
1948            .transaction(keys::DEPENDENT_SEND_QUEUE)
1949            .with_mode(TransactionMode::Readwrite)
1950            .build()?;
1951
1952        let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1953
1954        // We store an encoded vector of the dependent requests.
1955        // Reload the previous vector for this room.
1956        if let Some(val) = obj.get(&encoded_key).await? {
1957            let mut prev = self.deserialize_value::<Vec<DependentQueuedRequest>>(&val)?;
1958            if let Some(pos) = prev.iter().position(|item| item.own_transaction_id == *txn_id) {
1959                prev.remove(pos);
1960
1961                if prev.is_empty() {
1962                    obj.delete(&encoded_key).build()?;
1963                } else {
1964                    obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1965                }
1966
1967                tx.commit().await?;
1968                return Ok(true);
1969            }
1970        }
1971
1972        Ok(false)
1973    }
1974
1975    async fn load_dependent_queued_requests(
1976        &self,
1977        room_id: &RoomId,
1978    ) -> Result<Vec<DependentQueuedRequest>> {
1979        let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1980
1981        // We store an encoded vector of the dependent requests.
1982        let prev = self
1983            .inner
1984            .transaction(keys::DEPENDENT_SEND_QUEUE)
1985            .with_mode(TransactionMode::Readwrite)
1986            .build()?
1987            .object_store(keys::DEPENDENT_SEND_QUEUE)?
1988            .get(&encoded_key)
1989            .await?;
1990
1991        prev.map_or_else(
1992            || Ok(Vec::new()),
1993            |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1994        )
1995    }
1996
1997    async fn upsert_thread_subscriptions(
1998        &self,
1999        updates: Vec<(&RoomId, &EventId, StoredThreadSubscription)>,
2000    ) -> Result<()> {
2001        let tx = self
2002            .inner
2003            .transaction(keys::THREAD_SUBSCRIPTIONS)
2004            .with_mode(TransactionMode::Readwrite)
2005            .build()?;
2006        let obj = tx.object_store(keys::THREAD_SUBSCRIPTIONS)?;
2007
2008        for (room_id, thread_id, subscription) in updates {
2009            let encoded_key = self.encode_key(keys::THREAD_SUBSCRIPTIONS, (room_id, thread_id));
2010            let mut new = PersistedThreadSubscription::from(subscription);
2011
2012            // See if there's a previous subscription.
2013            if let Some(previous_value) = obj.get(&encoded_key).await? {
2014                let previous: PersistedThreadSubscription =
2015                    self.deserialize_value(&previous_value)?;
2016
2017                // If the previous status is the same as the new one, don't do anything.
2018                if new == previous {
2019                    continue;
2020                }
2021                if !compare_thread_subscription_bump_stamps(
2022                    previous.bump_stamp,
2023                    &mut new.bump_stamp,
2024                ) {
2025                    continue;
2026                }
2027            }
2028
2029            let serialized_value = self.serialize_value(&new);
2030            obj.put(&serialized_value?).with_key(encoded_key).build()?;
2031        }
2032
2033        tx.commit().await?;
2034
2035        Ok(())
2036    }
2037
2038    async fn load_thread_subscription(
2039        &self,
2040        room: &RoomId,
2041        thread_id: &EventId,
2042    ) -> Result<Option<StoredThreadSubscription>> {
2043        let encoded_key = self.encode_key(keys::THREAD_SUBSCRIPTIONS, (room, thread_id));
2044
2045        let js_value = self
2046            .inner
2047            .transaction(keys::THREAD_SUBSCRIPTIONS)
2048            .with_mode(TransactionMode::Readonly)
2049            .build()?
2050            .object_store(keys::THREAD_SUBSCRIPTIONS)?
2051            .get(&encoded_key)
2052            .await?;
2053
2054        let Some(js_value) = js_value else {
2055            // We didn't have a previous subscription for this thread.
2056            return Ok(None);
2057        };
2058
2059        let sub: PersistedThreadSubscription = self.deserialize_value(&js_value)?;
2060
2061        let status = ThreadSubscriptionStatus::from_str(&sub.status).map_err(|_| {
2062            StoreError::InvalidData {
2063                details: format!(
2064                    "invalid thread status for room {room} and thread {thread_id}: {}",
2065                    sub.status
2066                ),
2067            }
2068        })?;
2069
2070        Ok(Some(StoredThreadSubscription { status, bump_stamp: sub.bump_stamp }))
2071    }
2072
2073    async fn remove_thread_subscription(&self, room: &RoomId, thread_id: &EventId) -> Result<()> {
2074        let encoded_key = self.encode_key(keys::THREAD_SUBSCRIPTIONS, (room, thread_id));
2075
2076        let transaction = self
2077            .inner
2078            .transaction(keys::THREAD_SUBSCRIPTIONS)
2079            .with_mode(TransactionMode::Readwrite)
2080            .build()?;
2081        transaction.object_store(keys::THREAD_SUBSCRIPTIONS)?.delete(&encoded_key).await?;
2082        transaction.commit().await?;
2083
2084        Ok(())
2085    }
2086
2087    async fn get_global_profile(&self, user_id: &UserId) -> Result<Option<UserProfile>> {
2088        let transaction = self
2089            .inner
2090            .transaction(keys::GLOBAL_PROFILES)
2091            .with_mode(TransactionMode::Readonly)
2092            .build()?;
2093        let store = transaction.object_store(keys::GLOBAL_PROFILES)?;
2094        let key = self.encode_key(keys::GLOBAL_PROFILES, user_id);
2095
2096        store.get(&key).await?.map(|f| self.deserialize_value(&f)).transpose()
2097    }
2098
2099    async fn get_global_profiles<'a>(
2100        &self,
2101        user_ids: &'a [OwnedUserId],
2102    ) -> Result<BTreeMap<&'a UserId, UserProfile>> {
2103        let transaction = self
2104            .inner
2105            .transaction(keys::GLOBAL_PROFILES)
2106            .with_mode(TransactionMode::Readonly)
2107            .build()?;
2108        let store = transaction.object_store(keys::GLOBAL_PROFILES)?;
2109
2110        let mut profiles = BTreeMap::new();
2111        for user_id in user_ids {
2112            let key = self.encode_key(keys::GLOBAL_PROFILES, user_id);
2113            if let Some(value) = store.get(&key).await? {
2114                profiles.insert(user_id.as_ref(), self.deserialize_value(&value)?);
2115            }
2116        }
2117
2118        Ok(profiles)
2119    }
2120
2121    #[allow(clippy::unused_async)]
2122    async fn optimize(&self) -> Result<()> {
2123        Ok(())
2124    }
2125
2126    #[allow(clippy::unused_async)]
2127    async fn get_size(&self) -> Result<Option<usize>> {
2128        Ok(None)
2129    }
2130
2131    #[allow(clippy::unused_async)]
2132    async fn close(&self) -> Result<()> {
2133        Ok(())
2134    }
2135
2136    #[allow(clippy::unused_async)]
2137    async fn reopen(&self) -> Result<()> {
2138        Ok(())
2139    }
2140});
2141
2142/// A room member.
2143#[derive(Debug, Serialize, Deserialize)]
2144struct RoomMember {
2145    user_id: OwnedUserId,
2146    membership: MembershipState,
2147}
2148
2149impl From<&SyncStateEvent<RoomMemberEventContent>> for RoomMember {
2150    fn from(event: &SyncStateEvent<RoomMemberEventContent>) -> Self {
2151        Self { user_id: event.state_key().clone(), membership: event.membership().clone() }
2152    }
2153}
2154
2155impl From<&StrippedRoomMemberEvent> for RoomMember {
2156    fn from(event: &StrippedRoomMemberEvent) -> Self {
2157        Self { user_id: event.state_key.clone(), membership: event.content.membership.clone() }
2158    }
2159}
2160
2161#[cfg(test)]
2162mod migration_tests {
2163    use assert_matches2::assert_matches;
2164    use matrix_sdk_base::store::{QueuedRequestKind, SerializableEventContent};
2165    use ruma::{
2166        OwnedRoomId, OwnedTransactionId, TransactionId,
2167        events::room::message::RoomMessageEventContent, room_id,
2168    };
2169    use serde::{Deserialize, Serialize};
2170
2171    use crate::state_store::PersistedQueuedRequest;
2172
2173    #[derive(Serialize, Deserialize)]
2174    struct OldPersistedQueuedRequest {
2175        room_id: OwnedRoomId,
2176        event: SerializableEventContent,
2177        transaction_id: OwnedTransactionId,
2178        is_wedged: bool,
2179    }
2180
2181    // We now persist an error when an event failed to send instead of just a
2182    // boolean. To support that, `PersistedQueueEvent` changed a bool to
2183    // Option<bool>, ensures that this work properly.
2184    #[test]
2185    fn test_migrating_persisted_queue_event_serialization() {
2186        let room_a_id = room_id!("!room_a:dummy.local");
2187        let transaction_id = TransactionId::new();
2188        let content =
2189            SerializableEventContent::new(&RoomMessageEventContent::text_plain("Hello").into())
2190                .unwrap();
2191
2192        let old_persisted_queue_event = OldPersistedQueuedRequest {
2193            room_id: room_a_id.to_owned(),
2194            event: content,
2195            transaction_id: transaction_id.clone(),
2196            is_wedged: true,
2197        };
2198
2199        let serialized_persisted = serde_json::to_vec(&old_persisted_queue_event).unwrap();
2200
2201        // Load it with the new version.
2202        let new_persisted: PersistedQueuedRequest =
2203            serde_json::from_slice(&serialized_persisted).unwrap();
2204
2205        assert_eq!(new_persisted.is_wedged, Some(true));
2206        assert!(new_persisted.error.is_none());
2207
2208        assert!(new_persisted.event.is_some());
2209        assert!(new_persisted.kind.is_none());
2210
2211        let queued = new_persisted.into_queued_request().unwrap();
2212        assert_matches!(queued.kind, QueuedRequestKind::Event { .. });
2213        assert_eq!(queued.transaction_id, transaction_id);
2214        assert!(queued.error.is_some());
2215    }
2216}
2217
2218#[cfg(all(test, target_family = "wasm"))]
2219mod tests {
2220    #[cfg(target_family = "wasm")]
2221    wasm_bindgen_test::wasm_bindgen_test_configure!(run_in_browser);
2222
2223    use matrix_sdk_base::statestore_integration_tests;
2224    use uuid::Uuid;
2225
2226    use super::{IndexeddbStateStore, Result};
2227
2228    async fn get_store() -> Result<IndexeddbStateStore> {
2229        let db_name = format!("test-state-plain-{}", Uuid::new_v4().as_hyphenated());
2230        Ok(IndexeddbStateStore::builder().name(db_name).build().await?)
2231    }
2232
2233    statestore_integration_tests!();
2234}
2235
2236#[cfg(all(test, target_family = "wasm"))]
2237mod encrypted_tests {
2238    #[cfg(target_family = "wasm")]
2239    wasm_bindgen_test::wasm_bindgen_test_configure!(run_in_browser);
2240
2241    use matrix_sdk_base::statestore_integration_tests;
2242    use uuid::Uuid;
2243
2244    use super::{IndexeddbStateStore, Result};
2245
2246    async fn get_store() -> Result<IndexeddbStateStore> {
2247        let db_name = format!("test-state-encrypted-{}", Uuid::new_v4().as_hyphenated());
2248        let passphrase = format!("some_passphrase-{}", Uuid::new_v4().as_hyphenated());
2249        Ok(IndexeddbStateStore::builder().name(db_name).passphrase(passphrase).build().await?)
2250    }
2251
2252    statestore_integration_tests!();
2253}