1#[cfg(feature = "experimental-encrypted-state-events")]
16use std::borrow::Borrow;
17use std::{
18 collections::{BTreeMap, HashMap, HashSet},
19 sync::Arc,
20 time::Duration,
21};
22
23use itertools::Itertools;
24#[cfg(feature = "experimental-send-custom-to-device")]
25use matrix_sdk_common::deserialized_responses::WithheldCode;
26use matrix_sdk_common::{
27 BoxFuture,
28 deserialized_responses::{
29 AlgorithmInfo, DecryptedRoomEvent, DeviceLinkProblem, EncryptionInfo, ForwarderInfo,
30 ProcessedToDeviceEvent, ToDeviceUnableToDecryptInfo, ToDeviceUnableToDecryptReason,
31 UnableToDecryptInfo, UnableToDecryptReason, UnsignedDecryptionResult,
32 UnsignedEventLocation, VerificationLevel, VerificationState,
33 },
34 locks::RwLock as StdRwLock,
35 timer,
36};
37#[cfg(feature = "experimental-encrypted-state-events")]
38use ruma::events::{AnyStateEventContent, StateEventContent};
39use ruma::{
40 DeviceId, DeviceKeyAlgorithm, MilliSecondsSinceUnixEpoch, OneTimeKeyAlgorithm, OwnedDeviceId,
41 OwnedDeviceKeyId, OwnedTransactionId, OwnedUserId, RoomId, TransactionId, UInt, UserId,
42 api::client::{
43 dehydrated_device::DehydratedDeviceData,
44 keys::{
45 claim_keys::v3::Request as KeysClaimRequest,
46 get_keys::v3::Response as KeysQueryResponse,
47 upload_keys::v3::{Request as UploadKeysRequest, Response as UploadKeysResponse},
48 upload_signatures::v3::Request as UploadSignaturesRequest,
49 },
50 sync::sync_events::DeviceLists,
51 },
52 assign,
53 events::{
54 AnyMessageLikeEvent, AnyMessageLikeEventContent, AnyTimelineEvent, AnyToDeviceEvent,
55 MessageLikeEventContent, secret::request::SecretName,
56 },
57 serde::{JsonObject, Raw},
58};
59use serde::Serialize;
60use serde_json::{Value, value::to_raw_value};
61use tokio::sync::Mutex;
62use tracing::{
63 Span, debug, enabled, error,
64 field::{debug, display},
65 info, instrument, trace, warn,
66};
67use vodozemac::{Curve25519PublicKey, Ed25519Signature, megolm::DecryptionError};
68
69#[cfg(feature = "experimental-push-secrets")]
70use crate::error::SecretPushError;
71#[cfg(feature = "experimental-send-custom-to-device")]
72use crate::session_manager::split_devices_for_share_strategy;
73#[cfg(feature = "experimental-x509-identity-verification")]
74use crate::x509::{RawX509Signer, RawX509Verifier, X509Signer, X509Verifier};
75use crate::{
76 CollectStrategy, CryptoStoreError, DecryptionSettings, DeviceData, LocalTrust,
77 RoomEventDecryptionResult, SignatureError, TrustRequirement,
78 backups::{BackupMachine, MegolmV1BackupKey},
79 dehydrated_devices::{DehydratedDevices, DehydrationError},
80 error::{EventError, MegolmError, MegolmResult, OlmError, OlmResult, SetRoomSettingsError},
81 gossiping::GossipMachine,
82 identities::{Device, IdentityManager, UserDevices, user::UserIdentity},
83 olm::{
84 Account, CrossSigningStatus, EncryptionSettings, IdentityKeys, InboundGroupSession,
85 KnownSenderData, OlmDecryptionInfo, PrivateCrossSigningIdentity, SenderData,
86 SenderDataFinder, SessionType, StaticAccountData,
87 },
88 session_manager::{GroupSessionManager, SessionManager},
89 store::{
90 CryptoStoreWrapper, DynCryptoStore, IntoCryptoStore, MemoryStore, Result as StoreResult,
91 SecretImportError, Store, StoreTransaction,
92 caches::StoreCache,
93 types::{
94 Changes, CrossSigningKeyExport, DeviceChanges, IdentityChanges, PendingChanges,
95 RoomKeyInfo, RoomSettings, StoredRoomKeyBundleData,
96 },
97 },
98 types::{
99 EventEncryptionAlgorithm, Signatures,
100 events::{
101 ToDeviceEvent, ToDeviceEvents,
102 olm_v1::{AnyDecryptedOlmEvent, DecryptedRoomKeyBundleEvent, DecryptedRoomKeyEvent},
103 room::encrypted::{
104 EncryptedEvent, EncryptedToDeviceEvent, RoomEncryptedEventContent,
105 RoomEventEncryptionScheme, SupportedEventEncryptionSchemes,
106 ToDeviceEncryptedEventContent,
107 },
108 room_key::{MegolmV1AesSha2Content, RoomKeyContent},
109 room_key_bundle::RoomKeyBundleContent,
110 room_key_withheld::{
111 MegolmV1AesSha2WithheldContent, RoomKeyWithheldContent, RoomKeyWithheldEvent,
112 },
113 },
114 requests::{
115 AnyIncomingResponse, KeysQueryRequest, OutgoingRequest, ToDeviceRequest,
116 UploadSigningKeysRequest,
117 },
118 },
119 utilities::timestamp_to_iso8601,
120 verification::{Verification, VerificationMachine, VerificationRequest},
121};
122
123#[derive(Debug, Serialize)]
124pub struct RawEncryptionResult {
126 pub content: Raw<RoomEncryptedEventContent>,
128 pub encryption_info: EncryptionInfo,
130}
131
132pub struct OlmMachineBuilder {
134 user_id: OwnedUserId,
136
137 device_id: OwnedDeviceId,
139
140 store: Option<Arc<DynCryptoStore>>,
143
144 custom_account: Option<vodozemac::olm::Account>,
147
148 #[cfg(feature = "experimental-x509-identity-verification")]
150 x509_verifier: Option<X509Verifier>,
151
152 #[cfg(feature = "experimental-x509-identity-verification")]
154 x509_signer: Option<X509Signer>,
155}
156
157impl OlmMachineBuilder {
158 pub fn new(user_id: &UserId, device_id: &DeviceId) -> Self {
161 Self {
162 user_id: user_id.to_owned(),
163 device_id: device_id.to_owned(),
164 store: None,
165 custom_account: None,
166 #[cfg(feature = "experimental-x509-identity-verification")]
167 x509_verifier: None,
168 #[cfg(feature = "experimental-x509-identity-verification")]
169 x509_signer: None,
170 }
171 }
172
173 pub fn with_crypto_store(mut self, store: impl IntoCryptoStore) -> Self {
181 self.store = Some(store.into_crypto_store());
182 self
183 }
184
185 pub fn with_custom_account(mut self, custom_account: Option<vodozemac::olm::Account>) -> Self {
195 self.custom_account = custom_account;
196 self
197 }
198
199 #[cfg(feature = "experimental-x509-identity-verification")]
202 pub fn with_x509_verifier(mut self, x509_verifier: Option<Arc<dyn RawX509Verifier>>) -> Self {
203 self.x509_verifier = x509_verifier.map(X509Verifier::new);
204 self
205 }
206
207 #[cfg(feature = "experimental-x509-identity-verification")]
210 pub fn with_x509_signer(mut self, x509_signer: Option<Arc<dyn RawX509Signer>>) -> Self {
211 self.x509_signer = x509_signer.map(X509Signer::new);
212 self
213 }
214
215 pub async fn build(self) -> Result<OlmMachine, CryptoStoreError> {
222 OlmMachine::from_builder(self).await
223 }
224}
225
226impl std::fmt::Debug for OlmMachineBuilder {
227 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
228 f.debug_struct("OlmMachineBuilder")
229 .field("user_id", &self.user_id)
230 .field("device_id", &self.device_id)
231 .finish_non_exhaustive()
232 }
233}
234
235#[derive(Clone)]
238pub struct OlmMachine {
239 pub(crate) inner: Arc<OlmMachineInner>,
240}
241
242pub struct OlmMachineInner {
243 user_id: OwnedUserId,
245 device_id: OwnedDeviceId,
247 user_identity: Arc<Mutex<PrivateCrossSigningIdentity>>,
252 store: Store,
256 session_manager: SessionManager,
258 pub(crate) group_session_manager: GroupSessionManager,
260 verification_machine: VerificationMachine,
263 pub(crate) key_request_machine: GossipMachine,
266 identity_manager: IdentityManager,
269 backup_machine: BackupMachine,
271}
272
273#[cfg(not(tarpaulin_include))]
274impl std::fmt::Debug for OlmMachine {
275 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
276 f.debug_struct("OlmMachine")
277 .field("user_id", &self.user_id())
278 .field("device_id", &self.device_id())
279 .finish()
280 }
281}
282
283impl OlmMachine {
284 const CURRENT_GENERATION_STORE_KEY: &'static str = "generation-counter";
285 const HAS_MIGRATED_VERIFICATION_LATCH: &'static str = "HAS_MIGRATED_VERIFICATION_LATCH";
286
287 pub async fn new(user_id: &UserId, device_id: &DeviceId) -> Self {
298 OlmMachineBuilder::new(user_id, device_id)
299 .build()
300 .await
301 .expect("Reading and writing to the memory store always succeeds")
302 }
303
304 pub(crate) async fn rehydrate(
305 &self,
306 pickle_key: &[u8; 32],
307 device_id: &DeviceId,
308 device_data: Raw<DehydratedDeviceData>,
309 ) -> Result<OlmMachine, DehydrationError> {
310 let account = Account::rehydrate(pickle_key, self.user_id(), device_id, device_data)?;
311 let static_account = account.static_data().clone();
312
313 let store =
314 Arc::new(CryptoStoreWrapper::new(self.user_id(), device_id, MemoryStore::new()));
315 let device = DeviceData::from_account(&account);
316 store.save_pending_changes(PendingChanges { account: Some(account) }).await?;
317 store
318 .save_changes(Changes {
319 devices: DeviceChanges { new: vec![device], ..Default::default() },
320 ..Default::default()
321 })
322 .await?;
323
324 let (verification_machine, store, identity_manager) = Self::new_helper_prelude(
325 store,
326 static_account,
327 self.store().private_identity(),
328 #[cfg(feature = "experimental-x509-identity-verification")]
329 self.store().x509_verifier().cloned(),
330 #[cfg(feature = "experimental-x509-identity-verification")]
331 self.store().x509_signer().cloned(),
332 );
333
334 Ok(Self::new_helper(
335 device_id,
336 store,
337 verification_machine,
338 identity_manager,
339 self.store().private_identity(),
340 None,
341 ))
342 }
343
344 fn new_helper_prelude(
345 store_wrapper: Arc<CryptoStoreWrapper>,
346 account: StaticAccountData,
347 user_identity: Arc<Mutex<PrivateCrossSigningIdentity>>,
348 #[cfg(feature = "experimental-x509-identity-verification")] x509_verifier: Option<
349 X509Verifier,
350 >,
351 #[cfg(feature = "experimental-x509-identity-verification")] x509_signer: Option<X509Signer>,
352 ) -> (VerificationMachine, Store, IdentityManager) {
353 let verification_machine =
354 VerificationMachine::new(account.clone(), user_identity.clone(), store_wrapper.clone());
355
356 let store = Store::new_with_x509(
357 account,
358 user_identity,
359 store_wrapper,
360 verification_machine.clone(),
361 #[cfg(feature = "experimental-x509-identity-verification")]
362 x509_verifier,
363 #[cfg(feature = "experimental-x509-identity-verification")]
364 x509_signer,
365 );
366
367 let identity_manager = IdentityManager::new(store.clone());
368
369 (verification_machine, store, identity_manager)
370 }
371
372 fn new_helper(
373 device_id: &DeviceId,
374 store: Store,
375 verification_machine: VerificationMachine,
376 identity_manager: IdentityManager,
377 user_identity: Arc<Mutex<PrivateCrossSigningIdentity>>,
378 maybe_backup_key: Option<MegolmV1BackupKey>,
379 ) -> Self {
380 let group_session_manager = GroupSessionManager::new(store.clone());
381
382 let users_for_key_claim = Arc::new(StdRwLock::new(BTreeMap::new()));
383 let key_request_machine = GossipMachine::new(
384 store.clone(),
385 identity_manager.clone(),
386 group_session_manager.session_cache(),
387 users_for_key_claim.clone(),
388 );
389
390 let session_manager =
391 SessionManager::new(users_for_key_claim, key_request_machine.clone(), store.clone());
392
393 let backup_machine = BackupMachine::new(store.clone(), maybe_backup_key);
394
395 let inner = Arc::new(OlmMachineInner {
396 user_id: store.user_id().to_owned(),
397 device_id: device_id.to_owned(),
398 user_identity,
399 store,
400 session_manager,
401 group_session_manager,
402 verification_machine,
403 key_request_machine,
404 identity_manager,
405 backup_machine,
406 });
407
408 Self { inner }
409 }
410
411 #[instrument(skip(builder), fields(user_id, device_id, ed25519_key, curve25519_key))]
412 pub(crate) async fn from_builder(builder: OlmMachineBuilder) -> StoreResult<Self> {
413 let OlmMachineBuilder {
414 user_id,
415 device_id,
416 store,
417 custom_account,
418 #[cfg(feature = "experimental-x509-identity-verification")]
419 x509_verifier,
420 #[cfg(feature = "experimental-x509-identity-verification")]
421 x509_signer,
422 } = builder;
423
424 let store = store.unwrap_or_else(|| MemoryStore::new().into_crypto_store());
425
426 let static_account = match store.load_account().await? {
427 Some(account) => {
428 if user_id != account.user_id()
429 || device_id != account.device_id()
430 || custom_account.is_some()
431 {
432 return Err(CryptoStoreError::MismatchedAccount {
433 expected: (account.user_id().to_owned(), account.device_id().to_owned()),
434 got: (user_id.to_owned(), device_id.to_owned()),
435 });
436 }
437
438 Span::current()
439 .record("ed25519_key", display(account.identity_keys().ed25519))
440 .record("curve25519_key", display(account.identity_keys().curve25519));
441 debug!("Restored an Olm account");
442
443 account.static_data().clone()
444 }
445
446 None => {
447 let account = if let Some(account) = custom_account {
448 Account::new_helper(account, &user_id, &device_id)
449 } else {
450 Account::with_device_id(&user_id, &device_id)
451 };
452
453 let static_account = account.static_data().clone();
454
455 Span::current()
456 .record("ed25519_key", display(account.identity_keys().ed25519))
457 .record("curve25519_key", display(account.identity_keys().curve25519));
458
459 let device = DeviceData::from_account(&account);
460
461 device.set_trust_state(LocalTrust::Verified);
465
466 let changes = Changes {
467 devices: DeviceChanges { new: vec![device], ..Default::default() },
468 ..Default::default()
469 };
470 store.save_changes(changes).await?;
471 store.save_pending_changes(PendingChanges { account: Some(account) }).await?;
472
473 debug!("Created a new Olm account");
474
475 static_account
476 }
477 };
478
479 let identity = match store.load_identity().await? {
480 Some(i) => {
481 let master_key = i
482 .master_public_key()
483 .await
484 .and_then(|m| m.get_first_key().map(|m| m.to_owned()));
485 debug!(?master_key, "Restored the cross signing identity");
486 i
487 }
488 None => {
489 debug!("Creating an empty cross signing identity stub");
490 PrivateCrossSigningIdentity::empty(&user_id)
491 }
492 };
493
494 let saved_keys = store.load_backup_keys().await?;
499 let maybe_backup_key = saved_keys.decryption_key.and_then(|k| {
500 if let Some(version) = saved_keys.backup_version {
501 let megolm_v1_backup_key = k.megolm_v1_public_key();
502 megolm_v1_backup_key.set_version(version);
503 Some(megolm_v1_backup_key)
504 } else {
505 None
506 }
507 });
508
509 let identity = Arc::new(Mutex::new(identity));
510 let store = Arc::new(CryptoStoreWrapper::new(&user_id, &device_id, store));
511
512 let (verification_machine, store, identity_manager) = Self::new_helper_prelude(
513 store,
514 static_account,
515 identity.clone(),
516 #[cfg(feature = "experimental-x509-identity-verification")]
517 x509_verifier,
518 #[cfg(feature = "experimental-x509-identity-verification")]
519 x509_signer,
520 );
521
522 Self::migration_post_verified_latch_support(&store, &identity_manager).await?;
525
526 Ok(Self::new_helper(
527 &device_id,
528 store,
529 verification_machine,
530 identity_manager,
531 identity,
532 maybe_backup_key,
533 ))
534 }
535
536 pub(crate) async fn migration_post_verified_latch_support(
544 store: &Store,
545 identity_manager: &IdentityManager,
546 ) -> Result<(), CryptoStoreError> {
547 let maybe_migrate_for_identity_verified_latch =
548 store.get_custom_value(Self::HAS_MIGRATED_VERIFICATION_LATCH).await?.is_none();
549
550 if maybe_migrate_for_identity_verified_latch {
551 identity_manager.mark_all_tracked_users_as_dirty(store.cache().await?).await?;
552
553 store.set_custom_value(Self::HAS_MIGRATED_VERIFICATION_LATCH, vec![0]).await?
554 }
555 Ok(())
556 }
557
558 pub fn store(&self) -> &Store {
560 &self.inner.store
561 }
562
563 pub fn user_id(&self) -> &UserId {
565 &self.inner.user_id
566 }
567
568 pub fn device_id(&self) -> &DeviceId {
570 &self.inner.device_id
571 }
572
573 pub fn device_creation_time(&self) -> MilliSecondsSinceUnixEpoch {
580 self.inner.store.static_account().creation_local_time()
581 }
582
583 pub fn identity_keys(&self) -> IdentityKeys {
585 let account = self.inner.store.static_account();
586 account.identity_keys()
587 }
588
589 pub async fn display_name(&self) -> StoreResult<Option<String>> {
591 self.store().device_display_name().await
592 }
593
594 pub async fn tracked_users(&self) -> StoreResult<HashSet<OwnedUserId>> {
599 let cache = self.store().cache().await?;
600 Ok(self.inner.identity_manager.key_query_manager.synced(&cache).await?.tracked_users())
601 }
602
603 #[cfg(feature = "automatic-room-key-forwarding")]
612 pub fn set_room_key_requests_enabled(&self, enable: bool) {
613 self.inner.key_request_machine.set_room_key_requests_enabled(enable)
614 }
615
616 pub fn are_room_key_requests_enabled(&self) -> bool {
621 self.inner.key_request_machine.are_room_key_requests_enabled()
622 }
623
624 #[cfg(feature = "automatic-room-key-forwarding")]
633 pub fn set_room_key_forwarding_enabled(&self, enable: bool) {
634 self.inner.key_request_machine.set_room_key_forwarding_enabled(enable)
635 }
636
637 pub fn is_room_key_forwarding_enabled(&self) -> bool {
641 self.inner.key_request_machine.is_room_key_forwarding_enabled()
642 }
643
644 pub async fn outgoing_requests(&self) -> StoreResult<Vec<OutgoingRequest>> {
652 let mut requests = Vec::new();
653
654 {
655 let store_cache = self.inner.store.cache().await?;
656 let account = store_cache.account().await?;
657 if let Some(r) = self.keys_for_upload(&account).await.map(|r| OutgoingRequest {
658 request_id: TransactionId::new(),
659 request: Arc::new(r.into()),
660 }) {
661 requests.push(r);
662 }
663 }
664
665 for request in self
666 .inner
667 .identity_manager
668 .users_for_key_query()
669 .await?
670 .into_iter()
671 .map(|(request_id, r)| OutgoingRequest { request_id, request: Arc::new(r.into()) })
672 {
673 requests.push(request);
674 }
675
676 requests.append(&mut self.inner.verification_machine.outgoing_messages());
677 requests.append(&mut self.inner.key_request_machine.outgoing_to_device_requests().await?);
678
679 Ok(requests)
680 }
681
682 pub fn query_keys_for_users<'a>(
703 &self,
704 users: impl IntoIterator<Item = &'a UserId>,
705 ) -> (OwnedTransactionId, KeysQueryRequest) {
706 self.inner.identity_manager.build_key_query_for_users(users)
707 }
708
709 pub async fn mark_request_as_sent<'a>(
719 &self,
720 request_id: &TransactionId,
721 response: impl Into<AnyIncomingResponse<'a>>,
722 ) -> OlmResult<()> {
723 match response.into() {
724 AnyIncomingResponse::KeysUpload(response) => {
725 Box::pin(self.receive_keys_upload_response(response)).await?;
726 }
727 AnyIncomingResponse::KeysQuery(response) => {
728 Box::pin(self.receive_keys_query_response(request_id, response)).await?;
729 }
730 AnyIncomingResponse::KeysClaim(response) => {
731 Box::pin(
732 self.inner.session_manager.receive_keys_claim_response(request_id, response),
733 )
734 .await?;
735 }
736 AnyIncomingResponse::ToDevice(_) => {
737 Box::pin(self.mark_to_device_request_as_sent(request_id)).await?;
738 }
739 AnyIncomingResponse::SigningKeysUpload(_) => {
740 Box::pin(self.receive_cross_signing_upload_response()).await?;
741 }
742 AnyIncomingResponse::SignatureUpload(_) => {
743 self.inner.verification_machine.mark_request_as_sent(request_id);
744 self.inner.key_request_machine.mark_outgoing_request_as_sent(request_id).await?;
745 }
746 AnyIncomingResponse::RoomMessage(_) => {
747 self.inner.verification_machine.mark_request_as_sent(request_id);
748 }
749 AnyIncomingResponse::KeysBackup(_) => {
750 Box::pin(self.inner.backup_machine.mark_request_as_sent(request_id)).await?;
751 }
752 }
753
754 Ok(())
755 }
756
757 async fn receive_cross_signing_upload_response(&self) -> StoreResult<()> {
759 let identity = self.inner.user_identity.lock().await;
760 identity.mark_as_shared();
761
762 let changes = Changes { private_identity: Some(identity.clone()), ..Default::default() };
763
764 self.store().save_changes(changes).await
765 }
766
767 pub async fn bootstrap_cross_signing(
786 &self,
787 reset: bool,
788 ) -> Result<CrossSigningBootstrapRequests, BootstrapCrossSigningError> {
789 let identity = self.inner.user_identity.lock().await.clone();
794
795 let (upload_signing_keys_req, upload_signatures_req) = if reset || identity.is_empty().await
796 {
797 info!("Creating new cross signing identity");
798
799 let (identity, upload_signing_keys_req, upload_signatures_req) = {
800 let cache = self.inner.store.cache().await?;
801 let account = cache.account().await?;
802 account
803 .bootstrap_cross_signing(
804 #[cfg(feature = "experimental-x509-identity-verification")]
805 self.inner.store.x509_signer().map(|s| s.raw()),
806 )
807 .await?
808 };
809
810 let public = identity.to_public_identity().await.expect(
811 "Couldn't create a public version of the identity from a new private identity",
812 );
813
814 *self.inner.user_identity.lock().await = identity.clone();
815
816 self.store()
817 .save_changes(Changes {
818 identities: IdentityChanges { new: vec![public.into()], ..Default::default() },
819 private_identity: Some(identity),
820 ..Default::default()
821 })
822 .await?;
823
824 (upload_signing_keys_req, upload_signatures_req)
825 } else {
826 info!("Trying to upload the existing cross signing identity");
827 let upload_signing_keys_req = identity.as_upload_request().await;
828
829 let upload_signatures_req = identity
831 .sign_account(self.inner.store.static_account())
832 .await
833 .expect("Can't sign device keys");
834
835 (upload_signing_keys_req, upload_signatures_req)
836 };
837
838 let upload_keys_req =
842 self.upload_device_keys().await?.map(|(_, request)| OutgoingRequest::from(request));
843
844 Ok(CrossSigningBootstrapRequests {
845 upload_signing_keys_req,
846 upload_keys_req,
847 upload_signatures_req,
848 })
849 }
850
851 pub async fn upload_device_keys(
863 &self,
864 ) -> StoreResult<Option<(OwnedTransactionId, UploadKeysRequest)>> {
865 let cache = self.store().cache().await?;
866 let account = cache.account().await?;
867
868 Ok(self.keys_for_upload(&account).await.map(|request| (TransactionId::new(), request)))
869 }
870
871 async fn receive_keys_upload_response(&self, response: &UploadKeysResponse) -> OlmResult<()> {
878 self.inner
879 .store
880 .with_transaction(async |tr| {
881 let account = tr.account().await?;
882 account.receive_keys_upload_response(response)?;
883 Ok(())
884 })
885 .await
886 }
887
888 #[instrument(skip_all)]
916 pub async fn get_missing_sessions(
917 &self,
918 users: impl Iterator<Item = &UserId>,
919 ) -> StoreResult<Option<(OwnedTransactionId, KeysClaimRequest)>> {
920 self.inner.session_manager.get_missing_sessions(users).await
921 }
922
923 async fn receive_keys_query_response(
932 &self,
933 request_id: &TransactionId,
934 response: &KeysQueryResponse,
935 ) -> OlmResult<(DeviceChanges, IdentityChanges)> {
936 self.inner.identity_manager.receive_keys_query_response(request_id, response).await
937 }
938
939 async fn keys_for_upload(&self, account: &Account) -> Option<UploadKeysRequest> {
948 let (mut device_keys, one_time_keys, fallback_keys) = account.keys_for_upload();
949
950 if let Some(device_keys) = &mut device_keys {
960 let private_identity = self.store().private_identity();
961 let guard = private_identity.lock().await;
962
963 if guard.status().await.is_complete() {
964 guard.sign_device_keys(device_keys).await.expect(
965 "We should be able to sign our device keys since we confirmed that we \
966 have a complete set of private cross-signing keys",
967 );
968 }
969 }
970
971 if device_keys.is_none() && one_time_keys.is_empty() && fallback_keys.is_empty() {
972 None
973 } else {
974 let device_keys = device_keys.map(|d| d.to_raw());
975
976 Some(assign!(UploadKeysRequest::new(), {
977 device_keys, one_time_keys, fallback_keys
978 }))
979 }
980 }
981
982 async fn decrypt_to_device_event(
1005 &self,
1006 transaction: &mut StoreTransaction,
1007 event: &EncryptedToDeviceEvent,
1008 changes: &mut Changes,
1009 decryption_settings: &DecryptionSettings,
1010 ) -> Result<OlmDecryptionInfo, DecryptToDeviceError> {
1011 let mut decrypted = transaction
1013 .account()
1014 .await?
1015 .decrypt_to_device_event(&self.inner.store, event, decryption_settings)
1016 .await?;
1017
1018 self.check_to_device_event_is_not_from_dehydrated_device(&decrypted, &event.sender).await?;
1020
1021 self.handle_decrypted_to_device_event(transaction.cache(), &mut decrypted, changes).await?;
1023
1024 Ok(decrypted)
1025 }
1026
1027 #[instrument(
1028 skip_all,
1029 fields(room_id = ? content.room_id, session_id, message_index, shared_history = content.shared_history)
1033 )]
1034 async fn handle_key(
1035 &self,
1036 sender_key: Curve25519PublicKey,
1037 event: &DecryptedRoomKeyEvent,
1038 content: &MegolmV1AesSha2Content,
1039 ) -> OlmResult<Option<InboundGroupSession>> {
1040 let session =
1041 InboundGroupSession::from_room_key_content(sender_key, event.keys.ed25519, content);
1042
1043 match session {
1044 Ok(mut session) => {
1045 Span::current().record("session_id", session.session_id());
1046 Span::current().record("message_index", session.first_known_index());
1047
1048 let sender_data =
1049 SenderDataFinder::find_using_event(self.store(), sender_key, event, &session)
1050 .await?;
1051 session.sender_data = sender_data;
1052
1053 Ok(self.store().merge_received_group_session(session).await?)
1054 }
1055 Err(e) => {
1056 Span::current().record("session_id", &content.session_id);
1057 warn!("Received a room key event which contained an invalid session key: {e}");
1058
1059 Ok(None)
1060 }
1061 }
1062 }
1063
1064 #[instrument(skip_all, fields(algorithm = ?event.content.algorithm()))]
1066 async fn add_room_key(
1067 &self,
1068 sender_key: Curve25519PublicKey,
1069 event: &DecryptedRoomKeyEvent,
1070 ) -> OlmResult<Option<InboundGroupSession>> {
1071 match &event.content {
1072 RoomKeyContent::MegolmV1AesSha2(content) => {
1073 self.handle_key(sender_key, event, content).await
1074 }
1075 #[cfg(feature = "experimental-algorithms")]
1076 RoomKeyContent::MegolmV2AesSha2(content) => {
1077 self.handle_key(sender_key, event, content).await
1078 }
1079 RoomKeyContent::Unknown(_) => {
1080 warn!("Received a room key with an unsupported algorithm");
1081 Ok(None)
1082 }
1083 }
1084 }
1085
1086 #[instrument()]
1088 async fn receive_room_key_bundle_data(
1089 &self,
1090 sender_key: Curve25519PublicKey,
1091 event: &DecryptedRoomKeyBundleEvent,
1092 changes: &mut Changes,
1093 ) -> OlmResult<()> {
1094 let Some(sender_device_keys) = &event.sender_device_keys else {
1095 warn!("Received a room key bundle with no sender device keys: ignoring");
1096 return Ok(());
1097 };
1098
1099 let sender_device_data =
1104 DeviceData::try_from(sender_device_keys).expect("failed to verify sender device keys");
1105 let sender_device = self.store().wrap_device_data(sender_device_data).await?;
1106
1107 changes.received_room_key_bundles.push(StoredRoomKeyBundleData {
1108 sender_user: event.sender.clone(),
1109 sender_data: SenderData::from_device(&sender_device),
1110 sender_key,
1111 bundle_data: event.content.clone(),
1112 });
1113 Ok(())
1114 }
1115
1116 fn add_withheld_info(&self, changes: &mut Changes, event: &RoomKeyWithheldEvent) {
1117 debug!(?event.content, "Processing `m.room_key.withheld` event");
1118
1119 if let RoomKeyWithheldContent::MegolmV1AesSha2(
1120 MegolmV1AesSha2WithheldContent::BlackListed(c)
1121 | MegolmV1AesSha2WithheldContent::Unverified(c)
1122 | MegolmV1AesSha2WithheldContent::Unauthorised(c)
1123 | MegolmV1AesSha2WithheldContent::Unavailable(c),
1124 ) = &event.content
1125 {
1126 changes
1127 .withheld_session_info
1128 .entry(c.room_id.to_owned())
1129 .or_default()
1130 .insert(c.session_id.to_owned(), event.to_owned().into());
1131 }
1132 }
1133
1134 #[cfg(test)]
1135 pub(crate) async fn create_outbound_group_session_with_defaults_test_helper(
1136 &self,
1137 room_id: &RoomId,
1138 ) -> OlmResult<()> {
1139 let (_, session) = self
1140 .inner
1141 .group_session_manager
1142 .create_outbound_group_session(
1143 room_id,
1144 EncryptionSettings::default(),
1145 SenderData::unknown(),
1146 )
1147 .await?;
1148
1149 self.store().save_inbound_group_sessions(&[session]).await?;
1150
1151 Ok(())
1152 }
1153
1154 #[cfg(test)]
1155 #[allow(dead_code)]
1156 pub(crate) async fn create_inbound_session_test_helper(
1157 &self,
1158 room_id: &RoomId,
1159 ) -> OlmResult<InboundGroupSession> {
1160 let (_, session) = self
1161 .inner
1162 .group_session_manager
1163 .create_outbound_group_session(
1164 room_id,
1165 EncryptionSettings::default(),
1166 SenderData::unknown(),
1167 )
1168 .await?;
1169
1170 Ok(session)
1171 }
1172
1173 pub async fn encrypt_room_event(
1190 &self,
1191 room_id: &RoomId,
1192 content: impl MessageLikeEventContent,
1193 ) -> MegolmResult<RawEncryptionResult> {
1194 let event_type = content.event_type().to_string();
1195 let content = Raw::new(&content)?.cast_unchecked();
1196 self.encrypt_room_event_raw(room_id, &event_type, &content).await
1197 }
1198
1199 pub async fn encrypt_room_event_raw(
1219 &self,
1220 room_id: &RoomId,
1221 event_type: &str,
1222 content: &Raw<AnyMessageLikeEventContent>,
1223 ) -> MegolmResult<RawEncryptionResult> {
1224 self.inner.group_session_manager.encrypt(room_id, event_type, content).await.map(|result| {
1225 RawEncryptionResult {
1226 content: result.content,
1227 encryption_info: self
1228 .own_encryption_info(result.algorithm, result.session_id.to_string()),
1229 }
1230 })
1231 }
1232
1233 fn own_encryption_info(
1234 &self,
1235 algorithm: EventEncryptionAlgorithm,
1236 session_id: String,
1237 ) -> EncryptionInfo {
1238 let identity_keys = self.identity_keys();
1239
1240 let algorithm_info = match algorithm {
1241 EventEncryptionAlgorithm::MegolmV1AesSha2 => AlgorithmInfo::MegolmV1AesSha2 {
1242 curve25519_key: identity_keys.curve25519.to_base64(),
1243 sender_claimed_keys: BTreeMap::from([(
1244 DeviceKeyAlgorithm::Ed25519,
1245 identity_keys.ed25519.to_base64(),
1246 )]),
1247 session_id: Some(session_id),
1248 },
1249 EventEncryptionAlgorithm::OlmV1Curve25519AesSha2 => {
1250 AlgorithmInfo::OlmV1Curve25519AesSha2 {
1251 curve25519_public_key_base64: identity_keys.curve25519.to_base64(),
1252 }
1253 }
1254 _ => unreachable!(
1255 "Only MegolmV1AesSha2 and OlmV1Curve25519AesSha2 are supported on this level"
1256 ),
1257 };
1258
1259 EncryptionInfo {
1260 sender: self.inner.user_id.clone(),
1261 sender_device: Some(self.inner.device_id.clone()),
1262 forwarder: None,
1263 algorithm_info,
1264 verification_state: VerificationState::Verified,
1265 }
1266 }
1267
1268 #[cfg(feature = "experimental-encrypted-state-events")]
1280 pub async fn encrypt_state_event<C, K>(
1281 &self,
1282 room_id: &RoomId,
1283 content: C,
1284 state_key: K,
1285 ) -> MegolmResult<Raw<RoomEncryptedEventContent>>
1286 where
1287 C: StateEventContent,
1288 C::StateKey: Borrow<K>,
1289 K: AsRef<str>,
1290 {
1291 let event_type = content.event_type().to_string();
1292 let content = Raw::new(&content)?.cast_unchecked();
1293 self.encrypt_state_event_raw(room_id, &event_type, state_key.as_ref(), &content).await
1294 }
1295
1296 #[cfg(feature = "experimental-encrypted-state-events")]
1315 pub async fn encrypt_state_event_raw(
1316 &self,
1317 room_id: &RoomId,
1318 event_type: &str,
1319 state_key: &str,
1320 content: &Raw<AnyStateEventContent>,
1321 ) -> MegolmResult<Raw<RoomEncryptedEventContent>> {
1322 self.inner
1323 .group_session_manager
1324 .encrypt_state(room_id, event_type, state_key, content)
1325 .await
1326 }
1327
1328 pub async fn discard_room_key(&self, room_id: &RoomId) -> StoreResult<bool> {
1339 self.inner.group_session_manager.invalidate_group_session(room_id).await
1340 }
1341
1342 pub async fn share_room_key(
1362 &self,
1363 room_id: &RoomId,
1364 users: impl Iterator<Item = &UserId>,
1365 encryption_settings: impl Into<EncryptionSettings>,
1366 ) -> OlmResult<Vec<Arc<ToDeviceRequest>>> {
1367 self.inner.group_session_manager.share_room_key(room_id, users, encryption_settings).await
1368 }
1369
1370 #[cfg(feature = "experimental-send-custom-to-device")]
1384 pub async fn encrypt_content_for_devices(
1385 &self,
1386 devices: Vec<DeviceData>,
1387 event_type: &str,
1388 content: &Value,
1389 share_strategy: CollectStrategy,
1390 ) -> OlmResult<(Vec<ToDeviceRequest>, Vec<(DeviceData, WithheldCode)>)> {
1391 let mut changes = Changes::default();
1392
1393 let (allowed_devices, mut blocked_devices) =
1394 split_devices_for_share_strategy(&self.inner.store, devices, share_strategy).await?;
1395
1396 let result = self
1397 .inner
1398 .group_session_manager
1399 .encrypt_content_for_devices(allowed_devices, event_type, content.clone(), &mut changes)
1400 .await;
1401
1402 if !changes.is_empty() {
1404 let session_count = changes.sessions.len();
1405
1406 self.inner.store.save_changes(changes).await?;
1407
1408 trace!(
1409 session_count = session_count,
1410 "Stored the changed sessions after encrypting a custom to-device event"
1411 );
1412 }
1413
1414 result.map(|(to_device_requests, mut withheld)| {
1415 withheld.append(&mut blocked_devices);
1416 (to_device_requests, withheld)
1417 })
1418 }
1419 pub async fn share_room_key_bundle_data(
1424 &self,
1425 user_id: &UserId,
1426 collect_strategy: &CollectStrategy,
1427 bundle_data: RoomKeyBundleContent,
1428 ) -> OlmResult<Vec<ToDeviceRequest>> {
1429 self.inner
1430 .group_session_manager
1431 .share_room_key_bundle_data(user_id, collect_strategy, bundle_data)
1432 .await
1433 }
1434
1435 #[deprecated(note = "Use OlmMachine::receive_verification_event instead", since = "0.7.0")]
1443 pub async fn receive_unencrypted_verification_event(
1444 &self,
1445 event: &AnyMessageLikeEvent,
1446 ) -> StoreResult<()> {
1447 self.inner.verification_machine.receive_any_event(event).await
1448 }
1449
1450 pub async fn receive_verification_event(&self, event: &AnyMessageLikeEvent) -> StoreResult<()> {
1463 self.inner.verification_machine.receive_any_event(event).await
1464 }
1465
1466 #[instrument(
1472 skip_all,
1473 fields(
1474 sender_key = ?decrypted.result.sender_key,
1475 event_type = decrypted.result.event.event_type(),
1476 ),
1477 )]
1478 async fn handle_decrypted_to_device_event(
1479 &self,
1480 cache: &StoreCache,
1481 decrypted: &mut OlmDecryptionInfo,
1482 changes: &mut Changes,
1483 ) -> OlmResult<()> {
1484 debug!(
1485 sender_device_keys =
1486 ?decrypted.result.event.sender_device_keys().map(|k| (k.curve25519_key(), k.ed25519_key())).unwrap_or((None, None)),
1487 "Received a decrypted to-device event",
1488 );
1489
1490 match &*decrypted.result.event {
1491 AnyDecryptedOlmEvent::RoomKey(e) => {
1492 let session = self.add_room_key(decrypted.result.sender_key, e).await?;
1493 decrypted.inbound_group_session = session;
1494 }
1495 AnyDecryptedOlmEvent::ForwardedRoomKey(e) => {
1496 let session = self
1497 .inner
1498 .key_request_machine
1499 .receive_forwarded_room_key(decrypted.result.sender_key, e)
1500 .await?;
1501 decrypted.inbound_group_session = session;
1502 }
1503 AnyDecryptedOlmEvent::SecretSend(e) => {
1504 let name = self
1505 .inner
1506 .key_request_machine
1507 .receive_secret_event(cache, decrypted.result.sender_key, e, changes)
1508 .await?;
1509
1510 if let Ok(ToDeviceEvents::SecretSend(mut e)) =
1513 decrypted.result.raw_event.deserialize_as()
1514 {
1515 e.content.secret_name = name;
1516 decrypted.result.raw_event = Raw::from_json(to_raw_value(&e)?);
1517 }
1518
1519 if enabled!(tracing::Level::DEBUG) {
1520 let cross_signing_status = self.cross_signing_status().await;
1521 let backup_enabled = self.backup_machine().enabled().await;
1522 debug!(
1523 ?cross_signing_status,
1524 backup_enabled, "Status after receiving secret event"
1525 );
1526 }
1527 }
1528 AnyDecryptedOlmEvent::Dummy(_) => {
1529 debug!("Received an `m.dummy` event");
1530 }
1531 AnyDecryptedOlmEvent::RoomKeyBundle(e) => {
1532 debug!("Received a room key bundle event {:?}", e);
1533 self.receive_room_key_bundle_data(decrypted.result.sender_key, e, changes).await?;
1534 }
1535 #[cfg(feature = "experimental-push-secrets")]
1536 AnyDecryptedOlmEvent::SecretPush(e) => {
1537 self.inner
1538 .key_request_machine
1539 .receive_secret_push_event(&decrypted.result.sender_key, e, changes)
1540 .await?;
1541 }
1542 AnyDecryptedOlmEvent::Custom(_) => {
1543 warn!("Received an unexpected encrypted to-device event");
1544 }
1545 }
1546
1547 Ok(())
1548 }
1549
1550 async fn handle_verification_event(&self, event: &ToDeviceEvents) {
1551 if let Err(e) = self.inner.verification_machine.receive_any_event(event).await {
1552 error!("Error handling a verification event: {e:?}");
1553 }
1554 }
1555
1556 async fn mark_to_device_request_as_sent(&self, request_id: &TransactionId) -> StoreResult<()> {
1558 self.inner.verification_machine.mark_request_as_sent(request_id);
1559 self.inner.key_request_machine.mark_outgoing_request_as_sent(request_id).await?;
1560 self.inner.group_session_manager.mark_request_as_sent(request_id).await?;
1561 self.inner.session_manager.mark_outgoing_request_as_sent(request_id);
1562 Ok(())
1563 }
1564
1565 pub fn get_verification(&self, user_id: &UserId, flow_id: &str) -> Option<Verification> {
1567 self.inner.verification_machine.get_verification(user_id, flow_id)
1568 }
1569
1570 pub fn get_verification_request(
1572 &self,
1573 user_id: &UserId,
1574 flow_id: impl AsRef<str>,
1575 ) -> Option<VerificationRequest> {
1576 self.inner.verification_machine.get_request(user_id, flow_id)
1577 }
1578
1579 pub fn get_verification_requests(&self, user_id: &UserId) -> Vec<VerificationRequest> {
1581 self.inner.verification_machine.get_requests(user_id)
1582 }
1583
1584 async fn handle_to_device_event(&self, changes: &mut Changes, event: &ToDeviceEvents) {
1589 use crate::types::events::ToDeviceEvents::*;
1590
1591 match event {
1592 RoomKeyRequest(e) => self.inner.key_request_machine.receive_incoming_key_request(e),
1598 SecretRequest(e) => self.inner.key_request_machine.receive_incoming_secret_request(e),
1599 RoomKeyWithheld(e) => self.add_withheld_info(changes, e),
1600 KeyVerificationAccept(..)
1601 | KeyVerificationCancel(..)
1602 | KeyVerificationKey(..)
1603 | KeyVerificationMac(..)
1604 | KeyVerificationRequest(..)
1605 | KeyVerificationReady(..)
1606 | KeyVerificationDone(..)
1607 | KeyVerificationStart(..) => {
1608 self.handle_verification_event(event).await;
1609 }
1610
1611 Custom(_) | Dummy(_) => {}
1613
1614 RoomEncrypted(_) => {}
1616
1617 SecretSend(_) | RoomKey(_) | ForwardedRoomKey(_) => {}
1620 }
1621 }
1622
1623 fn record_message_id(event: &Raw<AnyToDeviceEvent>) {
1624 use serde::Deserialize;
1625
1626 #[derive(Deserialize)]
1627 struct ContentStub<'a> {
1628 #[serde(borrow, rename = "org.matrix.msgid")]
1629 message_id: Option<&'a str>,
1630 }
1631 #[derive(Deserialize)]
1632 struct ToDeviceStub<'a> {
1633 sender: &'a str,
1634 #[serde(rename = "type")]
1635 event_type: &'a str,
1636 #[serde(borrow)]
1637 content: ContentStub<'a>,
1638 }
1639
1640 if let Ok(event) = event.deserialize_as_unchecked::<ToDeviceStub<'_>>() {
1641 Span::current().record("sender", event.sender);
1642 Span::current().record("event_type", event.event_type);
1643 Span::current().record("message_id", event.content.message_id);
1644 }
1645 }
1646
1647 #[instrument(skip_all, fields(sender, event_type, message_id))]
1655 async fn receive_to_device_event(
1656 &self,
1657 transaction: &mut StoreTransaction,
1658 changes: &mut Changes,
1659 raw_event: Raw<AnyToDeviceEvent>,
1660 decryption_settings: &DecryptionSettings,
1661 ) -> Option<ProcessedToDeviceEvent> {
1662 Self::record_message_id(&raw_event);
1663
1664 let event: ToDeviceEvents = match raw_event.deserialize_as() {
1665 Ok(e) => e,
1666 Err(e) => {
1667 warn!("Received an invalid to-device event: {e}");
1669 return Some(ProcessedToDeviceEvent::Invalid(raw_event));
1670 }
1671 };
1672
1673 debug!("Received a to-device event");
1674
1675 match event {
1676 ToDeviceEvents::RoomEncrypted(e) => {
1677 self.receive_encrypted_to_device_event(
1678 transaction,
1679 changes,
1680 raw_event,
1681 e,
1682 decryption_settings,
1683 )
1684 .await
1685 }
1686 e => {
1687 self.handle_to_device_event(changes, &e).await;
1688 Some(ProcessedToDeviceEvent::PlainText(raw_event))
1689 }
1690 }
1691 }
1692
1693 async fn receive_encrypted_to_device_event(
1707 &self,
1708 transaction: &mut StoreTransaction,
1709 changes: &mut Changes,
1710 mut raw_event: Raw<AnyToDeviceEvent>,
1711 e: ToDeviceEvent<ToDeviceEncryptedEventContent>,
1712 decryption_settings: &DecryptionSettings,
1713 ) -> Option<ProcessedToDeviceEvent> {
1714 let decrypted = match self
1715 .decrypt_to_device_event(transaction, &e, changes, decryption_settings)
1716 .await
1717 {
1718 Ok(decrypted) => decrypted,
1719 Err(DecryptToDeviceError::OlmError(err)) => {
1720 let reason = if let OlmError::UnverifiedSenderDevice = &err {
1721 ToDeviceUnableToDecryptReason::UnverifiedSenderDevice
1722 } else {
1723 ToDeviceUnableToDecryptReason::DecryptionFailure
1724 };
1725
1726 if let OlmError::SessionWedged(sender, curve_key) = err
1727 && let Err(e) =
1728 self.inner.session_manager.mark_device_as_wedged(&sender, curve_key).await
1729 {
1730 error!(
1731 error = ?e,
1732 "Couldn't mark device to be unwedged",
1733 );
1734 }
1735
1736 return Some(ProcessedToDeviceEvent::UnableToDecrypt {
1737 encrypted_event: raw_event,
1738 utd_info: ToDeviceUnableToDecryptInfo { reason },
1739 });
1740 }
1741 Err(DecryptToDeviceError::FromDehydratedDevice) => return None,
1742 };
1743
1744 match decrypted.session {
1747 SessionType::New(s) | SessionType::Existing(s) => {
1748 changes.sessions.push(s);
1749 }
1750 }
1751
1752 changes.message_hashes.push(decrypted.message_hash);
1753
1754 if let Some(group_session) = decrypted.inbound_group_session {
1755 changes.inbound_group_sessions.push(group_session);
1756 }
1757
1758 match decrypted.result.raw_event.deserialize_as() {
1759 Ok(event) => {
1760 self.handle_to_device_event(changes, &event).await;
1761
1762 raw_event = event
1763 .serialize_zeroized()
1764 .expect("Zeroizing and reserializing our events should always work")
1765 .cast();
1766 }
1767 Err(e) => {
1768 warn!("Received an invalid encrypted to-device event: {e}");
1769 raw_event = decrypted.result.raw_event;
1770 }
1771 }
1772
1773 Some(ProcessedToDeviceEvent::Decrypted {
1774 raw: raw_event,
1775 encryption_info: decrypted.result.encryption_info,
1776 })
1777 }
1778
1779 async fn check_to_device_event_is_not_from_dehydrated_device(
1782 &self,
1783 decrypted: &OlmDecryptionInfo,
1784 sender_user_id: &UserId,
1785 ) -> Result<(), DecryptToDeviceError> {
1786 if self.to_device_event_is_from_dehydrated_device(decrypted, sender_user_id).await? {
1787 warn!(
1788 sender = ?sender_user_id,
1789 session = ?decrypted.session,
1790 "Received a to-device event from a dehydrated device. This is unexpected: ignoring event"
1791 );
1792 Err(DecryptToDeviceError::FromDehydratedDevice)
1793 } else {
1794 Ok(())
1795 }
1796 }
1797
1798 async fn to_device_event_is_from_dehydrated_device(
1804 &self,
1805 decrypted: &OlmDecryptionInfo,
1806 sender_user_id: &UserId,
1807 ) -> OlmResult<bool> {
1808 if let Some(device_keys) = decrypted.result.event.sender_device_keys() {
1810 if device_keys.dehydrated.unwrap_or(false) {
1816 return Ok(true);
1817 }
1818 }
1823
1824 Ok(self
1826 .store()
1827 .get_device_from_curve_key(sender_user_id, decrypted.result.sender_key)
1828 .await?
1829 .is_some_and(|d| d.is_dehydrated()))
1830 }
1831
1832 #[instrument(skip_all)]
1850 pub async fn receive_sync_changes(
1851 &self,
1852 sync_changes: EncryptionSyncChanges<'_>,
1853 decryption_settings: &DecryptionSettings,
1854 ) -> OlmResult<(Vec<ProcessedToDeviceEvent>, Vec<RoomKeyInfo>)> {
1855 let mut store_transaction = self.inner.store.transaction().await;
1856
1857 let (events, changes) = self
1858 .preprocess_sync_changes(&mut store_transaction, sync_changes, decryption_settings)
1859 .await?;
1860
1861 let room_key_updates: Vec<_> =
1864 changes.inbound_group_sessions.iter().map(RoomKeyInfo::from).collect();
1865
1866 self.store().save_changes(changes).await?;
1867 store_transaction.commit().await?;
1868
1869 Ok((events, room_key_updates))
1870 }
1871
1872 pub(crate) async fn preprocess_sync_changes(
1890 &self,
1891 transaction: &mut StoreTransaction,
1892 sync_changes: EncryptionSyncChanges<'_>,
1893 decryption_settings: &DecryptionSettings,
1894 ) -> OlmResult<(Vec<ProcessedToDeviceEvent>, Changes)> {
1895 let mut events: Vec<ProcessedToDeviceEvent> = self
1897 .inner
1898 .verification_machine
1899 .garbage_collect()
1900 .iter()
1901 .map(|e| ProcessedToDeviceEvent::PlainText(e.clone()))
1905 .collect();
1906 let mut changes = Default::default();
1909
1910 {
1911 let account = transaction.account().await?;
1912 account.update_key_counts(
1913 sync_changes.one_time_keys_counts,
1914 sync_changes.unused_fallback_keys,
1915 )
1916 }
1917
1918 if let Err(e) = self
1919 .inner
1920 .identity_manager
1921 .receive_device_changes(
1922 transaction.cache(),
1923 sync_changes.changed_devices.changed.iter().map(|u| u.as_ref()),
1924 )
1925 .await
1926 {
1927 error!(error = ?e, "Error marking a tracked user as changed");
1928 }
1929
1930 for raw_event in sync_changes.to_device_events {
1931 let processed_event = Box::pin(self.receive_to_device_event(
1932 transaction,
1933 &mut changes,
1934 raw_event,
1935 decryption_settings,
1936 ))
1937 .await;
1938
1939 if let Some(processed_event) = processed_event {
1940 events.push(processed_event);
1941 }
1942 }
1943
1944 let changed_sessions = self
1945 .inner
1946 .key_request_machine
1947 .collect_incoming_key_requests(transaction.cache())
1948 .await?;
1949
1950 changes.sessions.extend(changed_sessions);
1951 changes.next_batch_token = sync_changes.next_batch_token;
1952
1953 Ok((events, changes))
1954 }
1955
1956 pub async fn request_room_key(
1973 &self,
1974 event: &Raw<EncryptedEvent>,
1975 room_id: &RoomId,
1976 ) -> MegolmResult<(Option<OutgoingRequest>, OutgoingRequest)> {
1977 let event = event.deserialize()?;
1978 self.inner.key_request_machine.request_key(room_id, &event).await
1979 }
1980
1981 async fn get_room_event_verification_state(
1994 &self,
1995 session: &InboundGroupSession,
1996 sender: &UserId,
1997 ) -> MegolmResult<(VerificationState, Option<OwnedDeviceId>)> {
1998 let sender_data = self.get_or_update_sender_data(session, sender).await?;
1999
2000 let (verification_state, device_id) = match sender_data.user_id() {
2009 Some(i) if i != sender => {
2010 (VerificationState::Unverified(VerificationLevel::MismatchedSender), None)
2011 }
2012
2013 Some(_) | None => {
2014 sender_data_to_verification_state(sender_data, session.has_been_imported())
2015 }
2016 };
2017
2018 Ok((verification_state, device_id))
2019 }
2020
2021 async fn get_or_update_sender_data(
2036 &self,
2037 session: &InboundGroupSession,
2038 sender: &UserId,
2039 ) -> MegolmResult<SenderData> {
2040 let sender_data = if session.sender_data.should_recalculate() {
2041 let calculated_sender_data = SenderDataFinder::find_using_curve_key(
2060 self.store(),
2061 session.sender_key(),
2062 sender,
2063 session,
2064 )
2065 .await?;
2066
2067 if calculated_sender_data.compare_trust_level(&session.sender_data).is_gt() {
2069 let mut new_session = session.clone();
2071 new_session.sender_data = calculated_sender_data.clone();
2072 self.store().save_inbound_group_sessions(&[new_session]).await?;
2073
2074 calculated_sender_data
2076 } else {
2077 session.sender_data.clone()
2079 }
2080 } else {
2081 session.sender_data.clone()
2082 };
2083
2084 Ok(sender_data)
2085 }
2086
2087 pub async fn query_missing_secrets_from_other_sessions(&self) -> StoreResult<bool> {
2112 let identity = self.inner.user_identity.lock().await;
2113 let mut secrets = identity.get_missing_secrets().await;
2114
2115 if self.store().load_backup_keys().await?.decryption_key.is_none() {
2116 secrets.push(SecretName::RecoveryKey);
2117 }
2118
2119 if secrets.is_empty() {
2120 debug!("No missing requests to query");
2121 return Ok(false);
2122 }
2123
2124 let secret_requests = GossipMachine::request_missing_secrets(self.user_id(), secrets);
2125
2126 let unsent_request = self.store().get_unsent_secret_requests().await?;
2128 let not_yet_requested = secret_requests
2129 .into_iter()
2130 .filter(|request| !unsent_request.iter().any(|unsent| unsent.info == request.info))
2131 .collect_vec();
2132
2133 if not_yet_requested.is_empty() {
2134 debug!("The missing secrets have already been requested");
2135 Ok(false)
2136 } else {
2137 debug!("Requesting missing secrets");
2138
2139 let changes = Changes { key_requests: not_yet_requested, ..Default::default() };
2140
2141 self.store().save_changes(changes).await?;
2142 Ok(true)
2143 }
2144 }
2145
2146 #[cfg(feature = "experimental-push-secrets")]
2154 pub async fn push_secret_to_verified_devices(
2155 &self,
2156 secret_name: SecretName,
2157 ) -> Result<HashMap<OwnedDeviceId, OlmError>, SecretPushError> {
2158 self.inner.key_request_machine.push_secret_to_verified_devices(secret_name).await
2159 }
2160
2161 async fn get_encryption_info(
2167 &self,
2168 session: &InboundGroupSession,
2169 sender: &UserId,
2170 ) -> MegolmResult<Arc<EncryptionInfo>> {
2171 let (verification_state, device_id) =
2172 self.get_room_event_verification_state(session, sender).await?;
2173
2174 Ok(Arc::new(EncryptionInfo {
2175 sender: sender.to_owned(),
2176 sender_device: device_id,
2177 forwarder: session.forwarder_data.as_ref().and_then(|data| {
2178 data.device_id().map(|device_id| ForwarderInfo {
2182 device_id: device_id.to_owned(),
2183 user_id: data.user_id().to_owned(),
2184 })
2185 }),
2186 algorithm_info: AlgorithmInfo::MegolmV1AesSha2 {
2187 curve25519_key: session.sender_key().to_base64(),
2188 sender_claimed_keys: session
2189 .signing_keys()
2190 .iter()
2191 .map(|(k, v)| (k.to_owned(), v.to_base64()))
2192 .collect(),
2193 session_id: Some(session.session_id().to_owned()),
2194 },
2195 verification_state,
2196 }))
2197 }
2198
2199 async fn decrypt_megolm_events(
2200 &self,
2201 room_id: &RoomId,
2202 event: &EncryptedEvent,
2203 content: &SupportedEventEncryptionSchemes<'_>,
2204 decryption_settings: &DecryptionSettings,
2205 ) -> MegolmResult<(JsonObject, Arc<EncryptionInfo>)> {
2206 let session =
2207 self.get_inbound_group_session_or_error(room_id, content.session_id()).await?;
2208
2209 Span::current().record("sender_key", debug(session.sender_key()));
2215
2216 let result = session.decrypt(event).await;
2217 match result {
2218 Ok((decrypted_event, _)) => {
2219 let encryption_info = self.get_encryption_info(&session, &event.sender).await?;
2220
2221 self.check_sender_trust_requirement(
2222 &session,
2223 &encryption_info,
2224 &decryption_settings.sender_device_trust_requirement,
2225 )?;
2226
2227 Ok((decrypted_event, encryption_info))
2228 }
2229 Err(error) => Err(
2230 if let MegolmError::Decryption(DecryptionError::UnknownMessageIndex(_, _)) = error {
2231 let withheld_code = self
2232 .inner
2233 .store
2234 .get_withheld_info(room_id, content.session_id())
2235 .await?
2236 .map(|e| e.content.withheld_code());
2237
2238 if withheld_code.is_some() {
2239 MegolmError::MissingRoomKey(withheld_code)
2241 } else {
2242 error
2243 }
2244 } else {
2245 error
2246 },
2247 ),
2248 }
2249 }
2250
2251 fn check_sender_trust_requirement(
2257 &self,
2258 session: &InboundGroupSession,
2259 encryption_info: &EncryptionInfo,
2260 trust_requirement: &TrustRequirement,
2261 ) -> MegolmResult<()> {
2262 trace!(
2263 verification_state = ?encryption_info.verification_state,
2264 ?trust_requirement, "check_sender_trust_requirement",
2265 );
2266
2267 let verification_level = match &encryption_info.verification_state {
2270 VerificationState::Verified => return Ok(()),
2271 VerificationState::Unverified(verification_level) => verification_level,
2272 };
2273
2274 let ok = match trust_requirement {
2275 TrustRequirement::Untrusted => true,
2276
2277 TrustRequirement::CrossSignedOrLegacy => {
2278 let legacy_session = match session.sender_data {
2284 SenderData::DeviceInfo { legacy_session, .. } => legacy_session,
2285 SenderData::UnknownDevice { legacy_session, .. } => legacy_session,
2286 _ => false,
2287 };
2288
2289 match (verification_level, legacy_session) {
2299 (VerificationLevel::UnverifiedIdentity, _) => true,
2301
2302 (VerificationLevel::UnsignedDevice, true) => true,
2304
2305 (VerificationLevel::None(_), true) => true,
2307
2308 (VerificationLevel::VerificationViolation, _)
2310 | (VerificationLevel::MismatchedSender, _)
2311 | (VerificationLevel::UnsignedDevice, false)
2312 | (VerificationLevel::None(_), false) => false,
2313 }
2314 }
2315
2316 TrustRequirement::CrossSigned => match verification_level {
2319 VerificationLevel::UnverifiedIdentity => true,
2320
2321 VerificationLevel::VerificationViolation
2322 | VerificationLevel::MismatchedSender
2323 | VerificationLevel::UnsignedDevice
2324 | VerificationLevel::None(_) => false,
2325 },
2326 };
2327
2328 if ok {
2329 Ok(())
2330 } else {
2331 Err(MegolmError::SenderIdentityNotTrusted(verification_level.clone()))
2332 }
2333 }
2334
2335 async fn get_inbound_group_session_or_error(
2340 &self,
2341 room_id: &RoomId,
2342 session_id: &str,
2343 ) -> MegolmResult<InboundGroupSession> {
2344 match self.store().get_inbound_group_session(room_id, session_id).await? {
2345 Some(session) => Ok(session),
2346 None => {
2347 let withheld_code = self
2348 .inner
2349 .store
2350 .get_withheld_info(room_id, session_id)
2351 .await?
2352 .map(|e| e.content.withheld_code());
2353 Err(MegolmError::MissingRoomKey(withheld_code))
2354 }
2355 }
2356 }
2357
2358 pub async fn try_decrypt_room_event(
2373 &self,
2374 raw_event: &Raw<EncryptedEvent>,
2375 room_id: &RoomId,
2376 decryption_settings: &DecryptionSettings,
2377 ) -> Result<RoomEventDecryptionResult, CryptoStoreError> {
2378 match self.decrypt_room_event_inner(raw_event, room_id, true, decryption_settings).await {
2379 Ok(decrypted) => Ok(RoomEventDecryptionResult::Decrypted(decrypted)),
2380 Err(err) => Ok(RoomEventDecryptionResult::UnableToDecrypt(megolm_error_to_utd_info(
2381 raw_event, err,
2382 )?)),
2383 }
2384 }
2385
2386 pub async fn decrypt_room_event(
2394 &self,
2395 event: &Raw<EncryptedEvent>,
2396 room_id: &RoomId,
2397 decryption_settings: &DecryptionSettings,
2398 ) -> MegolmResult<DecryptedRoomEvent> {
2399 self.decrypt_room_event_inner(event, room_id, true, decryption_settings).await
2400 }
2401
2402 #[instrument(name = "decrypt_room_event", skip_all, fields(?room_id, event_id, origin_server_ts, sender, algorithm, session_id, message_index, sender_key))]
2403 async fn decrypt_room_event_inner(
2404 &self,
2405 event: &Raw<EncryptedEvent>,
2406 room_id: &RoomId,
2407 decrypt_unsigned: bool,
2408 decryption_settings: &DecryptionSettings,
2409 ) -> MegolmResult<DecryptedRoomEvent> {
2410 let _timer = timer!(tracing::Level::TRACE, "_method");
2411
2412 let event = event.deserialize()?;
2413
2414 Span::current()
2415 .record("sender", debug(&event.sender))
2416 .record("event_id", debug(&event.event_id))
2417 .record(
2418 "origin_server_ts",
2419 timestamp_to_iso8601(event.origin_server_ts)
2420 .unwrap_or_else(|| "<out of range>".to_owned()),
2421 )
2422 .record("algorithm", debug(event.content.algorithm()));
2423
2424 let content: SupportedEventEncryptionSchemes<'_> = match &event.content.scheme {
2425 RoomEventEncryptionScheme::MegolmV1AesSha2(c) => {
2426 Span::current().record("sender_key", debug(c.sender_key));
2427 c.into()
2428 }
2429 #[cfg(feature = "experimental-algorithms")]
2430 RoomEventEncryptionScheme::MegolmV2AesSha2(c) => c.into(),
2431 RoomEventEncryptionScheme::Unknown(_) => {
2432 warn!("Received an encrypted room event with an unsupported algorithm");
2433 return Err(EventError::UnsupportedAlgorithm.into());
2434 }
2435 };
2436
2437 Span::current().record("session_id", content.session_id());
2438 Span::current().record("message_index", content.message_index());
2439
2440 let result =
2441 self.decrypt_megolm_events(room_id, &event, &content, decryption_settings).await;
2442
2443 if let Err(e) = &result {
2444 #[cfg(feature = "automatic-room-key-forwarding")]
2445 match e {
2446 MegolmError::MissingRoomKey(_)
2449 | MegolmError::Decryption(DecryptionError::UnknownMessageIndex(_, _)) => {
2450 self.inner
2451 .key_request_machine
2452 .create_outgoing_key_request(room_id, &event)
2453 .await?;
2454 }
2455 _ => {}
2456 }
2457
2458 warn!("Failed to decrypt a room event: {e}");
2459 }
2460
2461 let (mut decrypted_event, encryption_info) = result?;
2462
2463 let mut unsigned_encryption_info = None;
2464 if decrypt_unsigned {
2465 unsigned_encryption_info = self
2467 .decrypt_unsigned_events(&mut decrypted_event, room_id, decryption_settings)
2468 .await;
2469 }
2470
2471 let decrypted_event =
2472 serde_json::from_value::<Raw<AnyTimelineEvent>>(decrypted_event.into())?;
2473
2474 #[cfg(feature = "experimental-encrypted-state-events")]
2475 self.verify_packed_state_key(&event, &decrypted_event)?;
2476
2477 Ok(DecryptedRoomEvent { event: decrypted_event, encryption_info, unsigned_encryption_info })
2478 }
2479
2480 #[cfg(feature = "experimental-encrypted-state-events")]
2497 fn verify_packed_state_key(
2498 &self,
2499 original: &EncryptedEvent,
2500 decrypted: &Raw<AnyTimelineEvent>,
2501 ) -> MegolmResult<()> {
2502 use serde::Deserialize;
2503
2504 #[derive(Deserialize)]
2506 struct PayloadDeserializationHelper {
2507 state_key: Option<String>,
2508 #[serde(rename = "type")]
2509 event_type: String,
2510 }
2511
2512 let PayloadDeserializationHelper {
2514 state_key: inner_state_key,
2515 event_type: inner_event_type,
2516 } = decrypted
2517 .deserialize_as_unchecked()
2518 .map_err(|_| MegolmError::StateKeyVerificationFailed)?;
2519
2520 let (raw_state_key, inner_state_key) = match (&original.state_key, &inner_state_key) {
2522 (Some(raw_state_key), Some(inner_state_key)) => (raw_state_key, inner_state_key),
2523 (None, None) => return Ok(()),
2524 _ => return Err(MegolmError::StateKeyVerificationFailed),
2525 };
2526
2527 let (outer_event_type, outer_state_key) =
2529 raw_state_key.split_once(":").ok_or(MegolmError::StateKeyVerificationFailed)?;
2530
2531 if outer_event_type != inner_event_type {
2533 return Err(MegolmError::StateKeyVerificationFailed);
2534 }
2535
2536 if outer_state_key != inner_state_key {
2538 return Err(MegolmError::StateKeyVerificationFailed);
2539 }
2540 Ok(())
2541 }
2542
2543 async fn decrypt_unsigned_events(
2553 &self,
2554 main_event: &mut JsonObject,
2555 room_id: &RoomId,
2556 decryption_settings: &DecryptionSettings,
2557 ) -> Option<BTreeMap<UnsignedEventLocation, UnsignedDecryptionResult>> {
2558 let unsigned = main_event.get_mut("unsigned")?.as_object_mut()?;
2559 let mut unsigned_encryption_info: Option<
2560 BTreeMap<UnsignedEventLocation, UnsignedDecryptionResult>,
2561 > = None;
2562
2563 let location = UnsignedEventLocation::RelationsReplace;
2565 let replace = location.find_mut(unsigned);
2566 if let Some(decryption_result) =
2567 self.decrypt_unsigned_event(replace, room_id, decryption_settings).await
2568 {
2569 unsigned_encryption_info
2570 .get_or_insert_with(Default::default)
2571 .insert(location, decryption_result);
2572 }
2573
2574 let location = UnsignedEventLocation::RelationsThreadLatestEvent;
2577 let thread_latest_event = location.find_mut(unsigned);
2578 if let Some(decryption_result) =
2579 self.decrypt_unsigned_event(thread_latest_event, room_id, decryption_settings).await
2580 {
2581 unsigned_encryption_info
2582 .get_or_insert_with(Default::default)
2583 .insert(location, decryption_result);
2584 }
2585
2586 unsigned_encryption_info
2587 }
2588
2589 fn decrypt_unsigned_event<'a>(
2597 &'a self,
2598 event: Option<&'a mut Value>,
2599 room_id: &'a RoomId,
2600 decryption_settings: &'a DecryptionSettings,
2601 ) -> BoxFuture<'a, Option<UnsignedDecryptionResult>> {
2602 Box::pin(async move {
2603 let event = event?;
2604
2605 let is_encrypted = event
2606 .get("type")
2607 .and_then(|type_| type_.as_str())
2608 .is_some_and(|s| s == "m.room.encrypted");
2609 if !is_encrypted {
2610 return None;
2611 }
2612
2613 let raw_event = serde_json::from_value(event.clone()).ok()?;
2614 match self
2615 .decrypt_room_event_inner(&raw_event, room_id, false, decryption_settings)
2616 .await
2617 {
2618 Ok(decrypted_event) => {
2619 *event = serde_json::to_value(decrypted_event.event).ok()?;
2621 Some(UnsignedDecryptionResult::Decrypted(decrypted_event.encryption_info))
2622 }
2623 Err(err) => {
2624 let utd_info = megolm_error_to_utd_info(&raw_event, err).ok()?;
2629 Some(UnsignedDecryptionResult::UnableToDecrypt(utd_info))
2630 }
2631 }
2632 })
2633 }
2634
2635 pub async fn is_room_key_available(
2642 &self,
2643 event: &Raw<EncryptedEvent>,
2644 room_id: &RoomId,
2645 ) -> Result<bool, CryptoStoreError> {
2646 let event = event.deserialize()?;
2647
2648 let (session_id, message_index) = match &event.content.scheme {
2649 RoomEventEncryptionScheme::MegolmV1AesSha2(c) => {
2650 (&c.session_id, c.ciphertext.message_index())
2651 }
2652 #[cfg(feature = "experimental-algorithms")]
2653 RoomEventEncryptionScheme::MegolmV2AesSha2(c) => {
2654 (&c.session_id, c.ciphertext.message_index())
2655 }
2656 RoomEventEncryptionScheme::Unknown(_) => {
2657 return Ok(false);
2659 }
2660 };
2661
2662 Ok(self
2665 .store()
2666 .get_inbound_group_session(room_id, session_id)
2667 .await?
2668 .filter(|s| s.first_known_index() <= message_index)
2669 .is_some())
2670 }
2671
2672 #[instrument(skip(self, event), fields(event_id, sender, session_id))]
2685 pub async fn get_room_event_encryption_info(
2686 &self,
2687 event: &Raw<EncryptedEvent>,
2688 room_id: &RoomId,
2689 ) -> MegolmResult<Arc<EncryptionInfo>> {
2690 let event = event.deserialize()?;
2691
2692 let content: SupportedEventEncryptionSchemes<'_> = match &event.content.scheme {
2693 RoomEventEncryptionScheme::MegolmV1AesSha2(c) => c.into(),
2694 #[cfg(feature = "experimental-algorithms")]
2695 RoomEventEncryptionScheme::MegolmV2AesSha2(c) => c.into(),
2696 RoomEventEncryptionScheme::Unknown(_) => {
2697 return Err(EventError::UnsupportedAlgorithm.into());
2698 }
2699 };
2700
2701 Span::current()
2702 .record("sender", debug(&event.sender))
2703 .record("event_id", debug(&event.event_id))
2704 .record("session_id", content.session_id());
2705
2706 self.get_session_encryption_info(room_id, content.session_id(), &event.sender).await
2707 }
2708
2709 pub async fn get_session_encryption_info(
2724 &self,
2725 room_id: &RoomId,
2726 session_id: &str,
2727 sender: &UserId,
2728 ) -> MegolmResult<Arc<EncryptionInfo>> {
2729 let session = self.get_inbound_group_session_or_error(room_id, session_id).await?;
2730 self.get_encryption_info(&session, sender).await
2731 }
2732
2733 pub async fn update_tracked_users(
2751 &self,
2752 users: impl IntoIterator<Item = &UserId>,
2753 ) -> StoreResult<()> {
2754 self.inner.identity_manager.update_tracked_users(users).await
2755 }
2756
2757 pub async fn mark_all_tracked_users_as_dirty(&self) -> StoreResult<()> {
2762 self.inner
2763 .identity_manager
2764 .mark_all_tracked_users_as_dirty(self.inner.store.cache().await?)
2765 .await
2766 }
2767
2768 async fn wait_if_user_pending(
2769 &self,
2770 user_id: &UserId,
2771 timeout: Option<Duration>,
2772 ) -> StoreResult<()> {
2773 if let Some(timeout) = timeout {
2774 let cache = self.store().cache().await?;
2775 self.inner
2776 .identity_manager
2777 .key_query_manager
2778 .wait_if_user_key_query_pending(cache, timeout, user_id)
2779 .await?;
2780 }
2781 Ok(())
2782 }
2783
2784 #[instrument(skip(self))]
2814 pub async fn get_device(
2815 &self,
2816 user_id: &UserId,
2817 device_id: &DeviceId,
2818 timeout: Option<Duration>,
2819 ) -> StoreResult<Option<Device>> {
2820 self.wait_if_user_pending(user_id, timeout).await?;
2821 self.store().get_device(user_id, device_id).await
2822 }
2823
2824 #[instrument(skip(self))]
2838 pub async fn get_identity(
2839 &self,
2840 user_id: &UserId,
2841 timeout: Option<Duration>,
2842 ) -> StoreResult<Option<UserIdentity>> {
2843 self.wait_if_user_pending(user_id, timeout).await?;
2844 self.store().get_identity(user_id).await
2845 }
2846
2847 #[instrument(skip(self))]
2874 pub async fn get_user_devices(
2875 &self,
2876 user_id: &UserId,
2877 timeout: Option<Duration>,
2878 ) -> StoreResult<UserDevices> {
2879 self.wait_if_user_pending(user_id, timeout).await?;
2880 self.store().get_user_devices(user_id).await
2881 }
2882
2883 pub async fn cross_signing_status(&self) -> CrossSigningStatus {
2888 self.inner.user_identity.lock().await.status().await
2889 }
2890
2891 pub async fn export_cross_signing_keys(&self) -> StoreResult<Option<CrossSigningKeyExport>> {
2899 let master_key = self.store().export_secret(&SecretName::CrossSigningMasterKey).await?;
2900 let self_signing_key =
2901 self.store().export_secret(&SecretName::CrossSigningSelfSigningKey).await?;
2902 let user_signing_key =
2903 self.store().export_secret(&SecretName::CrossSigningUserSigningKey).await?;
2904
2905 Ok(if master_key.is_none() && self_signing_key.is_none() && user_signing_key.is_none() {
2906 None
2907 } else {
2908 Some(CrossSigningKeyExport { master_key, self_signing_key, user_signing_key })
2909 })
2910 }
2911
2912 pub async fn import_cross_signing_keys(
2917 &self,
2918 export: CrossSigningKeyExport,
2919 ) -> Result<CrossSigningStatus, SecretImportError> {
2920 self.store().import_cross_signing_keys(export).await
2921 }
2922
2923 async fn sign_with_master_key(
2924 &self,
2925 message: &str,
2926 ) -> Result<(OwnedDeviceKeyId, Ed25519Signature), SignatureError> {
2927 let identity = &*self.inner.user_identity.lock().await;
2928 let key_id = identity.master_key_id().await.ok_or(SignatureError::MissingSigningKey)?;
2929
2930 let signature = identity.sign(message).await?;
2931
2932 Ok((key_id, signature))
2933 }
2934
2935 pub async fn sign(&self, message: &str) -> Result<Signatures, CryptoStoreError> {
2941 let mut signatures = Signatures::new();
2942
2943 {
2944 let cache = self.inner.store.cache().await?;
2945 let account = cache.account().await?;
2946 let key_id = account.signing_key_id();
2947 let signature = account.sign(message);
2948 signatures.add_signature(self.user_id().to_owned(), key_id, signature);
2949 }
2950
2951 match self.sign_with_master_key(message).await {
2952 Ok((key_id, signature)) => {
2953 signatures.add_signature(self.user_id().to_owned(), key_id, signature);
2954 }
2955 Err(e) => {
2956 warn!(error = ?e, "Couldn't sign the message using the cross signing master key")
2957 }
2958 }
2959
2960 Ok(signatures)
2961 }
2962
2963 pub fn backup_machine(&self) -> &BackupMachine {
2968 &self.inner.backup_machine
2969 }
2970
2971 pub async fn initialize_crypto_store_generation(
2975 &self,
2976 generation: &Mutex<Option<u64>>,
2977 ) -> StoreResult<()> {
2978 let mut gen_guard = generation.lock().await;
2981
2982 let prev_generation =
2983 self.inner.store.get_custom_value(Self::CURRENT_GENERATION_STORE_KEY).await?;
2984
2985 let generation = match prev_generation {
2986 Some(val) => {
2987 u64::from_le_bytes(val.try_into().map_err(|_| {
2990 CryptoStoreError::InvalidLockGeneration("invalid format".to_owned())
2991 })?)
2992 .wrapping_add(1)
2993 }
2994 None => 0,
2995 };
2996
2997 tracing::debug!("Initialising crypto store generation at {generation}");
2998
2999 self.inner
3000 .store
3001 .set_custom_value(Self::CURRENT_GENERATION_STORE_KEY, generation.to_le_bytes().to_vec())
3002 .await?;
3003
3004 *gen_guard = Some(generation);
3005
3006 Ok(())
3007 }
3008
3009 pub async fn maintain_crypto_store_generation(
3034 &'_ self,
3035 generation: &Mutex<Option<u64>>,
3036 ) -> StoreResult<(bool, u64)> {
3037 let mut gen_guard = generation.lock().await;
3038
3039 let actual_gen = self
3045 .inner
3046 .store
3047 .get_custom_value(Self::CURRENT_GENERATION_STORE_KEY)
3048 .await?
3049 .ok_or_else(|| {
3050 CryptoStoreError::InvalidLockGeneration("counter missing in store".to_owned())
3051 })?;
3052
3053 let actual_gen =
3054 u64::from_le_bytes(actual_gen.try_into().map_err(|_| {
3055 CryptoStoreError::InvalidLockGeneration("invalid format".to_owned())
3056 })?);
3057
3058 let new_gen = match gen_guard.as_ref() {
3059 Some(expected_gen) => {
3060 if actual_gen == *expected_gen {
3061 return Ok((false, actual_gen));
3062 }
3063 actual_gen.max(*expected_gen).wrapping_add(1)
3065 }
3066 None => {
3067 actual_gen.wrapping_add(1)
3070 }
3071 };
3072
3073 tracing::debug!(
3074 "Crypto store generation mismatch: previously known was {:?}, actual is {:?}, next is {}",
3075 *gen_guard,
3076 actual_gen,
3077 new_gen
3078 );
3079
3080 *gen_guard = Some(new_gen);
3082
3083 self.inner
3085 .store
3086 .set_custom_value(Self::CURRENT_GENERATION_STORE_KEY, new_gen.to_le_bytes().to_vec())
3087 .await?;
3088
3089 Ok((true, new_gen))
3090 }
3091
3092 pub fn dehydrated_devices(&self) -> DehydratedDevices {
3094 DehydratedDevices { inner: self.to_owned() }
3095 }
3096
3097 pub async fn room_settings(&self, room_id: &RoomId) -> StoreResult<Option<RoomSettings>> {
3102 self.inner.store.get_room_settings(room_id).await
3105 }
3106
3107 pub async fn set_room_settings(
3118 &self,
3119 room_id: &RoomId,
3120 new_settings: &RoomSettings,
3121 ) -> Result<(), SetRoomSettingsError> {
3122 let store = &self.inner.store;
3123
3124 let _store_transaction = store.transaction().await;
3129
3130 let old_settings = store.get_room_settings(room_id).await?;
3131
3132 if let Some(old_settings) = old_settings {
3145 if old_settings != *new_settings {
3146 return Err(SetRoomSettingsError::EncryptionDowngrade);
3147 } else {
3148 return Ok(());
3150 }
3151 }
3152
3153 match new_settings.algorithm {
3155 EventEncryptionAlgorithm::MegolmV1AesSha2 => (),
3156
3157 #[cfg(feature = "experimental-algorithms")]
3158 EventEncryptionAlgorithm::MegolmV2AesSha2 => (),
3159
3160 _ => {
3161 warn!(
3162 ?room_id,
3163 "Rejecting invalid encryption algorithm {}", new_settings.algorithm
3164 );
3165 return Err(SetRoomSettingsError::InvalidSettings);
3166 }
3167 }
3168
3169 store
3171 .save_changes(Changes {
3172 room_settings: HashMap::from([(room_id.to_owned(), new_settings.clone())]),
3173 ..Default::default()
3174 })
3175 .await?;
3176
3177 Ok(())
3178 }
3179
3180 #[cfg(any(feature = "testing", test))]
3184 pub fn same_as(&self, other: &OlmMachine) -> bool {
3185 Arc::ptr_eq(&self.inner, &other.inner)
3186 }
3187
3188 #[cfg(any(feature = "testing", test))]
3190 pub async fn uploaded_key_count(&self) -> Result<u64, CryptoStoreError> {
3191 let cache = self.inner.store.cache().await?;
3192 let account = cache.account().await?;
3193 Ok(account.uploaded_key_count())
3194 }
3195
3196 #[cfg(test)]
3198 pub(crate) fn identity_manager(&self) -> &IdentityManager {
3199 &self.inner.identity_manager
3200 }
3201
3202 #[cfg(test)]
3204 pub(crate) fn key_for_has_migrated_verification_latch() -> &'static str {
3205 Self::HAS_MIGRATED_VERIFICATION_LATCH
3206 }
3207}
3208
3209fn sender_data_to_verification_state(
3210 sender_data: SenderData,
3211 session_has_been_imported: bool,
3212) -> (VerificationState, Option<OwnedDeviceId>) {
3213 match sender_data {
3214 SenderData::UnknownDevice { owner_check_failed: false, .. } => {
3215 let device_link_problem = if session_has_been_imported {
3216 DeviceLinkProblem::InsecureSource
3217 } else {
3218 DeviceLinkProblem::MissingDevice
3219 };
3220
3221 (VerificationState::Unverified(VerificationLevel::None(device_link_problem)), None)
3222 }
3223 SenderData::UnknownDevice { owner_check_failed: true, .. } => (
3224 VerificationState::Unverified(VerificationLevel::None(
3225 DeviceLinkProblem::InsecureSource,
3226 )),
3227 None,
3228 ),
3229 SenderData::DeviceInfo { device_keys, .. } => (
3230 VerificationState::Unverified(VerificationLevel::UnsignedDevice),
3231 Some(device_keys.device_id),
3232 ),
3233 SenderData::VerificationViolation(KnownSenderData { device_id, .. }) => {
3234 (VerificationState::Unverified(VerificationLevel::VerificationViolation), device_id)
3235 }
3236 SenderData::SenderUnverified(KnownSenderData { device_id, .. }) => {
3237 (VerificationState::Unverified(VerificationLevel::UnverifiedIdentity), device_id)
3238 }
3239 SenderData::SenderVerified(KnownSenderData { device_id, .. }) => {
3240 (VerificationState::Verified, device_id)
3241 }
3242 }
3243}
3244
3245#[derive(Debug, Clone)]
3248pub struct CrossSigningBootstrapRequests {
3249 pub upload_keys_req: Option<OutgoingRequest>,
3256
3257 pub upload_signing_keys_req: UploadSigningKeysRequest,
3261
3262 pub upload_signatures_req: UploadSignaturesRequest,
3267}
3268
3269#[derive(Debug, thiserror::Error)]
3275pub enum BootstrapCrossSigningError {
3276 #[error(transparent)]
3278 CryptoStore(#[from] CryptoStoreError),
3279
3280 #[error(transparent)]
3282 Signature(#[from] SignatureError),
3283}
3284
3285#[derive(Debug)]
3288pub struct EncryptionSyncChanges<'a> {
3289 pub to_device_events: Vec<Raw<AnyToDeviceEvent>>,
3291 pub changed_devices: &'a DeviceLists,
3294 pub one_time_keys_counts: &'a BTreeMap<OneTimeKeyAlgorithm, UInt>,
3296 pub unused_fallback_keys: Option<&'a [OneTimeKeyAlgorithm]>,
3298 pub next_batch_token: Option<String>,
3300}
3301
3302fn megolm_error_to_utd_info(
3310 raw_event: &Raw<EncryptedEvent>,
3311 error: MegolmError,
3312) -> Result<UnableToDecryptInfo, CryptoStoreError> {
3313 use MegolmError::*;
3314 let reason = match error {
3315 EventError(_) => UnableToDecryptReason::MalformedEncryptedEvent,
3316 Decode(_) => UnableToDecryptReason::MalformedEncryptedEvent,
3317 MissingRoomKey(maybe_withheld) => {
3318 UnableToDecryptReason::MissingMegolmSession { withheld_code: maybe_withheld }
3319 }
3320 Decryption(DecryptionError::UnknownMessageIndex(_, _)) => {
3321 UnableToDecryptReason::UnknownMegolmMessageIndex
3322 }
3323 Decryption(_) => UnableToDecryptReason::MegolmDecryptionFailure,
3324 JsonError(_) => UnableToDecryptReason::PayloadDeserializationFailure,
3325 MismatchedIdentityKeys(_) => UnableToDecryptReason::MismatchedIdentityKeys,
3326 SenderIdentityNotTrusted(level) => UnableToDecryptReason::SenderIdentityNotTrusted(level),
3327 #[cfg(feature = "experimental-encrypted-state-events")]
3328 StateKeyVerificationFailed => UnableToDecryptReason::StateKeyVerificationFailed,
3329
3330 Store(error) => Err(error)?,
3333 };
3334
3335 let session_id = raw_event.deserialize().ok().and_then(|ev| match ev.content.scheme {
3336 RoomEventEncryptionScheme::MegolmV1AesSha2(s) => Some(s.session_id),
3337 #[cfg(feature = "experimental-algorithms")]
3338 RoomEventEncryptionScheme::MegolmV2AesSha2(s) => Some(s.session_id),
3339 RoomEventEncryptionScheme::Unknown(_) => None,
3340 });
3341
3342 Ok(UnableToDecryptInfo { session_id, reason })
3343}
3344
3345#[derive(Debug, thiserror::Error)]
3355pub(crate) enum DecryptToDeviceError {
3356 #[error("An Olm error occurred meaning we failed to decrypt the event")]
3357 OlmError(#[from] OlmError),
3358
3359 #[error("The event was sent from a dehydrated device")]
3360 FromDehydratedDevice,
3361}
3362
3363impl From<CryptoStoreError> for DecryptToDeviceError {
3364 fn from(value: CryptoStoreError) -> Self {
3365 Self::OlmError(value.into())
3366 }
3367}
3368
3369#[cfg(test)]
3370impl From<DecryptToDeviceError> for OlmError {
3371 fn from(value: DecryptToDeviceError) -> Self {
3374 match value {
3375 DecryptToDeviceError::OlmError(olm_error) => olm_error,
3376 DecryptToDeviceError::FromDehydratedDevice => {
3377 panic!("Expected an OlmError but found FromDehydratedDevice")
3378 }
3379 }
3380 }
3381}
3382
3383#[cfg(test)]
3384pub(crate) mod test_helpers;
3385
3386#[cfg(test)]
3387pub(crate) mod tests;