Skip to main content

mail4agent_messenger/
core.rs

1//! [`MessengerCore`] — the single-writer kernel tying every earlier piece
2//! together (plan §3's `lib.rs`/`sync_engine.rs`, work pieces M13a+M13b):
3//! the `/sync` loop, to-device/device-list/room ingestion, Megolm decrypt,
4//! the receive-side [`MessengerCommand`] variants (M13a), and the per-room
5//! ordered send pipeline plus every room-membership/tag/typing command
6//! (M13b, this module's own "M13b send pipeline" section below).
7//!
8//! # Kernel API shape
9//!
10//! [`MessengerCore::open`] rebuilds a core from whatever a shell already
11//! has durably stored (sealed records) plus fresh, caller-supplied config.
12//! From there, a shell drives one tick at a time:
13//!
14//! 1. [`MessengerCore::releasable_requests`] — everything safe to send
15//!    right now (the flush-before-send barrier already applied).
16//! 2. Execute each, feed the result back via [`MessengerCore::on_response`]
17//!    or [`MessengerCore::on_transport_error`].
18//! 3. [`MessengerCore::take_flush_batch`]/[`MessengerCore::ack_flush`] —
19//!    persist whatever mutated, ack once durable.
20//! 4. [`MessengerCore::events`] — drain whatever became visible, to diff
21//!    against the adapter's own last-rendered snapshot.
22//!
23//! No function here calls a clock or an RNG itself (crate doc, `ids`/
24//! `outgoing_queue`'s own module docs): every entry point that needs time
25//! takes `now_ms` explicitly, and retry jitter comes from the
26//! [`crate::outgoing_queue::Jitter`] supplied to [`MessengerCore::open`].
27//!
28//! # Counters that must never repeat across a restart
29//!
30//! [`crate::ids::RequestId`]/[`crate::ids::TxnId`] are minted from one
31//! monotonic seed per device session (`ids`'s own module doc) — this
32//! module is that seed's one owner, persisted as [`Counters`]. A bump is
33//! written through [`crate::store::StateStore::save_counters`] in the same
34//! synchronous call that mints the id ([`MessengerCore::next_request_id`]),
35//! which is what lets the ordinary flush-before-send barrier
36//! (`crate::persist`'s module doc) cover it for free: the very next
37//! [`crate::outgoing_queue::OutgoingQueue::enqueue`] call this id feeds
38//! into calls [`crate::persist::FlushEpoch::seal_for_request`] itself,
39//! which folds the counter's own not-yet-batched dirty write into the same
40//! batch the new request's `required_seq` waits on. A reused id after a
41//! crash would be silently deduplicated by the server into the OLD
42//! event/request — the exact "resend a `/keys/claim`, waste one OTK" cost
43//! [`crate::outgoing_queue`]'s own module doc accepts is very different
44//! from resending a `TxnId` a room member already has under a different
45//! plaintext, which is why this restore rule is load-bearing, not
46//! defensive polish.
47//!
48//! [`MessengerCore::open`] restores each sequence to
49//! `max(persisted, 1 + highest seed still embedded in a pending request)`
50//! — see [`restore_counter`]'s own doc for why the `pending` half of that
51//! max matters at all given the write-before-mint discipline above (belt
52//! and suspenders, not the primary guarantee). Only [`crate::ids::RequestId`]
53//! has a structural "still pending" source to scan
54//! ([`crate::store::StateStore::pending_requests`] — every pending request
55//! carries its own id); [`crate::ids::TxnId`] does not (a `TxnId` is baked
56//! into a request's `path`/`body`, never a distinct field), so its restore
57//! relies solely on the persisted value. [`MessengerCore::next_txn_id`] is
58//! the send pipeline's own counterpart to
59//! [`MessengerCore::next_request_id`] — symmetric restore, symmetric
60//! write-before-mint discipline.
61//!
62//! # Sync ingestion order
63//!
64//! One `/sync` request is ever outstanding at a time
65//! ([`crate::outgoing_queue::Lane::Sync`]'s own cap), 30s long-poll,
66//! `since` = the persisted [`crate::store::StateStore::sync_token`] (initial
67//! sync when `None`). On a successful response, in order: (1) to-device
68//! events — Olm-decrypt, route `m.room_key` into
69//! [`crate::crypto::group_sessions::GroupSessionManager`], drop a
70//! replayed/otherwise-unrecoverable decrypt silently, keep an
71//! [`crate::crypto::olm_sessions::OlmDecryptError::UnknownSenderDevice`] in
72//! a small bounded retry list re-tried after the next `/keys/query`.
73//! [`crate::crypto::olm_sessions::OlmSessionManager::decrypt_to_device`]
74//! checks the sender's device BEFORE it creates or advances any session, so
75//! a refused event has consumed nothing (no one-time key, no ratchet step) and
76//! the retry decrypts it from scratch -- the room key inside the very first
77//! message from a peer whose device this account has not queried yet (the
78//! usual case: the to-device share arrives before any `/keys/query` for
79//! them) is not lost.
80//! A missing inbound Megolm session sends `m.room_key_request` to the
81//! event sender and to this account's other devices. A request from a
82//! joined or invited room member is answered with `m.forwarded_room_key`
83//! (the inbound session exported at its first known index) or, if this
84//! device only still has the outbound session, with `m.room_key`.
85//! `m.forwarded_room_key` is imported and bound to the original sender.
86//! `m.room_key.withheld` records the session so the request is not
87//! repeated. (2) `device_lists.changed`/
88//! `left` → [`crate::crypto::device_tracker::DeviceTracker`]; (3) OTK
89//! counts / unused fallback types →
90//! [`crate::crypto::account::OlmAccountState::on_sync_counts`], may enqueue
91//! `/keys/upload`; (4) any outdated tracked user → one `/keys/query`; (5)
92//! rooms join/invite/leave (a member seen joining an encrypted room is marked
93//! outdated and queried right after, step 5b) — state → [`crate::room::state::RoomState`],
94//! timeline → [`crate::room::timeline::Timeline::apply_timeline_batch`],
95//! Megolm events decrypted via
96//! [`crate::crypto::group_sessions::GroupSessionManager::decrypt_event`] →
97//! [`crate::room::timeline::Timeline::set_decrypted`] or
98//! `Undecryptable{reason}`; ephemeral typing/receipts; room account data;
99//! summary + unread; (6) global account data (`m.direct`) → later read by
100//! [`MessengerCore::room_kind`]; (7) the new sync token is persisted LAST,
101//! after everything it covers. When a room key newly arrives (step 1),
102//! every `Undecryptable` item across every loaded timeline whose session id
103//! now matches is retried automatically, in addition to the caller's own
104//! [`MessengerCommand::RetryDecryption`].
105
106use 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
134/// `/sync`'s own long-poll timeout (research doc §1.4's convention).
135const SYNC_TIMEOUT_MS: u64 = 30_000;
136/// `GET /rooms/{roomId}/messages`'s page size for
137/// [`MessengerCommand::LoadOlder`].
138const LOAD_OLDER_PAGE_SIZE: u32 = 50;
139/// How long a `typing: true` indicator asks the server to hold before
140/// auto-clearing it (client-server API's own `timeout` field).
141const TYPING_TIMEOUT_MS: u64 = 30_000;
142/// How many to-device events with an as-yet-[`OlmDecryptError::UnknownSenderDevice`]
143/// sender this core holds for a retry after the next `/keys/query`, oldest
144/// evicted first — same bounded-FIFO doctrine as
145/// [`crate::room::timeline`]'s own pending-relation queue.
146const UNKNOWN_SENDER_RETRY_CAPACITY: usize = 256;
147/// How often [`MessengerCommand::SetTyping`] actually sends a `typing:
148/// true` PUT while the caller keeps re-asserting it — a `false` (stop) is
149/// never debounced (M13b's own send-pipeline doc).
150const TYPING_DEBOUNCE_MS: i64 = 4_000;
151/// [`GroupSessionManager::chunk_recipients_for_send_to_device`]'s own cap
152/// already bounds one `sendToDevice` request's device count; this is the
153/// wire event type every such request carries when it is sharing a Megolm
154/// room key -- the OUTER to-device type is always `m.room.encrypted` (the
155/// body is already Olm ciphertext; the inner, only-visible-after-decrypt
156/// type is `m.room_key`), matching every other Olm-wrapped to-device send
157/// this crate builds (`crypto::olm_sessions`'s own module doc).
158const ROOM_KEY_SHARE_EVENT_TYPE: &str = "m.room.encrypted";
159
160/// This device's own [`crate::ids::RequestId`]/[`crate::ids::TxnId`]
161/// minting sequences — see this module's own doc, "Counters that must
162/// never repeat across a restart".
163#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
164pub struct Counters {
165    /// The next, never-yet-handed-out [`crate::ids::RequestId`] seed.
166    pub next_request_id: u64,
167    /// The next, never-yet-handed-out [`crate::ids::TxnId`] seed. Not
168    /// minted from anywhere in this piece yet (see this module's own doc)
169    /// — restored and persisted symmetrically with `next_request_id` so
170    /// M13b's send pipeline needs no new persistence of its own.
171    pub next_txn_id: u64,
172}
173
174/// Computes the value a counter must resume from after a restart: never
175/// below what was durably persisted, and never low enough to reissue an id
176/// a still-pending request already carries. `max_used_seed` is the highest
177/// seed [`MessengerCore::open`] found embedded in a request
178/// [`crate::store::StateStore::pending_requests`] is still tracking (`None`
179/// when nothing is pending, or when this counter has no such structural
180/// source to scan at all — see this module's own doc for
181/// [`crate::ids::TxnId`]'s case). This is a defensive second check, not the
182/// primary guarantee: under the write-before-mint discipline
183/// [`MessengerCore::next_request_id`] follows, the persisted value alone
184/// should already dominate every pending id.
185/// Txn seed for a fresh store: microseconds since the Unix epoch, so it is
186/// above every id this device could have minted from an earlier store.
187fn 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/// Static identity a shell supplies once at [`MessengerCore::open`] time —
202/// never re-derived, never changed for the lifetime of a device session
203/// (plan §4.1).
204#[derive(Clone, Debug, PartialEq, Eq)]
205pub struct CoreConfig {
206    /// This account's own user id.
207    pub user_id: UserId,
208    /// This device's own, server-assigned device id.
209    pub device_id: DeviceId,
210    /// The homeserver's own server name — not consumed by this crate's own
211    /// logic yet (every id this core builds, including
212    /// [`MessengerCommand::CreateRoom`]'s, is either read whole off the wire
213    /// or supplied by a caller who already has one; `CreateRoomKind` has no
214    /// alias field), kept on [`CoreConfig`] because a shell already has to
215    /// know it to reach the homeserver at all, and a still-later piece
216    /// (minting a room alias's local part into a full alias) will need it.
217    pub server_name: String,
218}
219
220/// Secret material a shell supplies once at [`MessengerCore::open`] time,
221/// alongside [`CoreConfig`]'s static identity fields. The session private
222/// key is not here: it is the Olm account
223/// [`crate::crypto::account::OlmAccountState::load_or_create`] generates
224/// (`mail4agent_vodozemac::olm::Account::new`) and pickles into the store.
225/// These two fields are optional 32-byte keys the caller already holds.
226/// `Default` leaves both absent. Nothing in this crate reads `backup_key`
227/// yet.
228#[derive(Default)]
229pub struct CoreSecrets {
230    /// This device's record-sealing key. `None` when this core is opened
231    /// with a non-sealing codec (every test in this crate uses
232    /// [`crate::store::InsecurePlainCodecForTests`] instead).
233    /// [`MessengerCore::open_sealed`] is the only consumer: it requires
234    /// this field to build its own [`SealedRecordCodec`] and fails loudly
235    /// if it is absent rather than silently falling back to an unsealed
236    /// store.
237    pub store_seal_key: Option<Zeroizing<[u8; 32]>>,
238    /// This account's server-side key-backup key, supplied by the caller.
239    /// Carried for a later Megolm session backup/restore; nothing in this
240    /// crate reads it yet.
241    pub backup_key: Option<Zeroizing<[u8; 32]>>,
242}
243
244/// Which `m.room.message` shape [`OutgoingMessage`] renders as (M13b send
245/// pipeline).
246#[derive(Clone, Copy, Debug, PartialEq, Eq)]
247pub enum MessageKind {
248    /// `msgtype: "m.text"`.
249    Text,
250    /// `msgtype: "m.notice"`.
251    Notice,
252    /// `msgtype: "m.emote"`.
253    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/// One message a caller wants sent via [`MessengerCommand::SendMessage`] —
267/// M13b's binding content shape: text/notice/emote, an optional reply
268/// pointer, and an optional edit target (which must name one of the
269/// caller's own already-sent events — [`MessengerCore::dispatch`] rejects
270/// an edit of someone else's event rather than silently sending a
271/// replacement the room's other members will refuse to accept, per
272/// `m.replace`'s own "must come from the original sender" rule,
273/// `crate::room::relations`'s own module doc). `reply_to` and `edit_of`
274/// are mutually exclusive at the wire level — if both are set, `edit_of`
275/// wins (an edit's own `m.relates_to` is `m.replace`, not a reply).
276#[derive(Clone, Debug, PartialEq, Eq)]
277pub struct OutgoingMessage {
278    /// Which `msgtype` this renders as.
279    pub kind: MessageKind,
280    /// The message's own plain-text body (the edited/new body, if this is
281    /// an edit).
282    pub body: String,
283    /// A rich-reply pointer (`m.in_reply_to`), if any.
284    pub reply_to: Option<EventId>,
285    /// The event this message replaces, if this is an edit — must be one
286    /// of the caller's own already-sent events (this module's own doc).
287    pub edit_of: Option<EventId>,
288}
289
290/// Which shape [`MessengerCommand::CreateRoom`] builds, mapped onto the
291/// server's own `POST /createRoom` field set (`visibility`/`is_direct`/
292/// `invite`/`name`/`topic` — the server derives `kind` itself from
293/// `visibility`+`is_direct`, per the server plan's §5; a client never
294/// sends a `preset`/`kind` field directly).
295#[derive(Clone, Debug, PartialEq, Eq)]
296pub enum CreateRoomKind {
297    /// A 1:1 conversation: `is_direct: true`, exactly one invitee. Also
298    /// merges the new room into this account's own `m.direct` account data
299    /// once the room is created (this module's own doc).
300    Dm {
301        /// The other party.
302        peer: UserId,
303    },
304    /// A private multi-person room.
305    Group {
306        /// The room's display name.
307        name: String,
308        /// Users to invite at creation time.
309        invite: Vec<UserId>,
310        /// `true` gives every member the `invite` power level (`0`);
311        /// `false` restricts inviting to a moderator/owner (`50`) — the
312        /// server plan's §5 `group` column.
313        members_can_invite: bool,
314    },
315    /// A public, announcement-style room. Still E2E-encrypted; "public"
316    /// means join_rule public (anyone can join), not plaintext.
317    Channel {
318        /// The room's display name.
319        name: String,
320        /// The room's topic, if any.
321        topic: Option<String>,
322    },
323}
324
325/// One command a shell/adapter dispatches into the core — both the M13a
326/// receive-side variants (`MarkRead`, `SetTyping`, `RetryDecryption`,
327/// `LoadOlder`) and the M13b send/room-membership ones. Only `PartialEq`
328/// (not `Eq`): [`MessengerCommand::SetTag`]'s `order: Option<f64>` cannot
329/// implement `Eq`.
330#[derive(Clone, Debug, PartialEq)]
331pub enum MessengerCommand {
332    /// Sends one message into a room's per-room ordered send pipeline
333    /// (this module's own doc) — encrypts first if the room is encrypted,
334    /// sharing the room key to any not-yet-shared device before the
335    /// message itself is sent. A [`crate::room::timeline::TimelineItem`]
336    /// local echo appears immediately, in [`crate::room::timeline::SendState::Sending`].
337    SendMessage {
338        /// The room to send into.
339        room_id: RoomId,
340        /// The message content.
341        message: OutgoingMessage,
342        /// A caller-supplied transaction id (deterministic tests); `None`
343        /// mints a fresh one from this device's own counter.
344        txn_id: Option<TxnId>,
345    },
346    /// Forwards one event of `from_room` into `to_room`, through the same
347    /// per-room send pipeline as [`MessengerCommand::SendMessage`] (so it is
348    /// Megolm-encrypted when `to_room` is encrypted).
349    ///
350    /// The forward rules live HERE, in the engine. The content is read from
351    /// the source timeline item's DECRYPTED content -- for an edited item the
352    /// latest edit's `m.new_content` -- with `m.relates_to`, `m.mentions`
353    /// and `m.new_content` removed and any earlier forward marker replaced
354    /// by this hop's own:
355    ///
356    /// - source is a private/encrypted room (DM/group): `"forwarded": true`
357    ///   and nothing else -- no original sender, no room id, no room name;
358    /// - source is a public channel (`RoomKind::Channel`):
359    ///   `"forwarded_from": {"room_id", "room_name"}` --
360    ///   no sender (public join does not make content plaintext).
361    ///
362    /// The event type `m.room.message` is kept. An unrecognized `msgtype`
363    /// passes through as the raw content object plus the marker.
364    ///
365    /// Refused with [`MessengerError::IntentRefused`], naming the event,
366    /// when the source item is missing, redacted, still sealed
367    /// (undecryptable or not decrypted yet), or not an `m.room.message`;
368    /// and when `to_room` is unknown to this core (its encryption state
369    /// cannot be known, and a plaintext send into it could leak an
370    /// encrypted room's content).
371    Forward {
372        /// The room the event is taken from.
373        from_room: RoomId,
374        /// The event being forwarded.
375        event_id: EventId,
376        /// The room it is sent into.
377        to_room: RoomId,
378        /// A caller-supplied transaction id; `None` mints a fresh one.
379        txn_id: Option<TxnId>,
380    },
381    /// Sends a plaintext `m.reaction` — never Megolm-encrypted, even in an
382    /// encrypted room (allowed there as metadata, per the server's own
383    /// rule this crate mirrors).
384    React {
385        /// The room the reaction is sent in.
386        room_id: RoomId,
387        /// The event being reacted to.
388        target: EventId,
389        /// The reaction's own key (typically an emoji).
390        key: String,
391    },
392    /// Redacts an event (`PUT /rooms/{roomId}/redact/{eventId}/{txnId}`).
393    Redact {
394        /// The room the target event lives in.
395        room_id: RoomId,
396        /// The event being redacted.
397        target: EventId,
398        /// Why, if given.
399        reason: Option<String>,
400    },
401    /// Re-queues a [`crate::room::timeline::SendState::Failed`] send at the
402    /// back of its room's own send queue, WITHOUT re-encrypting it if it
403    /// had already been encrypted (the exact same ciphertext, the exact
404    /// same `txn_id` — this module's own doc). A no-op if `txn_id` names
405    /// no currently-failed send in `room_id`.
406    RetrySend {
407        /// The room the failed send belongs to.
408        room_id: RoomId,
409        /// The failed send's own transaction id.
410        txn_id: TxnId,
411    },
412    /// Creates a room (`POST /createRoom`).
413    CreateRoom {
414        /// Which shape to create.
415        kind: CreateRoomKind,
416    },
417    /// Joins a room this account has been invited to (or a public
418    /// channel).
419    JoinRoom {
420        /// The room to join.
421        room_id: RoomId,
422    },
423    /// Leaves a room.
424    LeaveRoom {
425        /// The room to leave.
426        room_id: RoomId,
427    },
428    /// Invites a user to a room.
429    Invite {
430        /// The room to invite into.
431        room_id: RoomId,
432        /// The user being invited.
433        user_id: UserId,
434    },
435    /// Removes a member from a room (`POST /rooms/{roomId}/kick`) — the
436    /// next encrypted send in this room rotates its Megolm session
437    /// automatically, excluding them (`crypto::group_sessions`'s own
438    /// rotation-trigger doc).
439    Kick {
440        /// The room to remove them from.
441        room_id: RoomId,
442        /// The member being removed.
443        user_id: UserId,
444        /// Why, if given.
445        reason: Option<String>,
446    },
447    /// Adds (or updates) an `m.tag` room-account-data tag.
448    SetTag {
449        /// The room being tagged.
450        room_id: RoomId,
451        /// The tag's own name.
452        tag: String,
453        /// The tag's sort order among sibling tags, if any.
454        order: Option<f64>,
455    },
456    /// Removes a tag from a room.
457    RemoveTag {
458        /// The room the tag is removed from.
459        room_id: RoomId,
460        /// The tag's own name.
461        tag: String,
462    },
463    /// Sets this room's read-up-to marker to `event_id` (`POST
464    /// /rooms/{roomId}/read_markers`).
465    MarkRead {
466        /// The room being marked read.
467        room_id: RoomId,
468        /// The last event the user has seen.
469        event_id: EventId,
470    },
471    /// Starts or stops this device's own typing indicator in a room (`PUT
472    /// /rooms/{roomId}/typing/{userId}`).
473    SetTyping {
474        /// The room to set the indicator in.
475        room_id: RoomId,
476        /// `true` to start, `false` to clear.
477        typing: bool,
478    },
479    /// Re-attempts decrypting one currently-undecryptable timeline item
480    /// (e.g. after the user manually re-verifies a device) — the same
481    /// retry this core already runs automatically the moment a matching
482    /// room key arrives (this module's own doc).
483    RetryDecryption {
484        /// The room the item lives in.
485        room_id: RoomId,
486        /// The item's event id.
487        event_id: EventId,
488    },
489    /// Fetches one older page of a room's timeline (`GET
490    /// /rooms/{roomId}/messages?dir=b`), prepending it once the response
491    /// arrives. A no-op if this room has no further history to page into
492    /// (no gap, no back-pagination token).
493    LoadOlder {
494        /// The room to page further back into.
495        room_id: RoomId,
496    },
497    /// Sets one piece of this account's own global account data (`PUT
498    /// /user/{userId}/account_data/{type}`) -- optimistically applied to
499    /// [`MessengerCore::global_account_data`] in the same call (this
500    /// module's own doc: the account-data endpoint replaces the whole
501    /// value, so there is nothing to merge).
502    SetAccountData {
503        /// The account-data type being set.
504        event_type: String,
505        /// The full new value.
506        content: serde_json::Value,
507    },
508    /// Sets one piece of `room_id`'s own room-scoped account data (`PUT
509    /// /user/{userId}/rooms/{roomId}/account_data/{type}`) -- optimistically
510    /// applied to [`MessengerCore::room_account_data`] in the same call.
511    /// [`MessengerCommand::SetTag`]/`RemoveTag` are the `"m.tag"`-specific
512    /// counterpart to this general command and behave the same way. Every
513    /// account-data write travels on [`crate::outgoing_queue::Lane::
514    /// AccountData`]: strictly FIFO, one in flight, so two writes to the
515    /// same key can neither race nor land out of order.
516    SetRoomAccountData {
517        /// The room the account data is scoped to.
518        room_id: RoomId,
519        /// The account-data type being set.
520        event_type: String,
521        /// The full new value.
522        content: serde_json::Value,
523    },
524    /// Searches the public-room directory (`POST /publicRooms`). The
525    /// latest result set is read back through
526    /// [`MessengerCore::public_rooms_result`]; a response to a search
527    /// superseded by a later one is dropped (this module's own doc).
528    SearchPublicRooms {
529        /// The search term.
530        term: String,
531    },
532    /// Searches the user directory (`POST /user_directory/search`). The
533    /// latest result set is read back through
534    /// [`MessengerCore::user_search_result`]; same supersession rule as
535    /// [`MessengerCommand::SearchPublicRooms`].
536    SearchUsers {
537        /// The search term.
538        term: String,
539    },
540}
541
542/// One user-visible change a shell/adapter should re-read from a snapshot
543/// getter. See [`MessengerCore::events`] for delivery semantics and
544/// [`MessengerCore::change_counter`] for the cheap "did anything change at
545/// all" signal.
546#[derive(Clone, Debug, PartialEq, Eq)]
547pub enum MessengerEvent {
548    /// The room list itself changed shape this tick (a room appeared via
549    /// join/invite/leave, or an already-known room's own state changed).
550    RoomsChanged,
551    /// `room_id`'s timeline gained or updated an item.
552    TimelineChanged {
553        /// The affected room.
554        room_id: RoomId,
555    },
556    /// `room_id`'s typing-user set changed.
557    TypingChanged {
558        /// The affected room.
559        room_id: RoomId,
560    },
561    /// `room_id`'s read-receipt set changed.
562    ReceiptsChanged {
563        /// The affected room.
564        room_id: RoomId,
565    },
566    /// `room_id`'s server-reported unread count changed.
567    UnreadChanged {
568        /// The affected room.
569        room_id: RoomId,
570    },
571    /// A previously-known device's Curve25519/Ed25519 keys changed — the
572    /// new keys were applied (`verified` cleared). Outbound Megolm shares
573    /// to that device are forgotten so the next send re-shares the room key.
574    DeviceKeyChanged {
575        /// The device's owner.
576        user_id: UserId,
577        /// The device whose keys changed.
578        device_id: DeviceId,
579    },
580    /// [`MessengerCore::public_rooms_result`] changed (a
581    /// [`MessengerCommand::SearchPublicRooms`] response landed and was not
582    /// superseded).
583    PublicRoomsChanged,
584    /// [`MessengerCore::user_search_result`] changed (a
585    /// [`MessengerCommand::SearchUsers`] response landed and was not
586    /// superseded).
587    UserSearchChanged,
588}
589
590/// One Megolm room event this core could not decrypt yet, kept so a later
591/// room-key arrival or an explicit [`MessengerCommand::RetryDecryption`]
592/// can re-attempt it without re-fetching anything — see this module's own
593/// "Sync ingestion order" doc.
594#[derive(Clone)]
595struct PendingDecryptItem {
596    session_id: String,
597    sender: UserId,
598    origin_server_ts: i64,
599    content: MegolmEncryptedContent,
600}
601
602/// `GET /rooms/{roomId}/messages`'s response shape, the parts this module
603/// reads (research doc §1.5) — `chunk` arrives newest-first for a `dir=b`
604/// page, reversed by [`MessengerCore::ingest_room_messages_response`]
605/// before [`crate::room::timeline::Timeline::prepend_back_page`] (which
606/// wants oldest-first, matching its own doc).
607#[derive(Deserialize)]
608struct RoomMessagesResponseBody {
609    #[serde(default)]
610    chunk: Vec<RawEvent>,
611    #[serde(default)]
612    end: Option<String>,
613}
614
615/// `PUT .../send/...`'s own response shape — the parts this module reads.
616#[derive(Deserialize)]
617struct SendResponseBody {
618    event_id: EventId,
619}
620
621/// `POST /createRoom`'s own response shape — the parts this module reads.
622#[derive(Deserialize)]
623struct CreateRoomResponseBody {
624    room_id: RoomId,
625}
626
627/// One row of a `POST /publicRooms` result -- read back through
628/// [`MessengerCore::public_rooms_result`].
629#[derive(Clone, Debug, PartialEq)]
630pub struct PublicRoomsResultEntry {
631    /// The room's own id.
632    pub room_id: RoomId,
633    /// The room's display name, if it set one.
634    pub name: Option<String>,
635    /// The room's topic, if it set one.
636    pub topic: Option<String>,
637    /// The room's own joined-member count.
638    pub num_joined_members: u64,
639}
640
641/// `POST /publicRooms`'s own response shape — the parts this module reads.
642#[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/// One row of a `POST /user_directory/search` result -- read back through
660/// [`MessengerCore::user_search_result`].
661#[derive(Clone, Debug, PartialEq)]
662pub struct UserDirectoryResultEntry {
663    /// The user's own id.
664    pub user_id: UserId,
665    /// The user's display name, if they set one.
666    pub display_name: Option<String>,
667}
668
669/// `POST /user_directory/search`'s own response shape — the parts this
670/// module reads.
671#[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/// One room-scoped send's own content, before it is turned into wire bytes
685/// (M13b send pipeline). A reaction/redaction never enters the key-sharing
686/// half of the pipeline at all — only a room event ([`SendPayload::Message`],
687/// [`SendPayload::Event`]) can be Megolm-encrypted.
688#[derive(Clone)]
689enum SendPayload {
690    /// A composed message — goes through the key-query/claim/share phases
691    /// when its room is encrypted.
692    Message(OutgoingMessage),
693    /// Any other room event whose wire type and plaintext content are
694    /// already final (a forward) — walks the same phases as
695    /// [`SendPayload::Message`].
696    Event {
697        /// The wire event type (`m.room.message`).
698        event_type: String,
699        /// The event's plaintext content.
700        content: serde_json::Value,
701    },
702    /// A plaintext `m.reaction` — sent as-is, in any room.
703    Reaction {
704        /// The event being reacted to.
705        target: EventId,
706        /// The reaction's own key.
707        key: String,
708    },
709    /// A redaction (`PUT .../redact/...`) — sent as-is, in any room.
710    Redaction {
711        /// The event being redacted.
712        target: EventId,
713        /// Why, if given.
714        reason: Option<String>,
715    },
716}
717
718/// A room event as it goes on the wire before any encryption — what
719/// [`SendPayload::room_event`] derives from a [`SendPayload::Message`] or a
720/// [`SendPayload::Event`].
721struct RoomEventWire {
722    event_type: String,
723    content: serde_json::Value,
724    /// The cleartext `m.relates_to` an encrypted envelope must copy
725    /// ([`message_relates_to`]).
726    relates_to: Option<RelatesTo>,
727}
728
729impl SendPayload {
730    /// The room-event shape of this payload; `None` for a reaction or a
731    /// redaction (never encrypted, never echoed into the timeline).
732    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/// The exact wire request a [`PendingSend`] will make once it reaches
752/// [`SendPhase::AwaitingSend`] — computed once (either directly, for a
753/// plaintext send, or after Megolm-encrypting a [`SendPayload::Message`])
754/// and cached from then on, so a later [`MessengerCommand::RetrySend`]
755/// never re-derives it (never re-encrypts — this module's own doc).
756#[derive(Clone)]
757enum WireSend {
758    /// `PUT /rooms/{roomId}/send/{eventType}/{txnId}`.
759    Event {
760        /// The wire event type (`m.room.message`, `m.room.encrypted`, or
761        /// `m.reaction`).
762        event_type: String,
763        /// The event's own content.
764        content: serde_json::Value,
765    },
766    /// `PUT /rooms/{roomId}/redact/{eventId}/{txnId}`.
767    Redact {
768        /// The event being redacted.
769        target: EventId,
770        /// Why, if given.
771        reason: Option<String>,
772    },
773}
774
775/// Where one room-scoped send currently sits in its own pipeline (M13b's
776/// module doc: key-query -> key-claim -> key-share -> the room event
777/// itself, in that order, for an encrypted [`SendPayload::Message`];
778/// straight to [`SendPhase::AwaitingSend`] for everything else).
779#[derive(Clone)]
780enum SendPhase {
781    /// Waiting for an earlier send in the same room to finish first —
782    /// [`MessengerCore::start_send_pipeline`] has not run for this send
783    /// yet.
784    Queued,
785    /// Waiting on one `/keys/query` this pipeline itself minted, for a
786    /// member whose device list is outdated or was never tracked at all —
787    /// [`MessengerCore::send_request_owner`] is what routes the eventual
788    /// response back to this send, so this phase itself carries no request
789    /// id of its own.
790    AwaitingKeysQuery,
791    /// Waiting on one `/keys/claim`, for recipient devices this account
792    /// has no Olm session with yet.
793    AwaitingKeysClaim {
794        /// Exactly the devices that request claimed a key for — needed to
795        /// apply the response ([`OlmSessionManager::on_keys_claim_response`]
796        /// takes the original request list, not just the response body).
797        devices: Vec<StoredDevice>,
798    },
799    /// Waiting on every `sendToDevice` request sharing the room key with
800    /// its not-yet-shared recipients — the room event itself is withheld
801    /// until every one of these succeeds (this module's own doc: "WAIT for
802    /// those requests to succeed before sending the room event").
803    AwaitingKeyShare {
804        /// Every still-outstanding share request's own id.
805        pending: BTreeSet<RequestId>,
806    },
807    /// Waiting on the room event/redaction itself.
808    AwaitingSend,
809    /// Terminal: the send (or one of its pipeline steps) failed outright.
810    /// Kept (rather than removed) so [`MessengerCommand::RetrySend`] can
811    /// find it again, and so [`MessengerCore::send_failure_reason`] can
812    /// surface why.
813    Failed {
814        /// The Matrix `errcode` the failing response carried.
815        errcode: String,
816    },
817}
818
819/// One room-scoped send this core is tracking end to end — from
820/// [`MessengerCommand::SendMessage`]/`React`/`Redact` through to a
821/// terminal [`SendPhase::AwaitingSend`] response. See this module's own
822/// doc for the pipeline every [`SendPayload::Message`] in an encrypted
823/// room walks through.
824#[derive(Clone)]
825struct PendingSend {
826    /// The room this send belongs to — also this send's own place in
827    /// [`MessengerCore::pending_sends`]'s per-room FIFO.
828    room_id: RoomId,
829    /// What is being sent.
830    payload: SendPayload,
831    /// Where this send currently sits in its own pipeline.
832    phase: SendPhase,
833    /// The exact wire request, once computed — see [`WireSend`]'s own doc.
834    wire: Option<WireSend>,
835}
836
837/// One account-data slot: `(room, event type)`; `None` is this account's own
838/// global account data.
839type AccountDataKey = (Option<RoomId>, String);
840
841/// Bookkeeping for one [`AccountDataKey`] that has account-data writes
842/// queued or in flight (an entry exists only while `pending > 0`).
843///
844/// A write updates the cached value optimistically. Without this guard a
845/// `/sync` that echoes an OLDER write of the same key would overwrite that
846/// optimistic value, and the next edit would build on the stale copy and
847/// silently drop a change. While writes are outstanding, `/sync` only
848/// records what the server reports (`server_value`); the optimistic cache
849/// stays authoritative. When the last write settles: a success keeps the
850/// optimistic value (it is what the server now holds; the next sync
851/// confirms it), a terminal failure restores `server_value`.
852struct AccountDataGuard {
853    /// Writes for this key queued or in flight.
854    pending: u32,
855    /// The last value the server is known to hold for the key: whatever
856    /// the cache held when the first outstanding write was made, then
857    /// whatever `/sync` reported since. `None`: nothing known (never seen).
858    server_value: Option<serde_json::Value>,
859}
860
861/// The `(room, event type)` an account-data PUT `path` writes, or `None` if
862/// `path` is not an account-data path.
863fn 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
880/// The single-writer kernel — see this module's own doc.
881pub struct MessengerCore<C: RecordCodec> {
882    config: CoreConfig,
883    store: Store<C>,
884    account: OlmAccountState,
885    /// Automatic cross-signing state (M3), loaded lazily on the first sync.
886    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    /// The latest [`MessengerCommand::SearchPublicRooms`] result set, read
902    /// back through [`MessengerCore::public_rooms_result`].
903    public_rooms_result: Vec<PublicRoomsResultEntry>,
904    /// The [`RequestId`] of the most recently dispatched
905    /// [`MessengerCommand::SearchPublicRooms`] -- a response naming any
906    /// other id is a superseded search and is dropped (this module's own
907    /// doc).
908    latest_public_rooms_request: Option<RequestId>,
909    /// The latest [`MessengerCommand::SearchUsers`] result set, read back
910    /// through [`MessengerCore::user_search_result`].
911    user_search_result: Vec<UserDirectoryResultEntry>,
912    /// Same supersession tracking as `latest_public_rooms_request`, for
913    /// [`MessengerCommand::SearchUsers`].
914    latest_user_search_request: Option<RequestId>,
915
916    /// Keys with account-data writes outstanding -- see
917    /// [`AccountDataGuard`]. In-memory only: rebuilt at
918    /// [`MessengerCore::open`] from the persisted request queue.
919    account_data_guard: BTreeMap<AccountDataKey, AccountDataGuard>,
920    /// Which key each outstanding account-data write belongs to.
921    account_data_write_keys: BTreeMap<RequestId, AccountDataKey>,
922
923    pending_decrypt: BTreeMap<RoomId, BTreeMap<EventId, PendingDecryptItem>>,
924    /// `(room, session)` pairs this process already asked for. Memory-only:
925    /// a restart asks again, which is what recovers a share the previous
926    /// process acked and then lost.
927    key_requests_sent: BTreeSet<(RoomId, String)>,
928    /// M4: standalone `/keys/claim` calls made so a joiner can reach old
929    /// members (id -> devices claimed).
930    history_claims: BTreeMap<RequestId, Vec<StoredDevice>>,
931    /// M4: devices a history claim was already attempted for this process.
932    history_claim_tried: BTreeSet<(UserId, DeviceId)>,
933    /// Sessions a peer refused with `m.room_key.withheld`. Stops the retry.
934    withheld_sessions: BTreeSet<(RoomId, String)>,
935    /// `(requester, device, request_id)` already answered with a share.
936    key_requests_answered: BTreeSet<(UserId, DeviceId, String)>,
937    unknown_sender_retry: VecDeque<ToDeviceEvent>,
938
939    /// The last `now_ms` a [`MessengerCommand::SetTyping { typing: true,
940    /// .. }`] actually sent a PUT at, per room — the debounce window
941    /// [`TYPING_DEBOUNCE_MS`] is measured from (this module's own doc). No
942    /// entry means either never sent, or explicitly stopped.
943    typing_debounce: BTreeMap<RoomId, i64>,
944    /// Every room-scoped send's own transaction id, in FIFO send order —
945    /// only the front of each room's queue is active (this module's own
946    /// "one message's pipeline completes before the next ... is
947    /// encrypted" doc); [`MessengerCore::finish_active_send`] pops it and
948    /// starts the next one.
949    pending_sends: BTreeMap<RoomId, VecDeque<TxnId>>,
950    /// Every [`PendingSend`] this core is tracking, keyed by its own
951    /// `txn_id` — kept even after it leaves `pending_sends` (a terminal
952    /// [`SendPhase::Failed`] stays here so [`MessengerCommand::RetrySend`]
953    /// can find it again).
954    send_state: BTreeMap<TxnId, PendingSend>,
955    /// Which [`PendingSend`] a still-pending request (a key-query,
956    /// key-claim, key-share, or the send/redact itself) this pipeline
957    /// minted belongs to — same reasoning as `pending_kinds`.
958    send_request_owner: BTreeMap<RequestId, TxnId>,
959    /// Which [`CreateRoomKind`] a still-pending `POST /createRoom` request
960    /// was building, so its response can merge `m.direct` for a
961    /// [`CreateRoomKind::Dm`] (this module's own doc).
962    pending_create_room: BTreeMap<RequestId, CreateRoomKind>,
963
964    /// Every [`OutgoingRequestKind`] this core itself minted and is still
965    /// waiting on, keyed by the [`RequestId`] it enqueued it under —
966    /// [`crate::outgoing_queue::OutgoingQueue`] does not expose a way to
967    /// look this back up given only an id, so this core tracks it
968    /// alongside instead (see [`MessengerCore::enqueue_request`]).
969    pending_kinds: BTreeMap<RequestId, OutgoingRequestKind>,
970    /// Which room a still-pending `RoomMessages` request was paginating —
971    /// same reasoning as `pending_kinds`.
972    pending_room_messages: BTreeMap<RequestId, RoomId>,
973    /// The currently outstanding `/sync` request's id, if any — this
974    /// core's own single-flight tracking (this module's own doc).
975    sync_request_id: Option<RequestId>,
976    /// Whether a `/sync` request has completed successfully in THIS core's
977    /// lifetime. Until then the next `/sync` is minted with `timeout=0` (an
978    /// immediate catch-up answer, the matrix-js-sdk convention); only after
979    /// it does the core long-poll.
980    sync_caught_up: bool,
981    /// [`FlushBatch`]es [`crate::outgoing_queue::OutgoingQueue::enqueue`]
982    /// already drained from the store's own pending set (to seal a new
983    /// request's [`crate::persist::RequiredSeq`]) but that a shell has not
984    /// yet drained via [`MessengerCore::take_flush_batch`] — see
985    /// [`MessengerCore::enqueue_request`]'s own doc for why this can't
986    /// simply be re-derived from `store.take_flush_batch()` later.
987    pending_flush_batches: VecDeque<FlushBatch>,
988
989    events: Vec<MessengerEvent>,
990    change_counter: u64,
991    /// The last response this core could not ingest (a wire shape it does
992    /// not decode, say), for [`MessengerCore::take_ingest_error`]. The
993    /// response is otherwise dropped, so without this a shell would see
994    /// nothing but a `/sync` loop that never advances.
995    ingest_error: Option<String>,
996}
997
998impl<C: RecordCodec> MessengerCore<C> {
999    /// Rebuilds a core from whatever a shell already read back from
1000    /// durable storage (`records`, possibly empty for a brand-new device),
1001    /// `config`, `secrets` (caller-supplied [`CoreSecrets`]: `store_seal_key`
1002    /// is consumed by [`MessengerCore::open_sealed`] before it reaches here,
1003    /// and `backup_key` is not read by this kernel), and a `jitter` source for
1004    /// [`crate::outgoing_queue::OutgoingQueue`]'s own retry backoff. `now_ms`
1005    /// is accepted for API symmetry with every other entry point (this
1006    /// crate's own "no clock inside" rule) but not read by anything in this
1007    /// piece's own restore logic yet — timeline content itself is never
1008    /// persisted (re-fetched off the next `/sync`/back-page instead, this
1009    /// module's own doc), so there is nothing time-sensitive to restore
1010    /// here before a later piece needs one.
1011    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        // The seal key lives in `codec` (`open_sealed` builds a
1020        // `SealedRecordCodec` from `store_seal_key`). `backup_key` stays on
1021        // `CoreSecrets` for the caller; this kernel does not read it.
1022        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        // A `/sync` restored from the store carries the timeout it was minted
1041        // with (a 30 s long-poll): replaying it would stall a resumed session
1042        // for that long before anything is known. It has no side effect to
1043        // lose, so drop it (its id already counted toward the seed above);
1044        // `ensure_sync_enqueued` mints a fresh catch-up sync from the
1045        // persisted token.
1046        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            // Restoration actually advanced a value beyond what was on
1054            // disk (the defensive "pending" half of `restore_counter`
1055            // fired) -- persist the advanced value immediately so it, in
1056            // turn, can never be handed out a second time by a future
1057            // restart.
1058            store.save_counters(counters)?;
1059        }
1060        // A store lost and recreated for a device that already exists
1061        // server-side (the bearer survived in the keychain) must not restart
1062        // txn ids low: the server would dedup each send into an OLD event
1063        // and report it as sent. Floor the in-memory counter at the wall
1064        // clock; `next_txn_id` persists it before the first id is used.
1065        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    /// Rebuilds [`MessengerCore::account_data_guard`] (and the optimistic
1116    /// cache it protects) from account-data writes that were still pending
1117    /// when the core last stopped: each persisted request counts as one
1118    /// outstanding write for its key, and the newest request's body -- the
1119    /// value the server will hold once the queue drains -- is the cached
1120    /// value again. The server's own current value is unknown until
1121    /// `/sync` reports it.
1122    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    /// Mints the next [`RequestId`], persisting the advanced [`Counters`]
1139    /// in the same call — see this module's own doc for why this ordering
1140    /// is what makes the flush-before-send barrier cover it for free.
1141    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    /// Mints the next [`TxnId`], persisting the advanced [`Counters`] in
1149    /// the same call — symmetric with [`MessengerCore::next_request_id`]
1150    /// (this module's own doc's "Counters that must never repeat across a
1151    /// restart" section).
1152    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    /// Enqueues `request` on `lane`, remembering its kind so a later
1160    /// [`MessengerCore::on_response`] can route the response without
1161    /// [`crate::outgoing_queue::OutgoingQueue`] needing a by-id lookup of
1162    /// its own.
1163    ///
1164    /// [`crate::outgoing_queue::OutgoingQueue::enqueue`] already drains
1165    /// whatever was pending into a [`FlushBatch`] the instant it needs a
1166    /// [`crate::persist::RequiredSeq`] to seal the new request against
1167    /// (`crate::persist`'s own doc) — that batch is gone from the store's
1168    /// own pending set the moment it is produced, so if this method
1169    /// discarded it here, [`MessengerCore::take_flush_batch`] would never
1170    /// see it again and the request (and everything after it in the same
1171    /// lane) would stay unreleasable forever. [`MessengerCore::pending_flush_batches`]
1172    /// is where it waits instead, until a shell actually drains it.
1173    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    // ---------------------------------------------------------------
1193    // M13b send pipeline
1194    // ---------------------------------------------------------------
1195
1196    /// Starts (or queues) one room-scoped send — the shared entry point
1197    /// for [`MessengerCommand::SendMessage`]/`Forward`/`React`/
1198    /// `Redact`. Only a room event ([`SendPayload::Message`],
1199    /// [`SendPayload::Event`]) gets an immediate timeline local echo (this
1200    /// module's own doc): appended as [`SendState::LocalEcho`] then
1201    /// immediately moved to [`SendState::Sending`], matching the plan's
1202    /// own wording ("local echo in the timeline with `SendState::Sending`
1203    /// -> `MessengerEvent::TimelineChanged`") — both happen inside this one
1204    /// call, so a caller never observes the intermediate `LocalEcho` value.
1205    /// Sending a message also stops this device's own typing indicator in
1206    /// the room, unconditionally (module doc: "false on send").
1207    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    /// Re-queues a [`SendPhase::Failed`] send at the back of its room's own
1254    /// FIFO, reusing its already-computed [`WireSend`] (never re-encrypting
1255    /// — [`MessengerCommand::RetrySend`]'s own doc). A no-op if `txn_id`
1256    /// names no currently-failed send in `room_id`.
1257    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    /// Advances `txn_id` (which must currently be the front of its room's
1279    /// own FIFO) to its next pipeline step. If [`PendingSend::wire`] is
1280    /// already computed (a retry, or any send re-entering this function
1281    /// after an earlier step completed), it goes straight to
1282    /// [`MessengerCore::begin_awaiting_send`] — never re-derived. Otherwise:
1283    /// a room event ([`SendPayload::Message`]/[`SendPayload::Event`]) in a
1284    /// currently-encrypted room enters the key pipeline; everything else
1285    /// builds its plaintext wire body directly.
1286    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    /// Every user this account must consider a Megolm recipient for
1306    /// `room_id`: every currently join-or-invite member, plus this
1307    /// account's own user id (its OTHER devices also need the room key —
1308    /// `crypto::group_sessions`'s own module doc).
1309    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    /// Marks any of `member_ids` this account has never tracked at all as
1327    /// newly tracked-and-outdated ([`DeviceTracker::on_device_lists`]),
1328    /// then returns every one of `member_ids` that is currently outdated
1329    /// (whether it already was, or just became so) — the send pipeline's
1330    /// own "any member's device list is outdated or unknown" check.
1331    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    /// Every device this store currently has on file for any of
1350    /// `member_ids`.
1351    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    /// `known_devices`, filtered down to the devices this pipeline must
1360    /// actually reach: a current member's, not blocked, not this device's
1361    /// own.
1362    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    /// Step 1 of the encrypted-message pipeline (this module's own doc):
1372    /// if any recipient's device list is outdated or was never tracked at
1373    /// all, mint one `/keys/query` and wait for it; otherwise fall through
1374    /// to step 2.
1375    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    /// Step 2: if any recipient candidate has no established Olm session
1391    /// yet, mint one `/keys/claim` and wait for it; otherwise fall through
1392    /// to step 3.
1393    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    /// Step 3: ensures a current outbound Megolm session, fans out the
1413    /// room key (Olm to-device) to every not-yet-shared, currently-reachable
1414    /// device, and waits for every one of those `sendToDevice` calls to
1415    /// succeed before the room event itself is ever encrypted or sent
1416    /// (module doc: a share failure fails the whole message, not just the
1417    /// unreachable device). A device this account still has no Olm session
1418    /// with after step 2 (a refused/unanswered `/keys/claim`) is left out
1419    /// of the fan-out rather than failing the send outright — the same
1420    /// "can't reach everyone, degrade rather than block" doctrine this
1421    /// crate already applies to an unresolved UTD. If nothing needs
1422    /// sharing at all (already fully shared, or nobody reachable), falls
1423    /// straight through to encrypting and sending.
1424    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    /// Step 4: Megolm-encrypts a room event under `room_id`'s (now current)
1490    /// outbound session — the inner type is the event's own (`m.room.message`,
1491    /// `m.sticker`), the outer is always `m.room.encrypted` — and caches the
1492    /// resulting [`WireSend::Event`]. A no-op for any other payload kind
1493    /// (never reached for one, since only a room event in an encrypted room
1494    /// enters the key pipeline at all).
1495    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    /// Builds the plaintext [`WireSend`] for a room event
1517    /// ([`SendPayload::Message`]/[`SendPayload::Event`]) in an unencrypted
1518    /// room, a [`SendPayload::Reaction`] (any room), or a
1519    /// [`SendPayload::Redaction`] (any room) — the cases that never touch the
1520    /// key pipeline at all.
1521    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    /// Mints (or re-mints, for a retry) the actual `PUT .../send/...` or
1542    /// `PUT .../redact/...` request from [`PendingSend::wire`], on this
1543    /// room's own [`Lane::Room`] (preserves the room's own send order —
1544    /// this crate's `outgoing_queue`'s own doc).
1545    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    /// Routes one successful response for a request this send pipeline
1567    /// itself minted onward to whatever step comes next. A no-op for any
1568    /// `kind` this pipeline never mints on its own behalf (defensive —
1569    /// every `send_request_owner` entry is inserted alongside a specific
1570    /// kind, this match is exhaustive over that set).
1571    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    /// The room event/redaction itself succeeded: for a room event
1616    /// ([`SendPayload::Message`]/[`SendPayload::Event`]), attaches the
1617    /// server-confirmed event id to its echo (module doc: "200 -> store `event_id` on the echo
1618    /// (`SendState::Sent`)"); either way, removes this send's own
1619    /// bookkeeping and starts the room's next queued send, if any.
1620    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    /// A step in `txn_id`'s own pipeline failed outright (a non-retryable
1636    /// 4xx anywhere along it): marks a room event's echo
1637    /// [`SendState::Failed`], moves this send to the terminal
1638    /// [`SendPhase::Failed`] (kept, not removed — [`MessengerCommand::RetrySend`]
1639    /// needs it), and starts the room's next queued send, if any.
1640    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    /// Pops `txn_id` off the front of `room_id`'s own send FIFO (if that is
1654    /// where it still is) and starts the next queued send's pipeline, if
1655    /// any — the "one message's pipeline completes before the next ... is
1656    /// encrypted" ordering rule, enforced structurally by never starting
1657    /// entry N+1 before entry N reaches a terminal state.
1658    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    /// A shell calls this once per tick, executes each, and feeds the
1673    /// result back via [`MessengerCore::on_response`]/
1674    /// [`MessengerCore::on_transport_error`]. Also where the next `/sync`
1675    /// request is minted, if none is currently outstanding (this module's
1676    /// own doc).
1677    ///
1678    /// **One-shot per request: never discard a non-empty result.** Every
1679    /// [`OutgoingRequest`] this call returns is marked in flight by
1680    /// [`crate::outgoing_queue::OutgoingQueue::releasable`] as a side
1681    /// effect of being returned — it will not be offered again until a
1682    /// matching [`MessengerCore::on_response`]/
1683    /// [`MessengerCore::on_transport_error`] call clears that flag. A
1684    /// caller that calls this method purely to trigger the mint side
1685    /// effect (e.g. to force a freshly-dirtied `/sync` request's own flush
1686    /// epoch open before deciding whether to act) and then throws the
1687    /// return value away strands every request in that batch permanently
1688    /// in flight, since nothing else will ever release it. Call this
1689    /// exactly once per tick, and always drive every request it returns
1690    /// through to a response.
1691    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        // The first `/sync` of a session must never long-poll. An initial
1703        // sync (no `since`) is answered at once by Synapse whatever `timeout`
1704        // says, but a server that treats a brand-new, still-empty account's
1705        // initial view as "nothing new" holds it for the full timeout -- and
1706        // this core publishes its device keys only once that first response
1707        // has landed. A resumed session (persisted `since`) has the same
1708        // need: the catch-up answer must not wait for news that may never
1709        // come. Live run, 2026-09-29: fresh accounts sat 30 s before their
1710        // keys went up.
1711        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    /// Feeds one HTTP response back in, keyed by the id the original
1719    /// [`OutgoingRequest`] carried. Returns whatever became visible as a
1720    /// direct result of processing it (equivalent to calling
1721    /// [`MessengerCore::events`] immediately afterwards — there is exactly
1722    /// one underlying event queue behind both).
1723    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            // A parse/logic failure against an otherwise-successful HTTP
1756            // response is dropped, not applied: the one invariant that
1757            // matters -- never advancing past a partially-applied `/sync` --
1758            // holds structurally, since `ingest_sync` only persists the new
1759            // sync token as its very last step. The failure is kept for
1760            // [`MessengerCore::take_ingest_error`] so a shell can show it.
1761            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        // A refused `/keys/upload` must not look like a quiet success: the
1785        // account stays unpublished and the shell can surface the errcode
1786        // instead of reporting mail as sent under a rejected identity.
1787        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    /// Feeds a transport-level failure back in (a timeout, a connection
1797    /// reset, ...) — always resolves to a retry with this core's own
1798    /// backoff (`crate::outgoing_queue`'s own doc); a no-op if `request_id`
1799    /// is not (or no longer) pending.
1800    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            // A response to a search this account has already moved past (a
1880            // later `SearchPublicRooms` was dispatched since) fails the
1881            // guard and is dropped -- only the LATEST search's own result
1882            // set is ever visible (this module's own doc).
1883            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            // Same supersession rule as `PublicRooms` above.
1898            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    /// A `POST /createRoom` succeeded. For [`CreateRoomKind::Dm`], also
1913    /// merges the new room into this account's own `m.direct` account data
1914    /// (M13b send-pipeline doc: "DM also merges `m.direct` account data")
1915    /// — every other kind needs nothing further here (room state/timeline
1916    /// populate off the next `/sync`, same as any other room).
1917    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    /// Merges `room_id` into this account's own `m.direct` account data
1924    /// under `peer`, optimistically, then sends the full updated value
1925    /// (the account-data endpoint replaces the whole value, same as
1926    /// [`MessengerCore::enqueue_tag_update`]'s own doc). Shared by
1927    /// [`MessengerCore::on_create_room_response`] (the CREATOR side of a
1928    /// [`CreateRoomKind::Dm`]) and [`MessengerCommand::JoinRoom`]'s own
1929    /// dispatch arm (the INVITEE side, this module's own doc: a DM
1930    /// invite's `is_direct` flag is not carried forward once THIS
1931    /// account's own join overwrites that membership event, so `m.direct`
1932    /// is the only durable "this is a DM" signal this account keeps of
1933    /// its own past that point).
1934    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    /// Sets (`Some`) or removes (`None`) the cached value for `key`,
1952    /// keeping the typed `m.direct` view in step with it.
1953    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    /// The one path every account-data write takes: builds the PUT, applies
1976    /// `value` to the cache optimistically (the next edit builds on it),
1977    /// guards the key against stale `/sync` echoes for as long as the
1978    /// write is outstanding ([`AccountDataGuard`]), and queues it on the
1979    /// strictly ordered [`Lane::AccountData`].
1980    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    /// One outstanding write for `key` reached its end (`succeeded`: a 2xx;
2004    /// otherwise a terminal failure). Once the last one settles, a failure
2005    /// restores the server's value over the now-wrong optimistic cache and
2006    /// tells the UI to re-render; a success keeps the optimistic value.
2007    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    /// Folds one account-data event `/sync` reported into the cache -- or,
2021    /// while writes for the same key are outstanding, only into that key's
2022    /// [`AccountDataGuard::server_value`].
2023    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    /// Drains everything dirtied since the previous call into one
2033    /// [`FlushBatch`] for a shell to write durably.
2034    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    /// Acknowledges that batch `id` has been durably written.
2042    pub fn ack_flush(&mut self, id: u64) {
2043        self.store.ack_flush(id);
2044    }
2045
2046    /// Drains every [`MessengerEvent`] queued since the last call to this
2047    /// method or to [`MessengerCore::on_response`].
2048    pub fn events(&mut self) -> Vec<MessengerEvent> {
2049        self.drain_events()
2050    }
2051
2052    /// The last response this core failed to ingest, once (cleared by
2053    /// reading it) -- see the `ingest_error` field's own doc.
2054    pub fn take_ingest_error(&mut self) -> Option<String> {
2055        self.ingest_error.take()
2056    }
2057
2058    /// Monotonically increasing count of visible changes — a cheap
2059    /// "did anything change at all since I last checked" signal for a
2060    /// provider that only wants to know whether to re-read its snapshot
2061    /// getters, without keeping its own copy of every [`MessengerEvent`].
2062    pub fn change_counter(&self) -> u64 {
2063        self.change_counter
2064    }
2065
2066    /// Every room this core currently holds state for.
2067    pub fn room_ids(&self) -> impl Iterator<Item = &RoomId> {
2068        self.rooms.keys()
2069    }
2070
2071    /// `room_id`'s current state, if this core has one.
2072    pub fn room_state(&self, room_id: &RoomId) -> Option<&RoomState> {
2073        self.rooms.get(room_id)
2074    }
2075
2076    /// `room_id`'s current timeline, if this core has one.
2077    pub fn timeline(&self, room_id: &RoomId) -> Option<&Timeline> {
2078        self.timelines.get(room_id)
2079    }
2080
2081    /// `room_id`'s derived [`RoomKind`] (plan manager decision #5's DM/
2082    /// channel/group precedence), using this core's current global
2083    /// `m.direct` account data.
2084    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    /// Every user currently typing in `room_id`, per the most recent
2089    /// `/sync` ephemeral batch — empty if none, or if this core has never
2090    /// seen a typing event for this room.
2091    pub fn typing_users(&self, room_id: &RoomId) -> &[UserId] {
2092        self.typing.get(room_id).map(Vec::as_slice).unwrap_or(&[])
2093    }
2094
2095    /// `room_id`'s most recently received `m.receipt` account data, if any.
2096    pub fn receipts(&self, room_id: &RoomId) -> Option<&ReceiptContent> {
2097        self.receipts.get(room_id)
2098    }
2099
2100    /// One piece of `room_id`'s room-scoped account data (e.g.
2101    /// `"m.fully_read"`), if this core has seen one.
2102    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    /// One piece of this account's own global account data, if this core
2107    /// has seen one.
2108    pub fn global_account_data(&self, event_type: &str) -> Option<&serde_json::Value> {
2109        self.global_account_data.get(event_type)
2110    }
2111
2112    /// How many missing Megolm sessions this process has already requested.
2113    pub fn key_request_count(&self) -> usize {
2114        self.key_requests_sent.len()
2115    }
2116
2117    /// This account's own user id, exactly as configured at
2118    /// [`MessengerCore::open`] -- a UI-facing adapter needs this to tell
2119    /// its own account apart from every other room member (a DM's peer, a
2120    /// reaction's `by_me`, a typing-user list minus self, ...).
2121    pub fn user_id(&self) -> &UserId {
2122        &self.config.user_id
2123    }
2124
2125    /// Whether this core currently has at least one verified device on
2126    /// file for `user_id` (`/keys/query` results already persisted by
2127    /// [`crate::crypto::device_tracker::DeviceTracker`]) -- a UI-facing
2128    /// adapter's own "has a messaging key" check for a member list, without
2129    /// duplicating that module's storage format.
2130    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    /// The latest [`MessengerCommand::SearchPublicRooms`] result set, empty
2135    /// until one has landed.
2136    pub fn public_rooms_result(&self) -> &[PublicRoomsResultEntry] {
2137        &self.public_rooms_result
2138    }
2139
2140    /// The latest [`MessengerCommand::SearchUsers`] result set, empty until
2141    /// one has landed.
2142    pub fn user_search_result(&self) -> &[UserDirectoryResultEntry] {
2143        &self.user_search_result
2144    }
2145
2146    /// The Matrix `errcode` a failed send (`txn_id`) ended with, if
2147    /// `txn_id` currently names a [`crate::room::timeline::SendState::Failed`]
2148    /// send this core is still tracking — `None` for a send that never
2149    /// failed, already succeeded, or that this core has never seen.
2150    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    /// The "in" edge — mailbox command (single-writer-core doctrine's Law
2158    /// 3). See [`MessengerCommand`]'s own doc for every variant this piece
2159    /// implements. `now_ms` is this call's own clock reading (this crate's
2160    /// "no clock inside" rule) — used for a fresh send's local-echo
2161    /// timestamp, [`MessengerCommand::SetTyping`]'s debounce window, and
2162    /// Megolm rotation timing.
2163    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                // A DM invite's `is_direct` flag lives on this account's own
2200                // invite membership event and is overwritten by the join
2201                // itself, so accepting it records the room in this account's
2202                // own `m.direct` (same as the creator side does on
2203                // `/createRoom`) -- otherwise the invitee's room kind falls
2204                // back to Group once the join lands.
2205                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                    // A known room with no pagination token yet -- a resumed
2281                    // session, whose timelines are not persisted and whose
2282                    // first `/sync` carries only what is new -- pages back
2283                    // from its own sync position (a `/sync` `next_batch` token
2284                    // is a valid `/messages` `from`).
2285                    .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    /// `POST /createRoom`'s own body per [`CreateRoomKind`] (server plan
2324    /// §5: the server derives `kind`/power levels from `visibility`+
2325    /// `is_direct` itself — a client only ever sends those two fields plus
2326    /// `invite`/`name`/`topic`).
2327    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    /// `room_id`'s current `m.tag` content, or an empty one if this core
2360    /// has never seen one — [`MessengerCommand::SetTag`]/`RemoveTag` both
2361    /// mutate a full copy of this (the account-data endpoint replaces the
2362    /// whole value, there is no partial-update verb).
2363    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    /// Optimistic and guarded like every account-data write
2370    /// ([`MessengerCore::write_account_data`]): the very next tag edit must
2371    /// build on THIS content, not on whatever `/sync` last echoed, or two
2372    /// back-to-back edits would each start from the same stale copy and the
2373    /// later PUT would drop the earlier one's change.
2374    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    /// The [`SendPayload::Event`] a [`MessengerCommand::Forward`] sends: the
2380    /// source item's current decrypted content, stripped and marked
2381    /// ([`forwardable_content`]) per that command's own rules, or the
2382    /// refusal naming `event_id`.
2383    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    /// The marker a forward out of `from_room` carries: the bare `forwarded:
2408    /// true` unless that room is a public channel, in which case
2409    /// `forwarded_from` names it (Slack-like attribution). Channels are
2410    /// still E2E; an unknown room is treated as private -- fail-safe.
2411    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    /// `room_id`'s name as a forward attribution: its `m.room.name`, or a
2422    /// generic label by kind. Never the members' names an unnamed room's
2423    /// chat-list row would list -- an attribution goes to third parties.
2424    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    /// Re-attempts decrypting one currently-pending Megolm item, applying
2437    /// the result to its timeline on success. Returns `true` iff the item
2438    /// was found and successfully decrypted (i.e. something a caller should
2439    /// treat as a visible change) — `false` for an unknown item or a
2440    /// decrypt that still fails.
2441    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    /// One backward page from the tip, not from the stored sync token.
2469    /// Used when a push names an event the timeline does not have as text.
2470    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    /// Decrypts Megolm rows already on the timeline. `pending_decrypt` is
2479    /// memory-only, so a restarted client still holds the ciphertext and
2480    /// would otherwise never turn it into text the leader trigger can see.
2481    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    /// Applies one successfully-decrypted Megolm plaintext to `room_id`'s
2552    /// timeline entry for `event_id` — as an ordinary message
2553    /// ([`Timeline::set_decrypted`]), or, when the plaintext itself carries
2554    /// an `m.annotation`/`m.replace` relation, via
2555    /// [`Timeline::apply_decrypted_relation`] instead (`room::timeline`'s
2556    /// own doc: an encrypted relation's real content — its `m.new_content`
2557    /// — only exists post-decrypt, unlike its `m.relates_to` pointer, which
2558    /// the outer envelope already copies in cleartext; a still-ciphertext
2559    /// event is therefore always given an ordinary placeholder row first,
2560    /// and this is the one place that later resolves what it actually was).
2561    /// Shared by forward decrypt ([`MessengerCore::decrypt_new_timeline_events`])
2562    /// and [`MessengerCore::retry_decrypt_one`].
2563    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        // 1. to-device.
2601        for event in &response.to_device.events {
2602            self.process_to_device_event(event, &mut newly_received_sessions)?;
2603        }
2604
2605        // 2. device_lists.
2606        DeviceTracker::on_device_lists(&mut self.store, &response.device_lists.changed, &response.device_lists.left)?;
2607
2608        // 3. OTK counts / unused fallback types.
2609        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        // 4. any outdated tracked user -> one keys/query.
2620        let outdated_before_rooms = self.outdated_tracked_users()?;
2621        if !outdated_before_rooms.is_empty() {
2622            self.enqueue_keys_query()?;
2623        }
2624
2625        // 5. rooms.
2626        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        // 5b. A member who joined an encrypted room was marked outdated while
2641        // the room was ingested (`track_joined_members`): query them now, not
2642        // at the next sync or send.
2643        if !self.outdated_tracked_users()?.is_subset(&outdated_before_rooms) {
2644            self.enqueue_keys_query()?;
2645        }
2646
2647        // 6. global account data.
2648        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        // 7. the new sync token is persisted LAST -- see this module's own
2657        // doc for why every earlier `?` above matters for this ordering.
2658        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            // Every other reason (a replay, a malformed envelope, a
2681            // payload-binding mismatch, ...) is dropped silently -- this
2682            // module's own doc, and `crypto::olm_sessions`'s: a redelivered
2683            // to-device message after a crash is legitimate, and this
2684            // crate has no logging dependency to report a genuine attack
2685            // attempt through instead.
2686            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    /// Who a forwarded Megolm session is bound to. The Olm sender is
2744    /// authenticated. The session is bound to the device named by
2745    /// `sender_key` when that device is the forwarder, or when the
2746    /// forwarder is this account (another of our devices). A room peer
2747    /// cannot install a session bound to somebody else.
2748    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    /// Re-share a session this device holds with the device that asked.
2781    /// Inbound export (first known index) is preferred: the outbound
2782    /// `session_key` is the current ratchet and cannot decrypt earlier
2783    /// messages.
2784    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    /// The device a key request should be encrypted back to. Prefers the
2845    /// stored device with `requesting_device_id`, then the Olm sender's
2846    /// curve key. A blocked device is not answered.
2847    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    /// Ask the event sender and this account's other devices for one
2864    /// missing Megolm session. Does nothing when this process already
2865    /// asked, or when the peer withheld it. A send that reaches no device
2866    /// stays unmarked so a later drive tries again.
2867    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        // M4: reach devices we share no Olm session with yet (we just joined
2898        // and never talked to them) by claiming one-time keys first; the
2899        // claim response re-runs this request.
2900        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    /// Olm-encrypt `content` as `inner_type` to every device that already
2941    /// has a session, and enqueue one `sendToDevice` per chunk. Devices
2942    /// with no Olm session are skipped. Returns how many devices were
2943    /// addressed.
2944    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    /// Every tracked user whose device list is currently outdated.
3008    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    /// Enqueues one `/keys/query` for every outdated tracked user (nothing if none is).
3013    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    /// Marks every other member that `events` show joining `room_id` (an
3022    /// encrypted room) as a tracked user with an outdated device list. The
3023    /// server answers `/keys/query` only for users who share a JOINED room with
3024    /// the caller and reports a device list in `device_lists.changed` only when
3025    /// it changes -- never when a user merely starts sharing a room. So a peer
3026    /// asked about while still only invited comes back with no devices and
3027    /// would otherwise stay "up to date" with none after joining, and no later
3028    /// message would ever be shared with them.
3029    /// M4: re-run key requests for every undecryptable (MissingSession) row.
3030    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    /// M4 policy: forward old room keys only when the room's history is
3042    /// shared (our rooms default to `shared`); a `joined`/`invited` room keeps
3043    /// pre-join history away from later members.
3044    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    /// M3: one step of automatic cross-signing. Keys first, then this
3052    /// device's self-signature. At most a few attempts per process.
3053    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    /// Attempts to decrypt every well-formed Megolm `m.room.encrypted`
3175    /// event in `events` (already applied to `room_id`'s timeline as
3176    /// `Encrypted{session_id}` placeholders by the caller), patching each
3177    /// one in place via [`MessengerCore::apply_decrypted_plaintext`] — shared
3178    /// by forward sync ([`MessengerCore::ingest_joined_room`]/
3179    /// [`MessengerCore::ingest_left_room`]) and back-pagination
3180    /// ([`MessengerCore::ingest_room_messages_response`]).
3181    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                    // Stored as the coarse `UtdReason` variant's own name
3209                    // (never the full, per-error-variant `Display` message)
3210                    // -- this is the only place a UI-facing adapter can
3211                    // recover which of the crate's own [`crate::crypto::
3212                    // group_sessions::UtdReason`] variants this failure
3213                    // maps to (`GroupDecryptError::utd_reason`'s own doc),
3214                    // since [`ItemContent::Undecryptable`] carries a plain
3215                    // `String`, not the enum itself.
3216                    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    /// The production entry point (M14): opens a core whose store seals
3247    /// every record with AES-256-GCM under `secrets.store_seal_key` — the
3248    /// constructor a shell calls, in place of the generic
3249    /// [`MessengerCore::open`] this crate's own tests use with
3250    /// [`crate::store::InsecurePlainCodecForTests`].
3251    ///
3252    /// Fails with [`MessengerError::Crypto`] if `secrets.store_seal_key` is
3253    /// `None`. The caller supplies that key. There is no unsealed fallback
3254    /// here by design (an unsealed messenger record store is exactly what
3255    /// [`SealedRecordCodec`] exists to keep off disk).
3256    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
3274/// A stable, always-well-formed placeholder [`EventId`] for a synthetic
3275/// [`RawEvent`] built from an invite's [`StrippedStateEvent`] (which itself
3276/// carries no event id) — [`RoomState::apply_state_event`] never reads
3277/// `event_id` at all, so its exact value is irrelevant; it only has to
3278/// parse.
3279fn stripped_placeholder_event_id() -> EventId {
3280    EventId::parse("$stripped-preview").expect("this literal is a well-formed, hardcoded $-prefixed opaque id")
3281}
3282
3283/// Reshapes one invite-preview [`StrippedStateEvent`] into the
3284/// [`RawEvent`] shape [`RoomState::apply_state_event`] accepts —
3285/// `origin_server_ts`/`unsigned`/`event_id` are unused by that method for a
3286/// state event, so a placeholder is enough (this function's own doc).
3287fn 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
3299/// `content` as a forward re-sends it: `m.relates_to` (a reply or edit
3300/// pointer into the source room), `m.mentions` and `m.new_content` removed,
3301/// any forward marker the source itself carried replaced by this hop's own
3302/// (`marker_key: marker`) -- so an earlier hop's origin never travels on.
3303/// Everything else (including an unrecognized `msgtype`'s own fields) passes
3304/// through untouched.
3305fn 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
3315/// The local-echo [`ItemContent`] for a not-yet-sent [`OutgoingMessage`] —
3316/// always the plain new body, regardless of `reply_to`/`edit_of` (this
3317/// crate's timeline has no dedicated "echo of an edit" rendering; the
3318/// confirmed event, once it arrives back off `/sync`, carries the full
3319/// `m.relates_to`-aware rendering via the ordinary
3320/// [`crate::room::relations`] aggregation — this module's own doc).
3321fn 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
3330/// The `m.room.message` wire content for `message` — plain
3331/// `{msgtype, body}` for an ordinary send; a reply adds the legacy
3332/// `m.in_reply_to` pointer; an edit instead prefixes the outer `body` with
3333/// `"* "` (the client-side fallback convention for a client that does not
3334/// understand `m.replace`) and carries the real replacement under
3335/// `m.new_content`, with `m.relates_to: m.replace` naming the edited event
3336/// (research doc §1.3; `reply_to` and `edit_of` are mutually exclusive at
3337/// the wire level — [`OutgoingMessage`]'s own doc).
3338fn 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
3352/// The outer, unencrypted-envelope `m.relates_to` a Megolm-encrypted
3353/// `message` must carry alongside its ciphertext (`crate::wire::events::
3354/// MegolmEncryptedContent`'s own doc: copied verbatim from the cleartext
3355/// plaintext's own relation, so the server can index it without ever
3356/// decrypting). `None` for a plain send with neither a reply nor an edit.
3357fn 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    /// Drains and acks every currently-pending flush batch -- the same
3398    /// "persist, then ack" tick a real shell performs once per loop
3399    /// iteration. Tests call this explicitly wherever they need a
3400    /// mutation's own flush epoch satisfied before checking
3401    /// [`MessengerCore::releasable_requests`] (the flush-before-send
3402    /// barrier, `crate::persist`'s own doc, applies just as much to a
3403    /// counter bump as to an Olm/Megolm ratchet advance).
3404    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        // The restore rule itself: a request still pending with a HIGHER
3424        // seed than whatever was last durably persisted must push the
3425        // resumed counter past it too, so a freshly minted id can never
3426        // collide with one already embedded in an in-flight request.
3427        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        // A stale, lower counter value -- as if the process crashed after
3439        // minting seed 5 but before this record was ever written.
3440        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        // The first call mints and enqueues the sync request, which itself
3464        // dirties the counter record -- not releasable until that flush is
3465        // acked (the flush-before-send barrier). The second call, after
3466        // acking, actually returns it.
3467        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        // Calling again before any response arrives must not mint a
3495        // second `/sync` -- the lane is occupied and this core's own
3496        // `sync_request_id` tracking prevents a duplicate enqueue too.
3497        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        // The next `/sync` is minted lazily (then needs its own flush round
3522        // before it is releasable) -- and only now does it long-poll.
3523        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        // Captured BEFORE encrypting: Megolm's ratchet only advances
3545        // forward, so a `session_key` exported AFTER `encrypt` would only
3546        // let a recipient decrypt messages from that later point onward,
3547        // never the one encrypted just before it.
3548        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        // The event arrives with no session known yet -- it lands as
3570        // `Undecryptable`. `apply_timeline_batch` first, exactly like
3571        // `ingest_joined_room` does, since `decrypt_new_timeline_events`
3572        // only ever patches an item [`Timeline::set_decrypted`] already
3573        // knows about -- it never inserts one itself.
3574        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        // The room key now arrives (e.g. via a to-device `m.room_key`) --
3585        // simulated directly against the manager, mirroring what
3586        // `route_decrypted_to_device` itself does.
3587        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        // Alice's own store/account -- the sender.
3610        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        // Signed once, up front -- see `unknown_sender_device_is_retried_after_keys_query`'s
3615        // own comment on this same call for why the order matters here.
3616        let alice_device_keys =
3617            alice_account.device_keys_json(&alice_user, &alice_device).expect("sign alice's device_keys");
3618
3619        // Bob's core is the one under test.
3620        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        // Bob learns alice's device ahead of time, so the FIRST decrypt
3648        // below is a genuine, fully-validated success (not merely an
3649        // `Ok(())` a caller can't tell apart from a silently-dropped
3650        // failure) -- this test's whole point is a real replay, not an
3651        // incidental one.
3652        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        // Simulate a restart: flush, ack, and rebuild the store/account
3700        // from just the durable records -- the exact scenario a redelivery
3701        // after a crash looks like.
3702        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        // The exact same to-device event, redelivered -- must not error the
3709        // whole ingestion path, and must not resurrect a second plaintext
3710        // anywhere a caller could observe.
3711        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        // Alice's own store/account -- the sender.
3722        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        // Signed once now (not only later, right before building the
3727        // `/keys/query` response) -- this is the exact same signed
3728        // `device_keys` payload either way, computed here so bob's own
3729        // `/keys/claim` handshake below observes it already having
3730        // happened, matching the sequencing every other multi-device e2e
3731        // test in this crate (`crypto::olm_sessions`'s own
3732        // `olm_session_rejects_out_of_order_ratchet_message`) already uses.
3733        let alice_device_keys =
3734            alice_account.device_keys_json(&alice_user, &alice_device).expect("sign alice's device_keys");
3735
3736        // Bob's own store/account -- a genuinely distinct user (not this
3737        // module's shared `user_id()`/`device_id()` test fixture, which is
3738        // "alice" for every OTHER test in this file) -- built standalone
3739        // first (same order as
3740        // `redelivered_to_device_after_restart_is_dropped_silently`) and
3741        // only assigned into a `MessengerCore` once the whole Olm handshake
3742        // is done.
3743        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        // Bob does not know alice's device yet: the decrypt is refused with
3795        // `UnknownSenderDevice` before any session is created, and the event
3796        // is queued for a retry.
3797        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        // Bob learns alice's device via a `/keys/query` response --
3807        // `MessengerCore::handle_terminal_success`'s own `KeysQuery` arm
3808        // calls `retry_unknown_sender_queue` exactly once, right after
3809        // applying the response.
3810        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        // Unknown room and no timeline: nothing to page.
3856        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        // The room is known (its state record was replayed) but its timeline is empty.
3861        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    /// A `/sync` body: `!r:example.org` where bob is seen joining, with an
3870    /// `m.room.encryption` state event only if `encrypted`.
3871    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        // `events()`/`on_response`'s own inline return share one queue --
3974        // nothing left to drain right after.
3975        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        // A second search supersedes the first before the first's response
4077        // ever lands.
4078        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        // The stale response (naming the FIRST request id) is dropped.
4087        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        // The current response is applied.
4092        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        // Optimistic: the cache already holds both tags, before any round trip.
4113        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    /// Answers the outstanding `/sync` request with `body` (minting and
4173    /// flushing it first if needed) and returns every OTHER request the
4174    /// rounds released, unanswered -- the tests below keep account-data
4175    /// writes pending on purpose.
4176    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        // The server echoes the FIRST write only, while both are outstanding.
4220        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        // Both writes settle; the optimistic value survives (it is what the
4232        // server holds now) and the guard is gone.
4233        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        // With nothing outstanding, sync applies normally again.
4247        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        // A failure while a later write is still outstanding leaves the cache
4283        // alone (the later write carries the full content anyway); once that
4284        // one succeeds its value stands.
4285        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        // Everything dirtied so far, including both pending PUT records.
4324        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    // -----------------------------------------------------------------
4346    // Stickers and forwarding
4347    // -----------------------------------------------------------------
4348
4349    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    /// Registers `room_id` in `core`: Megolm-encrypted or not, optionally
4360    /// named, with an (empty) timeline.
4361    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    /// Marks `room_id` as a public channel for [`RoomState::derive_room_kind`]
4376    /// (join_rule public + raised `events_default`).
4377    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    /// An event that arrived in `room_id` as plaintext.
4421    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    /// A Megolm envelope that arrived in `room_id` and has not been opened.
4427    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    /// A Megolm envelope that arrived in `room_id` and decrypted to
4433    /// `{event_type, content}`.
4434    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    /// Flushes, releases, and returns the `(path, body)` of the one
4449    /// `RoomSend` the core just queued -- after answering it with a fresh
4450    /// event id, so the room's next send is free to start.
4451    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        // Plaintext room: the edit is an ordinary event folded onto its target.
4617        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        // Encrypted room: the edit is only recognised once decrypted.
4637        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}