1use crate::crypto::account::OlmAccountState;
107use crate::crypto::cross_signing::CrossSigningState;
108use crate::crypto::device_tracker::{DeviceTracker, StoredDevice};
109use crate::crypto::group_sessions::{GroupDecryptError, GroupSessionManager, RoomEventPlaintext};
110use crate::crypto::olm_sessions::{DecryptedToDevice, OlmDecryptError, OlmSessionManager};
111use crate::error::MessengerError;
112use crate::ids::{DeviceId, EventId, RequestId, RoomId, TxnId, UserId};
113use crate::outgoing_queue::{Jitter, Lane, OutgoingQueue, PendingRequest, ResponseOutcome};
114use crate::persist::{FlushBatch, SealedRecord};
115use crate::room::state::RoomState;
116use crate::room::timeline::{
117 interpret_content, ForwardRefusal, ItemContent, SendState, Timeline, KEY_FORWARDED, KEY_FORWARDED_FROM,
118};
119use crate::room::RoomKind;
120use crate::store::sealed::SealedRecordCodec;
121use crate::store::{CryptoStore, RecordCodec, StateStore, Store};
122use crate::wire::events::{
123 DirectContent, ForwardedRoomKeyContent, InReplyTo, MegolmEncryptedContent, Membership, RawEvent, ReceiptContent,
124 RelatesTo, RoomEncryptedContent, RoomKeyContent, RoomKeyRequestAction, RoomKeyRequestBody, RoomKeyRequestContent,
125 RoomKeyWithheldContent, StrippedStateEvent, TagContent, TagInfo, TextLikeMessageContent, ToDeviceEvent,
126 TypingContent, Unsigned,
127};
128use crate::wire::sync::{parse_sync_response, InvitedRoom, JoinedRoom, LeftRoom};
129use crate::wire::{percent_decode_segment, HttpResponseDescriptor, OutgoingRequest, OutgoingRequestKind};
130use serde::{Deserialize, Serialize};
131use std::collections::{BTreeMap, BTreeSet, VecDeque};
132use zeroize::Zeroizing;
133
134const SYNC_TIMEOUT_MS: u64 = 30_000;
136const LOAD_OLDER_PAGE_SIZE: u32 = 50;
139const TYPING_TIMEOUT_MS: u64 = 30_000;
142const UNKNOWN_SENDER_RETRY_CAPACITY: usize = 256;
147const TYPING_DEBOUNCE_MS: i64 = 4_000;
151const ROOM_KEY_SHARE_EVENT_TYPE: &str = "m.room.encrypted";
159
160#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
164pub struct Counters {
165 pub next_request_id: u64,
167 pub next_txn_id: u64,
172}
173
174fn fresh_txn_seed() -> u64 {
188 std::time::SystemTime::now()
189 .duration_since(std::time::UNIX_EPOCH)
190 .map(|d| d.as_micros() as u64)
191 .unwrap_or(0)
192}
193
194fn restore_counter(persisted_next: u64, max_used_seed: Option<u64>) -> u64 {
195 match max_used_seed {
196 Some(seed) => persisted_next.max(seed.saturating_add(1)),
197 None => persisted_next,
198 }
199}
200
201#[derive(Clone, Debug, PartialEq, Eq)]
205pub struct CoreConfig {
206 pub user_id: UserId,
208 pub device_id: DeviceId,
210 pub server_name: String,
218}
219
220#[derive(Default)]
229pub struct CoreSecrets {
230 pub store_seal_key: Option<Zeroizing<[u8; 32]>>,
238 pub backup_key: Option<Zeroizing<[u8; 32]>>,
242}
243
244#[derive(Clone, Copy, Debug, PartialEq, Eq)]
247pub enum MessageKind {
248 Text,
250 Notice,
252 Emote,
254}
255
256impl MessageKind {
257 fn msgtype(self) -> &'static str {
258 match self {
259 MessageKind::Text => "m.text",
260 MessageKind::Notice => "m.notice",
261 MessageKind::Emote => "m.emote",
262 }
263 }
264}
265
266#[derive(Clone, Debug, PartialEq, Eq)]
277pub struct OutgoingMessage {
278 pub kind: MessageKind,
280 pub body: String,
283 pub reply_to: Option<EventId>,
285 pub edit_of: Option<EventId>,
288}
289
290#[derive(Clone, Debug, PartialEq, Eq)]
296pub enum CreateRoomKind {
297 Dm {
301 peer: UserId,
303 },
304 Group {
306 name: String,
308 invite: Vec<UserId>,
310 members_can_invite: bool,
314 },
315 Channel {
318 name: String,
320 topic: Option<String>,
322 },
323}
324
325#[derive(Clone, Debug, PartialEq)]
331pub enum MessengerCommand {
332 SendMessage {
338 room_id: RoomId,
340 message: OutgoingMessage,
342 txn_id: Option<TxnId>,
345 },
346 Forward {
372 from_room: RoomId,
374 event_id: EventId,
376 to_room: RoomId,
378 txn_id: Option<TxnId>,
380 },
381 React {
385 room_id: RoomId,
387 target: EventId,
389 key: String,
391 },
392 Redact {
394 room_id: RoomId,
396 target: EventId,
398 reason: Option<String>,
400 },
401 RetrySend {
407 room_id: RoomId,
409 txn_id: TxnId,
411 },
412 CreateRoom {
414 kind: CreateRoomKind,
416 },
417 JoinRoom {
420 room_id: RoomId,
422 },
423 LeaveRoom {
425 room_id: RoomId,
427 },
428 Invite {
430 room_id: RoomId,
432 user_id: UserId,
434 },
435 Kick {
440 room_id: RoomId,
442 user_id: UserId,
444 reason: Option<String>,
446 },
447 SetTag {
449 room_id: RoomId,
451 tag: String,
453 order: Option<f64>,
455 },
456 RemoveTag {
458 room_id: RoomId,
460 tag: String,
462 },
463 MarkRead {
466 room_id: RoomId,
468 event_id: EventId,
470 },
471 SetTyping {
474 room_id: RoomId,
476 typing: bool,
478 },
479 RetryDecryption {
484 room_id: RoomId,
486 event_id: EventId,
488 },
489 LoadOlder {
494 room_id: RoomId,
496 },
497 SetAccountData {
503 event_type: String,
505 content: serde_json::Value,
507 },
508 SetRoomAccountData {
517 room_id: RoomId,
519 event_type: String,
521 content: serde_json::Value,
523 },
524 SearchPublicRooms {
529 term: String,
531 },
532 SearchUsers {
537 term: String,
539 },
540}
541
542#[derive(Clone, Debug, PartialEq, Eq)]
547pub enum MessengerEvent {
548 RoomsChanged,
551 TimelineChanged {
553 room_id: RoomId,
555 },
556 TypingChanged {
558 room_id: RoomId,
560 },
561 ReceiptsChanged {
563 room_id: RoomId,
565 },
566 UnreadChanged {
568 room_id: RoomId,
570 },
571 DeviceKeyChanged {
575 user_id: UserId,
577 device_id: DeviceId,
579 },
580 PublicRoomsChanged,
584 UserSearchChanged,
588}
589
590#[derive(Clone)]
595struct PendingDecryptItem {
596 session_id: String,
597 sender: UserId,
598 origin_server_ts: i64,
599 content: MegolmEncryptedContent,
600}
601
602#[derive(Deserialize)]
608struct RoomMessagesResponseBody {
609 #[serde(default)]
610 chunk: Vec<RawEvent>,
611 #[serde(default)]
612 end: Option<String>,
613}
614
615#[derive(Deserialize)]
617struct SendResponseBody {
618 event_id: EventId,
619}
620
621#[derive(Deserialize)]
623struct CreateRoomResponseBody {
624 room_id: RoomId,
625}
626
627#[derive(Clone, Debug, PartialEq)]
630pub struct PublicRoomsResultEntry {
631 pub room_id: RoomId,
633 pub name: Option<String>,
635 pub topic: Option<String>,
637 pub num_joined_members: u64,
639}
640
641#[derive(Deserialize)]
643struct PublicRoomsResponseBody {
644 #[serde(default)]
645 chunk: Vec<PublicRoomsChunkEntry>,
646}
647
648#[derive(Deserialize)]
649struct PublicRoomsChunkEntry {
650 room_id: RoomId,
651 #[serde(default)]
652 name: Option<String>,
653 #[serde(default)]
654 topic: Option<String>,
655 #[serde(default)]
656 num_joined_members: u64,
657}
658
659#[derive(Clone, Debug, PartialEq)]
662pub struct UserDirectoryResultEntry {
663 pub user_id: UserId,
665 pub display_name: Option<String>,
667}
668
669#[derive(Deserialize)]
672struct UserDirectorySearchResponseBody {
673 #[serde(default)]
674 results: Vec<UserDirectorySearchResultEntry>,
675}
676
677#[derive(Deserialize)]
678struct UserDirectorySearchResultEntry {
679 user_id: UserId,
680 #[serde(default)]
681 display_name: Option<String>,
682}
683
684#[derive(Clone)]
689enum SendPayload {
690 Message(OutgoingMessage),
693 Event {
697 event_type: String,
699 content: serde_json::Value,
701 },
702 Reaction {
704 target: EventId,
706 key: String,
708 },
709 Redaction {
711 target: EventId,
713 reason: Option<String>,
715 },
716}
717
718struct RoomEventWire {
722 event_type: String,
723 content: serde_json::Value,
724 relates_to: Option<RelatesTo>,
727}
728
729impl SendPayload {
730 fn room_event(&self) -> Option<RoomEventWire> {
733 match self {
734 SendPayload::Message(message) => Some(RoomEventWire {
735 event_type: "m.room.message".to_string(),
736 content: message_wire_content(message),
737 relates_to: message_relates_to(message),
738 }),
739 SendPayload::Event { event_type, content } => {
740 Some(RoomEventWire { event_type: event_type.clone(), content: content.clone(), relates_to: None })
741 }
742 SendPayload::Reaction { .. } | SendPayload::Redaction { .. } => None,
743 }
744 }
745
746 fn is_room_event(&self) -> bool {
747 matches!(self, SendPayload::Message(_) | SendPayload::Event { .. })
748 }
749}
750
751#[derive(Clone)]
757enum WireSend {
758 Event {
760 event_type: String,
763 content: serde_json::Value,
765 },
766 Redact {
768 target: EventId,
770 reason: Option<String>,
772 },
773}
774
775#[derive(Clone)]
780enum SendPhase {
781 Queued,
785 AwaitingKeysQuery,
791 AwaitingKeysClaim {
794 devices: Vec<StoredDevice>,
798 },
799 AwaitingKeyShare {
804 pending: BTreeSet<RequestId>,
806 },
807 AwaitingSend,
809 Failed {
814 errcode: String,
816 },
817}
818
819#[derive(Clone)]
825struct PendingSend {
826 room_id: RoomId,
829 payload: SendPayload,
831 phase: SendPhase,
833 wire: Option<WireSend>,
835}
836
837type AccountDataKey = (Option<RoomId>, String);
840
841struct AccountDataGuard {
853 pending: u32,
855 server_value: Option<serde_json::Value>,
859}
860
861fn account_data_key_from_path(path: &str) -> Option<AccountDataKey> {
864 let segments: Vec<&str> = path.split('/').collect();
865 let n = segments.len();
866 if n < 4 || segments[n - 2] != "account_data" {
867 return None;
868 }
869 let event_type = percent_decode_segment(segments[n - 1])?;
870 if segments[n - 4] == "user" {
871 return Some((None, event_type));
872 }
873 if n >= 6 && segments[n - 4] == "rooms" && segments[n - 6] == "user" {
874 let room_id = RoomId::parse(percent_decode_segment(segments[n - 3])?).ok()?;
875 return Some((Some(room_id), event_type));
876 }
877 None
878}
879
880pub struct MessengerCore<C: RecordCodec> {
882 config: CoreConfig,
883 store: Store<C>,
884 account: OlmAccountState,
885 xsign: Option<CrossSigningState>,
887 xsign_inflight: bool,
888 xsign_attempts: u8,
889 outgoing: OutgoingQueue,
890 counters: Counters,
891 jitter: Box<dyn Jitter>,
892
893 rooms: BTreeMap<RoomId, RoomState>,
894 timelines: BTreeMap<RoomId, Timeline>,
895 direct_account_data: Option<DirectContent>,
896 global_account_data: BTreeMap<String, serde_json::Value>,
897 room_account_data: BTreeMap<RoomId, BTreeMap<String, serde_json::Value>>,
898 typing: BTreeMap<RoomId, Vec<UserId>>,
899 receipts: BTreeMap<RoomId, ReceiptContent>,
900
901 public_rooms_result: Vec<PublicRoomsResultEntry>,
904 latest_public_rooms_request: Option<RequestId>,
909 user_search_result: Vec<UserDirectoryResultEntry>,
912 latest_user_search_request: Option<RequestId>,
915
916 account_data_guard: BTreeMap<AccountDataKey, AccountDataGuard>,
920 account_data_write_keys: BTreeMap<RequestId, AccountDataKey>,
922
923 pending_decrypt: BTreeMap<RoomId, BTreeMap<EventId, PendingDecryptItem>>,
924 key_requests_sent: BTreeSet<(RoomId, String)>,
928 history_claims: BTreeMap<RequestId, Vec<StoredDevice>>,
931 history_claim_tried: BTreeSet<(UserId, DeviceId)>,
933 withheld_sessions: BTreeSet<(RoomId, String)>,
935 key_requests_answered: BTreeSet<(UserId, DeviceId, String)>,
937 unknown_sender_retry: VecDeque<ToDeviceEvent>,
938
939 typing_debounce: BTreeMap<RoomId, i64>,
944 pending_sends: BTreeMap<RoomId, VecDeque<TxnId>>,
950 send_state: BTreeMap<TxnId, PendingSend>,
955 send_request_owner: BTreeMap<RequestId, TxnId>,
959 pending_create_room: BTreeMap<RequestId, CreateRoomKind>,
963
964 pending_kinds: BTreeMap<RequestId, OutgoingRequestKind>,
970 pending_room_messages: BTreeMap<RequestId, RoomId>,
973 sync_request_id: Option<RequestId>,
976 sync_caught_up: bool,
981 pending_flush_batches: VecDeque<FlushBatch>,
988
989 events: Vec<MessengerEvent>,
990 change_counter: u64,
991 ingest_error: Option<String>,
996}
997
998impl<C: RecordCodec> MessengerCore<C> {
999 pub fn open(
1012 records: impl IntoIterator<Item = SealedRecord>,
1013 codec: C,
1014 config: CoreConfig,
1015 secrets: CoreSecrets,
1016 _now_ms: i64,
1017 jitter: Box<dyn Jitter>,
1018 ) -> Result<Self, MessengerError> {
1019 let _ = secrets;
1023 let mut store = Store::load(records, codec, config.device_id.clone())?;
1024 let mut outgoing = OutgoingQueue::load(&store)?;
1025 let account = OlmAccountState::load_or_create(&mut store)?;
1026
1027 let mut rooms = BTreeMap::new();
1028 let mut timelines = BTreeMap::new();
1029 for room_id in store.room_ids()? {
1030 if let Some(bytes) = store.room_state(&room_id)? {
1031 let state: RoomState = serde_json::from_slice(bytes)
1032 .map_err(|source| MessengerError::Crypto(format!("decode room state for {room_id}: {source}")))?;
1033 rooms.insert(room_id.clone(), state);
1034 timelines.insert(room_id, Timeline::new());
1035 }
1036 }
1037
1038 let pending_requests = store.pending_requests()?;
1039 let max_pending_request_seed = pending_requests.iter().filter_map(|req| req.request.id.as_seed()).max();
1040 outgoing.discard_kind(&mut store, OutgoingRequestKind::Sync)?;
1047 let persisted_counters = store.counters()?.unwrap_or_default();
1048 let mut counters = Counters {
1049 next_request_id: restore_counter(persisted_counters.next_request_id, max_pending_request_seed),
1050 next_txn_id: restore_counter(persisted_counters.next_txn_id, None),
1051 };
1052 if counters != persisted_counters {
1053 store.save_counters(counters)?;
1059 }
1060 counters.next_txn_id = counters.next_txn_id.max(fresh_txn_seed());
1066
1067 let mut core = Self {
1068 config,
1069 store,
1070 account,
1071 outgoing,
1072 counters,
1073 jitter,
1074 rooms,
1075 timelines,
1076 direct_account_data: None,
1077 global_account_data: BTreeMap::new(),
1078 room_account_data: BTreeMap::new(),
1079 typing: BTreeMap::new(),
1080 receipts: BTreeMap::new(),
1081 public_rooms_result: Vec::new(),
1082 latest_public_rooms_request: None,
1083 user_search_result: Vec::new(),
1084 latest_user_search_request: None,
1085 account_data_guard: BTreeMap::new(),
1086 account_data_write_keys: BTreeMap::new(),
1087 pending_decrypt: BTreeMap::new(),
1088 key_requests_sent: BTreeSet::new(),
1089 history_claims: BTreeMap::new(),
1090 history_claim_tried: BTreeSet::new(),
1091 withheld_sessions: BTreeSet::new(),
1092 key_requests_answered: BTreeSet::new(),
1093 unknown_sender_retry: VecDeque::new(),
1094 typing_debounce: BTreeMap::new(),
1095 pending_sends: BTreeMap::new(),
1096 send_state: BTreeMap::new(),
1097 send_request_owner: BTreeMap::new(),
1098 pending_create_room: BTreeMap::new(),
1099 pending_kinds: BTreeMap::new(),
1100 pending_room_messages: BTreeMap::new(),
1101 sync_request_id: None,
1102 sync_caught_up: false,
1103 pending_flush_batches: VecDeque::new(),
1104 events: Vec::new(),
1105 change_counter: 0,
1106 ingest_error: None,
1107 xsign: None,
1108 xsign_inflight: false,
1109 xsign_attempts: 0,
1110 };
1111 core.rebuild_account_data_guard(&pending_requests);
1112 Ok(core)
1113 }
1114
1115 fn rebuild_account_data_guard(&mut self, pending_requests: &[PendingRequest]) {
1123 let mut writes: Vec<&PendingRequest> = pending_requests
1124 .iter()
1125 .filter(|record| matches!(record.request.kind, OutgoingRequestKind::AccountData | OutgoingRequestKind::RoomAccountData))
1126 .collect();
1127 writes.sort_by(|a, b| a.request.id.cmp(&b.request.id));
1128 for record in writes {
1129 let Some(key) = account_data_key_from_path(&record.request.path) else { continue };
1130 let Some(body) = record.request.body.clone() else { continue };
1131 let guard = self.account_data_guard.entry(key.clone()).or_insert(AccountDataGuard { pending: 0, server_value: None });
1132 guard.pending += 1;
1133 self.store_account_data_cache(&key, Some(body));
1134 self.account_data_write_keys.insert(record.request.id.clone(), key);
1135 }
1136 }
1137
1138 fn next_request_id(&mut self) -> Result<RequestId, MessengerError> {
1142 let seed = self.counters.next_request_id;
1143 self.counters.next_request_id = seed.saturating_add(1);
1144 self.store.save_counters(self.counters)?;
1145 Ok(RequestId::next(seed))
1146 }
1147
1148 fn next_txn_id(&mut self) -> Result<TxnId, MessengerError> {
1153 let seed = self.counters.next_txn_id;
1154 self.counters.next_txn_id = seed.saturating_add(1);
1155 self.store.save_counters(self.counters)?;
1156 Ok(TxnId::new(seed))
1157 }
1158
1159 fn enqueue_request(&mut self, request: OutgoingRequest, lane: Lane) -> Result<(), MessengerError> {
1174 let id = request.id.clone();
1175 let kind = request.kind;
1176 if let Some(batch) = self.outgoing.enqueue(&mut self.store, request, lane)? {
1177 self.pending_flush_batches.push_back(batch);
1178 }
1179 self.pending_kinds.insert(id, kind);
1180 Ok(())
1181 }
1182
1183 fn emit(&mut self, event: MessengerEvent) {
1184 self.change_counter = self.change_counter.wrapping_add(1);
1185 self.events.push(event);
1186 }
1187
1188 fn drain_events(&mut self) -> Vec<MessengerEvent> {
1189 std::mem::take(&mut self.events)
1190 }
1191
1192 fn start_new_send(
1208 &mut self,
1209 room_id: RoomId,
1210 payload: SendPayload,
1211 caller_txn_id: Option<TxnId>,
1212 now_ms: i64,
1213 ) -> Result<(), MessengerError> {
1214 let txn_id = match caller_txn_id {
1215 Some(txn_id) => txn_id,
1216 None => self.next_txn_id()?,
1217 };
1218
1219 if let Some(wire) = payload.room_event() {
1220 if self.typing_debounce.remove(&room_id).is_some() {
1221 if let Ok(request_id) = self.next_request_id() {
1222 let request = OutgoingRequest::typing(request_id, &room_id, &self.config.user_id, false, None);
1223 let _ = self.enqueue_request(request, Lane::Other);
1224 }
1225 }
1226 let content = match &payload {
1227 SendPayload::Message(message) => message_local_echo_content(message),
1228 _ => interpret_content(&wire.event_type, None, &wire.content),
1229 };
1230 let timeline = self.timelines.entry(room_id.clone()).or_default();
1231 timeline.push_local_echo(
1232 txn_id.clone(),
1233 self.config.user_id.clone(),
1234 now_ms,
1235 &wire.event_type,
1236 content,
1237 wire.content,
1238 );
1239 timeline.mark_send_state(&txn_id, SendState::Sending);
1240 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
1241 }
1242
1243 self.send_state.insert(txn_id.clone(), PendingSend { room_id: room_id.clone(), payload, phase: SendPhase::Queued, wire: None });
1244 let queue = self.pending_sends.entry(room_id).or_default();
1245 let was_empty = queue.is_empty();
1246 queue.push_back(txn_id.clone());
1247 if was_empty {
1248 self.start_send_pipeline(&txn_id, now_ms)?;
1249 }
1250 Ok(())
1251 }
1252
1253 fn retry_send(&mut self, room_id: &RoomId, txn_id: &TxnId, now_ms: i64) -> Result<(), MessengerError> {
1258 let Some(pending) = self.send_state.get(txn_id) else { return Ok(()) };
1259 if pending.room_id != *room_id || !matches!(pending.phase, SendPhase::Failed { .. }) {
1260 return Ok(());
1261 }
1262 self.set_send_phase(txn_id, SendPhase::Queued);
1263 let queue = self.pending_sends.entry(room_id.clone()).or_default();
1264 let was_empty = queue.is_empty();
1265 queue.push_back(txn_id.clone());
1266 if was_empty {
1267 self.start_send_pipeline(txn_id, now_ms)?;
1268 }
1269 Ok(())
1270 }
1271
1272 fn set_send_phase(&mut self, txn_id: &TxnId, phase: SendPhase) {
1273 if let Some(entry) = self.send_state.get_mut(txn_id) {
1274 entry.phase = phase;
1275 }
1276 }
1277
1278 fn start_send_pipeline(&mut self, txn_id: &TxnId, now_ms: i64) -> Result<(), MessengerError> {
1287 let Some(pending) = self.send_state.get(txn_id) else { return Ok(()) };
1288 let room_id = pending.room_id.clone();
1289
1290 if pending.wire.is_some() {
1291 return self.begin_awaiting_send(txn_id, &room_id);
1292 }
1293
1294 let is_room_event = pending.payload.is_room_event();
1295 let encrypted_room = self.rooms.get(&room_id).is_some_and(|room| room.encryption.is_some());
1296
1297 if is_room_event && encrypted_room {
1298 self.begin_key_pipeline(txn_id, &room_id, now_ms)
1299 } else {
1300 self.build_plaintext_wire_body(txn_id);
1301 self.begin_awaiting_send(txn_id, &room_id)
1302 }
1303 }
1304
1305 fn encrypted_room_member_ids(&self, room_id: &RoomId) -> BTreeSet<UserId> {
1310 let mut ids: BTreeSet<UserId> = self
1311 .rooms
1312 .get(room_id)
1313 .map(|state| {
1314 state
1315 .members
1316 .iter()
1317 .filter(|(_, member)| matches!(member.membership, Membership::Join | Membership::Invite))
1318 .map(|(user_id, _)| user_id.clone())
1319 .collect()
1320 })
1321 .unwrap_or_default();
1322 ids.insert(self.config.user_id.clone());
1323 ids
1324 }
1325
1326 fn mark_outdated_or_untracked(&mut self, member_ids: &BTreeSet<UserId>) -> Result<Vec<UserId>, MessengerError> {
1332 let tracked: BTreeMap<UserId, bool> = self.store.tracked_users()?.into_iter().collect();
1333 let mut outdated = Vec::new();
1334 let mut newly_tracked = Vec::new();
1335 for user_id in member_ids {
1336 match tracked.get(user_id) {
1337 Some(true) => outdated.push(user_id.clone()),
1338 Some(false) => {}
1339 None => newly_tracked.push(user_id.clone()),
1340 }
1341 }
1342 if !newly_tracked.is_empty() {
1343 DeviceTracker::on_device_lists(&mut self.store, &newly_tracked, &[])?;
1344 outdated.extend(newly_tracked);
1345 }
1346 Ok(outdated)
1347 }
1348
1349 fn known_devices_for(&self, member_ids: &BTreeSet<UserId>) -> Result<Vec<StoredDevice>, MessengerError> {
1352 let mut all = Vec::new();
1353 for user_id in member_ids {
1354 all.extend(DeviceTracker::devices_for_user(&self.store, user_id)?);
1355 }
1356 Ok(all)
1357 }
1358
1359 fn candidate_recipient_devices(&self, member_ids: &BTreeSet<UserId>, known_devices: &[StoredDevice]) -> Vec<StoredDevice> {
1363 known_devices
1364 .iter()
1365 .filter(|device| member_ids.contains(&device.user_id) && !device.blocked)
1366 .filter(|device| !(device.user_id == self.config.user_id && device.device_id == self.config.device_id))
1367 .cloned()
1368 .collect()
1369 }
1370
1371 fn begin_key_pipeline(&mut self, txn_id: &TxnId, room_id: &RoomId, now_ms: i64) -> Result<(), MessengerError> {
1376 let members = self.encrypted_room_member_ids(room_id);
1377 let outdated = self.mark_outdated_or_untracked(&members)?;
1378 if !outdated.is_empty() {
1379 let request_id = self.next_request_id()?;
1380 if let Some(request) = DeviceTracker::keys_query_request(&self.store, request_id.clone())? {
1381 self.send_request_owner.insert(request_id.clone(), txn_id.clone());
1382 self.enqueue_request(request, Lane::Other)?;
1383 self.set_send_phase(txn_id, SendPhase::AwaitingKeysQuery);
1384 return Ok(());
1385 }
1386 }
1387 self.begin_keys_claim_or_share(txn_id, room_id, now_ms)
1388 }
1389
1390 fn begin_keys_claim_or_share(&mut self, txn_id: &TxnId, room_id: &RoomId, now_ms: i64) -> Result<(), MessengerError> {
1394 let members = self.encrypted_room_member_ids(room_id);
1395 let known_devices = self.known_devices_for(&members)?;
1396 let candidates = self.candidate_recipient_devices(&members, &known_devices);
1397 let missing: Vec<StoredDevice> =
1398 OlmSessionManager::sessions_missing_for(&self.store, &candidates)?.into_iter().cloned().collect();
1399 if !missing.is_empty() {
1400 let request_id = self.next_request_id()?;
1401 let refs: Vec<&StoredDevice> = missing.iter().collect();
1402 if let Some(request) = OlmSessionManager::keys_claim_request(request_id.clone(), &refs) {
1403 self.send_request_owner.insert(request_id.clone(), txn_id.clone());
1404 self.enqueue_request(request, Lane::Other)?;
1405 self.set_send_phase(txn_id, SendPhase::AwaitingKeysClaim { devices: missing });
1406 return Ok(());
1407 }
1408 }
1409 self.begin_share_or_encrypt(txn_id, room_id, now_ms)
1410 }
1411
1412 fn begin_share_or_encrypt(&mut self, txn_id: &TxnId, room_id: &RoomId, now_ms: i64) -> Result<(), MessengerError> {
1425 let Some(encryption) = self.rooms.get(room_id).and_then(|room| room.encryption.clone()) else {
1426 self.build_plaintext_wire_body(txn_id);
1427 return self.begin_awaiting_send(txn_id, room_id);
1428 };
1429 let members = self.encrypted_room_member_ids(room_id);
1430 let known_devices = self.known_devices_for(&members)?;
1431 let update = GroupSessionManager::ensure_outbound_session(
1432 &mut self.store,
1433 room_id,
1434 &self.config.user_id,
1435 &self.config.device_id,
1436 &encryption,
1437 &members,
1438 &known_devices,
1439 now_ms,
1440 )?;
1441 let own_keys = self.account.identity_keys();
1442 GroupSessionManager::adopt_own_outbound_session(
1443 &mut self.store,
1444 room_id,
1445 &self.config.user_id,
1446 own_keys.curve25519,
1447 own_keys.ed25519,
1448 &update,
1449 )?;
1450
1451 if let Some(room_key_content) = update.room_key_content.clone() {
1452 let missing_ids: BTreeSet<(UserId, DeviceId)> = OlmSessionManager::sessions_missing_for(&self.store, &update.new_recipients)?
1453 .into_iter()
1454 .map(|device| (device.user_id.clone(), device.device_id.clone()))
1455 .collect();
1456 let sessioned: Vec<StoredDevice> = update
1457 .new_recipients
1458 .into_iter()
1459 .filter(|device| !missing_ids.contains(&(device.user_id.clone(), device.device_id.clone())))
1460 .collect();
1461
1462 let mut pending_ids = BTreeSet::new();
1463 for chunk in GroupSessionManager::chunk_recipients_for_send_to_device(&sessioned) {
1464 let body = GroupSessionManager::build_room_key_send_to_device_body(
1465 &mut self.store,
1466 &self.account,
1467 &self.config.user_id,
1468 &self.config.device_id,
1469 chunk,
1470 &room_key_content,
1471 )?;
1472 let request_id = self.next_request_id()?;
1473 let share_txn = self.next_txn_id()?;
1474 let request = OutgoingRequest::send_to_device(request_id.clone(), ROOM_KEY_SHARE_EVENT_TYPE, &share_txn, body);
1475 self.send_request_owner.insert(request_id.clone(), txn_id.clone());
1476 self.enqueue_request(request, Lane::ToDevice)?;
1477 pending_ids.insert(request_id);
1478 }
1479 if !pending_ids.is_empty() {
1480 self.set_send_phase(txn_id, SendPhase::AwaitingKeyShare { pending: pending_ids });
1481 return Ok(());
1482 }
1483 }
1484
1485 self.finish_encrypt_and_send(txn_id, room_id)?;
1486 self.begin_awaiting_send(txn_id, room_id)
1487 }
1488
1489 fn finish_encrypt_and_send(&mut self, txn_id: &TxnId, room_id: &RoomId) -> Result<(), MessengerError> {
1496 let Some(wire) = self.send_state.get(txn_id).and_then(|pending| pending.payload.room_event()) else {
1497 return Ok(());
1498 };
1499 let sender_curve = self.account.identity_keys().curve25519.to_base64();
1500 let encrypted = GroupSessionManager::encrypt_event(
1501 &mut self.store,
1502 room_id,
1503 &sender_curve,
1504 &self.config.device_id,
1505 &wire.event_type,
1506 wire.content,
1507 wire.relates_to,
1508 )?;
1509 let wire_content = serde_json::to_value(RoomEncryptedContent::Megolm(encrypted))?;
1510 if let Some(entry) = self.send_state.get_mut(txn_id) {
1511 entry.wire = Some(WireSend::Event { event_type: "m.room.encrypted".to_string(), content: wire_content });
1512 }
1513 Ok(())
1514 }
1515
1516 fn build_plaintext_wire_body(&mut self, txn_id: &TxnId) {
1522 let Some(pending) = self.send_state.get(txn_id) else { return };
1523 let wire = match &pending.payload {
1524 SendPayload::Message(_) | SendPayload::Event { .. } => match pending.payload.room_event() {
1525 Some(room_event) => WireSend::Event { event_type: room_event.event_type, content: room_event.content },
1526 None => return,
1527 },
1528 SendPayload::Reaction { target, key } => WireSend::Event {
1529 event_type: "m.reaction".to_string(),
1530 content: serde_json::json!({
1531 "m.relates_to": { "rel_type": "m.annotation", "event_id": target.as_str(), "key": key }
1532 }),
1533 },
1534 SendPayload::Redaction { target, reason } => WireSend::Redact { target: target.clone(), reason: reason.clone() },
1535 };
1536 if let Some(entry) = self.send_state.get_mut(txn_id) {
1537 entry.wire = Some(wire);
1538 }
1539 }
1540
1541 fn begin_awaiting_send(&mut self, txn_id: &TxnId, room_id: &RoomId) -> Result<(), MessengerError> {
1546 let Some(wire) = self.send_state.get(txn_id).and_then(|pending| pending.wire.clone()) else { return Ok(()) };
1547 let request_id = self.next_request_id()?;
1548 let request = match &wire {
1549 WireSend::Event { event_type, content } => {
1550 OutgoingRequest::room_send(request_id.clone(), room_id, event_type, txn_id, content.clone())
1551 }
1552 WireSend::Redact { target, reason } => {
1553 let body = match reason {
1554 Some(reason) => serde_json::json!({ "reason": reason }),
1555 None => serde_json::json!({}),
1556 };
1557 OutgoingRequest::room_redact(request_id.clone(), room_id, target, txn_id, body)
1558 }
1559 };
1560 self.send_request_owner.insert(request_id, txn_id.clone());
1561 self.enqueue_request(request, Lane::Room(room_id.clone()))?;
1562 self.set_send_phase(txn_id, SendPhase::AwaitingSend);
1563 Ok(())
1564 }
1565
1566 fn advance_send_on_success(
1572 &mut self,
1573 txn_id: &TxnId,
1574 request_id: &RequestId,
1575 kind: Option<OutgoingRequestKind>,
1576 response: &HttpResponseDescriptor,
1577 now_ms: i64,
1578 ) {
1579 match kind {
1580 Some(OutgoingRequestKind::KeysQuery) => {
1581 let Some(room_id) = self.send_state.get(txn_id).map(|pending| pending.room_id.clone()) else { return };
1582 let _ = self.begin_keys_claim_or_share(txn_id, &room_id, now_ms);
1583 }
1584 Some(OutgoingRequestKind::KeysClaim) => {
1585 let Some(pending) = self.send_state.get(txn_id) else { return };
1586 let room_id = pending.room_id.clone();
1587 if let SendPhase::AwaitingKeysClaim { devices, .. } = &pending.phase {
1588 let devices = devices.clone();
1589 let refs: Vec<&StoredDevice> = devices.iter().collect();
1590 let _ = OlmSessionManager::on_keys_claim_response(&mut self.store, &self.account, &refs, &response.body);
1591 }
1592 let _ = self.begin_share_or_encrypt(txn_id, &room_id, now_ms);
1593 }
1594 Some(OutgoingRequestKind::SendToDevice) => {
1595 let mut remaining = None;
1596 if let Some(entry) = self.send_state.get_mut(txn_id) {
1597 if let SendPhase::AwaitingKeyShare { pending } = &mut entry.phase {
1598 pending.remove(request_id);
1599 remaining = Some(pending.len());
1600 }
1601 }
1602 if remaining == Some(0) {
1603 let Some(room_id) = self.send_state.get(txn_id).map(|pending| pending.room_id.clone()) else { return };
1604 let _ = self.finish_encrypt_and_send(txn_id, &room_id);
1605 let _ = self.begin_awaiting_send(txn_id, &room_id);
1606 }
1607 }
1608 Some(OutgoingRequestKind::RoomSend) | Some(OutgoingRequestKind::RoomRedact) => {
1609 self.complete_send_success(txn_id, response, now_ms);
1610 }
1611 _ => {}
1612 }
1613 }
1614
1615 fn complete_send_success(&mut self, txn_id: &TxnId, response: &HttpResponseDescriptor, now_ms: i64) {
1621 let Some(pending) = self.send_state.get(txn_id) else { return };
1622 let room_id = pending.room_id.clone();
1623 if pending.payload.is_room_event() {
1624 if let Ok(body) = serde_json::from_slice::<SendResponseBody>(&response.body) {
1625 if let Some(timeline) = self.timelines.get_mut(&room_id) {
1626 timeline.set_echo_event_id(txn_id, body.event_id);
1627 }
1628 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
1629 }
1630 }
1631 self.send_state.remove(txn_id);
1632 self.finish_active_send(&room_id, txn_id, now_ms);
1633 }
1634
1635 fn fail_send(&mut self, txn_id: &TxnId, errcode: String, now_ms: i64) {
1641 let Some(pending) = self.send_state.get(txn_id) else { return };
1642 let room_id = pending.room_id.clone();
1643 if pending.payload.is_room_event() {
1644 if let Some(timeline) = self.timelines.get_mut(&room_id) {
1645 timeline.mark_send_state(txn_id, SendState::Failed { reason: errcode.clone() });
1646 }
1647 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
1648 }
1649 self.set_send_phase(txn_id, SendPhase::Failed { errcode });
1650 self.finish_active_send(&room_id, txn_id, now_ms);
1651 }
1652
1653 fn finish_active_send(&mut self, room_id: &RoomId, txn_id: &TxnId, now_ms: i64) {
1659 if let Some(queue) = self.pending_sends.get_mut(room_id) {
1660 if queue.front() == Some(txn_id) {
1661 queue.pop_front();
1662 } else {
1663 queue.retain(|id| id != txn_id);
1664 }
1665 }
1666 let next = self.pending_sends.get(room_id).and_then(|queue| queue.front().cloned());
1667 if let Some(next_txn) = next {
1668 let _ = self.start_send_pipeline(&next_txn, now_ms);
1669 }
1670 }
1671
1672 pub fn releasable_requests(&mut self, now_ms: i64) -> Vec<OutgoingRequest> {
1692 self.ensure_sync_enqueued();
1693 self.outgoing.releasable(&self.store, now_ms)
1694 }
1695
1696 fn ensure_sync_enqueued(&mut self) {
1697 if self.sync_request_id.is_some() {
1698 return;
1699 }
1700 let Ok(request_id) = self.next_request_id() else { return };
1701 let since = self.store.sync_token().ok().flatten().map(str::to_string);
1702 let timeout_ms = if since.is_some() && self.sync_caught_up { SYNC_TIMEOUT_MS } else { 0 };
1712 let request = OutgoingRequest::sync(request_id.clone(), since.as_deref(), Some(timeout_ms));
1713 if self.enqueue_request(request, Lane::Sync).is_ok() {
1714 self.sync_request_id = Some(request_id);
1715 }
1716 }
1717
1718 pub fn on_response(&mut self, request_id: RequestId, resp: HttpResponseDescriptor, now_ms: i64) -> Vec<MessengerEvent> {
1724 let kind = self.pending_kinds.get(&request_id).copied();
1725 let room_messages_room = self.pending_room_messages.get(&request_id).cloned();
1726 let create_room_kind = self.pending_create_room.get(&request_id).cloned();
1727 let send_owner = self.send_request_owner.get(&request_id).cloned();
1728
1729 let outcome = self.outgoing.on_response(&mut self.store, &request_id, resp, now_ms, self.jitter.as_mut());
1730 let (is_terminal, terminal_response, failed_errcode) = match outcome {
1731 Ok(Some(ResponseOutcome::Done(response))) => (true, Some(response), None),
1732 Ok(Some(ResponseOutcome::Failed { errcode, .. })) => (true, None, Some(errcode)),
1733 _ => (false, None, None),
1734 };
1735
1736 if is_terminal {
1737 self.pending_kinds.remove(&request_id);
1738 if matches!(kind, Some(OutgoingRequestKind::Sync)) {
1739 self.sync_request_id = None;
1740 if terminal_response.is_some() {
1741 self.sync_caught_up = true;
1742 }
1743 }
1744 if room_messages_room.is_some() {
1745 self.pending_room_messages.remove(&request_id);
1746 }
1747 self.pending_create_room.remove(&request_id);
1748 self.send_request_owner.remove(&request_id);
1749 if let Some(key) = self.account_data_write_keys.remove(&request_id) {
1750 self.settle_account_data_write(&key, terminal_response.is_some());
1751 }
1752 }
1753
1754 if let Some(response) = &terminal_response {
1755 if let Err(error) = self.handle_terminal_success(&request_id, kind, room_messages_room, create_room_kind, response) {
1762 self.ingest_error = Some(match kind {
1763 Some(kind) => format!("{kind:?} response not ingested: {error}"),
1764 None => format!("response not ingested: {error}"),
1765 });
1766 }
1767 }
1768
1769 if is_terminal
1770 && terminal_response.is_none()
1771 && matches!(kind, Some(OutgoingRequestKind::SigningKeysUpload) | Some(OutgoingRequestKind::SignaturesUpload))
1772 {
1773 self.xsign_inflight = false;
1774 }
1775
1776 if let Some(txn_id) = send_owner {
1777 if let Some(response) = &terminal_response {
1778 self.advance_send_on_success(&txn_id, &request_id, kind, response, now_ms);
1779 } else if let Some(errcode) = failed_errcode.clone() {
1780 self.fail_send(&txn_id, errcode, now_ms);
1781 }
1782 }
1783
1784 if matches!(kind, Some(OutgoingRequestKind::KeysUpload)) {
1788 if let Some(errcode) = failed_errcode {
1789 self.ingest_error = Some(format!("keys/upload failed: {errcode}"));
1790 }
1791 }
1792
1793 self.drain_events()
1794 }
1795
1796 pub fn on_transport_error(&mut self, request_id: &RequestId, now_ms: i64) {
1801 let _ = self.outgoing.on_transport_error(request_id, now_ms, self.jitter.as_mut());
1802 }
1803
1804 fn handle_terminal_success(
1805 &mut self,
1806 request_id: &RequestId,
1807 kind: Option<OutgoingRequestKind>,
1808 room_messages_room: Option<RoomId>,
1809 create_room_kind: Option<CreateRoomKind>,
1810 response: &HttpResponseDescriptor,
1811 ) -> Result<(), MessengerError> {
1812 match kind {
1813 Some(OutgoingRequestKind::Sync) => {
1814 let parsed = parse_sync_response(&response.body)?;
1815 self.ingest_sync(&parsed)?;
1816 }
1817 Some(OutgoingRequestKind::KeysUpload) => {
1818 self.account.on_keys_upload_response(&mut self.store)?;
1819 }
1820 Some(OutgoingRequestKind::KeysClaim) if self.history_claims.contains_key(request_id) => {
1821 if let Some(devices) = self.history_claims.remove(request_id) {
1822 let refs: Vec<&StoredDevice> = devices.iter().collect();
1823 let _ = OlmSessionManager::on_keys_claim_response(&mut self.store, &self.account, &refs, &response.body);
1824 }
1825 self.retry_missing_room_keys();
1826 }
1827 Some(OutgoingRequestKind::SigningKeysUpload) => {
1828 self.xsign_inflight = false;
1829 if let Some(state) = self.xsign.as_mut() {
1830 state.uploaded = true;
1831 let snapshot = state.clone();
1832 snapshot.save(&mut self.store)?;
1833 }
1834 self.maybe_enqueue_cross_signing()?;
1835 }
1836 Some(OutgoingRequestKind::SignaturesUpload) => {
1837 self.xsign_inflight = false;
1838 let failures_empty = serde_json::from_slice::<serde_json::Value>(&response.body)
1839 .ok()
1840 .and_then(|v| v.get("failures").and_then(|f| f.as_object()).map(|f| f.is_empty()))
1841 .unwrap_or(true);
1842 if failures_empty {
1843 if let Some(state) = self.xsign.as_mut() {
1844 state.device_signed = true;
1845 let snapshot = state.clone();
1846 snapshot.save(&mut self.store)?;
1847 }
1848 }
1849 }
1850 Some(OutgoingRequestKind::KeysQuery) => {
1851 let outcome = DeviceTracker::on_keys_query_response(&mut self.store, &response.body)?;
1852 if !outcome.key_changes.is_empty() {
1853 let room_ids: Vec<RoomId> = self.rooms.keys().cloned().collect();
1854 for change in &outcome.key_changes {
1855 let _ = GroupSessionManager::forget_shared_device(
1856 &mut self.store,
1857 &room_ids,
1858 &change.user_id,
1859 &change.device_id,
1860 );
1861 self.emit(MessengerEvent::DeviceKeyChanged {
1862 user_id: change.user_id.clone(),
1863 device_id: change.device_id.clone(),
1864 });
1865 }
1866 }
1867 self.retry_unknown_sender_queue()?;
1868 }
1869 Some(OutgoingRequestKind::RoomMessages) => {
1870 if let Some(room_id) = room_messages_room {
1871 self.ingest_room_messages_response(&room_id, &response.body)?;
1872 }
1873 }
1874 Some(OutgoingRequestKind::CreateRoom) => {
1875 if let Some(kind) = create_room_kind {
1876 self.on_create_room_response(kind, response)?;
1877 }
1878 }
1879 Some(OutgoingRequestKind::PublicRooms) if self.latest_public_rooms_request.as_ref() == Some(request_id) => {
1884 let body: PublicRoomsResponseBody = serde_json::from_slice(&response.body)?;
1885 self.public_rooms_result = body
1886 .chunk
1887 .into_iter()
1888 .map(|entry| PublicRoomsResultEntry {
1889 room_id: entry.room_id,
1890 name: entry.name,
1891 topic: entry.topic,
1892 num_joined_members: entry.num_joined_members,
1893 })
1894 .collect();
1895 self.emit(MessengerEvent::PublicRoomsChanged);
1896 }
1897 Some(OutgoingRequestKind::UserDirectorySearch) if self.latest_user_search_request.as_ref() == Some(request_id) => {
1899 let body: UserDirectorySearchResponseBody = serde_json::from_slice(&response.body)?;
1900 self.user_search_result = body
1901 .results
1902 .into_iter()
1903 .map(|entry| UserDirectoryResultEntry { user_id: entry.user_id, display_name: entry.display_name })
1904 .collect();
1905 self.emit(MessengerEvent::UserSearchChanged);
1906 }
1907 _ => {}
1908 }
1909 Ok(())
1910 }
1911
1912 fn on_create_room_response(&mut self, kind: CreateRoomKind, response: &HttpResponseDescriptor) -> Result<(), MessengerError> {
1918 let CreateRoomKind::Dm { peer } = kind else { return Ok(()) };
1919 let body: CreateRoomResponseBody = serde_json::from_slice(&response.body)?;
1920 self.merge_direct_account_data(body.room_id, peer)
1921 }
1922
1923 fn merge_direct_account_data(&mut self, room_id: RoomId, peer: UserId) -> Result<(), MessengerError> {
1935 let mut direct = self.direct_account_data.clone().unwrap_or_default();
1936 let rooms_for_peer = direct.0.entry(peer).or_default();
1937 if !rooms_for_peer.contains(&room_id) {
1938 rooms_for_peer.push(room_id);
1939 }
1940 let value = serde_json::to_value(&direct)?;
1941 self.write_account_data(None, "m.direct".to_string(), value)
1942 }
1943
1944 fn cached_account_data(&self, key: &AccountDataKey) -> Option<&serde_json::Value> {
1945 match &key.0 {
1946 Some(room_id) => self.room_account_data(room_id, &key.1),
1947 None => self.global_account_data(&key.1),
1948 }
1949 }
1950
1951 fn store_account_data_cache(&mut self, key: &AccountDataKey, value: Option<serde_json::Value>) {
1954 if key.0.is_none() && key.1 == "m.direct" {
1955 self.direct_account_data = value.as_ref().and_then(|v| serde_json::from_value::<DirectContent>(v.clone()).ok());
1956 }
1957 match (&key.0, value) {
1958 (None, Some(v)) => {
1959 self.global_account_data.insert(key.1.clone(), v);
1960 }
1961 (None, None) => {
1962 self.global_account_data.remove(&key.1);
1963 }
1964 (Some(room_id), Some(v)) => {
1965 self.room_account_data.entry(room_id.clone()).or_default().insert(key.1.clone(), v);
1966 }
1967 (Some(room_id), None) => {
1968 if let Some(by_type) = self.room_account_data.get_mut(room_id) {
1969 by_type.remove(&key.1);
1970 }
1971 }
1972 }
1973 }
1974
1975 fn write_account_data(&mut self, room_id: Option<RoomId>, event_type: String, value: serde_json::Value) -> Result<(), MessengerError> {
1981 let request_id = self.next_request_id()?;
1982 let request = match &room_id {
1983 Some(room_id) => OutgoingRequest::room_account_data(request_id.clone(), &self.config.user_id, room_id, &event_type, value.clone()),
1984 None => OutgoingRequest::account_data(request_id.clone(), &self.config.user_id, &event_type, value.clone()),
1985 };
1986 let key: AccountDataKey = (room_id, event_type);
1987 let known = self.cached_account_data(&key).cloned();
1988 let guard = self.account_data_guard.entry(key.clone()).or_insert(AccountDataGuard { pending: 0, server_value: known });
1989 guard.pending += 1;
1990 self.store_account_data_cache(&key, Some(value));
1991 match self.enqueue_request(request, Lane::AccountData) {
1992 Ok(()) => {
1993 self.account_data_write_keys.insert(request_id, key);
1994 Ok(())
1995 }
1996 Err(error) => {
1997 self.settle_account_data_write(&key, false);
1998 Err(error)
1999 }
2000 }
2001 }
2002
2003 fn settle_account_data_write(&mut self, key: &AccountDataKey, succeeded: bool) {
2008 let Some(guard) = self.account_data_guard.get_mut(key) else { return };
2009 guard.pending = guard.pending.saturating_sub(1);
2010 if guard.pending > 0 {
2011 return;
2012 }
2013 let Some(guard) = self.account_data_guard.remove(key) else { return };
2014 if !succeeded {
2015 self.store_account_data_cache(key, guard.server_value);
2016 self.emit(MessengerEvent::RoomsChanged);
2017 }
2018 }
2019
2020 fn apply_synced_account_data(&mut self, room_id: Option<&RoomId>, event_type: &str, content: serde_json::Value) {
2024 let key: AccountDataKey = (room_id.cloned(), event_type.to_string());
2025 if let Some(guard) = self.account_data_guard.get_mut(&key) {
2026 guard.server_value = Some(content);
2027 return;
2028 }
2029 self.store_account_data_cache(&key, Some(content));
2030 }
2031
2032 pub fn take_flush_batch(&mut self) -> Option<FlushBatch> {
2035 if let Some(batch) = self.pending_flush_batches.pop_front() {
2036 return Some(batch);
2037 }
2038 self.store.take_flush_batch()
2039 }
2040
2041 pub fn ack_flush(&mut self, id: u64) {
2043 self.store.ack_flush(id);
2044 }
2045
2046 pub fn events(&mut self) -> Vec<MessengerEvent> {
2049 self.drain_events()
2050 }
2051
2052 pub fn take_ingest_error(&mut self) -> Option<String> {
2055 self.ingest_error.take()
2056 }
2057
2058 pub fn change_counter(&self) -> u64 {
2063 self.change_counter
2064 }
2065
2066 pub fn room_ids(&self) -> impl Iterator<Item = &RoomId> {
2068 self.rooms.keys()
2069 }
2070
2071 pub fn room_state(&self, room_id: &RoomId) -> Option<&RoomState> {
2073 self.rooms.get(room_id)
2074 }
2075
2076 pub fn timeline(&self, room_id: &RoomId) -> Option<&Timeline> {
2078 self.timelines.get(room_id)
2079 }
2080
2081 pub fn room_kind(&self, room_id: &RoomId) -> Option<RoomKind> {
2085 self.rooms.get(room_id).map(|state| state.derive_room_kind(room_id, self.direct_account_data.as_ref()))
2086 }
2087
2088 pub fn typing_users(&self, room_id: &RoomId) -> &[UserId] {
2092 self.typing.get(room_id).map(Vec::as_slice).unwrap_or(&[])
2093 }
2094
2095 pub fn receipts(&self, room_id: &RoomId) -> Option<&ReceiptContent> {
2097 self.receipts.get(room_id)
2098 }
2099
2100 pub fn room_account_data(&self, room_id: &RoomId, event_type: &str) -> Option<&serde_json::Value> {
2103 self.room_account_data.get(room_id).and_then(|by_type| by_type.get(event_type))
2104 }
2105
2106 pub fn global_account_data(&self, event_type: &str) -> Option<&serde_json::Value> {
2109 self.global_account_data.get(event_type)
2110 }
2111
2112 pub fn key_request_count(&self) -> usize {
2114 self.key_requests_sent.len()
2115 }
2116
2117 pub fn user_id(&self) -> &UserId {
2122 &self.config.user_id
2123 }
2124
2125 pub fn has_tracked_device(&self, user_id: &UserId) -> bool {
2131 DeviceTracker::devices_for_user(&self.store, user_id).map(|devices| !devices.is_empty()).unwrap_or(false)
2132 }
2133
2134 pub fn public_rooms_result(&self) -> &[PublicRoomsResultEntry] {
2137 &self.public_rooms_result
2138 }
2139
2140 pub fn user_search_result(&self) -> &[UserDirectoryResultEntry] {
2143 &self.user_search_result
2144 }
2145
2146 pub fn send_failure_reason(&self, txn_id: &TxnId) -> Option<&str> {
2151 match self.send_state.get(txn_id).map(|pending| &pending.phase) {
2152 Some(SendPhase::Failed { errcode }) => Some(errcode.as_str()),
2153 _ => None,
2154 }
2155 }
2156
2157 pub fn dispatch(&mut self, cmd: MessengerCommand, now_ms: i64) -> Result<(), MessengerError> {
2164 match cmd {
2165 MessengerCommand::SendMessage { room_id, message, txn_id } => {
2166 if let Some(edit_of) = &message.edit_of {
2167 let owned = self
2168 .timelines
2169 .get(&room_id)
2170 .and_then(|timeline| timeline.item_by_event_id(edit_of))
2171 .is_some_and(|item| item.sender == self.config.user_id);
2172 if !owned {
2173 return Err(MessengerError::EditNotOwned { event_id: edit_of.clone() });
2174 }
2175 }
2176 self.start_new_send(room_id, SendPayload::Message(message), txn_id, now_ms)?;
2177 }
2178 MessengerCommand::Forward { from_room, event_id, to_room, txn_id } => {
2179 let payload = self.build_forward_payload(&from_room, &event_id, &to_room)?;
2180 self.start_new_send(to_room, payload, txn_id, now_ms)?;
2181 }
2182 MessengerCommand::React { room_id, target, key } => {
2183 self.start_new_send(room_id, SendPayload::Reaction { target, key }, None, now_ms)?;
2184 }
2185 MessengerCommand::Redact { room_id, target, reason } => {
2186 self.start_new_send(room_id, SendPayload::Redaction { target, reason }, None, now_ms)?;
2187 }
2188 MessengerCommand::RetrySend { room_id, txn_id } => {
2189 self.retry_send(&room_id, &txn_id, now_ms)?;
2190 }
2191 MessengerCommand::CreateRoom { kind } => {
2192 self.dispatch_create_room(kind)?;
2193 }
2194 MessengerCommand::JoinRoom { room_id } => {
2195 let request_id = self.next_request_id()?;
2196 let request = OutgoingRequest::join_room(request_id, room_id.as_str());
2197 self.enqueue_request(request, Lane::Other)?;
2198
2199 let dm_peer = self.rooms.get(&room_id).and_then(|room| {
2206 let invited_as_direct = room.members.get(&self.config.user_id).is_some_and(|member| member.is_direct);
2207 if !invited_as_direct {
2208 return None;
2209 }
2210 room.members.keys().find(|user_id| **user_id != self.config.user_id).cloned()
2211 });
2212 if let Some(peer) = dm_peer {
2213 self.merge_direct_account_data(room_id, peer)?;
2214 }
2215 }
2216 MessengerCommand::LeaveRoom { room_id } => {
2217 let request_id = self.next_request_id()?;
2218 let request = OutgoingRequest::leave_room(request_id, &room_id);
2219 self.enqueue_request(request, Lane::Other)?;
2220 }
2221 MessengerCommand::Invite { room_id, user_id } => {
2222 let request_id = self.next_request_id()?;
2223 let request = OutgoingRequest::invite(request_id, &room_id, &user_id);
2224 self.enqueue_request(request, Lane::Other)?;
2225 }
2226 MessengerCommand::Kick { room_id, user_id, reason } => {
2227 let request_id = self.next_request_id()?;
2228 let request = OutgoingRequest::kick(request_id, &room_id, &user_id, reason.as_deref());
2229 self.enqueue_request(request, Lane::Other)?;
2230 }
2231 MessengerCommand::SetTag { room_id, tag, order } => {
2232 let mut content = self.current_tag_content(&room_id);
2233 content.tags.insert(tag, TagInfo { order });
2234 self.enqueue_tag_update(room_id, content)?;
2235 }
2236 MessengerCommand::RemoveTag { room_id, tag } => {
2237 let mut content = self.current_tag_content(&room_id);
2238 content.tags.remove(&tag);
2239 self.enqueue_tag_update(room_id, content)?;
2240 }
2241 MessengerCommand::MarkRead { room_id, event_id } => {
2242 let request_id = self.next_request_id()?;
2243 let body = serde_json::json!({ "m.fully_read": event_id.as_str(), "m.read": event_id.as_str() });
2244 let request = OutgoingRequest::read_markers(request_id, &room_id, body);
2245 self.enqueue_request(request, Lane::Other)?;
2246 }
2247 MessengerCommand::SetTyping { room_id, typing } => {
2248 if typing {
2249 let should_send = match self.typing_debounce.get(&room_id) {
2250 Some(&last_sent_ms) => now_ms.saturating_sub(last_sent_ms) >= TYPING_DEBOUNCE_MS,
2251 None => true,
2252 };
2253 if !should_send {
2254 return Ok(());
2255 }
2256 self.typing_debounce.insert(room_id.clone(), now_ms);
2257 } else {
2258 self.typing_debounce.remove(&room_id);
2259 }
2260 let request_id = self.next_request_id()?;
2261 let timeout_ms = typing.then_some(TYPING_TIMEOUT_MS);
2262 let request = OutgoingRequest::typing(request_id, &room_id, &self.config.user_id, typing, timeout_ms);
2263 self.enqueue_request(request, Lane::Other)?;
2264 }
2265 MessengerCommand::RetryDecryption { room_id, event_id } => {
2266 if self.retry_decrypt_one(&room_id, &event_id) {
2267 self.emit(MessengerEvent::TimelineChanged { room_id });
2268 }
2269 }
2270 MessengerCommand::LoadOlder { room_id } => {
2271 let from = self
2272 .timelines
2273 .get(&room_id)
2274 .and_then(|timeline| {
2275 timeline
2276 .older_token()
2277 .map(str::to_string)
2278 .or_else(|| timeline.gap().and_then(|gap| gap.prev_batch.clone()))
2279 })
2280 .or_else(|| {
2286 if self.rooms.contains_key(&room_id) {
2287 self.store.sync_token().ok().flatten().map(str::to_string)
2288 } else {
2289 None
2290 }
2291 });
2292 let Some(from) = from else { return Ok(()) };
2293 let request_id = self.next_request_id()?;
2294 let request =
2295 OutgoingRequest::room_messages(request_id.clone(), &room_id, &from, "b", Some(LOAD_OLDER_PAGE_SIZE));
2296 self.enqueue_request(request, Lane::Room(room_id.clone()))?;
2297 self.pending_room_messages.insert(request_id, room_id);
2298 }
2299 MessengerCommand::SetAccountData { event_type, content } => {
2300 self.write_account_data(None, event_type, content)?;
2301 }
2302 MessengerCommand::SetRoomAccountData { room_id, event_type, content } => {
2303 self.write_account_data(Some(room_id), event_type, content)?;
2304 }
2305 MessengerCommand::SearchPublicRooms { term } => {
2306 let request_id = self.next_request_id()?;
2307 let body = serde_json::json!({ "filter": { "generic_search_term": term }, "limit": 50 });
2308 let request = OutgoingRequest::public_rooms(request_id.clone(), body);
2309 self.latest_public_rooms_request = Some(request_id);
2310 self.enqueue_request(request, Lane::Other)?;
2311 }
2312 MessengerCommand::SearchUsers { term } => {
2313 let request_id = self.next_request_id()?;
2314 let body = serde_json::json!({ "search_term": term, "limit": 20 });
2315 let request = OutgoingRequest::user_directory_search(request_id.clone(), body);
2316 self.latest_user_search_request = Some(request_id);
2317 self.enqueue_request(request, Lane::Other)?;
2318 }
2319 }
2320 Ok(())
2321 }
2322
2323 fn dispatch_create_room(&mut self, kind: CreateRoomKind) -> Result<(), MessengerError> {
2328 let body = match &kind {
2329 CreateRoomKind::Dm { peer } => serde_json::json!({
2330 "is_direct": true,
2331 "invite": [peer.as_str()],
2332 }),
2333 CreateRoomKind::Group { name, invite, members_can_invite } => serde_json::json!({
2334 "visibility": "private",
2335 "is_direct": false,
2336 "name": name,
2337 "invite": invite.iter().map(UserId::as_str).collect::<Vec<_>>(),
2338 "members_can_invite": members_can_invite,
2339 }),
2340 CreateRoomKind::Channel { name, topic } => {
2341 let mut value = serde_json::json!({
2342 "visibility": "public",
2343 "is_direct": false,
2344 "name": name,
2345 });
2346 if let Some(topic) = topic {
2347 value["topic"] = serde_json::Value::String(topic.clone());
2348 }
2349 value
2350 }
2351 };
2352 let request_id = self.next_request_id()?;
2353 let request = OutgoingRequest::create_room(request_id.clone(), body);
2354 self.pending_create_room.insert(request_id.clone(), kind);
2355 self.enqueue_request(request, Lane::Other)?;
2356 Ok(())
2357 }
2358
2359 fn current_tag_content(&self, room_id: &RoomId) -> TagContent {
2364 self.room_account_data(room_id, "m.tag")
2365 .and_then(|value| serde_json::from_value(value.clone()).ok())
2366 .unwrap_or_default()
2367 }
2368
2369 fn enqueue_tag_update(&mut self, room_id: RoomId, content: TagContent) -> Result<(), MessengerError> {
2375 let value = serde_json::to_value(&content)?;
2376 self.write_account_data(Some(room_id), "m.tag".to_string(), value)
2377 }
2378
2379 fn build_forward_payload(
2384 &self,
2385 from_room: &RoomId,
2386 event_id: &EventId,
2387 to_room: &RoomId,
2388 ) -> Result<SendPayload, MessengerError> {
2389 let refuse = |reason: &dyn std::fmt::Display| {
2390 MessengerError::IntentRefused(format!("cannot forward {event_id} from {from_room}: {reason}"))
2391 };
2392 if !self.rooms.contains_key(to_room) {
2393 return Err(refuse(&format!("the target room {to_room} is unknown")));
2394 }
2395 let source = match self.timelines.get(from_room) {
2396 Some(timeline) => timeline.forward_source(event_id),
2397 None => Err(ForwardRefusal::Missing),
2398 }
2399 .map_err(|reason| refuse(&reason))?;
2400 let (marker_key, marker) = self.forward_marker(from_room);
2401 Ok(SendPayload::Event {
2402 event_type: source.event_type,
2403 content: forwardable_content(source.content, marker_key, marker),
2404 })
2405 }
2406
2407 fn forward_marker(&self, from_room: &RoomId) -> (&'static str, serde_json::Value) {
2412 let public_channel = self.room_kind(from_room) == Some(RoomKind::Channel);
2413 if public_channel {
2414 let room_name = self.forward_room_name(from_room);
2415 (KEY_FORWARDED_FROM, serde_json::json!({ "room_id": from_room.as_str(), "room_name": room_name }))
2416 } else {
2417 (KEY_FORWARDED, serde_json::Value::Bool(true))
2418 }
2419 }
2420
2421 fn forward_room_name(&self, room_id: &RoomId) -> String {
2425 if let Some(name) = self.rooms.get(room_id).and_then(|room| room.name.clone()) {
2426 return name;
2427 }
2428 match self.room_kind(room_id) {
2429 Some(RoomKind::Dm) => "Direct",
2430 Some(RoomKind::Channel) => "Channel",
2431 Some(RoomKind::Group) | None => "Group",
2432 }
2433 .to_string()
2434 }
2435
2436 fn retry_decrypt_one(&mut self, room_id: &RoomId, event_id: &EventId) -> bool {
2442 let Some(item) = self.pending_decrypt.get(room_id).and_then(|by_event| by_event.get(event_id)).cloned() else {
2443 return false;
2444 };
2445 match GroupSessionManager::decrypt_event(
2446 &mut self.store,
2447 room_id,
2448 event_id,
2449 item.origin_server_ts,
2450 &item.sender,
2451 &item.content,
2452 ) {
2453 Ok(plaintext) => {
2454 self.apply_decrypted_plaintext(room_id, event_id, &item.sender, item.origin_server_ts, &plaintext);
2455 }
2456 Err(GroupDecryptError::MissingSession { .. }) => {
2457 let _ = self.request_missing_room_key(room_id, &item.sender, &item.content);
2458 return false;
2459 }
2460 Err(_) => return false,
2461 }
2462 if let Some(by_event) = self.pending_decrypt.get_mut(room_id) {
2463 by_event.remove(event_id);
2464 }
2465 true
2466 }
2467
2468 pub fn pull_latest_page(&mut self, room_id: RoomId) -> Result<(), MessengerError> {
2471 let request_id = self.next_request_id()?;
2472 let request = OutgoingRequest::room_messages_latest(request_id.clone(), &room_id, 20);
2473 self.enqueue_request(request, Lane::Room(room_id.clone()))?;
2474 self.pending_room_messages.insert(request_id, room_id);
2475 Ok(())
2476 }
2477
2478 pub fn decrypt_loaded_timeline(&mut self) {
2482 let pending: Vec<(RoomId, Vec<EventId>)> = self
2483 .pending_decrypt
2484 .iter()
2485 .map(|(room_id, by_event)| (room_id.clone(), by_event.keys().cloned().collect()))
2486 .collect();
2487 for (room_id, event_ids) in pending {
2488 for event_id in event_ids {
2489 self.retry_decrypt_one(&room_id, &event_id);
2490 }
2491 }
2492 let rooms: Vec<RoomId> = self.timelines.keys().cloned().collect();
2493 for room_id in rooms {
2494 let sealed: Vec<(EventId, i64, UserId, MegolmEncryptedContent)> = {
2495 let Some(timeline) = self.timelines.get(&room_id) else {
2496 continue;
2497 };
2498 timeline
2499 .items()
2500 .iter()
2501 .filter_map(|item| {
2502 let event_id = item.event_id.clone()?;
2503 let (session_id, ciphertext) = match &item.content {
2504 ItemContent::Encrypted { session_id, ciphertext_b64 } => {
2505 (session_id.clone(), ciphertext_b64.clone())
2506 }
2507 _ => return None,
2508 };
2509 Some((
2510 event_id,
2511 item.origin_server_ts,
2512 item.sender.clone(),
2513 MegolmEncryptedContent {
2514 ciphertext,
2515 session_id,
2516 sender_key: None,
2517 device_id: None,
2518 relates_to: None,
2519 },
2520 ))
2521 })
2522 .collect()
2523 };
2524 for (event_id, origin_server_ts, sender, content) in sealed {
2525 let plaintext = match GroupSessionManager::decrypt_event(
2526 &mut self.store,
2527 &room_id,
2528 &event_id,
2529 origin_server_ts,
2530 &sender,
2531 &content,
2532 ) {
2533 Ok(plaintext) => plaintext,
2534 Err(GroupDecryptError::MissingSession { .. }) => {
2535 let _ = self.request_missing_room_key(&room_id, &sender, &content);
2536 continue;
2537 }
2538 Err(_) => continue,
2539 };
2540 self.apply_decrypted_plaintext(
2541 &room_id,
2542 &event_id,
2543 &sender,
2544 origin_server_ts,
2545 &plaintext,
2546 );
2547 }
2548 }
2549 }
2550
2551 fn apply_decrypted_plaintext(
2564 &mut self,
2565 room_id: &RoomId,
2566 event_id: &EventId,
2567 sender: &UserId,
2568 origin_server_ts: i64,
2569 plaintext: &RoomEventPlaintext,
2570 ) {
2571 let Some(timeline) = self.timelines.get_mut(room_id) else { return };
2572 if let Some(relation) = RelatesTo::from_content(&plaintext.content) {
2573 if matches!(relation, RelatesTo::Annotation { .. } | RelatesTo::Replace { .. }) {
2574 let new_content = plaintext.content.get("m.new_content").cloned();
2575 timeline.apply_decrypted_relation(event_id, sender.clone(), origin_server_ts, relation, new_content);
2576 return;
2577 }
2578 }
2579 timeline.set_decrypted_event(event_id, &plaintext.event_type, &plaintext.content);
2580 }
2581
2582 fn retry_pending_for_session(&mut self, room_id: &RoomId, session_id: &str) {
2583 let Some(by_event) = self.pending_decrypt.get(room_id) else { return };
2584 let candidates: Vec<EventId> =
2585 by_event.iter().filter(|(_, item)| item.session_id == session_id).map(|(event_id, _)| event_id.clone()).collect();
2586 let mut changed = false;
2587 for event_id in candidates {
2588 if self.retry_decrypt_one(room_id, &event_id) {
2589 changed = true;
2590 }
2591 }
2592 if changed {
2593 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
2594 }
2595 }
2596
2597 fn ingest_sync(&mut self, response: &crate::wire::sync::SyncResponse) -> Result<(), MessengerError> {
2598 let mut newly_received_sessions: Vec<(RoomId, String)> = Vec::new();
2599
2600 for event in &response.to_device.events {
2602 self.process_to_device_event(event, &mut newly_received_sessions)?;
2603 }
2604
2605 DeviceTracker::on_device_lists(&mut self.store, &response.device_lists.changed, &response.device_lists.left)?;
2607
2608 let published_otk_count = response.device_one_time_keys_count.get("signed_curve25519").copied().unwrap_or(0);
2610 self.account.on_sync_counts(&mut self.store, published_otk_count, &response.device_unused_fallback_key_types)?;
2611 let keys_upload_request_id = self.next_request_id()?;
2612 if let Some(request) =
2613 self.account.keys_upload_request(keys_upload_request_id, &self.config.user_id, &self.config.device_id)?
2614 {
2615 self.enqueue_request(request, Lane::Other)?;
2616 }
2617 self.maybe_enqueue_cross_signing()?;
2618
2619 let outdated_before_rooms = self.outdated_tracked_users()?;
2621 if !outdated_before_rooms.is_empty() {
2622 self.enqueue_keys_query()?;
2623 }
2624
2625 let rooms_touched =
2627 !response.rooms.join.is_empty() || !response.rooms.invite.is_empty() || !response.rooms.leave.is_empty();
2628 for (room_id, joined) in &response.rooms.join {
2629 self.ingest_joined_room(room_id, joined)?;
2630 }
2631 for (room_id, invited) in &response.rooms.invite {
2632 self.ingest_invited_room(room_id, invited)?;
2633 }
2634 for (room_id, left) in &response.rooms.leave {
2635 self.ingest_left_room(room_id, left)?;
2636 }
2637 if rooms_touched {
2638 self.emit(MessengerEvent::RoomsChanged);
2639 }
2640 if !self.outdated_tracked_users()?.is_subset(&outdated_before_rooms) {
2644 self.enqueue_keys_query()?;
2645 }
2646
2647 for event in &response.account_data.events {
2649 self.apply_synced_account_data(None, &event.event_type, event.content.clone());
2650 }
2651
2652 for (room_id, session_id) in &newly_received_sessions {
2653 self.retry_pending_for_session(room_id, session_id);
2654 }
2655
2656 self.store.save_sync_token(response.next_batch.clone())?;
2659 Ok(())
2660 }
2661
2662 fn process_to_device_event(
2663 &mut self,
2664 event: &ToDeviceEvent,
2665 newly_received_sessions: &mut Vec<(RoomId, String)>,
2666 ) -> Result<(), MessengerError> {
2667 if event.event_type != "m.room.encrypted" {
2668 return Ok(());
2669 }
2670 let Ok(RoomEncryptedContent::Olm(content)) = serde_json::from_value::<RoomEncryptedContent>(event.content.clone())
2671 else {
2672 return Ok(());
2673 };
2674 match OlmSessionManager::decrypt_to_device(&mut self.store, &mut self.account, &self.config.user_id, &event.sender, &content) {
2675 Ok(decrypted) => self.route_decrypted_to_device(decrypted, newly_received_sessions),
2676 Err(OlmDecryptError::UnknownSenderDevice { .. }) => {
2677 self.queue_unknown_sender_retry(event.clone());
2678 Ok(())
2679 }
2680 Err(_) => Ok(()),
2687 }
2688 }
2689
2690 fn route_decrypted_to_device(
2691 &mut self,
2692 decrypted: DecryptedToDevice,
2693 newly_received_sessions: &mut Vec<(RoomId, String)>,
2694 ) -> Result<(), MessengerError> {
2695 match decrypted.event_type.as_str() {
2696 "m.room_key" => {
2697 if let Ok(content) = serde_json::from_value::<RoomKeyContent>(decrypted.content.clone()) {
2698 if content.algorithm == "m.megolm.v1.aes-sha2" {
2699 GroupSessionManager::receive_room_key(
2700 &mut self.store,
2701 &decrypted.sender,
2702 decrypted.sender_device_curve25519,
2703 decrypted.sender_ed25519,
2704 &content,
2705 )?;
2706 newly_received_sessions.push((content.room_id, content.session_id));
2707 }
2708 }
2709 }
2710 "m.forwarded_room_key" => {
2711 if let Ok(content) = serde_json::from_value::<ForwardedRoomKeyContent>(decrypted.content.clone()) {
2712 if let Some((user_id, curve, ed)) = self.forward_binding(&decrypted, &content)? {
2713 let stored = GroupSessionManager::receive_forwarded_room_key(
2714 &mut self.store,
2715 &user_id,
2716 curve,
2717 ed,
2718 &content,
2719 )?;
2720 if stored {
2721 newly_received_sessions.push((content.room_id, content.session_id));
2722 }
2723 }
2724 }
2725 }
2726 "m.room_key_request" => {
2727 if let Ok(content) = serde_json::from_value::<RoomKeyRequestContent>(decrypted.content.clone()) {
2728 self.answer_room_key_request(&decrypted, &content)?;
2729 }
2730 }
2731 "m.room_key.withheld" => {
2732 if let Ok(content) = serde_json::from_value::<RoomKeyWithheldContent>(decrypted.content.clone()) {
2733 if let (Some(room_id), Some(session_id)) = (content.room_id, content.session_id) {
2734 self.withheld_sessions.insert((room_id, session_id));
2735 }
2736 }
2737 }
2738 _ => {}
2739 }
2740 Ok(())
2741 }
2742
2743 fn forward_binding(
2749 &self,
2750 decrypted: &DecryptedToDevice,
2751 content: &ForwardedRoomKeyContent,
2752 ) -> Result<Option<(UserId, mail4agent_vodozemac::Curve25519PublicKey, mail4agent_vodozemac::Ed25519PublicKey)>, MessengerError>
2753 {
2754 let members = self.encrypted_room_member_ids(&content.room_id);
2755 let forwarder_is_us = decrypted.sender == self.config.user_id;
2756 if !forwarder_is_us && !members.contains(&decrypted.sender) {
2757 return Ok(None);
2758 }
2759 for user_id in &members {
2760 for device in DeviceTracker::devices_for_user(&self.store, user_id)? {
2761 if device.blocked || device.curve25519.to_base64() != content.sender_key {
2762 continue;
2763 }
2764 let same_user = decrypted.sender == device.user_id;
2765 if same_user || forwarder_is_us {
2766 return Ok(Some((device.user_id, device.curve25519, device.ed25519)));
2767 }
2768 }
2769 }
2770 if decrypted.sender_device_curve25519.to_base64() == content.sender_key {
2771 return Ok(Some((
2772 decrypted.sender.clone(),
2773 decrypted.sender_device_curve25519,
2774 decrypted.sender_ed25519,
2775 )));
2776 }
2777 Ok(None)
2778 }
2779
2780 fn answer_room_key_request(
2785 &mut self,
2786 decrypted: &DecryptedToDevice,
2787 content: &RoomKeyRequestContent,
2788 ) -> Result<(), MessengerError> {
2789 if matches!(content.action, RoomKeyRequestAction::RequestCancellation) {
2790 self.key_requests_answered.remove(&(
2791 decrypted.sender.clone(),
2792 content.requesting_device_id.clone(),
2793 content.request_id.clone(),
2794 ));
2795 return Ok(());
2796 }
2797 if !matches!(content.action, RoomKeyRequestAction::Request) {
2798 return Ok(());
2799 }
2800 let Some(body) = content.body.as_ref() else { return Ok(()) };
2801 if body.algorithm != "m.megolm.v1.aes-sha2" {
2802 return Ok(());
2803 }
2804 let members = self.encrypted_room_member_ids(&body.room_id);
2805 if !members.contains(&decrypted.sender) {
2806 return Ok(());
2807 }
2808 if decrypted.sender != self.config.user_id && !self.history_forwarding_allowed(&body.room_id) {
2809 return Ok(());
2810 }
2811 if decrypted.sender == self.config.user_id && content.requesting_device_id == self.config.device_id {
2812 return Ok(());
2813 }
2814 let dedup = (decrypted.sender.clone(), content.requesting_device_id.clone(), content.request_id.clone());
2815 if self.key_requests_answered.contains(&dedup) {
2816 return Ok(());
2817 }
2818 let Some(target) = self.requester_device(decrypted, &content.requesting_device_id)? else {
2819 return Ok(());
2820 };
2821 let our_curve = self.account.identity_keys().curve25519.to_base64();
2822 let (inner_type, payload) =
2823 if let Some(forwarded) = GroupSessionManager::forwarded_room_key_content(
2824 &self.store,
2825 &body.room_id,
2826 &body.session_id,
2827 &our_curve,
2828 )? {
2829 ("m.forwarded_room_key", serde_json::to_value(&forwarded)?)
2830 } else if let Some(room_key) =
2831 GroupSessionManager::outbound_room_key_content(&self.store, &body.room_id, &body.session_id)?
2832 {
2833 ("m.room_key", serde_json::to_value(&room_key)?)
2834 } else {
2835 return Ok(());
2836 };
2837 let sent = self.enqueue_olm_to_devices(std::slice::from_ref(&target), inner_type, payload)?;
2838 if sent > 0 {
2839 self.key_requests_answered.insert(dedup);
2840 }
2841 Ok(())
2842 }
2843
2844 fn requester_device(
2848 &self,
2849 decrypted: &DecryptedToDevice,
2850 requesting_device_id: &DeviceId,
2851 ) -> Result<Option<StoredDevice>, MessengerError> {
2852 let devices = DeviceTracker::devices_for_user(&self.store, &decrypted.sender)?;
2853 if let Some(device) = devices.iter().find(|device| &device.device_id == requesting_device_id) {
2854 return Ok((!device.blocked).then(|| device.clone()));
2855 }
2856 let curve = decrypted.sender_device_curve25519.to_base64();
2857 if let Some(device) = devices.iter().find(|device| device.curve25519.to_base64() == curve) {
2858 return Ok((!device.blocked).then(|| device.clone()));
2859 }
2860 Ok(None)
2861 }
2862
2863 fn request_missing_room_key(
2868 &mut self,
2869 room_id: &RoomId,
2870 sender: &UserId,
2871 content: &MegolmEncryptedContent,
2872 ) -> Result<(), MessengerError> {
2873 let key = (room_id.clone(), content.session_id.clone());
2874 if self.key_requests_sent.contains(&key) || self.withheld_sessions.contains(&key) {
2875 return Ok(());
2876 }
2877 let sender_key = match content.sender_key.clone() {
2878 Some(sender_key) => sender_key,
2879 None => {
2880 let devices = DeviceTracker::devices_for_user(&self.store, sender)?;
2881 let found = content
2882 .device_id
2883 .as_ref()
2884 .and_then(|device_id| devices.iter().find(|device| &device.device_id == device_id))
2885 .or_else(|| devices.iter().find(|device| !device.blocked));
2886 match found {
2887 Some(device) => device.curve25519.to_base64(),
2888 None => return Ok(()),
2889 }
2890 }
2891 };
2892 let mut targets = DeviceTracker::devices_for_user(&self.store, sender)?;
2893 targets.extend(DeviceTracker::devices_for_user(&self.store, &self.config.user_id)?);
2894 if targets.is_empty() {
2895 return Ok(());
2896 }
2897 let missing: Vec<StoredDevice> = OlmSessionManager::sessions_missing_for(&self.store, &targets)?
2901 .into_iter()
2902 .filter(|device| {
2903 !device.blocked
2904 && !(device.user_id == self.config.user_id && device.device_id == self.config.device_id)
2905 && !self.history_claim_tried.contains(&(device.user_id.clone(), device.device_id.clone()))
2906 })
2907 .cloned()
2908 .collect();
2909 if !missing.is_empty() {
2910 let claim_id = self.next_request_id()?;
2911 let refs: Vec<&StoredDevice> = missing.iter().collect();
2912 if let Some(request) = OlmSessionManager::keys_claim_request(claim_id.clone(), &refs) {
2913 for device in &missing {
2914 self.history_claim_tried.insert((device.user_id.clone(), device.device_id.clone()));
2915 }
2916 self.history_claims.insert(claim_id, missing);
2917 self.enqueue_request(request, Lane::Other)?;
2918 return Ok(());
2919 }
2920 }
2921 let request_id = self.next_txn_id()?;
2922 let body = RoomKeyRequestContent {
2923 action: RoomKeyRequestAction::Request,
2924 body: Some(RoomKeyRequestBody {
2925 algorithm: "m.megolm.v1.aes-sha2".to_string(),
2926 room_id: room_id.clone(),
2927 session_id: content.session_id.clone(),
2928 sender_key,
2929 }),
2930 request_id: request_id.as_str().to_string(),
2931 requesting_device_id: self.config.device_id.clone(),
2932 };
2933 let sent = self.enqueue_olm_to_devices(&targets, "m.room_key_request", serde_json::to_value(&body)?)?;
2934 if sent > 0 {
2935 self.key_requests_sent.insert(key);
2936 }
2937 Ok(())
2938 }
2939
2940 fn enqueue_olm_to_devices(
2945 &mut self,
2946 devices: &[StoredDevice],
2947 inner_type: &str,
2948 content: serde_json::Value,
2949 ) -> Result<usize, MessengerError> {
2950 let mut unique = Vec::new();
2951 for device in devices {
2952 if device.blocked || (device.user_id == self.config.user_id && device.device_id == self.config.device_id) {
2953 continue;
2954 }
2955 if unique.iter().any(|have: &StoredDevice| have.user_id == device.user_id && have.device_id == device.device_id)
2956 {
2957 continue;
2958 }
2959 unique.push(device.clone());
2960 }
2961 let missing = OlmSessionManager::sessions_missing_for(&self.store, &unique)?;
2962 let missing_ids: BTreeSet<(UserId, DeviceId)> =
2963 missing.into_iter().map(|device| (device.user_id.clone(), device.device_id.clone())).collect();
2964 let ready: Vec<StoredDevice> = unique
2965 .into_iter()
2966 .filter(|device| !missing_ids.contains(&(device.user_id.clone(), device.device_id.clone())))
2967 .collect();
2968 let mut sent = 0;
2969 for chunk in GroupSessionManager::chunk_recipients_for_send_to_device(&ready) {
2970 let body = GroupSessionManager::build_encrypted_to_device_body(
2971 &mut self.store,
2972 &self.account,
2973 &self.config.user_id,
2974 &self.config.device_id,
2975 chunk,
2976 inner_type,
2977 content.clone(),
2978 )?;
2979 let request_id = self.next_request_id()?;
2980 let share_txn = self.next_txn_id()?;
2981 let request = OutgoingRequest::send_to_device(request_id, ROOM_KEY_SHARE_EVENT_TYPE, &share_txn, body);
2982 self.enqueue_request(request, Lane::ToDevice)?;
2983 sent += chunk.len();
2984 }
2985 Ok(sent)
2986 }
2987
2988 fn queue_unknown_sender_retry(&mut self, event: ToDeviceEvent) {
2989 self.unknown_sender_retry.push_back(event);
2990 while self.unknown_sender_retry.len() > UNKNOWN_SENDER_RETRY_CAPACITY {
2991 self.unknown_sender_retry.pop_front();
2992 }
2993 }
2994
2995 fn retry_unknown_sender_queue(&mut self) -> Result<(), MessengerError> {
2996 let pending: Vec<ToDeviceEvent> = std::mem::take(&mut self.unknown_sender_retry).into_iter().collect();
2997 let mut newly_received_sessions = Vec::new();
2998 for event in pending {
2999 self.process_to_device_event(&event, &mut newly_received_sessions)?;
3000 }
3001 for (room_id, session_id) in &newly_received_sessions {
3002 self.retry_pending_for_session(room_id, session_id);
3003 }
3004 Ok(())
3005 }
3006
3007 fn outdated_tracked_users(&self) -> Result<BTreeSet<UserId>, MessengerError> {
3009 Ok(self.store.tracked_users()?.into_iter().filter_map(|(user_id, outdated)| outdated.then_some(user_id)).collect())
3010 }
3011
3012 fn enqueue_keys_query(&mut self) -> Result<(), MessengerError> {
3014 let request_id = self.next_request_id()?;
3015 if let Some(request) = DeviceTracker::keys_query_request(&self.store, request_id)? {
3016 self.enqueue_request(request, Lane::Other)?;
3017 }
3018 Ok(())
3019 }
3020
3021 fn retry_missing_room_keys(&mut self) {
3031 let work: Vec<(RoomId, EventId)> = self
3032 .pending_decrypt
3033 .iter()
3034 .flat_map(|(room, by_event)| by_event.keys().map(move |event| (room.clone(), event.clone())))
3035 .collect();
3036 for (room_id, event_id) in work {
3037 let _ = self.retry_decrypt_one(&room_id, &event_id);
3038 }
3039 }
3040
3041 fn history_forwarding_allowed(&self, room_id: &RoomId) -> bool {
3045 !matches!(
3046 self.rooms.get(room_id).and_then(|room| room.history_visibility.clone()),
3047 Some(crate::wire::events::HistoryVisibility::Joined) | Some(crate::wire::events::HistoryVisibility::Invited)
3048 )
3049 }
3050
3051 fn maybe_enqueue_cross_signing(&mut self) -> Result<(), MessengerError> {
3054 if self.xsign_inflight || self.xsign_attempts >= 4 {
3055 return Ok(());
3056 }
3057 if self.xsign.is_none() {
3058 self.xsign = Some(CrossSigningState::load_or_create(&mut self.store)?);
3059 }
3060 let Some(state) = self.xsign.clone() else { return Ok(()) };
3061 let user_id = self.config.user_id.clone();
3062 let request_id;
3063 let request = if !state.uploaded {
3064 request_id = self.next_request_id()?;
3065 OutgoingRequest::signing_keys_upload(request_id, state.upload_body(&user_id)?)
3066 } else if !state.device_signed {
3067 let device_id = self.config.device_id.clone();
3068 let device_keys = self.account.device_keys_json(&user_id, &device_id)?;
3069 request_id = self.next_request_id()?;
3070 OutgoingRequest::signatures_upload(request_id, state.device_signature_body(&user_id, &device_id, device_keys)?)
3071 } else {
3072 return Ok(());
3073 };
3074 self.xsign_attempts += 1;
3075 self.xsign_inflight = true;
3076 self.enqueue_request(request, Lane::Other)
3077 }
3078
3079 fn track_joined_members<'a>(&mut self, room_id: &RoomId, events: impl Iterator<Item = &'a RawEvent>) -> Result<(), MessengerError> {
3080 if self.rooms.get(room_id).is_none_or(|room| room.encryption.is_none()) {
3081 return Ok(());
3082 }
3083 let joined: Vec<UserId> = events
3084 .filter(|event| event.event_type == "m.room.member")
3085 .filter(|event| event.content.get("membership").and_then(serde_json::Value::as_str) == Some("join"))
3086 .filter_map(|event| event.state_key.as_deref().and_then(|state_key| UserId::parse(state_key).ok()))
3087 .filter(|user_id| *user_id != self.config.user_id)
3088 .collect();
3089 if joined.is_empty() {
3090 return Ok(());
3091 }
3092 DeviceTracker::on_device_lists(&mut self.store, &joined, &[])
3093 }
3094
3095 fn ingest_joined_room(&mut self, room_id: &RoomId, joined: &JoinedRoom) -> Result<(), MessengerError> {
3096 for event in joined.state.events.iter().chain(joined.timeline.events.iter()) {
3097 self.rooms.entry(room_id.clone()).or_default().apply_state_event(event)?;
3098 }
3099
3100 self.track_joined_members(room_id, joined.state.events.iter().chain(joined.timeline.events.iter()))?;
3101
3102 let unread_before = self.rooms.get(room_id).map(RoomState::unread_count);
3103 self.rooms.entry(room_id.clone()).or_default().apply_summary(&joined.summary, &joined.unread_notifications);
3104 let unread_after = self.rooms.get(room_id).map(RoomState::unread_count);
3105 if unread_before != unread_after {
3106 self.emit(MessengerEvent::UnreadChanged { room_id: room_id.clone() });
3107 }
3108
3109 for event in &joined.account_data.events {
3110 self.apply_synced_account_data(Some(room_id), &event.event_type, event.content.clone());
3111 }
3112
3113 if let Some(typing_event) = joined.ephemeral.events.iter().find(|event| event.event_type == "m.typing") {
3114 if let Ok(typing) = serde_json::from_value::<TypingContent>(typing_event.content.clone()) {
3115 self.typing.insert(room_id.clone(), typing.user_ids);
3116 self.emit(MessengerEvent::TypingChanged { room_id: room_id.clone() });
3117 }
3118 }
3119 if let Some(receipt_event) = joined.ephemeral.events.iter().find(|event| event.event_type == "m.receipt") {
3120 if let Ok(receipts) = serde_json::from_value::<ReceiptContent>(receipt_event.content.clone()) {
3121 self.receipts.insert(room_id.clone(), receipts);
3122 self.emit(MessengerEvent::ReceiptsChanged { room_id: room_id.clone() });
3123 }
3124 }
3125
3126 if !joined.timeline.events.is_empty() {
3127 self.timelines.entry(room_id.clone()).or_default().apply_timeline_batch(
3128 &joined.timeline.events,
3129 joined.timeline.limited,
3130 joined.timeline.prev_batch.clone(),
3131 );
3132 self.decrypt_new_timeline_events(room_id, &joined.timeline.events);
3133 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
3134 }
3135
3136 self.persist_room_state(room_id)
3137 }
3138
3139 fn ingest_invited_room(&mut self, room_id: &RoomId, invited: &InvitedRoom) -> Result<(), MessengerError> {
3140 for stripped in &invited.invite_state.events {
3141 let synthetic = stripped_state_to_raw_event(stripped);
3142 self.rooms.entry(room_id.clone()).or_default().apply_state_event(&synthetic)?;
3143 }
3144 self.persist_room_state(room_id)
3145 }
3146
3147 fn ingest_left_room(&mut self, room_id: &RoomId, left: &LeftRoom) -> Result<(), MessengerError> {
3148 for event in left.state.events.iter().chain(left.timeline.events.iter()) {
3149 self.rooms.entry(room_id.clone()).or_default().apply_state_event(event)?;
3150 }
3151 for event in &left.account_data.events {
3152 self.apply_synced_account_data(Some(room_id), &event.event_type, event.content.clone());
3153 }
3154 if !left.timeline.events.is_empty() {
3155 self.timelines.entry(room_id.clone()).or_default().apply_timeline_batch(
3156 &left.timeline.events,
3157 left.timeline.limited,
3158 left.timeline.prev_batch.clone(),
3159 );
3160 self.decrypt_new_timeline_events(room_id, &left.timeline.events);
3161 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
3162 }
3163 self.persist_room_state(room_id)
3164 }
3165
3166 fn persist_room_state(&mut self, room_id: &RoomId) -> Result<(), MessengerError> {
3167 let Some(state) = self.rooms.get(room_id) else { return Ok(()) };
3168 let bytes = serde_json::to_vec(state)
3169 .map_err(|source| MessengerError::Crypto(format!("encode room state for {room_id}: {source}")))?;
3170 self.store.save_room_state(room_id, bytes)?;
3171 Ok(())
3172 }
3173
3174 fn decrypt_new_timeline_events(&mut self, room_id: &RoomId, events: &[RawEvent]) {
3182 for event in events {
3183 if event.event_type != "m.room.encrypted" || event.state_key.is_some() || event.unsigned.redacted_because.is_some() {
3184 continue;
3185 }
3186 let Ok(RoomEncryptedContent::Megolm(content)) = serde_json::from_value::<RoomEncryptedContent>(event.content.clone())
3187 else {
3188 continue;
3189 };
3190 match GroupSessionManager::decrypt_event(
3191 &mut self.store,
3192 room_id,
3193 &event.event_id,
3194 event.origin_server_ts,
3195 &event.sender,
3196 &content,
3197 ) {
3198 Ok(plaintext) => {
3199 self.apply_decrypted_plaintext(room_id, &event.event_id, &event.sender, event.origin_server_ts, &plaintext);
3200 if let Some(by_event) = self.pending_decrypt.get_mut(room_id) {
3201 by_event.remove(&event.event_id);
3202 }
3203 }
3204 Err(err) => {
3205 if matches!(err, GroupDecryptError::MissingSession { .. }) {
3206 let _ = self.request_missing_room_key(room_id, &event.sender, &content);
3207 }
3208 let reason = format!("{:?}", err.utd_reason());
3217 if let Some(timeline) = self.timelines.get_mut(room_id) {
3218 timeline.set_decrypted(&event.event_id, ItemContent::Undecryptable { reason });
3219 }
3220 self.pending_decrypt.entry(room_id.clone()).or_default().insert(
3221 event.event_id.clone(),
3222 PendingDecryptItem {
3223 session_id: content.session_id.clone(),
3224 sender: event.sender.clone(),
3225 origin_server_ts: event.origin_server_ts,
3226 content,
3227 },
3228 );
3229 }
3230 }
3231 }
3232 }
3233
3234 fn ingest_room_messages_response(&mut self, room_id: &RoomId, body: &[u8]) -> Result<(), MessengerError> {
3235 let parsed: RoomMessagesResponseBody = serde_json::from_slice(body)?;
3236 let mut events = parsed.chunk;
3237 events.reverse();
3238 self.timelines.entry(room_id.clone()).or_default().prepend_back_page(&events, parsed.end);
3239 self.decrypt_new_timeline_events(room_id, &events);
3240 self.emit(MessengerEvent::TimelineChanged { room_id: room_id.clone() });
3241 self.persist_room_state(room_id)
3242 }
3243}
3244
3245impl MessengerCore<SealedRecordCodec> {
3246 pub fn open_sealed(
3257 records: impl IntoIterator<Item = SealedRecord>,
3258 secrets: CoreSecrets,
3259 config: CoreConfig,
3260 now_ms: i64,
3261 jitter: Box<dyn Jitter>,
3262 ) -> Result<Self, MessengerError> {
3263 let seal_key = secrets.store_seal_key.clone().ok_or_else(|| {
3264 MessengerError::Crypto(
3265 "MessengerCore::open_sealed requires CoreSecrets::store_seal_key"
3266 .to_string(),
3267 )
3268 })?;
3269 let codec = SealedRecordCodec::new(seal_key);
3270 Self::open(records, codec, config, secrets, now_ms, jitter)
3271 }
3272}
3273
3274fn stripped_placeholder_event_id() -> EventId {
3280 EventId::parse("$stripped-preview").expect("this literal is a well-formed, hardcoded $-prefixed opaque id")
3281}
3282
3283fn stripped_state_to_raw_event(stripped: &StrippedStateEvent) -> RawEvent {
3288 RawEvent {
3289 event_id: stripped_placeholder_event_id(),
3290 event_type: stripped.event_type.clone(),
3291 sender: stripped.sender.clone(),
3292 origin_server_ts: 0,
3293 state_key: Some(stripped.state_key.clone()),
3294 content: stripped.content.clone(),
3295 unsigned: Unsigned::default(),
3296 }
3297}
3298
3299fn forwardable_content(mut content: serde_json::Value, marker_key: &str, marker: serde_json::Value) -> serde_json::Value {
3306 if let Some(object) = content.as_object_mut() {
3307 for key in ["m.relates_to", "m.mentions", "m.new_content", KEY_FORWARDED, KEY_FORWARDED_FROM] {
3308 object.remove(key);
3309 }
3310 object.insert(marker_key.to_string(), marker);
3311 }
3312 content
3313}
3314
3315fn message_local_echo_content(message: &OutgoingMessage) -> ItemContent {
3322 let body = TextLikeMessageContent { body: message.body.clone(), format: None, formatted_body: None };
3323 match message.kind {
3324 MessageKind::Text => ItemContent::Text(body),
3325 MessageKind::Notice => ItemContent::Notice(body),
3326 MessageKind::Emote => ItemContent::Emote(body),
3327 }
3328}
3329
3330fn message_wire_content(message: &OutgoingMessage) -> serde_json::Value {
3339 let msgtype = message.kind.msgtype();
3340 let new_content = serde_json::json!({ "msgtype": msgtype, "body": message.body });
3341 let mut content = new_content.clone();
3342 if let Some(edit_of) = &message.edit_of {
3343 content["body"] = serde_json::Value::String(format!("* {}", message.body));
3344 content["m.new_content"] = new_content;
3345 content["m.relates_to"] = serde_json::json!({ "rel_type": "m.replace", "event_id": edit_of.as_str() });
3346 } else if let Some(reply_to) = &message.reply_to {
3347 content["m.relates_to"] = serde_json::json!({ "m.in_reply_to": { "event_id": reply_to.as_str() } });
3348 }
3349 content
3350}
3351
3352fn message_relates_to(message: &OutgoingMessage) -> Option<RelatesTo> {
3358 if let Some(edit_of) = &message.edit_of {
3359 Some(RelatesTo::Replace { event_id: edit_of.clone() })
3360 } else {
3361 message.reply_to.clone().map(|event_id| RelatesTo::InReplyTo(InReplyTo { event_id }))
3362 }
3363}
3364
3365#[cfg(test)]
3366mod tests {
3367 use super::*;
3368 use crate::room::timeline::Forwarded;
3369 use crate::store::InsecurePlainCodecForTests;
3370 use crate::wire::HttpMethod;
3371
3372 struct FixedJitter(f64);
3373
3374 impl Jitter for FixedJitter {
3375 fn next_unit(&mut self) -> f64 {
3376 self.0
3377 }
3378 }
3379
3380 fn device_id() -> DeviceId {
3381 DeviceId::parse("DEV1").expect("valid device id")
3382 }
3383
3384 fn user_id() -> UserId {
3385 UserId::parse("@alice:example.org").expect("valid user id")
3386 }
3387
3388 fn test_config() -> CoreConfig {
3389 CoreConfig { user_id: user_id(), device_id: device_id(), server_name: "example.org".to_string() }
3390 }
3391
3392 fn open_fresh_core() -> MessengerCore<InsecurePlainCodecForTests> {
3393 MessengerCore::open(Vec::new(), InsecurePlainCodecForTests, test_config(), CoreSecrets::default(), 0, Box::new(FixedJitter(0.0)))
3394 .expect("open succeeds on a brand-new device")
3395 }
3396
3397 fn flush_and_ack(core: &mut MessengerCore<InsecurePlainCodecForTests>) {
3405 while let Some(batch) = core.take_flush_batch() {
3406 core.ack_flush(batch.id);
3407 }
3408 }
3409
3410 #[test]
3411 fn counters_survive_restart_and_never_repeat_txn_ids() {
3412 let device = device_id();
3413 let mut store = Store::new(device.clone(), InsecurePlainCodecForTests);
3414 store.save_counters(Counters { next_request_id: 3, next_txn_id: 7 }).expect("save");
3415 let batch = store.take_flush_batch().expect("counters dirtied the store");
3416 store.ack_flush(batch.id);
3417
3418 let reloaded = Store::load(batch.records, InsecurePlainCodecForTests, device).expect("reload");
3419 let restored = reloaded.counters().expect("no error").expect("counters were persisted");
3420 assert_eq!(restored.next_txn_id, 7, "a restart resumes the txn-id sequence exactly where it left off");
3421 assert_eq!(restored.next_request_id, 3);
3422
3423 assert_eq!(restore_counter(3, Some(10)), 11, "resumes past the highest pending id, never repeating it");
3428 assert_eq!(restore_counter(3, None), 3, "nothing pending: the persisted value alone is authoritative");
3429 }
3430
3431 #[test]
3432 fn open_restores_the_request_counter_past_a_still_pending_requests_id() {
3433 let device = device_id();
3434 let mut store = Store::new(device.clone(), InsecurePlainCodecForTests);
3435 let mut queue = OutgoingQueue::new();
3436 let pending = OutgoingRequest::sync(RequestId::next(5), None, None);
3437 queue.enqueue(&mut store, pending, Lane::Sync).expect("enqueue");
3438 store.save_counters(Counters { next_request_id: 2, next_txn_id: 0 }).expect("save");
3441 let batch = store.take_flush_batch().expect("pending");
3442 store.ack_flush(batch.id);
3443
3444 let core = MessengerCore::open(
3445 batch.records,
3446 InsecurePlainCodecForTests,
3447 test_config(),
3448 CoreSecrets::default(),
3449 0,
3450 Box::new(FixedJitter(0.0)),
3451 )
3452 .expect("open succeeds");
3453 assert_eq!(
3454 core.counters.next_request_id, 6,
3455 "resumes past the highest id already in flight, not the stale persisted value"
3456 );
3457 }
3458
3459 #[test]
3460 fn sync_token_persisted_after_everything_it_covers() {
3461 let mut core = open_fresh_core();
3462 flush_and_ack(&mut core);
3463 core.releasable_requests(0);
3468 flush_and_ack(&mut core);
3469 let requests = core.releasable_requests(0);
3470 let sync_request = requests.iter().find(|r| r.kind == OutgoingRequestKind::Sync).expect("a sync request was enqueued");
3471
3472 assert!(core.store.sync_token().expect("no error").is_none(), "no token before the first response");
3473
3474 let body = serde_json::json!({
3475 "next_batch": "s1",
3476 "device_one_time_keys_count": { "signed_curve25519": 0 },
3477 "device_unused_fallback_key_types": []
3478 })
3479 .to_string();
3480 core.on_response(sync_request.id.clone(), HttpResponseDescriptor { status: 200, body: body.into_bytes() }, 0);
3481
3482 assert_eq!(core.store.sync_token().expect("no error"), Some("s1"));
3483 }
3484
3485 #[test]
3486 fn only_one_sync_in_flight() {
3487 let mut core = open_fresh_core();
3488 flush_and_ack(&mut core);
3489 core.releasable_requests(0);
3490 flush_and_ack(&mut core);
3491 let first = core.releasable_requests(0);
3492 assert_eq!(first.iter().filter(|r| r.kind == OutgoingRequestKind::Sync).count(), 1);
3493
3494 let second = core.releasable_requests(0);
3498 assert!(second.iter().all(|r| r.kind != OutgoingRequestKind::Sync), "no second sync request released");
3499 }
3500
3501 #[test]
3502 fn first_sync_of_a_session_never_long_polls_then_the_loop_does() {
3503 let mut core = open_fresh_core();
3504 flush_and_ack(&mut core);
3505 core.releasable_requests(0);
3506 flush_and_ack(&mut core);
3507 let first = core.releasable_requests(0);
3508 let first_sync = first.iter().find(|r| r.kind == OutgoingRequestKind::Sync).expect("a sync request was enqueued");
3509 assert!(first_sync.query.contains(&("timeout".to_string(), "0".to_string())), "the initial sync must answer at once");
3510 assert!(first_sync.query.iter().all(|(key, _)| key != "since"), "a fresh device has no since token");
3511
3512 let body = serde_json::json!({
3513 "next_batch": "s1",
3514 "device_one_time_keys_count": { "signed_curve25519": 0 },
3515 "device_unused_fallback_key_types": []
3516 })
3517 .to_string();
3518 core.on_response(first_sync.id.clone(), HttpResponseDescriptor { status: 200, body: body.into_bytes() }, 0);
3519 flush_and_ack(&mut core);
3520
3521 core.releasable_requests(0);
3524 flush_and_ack(&mut core);
3525 let second = core.releasable_requests(0);
3526 let second_sync = second.iter().find(|r| r.kind == OutgoingRequestKind::Sync).expect("the follow-up sync was enqueued");
3527 assert!(second_sync.query.contains(&("since".to_string(), "s1".to_string())));
3528 assert!(second_sync.query.contains(&("timeout".to_string(), SYNC_TIMEOUT_MS.to_string())), "once caught up the loop long-polls");
3529 }
3530
3531 #[test]
3532 fn room_key_arrival_retries_undecryptable_items() {
3533 use crate::crypto::group_sessions::GroupSessionManager;
3534 use mail4agent_vodozemac::megolm::{GroupSession, SessionConfig};
3535 use mail4agent_vodozemac::olm::Account;
3536
3537 let mut core = open_fresh_core();
3538 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
3539 let sender = UserId::parse("@bob:example.org").expect("valid user id");
3540 let identity = Account::new().identity_keys();
3541
3542 let mut group_session = GroupSession::new(SessionConfig::version_1());
3543 let session_id = group_session.session_id();
3544 let session_key = group_session.session_key();
3549 let plaintext = crate::crypto::group_sessions::RoomEventPlaintext {
3550 event_type: "m.room.message".to_string(),
3551 content: serde_json::json!({ "msgtype": "m.text", "body": "hi" }),
3552 room_id: room_id.clone(),
3553 };
3554 let message = group_session.encrypt(serde_json::to_vec(&plaintext).expect("valid JSON"));
3555
3556 let event: RawEvent = serde_json::from_value(serde_json::json!({
3557 "event_id": "$evt:example.org",
3558 "type": "m.room.encrypted",
3559 "sender": sender.as_str(),
3560 "origin_server_ts": 1,
3561 "content": {
3562 "algorithm": "m.megolm.v1.aes-sha2",
3563 "ciphertext": message.to_base64(),
3564 "session_id": session_id,
3565 }
3566 }))
3567 .expect("valid event");
3568
3569 let mut timeline = Timeline::new();
3575 timeline.apply_timeline_batch(std::slice::from_ref(&event), false, None);
3576 core.timelines.insert(room_id.clone(), timeline);
3577 core.decrypt_new_timeline_events(&room_id, std::slice::from_ref(&event));
3578 {
3579 let timeline = core.timeline(&room_id).expect("timeline exists");
3580 let item = timeline.item_by_event_id(&EventId::parse("$evt:example.org").expect("valid event id")).expect("item present");
3581 assert!(matches!(item.content, ItemContent::Undecryptable { .. }));
3582 }
3583
3584 let content = RoomKeyContent {
3588 algorithm: "m.megolm.v1.aes-sha2".to_string(),
3589 room_id: room_id.clone(),
3590 session_id: session_id.clone(),
3591 session_key: session_key.to_base64(),
3592 };
3593 GroupSessionManager::receive_room_key(&mut core.store, &sender, identity.curve25519, identity.ed25519, &content)
3594 .expect("receive room key");
3595 core.retry_pending_for_session(&room_id, &session_id);
3596
3597 let timeline = core.timeline(&room_id).expect("timeline exists");
3598 let item = timeline.item_by_event_id(&EventId::parse("$evt:example.org").expect("valid event id")).expect("item present");
3599 match &item.content {
3600 ItemContent::Text(text) => assert_eq!(text.body, "hi"),
3601 other => panic!("expected the retried decrypt to succeed, got {other:?}"),
3602 }
3603 }
3604
3605 #[test]
3606 fn redelivered_to_device_after_restart_is_dropped_silently() {
3607 use crate::crypto::olm_sessions::OlmSessionManager;
3608
3609 let alice_user = UserId::parse("@alice:example.org").expect("valid user id");
3611 let alice_device = DeviceId::parse("ALICEDEV").expect("valid device id");
3612 let mut alice_store = Store::new(alice_device.clone(), InsecurePlainCodecForTests);
3613 let alice_account = OlmAccountState::load_or_create(&mut alice_store).expect("create alice's account");
3614 let alice_device_keys =
3617 alice_account.device_keys_json(&alice_user, &alice_device).expect("sign alice's device_keys");
3618
3619 let mut bob_store = Store::new(device_id(), InsecurePlainCodecForTests);
3621 let mut bob_account = OlmAccountState::load_or_create(&mut bob_store).expect("create bob's account");
3622 bob_account.on_sync_counts(&mut bob_store, 0, &["signed_curve25519".to_string()]).expect("top up bob's OTKs");
3623 let upload = bob_account
3624 .keys_upload_request(RequestId::next(0), &user_id(), &device_id())
3625 .expect("builds a request")
3626 .expect("something to upload");
3627 let one_time_keys = upload.body.expect("upload has a body");
3628 let one_time_keys = one_time_keys.get("one_time_keys").and_then(serde_json::Value::as_object).expect("OTKs present");
3629 let (key_id, key_value) = one_time_keys.iter().next().expect("at least one OTK");
3630 let claim_body = serde_json::to_vec(&serde_json::json!({
3631 "one_time_keys": { user_id().as_str(): { device_id().as_str(): { key_id: key_value } } }
3632 }))
3633 .expect("valid JSON");
3634 let bob_device_stored = crate::crypto::device_tracker::StoredDevice {
3635 user_id: user_id(),
3636 device_id: device_id(),
3637 curve25519: bob_account.identity_keys().curve25519,
3638 ed25519: bob_account.identity_keys().ed25519,
3639 algorithms: vec!["m.olm.v1.curve25519-aes-sha2".to_string(), "m.megolm.v1.aes-sha2".to_string()],
3640 display_name: None,
3641 verified: false,
3642 blocked: false,
3643 };
3644 OlmSessionManager::on_keys_claim_response(&mut alice_store, &alice_account, &[&bob_device_stored], &claim_body)
3645 .expect("alice establishes a session with bob");
3646
3647 DeviceTracker::on_keys_query_response(
3653 &mut bob_store,
3654 &serde_json::to_vec(&serde_json::json!({
3655 "device_keys": { alice_user.as_str(): { alice_device.as_str(): alice_device_keys } }
3656 }))
3657 .expect("valid JSON"),
3658 )
3659 .expect("bob learns alice's device");
3660
3661 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
3662 let group_session = mail4agent_vodozemac::megolm::GroupSession::new(mail4agent_vodozemac::megolm::SessionConfig::version_1());
3663 let room_key_content = RoomKeyContent {
3664 algorithm: "m.megolm.v1.aes-sha2".to_string(),
3665 room_id,
3666 session_id: group_session.session_id(),
3667 session_key: group_session.session_key().to_base64(),
3668 };
3669 let encrypted = OlmSessionManager::encrypt_to_device(
3670 &mut alice_store,
3671 &alice_account,
3672 &alice_user,
3673 &alice_device,
3674 &bob_device_stored,
3675 "m.room_key",
3676 serde_json::to_value(&room_key_content).expect("valid JSON"),
3677 )
3678 .expect("encrypt to bob");
3679 let to_device_event = ToDeviceEvent {
3680 sender: alice_user.clone(),
3681 event_type: "m.room.encrypted".to_string(),
3682 content: serde_json::to_value(RoomEncryptedContent::Olm(encrypted)).expect("valid JSON"),
3683 };
3684
3685 let mut core =
3686 MessengerCore::open(Vec::new(), InsecurePlainCodecForTests, test_config(), CoreSecrets::default(), 0, Box::new(FixedJitter(0.0)))
3687 .expect("open succeeds");
3688 core.store = bob_store;
3689 core.account = bob_account;
3690
3691 let mut sessions = Vec::new();
3692 core.process_to_device_event(&to_device_event, &mut sessions).expect("first decrypt succeeds");
3693 assert_eq!(
3694 sessions,
3695 vec![(room_key_content.room_id.clone(), room_key_content.session_id.clone())],
3696 "the first, genuine delivery is fully validated and accepted"
3697 );
3698
3699 let batch = core.store.take_flush_batch().expect("decrypting dirtied the store");
3703 core.store.ack_flush(batch.id);
3704 let reloaded_store = Store::load(batch.records, InsecurePlainCodecForTests, device_id()).expect("reload succeeds");
3705 core.store = reloaded_store;
3706 core.account = OlmAccountState::load_or_create(&mut core.store).expect("reload reuses the persisted account");
3707
3708 let mut sessions_after_restart = Vec::new();
3712 let result = core.process_to_device_event(&to_device_event, &mut sessions_after_restart);
3713 assert!(result.is_ok(), "a redelivered to-device event is dropped, not surfaced as an ingestion error");
3714 assert!(sessions_after_restart.is_empty(), "the replay never produces a second accepted room key");
3715 }
3716
3717 #[test]
3718 fn unknown_sender_device_is_retried_after_keys_query() {
3719 use crate::crypto::olm_sessions::OlmSessionManager;
3720
3721 let alice_user = UserId::parse("@alice:example.org").expect("valid user id");
3723 let alice_device = DeviceId::parse("ALICEDEV").expect("valid device id");
3724 let mut alice_store = Store::new(alice_device.clone(), InsecurePlainCodecForTests);
3725 let alice_account = OlmAccountState::load_or_create(&mut alice_store).expect("create alice's account");
3726 let alice_device_keys =
3734 alice_account.device_keys_json(&alice_user, &alice_device).expect("sign alice's device_keys");
3735
3736 let bob_user = UserId::parse("@bob:example.org").expect("valid user id");
3744 let bob_device = DeviceId::parse("BOBDEV").expect("valid device id");
3745 let mut bob_store = Store::new(bob_device.clone(), InsecurePlainCodecForTests);
3746 let mut bob_account = OlmAccountState::load_or_create(&mut bob_store).expect("create bob's account");
3747 bob_account.on_sync_counts(&mut bob_store, 0, &["signed_curve25519".to_string()]).expect("top up bob's OTKs");
3748 let upload = bob_account
3749 .keys_upload_request(RequestId::next(0), &bob_user, &bob_device)
3750 .expect("builds a request")
3751 .expect("something to upload");
3752 let one_time_keys = upload.body.expect("upload has a body");
3753 let one_time_keys = one_time_keys.get("one_time_keys").and_then(serde_json::Value::as_object).expect("OTKs present");
3754 let (key_id, key_value) = one_time_keys.iter().next().expect("at least one OTK");
3755 let claim_body = serde_json::to_vec(&serde_json::json!({
3756 "one_time_keys": { bob_user.as_str(): { bob_device.as_str(): { key_id: key_value } } }
3757 }))
3758 .expect("valid JSON");
3759 let bob_device_stored = crate::crypto::device_tracker::StoredDevice {
3760 user_id: bob_user.clone(),
3761 device_id: bob_device.clone(),
3762 curve25519: bob_account.identity_keys().curve25519,
3763 ed25519: bob_account.identity_keys().ed25519,
3764 algorithms: vec!["m.olm.v1.curve25519-aes-sha2".to_string(), "m.megolm.v1.aes-sha2".to_string()],
3765 display_name: None,
3766 verified: false,
3767 blocked: false,
3768 };
3769 OlmSessionManager::on_keys_claim_response(&mut alice_store, &alice_account, &[&bob_device_stored], &claim_body)
3770 .expect("alice establishes a session with bob");
3771 let encrypted = OlmSessionManager::encrypt_to_device(
3772 &mut alice_store,
3773 &alice_account,
3774 &alice_user,
3775 &alice_device,
3776 &bob_device_stored,
3777 "m.dummy",
3778 serde_json::json!({}),
3779 )
3780 .expect("encrypt to bob");
3781 let to_device_event = ToDeviceEvent {
3782 sender: alice_user.clone(),
3783 event_type: "m.room.encrypted".to_string(),
3784 content: serde_json::to_value(RoomEncryptedContent::Olm(encrypted)).expect("valid JSON"),
3785 };
3786
3787 let bob_config = CoreConfig { user_id: bob_user.clone(), device_id: bob_device.clone(), server_name: "example.org".to_string() };
3788 let mut core =
3789 MessengerCore::open(Vec::new(), InsecurePlainCodecForTests, bob_config, CoreSecrets::default(), 0, Box::new(FixedJitter(0.0)))
3790 .expect("open succeeds");
3791 core.store = bob_store;
3792 core.account = bob_account;
3793
3794 let mut sessions = Vec::new();
3798 core.process_to_device_event(&to_device_event, &mut sessions).expect("queues for retry, does not error");
3799 assert_eq!(core.unknown_sender_retry.len(), 1, "queued for a retry after the next keys/query");
3800 let alice_curve = alice_account.identity_keys().curve25519.to_base64();
3801 assert!(
3802 core.store.olm_sessions_for_device(&alice_curve).expect("no error").is_empty(),
3803 "the refused attempt created no Olm session (it would have consumed the one-time key)"
3804 );
3805
3806 let query_body = serde_json::to_vec(&serde_json::json!({
3811 "device_keys": { alice_user.as_str(): { alice_device.as_str(): alice_device_keys } }
3812 }))
3813 .expect("valid JSON");
3814 DeviceTracker::on_keys_query_response(&mut core.store, &query_body).expect("bob learns alice's device");
3815 core.retry_unknown_sender_queue().expect("no error");
3816 assert!(core.unknown_sender_retry.is_empty(), "the queued entry is drained by the retry attempt");
3817 assert_eq!(
3818 core.store.olm_sessions_for_device(&alice_curve).expect("no error").len(),
3819 1,
3820 "the retry decrypted the very same pre-key message and established the inbound session"
3821 );
3822 }
3823
3824 #[test]
3825 fn dispatch_load_older_produces_an_outgoing_request() {
3826 let mut core = open_fresh_core();
3827 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
3828 let mut timeline = Timeline::new();
3829 let event: RawEvent = serde_json::from_value(serde_json::json!({
3830 "event_id": "$1:example.org",
3831 "type": "m.room.message",
3832 "sender": "@alice:example.org",
3833 "origin_server_ts": 1,
3834 "content": { "msgtype": "m.text", "body": "hi" }
3835 }))
3836 .expect("valid event");
3837 timeline.apply_timeline_batch(&[event], true, Some("t1".to_string()));
3838 core.timelines.insert(room_id.clone(), timeline);
3839 flush_and_ack(&mut core);
3840
3841 core.dispatch(MessengerCommand::LoadOlder { room_id: room_id.clone() }, 0).expect("dispatch succeeds");
3842 flush_and_ack(&mut core);
3843 let requests = core.releasable_requests(0);
3844 let load_older = requests.iter().find(|r| r.kind == OutgoingRequestKind::RoomMessages);
3845 assert!(load_older.is_some(), "a RoomMessages request was enqueued using the room's own gap token");
3846 }
3847
3848 #[test]
3849 fn load_older_on_a_resumed_room_pages_back_from_the_sync_token() {
3850 let mut core = open_fresh_core();
3851 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
3852 core.store.save_sync_token("s42_7".to_string()).expect("save the sync token");
3853 flush_and_ack(&mut core);
3854
3855 core.dispatch(MessengerCommand::LoadOlder { room_id: room_id.clone() }, 0).expect("dispatch succeeds");
3857 flush_and_ack(&mut core);
3858 assert!(core.releasable_requests(0).iter().all(|r| r.kind != OutgoingRequestKind::RoomMessages));
3859
3860 core.rooms.insert(room_id.clone(), RoomState::default());
3862 core.dispatch(MessengerCommand::LoadOlder { room_id: room_id.clone() }, 0).expect("dispatch succeeds");
3863 flush_and_ack(&mut core);
3864 let requests = core.releasable_requests(0);
3865 let load_older = requests.iter().find(|r| r.kind == OutgoingRequestKind::RoomMessages).expect("a RoomMessages request was enqueued");
3866 assert!(load_older.query.contains(&("from".to_string(), "s42_7".to_string())), "pages back from the sync token: {:?}", load_older.query);
3867 }
3868
3869 fn bob_joins_room_body(encrypted: bool) -> serde_json::Value {
3872 let mut state = Vec::new();
3873 if encrypted {
3874 state.push(serde_json::json!({
3875 "event_id": "$enc:example.org", "type": "m.room.encryption", "sender": "@alice:example.org",
3876 "origin_server_ts": 1, "state_key": "", "content": { "algorithm": "m.megolm.v1.aes-sha2" }
3877 }));
3878 }
3879 serde_json::json!({
3880 "next_batch": "s1",
3881 "rooms": { "join": { "!r:example.org": {
3882 "state": { "events": state },
3883 "timeline": { "events": [{
3884 "event_id": "$join:example.org", "type": "m.room.member", "sender": "@bob:example.org",
3885 "origin_server_ts": 2, "state_key": "@bob:example.org", "content": { "membership": "join" }
3886 }] }
3887 } } }
3888 })
3889 }
3890
3891 #[test]
3892 fn a_member_joining_an_encrypted_room_is_tracked_and_queried_right_away() {
3893 let mut core = open_fresh_core();
3894 let bob = UserId::parse("@bob:example.org").expect("valid user id");
3895 deliver_sync(&mut core, bob_joins_room_body(true));
3896 let tracked = core.store.tracked_users().expect("no error");
3897 assert!(tracked.iter().any(|(user_id, outdated)| *user_id == bob && *outdated), "bob is tracked with an outdated device list: {tracked:?}");
3898
3899 flush_and_ack(&mut core);
3900 let requests = core.releasable_requests(0);
3901 let query = requests.iter().find(|r| r.kind == OutgoingRequestKind::KeysQuery).expect("a /keys/query was enqueued for the new member");
3902 assert!(query.body.as_ref().is_some_and(|body| body["device_keys"].get(bob.as_str()).is_some()), "the query names bob: {:?}", query.body);
3903 }
3904
3905 #[test]
3906 fn a_member_joining_a_plaintext_room_is_not_tracked() {
3907 let mut core = open_fresh_core();
3908 deliver_sync(&mut core, bob_joins_room_body(false));
3909 assert!(core.store.tracked_users().expect("no error").is_empty());
3910 flush_and_ack(&mut core);
3911 assert!(core.releasable_requests(0).iter().all(|r| r.kind != OutgoingRequestKind::KeysQuery));
3912 }
3913
3914 #[test]
3915 fn dispatch_mark_read_and_set_typing_produce_outgoing_requests() {
3916 let mut core = open_fresh_core();
3917 flush_and_ack(&mut core);
3918 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
3919 let event_id = EventId::parse("$1:example.org").expect("valid event id");
3920
3921 core.dispatch(MessengerCommand::MarkRead { room_id: room_id.clone(), event_id }, 0).expect("dispatch succeeds");
3922 core.dispatch(MessengerCommand::SetTyping { room_id: room_id.clone(), typing: true }, 0).expect("dispatch succeeds");
3923 flush_and_ack(&mut core);
3924
3925 let requests = core.releasable_requests(0);
3926 assert!(requests.iter().any(|r| r.kind == OutgoingRequestKind::ReadMarkers));
3927 assert!(requests.iter().any(|r| r.kind == OutgoingRequestKind::Typing));
3928 }
3929
3930 #[test]
3931 fn snapshot_reflects_the_latest_applied_sync() {
3932 let mut core = open_fresh_core();
3933 flush_and_ack(&mut core);
3934 core.releasable_requests(0);
3935 flush_and_ack(&mut core);
3936 let requests = core.releasable_requests(0);
3937 let sync_request = requests.iter().find(|r| r.kind == OutgoingRequestKind::Sync).expect("sync enqueued");
3938
3939 let body = serde_json::json!({
3940 "next_batch": "s1",
3941 "rooms": {
3942 "join": {
3943 "!r:example.org": {
3944 "state": { "events": [] },
3945 "timeline": {
3946 "events": [
3947 {
3948 "event_id": "$1:example.org",
3949 "type": "m.room.message",
3950 "sender": "@alice:example.org",
3951 "origin_server_ts": 1,
3952 "content": { "msgtype": "m.text", "body": "hi" }
3953 }
3954 ]
3955 },
3956 "unread_notifications": { "notification_count": 1, "highlight_count": 0 }
3957 }
3958 }
3959 }
3960 })
3961 .to_string();
3962 let events = core.on_response(sync_request.id.clone(), HttpResponseDescriptor { status: 200, body: body.into_bytes() }, 0);
3963
3964 assert!(events.contains(&MessengerEvent::RoomsChanged));
3965 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
3966 assert!(events.contains(&MessengerEvent::TimelineChanged { room_id: room_id.clone() }));
3967 assert!(events.contains(&MessengerEvent::UnreadChanged { room_id: room_id.clone() }));
3968
3969 let timeline = core.timeline(&room_id).expect("timeline present");
3970 assert_eq!(timeline.items().len(), 1);
3971 assert_eq!(core.room_state(&room_id).expect("room present").unread_count(), 1);
3972 assert_eq!(core.change_counter(), events.len() as u64);
3973 assert!(core.events().is_empty());
3976 }
3977
3978 #[test]
3979 fn dispatch_set_account_data_updates_optimistically_and_sends_a_put() {
3980 let mut core = open_fresh_core();
3981 flush_and_ack(&mut core);
3982 let content = serde_json::json!({ "note": "keep" });
3983 core.dispatch(
3984 MessengerCommand::SetAccountData { event_type: "com.example.prefs".to_string(), content: content.clone() },
3985 0,
3986 )
3987 .expect("dispatch succeeds");
3988
3989 assert_eq!(
3990 core.global_account_data("com.example.prefs"),
3991 Some(&content),
3992 "SetAccountData updates the core's own cache before any round trip"
3993 );
3994
3995 flush_and_ack(&mut core);
3996 let requests = core.releasable_requests(0);
3997 let request = requests.iter().find(|r| r.kind == OutgoingRequestKind::AccountData).expect("an AccountData request was enqueued");
3998 assert_eq!(request.method, HttpMethod::Put);
3999 assert_eq!(request.path, "/_matrix/client/v3/user/%40alice%3Aexample.org/account_data/com.example.prefs");
4000 assert_eq!(request.body.as_ref(), Some(&content));
4001 }
4002
4003 #[test]
4004 fn dispatch_set_room_account_data_updates_optimistically_and_sends_a_put() {
4005 let mut core = open_fresh_core();
4006 flush_and_ack(&mut core);
4007 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
4008 let content = serde_json::json!({ "flag": true });
4009 core.dispatch(
4010 MessengerCommand::SetRoomAccountData {
4011 room_id: room_id.clone(),
4012 event_type: "com.example.flag".to_string(),
4013 content: content.clone(),
4014 },
4015 0,
4016 )
4017 .expect("dispatch succeeds");
4018
4019 assert_eq!(
4020 core.room_account_data(&room_id, "com.example.flag"),
4021 Some(&content),
4022 "SetRoomAccountData updates the core's own cache before any round trip"
4023 );
4024
4025 flush_and_ack(&mut core);
4026 let requests = core.releasable_requests(0);
4027 let request =
4028 requests.iter().find(|r| r.kind == OutgoingRequestKind::RoomAccountData).expect("a RoomAccountData request was enqueued");
4029 assert_eq!(request.method, HttpMethod::Put);
4030 assert_eq!(
4031 request.path,
4032 "/_matrix/client/v3/user/%40alice%3Aexample.org/rooms/%21r%3Aexample.org/account_data/com.example.flag"
4033 );
4034 assert_eq!(request.body.as_ref(), Some(&content));
4035 }
4036
4037 #[test]
4038 fn dispatch_search_public_rooms_sends_the_expected_request() {
4039 let mut core = open_fresh_core();
4040 flush_and_ack(&mut core);
4041 core.dispatch(MessengerCommand::SearchPublicRooms { term: "trading".to_string() }, 0).expect("dispatch succeeds");
4042 flush_and_ack(&mut core);
4043 let requests = core.releasable_requests(0);
4044 let request = requests.iter().find(|r| r.kind == OutgoingRequestKind::PublicRooms).expect("a PublicRooms request was enqueued");
4045 assert_eq!(request.method, HttpMethod::Post);
4046 assert_eq!(request.path, "/_matrix/client/v3/publicRooms");
4047 assert_eq!(request.body, Some(serde_json::json!({ "filter": { "generic_search_term": "trading" }, "limit": 50 })));
4048 }
4049
4050 #[test]
4051 fn dispatch_search_users_sends_the_expected_request() {
4052 let mut core = open_fresh_core();
4053 flush_and_ack(&mut core);
4054 core.dispatch(MessengerCommand::SearchUsers { term: "bob".to_string() }, 0).expect("dispatch succeeds");
4055 flush_and_ack(&mut core);
4056 let requests = core.releasable_requests(0);
4057 let request =
4058 requests.iter().find(|r| r.kind == OutgoingRequestKind::UserDirectorySearch).expect("a UserDirectorySearch request was enqueued");
4059 assert_eq!(request.method, HttpMethod::Post);
4060 assert_eq!(request.path, "/_matrix/client/v3/user_directory/search");
4061 assert_eq!(request.body, Some(serde_json::json!({ "search_term": "bob", "limit": 20 })));
4062 }
4063
4064 #[test]
4065 fn public_rooms_response_superseded_by_a_later_search_is_dropped() {
4066 let mut core = open_fresh_core();
4067 flush_and_ack(&mut core);
4068 core.dispatch(MessengerCommand::SearchPublicRooms { term: "first".to_string() }, 0).expect("dispatch succeeds");
4069 flush_and_ack(&mut core);
4070 let first_request = core
4071 .releasable_requests(0)
4072 .into_iter()
4073 .find(|r| r.kind == OutgoingRequestKind::PublicRooms)
4074 .expect("the first search's own request");
4075
4076 core.dispatch(MessengerCommand::SearchPublicRooms { term: "second".to_string() }, 0).expect("dispatch succeeds");
4079 flush_and_ack(&mut core);
4080 let second_request = core
4081 .releasable_requests(0)
4082 .into_iter()
4083 .find(|r| r.kind == OutgoingRequestKind::PublicRooms)
4084 .expect("the second search's own request");
4085
4086 let stale_body = serde_json::json!({ "chunk": [{ "room_id": "!stale:example.org", "num_joined_members": 1 }] });
4088 core.on_response(first_request.id, HttpResponseDescriptor { status: 200, body: serde_json::to_vec(&stale_body).expect("json") }, 0);
4089 assert!(core.public_rooms_result().is_empty(), "a response to a superseded search is dropped");
4090
4091 let fresh_body = serde_json::json!({ "chunk": [{ "room_id": "!fresh:example.org", "num_joined_members": 2 }] });
4093 core.on_response(second_request.id, HttpResponseDescriptor { status: 200, body: serde_json::to_vec(&fresh_body).expect("json") }, 0);
4094 assert_eq!(core.public_rooms_result().len(), 1);
4095 assert_eq!(core.public_rooms_result()[0].room_id, RoomId::parse("!fresh:example.org").expect("valid room id"));
4096 }
4097
4098 fn is_account_data_write(request: &OutgoingRequest) -> bool {
4099 matches!(request.kind, OutgoingRequestKind::AccountData | OutgoingRequestKind::RoomAccountData)
4100 }
4101
4102 #[test]
4103 fn pin_then_archive_back_to_back_keeps_both_tags() {
4104 let mut core = open_fresh_core();
4105 flush_and_ack(&mut core);
4106 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
4107 for tag in ["m.favourite", "u.example.archived"] {
4108 core.dispatch(MessengerCommand::SetTag { room_id: room_id.clone(), tag: tag.to_string(), order: None }, 0)
4109 .expect("dispatch succeeds");
4110 }
4111
4112 let cached = core.room_account_data(&room_id, "m.tag").expect("m.tag cached");
4114 assert!(cached["tags"].get("m.favourite").is_some());
4115 assert!(cached["tags"].get("u.example.archived").is_some());
4116
4117 flush_and_ack(&mut core);
4118 let first_release = core.releasable_requests(0);
4119 let first_writes: Vec<&OutgoingRequest> = first_release.iter().filter(|r| is_account_data_write(r)).collect();
4120 assert_eq!(first_writes.len(), 1, "the second PUT is not released while the first is in flight");
4121 let first_body = first_writes[0].body.clone().expect("a body");
4122 assert!(first_body["tags"].get("m.favourite").is_some());
4123 assert!(first_body["tags"].get("u.example.archived").is_none(), "the first PUT carries only what existed when it was built");
4124
4125 core.on_response(first_writes[0].id.clone(), HttpResponseDescriptor { status: 200, body: b"{}".to_vec() }, 0);
4126 flush_and_ack(&mut core);
4127 let second_release = core.releasable_requests(0);
4128 let second = second_release.iter().find(|r| is_account_data_write(r)).expect("the second PUT is released after the first completed");
4129 let second_body = second.body.clone().expect("a body");
4130 assert!(second_body["tags"].get("m.favourite").is_some(), "the later PUT still carries the earlier tag");
4131 assert!(second_body["tags"].get("u.example.archived").is_some());
4132 }
4133
4134 #[test]
4135 fn account_data_lane_is_fifo_single_flight() {
4136 let mut core = open_fresh_core();
4137 flush_and_ack(&mut core);
4138 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
4139 core.dispatch(MessengerCommand::SetAccountData { event_type: "org.t.one".to_string(), content: serde_json::json!({}) }, 0)
4140 .expect("dispatch succeeds");
4141 core.dispatch(
4142 MessengerCommand::SetRoomAccountData {
4143 room_id: room_id.clone(),
4144 event_type: "org.t.two".to_string(),
4145 content: serde_json::json!({}),
4146 },
4147 0,
4148 )
4149 .expect("dispatch succeeds");
4150 core.dispatch(MessengerCommand::SetAccountData { event_type: "org.t.three".to_string(), content: serde_json::json!({}) }, 0)
4151 .expect("dispatch succeeds");
4152 core.dispatch(MessengerCommand::RemoveTag { room_id: room_id.clone(), tag: "m.favourite".to_string() }, 0)
4153 .expect("dispatch succeeds");
4154 flush_and_ack(&mut core);
4155
4156 let mut released_paths = Vec::new();
4157 for _ in 0..4 {
4158 let release = core.releasable_requests(0);
4159 let writes: Vec<&OutgoingRequest> = release.iter().filter(|r| is_account_data_write(r)).collect();
4160 assert_eq!(writes.len(), 1, "exactly one account-data write in flight at a time");
4161 released_paths.push(writes[0].path.clone());
4162 core.on_response(writes[0].id.clone(), HttpResponseDescriptor { status: 200, body: b"{}".to_vec() }, 0);
4163 flush_and_ack(&mut core);
4164 }
4165 assert!(released_paths[0].ends_with("/account_data/org.t.one"));
4166 assert!(released_paths[1].ends_with("/account_data/org.t.two"));
4167 assert!(released_paths[2].ends_with("/account_data/org.t.three"));
4168 assert!(released_paths[3].ends_with("/account_data/m.tag"));
4169 assert!(core.releasable_requests(0).iter().all(|r| !is_account_data_write(r)), "the lane is drained");
4170 }
4171
4172 fn deliver_sync(core: &mut MessengerCore<InsecurePlainCodecForTests>, body: serde_json::Value) -> Vec<OutgoingRequest> {
4177 let mut others = Vec::new();
4178 for _ in 0..3 {
4179 flush_and_ack(core);
4180 let mut sync_request = None;
4181 for request in core.releasable_requests(0) {
4182 if request.kind == OutgoingRequestKind::Sync && sync_request.is_none() {
4183 sync_request = Some(request);
4184 } else {
4185 others.push(request);
4186 }
4187 }
4188 if let Some(sync_request) = sync_request {
4189 let body = serde_json::to_vec(&body).expect("valid JSON");
4190 core.on_response(sync_request.id, HttpResponseDescriptor { status: 200, body }, 0);
4191 return others;
4192 }
4193 }
4194 panic!("no sync request became releasable");
4195 }
4196
4197 fn tag_sync_body(room_id: &str, tags: serde_json::Value, next_batch: &str) -> serde_json::Value {
4198 let mut join = serde_json::Map::new();
4199 join.insert(
4200 room_id.to_string(),
4201 serde_json::json!({ "account_data": { "events": [{ "type": "m.tag", "content": { "tags": tags } }] } }),
4202 );
4203 serde_json::json!({ "next_batch": next_batch, "rooms": { "join": serde_json::Value::Object(join) } })
4204 }
4205
4206 fn ok_response() -> HttpResponseDescriptor {
4207 HttpResponseDescriptor { status: 200, body: b"{}".to_vec() }
4208 }
4209
4210 #[test]
4211 fn sync_echo_of_an_older_tag_write_does_not_clobber_a_newer_optimistic_value() {
4212 let mut core = open_fresh_core();
4213 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
4214 for tag in ["m.favourite", "u.example.archived"] {
4215 core.dispatch(MessengerCommand::SetTag { room_id: room_id.clone(), tag: tag.to_string(), order: None }, 0)
4216 .expect("dispatch succeeds");
4217 }
4218
4219 let released = deliver_sync(&mut core, tag_sync_body("!r:example.org", serde_json::json!({ "m.favourite": {} }), "s1"));
4221 let cached = core.room_account_data(&room_id, "m.tag").expect("m.tag cached");
4222 assert!(cached["tags"].get("u.example.archived").is_some(), "the older echo must not clobber the newer optimistic value");
4223 assert!(cached["tags"].get("m.favourite").is_some());
4224 let key = (Some(room_id.clone()), "m.tag".to_string());
4225 assert_eq!(
4226 core.account_data_guard.get(&key).expect("writes are outstanding").server_value,
4227 Some(serde_json::json!({ "tags": { "m.favourite": {} } })),
4228 "the echo is remembered as the server's value"
4229 );
4230
4231 let first = released.iter().find(|r| r.kind == OutgoingRequestKind::RoomAccountData).expect("first write released");
4234 core.on_response(first.id.clone(), ok_response(), 0);
4235 flush_and_ack(&mut core);
4236 let second = core
4237 .releasable_requests(0)
4238 .into_iter()
4239 .find(|r| r.kind == OutgoingRequestKind::RoomAccountData)
4240 .expect("second write released after the first");
4241 core.on_response(second.id, ok_response(), 0);
4242 assert!(core.account_data_guard.is_empty());
4243 let cached = core.room_account_data(&room_id, "m.tag").expect("m.tag cached");
4244 assert!(cached["tags"].get("m.favourite").is_some() && cached["tags"].get("u.example.archived").is_some());
4245
4246 deliver_sync(
4248 &mut core,
4249 tag_sync_body("!r:example.org", serde_json::json!({ "m.favourite": {}, "u.example.archived": {} }), "s2"),
4250 );
4251 let cached = core.room_account_data(&room_id, "m.tag").expect("m.tag cached");
4252 assert!(cached["tags"].get("u.example.archived").is_some());
4253 }
4254
4255 #[test]
4256 fn terminal_account_data_failure_restores_the_server_value() {
4257 let mut core = open_fresh_core();
4258 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
4259 deliver_sync(&mut core, tag_sync_body("!r:example.org", serde_json::json!({ "m.favourite": {} }), "s1"));
4260
4261 core.dispatch(MessengerCommand::SetTag { room_id: room_id.clone(), tag: "u.example.archived".to_string(), order: None }, 0)
4262 .expect("dispatch succeeds");
4263 assert!(core.room_account_data(&room_id, "m.tag").expect("cached")["tags"].get("u.example.archived").is_some());
4264 flush_and_ack(&mut core);
4265 let write = core
4266 .releasable_requests(0)
4267 .into_iter()
4268 .find(|r| r.kind == OutgoingRequestKind::RoomAccountData)
4269 .expect("the write is released");
4270
4271 let counter_before = core.change_counter();
4272 let forbidden =
4273 HttpResponseDescriptor { status: 403, body: br#"{"errcode":"M_FORBIDDEN","error":"nope"}"#.to_vec() };
4274 let events = core.on_response(write.id, forbidden, 0);
4275 assert!(core.account_data_guard.is_empty());
4276 let restored = core.room_account_data(&room_id, "m.tag").expect("the server's value is back");
4277 assert!(restored["tags"].get("u.example.archived").is_none(), "the failed write's optimistic value is gone");
4278 assert!(restored["tags"].get("m.favourite").is_some());
4279 assert!(events.contains(&MessengerEvent::RoomsChanged));
4280 assert!(core.change_counter() > counter_before, "the UI is told to re-render");
4281
4282 for tag in ["u.a", "u.b"] {
4286 core.dispatch(MessengerCommand::SetTag { room_id: room_id.clone(), tag: tag.to_string(), order: None }, 0)
4287 .expect("dispatch succeeds");
4288 }
4289 flush_and_ack(&mut core);
4290 let first = core
4291 .releasable_requests(0)
4292 .into_iter()
4293 .find(|r| r.kind == OutgoingRequestKind::RoomAccountData)
4294 .expect("first write released");
4295 let forbidden =
4296 HttpResponseDescriptor { status: 403, body: br#"{"errcode":"M_FORBIDDEN","error":"nope"}"#.to_vec() };
4297 core.on_response(first.id, forbidden, 0);
4298 let cached = core.room_account_data(&room_id, "m.tag").expect("cached");
4299 assert!(cached["tags"].get("u.a").is_some() && cached["tags"].get("u.b").is_some(), "cache untouched while a write is pending");
4300 flush_and_ack(&mut core);
4301 let second = core
4302 .releasable_requests(0)
4303 .into_iter()
4304 .find(|r| r.kind == OutgoingRequestKind::RoomAccountData)
4305 .expect("second write released");
4306 core.on_response(second.id, ok_response(), 0);
4307 let cached = core.room_account_data(&room_id, "m.tag").expect("cached");
4308 assert!(cached["tags"].get("u.b").is_some(), "the last write succeeded: its optimistic value stands");
4309 }
4310
4311 #[test]
4312 fn pending_account_data_writes_rebuild_their_guard_after_a_restart() {
4313 let mut core = open_fresh_core();
4314 let room_id = RoomId::parse("!r:example.org").expect("valid room id");
4315 core.dispatch(MessengerCommand::SetTag { room_id: room_id.clone(), tag: "m.favourite".to_string(), order: None }, 0)
4316 .expect("dispatch succeeds");
4317 core.dispatch(
4318 MessengerCommand::SetAccountData { event_type: "org.t.one".to_string(), content: serde_json::json!({ "n": 1 }) },
4319 0,
4320 )
4321 .expect("dispatch succeeds");
4322
4323 let mut records = Vec::new();
4325 while let Some(batch) = core.take_flush_batch() {
4326 records.extend(batch.records.clone());
4327 core.ack_flush(batch.id);
4328 }
4329 let reopened = MessengerCore::open(
4330 records,
4331 InsecurePlainCodecForTests,
4332 test_config(),
4333 CoreSecrets::default(),
4334 0,
4335 Box::new(FixedJitter(0.0)),
4336 )
4337 .expect("reopen succeeds");
4338
4339 assert_eq!(reopened.account_data_guard.len(), 2, "one guard per key with a pending write");
4340 assert!(reopened.account_data_guard.values().all(|guard| guard.pending == 1 && guard.server_value.is_none()));
4341 assert_eq!(reopened.global_account_data("org.t.one"), Some(&serde_json::json!({ "n": 1 })));
4342 assert!(reopened.room_account_data(&room_id, "m.tag").expect("cached")["tags"].get("m.favourite").is_some());
4343 }
4344
4345 type TestCore = MessengerCore<InsecurePlainCodecForTests>;
4350
4351 fn room(id: &str) -> RoomId {
4352 RoomId::parse(id).expect("valid room id")
4353 }
4354
4355 fn event(id: &str) -> EventId {
4356 EventId::parse(id).expect("valid event id")
4357 }
4358
4359 fn add_room(core: &mut TestCore, room_id: &RoomId, encrypted: bool, name: Option<&str>) {
4362 let mut state = RoomState::new();
4363 if encrypted {
4364 state.encryption = Some(crate::wire::events::RoomEncryptionContent {
4365 algorithm: "m.megolm.v1.aes-sha2".to_string(),
4366 rotation_period_ms: 604_800_000,
4367 rotation_period_msgs: 100,
4368 });
4369 }
4370 state.name = name.map(str::to_string);
4371 core.rooms.insert(room_id.clone(), state);
4372 core.timelines.entry(room_id.clone()).or_default();
4373 }
4374
4375 fn mark_as_channel(core: &mut TestCore, room_id: &RoomId) {
4378 let Some(state) = core.rooms.get_mut(room_id) else { return };
4379 state
4380 .apply_state_event(&state_event_raw(
4381 "m.room.join_rules",
4382 "",
4383 "@alice:example.org",
4384 serde_json::json!({ "join_rule": "public" }),
4385 ))
4386 .expect("join_rules");
4387 state
4388 .apply_state_event(&state_event_raw(
4389 "m.room.power_levels",
4390 "",
4391 "@alice:example.org",
4392 serde_json::json!({ "events_default": 50, "users_default": 0 }),
4393 ))
4394 .expect("power_levels");
4395 }
4396
4397 fn state_event_raw(event_type: &str, state_key: &str, sender: &str, content: serde_json::Value) -> RawEvent {
4398 serde_json::from_value(serde_json::json!({
4399 "event_id": format!("${event_type}:example.org"),
4400 "type": event_type,
4401 "sender": sender,
4402 "origin_server_ts": 1,
4403 "state_key": state_key,
4404 "content": content,
4405 }))
4406 .expect("valid raw event")
4407 }
4408
4409 fn raw_event(event_id: &str, event_type: &str, sender: &str, content: serde_json::Value) -> RawEvent {
4410 serde_json::from_value(serde_json::json!({
4411 "event_id": event_id,
4412 "type": event_type,
4413 "sender": sender,
4414 "origin_server_ts": 1,
4415 "content": content,
4416 }))
4417 .expect("valid raw event")
4418 }
4419
4420 fn seed_plain(core: &mut TestCore, room_id: &RoomId, event_id: &str, sender: &str, event_type: &str, content: serde_json::Value) {
4422 let raw = raw_event(event_id, event_type, sender, content);
4423 core.timelines.entry(room_id.clone()).or_default().apply_timeline_batch(&[raw], false, None);
4424 }
4425
4426 fn seed_sealed(core: &mut TestCore, room_id: &RoomId, event_id: &str, sender: &str) {
4428 let content = serde_json::json!({ "algorithm": "m.megolm.v1.aes-sha2", "ciphertext": "AAAA", "session_id": "s1" });
4429 seed_plain(core, room_id, event_id, sender, "m.room.encrypted", content);
4430 }
4431
4432 fn seed_decrypted(
4435 core: &mut TestCore,
4436 room_id: &RoomId,
4437 event_id_str: &str,
4438 sender: &str,
4439 event_type: &str,
4440 content: serde_json::Value,
4441 ) {
4442 seed_sealed(core, room_id, event_id_str, sender);
4443 if let Some(timeline) = core.timelines.get_mut(room_id) {
4444 timeline.set_decrypted_event(&event(event_id_str), event_type, &content);
4445 }
4446 }
4447
4448 fn sent_room_event(core: &mut TestCore) -> (String, serde_json::Value) {
4452 flush_and_ack(core);
4453 let request = core
4454 .releasable_requests(0)
4455 .into_iter()
4456 .find(|request| request.kind == OutgoingRequestKind::RoomSend)
4457 .expect("a room send is releasable");
4458 let confirmed = format!("$sent{}:example.org", core.counters.next_request_id);
4459 let body = serde_json::json!({ "event_id": confirmed }).to_string().into_bytes();
4460 core.on_response(request.id.clone(), HttpResponseDescriptor { status: 200, body }, 0);
4461 (request.path, request.body.expect("a room send has a body"))
4462 }
4463
4464 fn forward(core: &mut TestCore, from: &RoomId, event_id: &str, to: &RoomId) -> Result<(), MessengerError> {
4465 core.dispatch(
4466 MessengerCommand::Forward { from_room: from.clone(), event_id: event(event_id), to_room: to.clone(), txn_id: None },
4467 0,
4468 )
4469 }
4470
4471 fn last_item(core: &TestCore, room_id: &RoomId) -> crate::room::timeline::TimelineItem {
4472 core.timeline(room_id).and_then(|timeline| timeline.items().last()).expect("the room has an item").clone()
4473 }
4474
4475 #[test]
4476 fn forward_from_encrypted_room_carries_only_the_forwarded_flag() {
4477 let mut core = open_fresh_core();
4478 let dm = room("!dm:example.org");
4479 let target = room("!target:example.org");
4480 add_room(&mut core, &dm, true, Some("Secret DM name"));
4481 add_room(&mut core, &target, false, None);
4482 seed_decrypted(
4483 &mut core,
4484 &dm,
4485 "$src:example.org",
4486 "@bob:example.org",
4487 "m.room.message",
4488 serde_json::json!({
4489 "msgtype": "m.text",
4490 "body": "quarterly numbers",
4491 "m.relates_to": { "m.in_reply_to": { "event_id": "$parent:example.org" } },
4492 "m.mentions": { "user_ids": ["@carol:example.org"] },
4493 }),
4494 );
4495
4496 forward(&mut core, &dm, "$src:example.org", &target).expect("a decrypted message is forwardable");
4497
4498 let (path, body) = sent_room_event(&mut core);
4499 assert!(path.contains("/send/m.room.message/"), "the event type is kept: {path}");
4500 assert_eq!(
4501 body,
4502 serde_json::json!({ "msgtype": "m.text", "body": "quarterly numbers", "forwarded": true }),
4503 "the text plus the bare flag; the reply pointer and the mentions are stripped"
4504 );
4505 let wire = body.to_string();
4506 for leaked in ["bob", "carol", "!dm", "Secret DM name", "$parent"] {
4507 assert!(!wire.contains(leaked), "{leaked:?} must not travel with a forward out of an encrypted room: {wire}");
4508 }
4509 let echo = last_item(&core, &target);
4510 assert_eq!(echo.forwarded, Some(Forwarded::Hidden), "the forwarder's own echo shows the same marker");
4511 }
4512
4513 #[test]
4514 fn forward_from_public_channel_carries_room_id_and_name_no_sender() {
4515 let mut core = open_fresh_core();
4516 let channel = room("!chan:example.org");
4517 let target = room("!target:example.org");
4518 add_room(&mut core, &channel, true, Some("Announcements"));
4519 mark_as_channel(&mut core, &channel);
4520 add_room(&mut core, &target, false, None);
4521 seed_plain(
4522 &mut core,
4523 &channel,
4524 "$post:example.org",
4525 "@dave:example.org",
4526 "m.room.message",
4527 serde_json::json!({
4528 "msgtype": "m.text",
4529 "body": "BTC breaks out",
4530 "m.mentions": { "room": true },
4531 "m.relates_to": { "m.in_reply_to": { "event_id": "$older:example.org" } },
4532 }),
4533 );
4534
4535 forward(&mut core, &channel, "$post:example.org", &target).expect("a channel post is forwardable");
4536
4537 let (_, body) = sent_room_event(&mut core);
4538 assert_eq!(
4539 body,
4540 serde_json::json!({
4541 "msgtype": "m.text",
4542 "body": "BTC breaks out",
4543 "forwarded_from": { "room_id": "!chan:example.org", "room_name": "Announcements" },
4544 })
4545 );
4546 assert!(!body.to_string().contains("dave"), "the original sender never travels");
4547 assert_eq!(
4548 last_item(&core, &target).forwarded,
4549 Some(Forwarded::Channel { room_id: channel.clone(), room_name: "Announcements".to_string() })
4550 );
4551 }
4552
4553 #[test]
4554 fn forward_from_an_unnamed_room_never_lists_member_names() {
4555 let mut core = open_fresh_core();
4556 let channel = room("!chan:example.org");
4557 let target = room("!target:example.org");
4558 add_room(&mut core, &channel, true, None);
4559 mark_as_channel(&mut core, &channel);
4560 if let Some(state) = core.rooms.get_mut(&channel) {
4561 state.members.insert(
4562 UserId::parse("@erin:example.org").expect("valid user id"),
4563 crate::room::state::MemberState {
4564 membership: Membership::Join,
4565 displayname: Some("Erin Private".to_string()),
4566 is_direct: false,
4567 },
4568 );
4569 }
4570 add_room(&mut core, &target, false, None);
4571 seed_plain(&mut core, &channel, "$post:example.org", "@erin:example.org", "m.room.message", serde_json::json!({ "msgtype": "m.text", "body": "hi" }));
4572
4573 forward(&mut core, &channel, "$post:example.org", &target).expect("forwardable");
4574
4575 let (_, body) = sent_room_event(&mut core);
4576 let name = body["forwarded_from"]["room_name"].as_str().expect("a name is sent");
4577 assert!(!name.contains("Erin"), "an unnamed room is attributed by a generic label, not by its members: {name}");
4578 }
4579
4580 #[test]
4581 fn forward_replaces_the_sources_own_forward_marker() {
4582 let mut core = open_fresh_core();
4583 let dm = room("!dm:example.org");
4584 let target = room("!target:example.org");
4585 add_room(&mut core, &dm, true, None);
4586 add_room(&mut core, &target, false, None);
4587 seed_decrypted(
4588 &mut core,
4589 &dm,
4590 "$src:example.org",
4591 "@bob:example.org",
4592 "m.room.message",
4593 serde_json::json!({
4594 "msgtype": "m.text",
4595 "body": "second hop",
4596 "forwarded_from": { "room_id": "!older:example.org", "room_name": "Older channel" },
4597 }),
4598 );
4599
4600 forward(&mut core, &dm, "$src:example.org", &target).expect("forwardable");
4601
4602 let (_, body) = sent_room_event(&mut core);
4603 assert_eq!(body, serde_json::json!({ "msgtype": "m.text", "body": "second hop", "forwarded": true }));
4604 }
4605
4606 #[test]
4607 fn forward_of_edited_item_sends_the_latest_text() {
4608 let mut core = open_fresh_core();
4609 let channel = room("!chan:example.org");
4610 let dm = room("!dm:example.org");
4611 let target = room("!target:example.org");
4612 add_room(&mut core, &channel, false, Some("Announcements"));
4613 add_room(&mut core, &dm, true, None);
4614 add_room(&mut core, &target, false, None);
4615
4616 seed_plain(&mut core, &channel, "$orig:example.org", "@dave:example.org", "m.room.message", serde_json::json!({ "msgtype": "m.text", "body": "typo text" }));
4618 seed_plain(
4619 &mut core,
4620 &channel,
4621 "$edit:example.org",
4622 "@dave:example.org",
4623 "m.room.message",
4624 serde_json::json!({
4625 "msgtype": "m.text",
4626 "body": "* fixed text",
4627 "m.new_content": { "msgtype": "m.text", "body": "fixed text" },
4628 "m.relates_to": { "rel_type": "m.replace", "event_id": "$orig:example.org" },
4629 }),
4630 );
4631 forward(&mut core, &channel, "$orig:example.org", &target).expect("forwardable");
4632 let (_, body) = sent_room_event(&mut core);
4633 assert_eq!(body["body"], "fixed text", "the latest edit's text, not the original and not the `* ` fallback");
4634 assert!(body.get("m.new_content").is_none() && body.get("m.relates_to").is_none());
4635
4636 seed_decrypted(&mut core, &dm, "$secret:example.org", "@bob:example.org", "m.room.message", serde_json::json!({ "msgtype": "m.text", "body": "typo secret" }));
4638 seed_sealed(&mut core, &dm, "$secret-edit:example.org", "@bob:example.org");
4639 let applied = core.timelines.get_mut(&dm).expect("timeline").apply_decrypted_relation(
4640 &event("$secret-edit:example.org"),
4641 UserId::parse("@bob:example.org").expect("valid user id"),
4642 2,
4643 RelatesTo::Replace { event_id: event("$secret:example.org") },
4644 Some(serde_json::json!({ "msgtype": "m.text", "body": "fixed secret" })),
4645 );
4646 assert!(applied, "the decrypted edit folds onto its target");
4647 forward(&mut core, &dm, "$secret:example.org", &target).expect("forwardable");
4648 let (_, body) = sent_room_event(&mut core);
4649 assert_eq!(body, serde_json::json!({ "msgtype": "m.text", "body": "fixed secret", "forwarded": true }));
4650 }
4651
4652 #[test]
4653 fn forward_of_unknown_msgtype_keeps_its_fields() {
4654 let mut core = open_fresh_core();
4655 let dm = room("!dm:example.org");
4656 let target = room("!target:example.org");
4657 add_room(&mut core, &dm, true, None);
4658 add_room(&mut core, &target, false, None);
4659 let card = serde_json::json!({
4660 "msgtype": "com.example.card",
4661 "body": "opaque",
4662 "note": "kept raw",
4663 });
4664 seed_decrypted(&mut core, &dm, "$card:example.org", "@bob:example.org", "m.room.message", card.clone());
4665
4666 forward(&mut core, &dm, "$card:example.org", &target).expect("an unrecognized msgtype is still an m.room.message");
4667
4668 let (path, body) = sent_room_event(&mut core);
4669 assert!(path.contains("/send/m.room.message/"));
4670 let mut expected = card;
4671 expected["forwarded"] = serde_json::Value::Bool(true);
4672 assert_eq!(body, expected, "unrecognized msgtype fields pass through plus the marker");
4673 assert!(matches!(last_item(&core, &target).content, ItemContent::Unknown), "the echo does not invent a card type");
4674 }
4675
4676 #[test]
4677 fn forward_refuses_redacted_sealed_unknown_and_missing() {
4678 let mut core = open_fresh_core();
4679 let dm = room("!dm:example.org");
4680 let target = room("!target:example.org");
4681 add_room(&mut core, &dm, true, None);
4682 add_room(&mut core, &target, false, None);
4683
4684 seed_decrypted(&mut core, &dm, "$redacted:example.org", "@bob:example.org", "m.room.message", serde_json::json!({ "msgtype": "m.text", "body": "gone" }));
4685 core.timelines.get_mut(&dm).expect("timeline").apply_redaction(&event("$redacted:example.org"));
4686 seed_sealed(&mut core, &dm, "$sealed:example.org", "@bob:example.org");
4687 seed_sealed(&mut core, &dm, "$utd:example.org", "@bob:example.org");
4688 core.timelines
4689 .get_mut(&dm)
4690 .expect("timeline")
4691 .set_decrypted(&event("$utd:example.org"), ItemContent::Undecryptable { reason: "MissingSession".to_string() });
4692 seed_plain(&mut core, &dm, "$call:example.org", "@bob:example.org", "m.call.invite", serde_json::json!({ "call_id": "c1" }));
4693
4694 for id in ["$redacted:example.org", "$sealed:example.org", "$utd:example.org", "$call:example.org", "$missing:example.org"] {
4695 match forward(&mut core, &dm, id, &target) {
4696 Err(MessengerError::IntentRefused(reason)) => {
4697 assert!(reason.contains(id), "the refusal names the event {id}: {reason}");
4698 }
4699 other => panic!("forwarding {id} must be refused, got {other:?}"),
4700 }
4701 }
4702 match forward(&mut core, &room("!nowhere:example.org"), "$redacted:example.org", &target) {
4703 Err(MessengerError::IntentRefused(reason)) => assert!(reason.contains("$redacted:example.org"), "{reason}"),
4704 other => panic!("an unknown source room means a missing event, got {other:?}"),
4705 }
4706 match forward(&mut core, &dm, "$sealed:example.org", &room("!nowhere:example.org")) {
4707 Err(MessengerError::IntentRefused(reason)) => assert!(reason.contains("unknown"), "{reason}"),
4708 other => panic!("an unknown target room is refused, got {other:?}"),
4709 }
4710
4711 assert!(core.send_state.is_empty(), "a refused forward queues nothing");
4712 assert!(core.timeline(&target).is_some_and(|timeline| timeline.items().is_empty()), "and leaves no echo");
4713 }
4714}