1use 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 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 pub const ROOM_SEND_QUEUE: &str = "room_send_queue";
176 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 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 pub const STORE_KEY: &str = "store_key";
216}
217
218pub use keys::ALL_STORES;
219use matrix_sdk_base::store::QueueWedgeError;
220
221fn 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
232fn 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#[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 pub fn name(mut self, value: String) -> Self {
293 self.name = Some(value);
294 self
295 }
296
297 pub fn passphrase(mut self, value: String) -> Self {
301 self.passphrase = Some(value);
302 self
303 }
304
305 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 pub fn builder() -> IndexeddbStateStoreBuilder {
346 IndexeddbStateStoreBuilder::new()
347 }
348
349 pub fn version(&self) -> u32 {
351 self.inner.version() as u32
352 }
353
354 pub fn meta_version(&self) -> u32 {
356 self.meta.version() as u32
357 }
358
359 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 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 fn serialize_value(&self, event: &impl Serialize) -> Result<JsValue> {
393 serialize_value(self.store_cipher.as_deref(), event)
394 }
395
396 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 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 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 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#[derive(Serialize, Deserialize)]
526struct PersistedQueuedRequest {
527 pub room_id: OwnedRoomId,
529
530 kind: Option<QueuedRequestKind>,
533 transaction_id: OwnedTransactionId,
534
535 pub error: Option<QueueWedgeError>,
536
537 priority: Option<usize>,
538
539 #[serde(default = "created_now")]
541 created_at: MilliSecondsSinceUnixEpoch,
542
543 is_wedged: Option<bool>,
547
548 event: Option<SerializableEventContent>,
549}
550
551fn created_now() -> MilliSecondsSinceUnixEpoch {
552 MilliSecondsSinceUnixEpoch::now()
553}
554
555impl PersistedQueuedRequest {
556 fn into_queued_request(self) -> Option<QueuedRequest> {
557 let kind = self.kind.or_else(|| self.event.map(QueuedRequestKind::from))?;
558
559 let error = match self.is_wedged {
560 Some(true) => {
561 Some(QueueWedgeError::GenericApiError {
563 msg: "local echo failed to send in a previous session".into(),
564 })
565 }
566 _ => self.error,
567 };
568
569 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#[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 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 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 receipt_thread: &ReceiptThread,
1459 user_id: &UserId,
1460 ) -> Result<Option<(OwnedEventId, Receipt)>> {
1461 let key = match receipt_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 receipt_thread: &ReceiptThread,
1482 event_id: &EventId,
1483 ) -> Result<Vec<(OwnedUserId, Receipt)>> {
1484 let range = match receipt_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 let direct_stores = [keys::ROOM_INFOS, keys::ROOM_SEND_QUEUE, keys::DEPENDENT_SEND_QUEUE];
1549
1550 let prefixed_stores = [
1553 keys::PROFILES,
1554 keys::DISPLAY_NAMES,
1555 keys::USER_IDS,
1556 keys::ROOM_STATE,
1557 keys::ROOM_ACCOUNT_DATA,
1558 keys::ROOM_EVENT_RECEIPTS,
1559 keys::ROOM_USER_RECEIPTS,
1560 keys::STRIPPED_ROOM_STATE,
1561 keys::STRIPPED_USER_IDS,
1562 keys::THREAD_SUBSCRIPTIONS,
1563 ];
1564
1565 let all_stores = {
1566 let mut v = Vec::new();
1567 v.extend(prefixed_stores);
1568 v.extend(direct_stores);
1569 v
1570 };
1571
1572 let tx =
1573 self.inner.transaction(all_stores).with_mode(TransactionMode::Readwrite).build()?;
1574
1575 for store_name in direct_stores {
1576 tx.object_store(store_name)?.delete(&self.encode_key(store_name, room_id)).build()?;
1577 }
1578
1579 for store_name in prefixed_stores {
1580 let store = tx.object_store(store_name)?;
1581 let range = self.encode_to_range(store_name, room_id);
1582 for key in store.get_all_keys::<JsValue>().with_query(&range).await? {
1583 store.delete(&key?).build()?;
1584 }
1585 }
1586
1587 tx.commit().await.map_err(|e| e.into())
1588 }
1589
1590 async fn get_user_ids(
1591 &self,
1592 room_id: &RoomId,
1593 memberships: RoomMemberships,
1594 ) -> Result<Vec<OwnedUserId>> {
1595 let ids = self.get_user_ids_inner(room_id, memberships, true).await?;
1596 if !ids.is_empty() {
1597 return Ok(ids);
1598 }
1599 self.get_user_ids_inner(room_id, memberships, false).await
1600 }
1601
1602 async fn save_send_queue_request(
1603 &self,
1604 room_id: &RoomId,
1605 transaction_id: OwnedTransactionId,
1606 created_at: MilliSecondsSinceUnixEpoch,
1607 kind: QueuedRequestKind,
1608 priority: usize,
1609 ) -> Result<()> {
1610 let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1611
1612 let tx = self
1613 .inner
1614 .transaction(keys::ROOM_SEND_QUEUE)
1615 .with_mode(TransactionMode::Readwrite)
1616 .build()?;
1617
1618 let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1619
1620 let prev = obj.get(&encoded_key).await?;
1625
1626 let mut prev = prev.map_or_else(
1627 || Ok(Vec::new()),
1628 |val| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val),
1629 )?;
1630
1631 prev.push(PersistedQueuedRequest {
1633 room_id: room_id.to_owned(),
1634 kind: Some(kind),
1635 transaction_id,
1636 error: None,
1637 is_wedged: None,
1638 event: None,
1639 priority: Some(priority),
1640 created_at,
1641 });
1642
1643 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1645
1646 tx.commit().await?;
1647
1648 Ok(())
1649 }
1650
1651 async fn update_send_queue_request(
1652 &self,
1653 room_id: &RoomId,
1654 transaction_id: &TransactionId,
1655 kind: QueuedRequestKind,
1656 ) -> Result<bool> {
1657 let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1658
1659 let tx = self
1660 .inner
1661 .transaction(keys::ROOM_SEND_QUEUE)
1662 .with_mode(TransactionMode::Readwrite)
1663 .build()?;
1664
1665 let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1666
1667 let prev = obj.get(&encoded_key).await?;
1672
1673 let mut prev = prev.map_or_else(
1674 || Ok(Vec::new()),
1675 |val| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val),
1676 )?;
1677
1678 if let Some(entry) = prev.iter_mut().find(|entry| entry.transaction_id == transaction_id) {
1680 entry.kind = Some(kind);
1681 entry.error = None;
1683 entry.is_wedged = None;
1685 entry.event = None;
1686
1687 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1689 tx.commit().await?;
1690
1691 Ok(true)
1692 } else {
1693 Ok(false)
1694 }
1695 }
1696
1697 async fn remove_send_queue_request(
1698 &self,
1699 room_id: &RoomId,
1700 transaction_id: &TransactionId,
1701 ) -> Result<bool> {
1702 let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1703
1704 let tx = self
1705 .inner
1706 .transaction([keys::ROOM_SEND_QUEUE, keys::DEPENDENT_SEND_QUEUE])
1707 .with_mode(TransactionMode::Readwrite)
1708 .build()?;
1709
1710 let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1711
1712 if let Some(val) = obj.get(&encoded_key).await? {
1717 let mut prev = self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val)?;
1718 if let Some(pos) = prev.iter().position(|item| item.transaction_id == transaction_id) {
1719 prev.remove(pos);
1720
1721 if prev.is_empty() {
1722 obj.delete(&encoded_key).build()?;
1723 } else {
1724 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1725 }
1726
1727 tx.commit().await?;
1728 return Ok(true);
1729 }
1730 }
1731
1732 Ok(false)
1733 }
1734
1735 async fn load_send_queue_requests(&self, room_id: &RoomId) -> Result<Vec<QueuedRequest>> {
1736 let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1737
1738 let prev = self
1741 .inner
1742 .transaction(keys::ROOM_SEND_QUEUE)
1743 .with_mode(TransactionMode::Readwrite)
1744 .build()?
1745 .object_store(keys::ROOM_SEND_QUEUE)?
1746 .get(&encoded_key)
1747 .await?;
1748
1749 let mut prev = prev.map_or_else(
1750 || Ok(Vec::new()),
1751 |val| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val),
1752 )?;
1753
1754 prev.sort_by_key(|item| Reverse(item.priority.unwrap_or(0)));
1756
1757 Ok(prev.into_iter().filter_map(PersistedQueuedRequest::into_queued_request).collect())
1758 }
1759
1760 async fn update_send_queue_request_status(
1761 &self,
1762 room_id: &RoomId,
1763 transaction_id: &TransactionId,
1764 error: Option<QueueWedgeError>,
1765 ) -> Result<()> {
1766 let encoded_key = self.encode_key(keys::ROOM_SEND_QUEUE, room_id);
1767
1768 let tx = self
1769 .inner
1770 .transaction(keys::ROOM_SEND_QUEUE)
1771 .with_mode(TransactionMode::Readwrite)
1772 .build()?;
1773
1774 let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1775
1776 if let Some(val) = obj.get(&encoded_key).await? {
1777 let mut prev = self.deserialize_value::<Vec<PersistedQueuedRequest>>(&val)?;
1778 if let Some(request) =
1779 prev.iter_mut().find(|item| item.transaction_id == transaction_id)
1780 {
1781 request.is_wedged = None;
1782 request.error = error;
1783 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1784 }
1785 }
1786
1787 tx.commit().await?;
1788
1789 Ok(())
1790 }
1791
1792 async fn load_rooms_with_unsent_requests(&self) -> Result<Vec<OwnedRoomId>> {
1793 let tx = self
1794 .inner
1795 .transaction(keys::ROOM_SEND_QUEUE)
1796 .with_mode(TransactionMode::Readwrite)
1797 .build()?;
1798
1799 let obj = tx.object_store(keys::ROOM_SEND_QUEUE)?;
1800
1801 let all_entries = obj
1802 .get_all()
1803 .await?
1804 .map(|item| self.deserialize_value::<Vec<PersistedQueuedRequest>>(&item?))
1805 .collect::<Result<Vec<Vec<PersistedQueuedRequest>>, _>>()?
1806 .into_iter()
1807 .flat_map(|vec| vec.into_iter().map(|item| item.room_id))
1808 .collect::<BTreeSet<_>>();
1809
1810 Ok(all_entries.into_iter().collect())
1811 }
1812
1813 async fn save_dependent_queued_request(
1814 &self,
1815 room_id: &RoomId,
1816 parent_txn_id: &TransactionId,
1817 own_txn_id: ChildTransactionId,
1818 created_at: MilliSecondsSinceUnixEpoch,
1819 content: DependentQueuedRequestKind,
1820 ) -> Result<()> {
1821 let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1822
1823 let tx = self
1824 .inner
1825 .transaction(keys::DEPENDENT_SEND_QUEUE)
1826 .with_mode(TransactionMode::Readwrite)
1827 .build()?;
1828
1829 let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1830
1831 let prev = obj.get(&encoded_key).await?;
1834
1835 let mut prev = prev.map_or_else(
1836 || Ok(Vec::new()),
1837 |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1838 )?;
1839
1840 prev.push(DependentQueuedRequest {
1842 kind: content,
1843 parent_transaction_id: parent_txn_id.to_owned(),
1844 own_transaction_id: own_txn_id,
1845 parent_key: None,
1846 created_at,
1847 });
1848
1849 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1851
1852 tx.commit().await?;
1853
1854 Ok(())
1855 }
1856
1857 async fn update_dependent_queued_request(
1858 &self,
1859 room_id: &RoomId,
1860 own_transaction_id: &ChildTransactionId,
1861 new_content: DependentQueuedRequestKind,
1862 ) -> Result<bool> {
1863 let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1864
1865 let tx = self
1866 .inner
1867 .transaction(keys::DEPENDENT_SEND_QUEUE)
1868 .with_mode(TransactionMode::Readwrite)
1869 .build()?;
1870
1871 let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1872
1873 let prev = obj.get(&encoded_key).await?;
1876
1877 let mut prev = prev.map_or_else(
1878 || Ok(Vec::new()),
1879 |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1880 )?;
1881
1882 let mut found = false;
1884 for entry in prev.iter_mut() {
1885 if entry.own_transaction_id == *own_transaction_id {
1886 found = true;
1887 entry.kind = new_content;
1888 break;
1889 }
1890 }
1891
1892 if found {
1893 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1894 tx.commit().await?;
1895 }
1896
1897 Ok(found)
1898 }
1899
1900 async fn mark_dependent_queued_requests_as_ready(
1901 &self,
1902 room_id: &RoomId,
1903 parent_txn_id: &TransactionId,
1904 parent_key: SentRequestKey,
1905 ) -> Result<usize> {
1906 let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1907
1908 let tx = self
1909 .inner
1910 .transaction(keys::DEPENDENT_SEND_QUEUE)
1911 .with_mode(TransactionMode::Readwrite)
1912 .build()?;
1913
1914 let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1915
1916 let prev = obj.get(&encoded_key).await?;
1919
1920 let mut prev = prev.map_or_else(
1921 || Ok(Vec::new()),
1922 |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1923 )?;
1924
1925 let mut num_updated = 0;
1927 for entry in prev.iter_mut().filter(|entry| entry.parent_transaction_id == parent_txn_id) {
1928 entry.parent_key = Some(parent_key.clone());
1929 num_updated += 1;
1930 }
1931
1932 if num_updated > 0 {
1933 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1934 tx.commit().await?;
1935 }
1936
1937 Ok(num_updated)
1938 }
1939
1940 async fn remove_dependent_queued_request(
1941 &self,
1942 room_id: &RoomId,
1943 txn_id: &ChildTransactionId,
1944 ) -> Result<bool> {
1945 let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1946
1947 let tx = self
1948 .inner
1949 .transaction(keys::DEPENDENT_SEND_QUEUE)
1950 .with_mode(TransactionMode::Readwrite)
1951 .build()?;
1952
1953 let obj = tx.object_store(keys::DEPENDENT_SEND_QUEUE)?;
1954
1955 if let Some(val) = obj.get(&encoded_key).await? {
1958 let mut prev = self.deserialize_value::<Vec<DependentQueuedRequest>>(&val)?;
1959 if let Some(pos) = prev.iter().position(|item| item.own_transaction_id == *txn_id) {
1960 prev.remove(pos);
1961
1962 if prev.is_empty() {
1963 obj.delete(&encoded_key).build()?;
1964 } else {
1965 obj.put(&self.serialize_value(&prev)?).with_key(encoded_key).build()?;
1966 }
1967
1968 tx.commit().await?;
1969 return Ok(true);
1970 }
1971 }
1972
1973 Ok(false)
1974 }
1975
1976 async fn load_dependent_queued_requests(
1977 &self,
1978 room_id: &RoomId,
1979 ) -> Result<Vec<DependentQueuedRequest>> {
1980 let encoded_key = self.encode_key(keys::DEPENDENT_SEND_QUEUE, room_id);
1981
1982 let prev = self
1984 .inner
1985 .transaction(keys::DEPENDENT_SEND_QUEUE)
1986 .with_mode(TransactionMode::Readwrite)
1987 .build()?
1988 .object_store(keys::DEPENDENT_SEND_QUEUE)?
1989 .get(&encoded_key)
1990 .await?;
1991
1992 prev.map_or_else(
1993 || Ok(Vec::new()),
1994 |val| self.deserialize_value::<Vec<DependentQueuedRequest>>(&val),
1995 )
1996 }
1997
1998 async fn upsert_thread_subscriptions(
1999 &self,
2000 updates: Vec<(&RoomId, &EventId, StoredThreadSubscription)>,
2001 ) -> Result<()> {
2002 let tx = self
2003 .inner
2004 .transaction(keys::THREAD_SUBSCRIPTIONS)
2005 .with_mode(TransactionMode::Readwrite)
2006 .build()?;
2007 let obj = tx.object_store(keys::THREAD_SUBSCRIPTIONS)?;
2008
2009 for (room_id, thread_id, subscription) in updates {
2010 let encoded_key = self.encode_key(keys::THREAD_SUBSCRIPTIONS, (room_id, thread_id));
2011 let mut new = PersistedThreadSubscription::from(subscription);
2012
2013 if let Some(previous_value) = obj.get(&encoded_key).await? {
2015 let previous: PersistedThreadSubscription =
2016 self.deserialize_value(&previous_value)?;
2017
2018 if new == previous {
2021 continue;
2022 }
2023 if !compare_thread_subscription_bump_stamps(
2024 previous.bump_stamp,
2025 &mut new.bump_stamp,
2026 ) {
2027 continue;
2028 }
2029 }
2030
2031 let serialized_value = self.serialize_value(&new);
2032 obj.put(&serialized_value?).with_key(encoded_key).build()?;
2033 }
2034
2035 tx.commit().await?;
2036
2037 Ok(())
2038 }
2039
2040 async fn load_thread_subscription(
2041 &self,
2042 room: &RoomId,
2043 thread_id: &EventId,
2044 ) -> Result<Option<StoredThreadSubscription>> {
2045 let encoded_key = self.encode_key(keys::THREAD_SUBSCRIPTIONS, (room, thread_id));
2046
2047 let js_value = self
2048 .inner
2049 .transaction(keys::THREAD_SUBSCRIPTIONS)
2050 .with_mode(TransactionMode::Readonly)
2051 .build()?
2052 .object_store(keys::THREAD_SUBSCRIPTIONS)?
2053 .get(&encoded_key)
2054 .await?;
2055
2056 let Some(js_value) = js_value else {
2057 return Ok(None);
2059 };
2060
2061 let sub: PersistedThreadSubscription = self.deserialize_value(&js_value)?;
2062
2063 let status = ThreadSubscriptionStatus::from_str(&sub.status).map_err(|_| {
2064 StoreError::InvalidData {
2065 details: format!(
2066 "invalid thread status for room {room} and thread {thread_id}: {}",
2067 sub.status
2068 ),
2069 }
2070 })?;
2071
2072 Ok(Some(StoredThreadSubscription { status, bump_stamp: sub.bump_stamp }))
2073 }
2074
2075 async fn remove_thread_subscription(&self, room: &RoomId, thread_id: &EventId) -> Result<()> {
2076 let encoded_key = self.encode_key(keys::THREAD_SUBSCRIPTIONS, (room, thread_id));
2077
2078 let transaction = self
2079 .inner
2080 .transaction(keys::THREAD_SUBSCRIPTIONS)
2081 .with_mode(TransactionMode::Readwrite)
2082 .build()?;
2083 transaction.object_store(keys::THREAD_SUBSCRIPTIONS)?.delete(&encoded_key).await?;
2084 transaction.commit().await?;
2085
2086 Ok(())
2087 }
2088
2089 async fn get_global_profile(&self, user_id: &UserId) -> Result<Option<UserProfile>> {
2090 let transaction = self
2091 .inner
2092 .transaction(keys::GLOBAL_PROFILES)
2093 .with_mode(TransactionMode::Readonly)
2094 .build()?;
2095 let store = transaction.object_store(keys::GLOBAL_PROFILES)?;
2096 let key = self.encode_key(keys::GLOBAL_PROFILES, user_id);
2097
2098 store.get(&key).await?.map(|f| self.deserialize_value(&f)).transpose()
2099 }
2100
2101 async fn get_global_profiles<'a>(
2102 &self,
2103 user_ids: &'a [OwnedUserId],
2104 ) -> Result<BTreeMap<&'a UserId, UserProfile>> {
2105 let transaction = self
2106 .inner
2107 .transaction(keys::GLOBAL_PROFILES)
2108 .with_mode(TransactionMode::Readonly)
2109 .build()?;
2110 let store = transaction.object_store(keys::GLOBAL_PROFILES)?;
2111
2112 let mut profiles = BTreeMap::new();
2113 for user_id in user_ids {
2114 let key = self.encode_key(keys::GLOBAL_PROFILES, user_id);
2115 if let Some(value) = store.get(&key).await? {
2116 profiles.insert(user_id.as_ref(), self.deserialize_value(&value)?);
2117 }
2118 }
2119
2120 Ok(profiles)
2121 }
2122
2123 #[allow(clippy::unused_async)]
2124 async fn optimize(&self) -> Result<()> {
2125 Ok(())
2126 }
2127
2128 #[allow(clippy::unused_async)]
2129 async fn get_size(&self) -> Result<Option<usize>> {
2130 Ok(None)
2131 }
2132
2133 #[allow(clippy::unused_async)]
2134 async fn close(&self) -> Result<()> {
2135 Ok(())
2136 }
2137
2138 #[allow(clippy::unused_async)]
2139 async fn reopen(&self) -> Result<()> {
2140 Ok(())
2141 }
2142});
2143
2144#[derive(Debug, Serialize, Deserialize)]
2146struct RoomMember {
2147 user_id: OwnedUserId,
2148 membership: MembershipState,
2149}
2150
2151impl From<&SyncStateEvent<RoomMemberEventContent>> for RoomMember {
2152 fn from(event: &SyncStateEvent<RoomMemberEventContent>) -> Self {
2153 Self { user_id: event.state_key().clone(), membership: event.membership().clone() }
2154 }
2155}
2156
2157impl From<&StrippedRoomMemberEvent> for RoomMember {
2158 fn from(event: &StrippedRoomMemberEvent) -> Self {
2159 Self { user_id: event.state_key.clone(), membership: event.content.membership.clone() }
2160 }
2161}
2162
2163#[cfg(test)]
2164mod migration_tests {
2165 use std::assert_matches;
2166
2167 use matrix_sdk_base::store::{QueuedRequestKind, SerializableEventContent};
2168 use ruma::{
2169 OwnedRoomId, OwnedTransactionId, TransactionId,
2170 events::room::message::RoomMessageEventContent, room_id,
2171 };
2172 use serde::{Deserialize, Serialize};
2173
2174 use crate::state_store::PersistedQueuedRequest;
2175
2176 #[derive(Serialize, Deserialize)]
2177 struct OldPersistedQueuedRequest {
2178 room_id: OwnedRoomId,
2179 event: SerializableEventContent,
2180 transaction_id: OwnedTransactionId,
2181 is_wedged: bool,
2182 }
2183
2184 #[test]
2188 fn test_migrating_persisted_queue_event_serialization() {
2189 let room_a_id = room_id!("!room_a:dummy.local");
2190 let transaction_id = TransactionId::new();
2191 let content =
2192 SerializableEventContent::new(&RoomMessageEventContent::text_plain("Hello").into())
2193 .unwrap();
2194
2195 let old_persisted_queue_event = OldPersistedQueuedRequest {
2196 room_id: room_a_id.to_owned(),
2197 event: content,
2198 transaction_id: transaction_id.clone(),
2199 is_wedged: true,
2200 };
2201
2202 let serialized_persisted = serde_json::to_vec(&old_persisted_queue_event).unwrap();
2203
2204 let new_persisted: PersistedQueuedRequest =
2206 serde_json::from_slice(&serialized_persisted).unwrap();
2207
2208 assert_eq!(new_persisted.is_wedged, Some(true));
2209 assert!(new_persisted.error.is_none());
2210
2211 assert!(new_persisted.event.is_some());
2212 assert!(new_persisted.kind.is_none());
2213
2214 let queued = new_persisted.into_queued_request().unwrap();
2215 assert_matches!(queued.kind, QueuedRequestKind::Event { .. });
2216 assert_eq!(queued.transaction_id, transaction_id);
2217 assert!(queued.error.is_some());
2218 }
2219}
2220
2221#[cfg(all(test, target_family = "wasm"))]
2222mod tests {
2223 #[cfg(target_family = "wasm")]
2224 wasm_bindgen_test::wasm_bindgen_test_configure!(run_in_browser);
2225
2226 use matrix_sdk_base::statestore_integration_tests;
2227 use uuid::Uuid;
2228
2229 use super::{IndexeddbStateStore, Result};
2230
2231 async fn get_store() -> Result<IndexeddbStateStore> {
2232 let db_name = format!("test-state-plain-{}", Uuid::new_v4().as_hyphenated());
2233 Ok(IndexeddbStateStore::builder().name(db_name).build().await?)
2234 }
2235
2236 statestore_integration_tests!();
2237}
2238
2239#[cfg(all(test, target_family = "wasm"))]
2240mod encrypted_tests {
2241 #[cfg(target_family = "wasm")]
2242 wasm_bindgen_test::wasm_bindgen_test_configure!(run_in_browser);
2243
2244 use matrix_sdk_base::statestore_integration_tests;
2245 use uuid::Uuid;
2246
2247 use super::{IndexeddbStateStore, Result};
2248
2249 async fn get_store() -> Result<IndexeddbStateStore> {
2250 let db_name = format!("test-state-encrypted-{}", Uuid::new_v4().as_hyphenated());
2251 let passphrase = format!("some_passphrase-{}", Uuid::new_v4().as_hyphenated());
2252 Ok(IndexeddbStateStore::builder().name(db_name).passphrase(passphrase).build().await?)
2253 }
2254
2255 statestore_integration_tests!();
2256}