Skip to main content

mail4agent_server/
store.rs

1//! Matrix-shaped event store — rooms/events/state/members/relations/
2//! receipts/account-data/txn-dedup/filters half of `messenger.db` (the
3//! devices/keys/backup half is `matrix_keys_store.rs`, a separate work
4//! item). See `the messenger protocol notes`
5//! §2 for the full DDL this module implements (the "Manager decisions on
6//! this plan" section at the top of that file overrides the body — this
7//! module follows those corrections, noted inline where they apply) and §1
8//! for the id-format rules.
9//!
10//! This server adopts Matrix's own event envelope, state-resolution model,
11//! and CS API shape for everything except federation and `/login` — CS API
12//! v1.19, room version 11. Crypto stays entirely client-side: every
13//! `content` blob this module stores is opaque JSON it never inspects
14//! except for `m.relates_to` (relation bookkeeping, kept in cleartext at
15//! the top level of `content` even for `m.room.encrypted`, per the Matrix
16//! spec's own accepted trade-off — plan §10 item 3) and the redaction
17//! allow-list (§2's `redact_event`, which strips by `event_type`, never by
18//! interpreting the payload's meaning).
19//!
20//! # Single-writer discipline
21//!
22//! Every function here takes an already-open [`Connection`]; the caller
23//! (`state.rs`, from P14 onward) holds it behind one `std::sync::Mutex`,
24//! the same discipline `db`/`social_db` already use — this module never
25//! locks anything itself and never calls `std::sync::Mutex` internally.
26//! That single-writer guarantee is what makes [`next_stream_id`] safe: it
27//! is always called inside the same transaction as the row it stamps, and
28//! there is never a second writer racing it.
29//!
30//! # Cross-database rule
31//!
32//! This module stores `user_id INTEGER`. The nick lives on
33//! `messenger_sessions`, which [`crate::nick`] reads and writes. The
34//! `matrix_users.nick` column is left in place and is not the source of
35//! truth. There is no identity database. `matrix_users` maps `user_id` to
36//! its Matrix id (`mxid`) once: `public_id` is immutable, so the mapping
37//! does not change.
38//!
39//! The one place a label does end up in this database is the `displayname`
40//! field of an `m.room.member` event's `content` — the Matrix convention,
41//! and what a client names a DM and lists members by. Callers stamp the
42//! identity database's effective label into every `join`/`invite` member
43//! event they write, and [`refresh_member_displayname`] re-stamps it into
44//! the user's existing member events when the label changes.
45
46use rusqlite::{params, Connection, OptionalExtension, Transaction};
47use std::collections::HashSet;
48use std::fmt;
49
50// ============================================================================
51// Server name, id formats
52// ============================================================================
53
54/// The one Matrix server name this deployment ever answers as — used
55/// everywhere an mxid/room id is formatted or parsed. Never a second
56/// literal (plan §1).
57static SERVER_NAME_CELL: std::sync::OnceLock<String> = std::sync::OnceLock::new();
58
59/// Set the homeserver name once, before any mxid or room id is minted.
60/// The default until this runs is `example.org`.
61pub fn set_matrix_server_name(name: impl Into<String>) -> Result<(), &'static str> {
62    let name = name.into();
63    // A hostname, or `host:port` (a Matrix server name may carry a port; tests and private
64    // networks use it).
65    let (host, port) = name.split_once(':').map_or((name.as_str(), None), |(h, p)| (h, Some(p)));
66    if host.is_empty() || host.contains('/') || host.contains(':') || port.is_some_and(|p| p.parse::<u16>().is_err()) {
67        return Err("server name must be a hostname, optionally with :port");
68    }
69    SERVER_NAME_CELL.set(name).map_err(|_| "server name already set")
70}
71
72static LOCAL_ALIASES: std::sync::OnceLock<Vec<String>> = std::sync::OnceLock::new();
73
74
75/// Local DNS aliases of the one homeserver are a server-name CHECK only: an mxid addressed to an enabled
76/// alias resolves to the same local account. No second homeserver exists and ids are always minted
77/// with [`matrix_server_name()`]. Off by default.
78///
79/// Accept `names` (hostnames) as local aliases of [`matrix_server_name()`] when
80/// parsing mxids. Call once at boot (the server reads `M4A_LOCAL_NAMES`,
81/// comma-separated). Empty or repeated calls are ignored.
82pub fn set_local_aliases(names: impl IntoIterator<Item = String>) {
83    let list: Vec<String> = names.into_iter().map(|n| n.trim().to_ascii_lowercase()).filter(|n| !n.is_empty() && !n.contains(':') && !n.contains('/')).collect();
84    let _ = LOCAL_ALIASES.set(list);
85}
86
87/// Whether `name` is the minted server name or an enabled local alias.
88pub fn is_local_server_name(name: &str) -> bool {
89    name == matrix_server_name() || LOCAL_ALIASES.get().is_some_and(|aliases| aliases.iter().any(|a| a.eq_ignore_ascii_case(name)))
90}
91
92/// Homeserver name used when minting and parsing local ids.
93pub fn matrix_server_name() -> &'static str {
94    SERVER_NAME_CELL.get_or_init(|| "example.org".to_string()).as_str()
95}
96
97/// Room version pinned for every room this server creates (plan §1) — v11
98/// still lists the creator explicitly in `m.room.power_levels.users`,
99/// unlike v12's "infinite implicit power" simplification, which this plan
100/// does not adopt.
101pub const MATRIX_ROOM_VERSION: &str = "11";
102
103/// Ceiling on the serialized `content` of any state/timeline event this
104/// server accepts (plan §3.8, mirroring `dm_db::DM_MAX_CIPHERTEXT_BYTES`) —
105/// checked by every route that takes client-supplied event content
106/// (`routes::matrix::rooms::put_state` first; messaging/to-device pieces
107/// reuse the same constant rather than minting their own).
108pub const MATRIX_EVENT_CONTENT_MAX_BYTES: usize = 16 * 1024;
109
110/// Localpart prefix reserved for a future Application Service (bridge)
111/// namespace (plan §10 item 6). `matrix_store`'s own id generators never
112/// mint a localpart starting with this — see [`ensure_matrix_user`] — so a
113/// later AS registration slots in without a retroactive id-collision audit.
114pub const RESERVED_LOCALPART_PREFIX: &str = "_bridge_";
115
116/// True if `localpart` is reserved for the future Application Service
117/// namespace (plan §10 item 6) and must never be minted as a real user's
118/// mxid localpart.
119pub fn is_reserved_localpart(localpart: &str) -> bool {
120    localpart.starts_with(RESERVED_LOCALPART_PREFIX)
121}
122
123/// Why parsing an mxid failed.
124#[derive(Debug, Clone, Copy, PartialEq, Eq)]
125pub enum MatrixIdError {
126    /// Did not start with `@`.
127    MissingSigil,
128    /// No `:server_name` suffix at all.
129    MissingServerName,
130    /// Addressed to a server other than [`matrix_server_name()`] — this
131    /// deployment does not federate, so such an id can never resolve to a
132    /// real local account.
133    ForeignServerName,
134}
135
136impl fmt::Display for MatrixIdError {
137    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
138        match self {
139            MatrixIdError::MissingSigil => write!(f, "mxid is missing its '@' sigil"),
140            MatrixIdError::MissingServerName => write!(f, "mxid is missing a ':server_name' suffix"),
141            MatrixIdError::ForeignServerName => write!(f, "mxid is addressed to a foreign server name"),
142        }
143    }
144}
145
146/// Format `public_id` (the identity database's immutable `users.public_id`)
147/// as this server's own mxid: `@<public_id>:example.org`.
148pub fn mxid_for_public_id(public_id: &str) -> String {
149    format!("@{public_id}:{}", matrix_server_name())
150}
151
152/// Parse `@localpart:server_name`, returning the localpart — refusing
153/// anything not addressed to [`matrix_server_name()`] (plan §1: no
154/// federation, so a foreign-server mxid can never be a real local user).
155pub fn public_id_from_mxid(mxid: &str) -> Result<&str, MatrixIdError> {
156    let rest = mxid.strip_prefix('@').ok_or(MatrixIdError::MissingSigil)?;
157    let (localpart, server_name) = rest.split_once(':').ok_or(MatrixIdError::MissingServerName)?;
158    if !is_local_server_name(server_name) {
159        return Err(MatrixIdError::ForeignServerName);
160    }
161    Ok(localpart)
162}
163
164/// Mint 16 random bytes, URL-safe base64, no padding (~22 chars) — the
165/// random component shared by room ids and event ids (plan §1).
166fn random_id_component() -> String {
167    use base64::engine::general_purpose::URL_SAFE_NO_PAD;
168    use base64::Engine;
169    use rand::Rng;
170    let bytes: [u8; 16] = rand::thread_rng().gen();
171    URL_SAFE_NO_PAD.encode(bytes)
172}
173
174/// Mint a fresh room id: `"!" + 22 random URL-safe-base64 chars + ":example.org"`.
175pub fn new_room_id() -> String {
176    format!("!{}:{}", random_id_component(), matrix_server_name())
177}
178
179/// Mint a fresh event id: `"$" + 22 random URL-safe-base64 chars` — opaque,
180/// no reference-hash content addressing (no federation to need it).
181pub fn new_event_id() -> String {
182    format!("${}", random_id_component())
183}
184
185// ============================================================================
186// Errors
187// ============================================================================
188
189/// Why a write into this store was refused, on top of a real database
190/// failure — same shape as `dm_db::MessageError`/`IdentityKeyError`.
191#[derive(Debug)]
192pub enum MatrixStoreError {
193    Db(rusqlite::Error),
194    /// An event's `content` was not valid JSON where this module needed to
195    /// read a field out of it (membership, power levels, `m.relates_to`).
196    Json(serde_json::Error),
197    /// [`ensure_matrix_user`] was asked to mint a reserved localpart (plan
198    /// §10 item 6) — see [`is_reserved_localpart`].
199    ReservedLocalpart,
200    /// A second `m.annotation` relation with the same sender, target, and
201    /// aggregation key as an existing one (plan §2/§9: maps to
202    /// `M_DUPLICATE_ANNOTATION`/400 at the route layer, a later piece).
203    DuplicateAnnotation,
204    /// A referenced event id does not exist in `events`.
205    UnknownEventId(String),
206    /// A referenced mxid does not exist in `matrix_users` — the caller must
207    /// [`ensure_matrix_user`] before naming that user in room state.
208    UnknownMxid(String),
209    /// An `m.room.member` event's `content.membership` was missing or not
210    /// one of `join`/`invite`/`leave`/`ban`.
211    InvalidMembership(String),
212    /// A relation's `m.relates_to` target (the `String` is the target event
213    /// id) does not exist, or exists in a different room than the relation
214    /// event being inserted — a relation can never point out of its own
215    /// room.
216    InvalidRelationTarget(String),
217    /// A referenced event (the `String` is its event id) exists but is in a
218    /// different room than the one the caller named — refused rather than
219    /// silently acted on, since acting on it would let a caller redact or
220    /// mark-read an event that was never a member of the room they claimed.
221    WrongRoom(String),
222    /// [`redact_event`] was asked to redact an event type this server never
223    /// allows to be redacted (the `String` is that `event_type`) — see
224    /// [`redact_event`]'s doc comment.
225    UnredactableEvent(String),
226    /// A room of the DAG layer (feature `f3-hash-ids`) refused an event: the auth rules or a
227    /// signature check said no. The message is the reason.
228    F3Rejected(String),
229}
230
231impl From<rusqlite::Error> for MatrixStoreError {
232    fn from(e: rusqlite::Error) -> Self {
233        MatrixStoreError::Db(e)
234    }
235}
236
237impl From<serde_json::Error> for MatrixStoreError {
238    fn from(e: serde_json::Error) -> Self {
239        MatrixStoreError::Json(e)
240    }
241}
242
243// ============================================================================
244// Small TEXT-backed enums (`as_str`/`from_wire_name`, matching
245// `social_db::Class`'s own convention)
246// ============================================================================
247
248/// Our own room-kind label (`rooms.kind`) — drives the power-level/history-
249/// visibility defaults of plan §5. Not a Matrix wire field.
250#[derive(Debug, Clone, Copy, PartialEq, Eq)]
251pub enum RoomKind {
252    Dm,
253    Group,
254    Channel,
255}
256
257impl RoomKind {
258    pub fn as_str(self) -> &'static str {
259        match self {
260            RoomKind::Dm => "dm",
261            RoomKind::Group => "group",
262            RoomKind::Channel => "channel",
263        }
264    }
265
266    pub fn from_wire_name(s: &str) -> Option<Self> {
267        match s {
268            "dm" => Some(RoomKind::Dm),
269            "group" => Some(RoomKind::Group),
270            "channel" => Some(RoomKind::Channel),
271            _ => None,
272        }
273    }
274}
275
276#[derive(Debug, Clone, Copy, PartialEq, Eq)]
277pub enum JoinRule {
278    Invite,
279    Public,
280}
281
282impl JoinRule {
283    pub fn as_str(self) -> &'static str {
284        match self {
285            JoinRule::Invite => "invite",
286            JoinRule::Public => "public",
287        }
288    }
289
290    pub fn from_wire_name(s: &str) -> Option<Self> {
291        match s {
292            "invite" => Some(JoinRule::Invite),
293            "public" => Some(JoinRule::Public),
294            _ => None,
295        }
296    }
297}
298
299#[derive(Debug, Clone, Copy, PartialEq, Eq)]
300pub enum HistoryVisibility {
301    Shared,
302    WorldReadable,
303    Invited,
304    Joined,
305}
306
307impl HistoryVisibility {
308    pub fn as_str(self) -> &'static str {
309        match self {
310            HistoryVisibility::Shared => "shared",
311            HistoryVisibility::WorldReadable => "world_readable",
312            HistoryVisibility::Invited => "invited",
313            HistoryVisibility::Joined => "joined",
314        }
315    }
316
317    pub fn from_wire_name(s: &str) -> Option<Self> {
318        match s {
319            "shared" => Some(HistoryVisibility::Shared),
320            "world_readable" => Some(HistoryVisibility::WorldReadable),
321            "invited" => Some(HistoryVisibility::Invited),
322            "joined" => Some(HistoryVisibility::Joined),
323            _ => None,
324        }
325    }
326}
327
328#[derive(Debug, Clone, Copy, PartialEq, Eq)]
329pub enum Membership {
330    Join,
331    Invite,
332    Leave,
333    Ban,
334}
335
336impl Membership {
337    pub fn as_str(self) -> &'static str {
338        match self {
339            Membership::Join => "join",
340            Membership::Invite => "invite",
341            Membership::Leave => "leave",
342            Membership::Ban => "ban",
343        }
344    }
345
346    pub fn from_wire_name(s: &str) -> Option<Self> {
347        match s {
348            "join" => Some(Membership::Join),
349            "invite" => Some(Membership::Invite),
350            "leave" => Some(Membership::Leave),
351            "ban" => Some(Membership::Ban),
352            _ => None,
353        }
354    }
355}
356
357#[derive(Debug, Clone, Copy, PartialEq, Eq)]
358pub enum ReceiptType {
359    Read,
360    ReadPrivate,
361}
362
363impl ReceiptType {
364    pub fn as_str(self) -> &'static str {
365        match self {
366            ReceiptType::Read => "m.read",
367            ReceiptType::ReadPrivate => "m.read.private",
368        }
369    }
370
371    pub fn from_wire_name(s: &str) -> Option<Self> {
372        match s {
373            "m.read" => Some(ReceiptType::Read),
374            "m.read.private" => Some(ReceiptType::ReadPrivate),
375            _ => None,
376        }
377    }
378}
379
380/// Decode a TEXT column this module itself always writes from one of the
381/// enums above — a `None` from `parse` means the row was written by
382/// something other than this module's own CRUD, which is a data-integrity
383/// bug, not a normal "not found" case, hence `InvalidColumnType` rather
384/// than silently defaulting.
385fn decode_enum<T>(idx: usize, column: &'static str, raw: &str, parse: fn(&str) -> Option<T>) -> rusqlite::Result<T> {
386    parse(raw).ok_or_else(|| rusqlite::Error::InvalidColumnType(idx, column.to_string(), rusqlite::types::Type::Text))
387}
388
389// ============================================================================
390// Schema
391// ============================================================================
392
393/// Create every table/index this module needs — idempotent (`IF NOT
394/// EXISTS`/`INSERT OR IGNORE` throughout), so it runs on every boot. Called
395/// from [`init_messenger_db`] and directly by this module's own tests
396/// against a plain in-memory connection (no SQLCipher key needed for schema
397/// creation) — same split `social_db::create_social_schema` uses.
398///
399/// Deviation from the plan's literal DDL text (§2, both noted inline as
400/// `-- DEVIATION`): `account_data.room_id` drops its `REFERENCES rooms(id)`
401/// — the manager's own correction makes `room_id` `NOT NULL DEFAULT ''` for
402/// global account data, and `''` never matches a real room id, so keeping
403/// the foreign key would make every global upsert fail once `PRAGMA
404/// foreign_keys=ON` is set (as [`init_messenger_db`] does).
405pub fn create_matrix_schema(conn: &Connection) -> rusqlite::Result<()> {
406    conn.execute_batch(
407        r#"
408        -- Global stream ordering — see this module's doc comment on the
409        -- single-writer guarantee that makes `UPDATE ... RETURNING` safe.
410        CREATE TABLE IF NOT EXISTS stream_counter (
411            id    INTEGER PRIMARY KEY CHECK (id = 1),
412            value INTEGER NOT NULL
413        );
414        INSERT OR IGNORE INTO stream_counter (id, value) VALUES (1, 0);
415
416        -- Federation F0: this server's signing keys and the verify keys
417        -- cached from remote servers. Secrets live only in the (encrypted) DB.
418        CREATE TABLE IF NOT EXISTS fed_signing_keys (
419            key_id     TEXT PRIMARY KEY,
420            secret     BLOB NOT NULL,
421            created_ms INTEGER NOT NULL,
422            retired_ms INTEGER
423        );
424        CREATE TABLE IF NOT EXISTS fed_remote_keys (
425            server_name    TEXT NOT NULL,
426            key_id         TEXT NOT NULL,
427            public_key     TEXT NOT NULL,
428            valid_until_ms INTEGER NOT NULL,
429            fetched_ms     INTEGER NOT NULL,
430            PRIMARY KEY (server_name, key_id)
431        );
432
433        -- user_id -> mxid, filled on first touch by ensure_matrix_user
434        -- (plan §2 manager decision: new table, not in the original DDL
435        -- text). public_id is immutable, so this mapping never changes.
436        CREATE TABLE IF NOT EXISTS matrix_users (
437            user_id    INTEGER PRIMARY KEY,
438            mxid       TEXT NOT NULL UNIQUE,
439            created_at TEXT NOT NULL,
440            nick       TEXT
441        );
442        CREATE UNIQUE INDEX IF NOT EXISTS idx_matrix_users_nick_lower
443            ON matrix_users(LOWER(nick)) WHERE nick IS NOT NULL;
444
445        -- Nick belongs to a session, not to matrix_users and not to the
446        -- device. device_id is only a mark. One device may have many sessions.
447        -- matrix_users.nick stays for old databases and is not read.
448        CREATE TABLE IF NOT EXISTS messenger_sessions (
449            session_id TEXT PRIMARY KEY,
450            user_id    INTEGER NOT NULL,
451            device_id  TEXT NOT NULL,
452            nick       TEXT NOT NULL
453        );
454        CREATE INDEX IF NOT EXISTS idx_messenger_sessions_user
455            ON messenger_sessions(user_id);
456        CREATE UNIQUE INDEX IF NOT EXISTS idx_messenger_sessions_nick_lower
457            ON messenger_sessions(LOWER(nick));
458
459        CREATE TABLE IF NOT EXISTS rooms (
460            id                  TEXT PRIMARY KEY,
461            kind                TEXT NOT NULL,
462            room_version        TEXT NOT NULL DEFAULT '11',
463            creator_user_id     INTEGER NOT NULL,
464            created_at          TEXT NOT NULL,
465            is_encrypted        INTEGER NOT NULL DEFAULT 0,
466            join_rule           TEXT NOT NULL DEFAULT 'invite',
467            history_visibility  TEXT NOT NULL DEFAULT 'shared',
468            dm_pair_key         TEXT UNIQUE,
469            legacy_dm_id        INTEGER UNIQUE
470        );
471        CREATE INDEX IF NOT EXISTS idx_rooms_kind ON rooms(kind);
472
473        CREATE TABLE IF NOT EXISTS events (
474            stream_id        INTEGER PRIMARY KEY,
475            event_id         TEXT NOT NULL UNIQUE,
476            room_id          TEXT NOT NULL REFERENCES rooms(id),
477            sender_user_id   INTEGER NOT NULL,
478            event_type       TEXT NOT NULL,
479            state_key        TEXT,
480            content          TEXT NOT NULL,
481            origin_server_ts INTEGER NOT NULL,
482            txn_id           TEXT,
483            redacts          TEXT REFERENCES events(event_id),
484            redacted_by      TEXT REFERENCES events(event_id)
485        );
486        CREATE INDEX IF NOT EXISTS idx_events_room_stream ON events(room_id, stream_id);
487        CREATE INDEX IF NOT EXISTS idx_events_room_type_state ON events(room_id, event_type, state_key);
488        CREATE INDEX IF NOT EXISTS idx_events_sender ON events(sender_user_id, stream_id);
489
490        CREATE TABLE IF NOT EXISTS current_state (
491            room_id    TEXT NOT NULL REFERENCES rooms(id),
492            event_type TEXT NOT NULL,
493            state_key  TEXT NOT NULL,
494            event_id   TEXT NOT NULL REFERENCES events(event_id),
495            PRIMARY KEY (room_id, event_type, state_key)
496        );
497
498        CREATE TABLE IF NOT EXISTS room_members (
499            room_id     TEXT NOT NULL REFERENCES rooms(id),
500            user_id     INTEGER NOT NULL,
501            membership  TEXT NOT NULL,
502            power_level INTEGER,
503            updated_at  TEXT NOT NULL,
504            PRIMARY KEY (room_id, user_id)
505        );
506        CREATE INDEX IF NOT EXISTS idx_room_members_user ON room_members(user_id, membership);
507
508        CREATE TABLE IF NOT EXISTS relations (
509            event_id  TEXT PRIMARY KEY REFERENCES events(event_id),
510            room_id   TEXT NOT NULL REFERENCES rooms(id),
511            rel_type  TEXT NOT NULL,
512            target_id TEXT NOT NULL REFERENCES events(event_id),
513            agg_key   TEXT
514        );
515        CREATE INDEX IF NOT EXISTS idx_relations_target ON relations(target_id, rel_type);
516
517        CREATE TABLE IF NOT EXISTS receipts (
518            room_id      TEXT NOT NULL REFERENCES rooms(id),
519            user_id      INTEGER NOT NULL,
520            receipt_type TEXT NOT NULL,
521            event_id     TEXT NOT NULL REFERENCES events(event_id),
522            ts           INTEGER NOT NULL,
523            stream_id    INTEGER NOT NULL,
524            PRIMARY KEY (room_id, user_id, receipt_type)
525        );
526        CREATE INDEX IF NOT EXISTS idx_receipts_room_stream ON receipts(room_id, stream_id);
527
528        -- DEVIATION from plan §2's literal text: room_id has no
529        -- `REFERENCES rooms(id)` — see this function's doc comment.
530        CREATE TABLE IF NOT EXISTS account_data (
531            user_id   INTEGER NOT NULL,
532            room_id   TEXT NOT NULL DEFAULT '',
533            data_type TEXT NOT NULL,
534            content   TEXT NOT NULL,
535            stream_id INTEGER NOT NULL,
536            PRIMARY KEY (user_id, room_id, data_type)
537        );
538        CREATE INDEX IF NOT EXISTS idx_account_data_user_stream ON account_data(user_id, stream_id);
539
540        CREATE TABLE IF NOT EXISTS txn_dedup (
541            user_id    INTEGER NOT NULL,
542            device_id  TEXT NOT NULL,
543            txn_id     TEXT NOT NULL,
544            event_id   TEXT REFERENCES events(event_id),
545            created_at TEXT NOT NULL,
546            PRIMARY KEY (user_id, device_id, txn_id)
547        );
548
549        CREATE TABLE IF NOT EXISTS filters (
550            id         INTEGER PRIMARY KEY AUTOINCREMENT,
551            user_id    INTEGER NOT NULL,
552            definition TEXT NOT NULL
553        );
554
555        -- legacy_dm_message_map removed (M2); drop_legacy_dm_scaffold_if_empty cleans old DBs
556        "#,
557    )?;
558    // Public plaintext store: own tables, created beside (never inside) the closed set.
559    crate::public_channels::create_public_schema(conn)?;
560    crate::public_forum::create_forum_schema(conn)?;
561    crate::media::create_media_schema(conn)?;
562    crate::fed_rooms::create_fed_schema(conn)?;
563    crate::dag_schema::create_dag_schema(conn)?;
564    crate::http::extras::create_schema(conn)?;
565    crate::identities::create_identities_schema(conn)
566}
567
568/// Parses the 32-byte database key given as 64 hex characters (the raw SQLCipher key).
569pub fn parse_db_key(key_hex: &str) -> Result<[u8; 32], String> {
570    let h = key_hex.trim();
571    if h.len() != 64 || !h.bytes().all(|b| b.is_ascii_hexdigit()) {
572        return Err("database key must be 64 hex characters (32 bytes)".into());
573    }
574    let mut key = [0u8; 32];
575    for (i, b) in key.iter_mut().enumerate() {
576        *b = u8::from_str_radix(&h[2 * i..2 * i + 2], 16).map_err(|_| "database key must be hex".to_string())?;
577    }
578    Ok(key)
579}
580
581/// The SQLCipher store configuration for `path` and a hex key (see [`parse_db_key`]).
582pub fn messenger_db_config(path: &str, key_hex: &str) -> Result<tesserax_store::DbConfig, String> {
583    let key = parse_db_key(key_hex)?;
584    Ok(tesserax_store::DbConfig::encrypted_native(path, std::sync::Arc::new(tesserax_store::keysource::StaticKeySource(key))))
585}
586
587/// Opens the messenger store (tesserax-store: SQLCipher, WAL, one writer) and makes sure every
588/// table exists. The key is applied first on every connection the engine opens.
589pub fn open_messenger_db(path: &str, key_hex: &str) -> Result<tesserax_store::Db, String> {
590    let cfg = messenger_db_config(path, key_hex)?;
591    let db = tesserax_store::Db::open(&cfg).map_err(|e| e.to_string())?;
592    db.blocking(|conn| ensure_schema(conn)).map_err(|e| e.to_string())?;
593    Ok(db)
594}
595
596/// Parallel read-only connections (same file, same key) for the store opened by [`open_messenger_db`].
597pub fn open_read_pool(path: &str, key_hex: &str, size: usize) -> Result<tesserax_store::ReadPool, String> {
598    let cfg = messenger_db_config(path, key_hex)?;
599    tesserax_store::ReadPoolConfig::from_config(cfg).pool_size(size.max(1)).open().map_err(|e| e.to_string())
600}
601
602/// Creates every messenger table that does not exist yet (idempotent).
603pub fn ensure_schema(conn: &Connection) -> rusqlite::Result<()> {
604    create_matrix_schema(conn)?;
605    crate::keys::create_matrix_keys_schema(conn)?;
606    crate::retention::create_retention_schema(conn)?;
607    crate::public_channels::create_public_schema(conn)?;
608    Ok(())
609}
610
611// ============================================================================
612// Stream ordering
613// ============================================================================
614
615/// Take the next value from the single global counter, inside `tx` — always
616/// called in the SAME transaction as the row it stamps (this module's one
617/// hard invariant; see the module doc's single-writer note on why this
618/// stays race-free with no additional locking). `pub(crate)` so
619/// `matrix_keys_store.rs` can stamp its own stream-ordered tables
620/// (`to_device_messages`, `device_list_changes`) off this same counter
621/// rather than minting a second one.
622pub(crate) fn next_stream_id(tx: &Transaction) -> rusqlite::Result<i64> {
623    tx.query_row("UPDATE stream_counter SET value = value + 1 WHERE id = 1 RETURNING value", [], |row| row.get(0))
624}
625
626/// The last stream id issued to any table — the counter's current value.
627/// Since every stream-ordered table shares this one counter, this is also
628/// the highest `stream_id` that exists anywhere in the database right now.
629pub fn max_stream_id(conn: &Connection) -> rusqlite::Result<i64> {
630    conn.query_row("SELECT value FROM stream_counter WHERE id = 1", [], |row| row.get(0))
631}
632
633// ============================================================================
634// Matrix users (user_id <-> mxid)
635// ============================================================================
636
637/// Register `user_id`'s mxid on first touch (idempotent — a later call for
638/// the same `user_id` is a no-op) and return it. Refuses to mint a
639/// [reserved][is_reserved_localpart] localpart.
640pub fn ensure_matrix_user(conn: &Connection, user_id: i64, public_id: &str, now: &str) -> Result<String, MatrixStoreError> {
641    if is_reserved_localpart(public_id) {
642        return Err(MatrixStoreError::ReservedLocalpart);
643    }
644    let mxid = mxid_for_public_id(public_id);
645    conn.execute(
646        "INSERT INTO matrix_users (user_id, mxid, created_at) VALUES (?1, ?2, ?3)
647         ON CONFLICT(user_id) DO NOTHING",
648        params![user_id, mxid, now],
649    )?;
650    Ok(mxid)
651}
652
653/// `user_id`'s mxid, if [`ensure_matrix_user`] has ever run for it.
654pub fn mxid_of(conn: &Connection, user_id: i64) -> rusqlite::Result<Option<String>> {
655    conn.query_row("SELECT mxid FROM matrix_users WHERE user_id = ?1", params![user_id], |row| row.get(0))
656        .optional()
657}
658
659/// The internal `user_id` behind `mxid`, if known.
660pub fn user_id_of(conn: &Connection, mxid: &str) -> rusqlite::Result<Option<i64>> {
661    conn.query_row("SELECT user_id FROM matrix_users WHERE mxid = ?1", params![mxid], |row| row.get(0))
662        .optional()
663}
664
665// ============================================================================
666// Rooms
667// ============================================================================
668
669#[derive(Debug, Clone, PartialEq)]
670pub struct Room {
671    pub id: String,
672    pub kind: RoomKind,
673    pub room_version: String,
674    pub creator_user_id: i64,
675    pub created_at: String,
676    pub is_encrypted: bool,
677    pub join_rule: JoinRule,
678    pub history_visibility: HistoryVisibility,
679    pub dm_pair_key: Option<String>,
680    pub legacy_dm_id: Option<i64>,
681}
682
683const ROOM_SELECT_COLUMNS: &str =
684    "id, kind, room_version, creator_user_id, created_at, is_encrypted, join_rule, history_visibility, dm_pair_key, legacy_dm_id";
685
686fn room_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Room> {
687    let kind_raw: String = row.get(1)?;
688    let join_rule_raw: String = row.get(6)?;
689    let history_visibility_raw: String = row.get(7)?;
690    Ok(Room {
691        id: row.get(0)?,
692        kind: decode_enum(1, "kind", &kind_raw, RoomKind::from_wire_name)?,
693        room_version: row.get(2)?,
694        creator_user_id: row.get(3)?,
695        created_at: row.get(4)?,
696        is_encrypted: row.get(5)?,
697        join_rule: decode_enum(6, "join_rule", &join_rule_raw, JoinRule::from_wire_name)?,
698        history_visibility: decode_enum(7, "history_visibility", &history_visibility_raw, HistoryVisibility::from_wire_name)?,
699        dm_pair_key: row.get(8)?,
700        legacy_dm_id: row.get(9)?,
701    })
702}
703
704/// The shared core of [`create_room`]/[`create_room_with_state`]: one INSERT
705/// into `rooms`. Takes a plain `&Connection` — callable as `insert_room_row(conn, ...)`
706/// from the standalone [`create_room`] or as `insert_room_row(&tx, ...)` from
707/// inside [`create_room_with_state`]'s transaction (`rusqlite::Transaction`
708/// derefs to `Connection`, so the same function body serves both, per this
709/// module's "single-event API and the batch share code" rule).
710#[allow(clippy::too_many_arguments)]
711fn insert_room_row(
712    conn: &Connection,
713    room_id: &str,
714    kind: RoomKind,
715    creator_user_id: i64,
716    created_at: &str,
717    is_encrypted: bool,
718    join_rule: JoinRule,
719    history_visibility: HistoryVisibility,
720    dm_pair_key: Option<&str>,
721    legacy_dm_id: Option<i64>,
722) -> rusqlite::Result<()> {
723    conn.execute(
724        "INSERT INTO rooms (id, kind, room_version, creator_user_id, created_at, is_encrypted, join_rule, history_visibility, dm_pair_key, legacy_dm_id)
725         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
726        params![
727            room_id,
728            kind.as_str(),
729            MATRIX_ROOM_VERSION,
730            creator_user_id,
731            created_at,
732            is_encrypted,
733            join_rule.as_str(),
734            history_visibility.as_str(),
735            dm_pair_key,
736            legacy_dm_id,
737        ],
738    )?;
739    Ok(())
740}
741
742/// Insert a room's metadata row only — no bootstrap state events (those go
743/// through [`apply_state_event`] separately, e.g. `m.room.create`,
744/// `m.room.member`, `m.room.power_levels`, per plan §5/§7). Kept as a
745/// standalone single-row entry point for a future non-transactional caller
746/// and for this module's own tests; [`create_room_with_state`] is the P5
747/// batch entry point everything under `routes::matrix::rooms` uses instead.
748#[allow(clippy::too_many_arguments)]
749pub fn create_room(
750    conn: &Connection,
751    room_id: &str,
752    kind: RoomKind,
753    creator_user_id: i64,
754    created_at: &str,
755    is_encrypted: bool,
756    join_rule: JoinRule,
757    history_visibility: HistoryVisibility,
758    dm_pair_key: Option<&str>,
759    legacy_dm_id: Option<i64>,
760) -> rusqlite::Result<()> {
761    insert_room_row(conn, room_id, kind, creator_user_id, created_at, is_encrypted, join_rule, history_visibility, dm_pair_key, legacy_dm_id)
762}
763
764pub fn get_room(conn: &Connection, room_id: &str) -> rusqlite::Result<Option<Room>> {
765    conn.query_row(&format!("SELECT {ROOM_SELECT_COLUMNS} FROM rooms WHERE id = ?1"), params![room_id], room_from_row)
766        .optional()
767}
768
769/// The room currently holding `pair_key` as its `dm_pair_key`, if any — the
770/// binary crate's `routes::matrix::rooms::create_room`'s DM-reuse lookup
771/// (plan P5 correction: "a second DM between the same pair returns the
772/// existing room id only if it is still a live DM for both").
773pub fn room_by_dm_pair_key(conn: &Connection, pair_key: &str) -> rusqlite::Result<Option<Room>> {
774    conn.query_row(&format!("SELECT {ROOM_SELECT_COLUMNS} FROM rooms WHERE dm_pair_key = ?1"), params![pair_key], room_from_row)
775        .optional()
776}
777
778/// The room already migrated from `dm_conversations.id = legacy_dm_id`, if
779/// any — the binary crate's own `matrix_migration` module's idempotency gate
780/// (plan §7 step 2: "For each `dm_conversations` row without a
781/// `rooms.legacy_dm_id` match").
782pub fn room_by_legacy_dm_id(conn: &Connection, legacy_dm_id: i64) -> rusqlite::Result<Option<Room>> {
783    conn.query_row(&format!("SELECT {ROOM_SELECT_COLUMNS} FROM rooms WHERE legacy_dm_id = ?1"), params![legacy_dm_id], room_from_row)
784        .optional()
785}
786
787/// Free `room_id`'s `dm_pair_key` slot (set it to `NULL`) — called when a
788/// prior DM between the same pair is no longer live (one party left), so a
789/// freshly created room between the same two users can claim that pair key
790/// without violating `rooms.dm_pair_key`'s `UNIQUE` constraint. The old room
791/// keeps every other field; only its claim on the pair key is released.
792pub fn clear_dm_pair_key(conn: &Connection, room_id: &str) -> rusqlite::Result<()> {
793    conn.execute("UPDATE rooms SET dm_pair_key = NULL WHERE id = ?1", params![room_id])?;
794    Ok(())
795}
796
797// ============================================================================
798// Events (timeline + state)
799// ============================================================================
800
801#[derive(Debug, Clone, PartialEq)]
802pub struct MatrixEvent {
803    pub stream_id: i64,
804    pub event_id: String,
805    pub room_id: String,
806    pub sender_user_id: i64,
807    pub event_type: String,
808    /// `None` for a timeline event; `Some("")` or more for a state event.
809    pub state_key: Option<String>,
810    pub content: String,
811    pub origin_server_ts: i64,
812    pub txn_id: Option<String>,
813    pub redacts: Option<String>,
814    pub redacted_by: Option<String>,
815}
816
817const EVENT_SELECT_COLUMNS: &str =
818    "stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts, txn_id, redacts, redacted_by";
819
820const EVENT_SELECT_COLUMNS_ALIASED: &str = "e.stream_id, e.event_id, e.room_id, e.sender_user_id, e.event_type, e.state_key, e.content, e.origin_server_ts, e.txn_id, e.redacts, e.redacted_by";
821
822fn event_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<MatrixEvent> {
823    Ok(MatrixEvent {
824        stream_id: row.get(0)?,
825        event_id: row.get(1)?,
826        room_id: row.get(2)?,
827        sender_user_id: row.get(3)?,
828        event_type: row.get(4)?,
829        state_key: row.get(5)?,
830        content: row.get(6)?,
831        origin_server_ts: row.get(7)?,
832        txn_id: row.get(8)?,
833        redacts: row.get(9)?,
834        redacted_by: row.get(10)?,
835    })
836}
837
838fn collect_events(rows: &mut rusqlite::Rows<'_>) -> rusqlite::Result<Vec<MatrixEvent>> {
839    let mut out = Vec::new();
840    while let Some(row) = rows.next()? {
841        out.push(event_from_row(row)?);
842    }
843    Ok(out)
844}
845
846/// The columns of one timeline event row, grouped for
847/// [`insert_timeline_event_in_tx`] (its callers all build one; the public
848/// writers spell the same fields out as separate arguments).
849#[derive(Debug, Clone, Copy)]
850pub(crate) struct TimelineEventRow<'a> {
851    pub(crate) event_id: &'a str,
852    pub(crate) room_id: &'a str,
853    pub(crate) sender_user_id: i64,
854    pub(crate) event_type: &'a str,
855    pub(crate) content: &'a str,
856    pub(crate) origin_server_ts: i64,
857    pub(crate) txn_id: Option<&'a str>,
858}
859
860/// The `&Transaction`-scoped core of [`insert_timeline_event`]. Room
861/// creation (`createRoom`) never sends a message, so [`create_room_with_state`]
862/// has no timeline events to insert and this function's only caller today is
863/// [`insert_timeline_event`] itself — it is split out anyway, matching
864/// [`apply_state_event_in_tx`]'s shape, so a future batch that DOES need to
865/// mix timeline and state events (e.g. the P12 migration's legacy-message
866/// replay) can call it directly instead of re-deriving a transaction-scoped
867/// version of this logic.
868fn insert_timeline_event_in_tx(tx: &Transaction, row: &TimelineEventRow<'_>) -> Result<MatrixEvent, MatrixStoreError> {
869    // A room of the DAG layer (feature `f3-hash-ids`) gets a signed event with a computed id,
870    // built and stored in this same transaction; every other room is untouched.
871    #[cfg(feature = "f3-hash-ids")]
872    if let Some(p) = crate::f3::prepare_local(tx, row.room_id, row.sender_user_id, row.event_type, None, row.content, row.origin_server_ts)? {
873        let row = TimelineEventRow { event_id: &p.event_id, content: &p.content, ..*row };
874        return insert_timeline_event_raw_in_tx(tx, &row);
875    }
876    insert_timeline_event_raw_in_tx(tx, row)
877}
878
879/// The plain write: the caller already chose the id (legacy rooms) or the DAG layer did.
880pub(crate) fn insert_timeline_event_raw_in_tx(tx: &Transaction, row: &TimelineEventRow<'_>) -> Result<MatrixEvent, MatrixStoreError> {
881    let TimelineEventRow { event_id, room_id, sender_user_id, event_type, content, origin_server_ts, txn_id } = *row;
882    let stream_id = next_stream_id(tx)?;
883    tx.execute(
884        "INSERT INTO events (stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts, txn_id)
885         VALUES (?1, ?2, ?3, ?4, ?5, NULL, ?6, ?7, ?8)",
886        params![stream_id, event_id, room_id, sender_user_id, event_type, content, origin_server_ts, txn_id],
887    )?;
888    populate_relations(tx, event_id, room_id, sender_user_id, content)?;
889    Ok(MatrixEvent {
890        stream_id,
891        event_id: event_id.to_string(),
892        room_id: room_id.to_string(),
893        sender_user_id,
894        event_type: event_type.to_string(),
895        state_key: None,
896        content: content.to_string(),
897        origin_server_ts,
898        txn_id: txn_id.map(str::to_string),
899        redacts: None,
900        redacted_by: None,
901    })
902}
903
904/// Insert one timeline event (`state_key` always `NULL`) inside its own
905/// transaction: mint the next stream id, insert the row, and — if
906/// `content` carries a top-level `m.relates_to` (cleartext even for
907/// `m.room.encrypted`, plan §10 item 3) — populate [`relations`]. Refusing
908/// a duplicate `m.annotation` rolls back the whole transaction, so a
909/// refused call leaves no partial `events` row. It never records a `txn_id`:
910/// a write that carries one goes through [`insert_timeline_event_deduped`],
911/// which stores it next to its dedup record in the same transaction.
912pub fn insert_timeline_event(
913    conn: &mut Connection,
914    event_id: &str,
915    room_id: &str,
916    sender_user_id: i64,
917    event_type: &str,
918    content: &str,
919    origin_server_ts: i64,
920) -> Result<MatrixEvent, MatrixStoreError> {
921    let tx = conn.transaction()?;
922    let row = TimelineEventRow { event_id, room_id, sender_user_id, event_type, content, origin_server_ts, txn_id: None };
923    let event = insert_timeline_event_in_tx(&tx, &row)?;
924    tx.commit()?;
925    Ok(event)
926}
927
928/// What [`insert_timeline_event_deduped`]/[`redact_event_deduped`] found: a
929/// brand-new write, or the SAME event a prior submission of this
930/// `(user, device, txn_id)` already produced (plan §4 `send`/`redact` rows:
931/// "idempotent per (user, device, txn_id): a repeat returns the SAME
932/// event_id with no second insert").
933#[derive(Debug, Clone, PartialEq)]
934pub enum DedupedWrite {
935    New(MatrixEvent),
936    Existing(MatrixEvent),
937}
938
939/// [`insert_timeline_event`], but dedup-checked and recorded in the SAME
940/// transaction as the insert (P6 binding rule) — a repeat of
941/// `(sender_user_id, device_id, txn_id)` returns the ORIGINAL event without a
942/// second `events` row, and this is race-free even under a hypothetical
943/// second writer (never actually possible under this module's single-writer
944/// discipline) because the lookup, the insert, and the dedup record all
945/// commit or roll back together. `txn_dedup_lookup`/`txn_dedup_record` take
946/// `&Connection`; passing `&tx` (a `Transaction`) works via `Deref` — see
947/// `apply_state_event_in_tx`'s own sibling functions for the same pattern.
948#[allow(clippy::too_many_arguments)]
949pub fn insert_timeline_event_deduped(
950    conn: &mut Connection,
951    device_id: &str,
952    txn_id: &str,
953    event_id: &str,
954    room_id: &str,
955    sender_user_id: i64,
956    event_type: &str,
957    content: &str,
958    origin_server_ts: i64,
959    now: &str,
960) -> Result<DedupedWrite, MatrixStoreError> {
961    let tx = conn.transaction()?;
962    if let TxnDedupEntry::Seen(existing_event_id) = txn_dedup_lookup(&tx, sender_user_id, device_id, txn_id)? {
963        let existing_event_id = existing_event_id.ok_or_else(|| MatrixStoreError::UnknownEventId(txn_id.to_string()))?;
964        let event = get_event(&tx, &existing_event_id)?.ok_or_else(|| MatrixStoreError::UnknownEventId(existing_event_id.clone()))?;
965        tx.commit()?;
966        return Ok(DedupedWrite::Existing(event));
967    }
968    let row = TimelineEventRow { event_id, room_id, sender_user_id, event_type, content, origin_server_ts, txn_id: Some(txn_id) };
969    let event = insert_timeline_event_in_tx(&tx, &row)?;
970    txn_dedup_record(&tx, sender_user_id, device_id, txn_id, Some(&event.event_id), now)?;
971    tx.commit()?;
972    Ok(DedupedWrite::New(event))
973}
974
975/// The shared core of [`apply_state_event`] and the P5 batch entry point
976/// [`create_room_with_state`] — every bootstrap state event a room creation
977/// inserts (`m.room.create`, both `m.room.member`s, `m.room.power_levels`,
978/// ...) goes through this SAME function, once per event, all inside the
979/// batch's one transaction, so a failure on any one of them rolls back the
980/// entire room creation (plan P5 correction: "Room creation is ONE
981/// transaction").
982#[allow(clippy::too_many_arguments)]
983fn apply_state_event_in_tx(
984    tx: &Transaction,
985    event_id: &str,
986    room_id: &str,
987    sender_user_id: i64,
988    event_type: &str,
989    state_key: &str,
990    content: &str,
991    origin_server_ts: i64,
992    now: &str,
993) -> Result<MatrixEvent, MatrixStoreError> {
994    // Rooms of the DAG layer (feature `f3-hash-ids`): signed event, computed id, DAG rows and the
995    // projection below all land in this one transaction.
996    #[cfg(feature = "f3-hash-ids")]
997    if let Some(p) = crate::f3::prepare_local(tx, room_id, sender_user_id, event_type, Some(state_key), content, origin_server_ts)? {
998        return apply_state_event_raw_in_tx(tx, &p.event_id, room_id, sender_user_id, event_type, state_key, &p.content, origin_server_ts, now);
999    }
1000    apply_state_event_raw_in_tx(tx, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts, now)
1001}
1002
1003/// The plain state write (event row, `current_state` slot, membership and power caches).
1004#[allow(clippy::too_many_arguments)]
1005pub(crate) fn apply_state_event_raw_in_tx(
1006    tx: &Transaction,
1007    event_id: &str,
1008    room_id: &str,
1009    sender_user_id: i64,
1010    event_type: &str,
1011    state_key: &str,
1012    content: &str,
1013    origin_server_ts: i64,
1014    now: &str,
1015) -> Result<MatrixEvent, MatrixStoreError> {
1016    let stream_id = next_stream_id(tx)?;
1017    tx.execute(
1018        "INSERT INTO events (stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts, txn_id)
1019         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, NULL)",
1020        params![stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts],
1021    )?;
1022    tx.execute(
1023        "INSERT INTO current_state (room_id, event_type, state_key, event_id) VALUES (?1, ?2, ?3, ?4)
1024         ON CONFLICT(room_id, event_type, state_key) DO UPDATE SET event_id = excluded.event_id",
1025        params![room_id, event_type, state_key, event_id],
1026    )?;
1027    if event_type == "m.room.member" {
1028        refresh_room_member(tx, room_id, state_key, content, now)?;
1029    }
1030    if event_type == "m.room.power_levels" {
1031        refresh_power_levels(tx, room_id, content)?;
1032    }
1033    Ok(MatrixEvent {
1034        stream_id,
1035        event_id: event_id.to_string(),
1036        room_id: room_id.to_string(),
1037        sender_user_id,
1038        event_type: event_type.to_string(),
1039        state_key: Some(state_key.to_string()),
1040        content: content.to_string(),
1041        origin_server_ts,
1042        txn_id: None,
1043        redacts: None,
1044        redacted_by: None,
1045    })
1046}
1047
1048/// An event row for a state event of the past: it is in the room's history, it fills no
1049/// `current_state` slot. Used for events fetched by backfill or get_missing_events.
1050#[cfg_attr(not(feature = "f3-hash-ids"), allow(dead_code))]
1051#[allow(clippy::too_many_arguments)]
1052pub(crate) fn insert_past_state_row(tx: &Transaction, event_id: &str, room_id: &str, sender_user_id: i64, event_type: &str, state_key: &str, content: &str, origin_server_ts: i64) -> Result<(), MatrixStoreError> {
1053    let stream_id = next_stream_id(tx)?;
1054    tx.execute(
1055        "INSERT INTO events (stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts, txn_id) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, NULL)",
1056        params![stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts],
1057    )?;
1058    Ok(())
1059}
1060
1061/// Point a `current_state` slot at `event_id` (an already stored state event) and refresh the
1062/// membership and power-level caches from that event's content. Used by the DAG layer when state
1063/// resolution picks a different winner than the event that was written last.
1064#[cfg_attr(not(feature = "f3-hash-ids"), allow(dead_code))]
1065pub(crate) fn set_current_state_slot(tx: &Transaction, room_id: &str, event_type: &str, state_key: &str, event_id: &str, now: &str) -> Result<(), MatrixStoreError> {
1066    let content: String = tx.query_row("SELECT content FROM events WHERE event_id = ?1", params![event_id], |r| r.get(0))?;
1067    tx.execute(
1068        "INSERT INTO current_state (room_id, event_type, state_key, event_id) VALUES (?1, ?2, ?3, ?4)
1069         ON CONFLICT(room_id, event_type, state_key) DO UPDATE SET event_id = excluded.event_id",
1070        params![room_id, event_type, state_key, event_id],
1071    )?;
1072    if event_type == "m.room.member" {
1073        refresh_room_member(tx, room_id, state_key, &content, now)?;
1074    }
1075    if event_type == "m.room.power_levels" {
1076        refresh_power_levels(tx, room_id, &content)?;
1077    }
1078    Ok(())
1079}
1080
1081/// One state event for [`apply_state_event`]: its identity, the state slot
1082/// it fills, its content, and the two clocks (`origin_server_ts` in
1083/// milliseconds for the event row, `now` as RFC 3339 for the membership
1084/// cache's `updated_at`).
1085#[derive(Debug, Clone, Copy)]
1086pub struct StateEventWrite<'a> {
1087    pub event_id: &'a str,
1088    pub room_id: &'a str,
1089    pub sender_user_id: i64,
1090    pub event_type: &'a str,
1091    pub state_key: &'a str,
1092    pub content: &'a str,
1093    pub origin_server_ts: i64,
1094    pub now: &'a str,
1095}
1096
1097/// Insert one state event and refresh the two projections that read off it:
1098/// `current_state` (always) and, for `m.room.member`/`m.room.power_levels`,
1099/// the denormalized `room_members` cache (plan §2). All inside one
1100/// transaction with the event row.
1101pub fn apply_state_event(conn: &mut Connection, write: &StateEventWrite<'_>) -> Result<MatrixEvent, MatrixStoreError> {
1102    let tx = conn.transaction()?;
1103    let event = apply_state_event_in_tx(
1104        &tx,
1105        write.event_id,
1106        write.room_id,
1107        write.sender_user_id,
1108        write.event_type,
1109        write.state_key,
1110        write.content,
1111        write.origin_server_ts,
1112        write.now,
1113    )?;
1114    tx.commit()?;
1115    Ok(event)
1116}
1117
1118/// What one [`refresh_member_displayname`] pass did.
1119#[derive(Debug, Clone, Default, PartialEq, Eq)]
1120pub struct DisplaynameRefresh {
1121    /// How many rooms got a new `m.room.member` event for the user.
1122    pub rooms_updated: usize,
1123    /// Every `user_id` that must be woken for those new events: each updated
1124    /// room's joined and invited members, the refreshed user included.
1125    pub affected_user_ids: HashSet<i64>,
1126}
1127
1128/// Re-stamp `displayname` into `user_id`'s own `m.room.member` event in every
1129/// room where they are currently `join`ed or `invite`d, so the other
1130/// members' clients see the new name. A room whose current member event
1131/// already carries exactly `displayname` is skipped (the pass is idempotent,
1132/// which is what lets the boot backfill call it for every user on every
1133/// start). The new event keeps the current content (an invite keeps its
1134/// `is_direct` marker) and only replaces `displayname`; a `join` refresh is
1135/// sent by the user, an `invite` refresh keeps the inviter as its sender —
1136/// [`stripped_invite_state`] reads the inviter off that field. `leave`/`ban`
1137/// rows are never touched.
1138///
1139/// A user with no [`matrix_users`] row (never used the messenger) or an empty
1140/// `displayname` is a no-op. All rooms land in ONE transaction: a failure
1141/// leaves every room exactly as it was. The label itself comes from the
1142/// identity database, which this module never reads — the caller resolves it
1143/// BEFORE taking the messenger connection (lock order `db` -> `social` ->
1144/// `messenger`).
1145pub fn refresh_member_displayname(
1146    conn: &mut Connection,
1147    user_id: i64,
1148    displayname: &str,
1149    now: &str,
1150    origin_server_ts: i64,
1151) -> Result<DisplaynameRefresh, MatrixStoreError> {
1152    let mut outcome = DisplaynameRefresh::default();
1153    if displayname.is_empty() {
1154        return Ok(outcome);
1155    }
1156    let Some(mxid) = mxid_of(conn, user_id)? else {
1157        return Ok(outcome);
1158    };
1159
1160    let tx = conn.transaction()?;
1161    for membership in [Membership::Join, Membership::Invite] {
1162        for room_id in rooms_for_user(&tx, user_id, Some(membership))? {
1163            let Some(current) = current_state_event(&tx, &room_id, "m.room.member", &mxid)? else {
1164                continue;
1165            };
1166            let mut content: serde_json::Value = serde_json::from_str(&current.content)?;
1167            if content.get("displayname").and_then(|v| v.as_str()) == Some(displayname) {
1168                continue;
1169            }
1170            let Some(fields) = content.as_object_mut() else {
1171                continue;
1172            };
1173            fields.insert("displayname".to_string(), serde_json::Value::String(displayname.to_string()));
1174
1175            let sender_user_id = if membership == Membership::Join { user_id } else { current.sender_user_id };
1176            apply_state_event_in_tx(
1177                &tx,
1178                &new_event_id(),
1179                &room_id,
1180                sender_user_id,
1181                "m.room.member",
1182                &mxid,
1183                &content.to_string(),
1184                origin_server_ts,
1185                now,
1186            )?;
1187            outcome.rooms_updated += 1;
1188            outcome.affected_user_ids.insert(user_id);
1189            for member in room_members(&tx, &room_id, None)? {
1190                if matches!(member.membership, Membership::Join | Membership::Invite) {
1191                    outcome.affected_user_ids.insert(member.user_id);
1192                }
1193            }
1194        }
1195    }
1196    tx.commit()?;
1197    Ok(outcome)
1198}
1199
1200/// Every `user_id` that has a [`matrix_users`] row — the boot backfill's
1201/// worklist of users whose member events may predate the `displayname`
1202/// stamping.
1203pub fn matrix_user_ids(conn: &Connection) -> rusqlite::Result<Vec<i64>> {
1204    let mut stmt = conn.prepare("SELECT user_id FROM matrix_users ORDER BY user_id")?;
1205    let rows = stmt.query_map([], |row| row.get(0))?;
1206    rows.collect()
1207}
1208
1209// ============================================================================
1210// Room creation batch (plan P5 correction: room + every bootstrap state
1211// event, in ONE transaction)
1212// ============================================================================
1213
1214/// One state event queued for [`create_room_with_state`] — the same shape
1215/// [`apply_state_event`] takes per-field, minus `room_id`/`origin_server_ts`/
1216/// `now` (shared by the whole batch, passed once to
1217/// [`create_room_with_state`] itself rather than repeated per event).
1218#[derive(Debug, Clone, PartialEq)]
1219pub struct NewStateEvent {
1220    pub event_id: String,
1221    pub sender_user_id: i64,
1222    pub event_type: String,
1223    pub state_key: String,
1224    pub content: String,
1225}
1226
1227/// The room-metadata half of [`create_room_with_state`]'s input — every
1228/// field [`insert_room_row`] needs, grouped into one value instead of nine
1229/// positional arguments (this function's own params stay lint-clean without
1230/// an `#[allow(clippy::too_many_arguments)]`, unlike the legacy single-row
1231/// [`create_room`]/[`insert_room_row`] this batch shares its INSERT with).
1232#[derive(Debug, Clone, Copy)]
1233pub struct RoomBootstrap<'a> {
1234    pub room_id: &'a str,
1235    pub kind: RoomKind,
1236    pub creator_user_id: i64,
1237    pub created_at: &'a str,
1238    pub is_encrypted: bool,
1239    pub join_rule: JoinRule,
1240    pub history_visibility: HistoryVisibility,
1241    pub dm_pair_key: Option<&'a str>,
1242    pub legacy_dm_id: Option<i64>,
1243}
1244
1245/// Create a room AND its full bootstrap event set in ONE transaction (plan
1246/// P5 correction) — the room row ([`insert_room_row`]), then every entry of
1247/// `state_events` in order via [`apply_state_event_in_tx`] (so, e.g., the
1248/// creator's own `m.room.member` MUST precede `m.room.power_levels` in
1249/// `state_events` for [`refresh_power_levels`] to see them as an existing
1250/// member to stamp — see that function's own doc comment), then commit. A
1251/// failure on the room insert OR on any one state event leaves NEITHER a
1252/// room row NOR any event row behind (rusqlite rolls back a `Transaction`
1253/// dropped without `commit()`, the same guarantee
1254/// [`populate_relations`]'s duplicate-annotation refusal already relies on).
1255pub fn create_room_with_state(
1256    conn: &mut Connection,
1257    bootstrap: RoomBootstrap<'_>,
1258    state_events: &[NewStateEvent],
1259    origin_server_ts: i64,
1260) -> Result<(Room, Vec<MatrixEvent>), MatrixStoreError> {
1261    let tx = conn.transaction()?;
1262    insert_room_row(
1263        &tx,
1264        bootstrap.room_id,
1265        bootstrap.kind,
1266        bootstrap.creator_user_id,
1267        bootstrap.created_at,
1268        bootstrap.is_encrypted,
1269        bootstrap.join_rule,
1270        bootstrap.history_visibility,
1271        bootstrap.dm_pair_key,
1272        bootstrap.legacy_dm_id,
1273    )?;
1274    // New closed rooms are DAG rooms when the layer is compiled in; rooms created before stay legacy.
1275    // Plaintext public channels live in the public store (`pub_events`) and stay legacy.
1276    #[cfg(feature = "f3-hash-ids")]
1277    if !(bootstrap.kind == RoomKind::Channel && !bootstrap.is_encrypted) {
1278        crate::f3::mark_room(&tx, bootstrap.room_id)?;
1279    }
1280
1281    let mut applied = Vec::with_capacity(state_events.len());
1282    for event in state_events {
1283        applied.push(apply_state_event_in_tx(
1284            &tx,
1285            &event.event_id,
1286            bootstrap.room_id,
1287            event.sender_user_id,
1288            &event.event_type,
1289            &event.state_key,
1290            &event.content,
1291            origin_server_ts,
1292            bootstrap.created_at,
1293        )?);
1294    }
1295    tx.commit()?;
1296
1297    Ok((
1298        Room {
1299            id: bootstrap.room_id.to_string(),
1300            kind: bootstrap.kind,
1301            room_version: MATRIX_ROOM_VERSION.to_string(),
1302            creator_user_id: bootstrap.creator_user_id,
1303            created_at: bootstrap.created_at.to_string(),
1304            is_encrypted: bootstrap.is_encrypted,
1305            join_rule: bootstrap.join_rule,
1306            history_visibility: bootstrap.history_visibility,
1307            dm_pair_key: bootstrap.dm_pair_key.map(str::to_string),
1308            legacy_dm_id: bootstrap.legacy_dm_id,
1309        },
1310        applied,
1311    ))
1312}
1313
1314/// Refresh `room_members` off an `m.room.member` state event. `state_key`
1315/// is the member's mxid (Matrix's own wire convention) — resolved to our
1316/// internal `user_id` via [`user_id_of`]; the member must already have been
1317/// through [`ensure_matrix_user`] or this refuses with
1318/// [`MatrixStoreError::UnknownMxid`]. `power_level` is left untouched on an
1319/// existing row (only [`refresh_power_levels`] ever sets it) and starts
1320/// `NULL` (falls back to `users_default`) on a brand-new one.
1321pub(crate) fn refresh_room_member(tx: &Transaction, room_id: &str, state_key: &str, content: &str, now: &str) -> Result<(), MatrixStoreError> {
1322    let value: serde_json::Value = serde_json::from_str(content)?;
1323    let membership_str = value
1324        .get("membership")
1325        .and_then(|v| v.as_str())
1326        .ok_or_else(|| MatrixStoreError::InvalidMembership("missing 'membership' field".to_string()))?;
1327    let membership =
1328        Membership::from_wire_name(membership_str).ok_or_else(|| MatrixStoreError::InvalidMembership(membership_str.to_string()))?;
1329    let user_id = user_id_of(tx, state_key)?.ok_or_else(|| MatrixStoreError::UnknownMxid(state_key.to_string()))?;
1330    tx.execute(
1331        "INSERT INTO room_members (room_id, user_id, membership, power_level, updated_at)
1332         VALUES (?1, ?2, ?3, NULL, ?4)
1333         ON CONFLICT(room_id, user_id) DO UPDATE SET membership = excluded.membership, updated_at = excluded.updated_at",
1334        params![room_id, user_id, membership.as_str(), now],
1335    )?;
1336    Ok(())
1337}
1338
1339/// Refresh `room_members.power_level` off an `m.room.power_levels` state
1340/// event's `users` map: every current member is reset to `NULL` (falls
1341/// back to `users_default`), then every `users` entry that resolves to a
1342/// known member is applied. An entry naming a user who is not (yet) a
1343/// member is silently skipped — their level will apply once they join,
1344/// resolved fresh from the room's current `m.room.power_levels` by whatever
1345/// reads it (a later piece; this table is a cache of the last-seen event,
1346/// not re-derived from an entry that never had a member row to land on).
1347pub(crate) fn refresh_power_levels(tx: &Transaction, room_id: &str, content: &str) -> Result<(), MatrixStoreError> {
1348    let value: serde_json::Value = serde_json::from_str(content)?;
1349    tx.execute("UPDATE room_members SET power_level = NULL WHERE room_id = ?1", params![room_id])?;
1350    if let Some(users) = value.get("users").and_then(|v| v.as_object()) {
1351        for (mxid, level) in users {
1352            let Some(level) = level.as_i64() else { continue };
1353            let Some(user_id) = user_id_of(tx, mxid)? else { continue };
1354            tx.execute(
1355                "UPDATE room_members SET power_level = ?1 WHERE room_id = ?2 AND user_id = ?3",
1356                params![level, room_id, user_id],
1357            )?;
1358        }
1359    }
1360    Ok(())
1361}
1362
1363// ============================================================================
1364// Power levels (m.room.power_levels content) — plan §5, used by
1365// `routes::matrix::auth`'s PowerCheck gate. Pure `serde_json::Value` logic,
1366// no identity/messenger connection needed, hence living here (a lib-crate
1367// module with no bin-crate dependency) rather than in the binary crate's own
1368// `routes/matrix/auth.rs` alongside the rest of that module's caller/device
1369// resolution.
1370// ============================================================================
1371
1372/// An action `m.room.power_levels` gates by a named threshold field, plus
1373/// `StateDefault` for "may send an otherwise-unlisted state event type".
1374#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1375pub enum PowerAction {
1376    Invite,
1377    Kick,
1378    Ban,
1379    Redact,
1380    StateDefault,
1381}
1382
1383impl PowerAction {
1384    fn field_and_default(self) -> (&'static str, i64) {
1385        match self {
1386            PowerAction::Invite => ("invite", 50),
1387            PowerAction::Kick => ("kick", 50),
1388            PowerAction::Ban => ("ban", 50),
1389            PowerAction::Redact => ("redact", 50),
1390            PowerAction::StateDefault => ("state_default", 50),
1391        }
1392    }
1393}
1394
1395/// `power_levels.users[mxid]`, falling back to `users_default` (Matrix
1396/// default `0`) when either field is missing.
1397pub fn user_level(power_levels: &serde_json::Value, mxid: &str) -> i64 {
1398    power_levels
1399        .get("users")
1400        .and_then(|users| users.get(mxid))
1401        .and_then(serde_json::Value::as_i64)
1402        .unwrap_or_else(|| power_levels.get("users_default").and_then(serde_json::Value::as_i64).unwrap_or(0))
1403}
1404
1405/// The power level required to send an event of `event_type`:
1406/// `power_levels.events[event_type]` if present, else `state_default`
1407/// (Matrix default `50`) for a state event or `events_default` (Matrix
1408/// default `0`) for a timeline event.
1409pub fn event_level(power_levels: &serde_json::Value, event_type: &str, is_state: bool) -> i64 {
1410    if let Some(level) = power_levels.get("events").and_then(|events| events.get(event_type)).and_then(serde_json::Value::as_i64) {
1411        return level;
1412    }
1413    let (key, default) = if is_state { ("state_default", 50) } else { ("events_default", 0) };
1414    power_levels.get(key).and_then(serde_json::Value::as_i64).unwrap_or(default)
1415}
1416
1417/// Whether `mxid`'s [`user_level`] reaches the threshold `action` names
1418/// (falling back to the Matrix default when the corresponding field is
1419/// missing — see [`PowerAction::field_and_default`]).
1420pub fn can(power_levels: &serde_json::Value, action: PowerAction, mxid: &str) -> bool {
1421    let (field, default) = action.field_and_default();
1422    let required = power_levels.get(field).and_then(serde_json::Value::as_i64).unwrap_or(default);
1423    user_level(power_levels, mxid) >= required
1424}
1425
1426/// Matrix's room-v11 membership-change authorization rule for acting on
1427/// ANOTHER user (manager review, 2026-09-24 — a real gap in the first cut of
1428/// [`can`] alone): reaching the flat `action` threshold via [`can`] is
1429/// necessary but not sufficient — the sender's own level must ALSO be
1430/// STRICTLY GREATER than the target's current level, so a level-50 admin
1431/// can never kick/ban/unban a level-100 owner or an equal-level peer, even
1432/// though 50 reaches the flat kick/ban threshold. `self_leave` is the one
1433/// exception the spec carves out: a member may always change their OWN
1434/// membership to `leave` (declining an invite, or a self-kick, is just a
1435/// leave) regardless of level — pass `true` only when `new_membership` for
1436/// this call is `leave` and let the identity check below decide whether it
1437/// actually applies.
1438pub fn can_act_on(power_levels: &serde_json::Value, action: PowerAction, sender_mxid: &str, target_mxid: &str, self_leave: bool) -> bool {
1439    if self_leave && sender_mxid == target_mxid {
1440        return true;
1441    }
1442    can(power_levels, action, sender_mxid) && user_level(power_levels, sender_mxid) > user_level(power_levels, target_mxid)
1443}
1444
1445/// The default an `m.room.power_levels` scalar field resolves to when
1446/// absent — `50` for every threshold field, `0` for the two `*_default`
1447/// fields (matches [`PowerAction::field_and_default`] for the fields that
1448/// enum names, plus the two it does not).
1449fn power_levels_scalar_default(key: &str) -> i64 {
1450    match key {
1451        "events_default" | "users_default" => 0,
1452        _ => 50,
1453    }
1454}
1455
1456/// Every top-level scalar field of `m.room.power_levels` this server
1457/// enforces a change rule on (manager review, 2026-09-24). `events[type]`
1458/// entries and `notifications.room` are handled separately below (they are
1459/// nested, not top-level scalars).
1460const POWER_LEVELS_SCALAR_KEYS: [&str; 7] = ["ban", "kick", "redact", "invite", "state_default", "events_default", "users_default"];
1461
1462/// Reject a scalar field change whose OLD or NEW value exceeds `sender_level`
1463/// — shared by every scalar/nested-scalar field [`validate_power_levels_change`]
1464/// checks (the top-level keys, each `events[type]` entry, and
1465/// `notifications.room` all follow the identical "both sides must be ≤ your
1466/// own level" rule; only `users` has the different, asymmetric rule).
1467fn reject_if_either_side_exceeds(old_value: Option<i64>, new_value: Option<i64>, sender_level: i64) -> Result<(), &'static str> {
1468    if old_value != new_value && (old_value.is_some_and(|v| v > sender_level) || new_value.is_some_and(|v| v > sender_level)) {
1469        return Err("cannot change a power-level field at or above your own level");
1470    }
1471    Ok(())
1472}
1473
1474/// The `m.room.power_levels` v11 change-authorization rules beyond the flat
1475/// `PUT state` PowerCheck already gating who may send this event type at
1476/// all (manager review, 2026-09-24 — a real gap: without this, any admin
1477/// who merely reaches `state_default` could PUT themselves straight to
1478/// `100`). `sender_mxid`'s level is read from `old` (their authority BEFORE
1479/// this change takes effect):
1480///
1481/// - Every top-level scalar field ([`POWER_LEVELS_SCALAR_KEYS`]) that
1482///   CHANGES: both the old and the new (Matrix-default-resolved) value must
1483///   be at or below the sender's level.
1484/// - Every `events[type]` entry and `notifications.room` that CHANGES:
1485///   same rule, but with NO invented default for a side where the key is
1486///   simply absent (an absent side is not compared against the sender's
1487///   level at all — only a side that is actually PRESENT and exceeds it
1488///   is refused).
1489/// - `users`: for every mxid whose EFFECTIVE level ([`user_level`], which
1490///   already falls back to `users_default`) differs between `old` and
1491///   `new`: a target OTHER than the sender may never be touched if their
1492///   CURRENT (old) level is already at or above the sender's own level
1493///   (even to lower them) — a peer can never act on an equal-or-higher
1494///   peer; AND no entry's NEW level may exceed the sender's own level
1495///   (the sender may lower their own level, but never raise it, and can
1496///   never promote anyone above themselves).
1497pub fn validate_power_levels_change(old: &serde_json::Value, new: &serde_json::Value, sender_mxid: &str) -> Result<(), &'static str> {
1498    let sender_level = user_level(old, sender_mxid);
1499
1500    for key in POWER_LEVELS_SCALAR_KEYS {
1501        let default = power_levels_scalar_default(key);
1502        let old_value = old.get(key).and_then(serde_json::Value::as_i64).unwrap_or(default);
1503        let new_value = new.get(key).and_then(serde_json::Value::as_i64).unwrap_or(default);
1504        reject_if_either_side_exceeds(Some(old_value), Some(new_value), sender_level)?;
1505    }
1506
1507    let old_events = old.get("events").and_then(serde_json::Value::as_object);
1508    let new_events = new.get("events").and_then(serde_json::Value::as_object);
1509    let mut event_type_keys: std::collections::BTreeSet<&str> = std::collections::BTreeSet::new();
1510    if let Some(map) = old_events {
1511        event_type_keys.extend(map.keys().map(String::as_str));
1512    }
1513    if let Some(map) = new_events {
1514        event_type_keys.extend(map.keys().map(String::as_str));
1515    }
1516    for event_type in event_type_keys {
1517        let old_value = old_events.and_then(|m| m.get(event_type)).and_then(serde_json::Value::as_i64);
1518        let new_value = new_events.and_then(|m| m.get(event_type)).and_then(serde_json::Value::as_i64);
1519        reject_if_either_side_exceeds(old_value, new_value, sender_level)?;
1520    }
1521
1522    let old_notif_room = old.get("notifications").and_then(|v| v.get("room")).and_then(serde_json::Value::as_i64);
1523    let new_notif_room = new.get("notifications").and_then(|v| v.get("room")).and_then(serde_json::Value::as_i64);
1524    reject_if_either_side_exceeds(old_notif_room, new_notif_room, sender_level)?;
1525
1526    let old_users = old.get("users").and_then(serde_json::Value::as_object);
1527    let new_users = new.get("users").and_then(serde_json::Value::as_object);
1528    let mut user_keys: std::collections::BTreeSet<&str> = std::collections::BTreeSet::new();
1529    if let Some(map) = old_users {
1530        user_keys.extend(map.keys().map(String::as_str));
1531    }
1532    if let Some(map) = new_users {
1533        user_keys.extend(map.keys().map(String::as_str));
1534    }
1535    for target_mxid in user_keys {
1536        let old_effective = user_level(old, target_mxid);
1537        let new_effective = user_level(new, target_mxid);
1538        if old_effective == new_effective {
1539            continue;
1540        }
1541        if target_mxid != sender_mxid && old_effective >= sender_level {
1542            return Err("cannot change the level of a user at or above your own level");
1543        }
1544        if new_effective > sender_level {
1545            return Err("cannot set a user's level above your own");
1546        }
1547    }
1548
1549    Ok(())
1550}
1551
1552// ============================================================================
1553// Stripped invite-room-state (plan P5 correction: `unsigned.invite_room_state`
1554// on an invite's own `m.room.member` event, computed at read/serialize time —
1555// see this function's own doc for why, and `routes::matrix::client_event_json`
1556// (mod.rs) for its one call site today).
1557// ============================================================================
1558
1559/// One `{content, state_key, type, sender}` Matrix `StrippedStateEvent` —
1560/// `sender` already resolved to an mxid, `content` already parsed JSON.
1561/// `pub` (beyond this module's own [`stripped_invite_state`]) for the
1562/// binary crate's `routes::matrix::sync` module: `GET /sync`'s
1563/// `invite_state.events` needs the SAME stripped shape for the invitee's
1564/// OWN `m.room.member` event, which [`stripped_invite_state`] does not
1565/// itself include (see that function's own doc — it strips the INVITER's
1566/// state, not the invitee's own invite event, since its other call site,
1567/// `routes::matrix::client_event_json`'s `unsigned.invite_room_state`, is
1568/// already attaching it TO that very event).
1569pub fn stripped_state_json(conn: &Connection, event: &MatrixEvent) -> Result<serde_json::Value, MatrixStoreError> {
1570    let sender = mxid_of(conn, event.sender_user_id)?.unwrap_or_default();
1571    let content: serde_json::Value = serde_json::from_str(&event.content)?;
1572    Ok(serde_json::json!({
1573        "content": content,
1574        "state_key": event.state_key.clone().unwrap_or_default(),
1575        "type": event.event_type,
1576        "sender": sender,
1577    }))
1578}
1579
1580/// The stripped state Matrix attaches as `unsigned.invite_room_state` on an
1581/// invite's own `m.room.member` event (name/join_rules/encryption/create,
1582/// plus the inviter's own membership event) — recomputed fresh every time
1583/// it is served, NEVER stored on the member event itself: `unsigned` is not
1584/// part of the canonical, signed event, and this server's `events.content`
1585/// column holds only the canonical, signed shape. Computing it at read time
1586/// also means it always reflects the room's CURRENT state (e.g. a rename
1587/// after the invite was sent), which is what a client opening an invite it
1588/// received a while ago actually wants to see. `inviter_user_id` is the
1589/// member event's own `sender_user_id` — the caller already has it.
1590pub fn stripped_invite_state(conn: &Connection, room_id: &str, inviter_user_id: i64) -> Result<Vec<serde_json::Value>, MatrixStoreError> {
1591    let mut out = Vec::new();
1592    for event_type in ["m.room.create", "m.room.join_rules", "m.room.encryption", "m.room.name"] {
1593        if let Some(event) = current_state_event(conn, room_id, event_type, "")? {
1594            out.push(stripped_state_json(conn, &event)?);
1595        }
1596    }
1597    if let Some(inviter_mxid) = mxid_of(conn, inviter_user_id)? {
1598        if let Some(event) = current_state_event(conn, room_id, "m.room.member", &inviter_mxid)? {
1599            out.push(stripped_state_json(conn, &event)?);
1600        }
1601    }
1602    Ok(out)
1603}
1604
1605pub fn get_event(conn: &Connection, event_id: &str) -> rusqlite::Result<Option<MatrixEvent>> {
1606    let closed = conn
1607        .query_row(&format!("SELECT {EVENT_SELECT_COLUMNS} FROM events WHERE event_id = ?1"), params![event_id], event_from_row)
1608        .optional()?;
1609    if closed.is_some() {
1610        return Ok(closed);
1611    }
1612    // Public plaintext store (own tables); absent table on a pre-cut DB is "not found".
1613    match crate::public_channels::get_event(conn, event_id) {
1614        Ok(found) => Ok(found),
1615        Err(rusqlite::Error::SqliteFailure(_, Some(msg))) if msg.contains("no such table") => Ok(None),
1616        Err(e) => Err(e),
1617    }
1618}
1619
1620pub fn current_state_event(conn: &Connection, room_id: &str, event_type: &str, state_key: &str) -> rusqlite::Result<Option<MatrixEvent>> {
1621    conn.query_row(
1622        &format!(
1623            "SELECT {EVENT_SELECT_COLUMNS_ALIASED} FROM current_state cs JOIN events e ON e.event_id = cs.event_id
1624             WHERE cs.room_id = ?1 AND cs.event_type = ?2 AND cs.state_key = ?3"
1625        ),
1626        params![room_id, event_type, state_key],
1627        event_from_row,
1628    )
1629    .optional()
1630}
1631
1632pub fn current_state_all(conn: &Connection, room_id: &str) -> rusqlite::Result<Vec<MatrixEvent>> {
1633    let mut stmt = conn.prepare(&format!(
1634        "SELECT {EVENT_SELECT_COLUMNS_ALIASED} FROM current_state cs JOIN events e ON e.event_id = cs.event_id WHERE cs.room_id = ?1"
1635    ))?;
1636    let mut rows = stmt.query(params![room_id])?;
1637    collect_events(&mut rows)
1638}
1639
1640/// The point-in-time projection of every `event_type` state event as of
1641/// `at_stream_id` — the latest such event per `state_key` with
1642/// `stream_id <= at_stream_id` — for `GET /rooms/{roomId}/members?at=`
1643/// (plan §3.3: a client fetching the full member list as of a given `/sync`
1644/// token, rather than the live [`current_state_all`] projection). A
1645/// `state_key` whose first event lands AFTER `at_stream_id` is correctly
1646/// absent (it did not exist yet at that point in the room's history).
1647pub fn state_events_of_type_at(conn: &Connection, room_id: &str, event_type: &str, at_stream_id: i64) -> rusqlite::Result<Vec<MatrixEvent>> {
1648    let mut stmt = conn.prepare(&format!(
1649        "SELECT {EVENT_SELECT_COLUMNS_ALIASED} FROM events e
1650         WHERE e.room_id = ?1 AND e.event_type = ?2 AND e.state_key IS NOT NULL AND e.stream_id <= ?3
1651           AND e.stream_id = (
1652             SELECT MAX(stream_id) FROM events e2
1653             WHERE e2.room_id = e.room_id AND e2.event_type = e.event_type AND e2.state_key = e.state_key AND e2.stream_id <= ?3
1654           )"
1655    ))?;
1656    let mut rows = stmt.query(params![room_id, event_type, at_stream_id])?;
1657    collect_events(&mut rows)
1658}
1659
1660/// Every `m.room.member` state event whose CURRENT value as of
1661/// `upto_inclusive` was itself set within `(since_exclusive,
1662/// upto_inclusive]` — `routes::matrix::sync`'s §3.3 clause-(b) lazy-load
1663/// "gap rule" (matrix-spec#942): a member who joined/left/changed profile
1664/// inside a truncated timeline gap must still be reported even when they
1665/// never sent anything into the visible window. Bounded to
1666/// `upto_inclusive` throughout (never the LIVE `current_state` projection,
1667/// which could have moved past a `/sync` build's own snapshot) so this
1668/// stays part of one consistent cut.
1669pub fn member_state_changed_in_window(conn: &Connection, room_id: &str, since_exclusive: i64, upto_inclusive: i64) -> rusqlite::Result<Vec<MatrixEvent>> {
1670    let mut stmt = conn.prepare(&format!(
1671        "SELECT {EVENT_SELECT_COLUMNS_ALIASED} FROM events e
1672         WHERE e.room_id = ?1 AND e.event_type = 'm.room.member' AND e.state_key IS NOT NULL
1673           AND e.stream_id = (
1674             SELECT MAX(e2.stream_id) FROM events e2
1675             WHERE e2.room_id = e.room_id AND e2.event_type = 'm.room.member' AND e2.state_key = e.state_key AND e2.stream_id <= ?3
1676           )
1677           AND e.stream_id > ?2"
1678    ))?;
1679    let mut rows = stmt.query(params![room_id, since_exclusive, upto_inclusive])?;
1680    collect_events(&mut rows)
1681}
1682
1683/// Every state event type OTHER than `m.room.member` whose CURRENT value as
1684/// of `upto_inclusive` was itself set within `(since_exclusive,
1685/// upto_inclusive]` — `routes::matrix::sync`'s per-room state delta for
1686/// every state type this server does not lazy-load ("every OTHER
1687/// state-event type changed since `since` ... always included in full").
1688/// `since_exclusive = 0` naturally reproduces every current state event of
1689/// these types (the initial-sync/`full_state` case), since every such
1690/// event's own `stream_id` is `> 0` — one code path serves both.
1691pub fn non_member_state_changed_in_window(conn: &Connection, room_id: &str, since_exclusive: i64, upto_inclusive: i64) -> rusqlite::Result<Vec<MatrixEvent>> {
1692    let mut stmt = conn.prepare(&format!(
1693        "SELECT {EVENT_SELECT_COLUMNS_ALIASED} FROM events e
1694         WHERE e.room_id = ?1 AND e.event_type != 'm.room.member' AND e.state_key IS NOT NULL
1695           AND e.stream_id = (
1696             SELECT MAX(e2.stream_id) FROM events e2
1697             WHERE e2.room_id = e.room_id AND e2.event_type = e.event_type AND e2.state_key = e.state_key AND e2.stream_id <= ?3
1698           )
1699           AND e.stream_id > ?2"
1700    ))?;
1701    let mut rows = stmt.query(params![room_id, since_exclusive, upto_inclusive])?;
1702    collect_events(&mut rows)
1703}
1704
1705/// The oldest-first timeline window strictly after `since_stream`, capped
1706/// at `limit` — the incremental-`/sync` and forward-`/messages` shape.
1707pub fn events_in_room_after(conn: &Connection, room_id: &str, since_stream: i64, limit: i64) -> rusqlite::Result<Vec<MatrixEvent>> {
1708    let mut stmt = conn.prepare(&format!(
1709        "SELECT {EVENT_SELECT_COLUMNS} FROM events WHERE room_id = ?1 AND stream_id > ?2 ORDER BY stream_id ASC LIMIT ?3"
1710    ))?;
1711    let mut rows = stmt.query(params![room_id, since_stream, limit])?;
1712    let mut out = collect_events(&mut rows)?;
1713    if crate::public_channels::is_public_room(conn, room_id)? {
1714        out.extend(crate::public_channels::events_after(conn, room_id, since_stream, limit)?);
1715        out.sort_by_key(|e| e.stream_id);
1716        out.truncate(limit.max(0) as usize);
1717    }
1718    Ok(out)
1719}
1720
1721/// The newest-first page strictly before `before_stream`, capped at
1722/// `limit` — `GET /messages?dir=b` and the initial-`/sync` "newest N"
1723/// window (§3.2: the caller reverses it if an ascending page is needed).
1724pub fn events_in_room_before(conn: &Connection, room_id: &str, before_stream: i64, limit: i64) -> rusqlite::Result<Vec<MatrixEvent>> {
1725    let mut stmt = conn.prepare(&format!(
1726        "SELECT {EVENT_SELECT_COLUMNS} FROM events WHERE room_id = ?1 AND stream_id < ?2 ORDER BY stream_id DESC LIMIT ?3"
1727    ))?;
1728    let mut rows = stmt.query(params![room_id, before_stream, limit])?;
1729    let mut out = collect_events(&mut rows)?;
1730    if crate::public_channels::is_public_room(conn, room_id)? {
1731        out.extend(crate::public_channels::events_before(conn, room_id, before_stream, limit)?);
1732        out.sort_by_key(|e| std::cmp::Reverse(e.stream_id));
1733        out.truncate(limit.max(0) as usize);
1734    }
1735    Ok(out)
1736}
1737
1738// ============================================================================
1739// Relations (m.annotation / m.replace / m.in_reply_to / m.thread)
1740// ============================================================================
1741
1742/// Populate `relations` from `content`'s top-level `m.relates_to`, if any —
1743/// a no-op if `content` is not JSON, has no `m.relates_to`, or the relation
1744/// shape is unrecognized. Two shapes are handled: the common
1745/// `{"rel_type":..., "event_id":..., "key": "..."}` (annotation/replace/
1746/// thread) and the nested `{"m.in_reply_to": {"event_id": "..."}}` reply
1747/// form.
1748///
1749/// The target must exist and be in the SAME room as `room_id` — refused
1750/// with [`MatrixStoreError::InvalidRelationTarget`] otherwise (a relation
1751/// can never point out of its own room, and a target that does not exist
1752/// would otherwise only surface as an opaque foreign-key `Db` error from
1753/// the `INSERT` below).
1754///
1755/// A second identical `m.annotation` (same sender, target, key) is refused
1756/// with [`MatrixStoreError::DuplicateAnnotation`] — checked BEFORE insert,
1757/// so the caller's whole transaction rolls back and no partial `events` row
1758/// survives the refusal. A REDACTED prior annotation does not count: once a
1759/// reaction is redacted, its sender is free to react with the same key
1760/// again (`e.redacted_by IS NULL` in the duplicate check).
1761fn populate_relations(tx: &Transaction, event_id: &str, room_id: &str, sender_user_id: i64, content: &str) -> Result<(), MatrixStoreError> {
1762    let Ok(value) = serde_json::from_str::<serde_json::Value>(content) else {
1763        return Ok(());
1764    };
1765    let Some(relates_to) = value.get("m.relates_to") else {
1766        return Ok(());
1767    };
1768
1769    let (rel_type, target_id, agg_key): (String, String, Option<String>) = if let Some(reply) = relates_to.get("m.in_reply_to") {
1770        match reply.get("event_id").and_then(|v| v.as_str()) {
1771            Some(target) => ("m.in_reply_to".to_string(), target.to_string(), None),
1772            None => return Ok(()),
1773        }
1774    } else {
1775        let rel_type = relates_to.get("rel_type").and_then(|v| v.as_str());
1776        let target = relates_to.get("event_id").and_then(|v| v.as_str());
1777        match (rel_type, target) {
1778            (Some(rt), Some(target)) => {
1779                let key = relates_to.get("key").and_then(|v| v.as_str()).map(str::to_string);
1780                (rt.to_string(), target.to_string(), key)
1781            }
1782            _ => return Ok(()),
1783        }
1784    };
1785
1786    let target_room: Option<String> = tx
1787        .query_row("SELECT room_id FROM events WHERE event_id = ?1", params![target_id], |row| row.get(0))
1788        .optional()?;
1789    match target_room {
1790        Some(ref found_room) if found_room == room_id => {}
1791        _ => return Err(MatrixStoreError::InvalidRelationTarget(target_id)),
1792    }
1793
1794    if rel_type == "m.annotation" {
1795        let duplicate: Option<i64> = tx
1796            .query_row(
1797                "SELECT 1 FROM relations r JOIN events e ON e.event_id = r.event_id
1798                 WHERE r.target_id = ?1 AND r.rel_type = 'm.annotation' AND r.agg_key IS ?2 AND e.sender_user_id = ?3
1799                    AND e.redacted_by IS NULL
1800                 LIMIT 1",
1801                params![target_id, agg_key, sender_user_id],
1802                |row| row.get(0),
1803            )
1804            .optional()?;
1805        if duplicate.is_some() {
1806            return Err(MatrixStoreError::DuplicateAnnotation);
1807        }
1808    }
1809
1810    tx.execute(
1811        "INSERT INTO relations (event_id, room_id, rel_type, target_id, agg_key) VALUES (?1, ?2, ?3, ?4, ?5)",
1812        params![event_id, room_id, rel_type, target_id, agg_key],
1813    )?;
1814    Ok(())
1815}
1816
1817/// Every event related to `target_event_id`, newest-first (Matrix's own
1818/// default order for `GET /relations`), optionally narrowed to one
1819/// `rel_type` and/or one `event_type`, strictly before `before_stream`
1820/// (pass [`i64::MAX`] for a first page), capped at `limit` — P6's
1821/// `routes::matrix::messaging::get_relations`.
1822pub fn relations_of(
1823    conn: &Connection,
1824    target_event_id: &str,
1825    rel_type: Option<&str>,
1826    event_type: Option<&str>,
1827    before_stream: i64,
1828    limit: i64,
1829) -> rusqlite::Result<Vec<MatrixEvent>> {
1830    // Anonymous `?` placeholders (not `?N`) — purely positional, so building
1831    // the SQL text and the bound-value list in the SAME order below can
1832    // never desync the way explicit numbered placeholders would if a
1833    // variant's text skipped a number.
1834    let mut sql = format!("SELECT {EVENT_SELECT_COLUMNS_ALIASED} FROM relations r JOIN events e ON e.event_id = r.event_id WHERE r.target_id = ? AND e.stream_id < ?");
1835    let mut values: Vec<&dyn rusqlite::ToSql> = vec![&target_event_id, &before_stream];
1836    if let Some(rt) = &rel_type {
1837        sql.push_str(" AND r.rel_type = ?");
1838        values.push(rt);
1839    }
1840    if let Some(et) = &event_type {
1841        sql.push_str(" AND e.event_type = ?");
1842        values.push(et);
1843    }
1844    sql.push_str(" ORDER BY e.stream_id DESC LIMIT ?");
1845    values.push(&limit);
1846
1847    let mut stmt = conn.prepare(&sql)?;
1848    let mut rows = stmt.query(values.as_slice())?;
1849    collect_events(&mut rows)
1850}
1851
1852// ============================================================================
1853// Redaction (room v11 content allow-list)
1854// ============================================================================
1855
1856/// Strip `content` to the room-version-11 redaction allow-list for
1857/// `event_type` (Matrix Client-Server API, room version 11 redaction
1858/// algorithm — `spec.matrix.org/v1.12/rooms/v11/#redactions`, cited by plan
1859/// §1/§2): `m.room.create` keeps everything; `m.room.member` keeps
1860/// `membership`, `join_authorised_via_users_server`, and
1861/// `third_party_invite.signed` (nested — only `signed` survives inside
1862/// `third_party_invite`); `m.room.join_rules` keeps `join_rule`, `allow`;
1863/// `m.room.power_levels` keeps `ban`, `events`, `events_default`,
1864/// `invite`, `kick`, `redact`, `state_default`, `users`, `users_default`;
1865/// `m.room.history_visibility` keeps `history_visibility`; every other
1866/// type is stripped to `{}`.
1867fn redact_content_per_v11(event_type: &str, content: &str) -> Result<String, MatrixStoreError> {
1868    let value: serde_json::Value = serde_json::from_str(content)?;
1869    let obj = value.as_object().cloned().unwrap_or_default();
1870    let mut kept = serde_json::Map::new();
1871
1872    match event_type {
1873        "m.room.create" => kept = obj,
1874        "m.room.member" => {
1875            for key in ["membership", "join_authorised_via_users_server"] {
1876                if let Some(v) = obj.get(key) {
1877                    kept.insert(key.to_string(), v.clone());
1878                }
1879            }
1880            if let Some(signed) = obj.get("third_party_invite").and_then(|v| v.get("signed")) {
1881                let mut third_party_invite = serde_json::Map::new();
1882                third_party_invite.insert("signed".to_string(), signed.clone());
1883                kept.insert("third_party_invite".to_string(), serde_json::Value::Object(third_party_invite));
1884            }
1885        }
1886        "m.room.join_rules" => {
1887            for key in ["join_rule", "allow"] {
1888                if let Some(v) = obj.get(key) {
1889                    kept.insert(key.to_string(), v.clone());
1890                }
1891            }
1892        }
1893        "m.room.power_levels" => {
1894            for key in [
1895                "ban",
1896                "events",
1897                "events_default",
1898                "invite",
1899                "kick",
1900                "redact",
1901                "state_default",
1902                "users",
1903                "users_default",
1904            ] {
1905                if let Some(v) = obj.get(key) {
1906                    kept.insert(key.to_string(), v.clone());
1907                }
1908            }
1909        }
1910        "m.room.history_visibility" => {
1911            if let Some(v) = obj.get("history_visibility") {
1912                kept.insert("history_visibility".to_string(), v.clone());
1913            }
1914        }
1915        _ => {}
1916    }
1917
1918    Ok(serde_json::Value::Object(kept).to_string())
1919}
1920
1921/// Redact `target_event_id`: store the `m.room.redaction` event itself
1922/// (`redacts = target_event_id`, and `content.redacts = target_event_id`,
1923/// the room-version-11 location clients read), set the target's
1924/// `redacted_by`, and strip the target's `content` per
1925/// [`redact_content_per_v11`]. All inside one transaction.
1926///
1927/// Refuses with:
1928/// - [`MatrixStoreError::UnknownEventId`] if `target_event_id` does not
1929///   exist;
1930/// - [`MatrixStoreError::WrongRoom`] if it exists but is in a different
1931///   room than `room_id`;
1932/// - [`MatrixStoreError::UnredactableEvent`] for `m.room.create` or
1933///   `m.room.encryption` — `m.room.create` is the room's own identity
1934///   anchor (its v11 allow-list already keeps every key, so "redacting" it
1935///   would be a no-op at best), and `m.room.encryption`'s v11 allow-list
1936///   has NO keys at all, meaning a redacted encryption event strips to
1937///   `{}` — indistinguishable from an unencrypted room to anything reading
1938///   current state. This server's rule is that encryption never turns off
1939///   once set (`rooms.is_encrypted`'s own doc comment says the same), so
1940///   the event that turned it on must never become redactable.
1941fn redact_event_in_tx(
1942    tx: &Transaction,
1943    room_id: &str,
1944    target_event_id: &str,
1945    redaction_event_id: &str,
1946    sender_user_id: i64,
1947    reason: Option<&str>,
1948    origin_server_ts: i64,
1949) -> Result<MatrixEvent, MatrixStoreError> {
1950    let target: Option<(String, String, String)> = tx
1951        .query_row(
1952            "SELECT event_type, content, room_id FROM events WHERE event_id = ?1",
1953            params![target_event_id],
1954            |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
1955        )
1956        .optional()?;
1957    let (target_type, target_content, target_room) = target.ok_or_else(|| MatrixStoreError::UnknownEventId(target_event_id.to_string()))?;
1958    if target_room != room_id {
1959        return Err(MatrixStoreError::WrongRoom(target_event_id.to_string()));
1960    }
1961    if target_type == "m.room.create" || target_type == "m.room.encryption" {
1962        return Err(MatrixStoreError::UnredactableEvent(target_type));
1963    }
1964
1965    let stream_id = next_stream_id(tx)?;
1966    // Room version 11 carries the redaction's target in `content.redacts`
1967    // (the `redacts` column keeps serving the older top-level shape).
1968    let mut redaction_content = serde_json::json!({ "redacts": target_event_id });
1969    if let Some(r) = reason {
1970        redaction_content["reason"] = serde_json::Value::String(r.to_string());
1971    }
1972    let redaction_content = redaction_content.to_string();
1973    tx.execute(
1974        "INSERT INTO events (stream_id, event_id, room_id, sender_user_id, event_type, state_key, content, origin_server_ts, txn_id, redacts)
1975         VALUES (?1, ?2, ?3, ?4, 'm.room.redaction', NULL, ?5, ?6, NULL, ?7)",
1976        params![stream_id, redaction_event_id, room_id, sender_user_id, redaction_content, origin_server_ts, target_event_id],
1977    )?;
1978
1979    let stripped_content = redact_content_per_v11(&target_type, &target_content)?;
1980    tx.execute(
1981        "UPDATE events SET content = ?1, redacted_by = ?2 WHERE event_id = ?3",
1982        params![stripped_content, redaction_event_id, target_event_id],
1983    )?;
1984
1985    Ok(MatrixEvent {
1986        stream_id,
1987        event_id: redaction_event_id.to_string(),
1988        room_id: room_id.to_string(),
1989        sender_user_id,
1990        event_type: "m.room.redaction".to_string(),
1991        state_key: None,
1992        content: redaction_content,
1993        origin_server_ts,
1994        txn_id: None,
1995        redacts: Some(target_event_id.to_string()),
1996        redacted_by: None,
1997    })
1998}
1999
2000pub fn redact_event(
2001    conn: &mut Connection,
2002    room_id: &str,
2003    target_event_id: &str,
2004    redaction_event_id: &str,
2005    sender_user_id: i64,
2006    reason: Option<&str>,
2007    origin_server_ts: i64,
2008) -> Result<MatrixEvent, MatrixStoreError> {
2009    let tx = conn.transaction()?;
2010    let event = redact_event_in_tx(&tx, room_id, target_event_id, redaction_event_id, sender_user_id, reason, origin_server_ts)?;
2011    tx.commit()?;
2012    Ok(event)
2013}
2014
2015/// The redaction [`redact_event_marked`] writes: which event, in which room,
2016/// by whom, under which new event id.
2017#[derive(Debug, Clone, Copy)]
2018pub struct Redaction<'a> {
2019    pub room_id: &'a str,
2020    pub target_event_id: &'a str,
2021    pub redaction_event_id: &'a str,
2022    pub sender_user_id: i64,
2023    pub reason: Option<&'a str>,
2024    pub origin_server_ts: i64,
2025}
2026
2027/// [`redact_event`], but merges `extra_content` into the resulting
2028/// `m.room.redaction` event's own content, in the SAME transaction as the
2029/// redaction itself — `routes::matrix::moderation`'s (P11) ONLY caller,
2030/// which stamps `{"org.example.site_moderation": true}` on a
2031/// site-moderator's redaction of a public-room message (plan §4's
2032/// `/api/matrix-admin/rooms/{roomId}/moderate` row) so every client can
2033/// render it distinctly from an ordinary member-initiated redaction. Every
2034/// other redaction path ([`redact_event`], [`redact_event_deduped`]) never
2035/// needs this and stays untouched.
2036pub fn redact_event_marked(
2037    conn: &mut Connection,
2038    redaction: &Redaction<'_>,
2039    extra_content: &serde_json::Value,
2040) -> Result<MatrixEvent, MatrixStoreError> {
2041    let Redaction { room_id, target_event_id, redaction_event_id, sender_user_id, reason, origin_server_ts } = *redaction;
2042    let tx = conn.transaction()?;
2043    let mut event = redact_event_in_tx(&tx, room_id, target_event_id, redaction_event_id, sender_user_id, reason, origin_server_ts)?;
2044    let mut content: serde_json::Value = serde_json::from_str(&event.content)?;
2045    if let (Some(content_obj), Some(extra_obj)) = (content.as_object_mut(), extra_content.as_object()) {
2046        for (key, value) in extra_obj {
2047            content_obj.insert(key.clone(), value.clone());
2048        }
2049    }
2050    let content_str = content.to_string();
2051    tx.execute("UPDATE events SET content = ?1 WHERE event_id = ?2", params![content_str, redaction_event_id])?;
2052    event.content = content_str;
2053    tx.commit()?;
2054    Ok(event)
2055}
2056
2057/// [`redact_event`], but dedup-checked and recorded in the SAME transaction
2058/// as the redaction (P6 binding rule, mirroring
2059/// [`insert_timeline_event_deduped`] — see its own doc for the race-freedom
2060/// argument).
2061#[allow(clippy::too_many_arguments)]
2062pub fn redact_event_deduped(
2063    conn: &mut Connection,
2064    device_id: &str,
2065    txn_id: &str,
2066    room_id: &str,
2067    target_event_id: &str,
2068    redaction_event_id: &str,
2069    sender_user_id: i64,
2070    reason: Option<&str>,
2071    origin_server_ts: i64,
2072    now: &str,
2073) -> Result<DedupedWrite, MatrixStoreError> {
2074    let tx = conn.transaction()?;
2075    if let TxnDedupEntry::Seen(existing_event_id) = txn_dedup_lookup(&tx, sender_user_id, device_id, txn_id)? {
2076        let existing_event_id = existing_event_id.ok_or_else(|| MatrixStoreError::UnknownEventId(txn_id.to_string()))?;
2077        let event = get_event(&tx, &existing_event_id)?.ok_or_else(|| MatrixStoreError::UnknownEventId(existing_event_id.clone()))?;
2078        tx.commit()?;
2079        return Ok(DedupedWrite::Existing(event));
2080    }
2081    let event = redact_event_in_tx(&tx, room_id, target_event_id, redaction_event_id, sender_user_id, reason, origin_server_ts)?;
2082    txn_dedup_record(&tx, sender_user_id, device_id, txn_id, Some(redaction_event_id), now)?;
2083    tx.commit()?;
2084    Ok(DedupedWrite::New(event))
2085}
2086
2087// ============================================================================
2088// Membership projection
2089// ============================================================================
2090
2091#[derive(Debug, Clone, PartialEq)]
2092pub struct RoomMember {
2093    pub room_id: String,
2094    pub user_id: i64,
2095    pub membership: Membership,
2096    pub power_level: Option<i64>,
2097    pub updated_at: String,
2098}
2099
2100fn room_member_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<RoomMember> {
2101    let membership_raw: String = row.get(2)?;
2102    Ok(RoomMember {
2103        room_id: row.get(0)?,
2104        user_id: row.get(1)?,
2105        membership: decode_enum(2, "membership", &membership_raw, Membership::from_wire_name)?,
2106        power_level: row.get(3)?,
2107        updated_at: row.get(4)?,
2108    })
2109}
2110
2111const ROOM_MEMBER_SELECT_COLUMNS: &str = "room_id, user_id, membership, power_level, updated_at";
2112
2113/// Every member row of `room_id`, optionally filtered to one `membership`.
2114pub fn room_members(conn: &Connection, room_id: &str, membership: Option<Membership>) -> rusqlite::Result<Vec<RoomMember>> {
2115    match membership {
2116        Some(m) => {
2117            let mut stmt = conn.prepare(&format!(
2118                "SELECT {ROOM_MEMBER_SELECT_COLUMNS} FROM room_members WHERE room_id = ?1 AND membership = ?2"
2119            ))?;
2120            let rows = stmt.query_map(params![room_id, m.as_str()], room_member_from_row)?;
2121            rows.collect()
2122        }
2123        None => {
2124            let mut stmt = conn.prepare(&format!("SELECT {ROOM_MEMBER_SELECT_COLUMNS} FROM room_members WHERE room_id = ?1"))?;
2125            let rows = stmt.query_map(params![room_id], room_member_from_row)?;
2126            rows.collect()
2127        }
2128    }
2129}
2130
2131/// `(room_id, user_id)`'s single `room_members` row, if any — the targeted
2132/// counterpart of [`room_members`] for a "is this one user a member, and
2133/// what membership/power level do they hold" check (every gate in
2134/// `routes::matrix::rooms` needs exactly this, not a full-room scan).
2135pub fn room_member(conn: &Connection, room_id: &str, user_id: i64) -> rusqlite::Result<Option<RoomMember>> {
2136    conn.query_row(
2137        &format!("SELECT {ROOM_MEMBER_SELECT_COLUMNS} FROM room_members WHERE room_id = ?1 AND user_id = ?2"),
2138        params![room_id, user_id],
2139        room_member_from_row,
2140    )
2141    .optional()
2142}
2143
2144/// Delete `user_id`'s `room_members` row in `room_id`, but ONLY if its
2145/// current membership is `leave` — `POST /rooms/{roomId}/forget`'s own gate
2146/// AND its effect are the same one-row conditional delete, so there is no
2147/// separate read-then-write race to worry about. Returns the number of rows
2148/// deleted (`0` if the row was absent or not currently `leave`), so the
2149/// route layer can tell a genuine forget from a no-op gate refusal.
2150pub fn forget_membership(conn: &Connection, room_id: &str, user_id: i64) -> rusqlite::Result<usize> {
2151    conn.execute(
2152        "DELETE FROM room_members WHERE room_id = ?1 AND user_id = ?2 AND membership = 'leave'",
2153        params![room_id, user_id],
2154    )
2155}
2156
2157/// Up to `limit` other members (`join`/`invite`) of `room_id`, oldest
2158/// membership-change first, excluding `exclude_user_id` — `/sync`'s
2159/// `m.heroes` (plan §3.4).
2160pub fn room_heroes(conn: &Connection, room_id: &str, exclude_user_id: i64, limit: i64) -> rusqlite::Result<Vec<i64>> {
2161    let mut stmt = conn.prepare(
2162        "SELECT user_id FROM room_members
2163         WHERE room_id = ?1 AND user_id != ?2 AND membership IN ('join', 'invite')
2164         ORDER BY updated_at ASC LIMIT ?3",
2165    )?;
2166    let rows = stmt.query_map(params![room_id, exclude_user_id, limit], |row| row.get(0))?;
2167    rows.collect()
2168}
2169
2170/// Every room id `user_id` has a `room_members` row in, optionally filtered
2171/// to one `membership` — the "which rooms is this user in" query behind
2172/// `/sync`'s room list and `on_credential_revoked`'s wake fan-out.
2173pub fn rooms_for_user(conn: &Connection, user_id: i64, membership: Option<Membership>) -> rusqlite::Result<Vec<String>> {
2174    match membership {
2175        Some(m) => {
2176            let mut stmt = conn.prepare("SELECT room_id FROM room_members WHERE user_id = ?1 AND membership = ?2")?;
2177            let rows = stmt.query_map(params![user_id, m.as_str()], |row| row.get(0))?;
2178            rows.collect()
2179        }
2180        None => {
2181            let mut stmt = conn.prepare("SELECT room_id FROM room_members WHERE user_id = ?1")?;
2182            let rows = stmt.query_map(params![user_id], |row| row.get(0))?;
2183            rows.collect()
2184        }
2185    }
2186}
2187
2188/// Every room in `room_ids` (typically the caller's currently joined rooms)
2189/// with at least one `events`, `receipts`, or room-scoped `account_data` row
2190/// whose `stream_id` falls in `(since_exclusive, upto_inclusive]` —
2191/// `routes::matrix::sync`'s changed-room prefilter: an incremental sync must
2192/// not run a room block builder's own dozen-odd queries against EVERY
2193/// joined room just to discover that most did not change (with a few
2194/// hundred rooms that is thousands of queries per wake, all held under the
2195/// single `messenger.db` mutex). This set is EXACT, not merely a safe
2196/// over-approximation: every field a room's `/sync` block can ever report
2197/// (`timeline`, `state`, per-room `account_data`, and `m.receipt`) derives
2198/// from exactly these three tables — the one exception, typing, is an
2199/// in-memory `crate::typing::TypingRegistry` fact this function has no visibility
2200/// into and the caller checks separately. Returns an empty set without
2201/// querying when `room_ids` is empty.
2202pub fn rooms_changed_in_window(
2203    conn: &Connection,
2204    room_ids: &[String],
2205    caller_user_id: i64,
2206    since_exclusive: i64,
2207    upto_inclusive: i64,
2208) -> rusqlite::Result<HashSet<String>> {
2209    let mut changed = HashSet::new();
2210    if room_ids.is_empty() {
2211        return Ok(changed);
2212    }
2213    let placeholders = vec!["?"; room_ids.len()].join(",");
2214
2215    // `events` and `receipts` are not scoped to `caller_user_id` at all
2216    // (any member's write counts), so both need the `room_id IN (...)`
2217    // restriction to the caller's own joined-room set.
2218    for table in ["events", "receipts", "pub_events"] {
2219        let sql = format!("SELECT DISTINCT room_id FROM {table} WHERE stream_id > ? AND stream_id <= ? AND room_id IN ({placeholders})");
2220        let mut stmt = conn.prepare(&sql)?;
2221        let mut bound: Vec<&dyn rusqlite::ToSql> = vec![&since_exclusive, &upto_inclusive];
2222        for room_id in room_ids {
2223            bound.push(room_id as &dyn rusqlite::ToSql);
2224        }
2225        let mut rows = stmt.query(bound.as_slice())?;
2226        while let Some(row) = rows.next()? {
2227            changed.insert(row.get::<_, String>(0)?);
2228        }
2229    }
2230
2231    // `account_data` rows are already scoped to ONE user — narrowing by
2232    // `user_id` first (its own indexed column, `idx_account_data_user_stream`)
2233    // is cheaper than another `room_id IN (...)` list, then intersect with
2234    // `room_ids` in Rust (global account data, `room_id = ''`, never
2235    // matches a real room id and is naturally excluded by that intersection).
2236    let room_id_set: HashSet<&str> = room_ids.iter().map(String::as_str).collect();
2237    let mut stmt = conn.prepare("SELECT DISTINCT room_id FROM account_data WHERE user_id = ?1 AND stream_id > ?2 AND stream_id <= ?3")?;
2238    let mut rows = stmt.query(params![caller_user_id, since_exclusive, upto_inclusive])?;
2239    while let Some(row) = rows.next()? {
2240        let room_id: String = row.get(0)?;
2241        if room_id_set.contains(room_id.as_str()) {
2242            changed.insert(room_id);
2243        }
2244    }
2245
2246    Ok(changed)
2247}
2248
2249/// `mxid`'s membership in `room_id` as of `at_stream_id`: the latest
2250/// `m.room.member` event for that state key with `stream_id <= at_stream_id`,
2251/// or `None` when there was none yet (or its content carries no known
2252/// membership). The point-in-time counterpart of [`room_member`] — the
2253/// device-list delta needs "what was this user's membership BEFORE the sync
2254/// window opened" to tell a newly shared room from a profile update.
2255pub fn membership_at(conn: &Connection, room_id: &str, mxid: &str, at_stream_id: i64) -> rusqlite::Result<Option<Membership>> {
2256    let content: Option<String> = conn
2257        .query_row(
2258            "SELECT content FROM events
2259             WHERE room_id = ?1 AND event_type = 'm.room.member' AND state_key = ?2 AND stream_id <= ?3
2260             ORDER BY stream_id DESC LIMIT 1",
2261            params![room_id, mxid, at_stream_id],
2262            |row| row.get(0),
2263        )
2264        .optional()?;
2265    Ok(content
2266        .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
2267        .and_then(|value| value.get("membership").and_then(|m| m.as_str()).and_then(Membership::from_wire_name)))
2268}
2269
2270/// Every room in `room_ids` that has at least one `m.room.member` event with
2271/// `stream_id` in `(from_exclusive, to_inclusive]` — ONE indexed query for
2272/// the whole room set, so the device-list delta only does per-room work for
2273/// rooms where a membership could have changed (an incremental `/sync` runs
2274/// it on every wake; see [`rooms_changed_in_window`] for the same rule).
2275/// Returns an empty vec without querying when `room_ids` is empty.
2276pub fn rooms_with_member_events_in_window(conn: &Connection, room_ids: &[String], from_exclusive: i64, to_inclusive: i64) -> rusqlite::Result<Vec<String>> {
2277    if room_ids.is_empty() {
2278        return Ok(Vec::new());
2279    }
2280    let placeholders = vec!["?"; room_ids.len()].join(",");
2281    let sql = format!(
2282        "SELECT DISTINCT room_id FROM events
2283         WHERE event_type = 'm.room.member' AND stream_id > ? AND stream_id <= ? AND room_id IN ({placeholders})"
2284    );
2285    let mut stmt = conn.prepare(&sql)?;
2286    let mut bound: Vec<&dyn rusqlite::ToSql> = vec![&from_exclusive, &to_inclusive];
2287    for room_id in room_ids {
2288        bound.push(room_id as &dyn rusqlite::ToSql);
2289    }
2290    let rows = stmt.query_map(bound.as_slice(), |row| row.get(0))?;
2291    rows.collect()
2292}
2293
2294/// Every mxid (state key) that has an `m.room.member` event in `room_id`
2295/// with `stream_id` in `(from_exclusive, to_inclusive]` — the candidates
2296/// whose membership MAY have changed inside the window; the caller compares
2297/// [`room_member`] against [`membership_at`] to find the real transitions.
2298pub fn member_state_keys_in_window(conn: &Connection, room_id: &str, from_exclusive: i64, to_inclusive: i64) -> rusqlite::Result<Vec<String>> {
2299    let mut stmt = conn.prepare(
2300        "SELECT DISTINCT state_key FROM events
2301         WHERE room_id = ?1 AND event_type = 'm.room.member' AND state_key IS NOT NULL AND stream_id > ?2 AND stream_id <= ?3",
2302    )?;
2303    let rows = stmt.query_map(params![room_id, from_exclusive, to_inclusive], |row| row.get(0))?;
2304    rows.collect()
2305}
2306
2307/// Every `user_id` whose `m.room.member` state transitioned to `leave` or
2308/// `ban` in one of `room_ids`, with `stream_id` in `(from_exclusive,
2309/// to_inclusive]` — `routes::matrix::keys`'s `GET /keys/changes`
2310/// `device_lists.left` set (plan §3.7): a user who no longer shares any room
2311/// with the caller. This function only finds the CANDIDATE departures inside
2312/// the given room set; the caller (which already knows its own current
2313/// shared-room membership) is responsible for the second half of the plan's
2314/// rule — excluding anyone who still shares SOME OTHER room with it today.
2315/// Bounded to `room_ids` (typically every room the caller is or was in)
2316/// rather than scanning every room this server has ever created. Returns an
2317/// empty vec without querying when `room_ids` is empty.
2318pub fn user_ids_with_leave_transition_in_rooms(
2319    conn: &Connection,
2320    room_ids: &[String],
2321    from_exclusive: i64,
2322    to_inclusive: i64,
2323) -> rusqlite::Result<Vec<i64>> {
2324    if room_ids.is_empty() {
2325        return Ok(Vec::new());
2326    }
2327    let placeholders = vec!["?"; room_ids.len()].join(",");
2328    let sql = format!(
2329        "SELECT DISTINCT e.state_key, e.content FROM events e
2330         WHERE e.event_type = 'm.room.member' AND e.state_key IS NOT NULL
2331           AND e.stream_id > ? AND e.stream_id <= ?
2332           AND e.room_id IN ({placeholders})"
2333    );
2334    let mut stmt = conn.prepare(&sql)?;
2335    let mut bound: Vec<&dyn rusqlite::ToSql> = vec![&from_exclusive, &to_inclusive];
2336    for room_id in room_ids {
2337        bound.push(room_id as &dyn rusqlite::ToSql);
2338    }
2339    let mut rows = stmt.query(bound.as_slice())?;
2340    let mut mxids = Vec::new();
2341    while let Some(row) = rows.next()? {
2342        let mxid: String = row.get(0)?;
2343        let content: String = row.get(1)?;
2344        if let Ok(value) = serde_json::from_str::<serde_json::Value>(&content) {
2345            if matches!(value.get("membership").and_then(|m| m.as_str()), Some("leave") | Some("ban")) {
2346                mxids.push(mxid);
2347            }
2348        }
2349    }
2350    let mut user_ids = Vec::new();
2351    for mxid in mxids {
2352        if let Some(user_id) = user_id_of(conn, &mxid)? {
2353            user_ids.push(user_id);
2354        }
2355    }
2356    Ok(user_ids)
2357}
2358
2359/// Timeline (non-state, non-redaction) events `sender_user_id` has sent into
2360/// any PRIVATE room (`join_rule = 'invite'`) with `origin_server_ts >=
2361/// since_ms` — `routes::matrix::messaging`'s per-sender daily send-rate
2362/// counter (plan §3.8: "mirroring `DM_MESSAGE_DAILY_LIMIT`... own copy of
2363/// the constant" — the constant itself lives in the route module per this
2364/// codebase's convention; this function is only the count query, since raw
2365/// SQL against `events`/`rooms` stays inside this module). State events
2366/// (`state_key IS NOT NULL`, sent via `/state`) and redactions are excluded
2367/// — this counts actual message-shaped sends, not every write a sender
2368/// makes.
2369pub fn private_room_messages_sent_since(conn: &Connection, sender_user_id: i64, since_ms: i64) -> rusqlite::Result<i64> {
2370    conn.query_row(
2371        "SELECT COUNT(*) FROM events e JOIN rooms r ON r.id = e.room_id
2372         WHERE e.sender_user_id = ?1 AND e.state_key IS NULL AND e.event_type != 'm.room.redaction'
2373           AND e.origin_server_ts >= ?2 AND r.join_rule = 'invite'",
2374        params![sender_user_id, since_ms],
2375        |row| row.get(0),
2376    )
2377}
2378
2379/// How many `(user_id, device_id)` sends/redacts were recorded since `since`
2380/// (an RFC-3339 timestamp, compared as a plain string against `created_at`
2381/// values this server ALWAYS writes in that same shape — safe without the
2382/// `datetime()` normalization `lookup_web_credential`'s own doc warns about,
2383/// which is specifically about comparing against SQLite's own `datetime()`
2384/// output, a different textual shape) — `routes::matrix::messaging`'s
2385/// short-window burst-rate counter (plan §3.8: "20 events / 10s per
2386/// device"). Reads `txn_dedup`, the one table that already carries a
2387/// per-DEVICE identity for a write (`events` itself has no `device_id`
2388/// column).
2389pub fn txn_dedup_count_since(conn: &Connection, user_id: i64, device_id: &str, since: &str) -> rusqlite::Result<i64> {
2390    conn.query_row(
2391        "SELECT COUNT(*) FROM txn_dedup WHERE user_id = ?1 AND device_id = ?2 AND created_at >= ?3",
2392        params![user_id, device_id, since],
2393        |row| row.get(0),
2394    )
2395}
2396
2397// ============================================================================
2398// History visibility (P6 binding rule, shared by `/messages`, `/event`,
2399// `/relations`, and exported for P10's `/sync`)
2400// ============================================================================
2401
2402/// How much of a room's timeline a caller may read, per
2403/// [`visible_upper_bound`]. `UpTo` is INCLUSIVE (a caller who left the room
2404/// still sees the `m.room.member` leave event itself, since that is the last
2405/// event they witnessed).
2406#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2407pub enum HistoryWindow {
2408    /// No read access to this room's timeline at all.
2409    Nothing,
2410    /// Every event, unbounded.
2411    All,
2412    /// Every event with `stream_id <= this`.
2413    UpTo(i64),
2414}
2415
2416impl HistoryWindow {
2417    pub fn contains(self, stream_id: i64) -> bool {
2418        match self {
2419            HistoryWindow::Nothing => false,
2420            HistoryWindow::All => true,
2421            HistoryWindow::UpTo(upto) => stream_id <= upto,
2422        }
2423    }
2424}
2425
2426/// The history-visibility rule every timeline-reading route enforces
2427/// (`routes::matrix::messaging`'s `/messages`, `/event`, `/relations`; P10's
2428/// `/sync` reuses this too):
2429///
2430/// - `world_readable`: ANY signed-in caller sees the WHOLE timeline
2431///   ([`HistoryWindow::All`]), regardless of their own membership — this is
2432///   the one case where a caller who was never a member of the room at all
2433///   still reads it.
2434/// - Every other `history_visibility` value (`shared`, and this server's
2435///   simplified handling of `invited`/`joined` — v1 does not implement their
2436///   finer per-event lower bound, see the module doc): a currently `join`ed
2437///   member sees the whole timeline ([`HistoryWindow::All`]); a member whose
2438///   CURRENT membership is `leave`/`ban` sees only up to the stream position
2439///   of that leave/ban event itself ([`HistoryWindow::UpTo`], read off the
2440///   `m.room.member` state event's own `stream_id` — it is, by definition,
2441///   the LATEST such event for that member, since that is what "current
2442///   membership" means); anyone else (never touched this room's membership
2443///   at all, or is merely `invite`d without ever having joined) sees
2444///   [`HistoryWindow::Nothing`].
2445pub fn visible_upper_bound(conn: &Connection, room: &Room, caller_user_id: i64) -> Result<HistoryWindow, MatrixStoreError> {
2446    if room.history_visibility == HistoryVisibility::WorldReadable {
2447        return Ok(HistoryWindow::All);
2448    }
2449    let Some(mxid) = mxid_of(conn, caller_user_id)? else {
2450        return Ok(HistoryWindow::Nothing);
2451    };
2452    let Some(member_event) = current_state_event(conn, &room.id, "m.room.member", &mxid)? else {
2453        return Ok(HistoryWindow::Nothing);
2454    };
2455    let content: serde_json::Value = serde_json::from_str(&member_event.content)?;
2456    match content.get("membership").and_then(|v| v.as_str()) {
2457        Some("join") => Ok(HistoryWindow::All),
2458        Some("leave") | Some("ban") => Ok(HistoryWindow::UpTo(member_event.stream_id)),
2459        _ => Ok(HistoryWindow::Nothing),
2460    }
2461}
2462
2463// ============================================================================
2464// Receipts
2465// ============================================================================
2466
2467#[derive(Debug, Clone, PartialEq)]
2468pub struct ReceiptRow {
2469    pub room_id: String,
2470    pub user_id: i64,
2471    pub receipt_type: ReceiptType,
2472    pub event_id: String,
2473    pub ts: i64,
2474    pub stream_id: i64,
2475}
2476
2477fn receipt_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ReceiptRow> {
2478    let receipt_type_raw: String = row.get(2)?;
2479    Ok(ReceiptRow {
2480        room_id: row.get(0)?,
2481        user_id: row.get(1)?,
2482        receipt_type: decode_enum(2, "receipt_type", &receipt_type_raw, ReceiptType::from_wire_name)?,
2483        event_id: row.get(3)?,
2484        ts: row.get(4)?,
2485        stream_id: row.get(5)?,
2486    })
2487}
2488
2489const RECEIPT_SELECT_COLUMNS: &str = "room_id, user_id, receipt_type, event_id, ts, stream_id";
2490
2491pub fn get_receipt(conn: &Connection, room_id: &str, user_id: i64, receipt_type: ReceiptType) -> rusqlite::Result<Option<ReceiptRow>> {
2492    conn.query_row(
2493        &format!("SELECT {RECEIPT_SELECT_COLUMNS} FROM receipts WHERE room_id = ?1 AND user_id = ?2 AND receipt_type = ?3"),
2494        params![room_id, user_id, receipt_type.as_str()],
2495        receipt_from_row,
2496    )
2497    .optional()
2498}
2499
2500/// Every receipt row in `room_id` whose own write-order `stream_id` (NOT the
2501/// TARGET event's position — see [`upsert_receipt`]'s own doc on that
2502/// distinction) falls in `(since_exclusive, upto_inclusive]` —
2503/// `routes::matrix::sync`'s per-room `m.receipt` ephemeral delta (plan
2504/// §3.6).
2505pub fn receipts_changed_in_room(conn: &Connection, room_id: &str, since_exclusive: i64, upto_inclusive: i64) -> rusqlite::Result<Vec<ReceiptRow>> {
2506    let mut stmt = conn.prepare(&format!(
2507        "SELECT {RECEIPT_SELECT_COLUMNS} FROM receipts WHERE room_id = ?1 AND stream_id > ?2 AND stream_id <= ?3"
2508    ))?;
2509    let rows = stmt.query_map(params![room_id, since_exclusive, upto_inclusive], receipt_from_row)?;
2510    rows.collect()
2511}
2512
2513/// Upsert `user_id`'s `receipt_type` receipt in `room_id`, refusing to move
2514/// it BACKWARDS: compared by the TARGET events' own `stream_id` (timeline
2515/// position), never by the receipt row's own `stream_id` column (which
2516/// orders receipt *writes* for `/sync`'s ephemeral-since-last-batch filter,
2517/// a different axis). Returns the resulting `stream_id` stamped on the
2518/// row — the existing one, unchanged, on a refused backward move.
2519///
2520/// Refuses with [`MatrixStoreError::UnknownEventId`] if `event_id` does not
2521/// exist, or [`MatrixStoreError::WrongRoom`] if it exists but is in a
2522/// different room than `room_id` — a receipt can never point at an event
2523/// outside the room it is filed in.
2524pub fn upsert_receipt(
2525    conn: &mut Connection,
2526    room_id: &str,
2527    user_id: i64,
2528    receipt_type: ReceiptType,
2529    event_id: &str,
2530    ts: i64,
2531) -> Result<i64, MatrixStoreError> {
2532    let tx = conn.transaction()?;
2533    let stream_id = upsert_receipt_in_tx(&tx, room_id, user_id, receipt_type, event_id, ts)?;
2534    tx.commit()?;
2535    Ok(stream_id)
2536}
2537
2538/// The `&Transaction`-scoped core of [`upsert_receipt`] — split out the same
2539/// way [`insert_timeline_event_in_tx`] is split from [`insert_timeline_event`],
2540/// so [`catch_up_dm_conversation`] (P15: bridging a live legacy DM send, or a
2541/// boot pass catching one up) can advance a reader's receipt in the SAME
2542/// transaction as the messages that reader's receipt targets, rather than a
2543/// second one — `upsert_receipt` itself opens its own transaction and so
2544/// cannot be nested inside a caller's own `Transaction` on the same
2545/// connection.
2546fn upsert_receipt_in_tx(
2547    tx: &Transaction,
2548    room_id: &str,
2549    user_id: i64,
2550    receipt_type: ReceiptType,
2551    event_id: &str,
2552    ts: i64,
2553) -> Result<i64, MatrixStoreError> {
2554    let target: Option<(i64, String)> = tx
2555        .query_row("SELECT stream_id, room_id FROM events WHERE event_id = ?1", params![event_id], |row| {
2556            Ok((row.get(0)?, row.get(1)?))
2557        })
2558        .optional()?;
2559    let Some((target_position, target_room)) = target else {
2560        // A read receipt on a public-store post is accepted and not stored:
2561        // the public store keeps no per-user read state.
2562        if let Some(public) = crate::public_channels::get_event(tx, event_id)? {
2563            if public.room_id != room_id {
2564                return Err(MatrixStoreError::WrongRoom(event_id.to_string()));
2565            }
2566            return Ok(public.stream_id);
2567        }
2568        return Err(MatrixStoreError::UnknownEventId(event_id.to_string()));
2569    };
2570    if target_room != room_id {
2571        return Err(MatrixStoreError::WrongRoom(event_id.to_string()));
2572    }
2573
2574    let existing: Option<(String, i64)> = tx
2575        .query_row(
2576            "SELECT event_id, stream_id FROM receipts WHERE room_id = ?1 AND user_id = ?2 AND receipt_type = ?3",
2577            params![room_id, user_id, receipt_type.as_str()],
2578            |row| Ok((row.get(0)?, row.get(1)?)),
2579        )
2580        .optional()?;
2581
2582    if let Some((existing_event_id, existing_stream_id)) = &existing {
2583        let existing_position: i64 =
2584            tx.query_row("SELECT stream_id FROM events WHERE event_id = ?1", params![existing_event_id], |row| row.get(0))?;
2585        if target_position <= existing_position {
2586            return Ok(*existing_stream_id);
2587        }
2588    }
2589
2590    let stream_id = next_stream_id(tx)?;
2591    tx.execute(
2592        "INSERT INTO receipts (room_id, user_id, receipt_type, event_id, ts, stream_id)
2593         VALUES (?1, ?2, ?3, ?4, ?5, ?6)
2594         ON CONFLICT(room_id, user_id, receipt_type) DO UPDATE SET
2595            event_id = excluded.event_id, ts = excluded.ts, stream_id = excluded.stream_id",
2596        params![room_id, user_id, receipt_type.as_str(), event_id, ts, stream_id],
2597    )?;
2598    Ok(stream_id)
2599}
2600
2601// ============================================================================
2602// Unread notifications (plan §3.5) — `/sync`'s `unread_notifications` block
2603// ============================================================================
2604
2605/// Message-like timeline event types that count toward
2606/// `unread_notifications.notification_count` (plan §3.5) — everything else
2607/// (state events, `m.reaction`, `m.room.redaction`) is excluded. A legacy DM
2608/// migrated from `dm_conversations` (plan §7) lands as
2609/// `org.example.legacy_dm`, which counts the same as a live message.
2610/// Includes historical `org.example.legacy_dm` so unread counts stay correct
2611/// until those timeline rows age out; new writes of that type are refused.
2612const NOTIFICATION_MESSAGE_TYPES_SQL: &str = "('m.room.message', 'm.room.encrypted', 'org.example.legacy_dm')";
2613
2614/// `unread_notifications.notification_count` for `user_id` in `room` (plan
2615/// §3.5): the count of message-like timeline events (see
2616/// [`NOTIFICATION_MESSAGE_TYPES_SQL`]) strictly after the FURTHER-AHEAD of
2617/// the caller's own `m.read`/`m.read.private` receipt — resolved by the
2618/// target events' own `stream_id` (timeline position), never by a receipt
2619/// row's own `stream_id` column (a different axis — see [`upsert_receipt`]'s
2620/// own doc) — and bounded above by [`visible_upper_bound`] (a `leave`/
2621/// `ban`'d member's count never includes anything past their own
2622/// departure). Excludes the caller's own sends (nobody is notified about
2623/// their own message) and, being restricted to
2624/// [`NOTIFICATION_MESSAGE_TYPES_SQL`], every state event and `m.reaction`.
2625/// Not a member of the room at all ([`HistoryWindow::Nothing`]) is always
2626/// `0` — there is nothing this caller could be notified about.
2627///
2628/// `highlight_count` is always `0` from this server (plan §3.5: no push-
2629/// rule/keyword-highlight engine in v1, encrypted or not) — that is a
2630/// constant the `/sync` response builder (P10) sets directly, not a second
2631/// query here.
2632pub fn notification_count(conn: &Connection, room: &Room, user_id: i64) -> Result<i64, MatrixStoreError> {
2633    if room.kind == RoomKind::Channel && !room.is_encrypted {
2634        return Ok(0); // public store: no server-side unread state
2635    }
2636    let upper_bound = match visible_upper_bound(conn, room, user_id)? {
2637        HistoryWindow::Nothing => return Ok(0),
2638        HistoryWindow::All => i64::MAX,
2639        HistoryWindow::UpTo(upper) => upper,
2640    };
2641
2642    let mut after_stream_id: i64 = 0;
2643    for receipt_type in [ReceiptType::Read, ReceiptType::ReadPrivate] {
2644        if let Some(receipt) = get_receipt(conn, &room.id, user_id, receipt_type)? {
2645            if let Some(target) = get_event(conn, &receipt.event_id)? {
2646                after_stream_id = after_stream_id.max(target.stream_id);
2647            }
2648        }
2649    }
2650
2651    let count = conn.query_row(
2652        &format!(
2653            "SELECT COUNT(*) FROM events
2654             WHERE room_id = ?1 AND state_key IS NULL AND sender_user_id != ?2
2655               AND stream_id > ?3 AND stream_id <= ?4
2656               AND event_type IN {NOTIFICATION_MESSAGE_TYPES_SQL}"
2657        ),
2658        params![room.id, user_id, after_stream_id, upper_bound],
2659        |row| row.get(0),
2660    )?;
2661    Ok(count)
2662}
2663
2664// ============================================================================
2665// Account data
2666// ============================================================================
2667
2668/// `room_id` value meaning "global account data" — see [`create_matrix_schema`]'s
2669/// doc comment on why this is `""`, never `NULL`.
2670pub const GLOBAL_ACCOUNT_DATA_ROOM: &str = "";
2671
2672#[derive(Debug, Clone, PartialEq)]
2673pub struct AccountDataRow {
2674    pub user_id: i64,
2675    pub room_id: String,
2676    pub data_type: String,
2677    pub content: String,
2678    pub stream_id: i64,
2679}
2680
2681fn account_data_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<AccountDataRow> {
2682    Ok(AccountDataRow {
2683        user_id: row.get(0)?,
2684        room_id: row.get(1)?,
2685        data_type: row.get(2)?,
2686        content: row.get(3)?,
2687        stream_id: row.get(4)?,
2688    })
2689}
2690
2691const ACCOUNT_DATA_SELECT_COLUMNS: &str = "user_id, room_id, data_type, content, stream_id";
2692
2693/// Upsert one account-data entry. Pass [`GLOBAL_ACCOUNT_DATA_ROOM`] for
2694/// global data (`m.direct`, `m.push_rules`, ...); a real room id for
2695/// per-room data.
2696pub fn upsert_account_data(conn: &mut Connection, user_id: i64, room_id: &str, data_type: &str, content: &str) -> Result<i64, MatrixStoreError> {
2697    let tx = conn.transaction()?;
2698    let stream_id = next_stream_id(&tx)?;
2699    tx.execute(
2700        "INSERT INTO account_data (user_id, room_id, data_type, content, stream_id) VALUES (?1, ?2, ?3, ?4, ?5)
2701         ON CONFLICT(user_id, room_id, data_type) DO UPDATE SET content = excluded.content, stream_id = excluded.stream_id",
2702        params![user_id, room_id, data_type, content, stream_id],
2703    )?;
2704    tx.commit()?;
2705    Ok(stream_id)
2706}
2707
2708pub fn get_account_data(conn: &Connection, user_id: i64, room_id: &str, data_type: &str) -> rusqlite::Result<Option<AccountDataRow>> {
2709    conn.query_row(
2710        &format!("SELECT {ACCOUNT_DATA_SELECT_COLUMNS} FROM account_data WHERE user_id = ?1 AND room_id = ?2 AND data_type = ?3"),
2711        params![user_id, room_id, data_type],
2712        account_data_from_row,
2713    )
2714    .optional()
2715}
2716
2717/// Every `(user_id, room_id)` account-data row changed strictly after
2718/// `since_stream` — the `/sync` delta shape for both global and per-room
2719/// account data (caller passes [`GLOBAL_ACCOUNT_DATA_ROOM`] or a real room
2720/// id).
2721pub fn account_data_since(conn: &Connection, user_id: i64, room_id: &str, since_stream: i64) -> rusqlite::Result<Vec<AccountDataRow>> {
2722    let mut stmt = conn.prepare(&format!(
2723        "SELECT {ACCOUNT_DATA_SELECT_COLUMNS} FROM account_data WHERE user_id = ?1 AND room_id = ?2 AND stream_id > ?3 ORDER BY stream_id ASC"
2724    ))?;
2725    let rows = stmt.query_map(params![user_id, room_id, since_stream], account_data_from_row)?;
2726    rows.collect()
2727}
2728
2729// ============================================================================
2730// Txn-id idempotency
2731// ============================================================================
2732
2733/// What [`txn_dedup_lookup`] found for a given `(user_id, device_id, txn_id)`.
2734#[derive(Debug, Clone, PartialEq, Eq)]
2735pub enum TxnDedupEntry {
2736    /// No prior submission recorded — the caller should proceed to create
2737    /// the event/to-device-send and then [`txn_dedup_record`] it.
2738    NotSeen,
2739    /// Already recorded. `Some(event_id)` for a send/redact that produced
2740    /// an event; `None` for a to-device send (which has no `events` row).
2741    Seen(Option<String>),
2742}
2743
2744/// Check whether `(user_id, device_id, txn_id)` was already handled,
2745/// without recording anything.
2746pub fn txn_dedup_lookup(conn: &Connection, user_id: i64, device_id: &str, txn_id: &str) -> rusqlite::Result<TxnDedupEntry> {
2747    let found: Option<Option<String>> = conn
2748        .query_row(
2749            "SELECT event_id FROM txn_dedup WHERE user_id = ?1 AND device_id = ?2 AND txn_id = ?3",
2750            params![user_id, device_id, txn_id],
2751            |row| row.get(0),
2752        )
2753        .optional()?;
2754    Ok(match found {
2755        None => TxnDedupEntry::NotSeen,
2756        Some(event_id) => TxnDedupEntry::Seen(event_id),
2757    })
2758}
2759
2760/// Record that `(user_id, device_id, txn_id)` has now been handled,
2761/// producing `event_id` (or `None` for a to-device send). Call only after
2762/// [`txn_dedup_lookup`] returned [`TxnDedupEntry::NotSeen`] — this does not
2763/// itself check for a race, per this module's single-writer discipline.
2764pub fn txn_dedup_record(conn: &Connection, user_id: i64, device_id: &str, txn_id: &str, event_id: Option<&str>, now: &str) -> rusqlite::Result<()> {
2765    conn.execute(
2766        "INSERT INTO txn_dedup (user_id, device_id, txn_id, event_id, created_at) VALUES (?1, ?2, ?3, ?4, ?5)",
2767        params![user_id, device_id, txn_id, event_id, now],
2768    )?;
2769    Ok(())
2770}
2771
2772/// The `txn_id` `(user_id, device_id)` used to produce `event_id`, if any —
2773/// [`client_event_json`](crate::routes::matrix::client_event_json)'s
2774/// `unsigned.transaction_id`, which Matrix reveals ONLY to the same device
2775/// that sent the event, never to any other viewer (including the sender's
2776/// OTHER devices). A miss (`None`) is the overwhelmingly common case (every
2777/// event not authored by this exact device on this exact send) and costs a
2778/// single indexed `txn_dedup` primary-key lookup.
2779pub fn txn_id_for_event(conn: &Connection, user_id: i64, device_id: &str, event_id: &str) -> rusqlite::Result<Option<String>> {
2780    conn.query_row(
2781        "SELECT txn_id FROM txn_dedup WHERE user_id = ?1 AND device_id = ?2 AND event_id = ?3",
2782        params![user_id, device_id, event_id],
2783        |row| row.get(0),
2784    )
2785    .optional()
2786}
2787
2788// ============================================================================
2789// Filters
2790// ============================================================================
2791
2792/// Store an opaque Filter JSON object, returning its new `filter_id`.
2793pub fn create_filter(conn: &Connection, user_id: i64, definition: &str) -> rusqlite::Result<i64> {
2794    conn.execute("INSERT INTO filters (user_id, definition) VALUES (?1, ?2)", params![user_id, definition])?;
2795    Ok(conn.last_insert_rowid())
2796}
2797
2798/// Fetch `user_id`'s own `filter_id`'s definition — scoped to `user_id` so
2799/// one account can never read another's stored filter by guessing an id.
2800pub fn get_filter(conn: &Connection, user_id: i64, filter_id: i64) -> rusqlite::Result<Option<String>> {
2801    conn.query_row(
2802        "SELECT definition FROM filters WHERE id = ?1 AND user_id = ?2",
2803        params![filter_id, user_id],
2804        |row| row.get(0),
2805    )
2806    .optional()
2807}
2808
2809// ============================================================================
2810// Legacy DM migration map
2811// ============================================================================
2812
2813pub fn insert_legacy_dm_message_map(conn: &Connection, legacy_message_id: i64, event_id: &str) -> rusqlite::Result<()> {
2814    conn.execute(
2815        "INSERT INTO legacy_dm_message_map (legacy_message_id, event_id) VALUES (?1, ?2)",
2816        params![legacy_message_id, event_id],
2817    )?;
2818    Ok(())
2819}
2820
2821pub fn legacy_dm_message_event_id(conn: &Connection, legacy_message_id: i64) -> rusqlite::Result<Option<String>> {
2822    conn.query_row(
2823        "SELECT event_id FROM legacy_dm_message_map WHERE legacy_message_id = ?1",
2824        params![legacy_message_id],
2825        |row| row.get(0),
2826    )
2827    .optional()
2828}
2829
2830/// The highest `legacy_message_id` already imported into `room_id`'s
2831/// `org.example.legacy_dm` timeline (P15: `matrix_migration`'s own
2832/// catch-up cursor). `legacy_dm_message_map` carries no `room_id` of its
2833/// own — `dm_messages.id` is a single autoincrement column shared by every
2834/// legacy conversation, not scoped per conversation — so this joins through
2835/// `events` to find only the rows imported into THIS room. Message ids are
2836/// ascending within one conversation and every migration/catch-up pass
2837/// imports every not-yet-mapped row up to "now", so `id > this` is exactly
2838/// that conversation's not-yet-imported set — the caller's own cheap
2839/// alternative to re-checking `legacy_dm_message_event_id` once per message.
2840/// `None` for a room with no legacy message imported yet (a conversation
2841/// that had zero messages at migration time).
2842pub fn highest_mapped_legacy_message_id(conn: &Connection, room_id: &str) -> rusqlite::Result<Option<i64>> {
2843    conn.query_row(
2844        "SELECT MAX(m.legacy_message_id) FROM legacy_dm_message_map m
2845         JOIN events e ON e.event_id = m.event_id
2846         WHERE e.room_id = ?1",
2847        params![room_id],
2848        |row| row.get(0),
2849    )
2850}
2851
2852// ============================================================================
2853// Legacy DM migration batch (P12) — the binary crate's own `matrix_migration`
2854// module is the only caller. Everything one `dm_conversations` row needs
2855// (room + bootstrap state events + every not-yet-imported message + read
2856// receipts + the `m.direct` account-data hint for both parties) lands in
2857// ONE transaction here, for the same reason `create_room_with_state` is a
2858// batch entry point rather than several independent calls (P5 correction):
2859// a crash partway through must never leave a room with some but not all of
2860// its own bootstrap state, or a room with some but not all of its messages.
2861// This reuses the SAME private `_in_tx` helpers `create_room_with_state`
2862// itself uses (`insert_room_row`, `apply_state_event_in_tx`,
2863// `insert_timeline_event_in_tx`), so a migrated room's rows are written by
2864// the exact same code paths a live `createRoom`/`send` call would use.
2865// ============================================================================
2866
2867/// One legacy `dm_messages` row queued for import (plan §7 step 5).
2868/// `event_id` is minted by the caller ([`crate::matrix_migration`] in the
2869/// binary crate) so it can be recorded in the caller's own bookkeeping (the
2870/// idempotency proof compares event ids across two runs) without a second
2871/// round trip back into this module. `content` is already the full
2872/// `org.example.legacy_dm` JSON body
2873/// (`{"legacy_message_id","legacy_conversation_id","nonce_b64","ciphertext_b64"}`).
2874///
2875/// # `legacy_conversation_id`
2876///
2877/// Both this message's own content AND the room's
2878/// `org.example.legacy_dm_key` state event content (plan §2 manager
2879/// decision 5) carry the legacy `dm_conversations.id` (as
2880/// `legacy_conversation_id`) — a client derives
2881/// `mlc_vault::dm::conversation_key` from that id plus the sender's and its
2882/// own keys, and has no other way to learn which legacy conversation a
2883/// migrated room came from. A `messenger.db` migrated by a build before this
2884/// field existed is missing it on every already-migrated row; that is a
2885/// dev-only concern (nothing has shipped) — recreate the database rather
2886/// than backfill it.
2887#[derive(Debug, Clone, PartialEq)]
2888pub struct LegacyDmMessageImport {
2889    pub legacy_message_id: i64,
2890    pub event_id: String,
2891    pub sender_user_id: i64,
2892    pub content: String,
2893    pub origin_server_ts: i64,
2894}
2895
2896/// One reader's `m.read` receipt to set once its target message has been
2897/// imported — `up_to_legacy_message_id` is the LAST `dm_messages.id` this
2898/// reader actually had `read_at` stamped for in `social.db` (never a later
2899/// one — plan P12 scope: "never over-count as read"). `ts_ms` is that same
2900/// message's own `read_at`, converted to epoch milliseconds.
2901#[derive(Debug, Clone, PartialEq)]
2902pub struct LegacyDmReadReceipt {
2903    pub reader_user_id: i64,
2904    pub up_to_legacy_message_id: i64,
2905    pub ts_ms: i64,
2906}
2907
2908/// One `m.direct` merge hint (plan P12 scope: "merge into their existing
2909/// `m.direct`, do not overwrite other entries") — `user_id`'s global
2910/// `m.direct` account data gains `peer_mxid` mapped to the migrated room id,
2911/// alongside whatever entries it already has.
2912#[derive(Debug, Clone, PartialEq)]
2913pub struct LegacyDmDirectHint {
2914    pub user_id: i64,
2915    pub peer_mxid: String,
2916}
2917
2918/// Everything [`migrate_dm_conversation`] needs beyond the room bootstrap
2919/// itself, grouped into one value for the same "stays lint-clean without
2920/// `#[allow(clippy::too_many_arguments)]`" reason [`RoomBootstrap`]'s own doc
2921/// comment states.
2922#[derive(Debug, Clone, Copy)]
2923pub struct DmMigrationExtras<'a> {
2924    pub messages: &'a [LegacyDmMessageImport],
2925    pub receipts: &'a [LegacyDmReadReceipt],
2926    pub direct_hints: &'a [LegacyDmDirectHint],
2927}
2928
2929/// What one [`migrate_dm_conversation`] call actually wrote.
2930#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
2931pub struct DmMigrationCounts {
2932    pub messages_imported: usize,
2933    pub receipts_set: usize,
2934}
2935
2936/// Merge `{peer_mxid: [room_id]}` into `user_id`'s existing GLOBAL `m.direct`
2937/// account data, inside `tx` — every other peer's entry already in the
2938/// object is left untouched, and a `room_id` already present under
2939/// `peer_mxid` is not duplicated.
2940fn merge_m_direct_in_tx(tx: &Transaction, user_id: i64, peer_mxid: &str, room_id: &str) -> Result<(), MatrixStoreError> {
2941    let existing: Option<String> = tx
2942        .query_row(
2943            "SELECT content FROM account_data WHERE user_id = ?1 AND room_id = ?2 AND data_type = 'm.direct'",
2944            params![user_id, GLOBAL_ACCOUNT_DATA_ROOM],
2945            |row| row.get(0),
2946        )
2947        .optional()?;
2948
2949    let mut direct: serde_json::Map<String, serde_json::Value> = match existing {
2950        Some(content) => serde_json::from_str(&content)?,
2951        None => serde_json::Map::new(),
2952    };
2953    let rooms_entry = direct.entry(peer_mxid.to_string()).or_insert_with(|| serde_json::Value::Array(Vec::new()));
2954    if !rooms_entry.is_array() {
2955        *rooms_entry = serde_json::Value::Array(Vec::new());
2956    }
2957    if let serde_json::Value::Array(list) = rooms_entry {
2958        if !list.iter().any(|v| v.as_str() == Some(room_id)) {
2959            list.push(serde_json::Value::String(room_id.to_string()));
2960        }
2961    }
2962
2963    let stream_id = next_stream_id(tx)?;
2964    tx.execute(
2965        "INSERT INTO account_data (user_id, room_id, data_type, content, stream_id) VALUES (?1, ?2, 'm.direct', ?3, ?4)
2966         ON CONFLICT(user_id, room_id, data_type) DO UPDATE SET content = excluded.content, stream_id = excluded.stream_id",
2967        params![user_id, GLOBAL_ACCOUNT_DATA_ROOM, serde_json::Value::Object(direct).to_string(), stream_id],
2968    )?;
2969    Ok(())
2970}
2971
2972/// The P12 migration's own batch entry point: `bootstrap`'s room row, then
2973/// every entry of `state_events` (identical shape to
2974/// [`create_room_with_state`]'s own loop — same order requirement: a
2975/// member's `m.room.member` must precede `m.room.power_levels` for
2976/// [`refresh_power_levels`] to see them, see that function's own doc), then
2977/// every `extras.messages` entry as an `org.example.legacy_dm`
2978/// timeline event plus its `legacy_dm_message_map` row, then
2979/// `extras.receipts`, then `extras.direct_hints` — all in ONE transaction.
2980/// `extras.messages` must already be filtered by the caller to exclude any
2981/// `legacy_message_id` already present in `legacy_dm_message_map` and
2982/// ordered ascending by `legacy_message_id`; this function does not
2983/// re-check either, since the binary crate's own `matrix_migration` module
2984/// already established the room does not exist yet via
2985/// [`room_by_legacy_dm_id`] before ever calling this.
2986pub fn migrate_dm_conversation(
2987    conn: &mut Connection,
2988    bootstrap: RoomBootstrap<'_>,
2989    state_events: &[NewStateEvent],
2990    bootstrap_origin_server_ts: i64,
2991    extras: DmMigrationExtras<'_>,
2992) -> Result<(Room, DmMigrationCounts), MatrixStoreError> {
2993    let tx = conn.transaction()?;
2994
2995    insert_room_row(
2996        &tx,
2997        bootstrap.room_id,
2998        bootstrap.kind,
2999        bootstrap.creator_user_id,
3000        bootstrap.created_at,
3001        bootstrap.is_encrypted,
3002        bootstrap.join_rule,
3003        bootstrap.history_visibility,
3004        bootstrap.dm_pair_key,
3005        bootstrap.legacy_dm_id,
3006    )?;
3007    for event in state_events {
3008        apply_state_event_in_tx(
3009            &tx,
3010            &event.event_id,
3011            bootstrap.room_id,
3012            event.sender_user_id,
3013            &event.event_type,
3014            &event.state_key,
3015            &event.content,
3016            bootstrap_origin_server_ts,
3017            bootstrap.created_at,
3018        )?;
3019    }
3020
3021    let mut event_id_by_legacy_id: std::collections::HashMap<i64, String> = std::collections::HashMap::new();
3022    for message in extras.messages {
3023        insert_timeline_event_in_tx(
3024            &tx,
3025            &TimelineEventRow {
3026                event_id: &message.event_id,
3027                room_id: bootstrap.room_id,
3028                sender_user_id: message.sender_user_id,
3029                event_type: "org.example.legacy_dm",
3030                content: &message.content,
3031                origin_server_ts: message.origin_server_ts,
3032                txn_id: None,
3033            },
3034        )?;
3035        insert_legacy_dm_message_map(&tx, message.legacy_message_id, &message.event_id)?;
3036        event_id_by_legacy_id.insert(message.legacy_message_id, message.event_id.clone());
3037    }
3038
3039    let mut receipts_set = 0usize;
3040    for receipt in extras.receipts {
3041        let Some(target_event_id) = event_id_by_legacy_id.get(&receipt.up_to_legacy_message_id) else { continue };
3042        let target_stream_id: i64 =
3043            tx.query_row("SELECT stream_id FROM events WHERE event_id = ?1", params![target_event_id], |row| row.get(0))?;
3044        tx.execute(
3045            "INSERT INTO receipts (room_id, user_id, receipt_type, event_id, ts, stream_id) VALUES (?1, ?2, 'm.read', ?3, ?4, ?5)",
3046            params![bootstrap.room_id, receipt.reader_user_id, target_event_id, receipt.ts_ms, target_stream_id],
3047        )?;
3048        receipts_set += 1;
3049    }
3050
3051    for hint in extras.direct_hints {
3052        merge_m_direct_in_tx(&tx, hint.user_id, &hint.peer_mxid, bootstrap.room_id)?;
3053    }
3054
3055    tx.commit()?;
3056
3057    Ok((
3058        Room {
3059            id: bootstrap.room_id.to_string(),
3060            kind: bootstrap.kind,
3061            room_version: MATRIX_ROOM_VERSION.to_string(),
3062            creator_user_id: bootstrap.creator_user_id,
3063            created_at: bootstrap.created_at.to_string(),
3064            is_encrypted: bootstrap.is_encrypted,
3065            join_rule: bootstrap.join_rule,
3066            history_visibility: bootstrap.history_visibility,
3067            dm_pair_key: bootstrap.dm_pair_key.map(str::to_string),
3068            legacy_dm_id: bootstrap.legacy_dm_id,
3069        },
3070        DmMigrationCounts { messages_imported: extras.messages.len(), receipts_set },
3071    ))
3072}
3073
3074/// A catch-up pass for an ALREADY migrated room (P15) — the binary crate's
3075/// own `matrix_migration` module calls this both from a live legacy DM
3076/// send's bridge (`routes::dm::create_message`, right after its own insert
3077/// commits) and from a boot pass over a conversation
3078/// [`room_by_legacy_dm_id`] already finds a room for. `messages` (already
3079/// filtered by the caller to exclude anything already in
3080/// `legacy_dm_message_map`, ordered ascending by `legacy_message_id`) lands
3081/// as new `org.example.legacy_dm` timeline events plus their
3082/// `legacy_dm_message_map` rows; then each of `receipts` advances that
3083/// reader's `m.read` receipt — never moved backwards
3084/// ([`upsert_receipt_in_tx`]'s own guarantee, the SAME one a live `/receipt`
3085/// call gets) — all in ONE transaction. Reuses the exact per-event helpers
3086/// [`migrate_dm_conversation`] itself uses, so a caught-up room's rows are
3087/// indistinguishable from ones a fresh migration (or a live Matrix send)
3088/// would have produced.
3089///
3090/// Each `receipt.up_to_legacy_message_id` is resolved via
3091/// [`legacy_dm_message_event_id`] against the FULL map (this transaction's
3092/// own just-inserted rows are visible to it too, same connection) — not just
3093/// this call's own `messages` — because the caller's own read-position query
3094/// considers EVERY read message in the conversation, including ones a
3095/// PRIOR pass already imported (a message read only after it was already
3096/// bridged must still move the receipt here; manager review, P15). A target
3097/// this catch-up's own `messages` did not just insert AND that is not yet in
3098/// the map at all (should not happen — a receipt only ever targets a message
3099/// that exists) is skipped rather than erroring, same defensive posture
3100/// [`migrate_dm_conversation`]'s own receipt loop takes.
3101pub fn catch_up_dm_conversation(
3102    conn: &mut Connection,
3103    room_id: &str,
3104    messages: &[LegacyDmMessageImport],
3105    receipts: &[LegacyDmReadReceipt],
3106) -> Result<DmMigrationCounts, MatrixStoreError> {
3107    let tx = conn.transaction()?;
3108    let counts = catch_up_dm_in_tx(&tx, room_id, messages, receipts)?;
3109    tx.commit()?;
3110    Ok(counts)
3111}
3112
3113/// The body of [`catch_up_dm_conversation`], inside the caller's transaction —
3114/// shared with [`adopt_dm_room_for_legacy`], so an adopted native room takes
3115/// its legacy messages and read positions through exactly the same code as
3116/// an already-migrated room's catch-up.
3117fn catch_up_dm_in_tx(
3118    tx: &Transaction,
3119    room_id: &str,
3120    messages: &[LegacyDmMessageImport],
3121    receipts: &[LegacyDmReadReceipt],
3122) -> Result<DmMigrationCounts, MatrixStoreError> {
3123    for message in messages {
3124        insert_timeline_event_in_tx(
3125            tx,
3126            &TimelineEventRow {
3127                event_id: &message.event_id,
3128                room_id,
3129                sender_user_id: message.sender_user_id,
3130                event_type: "org.example.legacy_dm",
3131                content: &message.content,
3132                origin_server_ts: message.origin_server_ts,
3133                txn_id: None,
3134            },
3135        )?;
3136        insert_legacy_dm_message_map(tx, message.legacy_message_id, &message.event_id)?;
3137    }
3138
3139    let mut receipts_set = 0usize;
3140    for receipt in receipts {
3141        let Some(target_event_id) = legacy_dm_message_event_id(tx, receipt.up_to_legacy_message_id)? else { continue };
3142        upsert_receipt_in_tx(tx, room_id, receipt.reader_user_id, ReceiptType::Read, &target_event_id, receipt.ts_ms)?;
3143        receipts_set += 1;
3144    }
3145
3146    Ok(DmMigrationCounts { messages_imported: messages.len(), receipts_set })
3147}
3148
3149/// What [`adopt_dm_room_for_legacy`] needs to bind a native DM room to a
3150/// legacy conversation, grouped for the same lint-clean reason
3151/// [`RoomBootstrap`]'s doc states.
3152#[derive(Debug, Clone, Copy)]
3153pub struct DmAdoption<'a> {
3154    /// The live native DM room that already holds the pair's `dm_pair_key`.
3155    pub room_id: &'a str,
3156    /// The legacy `dm_conversations.id` the room is bound to.
3157    pub legacy_dm_id: i64,
3158    /// The pair's `org.example.legacy_dm_key` state events. Each is
3159    /// written only when the room has no current event for its
3160    /// `(event_type, state_key)` yet — an existing key event is never
3161    /// overwritten.
3162    pub key_events: &'a [NewStateEvent],
3163    pub key_events_origin_server_ts: i64,
3164    /// RFC3339 stamp for the member/state bookkeeping of the key events.
3165    pub now: &'a str,
3166}
3167
3168/// Bind an EXISTING native DM room to legacy conversation
3169/// `adoption.legacy_dm_id`, instead of minting a second room for the same
3170/// pair (`rooms.dm_pair_key` is `UNIQUE`, so a second insert is refused):
3171/// sets `rooms.legacy_dm_id`, writes any missing legacy key state event,
3172/// then imports `messages` and `receipts` through the shared catch-up
3173/// body ([`catch_up_dm_in_tx`]) — all in ONE transaction, so a failure
3174/// leaves the room exactly as it was. `messages` follow
3175/// [`catch_up_dm_conversation`]'s contract (not yet in
3176/// `legacy_dm_message_map`, ascending by `legacy_message_id`).
3177///
3178/// Returns `Ok(None)`, having written nothing, when `room_id` is not a DM
3179/// room without a legacy binding (missing, another kind, or already bound to
3180/// a legacy conversation) — the caller decides what that means.
3181pub fn adopt_dm_room_for_legacy(
3182    conn: &mut Connection,
3183    adoption: DmAdoption<'_>,
3184    messages: &[LegacyDmMessageImport],
3185    receipts: &[LegacyDmReadReceipt],
3186) -> Result<Option<DmMigrationCounts>, MatrixStoreError> {
3187    let tx = conn.transaction()?;
3188
3189    let bound = tx.execute(
3190        "UPDATE rooms SET legacy_dm_id = ?1 WHERE id = ?2 AND kind = 'dm' AND legacy_dm_id IS NULL",
3191        params![adoption.legacy_dm_id, adoption.room_id],
3192    )?;
3193    if bound == 0 {
3194        return Ok(None);
3195    }
3196
3197    for event in adoption.key_events {
3198        if current_state_event(&tx, adoption.room_id, &event.event_type, &event.state_key)?.is_some() {
3199            continue;
3200        }
3201        apply_state_event_in_tx(
3202            &tx,
3203            &event.event_id,
3204            adoption.room_id,
3205            event.sender_user_id,
3206            &event.event_type,
3207            &event.state_key,
3208            &event.content,
3209            adoption.key_events_origin_server_ts,
3210            adoption.now,
3211        )?;
3212    }
3213
3214    let counts = catch_up_dm_in_tx(&tx, adoption.room_id, messages, receipts)?;
3215    tx.commit()?;
3216    Ok(Some(counts))
3217}
3218
3219// ============================================================================
3220// Public room directory (P8: `routes::matrix::account`'s `GET`/`POST
3221// /publicRooms`)
3222// ============================================================================
3223
3224/// One row of `GET /publicRooms`'s `chunk` — a public-`join_rule` room's
3225/// directory summary. `name`/`topic` are `None` when the room never had an
3226/// `m.room.name`/`m.room.topic` state event (plan §5: only a `channel`-kind
3227/// room is ever `join_rule = 'public'`, but this reads the column directly
3228/// rather than also filtering on `kind`, so it stays correct even if a
3229/// future piece ever makes a `group` room public).
3230#[derive(Debug, Clone, PartialEq)]
3231pub struct PublicRoomSummary {
3232    pub room_id: String,
3233    pub name: Option<String>,
3234    pub topic: Option<String>,
3235    pub num_joined_members: i64,
3236    pub world_readable: bool,
3237}
3238
3239/// One page of the public-room directory, ordered by `rooms.id` — a stable
3240/// pagination key (a room id never changes once minted). `after_room_id` is
3241/// the previous page's own last room id (`None` for the first page);
3242/// `search_term`, when non-empty, keeps only rooms whose current
3243/// `m.room.name` contains it (case-insensitive substring — Matrix's own
3244/// `generic_search_term`).
3245///
3246/// This scans every public room in Rust rather than pushing the search/
3247/// pagination interaction into SQL: a non-federated single server's public-
3248/// channel count is small and moderator-managed (channels are not something
3249/// a bot mass-creates), so an O(public rooms) scan per call is cheaper to
3250/// keep correct than a SQL query that would otherwise have to paginate
3251/// `LIMIT`-first and then discover a `search_term` filtered a whole page
3252/// down to nothing.
3253///
3254/// Returns `(page, has_more, total_room_count_estimate)` — `has_more` is
3255/// whether another page exists beyond this one (the route layer's own
3256/// `next_batch` decision); `total_room_count_estimate` is the count of every
3257/// public room regardless of `search_term`/pagination, matching the spec's
3258/// "estimate of the total number of public rooms" wording.
3259pub fn public_rooms_page(
3260    conn: &Connection,
3261    after_room_id: Option<&str>,
3262    limit: usize,
3263    search_term: Option<&str>,
3264) -> Result<(Vec<PublicRoomSummary>, bool, i64), MatrixStoreError> {
3265    let mut stmt = conn.prepare(
3266        "SELECT r.id, r.history_visibility,
3267                (SELECT COUNT(*) FROM room_members WHERE room_id = r.id AND membership = 'join')
3268         FROM rooms r WHERE r.join_rule = 'public' ORDER BY r.id ASC",
3269    )?;
3270    let mut rows = stmt.query([])?;
3271    let mut all = Vec::new();
3272    while let Some(row) = rows.next()? {
3273        let room_id: String = row.get(0)?;
3274        let history_visibility_raw: String = row.get(1)?;
3275        let num_joined_members: i64 = row.get(2)?;
3276        let world_readable = history_visibility_raw == HistoryVisibility::WorldReadable.as_str();
3277
3278        let name = current_state_event(conn, &room_id, "m.room.name", "")?
3279            .and_then(|e| serde_json::from_str::<serde_json::Value>(&e.content).ok())
3280            .and_then(|v| v.get("name").and_then(|n| n.as_str()).map(str::to_string));
3281        let topic = current_state_event(conn, &room_id, "m.room.topic", "")?
3282            .and_then(|e| serde_json::from_str::<serde_json::Value>(&e.content).ok())
3283            .and_then(|v| v.get("topic").and_then(|t| t.as_str()).map(str::to_string));
3284
3285        all.push(PublicRoomSummary { room_id, name, topic, num_joined_members, world_readable });
3286    }
3287    drop(rows);
3288    drop(stmt);
3289
3290    let total_room_count_estimate = all.len() as i64;
3291
3292    let filtered: Vec<PublicRoomSummary> = match search_term {
3293        Some(term) if !term.is_empty() => {
3294            let term_lower = term.to_lowercase();
3295            all.into_iter().filter(|room| room.name.as_deref().is_some_and(|n| n.to_lowercase().contains(&term_lower))).collect()
3296        }
3297        _ => all,
3298    };
3299
3300    let start = match after_room_id {
3301        Some(cursor) => filtered.iter().position(|r| r.room_id == cursor).map_or(0, |idx| idx + 1),
3302        None => 0,
3303    };
3304    let has_more = filtered.len() > start + limit;
3305    let page: Vec<PublicRoomSummary> = filtered.into_iter().skip(start).take(limit).collect();
3306    Ok((page, has_more, total_room_count_estimate))
3307}
3308
3309#[cfg(test)]
3310mod tests {
3311    use super::*;
3312
3313    const T0: &str = "2026-09-24T00:00:00+00:00";
3314    const ROOM: &str = "!testroom:example.org";
3315
3316    fn test_conn() -> Connection {
3317        let conn = Connection::open_in_memory().expect("in-memory sqlite");
3318        create_matrix_schema(&conn).expect("schema");
3319        conn
3320    }
3321
3322    fn ensure_legacy_dm_map_table(conn: &Connection) {
3323        conn.execute_batch(
3324            "CREATE TABLE IF NOT EXISTS legacy_dm_message_map (
3325                legacy_message_id INTEGER PRIMARY KEY,
3326                event_id TEXT NOT NULL UNIQUE REFERENCES events(event_id)
3327            );",
3328        )
3329        .expect("legacy map table for remaining unit tests");
3330    }
3331
3332    fn make_room(conn: &Connection) {
3333        create_room(conn, ROOM, RoomKind::Group, 1, T0, false, JoinRule::Invite, HistoryVisibility::Shared, None, None).expect("create room");
3334    }
3335
3336    // ---- P1 test 1: stream ordering ----
3337
3338    #[test]
3339    fn next_stream_id_is_strictly_monotonic_across_every_table() {
3340        let mut conn = test_conn();
3341        make_room(&conn);
3342
3343        let event = insert_timeline_event(&mut conn, "$event1", ROOM, 1, "m.room.message", "{}", 1000).expect("insert event");
3344        let account_data_stream = upsert_account_data(&mut conn, 1, GLOBAL_ACCOUNT_DATA_ROOM, "m.direct", "{}").expect("account data");
3345        let receipt_stream = upsert_receipt(&mut conn, ROOM, 1, ReceiptType::Read, "$event1", 1500).expect("receipt");
3346
3347        assert!(event.stream_id < account_data_stream, "event {} should precede account data {}", event.stream_id, account_data_stream);
3348        assert!(
3349            account_data_stream < receipt_stream,
3350            "account data {account_data_stream} should precede receipt {receipt_stream}"
3351        );
3352    }
3353
3354    // ---- P1 test 2: current_state replaces, events keeps history ----
3355
3356    #[test]
3357    fn apply_state_event_replaces_current_state_but_keeps_history_in_events() {
3358        let mut conn = test_conn();
3359        make_room(&conn);
3360
3361        apply_state_event(&mut conn, &StateEventWrite { event_id: "$e1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.name", state_key: "", content: r#"{"name":"first"}"#, origin_server_ts: 1000, now: T0 }).expect("apply first");
3362        apply_state_event(&mut conn, &StateEventWrite { event_id: "$e2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.name", state_key: "", content: r#"{"name":"second"}"#, origin_server_ts: 2000, now: T0 }).expect("apply second");
3363
3364        let current = current_state_event(&conn, ROOM, "m.room.name", "").expect("query").expect("row exists");
3365        assert_eq!(current.event_id, "$e2");
3366        assert_eq!(current.content, r#"{"name":"second"}"#);
3367
3368        let e1 = get_event(&conn, "$e1").expect("get e1").expect("row exists");
3369        let e2 = get_event(&conn, "$e2").expect("get e2").expect("row exists");
3370        assert_eq!(e1.content, r#"{"name":"first"}"#);
3371        assert_eq!(e2.content, r#"{"name":"second"}"#);
3372    }
3373
3374    // ---- P1 test 3: v11 redaction allow-list, table-driven ----
3375
3376    #[test]
3377    fn redact_event_strips_content_per_v11_allow_list_by_type() {
3378        let cases: Vec<(&str, &str, serde_json::Value)> = vec![
3379            (
3380                "m.room.member",
3381                r#"{"membership":"join","join_authorised_via_users_server":"@x:example.org","avatar_url":"mxc://x","displayname":"Bob","third_party_invite":{"display_name":"bob@example.com","signed":{"mxid":"@bob:example.org"}}}"#,
3382                serde_json::json!({
3383                    "membership": "join",
3384                    "join_authorised_via_users_server": "@x:example.org",
3385                    "third_party_invite": {"signed": {"mxid": "@bob:example.org"}}
3386                }),
3387            ),
3388            (
3389                "m.room.power_levels",
3390                r#"{"ban":50,"events":{},"events_default":0,"invite":50,"kick":50,"redact":50,"state_default":50,"users":{"@a:x":100},"users_default":0,"extra":"drop me"}"#,
3391                serde_json::json!({
3392                    "ban": 50, "events": {}, "events_default": 0, "invite": 50, "kick": 50,
3393                    "redact": 50, "state_default": 50, "users": {"@a:x": 100}, "users_default": 0
3394                }),
3395            ),
3396            (
3397                "m.room.history_visibility",
3398                r#"{"history_visibility":"shared","extra":"drop"}"#,
3399                serde_json::json!({"history_visibility": "shared"}),
3400            ),
3401            // `m.room.create` is NOT exercised end-to-end here: `redact_event`
3402            // refuses it outright (see `redact_event_refuses_create_and_
3403            // encryption_events`); its "keep everything" allow-list branch
3404            // is covered directly below, against `redact_content_per_v11`.
3405            (
3406                "m.room.message",
3407                r#"{"body":"hi","msgtype":"m.text"}"#,
3408                serde_json::json!({}),
3409            ),
3410        ];
3411
3412        for (event_type, content, expected) in cases {
3413            let mut conn = test_conn();
3414            make_room(&conn);
3415            insert_timeline_event(&mut conn, "$target", ROOM, 1, event_type, content, 1000).expect("insert target");
3416            redact_event(&mut conn, ROOM, "$target", "$redaction", 1, None, 2000).expect("redact");
3417
3418            let target = get_event(&conn, "$target").expect("get target").expect("row exists");
3419            let got: serde_json::Value = serde_json::from_str(&target.content).expect("parse stripped content");
3420            assert_eq!(got, expected, "event_type={event_type}");
3421            assert_eq!(target.redacted_by.as_deref(), Some("$redaction"), "event_type={event_type}");
3422        }
3423    }
3424
3425    #[test]
3426    fn redaction_content_carries_the_target_as_redacts_for_every_write_path() {
3427        let mut conn = test_conn();
3428        make_room(&conn);
3429        insert_timeline_event(&mut conn, "$t1", ROOM, 1, "m.room.message", "{}", 1000).expect("insert t1");
3430        insert_timeline_event(&mut conn, "$t2", ROOM, 1, "m.room.message", "{}", 1100).expect("insert t2");
3431        insert_timeline_event(&mut conn, "$t3", ROOM, 1, "m.room.message", "{}", 1200).expect("insert t3");
3432
3433        let plain = redact_event(&mut conn, ROOM, "$t1", "$r1", 1, None, 2000).expect("plain redact");
3434        assert_eq!(serde_json::from_str::<serde_json::Value>(&plain.content).expect("json"), serde_json::json!({ "redacts": "$t1" }));
3435        assert_eq!(plain.redacts.as_deref(), Some("$t1"));
3436
3437        let marked = redact_event_marked(
3438            &mut conn,
3439            &Redaction { room_id: ROOM, target_event_id: "$t2", redaction_event_id: "$r2", sender_user_id: 1, reason: Some("spam"), origin_server_ts: 2100 },
3440            &serde_json::json!({ "org.example.site_moderation": true }),
3441        )
3442            .expect("marked redact");
3443        assert_eq!(
3444            serde_json::from_str::<serde_json::Value>(&marked.content).expect("json"),
3445            serde_json::json!({ "redacts": "$t2", "reason": "spam", "org.example.site_moderation": true })
3446        );
3447
3448        let deduped = redact_event_deduped(&mut conn, "DEV1", "txn-1", ROOM, "$t3", "$r3", 1, None, 2200, T0).expect("deduped redact");
3449        let DedupedWrite::New(event) = deduped else { panic!("first write is new") };
3450        let stored = get_event(&conn, &event.event_id).expect("get").expect("row exists");
3451        assert_eq!(serde_json::from_str::<serde_json::Value>(&stored.content).expect("json")["redacts"], "$t3");
3452    }
3453
3454    #[test]
3455    fn redact_content_per_v11_keeps_everything_for_m_room_create() {
3456        let content = r#"{"room_version":"11","creator":"@a:example.org"}"#;
3457        let stripped = redact_content_per_v11("m.room.create", content).expect("strip");
3458        let got: serde_json::Value = serde_json::from_str(&stripped).expect("parse");
3459        assert_eq!(got, serde_json::json!({"room_version": "11", "creator": "@a:example.org"}));
3460    }
3461
3462    // ---- P1 test 4: txn dedup ----
3463
3464    #[test]
3465    fn txn_dedup_returns_the_same_event_id_on_a_repeated_txn_id() {
3466        let mut conn = test_conn();
3467        make_room(&conn);
3468
3469        assert_eq!(txn_dedup_lookup(&conn, 1, "DEV1", "txn-1").expect("lookup 1"), TxnDedupEntry::NotSeen);
3470
3471        let event = {
3472            let tx = conn.transaction().expect("tx");
3473            let row = TimelineEventRow {
3474                event_id: "$e1",
3475                room_id: ROOM,
3476                sender_user_id: 1,
3477                event_type: "m.room.message",
3478                content: "{}",
3479                origin_server_ts: 1000,
3480                txn_id: Some("txn-1"),
3481            };
3482            let event = insert_timeline_event_in_tx(&tx, &row).expect("insert");
3483            tx.commit().expect("commit");
3484            event
3485        };
3486        txn_dedup_record(&conn, 1, "DEV1", "txn-1", Some(&event.event_id), T0).expect("record");
3487
3488        assert_eq!(
3489            txn_dedup_lookup(&conn, 1, "DEV1", "txn-1").expect("lookup 2"),
3490            TxnDedupEntry::Seen(Some(event.event_id.clone()))
3491        );
3492
3493        let count: i64 = conn.query_row("SELECT COUNT(*) FROM events", [], |row| row.get(0)).expect("count");
3494        assert_eq!(count, 1, "a repeated txn_id must never create a second events row");
3495    }
3496
3497    // ---- P1 test 5: relations survive ciphertext content ----
3498
3499    #[test]
3500    fn relations_index_populated_from_cleartext_relates_to_even_when_content_is_ciphertext() {
3501        let mut conn = test_conn();
3502        make_room(&conn);
3503        insert_timeline_event(&mut conn, "$target", ROOM, 1, "m.room.message", "{}", 1000).expect("target");
3504
3505        let content = r#"{"algorithm":"m.megolm.v1.aes-sha2","ciphertext":"opaque-base64==","sender_key":"opaque","m.relates_to":{"rel_type":"m.annotation","event_id":"$target","key":"a"}}"#;
3506        let reaction = insert_timeline_event(&mut conn, "$reaction", ROOM, 2, "m.room.encrypted", content, 2000).expect("insert reaction");
3507
3508        let (rel_type, target_id, agg_key): (String, String, Option<String>) = conn
3509            .query_row(
3510                "SELECT rel_type, target_id, agg_key FROM relations WHERE event_id = ?1",
3511                params![reaction.event_id],
3512                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
3513            )
3514            .expect("relation row exists");
3515        assert_eq!(rel_type, "m.annotation");
3516        assert_eq!(target_id, "$target");
3517        assert_eq!(agg_key.as_deref(), Some("a"));
3518    }
3519
3520    // ---- extra test: global account data is unique ----
3521
3522    #[test]
3523    fn account_data_global_row_is_unique() {
3524        let mut conn = test_conn();
3525        upsert_account_data(&mut conn, 1, GLOBAL_ACCOUNT_DATA_ROOM, "m.direct", r#"{"v":1}"#).expect("first upsert");
3526        upsert_account_data(&mut conn, 1, GLOBAL_ACCOUNT_DATA_ROOM, "m.direct", r#"{"v":2}"#).expect("second upsert");
3527
3528        let count: i64 = conn
3529            .query_row(
3530                "SELECT COUNT(*) FROM account_data WHERE user_id = 1 AND room_id = ''",
3531                [],
3532                |row| row.get(0),
3533            )
3534            .expect("count");
3535        assert_eq!(count, 1, "two global upserts of the same type must leave one row");
3536
3537        let row = get_account_data(&conn, 1, GLOBAL_ACCOUNT_DATA_ROOM, "m.direct").expect("get").expect("row exists");
3538        assert_eq!(row.content, r#"{"v":2}"#);
3539    }
3540
3541    // ---- extra test: receipts never move backwards ----
3542
3543    #[test]
3544    fn receipt_never_moves_backwards() {
3545        let mut conn = test_conn();
3546        make_room(&conn);
3547        let e1 = insert_timeline_event(&mut conn, "$e1", ROOM, 1, "m.room.message", "{}", 1000).expect("e1");
3548        let e2 = insert_timeline_event(&mut conn, "$e2", ROOM, 1, "m.room.message", "{}", 2000).expect("e2");
3549
3550        upsert_receipt(&mut conn, ROOM, 9, ReceiptType::Read, &e2.event_id, 5000).expect("advance to e2");
3551        upsert_receipt(&mut conn, ROOM, 9, ReceiptType::Read, &e1.event_id, 6000).expect("attempted backward move is a no-op");
3552
3553        let receipt = get_receipt(&conn, ROOM, 9, ReceiptType::Read).expect("get").expect("row exists");
3554        assert_eq!(receipt.event_id, e2.event_id, "a receipt must never move back to an earlier event");
3555    }
3556
3557    // ---- P7 tests: notification_count (plan §3.5) ----
3558
3559    #[test]
3560    fn receipt_further_ahead_of_the_two_types_wins_for_notification_count() {
3561        let mut conn = test_conn();
3562        make_room(&conn);
3563        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000001", T0).expect("alice");
3564        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 900, now: T0 }).expect("alice joins");
3565
3566        insert_timeline_event(&mut conn, "$e1", ROOM, 2, "m.room.message", "{}", 1000).expect("e1");
3567        let e2 = insert_timeline_event(&mut conn, "$e2", ROOM, 2, "m.room.message", "{}", 2000).expect("e2");
3568        insert_timeline_event(&mut conn, "$e3", ROOM, 2, "m.room.message", "{}", 3000).expect("e3");
3569        let e4 = insert_timeline_event(&mut conn, "$e4", ROOM, 2, "m.room.message", "{}", 4000).expect("e4");
3570        let e5 = insert_timeline_event(&mut conn, "$e5", ROOM, 2, "m.room.message", "{}", 5000).expect("e5");
3571
3572        // m.read at e2, m.read.private FURTHER AHEAD at e4 — the private
3573        // receipt must win even though "private" has nothing to do with
3574        // ordering: notification_count takes the MAX of the two positions.
3575        upsert_receipt(&mut conn, ROOM, 1, ReceiptType::Read, &e2.event_id, 2500).expect("read receipt");
3576        upsert_receipt(&mut conn, ROOM, 1, ReceiptType::ReadPrivate, &e4.event_id, 4500).expect("private receipt further ahead");
3577
3578        let room = get_room(&conn, ROOM).expect("get room").expect("room exists");
3579        assert_eq!(notification_count(&conn, &room, 1).expect("count"), 1, "only $e5 is after the further-ahead receipt ($e4)");
3580
3581        // Reverse: m.read now further ahead than m.read.private — m.read
3582        // must win this time.
3583        upsert_receipt(&mut conn, ROOM, 1, ReceiptType::Read, &e5.event_id, 5500).expect("read receipt advances past e5");
3584        assert_eq!(notification_count(&conn, &room, 1).expect("count"), 0, "m.read now covers every message");
3585    }
3586
3587    #[test]
3588    fn notification_count_ignores_own_state_and_reactions() {
3589        let mut conn = test_conn();
3590        make_room(&conn);
3591        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000001", T0).expect("alice");
3592        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 900, now: T0 }).expect("alice joins");
3593
3594        // Bob's message-like sends — these count.
3595        insert_timeline_event(&mut conn, "$bob-msg", ROOM, 2, "m.room.message", "{}", 1000).expect("bob message");
3596        insert_timeline_event(&mut conn, "$bob-enc", ROOM, 2, "m.room.encrypted", "{}", 1100).expect("bob encrypted");
3597        insert_timeline_event(&mut conn, "$bob-legacy", ROOM, 2, "org.example.legacy_dm", "{}", 1200).expect("bob legacy dm");
3598
3599        // Excluded: alice's own message, a state event, and a reaction.
3600        insert_timeline_event(&mut conn, "$alice-msg", ROOM, 1, "m.room.message", "{}", 1300).expect("alice's own message");
3601        apply_state_event(&mut conn, &StateEventWrite { event_id: "$name", room_id: ROOM, sender_user_id: 2, event_type: "m.room.name", state_key: "", content: r#"{"name":"x"}"#, origin_server_ts: 1400, now: T0 }).expect("state event");
3602        insert_timeline_event(&mut conn, "$reaction", ROOM, 2, "m.reaction", "{}", 1500).expect("reaction");
3603
3604        let room = get_room(&conn, ROOM).expect("get room").expect("room exists");
3605        assert_eq!(notification_count(&conn, &room, 1).expect("count"), 3, "only bob's 3 message-like sends count");
3606    }
3607
3608    #[test]
3609    fn notification_count_is_zero_for_a_non_member() {
3610        let mut conn = test_conn();
3611        make_room(&conn);
3612        insert_timeline_event(&mut conn, "$e1", ROOM, 2, "m.room.message", "{}", 1000).expect("e1");
3613
3614        let room = get_room(&conn, ROOM).expect("get room").expect("room exists");
3615        assert_eq!(notification_count(&conn, &room, 999).expect("count"), 0, "a caller with no membership row at all sees nothing");
3616    }
3617
3618    // ---- extra test: duplicate annotation refused ----
3619
3620    #[test]
3621    fn duplicate_annotation_is_refused() {
3622        let mut conn = test_conn();
3623        make_room(&conn);
3624        insert_timeline_event(&mut conn, "$target", ROOM, 1, "m.room.message", "{}", 1000).expect("target");
3625
3626        let content = r#"{"m.relates_to":{"rel_type":"m.annotation","event_id":"$target","key":"a"}}"#;
3627        insert_timeline_event(&mut conn, "$react1", ROOM, 5, "m.reaction", content, 2000).expect("first reaction");
3628
3629        let err = insert_timeline_event(&mut conn, "$react2", ROOM, 5, "m.reaction", content, 3000).unwrap_err();
3630        assert!(matches!(err, MatrixStoreError::DuplicateAnnotation));
3631
3632        // The whole transaction must have rolled back — no partial events row.
3633        assert_eq!(get_event(&conn, "$react2").expect("get"), None);
3634
3635        // A DIFFERENT sender reacting with the same key is a distinct,
3636        // accepted annotation (not covered by the refusal).
3637        let other = insert_timeline_event(&mut conn, "$react3", ROOM, 6, "m.reaction", content, 4000).expect("different sender reacts");
3638        assert_eq!(other.event_id, "$react3");
3639    }
3640
3641    // ---- manager-review fixes (2026-09-24 follow-up) ----
3642
3643    #[test]
3644    fn reannotation_after_redaction_is_allowed() {
3645        let mut conn = test_conn();
3646        make_room(&conn);
3647        insert_timeline_event(&mut conn, "$target", ROOM, 1, "m.room.message", "{}", 1000).expect("target");
3648
3649        let content = r#"{"m.relates_to":{"rel_type":"m.annotation","event_id":"$target","key":"a"}}"#;
3650        let first = insert_timeline_event(&mut conn, "$react1", ROOM, 5, "m.reaction", content, 2000).expect("first reaction");
3651
3652        redact_event(&mut conn, ROOM, &first.event_id, "$redaction", 1, None, 2500).expect("redact the reaction");
3653
3654        // The same sender may now react with the same key again — the
3655        // redacted prior annotation no longer counts as a duplicate.
3656        let second = insert_timeline_event(&mut conn, "$react2", ROOM, 5, "m.reaction", content, 3000)
3657            .expect("re-annotation after redaction must succeed");
3658        assert_eq!(second.event_id, "$react2");
3659    }
3660
3661    #[test]
3662    fn relation_target_must_exist() {
3663        let mut conn = test_conn();
3664        make_room(&conn);
3665
3666        let content = r#"{"m.relates_to":{"rel_type":"m.annotation","event_id":"$missing","key":"a"}}"#;
3667        let err = insert_timeline_event(&mut conn, "$react1", ROOM, 5, "m.reaction", content, 2000).unwrap_err();
3668        assert!(matches!(err, MatrixStoreError::InvalidRelationTarget(ref id) if id == "$missing"));
3669    }
3670
3671    #[test]
3672    fn relation_target_must_be_in_the_same_room() {
3673        let mut conn = test_conn();
3674        make_room(&conn);
3675        let other_room = format!("!other:{}", matrix_server_name());
3676        create_room(&conn, &other_room, RoomKind::Group, 1, T0, false, JoinRule::Invite, HistoryVisibility::Shared, None, None)
3677            .expect("other room");
3678        insert_timeline_event(&mut conn, "$target", &other_room, 1, "m.room.message", "{}", 1000).expect("target in other room");
3679
3680        let content = r#"{"m.relates_to":{"rel_type":"m.annotation","event_id":"$target","key":"a"}}"#;
3681        let err = insert_timeline_event(&mut conn, "$react1", ROOM, 5, "m.reaction", content, 2000).unwrap_err();
3682        assert!(matches!(err, MatrixStoreError::InvalidRelationTarget(ref id) if id == "$target"));
3683    }
3684
3685    #[test]
3686    fn redact_event_refuses_a_target_in_a_different_room() {
3687        let mut conn = test_conn();
3688        make_room(&conn);
3689        let other_room = format!("!other:{}", matrix_server_name());
3690        create_room(&conn, &other_room, RoomKind::Group, 1, T0, false, JoinRule::Invite, HistoryVisibility::Shared, None, None)
3691            .expect("other room");
3692        insert_timeline_event(&mut conn, "$target", &other_room, 1, "m.room.message", "{}", 1000).expect("target in other room");
3693
3694        let err = redact_event(&mut conn, ROOM, "$target", "$redaction", 1, None, 2000).unwrap_err();
3695        assert!(matches!(err, MatrixStoreError::WrongRoom(ref id) if id == "$target"));
3696    }
3697
3698    #[test]
3699    fn redact_event_refuses_create_and_encryption_events() {
3700        for event_type in ["m.room.create", "m.room.encryption"] {
3701            let mut conn = test_conn();
3702            make_room(&conn);
3703            apply_state_event(&mut conn, &StateEventWrite { event_id: "$target", room_id: ROOM, sender_user_id: 1, event_type, state_key: "", content: r#"{"a":1}"#, origin_server_ts: 1000, now: T0 }).expect("apply state event");
3704
3705            let err = redact_event(&mut conn, ROOM, "$target", "$redaction", 1, None, 2000).unwrap_err();
3706            assert!(
3707                matches!(err, MatrixStoreError::UnredactableEvent(ref t) if t == event_type),
3708                "event_type={event_type}"
3709            );
3710        }
3711    }
3712
3713    #[test]
3714    fn upsert_receipt_refuses_a_target_in_a_different_room() {
3715        let mut conn = test_conn();
3716        make_room(&conn);
3717        let other_room = format!("!other:{}", matrix_server_name());
3718        create_room(&conn, &other_room, RoomKind::Group, 1, T0, false, JoinRule::Invite, HistoryVisibility::Shared, None, None)
3719            .expect("other room");
3720        let event = insert_timeline_event(&mut conn, "$e1", &other_room, 1, "m.room.message", "{}", 1000).expect("event in other room");
3721
3722        let err = upsert_receipt(&mut conn, ROOM, 9, ReceiptType::Read, &event.event_id, 5000).unwrap_err();
3723        assert!(matches!(err, MatrixStoreError::WrongRoom(ref id) if id == &event.event_id));
3724    }
3725
3726    // ---- extra test: mxid parsing refuses a foreign server ----
3727
3728    #[test]
3729    fn mxid_parse_refuses_foreign_server() {
3730        assert_eq!(public_id_from_mxid("@abc123:example.org"), Ok("abc123"));
3731        assert_eq!(public_id_from_mxid("@abc123:otherserver.example"), Err(MatrixIdError::ForeignServerName));
3732        // Local aliases: server-name check only (enabled once per process).
3733        let aliases = ["chat.example", "m4a.example.net", "m4a.example.org"];
3734        set_local_aliases(aliases.iter().map(|s| s.to_string()));
3735        for name in aliases {
3736            assert_eq!(public_id_from_mxid(&format!("@abc123:{name}")), Ok("abc123"));
3737        }
3738        assert_eq!(public_id_from_mxid("@abc123:evil.example"), Err(MatrixIdError::ForeignServerName));
3739        assert_eq!(mxid_for_public_id("abc123"), format!("@abc123:{}", matrix_server_name()), "minting never uses an alias");
3740        assert_eq!(public_id_from_mxid("abc123:example.org"), Err(MatrixIdError::MissingSigil));
3741        assert_eq!(public_id_from_mxid("@abc123"), Err(MatrixIdError::MissingServerName));
3742    }
3743
3744    // ---- supporting coverage for the CRUD surface the tests above don't already exercise ----
3745
3746    #[test]
3747    fn ensure_matrix_user_is_idempotent_and_refuses_a_reserved_localpart() {
3748        let conn = test_conn();
3749        let mxid = ensure_matrix_user(&conn, 1, "abc123", T0).expect("first ensure");
3750        assert_eq!(mxid, "@abc123:example.org");
3751        let mxid_again = ensure_matrix_user(&conn, 1, "abc123", T0).expect("second ensure is a no-op");
3752        assert_eq!(mxid_again, mxid);
3753        assert_eq!(mxid_of(&conn, 1).expect("mxid_of"), Some(mxid.clone()));
3754        assert_eq!(user_id_of(&conn, &mxid).expect("user_id_of"), Some(1));
3755
3756        let err = ensure_matrix_user(&conn, 2, "_bridge_evil", T0).unwrap_err();
3757        assert!(matches!(err, MatrixStoreError::ReservedLocalpart));
3758    }
3759
3760    #[test]
3761    fn apply_state_event_member_refreshes_room_members_and_power_levels() {
3762        let mut conn = test_conn();
3763        make_room(&conn);
3764        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000001", T0).expect("alice");
3765        let bob = ensure_matrix_user(&conn, 2, "bob000000000000000000000000002", T0).expect("bob");
3766
3767        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("alice joins");
3768        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"invite"}"#, origin_server_ts: 1100, now: T0 }).expect("bob invited");
3769
3770        let joined = room_members(&conn, ROOM, Some(Membership::Join)).expect("joined members");
3771        assert_eq!(joined.len(), 1);
3772        assert_eq!(joined[0].user_id, 1);
3773        assert_eq!(joined[0].power_level, None);
3774
3775        let power_levels_content = serde_json::json!({"users": {alice.clone(): 100}, "users_default": 0}).to_string();
3776        apply_state_event(&mut conn, &StateEventWrite { event_id: "$pl1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.power_levels", state_key: "", content: &power_levels_content, origin_server_ts: 1200, now: T0 }).expect("power levels");
3777
3778        let alice_row = room_members(&conn, ROOM, None)
3779            .expect("all members")
3780            .into_iter()
3781            .find(|m| m.user_id == 1)
3782            .expect("alice row");
3783        assert_eq!(alice_row.power_level, Some(100));
3784
3785        let bob_rooms = rooms_for_user(&conn, 2, Some(Membership::Invite)).expect("bob's invited rooms");
3786        assert_eq!(bob_rooms, vec![ROOM.to_string()]);
3787    }
3788
3789    #[test]
3790    fn events_in_room_after_and_before_page_in_the_documented_order() {
3791        let mut conn = test_conn();
3792        make_room(&conn);
3793        let e1 = insert_timeline_event(&mut conn, "$e1", ROOM, 1, "m.room.message", "{}", 1000).expect("e1");
3794        let e2 = insert_timeline_event(&mut conn, "$e2", ROOM, 1, "m.room.message", "{}", 2000).expect("e2");
3795        let e3 = insert_timeline_event(&mut conn, "$e3", ROOM, 1, "m.room.message", "{}", 3000).expect("e3");
3796
3797        let after = events_in_room_after(&conn, ROOM, e1.stream_id, 10).expect("after");
3798        assert_eq!(after.iter().map(|e| e.event_id.clone()).collect::<Vec<_>>(), vec![e2.event_id.clone(), e3.event_id.clone()]);
3799
3800        let before = events_in_room_before(&conn, ROOM, e3.stream_id, 10).expect("before");
3801        assert_eq!(before.iter().map(|e| e.event_id.clone()).collect::<Vec<_>>(), vec![e2.event_id.clone(), e1.event_id.clone()]);
3802
3803        assert_eq!(max_stream_id(&conn).expect("max"), e3.stream_id);
3804    }
3805
3806    #[test]
3807    fn txn_dedup_lookup_of_a_to_device_send_has_no_event_id() {
3808        let conn = test_conn();
3809        assert_eq!(txn_dedup_lookup(&conn, 1, "DEV1", "txn-td").expect("lookup"), TxnDedupEntry::NotSeen);
3810        txn_dedup_record(&conn, 1, "DEV1", "txn-td", None, T0).expect("record to-device send");
3811        assert_eq!(txn_dedup_lookup(&conn, 1, "DEV1", "txn-td").expect("lookup again"), TxnDedupEntry::Seen(None));
3812    }
3813
3814    #[test]
3815    fn filters_create_and_get_are_scoped_to_their_owner() {
3816        let conn = test_conn();
3817        let filter_id = create_filter(&conn, 1, r#"{"room":{"timeline":{"limit":20}}}"#).expect("create");
3818        assert_eq!(get_filter(&conn, 1, filter_id).expect("owner reads it"), Some(r#"{"room":{"timeline":{"limit":20}}}"#.to_string()));
3819        assert_eq!(get_filter(&conn, 2, filter_id).expect("a different user cannot"), None);
3820    }
3821
3822    #[test]
3823    fn legacy_dm_message_map_insert_and_get_round_trip() {
3824        let mut conn = test_conn();
3825        make_room(&conn);
3826        ensure_legacy_dm_map_table(&conn);
3827        let event = insert_timeline_event(&mut conn, "$legacy1", ROOM, 1, "org.example.legacy_dm", "{}", 1000).expect("insert");
3828        insert_legacy_dm_message_map(&conn, 42, &event.event_id).expect("map insert");
3829        assert_eq!(legacy_dm_message_event_id(&conn, 42).expect("map get"), Some(event.event_id));
3830        assert_eq!(legacy_dm_message_event_id(&conn, 999).expect("missing"), None);
3831    }
3832
3833    // ---- P4 test: power-level defaults ----
3834
3835    #[test]
3836    fn power_level_defaults_apply_when_fields_missing() {
3837        let pl = serde_json::json!({});
3838        assert_eq!(user_level(&pl, "@nobody:example.org"), 0, "users_default defaults to 0");
3839        assert_eq!(event_level(&pl, "m.room.message", false), 0, "events_default defaults to 0");
3840        assert_eq!(event_level(&pl, "m.room.name", true), 50, "state_default defaults to 50");
3841        for action in [PowerAction::Invite, PowerAction::Kick, PowerAction::Ban, PowerAction::Redact, PowerAction::StateDefault] {
3842            assert!(!can(&pl, action, "@nobody:example.org"), "level 0 must not reach the default 50 threshold for {action:?}");
3843        }
3844
3845        let pl_with_creator = serde_json::json!({ "users": { "@creator:example.org": 100 } });
3846        assert!(can(&pl_with_creator, PowerAction::Ban, "@creator:example.org"));
3847        assert!(can(&pl_with_creator, PowerAction::StateDefault, "@creator:example.org"));
3848    }
3849
3850    #[test]
3851    fn event_level_uses_the_events_type_override_before_falling_back_to_a_default() {
3852        let pl = serde_json::json!({ "events": { "m.room.name": 60 }, "events_default": 0, "state_default": 50 });
3853        assert_eq!(event_level(&pl, "m.room.name", true), 60, "an explicit events[type] override wins");
3854        assert_eq!(event_level(&pl, "m.room.topic", true), 50, "an unlisted state type falls back to state_default");
3855        assert_eq!(event_level(&pl, "m.room.message", false), 0, "an unlisted timeline type falls back to events_default");
3856    }
3857
3858    // ---- manager review 2026-09-24: can_act_on / validate_power_levels_change ----
3859
3860    #[test]
3861    fn can_act_on_requires_strictly_greater_level_except_self_leave() {
3862        let pl = serde_json::json!({
3863            "users": { "@owner:example.org": 100, "@admin:example.org": 50, "@peer:example.org": 50 },
3864            "kick": 50,
3865            "ban": 50,
3866        });
3867
3868        // A level-50 admin reaches the flat kick threshold (50) but must
3869        // NOT be able to kick the level-100 owner.
3870        assert!(!can_act_on(&pl, PowerAction::Kick, "@admin:example.org", "@owner:example.org", false));
3871        // Nor an equal-level peer.
3872        assert!(!can_act_on(&pl, PowerAction::Kick, "@admin:example.org", "@peer:example.org", false));
3873        // The owner CAN kick the admin (100 > 50).
3874        assert!(can_act_on(&pl, PowerAction::Kick, "@owner:example.org", "@admin:example.org", false));
3875        // A member may always leave (reject an invite / self-kick) regardless
3876        // of level, when self_leave is set and sender == target.
3877        assert!(can_act_on(&pl, PowerAction::Kick, "@admin:example.org", "@admin:example.org", true));
3878        // ...but NOT when the action does not actually result in a leave.
3879        assert!(!can_act_on(&pl, PowerAction::Ban, "@admin:example.org", "@admin:example.org", false));
3880    }
3881
3882    #[test]
3883    fn validate_power_levels_change_refuses_raising_self_above_own_level() {
3884        let old = serde_json::json!({ "users": { "@admin:example.org": 50 } });
3885        let new = serde_json::json!({ "users": { "@admin:example.org": 100 } });
3886        assert!(validate_power_levels_change(&old, &new, "@admin:example.org").is_err());
3887    }
3888
3889    #[test]
3890    fn validate_power_levels_change_refuses_demoting_a_peer_at_an_equal_level() {
3891        let old = serde_json::json!({ "users": { "@a:example.org": 50, "@b:example.org": 50 } });
3892        let new = serde_json::json!({ "users": { "@a:example.org": 50, "@b:example.org": 0 } });
3893        assert!(validate_power_levels_change(&old, &new, "@a:example.org").is_err());
3894    }
3895
3896    #[test]
3897    fn validate_power_levels_change_allows_demoting_self() {
3898        let old = serde_json::json!({ "users": { "@admin:example.org": 50 } });
3899        let new = serde_json::json!({ "users": { "@admin:example.org": 10 } });
3900        assert!(validate_power_levels_change(&old, &new, "@admin:example.org").is_ok());
3901    }
3902
3903    #[test]
3904    fn validate_power_levels_change_refuses_raising_events_default_above_own_level() {
3905        let old = serde_json::json!({ "users": { "@admin:example.org": 50 }, "events_default": 0 });
3906        let new = serde_json::json!({ "users": { "@admin:example.org": 50 }, "events_default": 60 });
3907        assert!(validate_power_levels_change(&old, &new, "@admin:example.org").is_err());
3908    }
3909
3910    #[test]
3911    fn validate_power_levels_change_allows_the_owner_changing_anything_up_to_their_own_level() {
3912        let old = serde_json::json!({ "users": { "@owner:example.org": 100, "@a:example.org": 50 } });
3913        let new = serde_json::json!({
3914            "users": { "@owner:example.org": 100, "@a:example.org": 90 },
3915            "ban": 100,
3916            "kick": 100,
3917            "events_default": 100,
3918        });
3919        assert!(validate_power_levels_change(&old, &new, "@owner:example.org").is_ok());
3920    }
3921
3922    #[test]
3923    fn validate_power_levels_change_refuses_a_scalar_field_change_above_own_level() {
3924        let old = serde_json::json!({ "users": { "@admin:example.org": 50 }, "ban": 50 });
3925        let new = serde_json::json!({ "users": { "@admin:example.org": 50 }, "ban": 75 });
3926        assert!(validate_power_levels_change(&old, &new, "@admin:example.org").is_err());
3927    }
3928
3929    #[test]
3930    fn validate_power_levels_change_ignores_unchanged_fields() {
3931        let old = serde_json::json!({ "users": { "@admin:example.org": 50 }, "ban": 50, "events": { "m.room.name": 40 } });
3932        let new = old.clone();
3933        assert!(validate_power_levels_change(&old, &new, "@admin:example.org").is_ok());
3934    }
3935
3936    // ---- P5 test: create_room_with_state is one transaction ----
3937
3938    fn bootstrap<'a>(room_id: &'a str, kind: RoomKind, creator: i64) -> RoomBootstrap<'a> {
3939        RoomBootstrap {
3940            room_id,
3941            kind,
3942            creator_user_id: creator,
3943            created_at: T0,
3944            is_encrypted: true,
3945            join_rule: if kind == RoomKind::Channel { JoinRule::Public } else { JoinRule::Invite },
3946            history_visibility: HistoryVisibility::Shared,
3947            dm_pair_key: None,
3948            legacy_dm_id: None,
3949        }
3950    }
3951
3952    fn state_event(event_id: &str, sender: i64, event_type: &str, state_key: &str, content: &str) -> NewStateEvent {
3953        NewStateEvent {
3954            event_id: event_id.to_string(),
3955            sender_user_id: sender,
3956            event_type: event_type.to_string(),
3957            state_key: state_key.to_string(),
3958            content: content.to_string(),
3959        }
3960    }
3961
3962    #[test]
3963    fn create_room_with_state_inserts_the_room_and_every_bootstrap_event_atomically() {
3964        let mut conn = test_conn();
3965        let creator_mxid = ensure_matrix_user(&conn, 1, "creator0000000000000000000001", T0).expect("creator");
3966        let room_id = "!batch:example.org";
3967
3968        let events = vec![
3969            state_event("$create", 1, "m.room.create", "", r#"{"room_version":"11"}"#),
3970            state_event("$m1", 1, "m.room.member", &creator_mxid, r#"{"membership":"join"}"#),
3971            state_event(
3972                "$pl",
3973                1,
3974                "m.room.power_levels",
3975                "",
3976                &serde_json::json!({"users": {creator_mxid.clone(): 100}, "users_default": 0}).to_string(),
3977            ),
3978        ];
3979
3980        let (room, applied) = create_room_with_state(&mut conn, bootstrap(room_id, RoomKind::Group, 1), &events, 1000).expect("create batch");
3981        assert_eq!(room.id, room_id);
3982        assert_eq!(applied.len(), 3);
3983        assert!(get_room(&conn, room_id).expect("get room").is_some());
3984        for event in &applied {
3985            assert!(get_event(&conn, &event.event_id).expect("get event").is_some());
3986        }
3987        let creator_row = room_member(&conn, room_id, 1).expect("member row").expect("row exists");
3988        assert_eq!(creator_row.power_level, Some(100), "power_levels applied after the member row existed");
3989    }
3990
3991    #[test]
3992    fn create_room_is_atomic_on_failure() {
3993        let mut conn = test_conn();
3994        let creator_mxid = ensure_matrix_user(&conn, 1, "creator0000000000000000000002", T0).expect("creator");
3995        let room_id = "!atomic:example.org";
3996
3997        let events = vec![
3998            state_event("$create", 1, "m.room.create", "", r#"{"room_version":"11"}"#),
3999            state_event("$m1", 1, "m.room.member", &creator_mxid, r#"{"membership":"join"}"#),
4000            // Never `ensure_matrix_user`'d — `refresh_room_member` must
4001            // refuse this with `UnknownMxid`, rolling back the WHOLE batch.
4002            state_event("$bad", 1, "m.room.member", "@ghost:example.org", r#"{"membership":"invite"}"#),
4003        ];
4004
4005        let err = create_room_with_state(&mut conn, bootstrap(room_id, RoomKind::Group, 1), &events, 1000).unwrap_err();
4006        assert!(matches!(err, MatrixStoreError::UnknownMxid(ref m) if m == "@ghost:example.org"));
4007
4008        assert_eq!(get_room(&conn, room_id).expect("get room"), None, "a failed batch must leave no room row");
4009        assert_eq!(get_event(&conn, "$create").expect("get"), None, "a failed batch must leave no event rows at all");
4010        assert_eq!(get_event(&conn, "$m1").expect("get"), None);
4011    }
4012
4013    // ---- P5 test: DM pair-key reuse frees a dead pair's slot ----
4014
4015    #[test]
4016    fn dm_pair_key_reuse_is_freed_once_the_room_is_not_reused() {
4017        let conn = test_conn();
4018        create_room(&conn, "!dm1:example.org", RoomKind::Dm, 1, T0, true, JoinRule::Invite, HistoryVisibility::Shared, Some("1:2"), None)
4019            .expect("first dm room");
4020        assert_eq!(room_by_dm_pair_key(&conn, "1:2").expect("lookup").map(|r| r.id), Some("!dm1:example.org".to_string()));
4021
4022        clear_dm_pair_key(&conn, "!dm1:example.org").expect("clear");
4023        assert_eq!(room_by_dm_pair_key(&conn, "1:2").expect("lookup after clear"), None);
4024
4025        // The freed pair key can now be claimed by a second room.
4026        create_room(&conn, "!dm2:example.org", RoomKind::Dm, 1, T0, true, JoinRule::Invite, HistoryVisibility::Shared, Some("1:2"), None)
4027            .expect("second dm room reuses the freed pair key");
4028        assert_eq!(room_by_dm_pair_key(&conn, "1:2").expect("lookup").map(|r| r.id), Some("!dm2:example.org".to_string()));
4029    }
4030
4031    // ---- P5 test: room_member / forget_membership ----
4032
4033    #[test]
4034    fn room_member_finds_the_one_row_forget_membership_deletes_only_when_left() {
4035        let mut conn = test_conn();
4036        make_room(&conn);
4037        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000099", T0).expect("alice");
4038        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("join");
4039
4040        assert_eq!(room_member(&conn, ROOM, 1).expect("member").map(|m| m.membership), Some(Membership::Join));
4041        assert_eq!(room_member(&conn, ROOM, 999).expect("no such member"), None);
4042
4043        // Still joined — forget must refuse (0 rows deleted), matching the
4044        // route-layer gate "Member with membership='leave'".
4045        assert_eq!(forget_membership(&conn, ROOM, 1).expect("forget while joined"), 0);
4046        assert!(room_member(&conn, ROOM, 1).expect("still a member").is_some());
4047
4048        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"leave"}"#, origin_server_ts: 2000, now: T0 }).expect("leave");
4049        assert_eq!(forget_membership(&conn, ROOM, 1).expect("forget after leaving"), 1);
4050        assert_eq!(room_member(&conn, ROOM, 1).expect("gone"), None);
4051    }
4052
4053    // ---- P5 test: state_events_of_type_at reconstructs a point-in-time projection ----
4054
4055    #[test]
4056    fn state_events_of_type_at_excludes_state_keys_created_after_the_cutoff() {
4057        let mut conn = test_conn();
4058        make_room(&conn);
4059        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000098", T0).expect("alice");
4060        let bob = ensure_matrix_user(&conn, 2, "bob0000000000000000000000000098", T0).expect("bob");
4061
4062        let e1 = apply_state_event(&mut conn, &StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("alice joins");
4063        let cutoff = e1.stream_id;
4064        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"invite"}"#, origin_server_ts: 1100, now: T0 }).expect("bob invited later");
4065
4066        let at_cutoff = state_events_of_type_at(&conn, ROOM, "m.room.member", cutoff).expect("at cutoff");
4067        assert_eq!(at_cutoff.len(), 1, "bob's invite lands strictly after the cutoff and must be excluded");
4068        assert_eq!(at_cutoff[0].event_id, "$m1");
4069
4070        let after_both = state_events_of_type_at(&conn, ROOM, "m.room.member", cutoff + 1).expect("after both");
4071        assert_eq!(after_both.len(), 2);
4072    }
4073
4074    // ---- P5 test: stripped_invite_state ----
4075
4076    #[test]
4077    fn stripped_invite_state_includes_room_basics_and_the_inviters_own_member_event() {
4078        let mut conn = test_conn();
4079        make_room(&conn);
4080        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000097", T0).expect("alice");
4081        let bob = ensure_matrix_user(&conn, 2, "bob0000000000000000000000000097", T0).expect("bob");
4082
4083        apply_state_event(&mut conn, &StateEventWrite { event_id: "$create", room_id: ROOM, sender_user_id: 1, event_type: "m.room.create", state_key: "", content: r#"{"room_version":"11"}"#, origin_server_ts: 900, now: T0 }).expect("create");
4084        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("alice joins");
4085        apply_state_event(&mut conn, &StateEventWrite { event_id: "$jr", room_id: ROOM, sender_user_id: 1, event_type: "m.room.join_rules", state_key: "", content: r#"{"join_rule":"invite"}"#, origin_server_ts: 1100, now: T0 }).expect("join rules");
4086        apply_state_event(&mut conn, &StateEventWrite { event_id: "$m2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"invite"}"#, origin_server_ts: 1200, now: T0 }).expect("bob invited");
4087
4088        let stripped = stripped_invite_state(&conn, ROOM, 1).expect("stripped state");
4089        let types: Vec<&str> = stripped.iter().map(|v| v["type"].as_str().expect("type")).collect();
4090        assert!(types.contains(&"m.room.create"));
4091        assert!(types.contains(&"m.room.join_rules"));
4092        assert!(!types.contains(&"m.room.encryption"), "no encryption event exists in this room");
4093
4094        let inviter_member = stripped
4095            .iter()
4096            .find(|v| v["type"] == "m.room.member" && v["state_key"] == alice)
4097            .expect("the inviter's own member event is included");
4098        assert_eq!(inviter_member["sender"], alice);
4099        assert_eq!(inviter_member["content"]["membership"], "join");
4100    }
4101
4102    // ---- P9 support: user_ids_with_leave_transition_in_rooms ----
4103
4104    #[test]
4105    fn user_ids_with_leave_transition_in_rooms_finds_only_leave_and_ban_inside_the_window() {
4106        let mut conn = test_conn();
4107        make_room(&conn);
4108        let alice = ensure_matrix_user(&conn, 2, "alice00000000000000000000000097", T0).expect("alice");
4109        let bob = ensure_matrix_user(&conn, 3, "bob0000000000000000000000000097", T0).expect("bob");
4110        let carol = ensure_matrix_user(&conn, 4, "carol0000000000000000000000097a", T0).expect("carol");
4111
4112        apply_state_event(&mut conn, &StateEventWrite { event_id: "$a1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("alice joins");
4113        let boundary = apply_state_event(&mut conn, &StateEventWrite { event_id: "$b1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"join"}"#, origin_server_ts: 1100, now: T0 })
4114            .expect("bob joins")
4115            .stream_id;
4116        apply_state_event(&mut conn, &StateEventWrite { event_id: "$a2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"leave"}"#, origin_server_ts: 1200, now: T0 }).expect("alice leaves");
4117        apply_state_event(&mut conn, &StateEventWrite { event_id: "$b2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"ban"}"#, origin_server_ts: 1300, now: T0 }).expect("bob banned");
4118        apply_state_event(&mut conn, &StateEventWrite { event_id: "$c1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &carol, content: r#"{"membership":"join"}"#, origin_server_ts: 1400, now: T0 }).expect("carol joins (not a departure)");
4119
4120        let mut left = user_ids_with_leave_transition_in_rooms(&conn, &[ROOM.to_string()], boundary, i64::MAX).expect("query");
4121        left.sort_unstable();
4122        assert_eq!(left, vec![2, 3], "alice (leave) and bob (ban) both count; carol's join does not");
4123
4124        let empty = user_ids_with_leave_transition_in_rooms(&conn, &[], 0, i64::MAX).expect("empty room set");
4125        assert!(empty.is_empty());
4126    }
4127
4128    // ---- P16 S-f: point-in-time membership and the window's member candidates ----
4129
4130    #[test]
4131    fn membership_at_reads_the_state_as_of_a_stream_position() {
4132        let mut conn = test_conn();
4133        make_room(&conn);
4134        let alice = ensure_matrix_user(&conn, 2, "alice00000000000000000000000081", T0).expect("alice");
4135        let bob = ensure_matrix_user(&conn, 3, "bob0000000000000000000000000081", T0).expect("bob");
4136
4137        let alice_joined = apply_state_event(&mut conn, &StateEventWrite { event_id: "$a1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("alice joins").stream_id;
4138        let bob_invited = apply_state_event(&mut conn, &StateEventWrite { event_id: "$b1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"invite"}"#, origin_server_ts: 1100, now: T0 }).expect("bob invited").stream_id;
4139        let bob_joined = apply_state_event(&mut conn, &StateEventWrite { event_id: "$b2", room_id: ROOM, sender_user_id: 3, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"join","displayname":"Bob"}"#, origin_server_ts: 1200, now: T0 }).expect("bob joins").stream_id;
4140
4141        assert_eq!(membership_at(&conn, ROOM, &bob, alice_joined).expect("query"), None, "bob had no member event yet");
4142        assert_eq!(membership_at(&conn, ROOM, &bob, bob_invited).expect("query"), Some(Membership::Invite));
4143        assert_eq!(membership_at(&conn, ROOM, &bob, bob_joined).expect("query"), Some(Membership::Join));
4144        assert_eq!(membership_at(&conn, ROOM, &alice, bob_invited).expect("query"), Some(Membership::Join));
4145
4146        let mut window = member_state_keys_in_window(&conn, ROOM, bob_invited, bob_joined).expect("window");
4147        window.sort();
4148        assert_eq!(window, vec![bob.clone()], "only bob has a member event after `bob_invited`");
4149        let mut all = member_state_keys_in_window(&conn, ROOM, 0, bob_joined).expect("whole history");
4150        all.sort();
4151        let mut expected = vec![alice, bob];
4152        expected.sort();
4153        assert_eq!(all, expected);
4154
4155        let rooms = [ROOM.to_string(), "!other:example.org".to_string()];
4156        assert_eq!(rooms_with_member_events_in_window(&conn, &rooms, bob_invited, bob_joined).expect("rooms"), vec![ROOM.to_string()]);
4157        assert!(rooms_with_member_events_in_window(&conn, &rooms, bob_joined, i64::MAX).expect("rooms after the last member event").is_empty());
4158        assert!(rooms_with_member_events_in_window(&conn, &[], 0, i64::MAX).expect("empty room set").is_empty());
4159    }
4160
4161    // ---- a changed label is re-stamped into the user's member events ----
4162
4163    const ROOM_A: &str = "!roomA:example.org";
4164    const ROOM_B: &str = "!roomB:example.org";
4165    const ROOM_C: &str = "!roomC:example.org";
4166
4167    fn make_room_with_kind(conn: &Connection, room_id: &str, kind: RoomKind) {
4168        create_room(conn, room_id, kind, 1, T0, false, JoinRule::Invite, HistoryVisibility::Shared, None, None).expect("create room");
4169    }
4170
4171    fn member_content_in(conn: &Connection, room_id: &str, mxid: &str) -> serde_json::Value {
4172        let event = current_state_event(conn, room_id, "m.room.member", mxid).expect("query").expect("member event exists");
4173        serde_json::from_str(&event.content).expect("member content is json")
4174    }
4175
4176    fn member_event_count(conn: &Connection, room_id: &str, mxid: &str) -> i64 {
4177        conn.query_row(
4178            "SELECT COUNT(*) FROM events WHERE room_id = ?1 AND event_type = 'm.room.member' AND state_key = ?2",
4179            params![room_id, mxid],
4180            |row| row.get(0),
4181        )
4182        .expect("count member events")
4183    }
4184
4185    #[test]
4186    fn refresh_member_displayname_restamps_every_joined_room_and_skips_a_room_the_user_left() {
4187        let mut conn = test_conn();
4188        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000091", T0).expect("alice");
4189        let bob = ensure_matrix_user(&conn, 2, "bob0000000000000000000000000091", T0).expect("bob");
4190        for room in [ROOM_A, ROOM_B, ROOM_C] {
4191            make_room_with_kind(&conn, room, RoomKind::Group);
4192            apply_state_event(&mut conn, &StateEventWrite { event_id: &new_event_id(), room_id: room, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join","displayname":"old_alice"}"#, origin_server_ts: 1000, now: T0 })
4193                .expect("alice joins");
4194            apply_state_event(&mut conn, &StateEventWrite { event_id: &new_event_id(), room_id: room, sender_user_id: 2, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"join","displayname":"bob_nick"}"#, origin_server_ts: 1100, now: T0 })
4195                .expect("bob joins");
4196        }
4197        apply_state_event(&mut conn, &StateEventWrite { event_id: &new_event_id(), room_id: ROOM_C, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"leave"}"#, origin_server_ts: 1200, now: T0 }).expect("alice leaves C");
4198
4199        let refresh = refresh_member_displayname(&mut conn, 1, "new_alice", T0, 5000).expect("refresh");
4200
4201        assert_eq!(refresh.rooms_updated, 2, "the two rooms alice is still joined in");
4202        for room in [ROOM_A, ROOM_B] {
4203            let content = member_content_in(&conn, room, &alice);
4204            assert_eq!(content["membership"], "join");
4205            assert_eq!(content["displayname"], "new_alice");
4206            assert_eq!(member_event_count(&conn, room, &alice), 2, "a NEW member event lands; the old one stays in history");
4207            let event = current_state_event(&conn, room, "m.room.member", &alice).expect("query").expect("exists");
4208            assert_eq!(event.sender_user_id, 1, "a join refresh is sent by the user themself");
4209            assert_eq!(event.origin_server_ts, 5000);
4210            assert_eq!(member_content_in(&conn, room, &bob)["displayname"], "bob_nick", "another member's event is never touched");
4211        }
4212        assert_eq!(member_content_in(&conn, ROOM_C, &alice), serde_json::json!({ "membership": "leave" }), "the left room gets no new event");
4213        assert_eq!(member_event_count(&conn, ROOM_C, &alice), 2, "join + leave, nothing more");
4214        assert_eq!(refresh.affected_user_ids, HashSet::from([1, 2]), "alice and bob are woken; room C's members are not part of it");
4215
4216        let repeat = refresh_member_displayname(&mut conn, 1, "new_alice", T0, 6000).expect("second pass");
4217        assert_eq!(repeat, DisplaynameRefresh::default(), "an up-to-date event is skipped, so a repeat pass writes nothing");
4218        assert_eq!(member_event_count(&conn, ROOM_A, &alice), 2);
4219    }
4220
4221    #[test]
4222    fn refresh_member_displayname_keeps_an_invite_events_sender_and_is_direct() {
4223        let mut conn = test_conn();
4224        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000092", T0).expect("alice");
4225        let bob = ensure_matrix_user(&conn, 2, "bob0000000000000000000000000092", T0).expect("bob");
4226        make_room_with_kind(&conn, ROOM_A, RoomKind::Dm);
4227        apply_state_event(&mut conn, &StateEventWrite { event_id: &new_event_id(), room_id: ROOM_A, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join","displayname":"alice_nick"}"#, origin_server_ts: 1000, now: T0 })
4228            .expect("alice joins");
4229        apply_state_event(
4230            &mut conn,
4231            &StateEventWrite {
4232                event_id: &new_event_id(),
4233                room_id: ROOM_A,
4234                sender_user_id: 1,
4235                event_type: "m.room.member",
4236                state_key: &bob,
4237                content: r#"{"membership":"invite","is_direct":true}"#,
4238                origin_server_ts: 1100,
4239                now: T0,
4240            },
4241        )
4242        .expect("alice invites bob");
4243
4244        let refresh = refresh_member_displayname(&mut conn, 2, "bob_nick", T0, 5000).expect("refresh");
4245
4246        assert_eq!(refresh.rooms_updated, 1);
4247        let content = member_content_in(&conn, ROOM_A, &bob);
4248        assert_eq!(content, serde_json::json!({ "membership": "invite", "is_direct": true, "displayname": "bob_nick" }));
4249        let event = current_state_event(&conn, ROOM_A, "m.room.member", &bob).expect("query").expect("exists");
4250        assert_eq!(event.sender_user_id, 1, "the inviter stays the sender: stripped invite state reads the inviter off it");
4251        assert_eq!(room_member(&conn, ROOM_A, 2).expect("query").expect("row").membership, Membership::Invite, "still an invitation");
4252        assert_eq!(refresh.affected_user_ids, HashSet::from([1, 2]));
4253    }
4254
4255    #[test]
4256    fn refresh_member_displayname_is_a_noop_without_a_matrix_user_or_a_label() {
4257        let mut conn = test_conn();
4258        let alice = ensure_matrix_user(&conn, 1, "alice00000000000000000000000093", T0).expect("alice");
4259        make_room_with_kind(&conn, ROOM_A, RoomKind::Group);
4260        apply_state_event(&mut conn, &StateEventWrite { event_id: &new_event_id(), room_id: ROOM_A, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("alice joins");
4261
4262        assert_eq!(refresh_member_displayname(&mut conn, 99, "ghost", T0, 5000).expect("unknown user"), DisplaynameRefresh::default());
4263        assert_eq!(refresh_member_displayname(&mut conn, 1, "", T0, 5000).expect("empty label"), DisplaynameRefresh::default());
4264        assert_eq!(member_event_count(&conn, ROOM_A, &alice), 1);
4265        assert_eq!(matrix_user_ids(&conn).expect("ids"), vec![1]);
4266    }
4267
4268    // ---- P16 S-d: adopting a native DM room for a legacy conversation ----
4269
4270    fn native_dm_room(conn: &Connection) {
4271        create_room(conn, ROOM, RoomKind::Dm, 1, T0, true, JoinRule::Invite, HistoryVisibility::Shared, Some("1:2"), None).expect("create dm room");
4272    }
4273
4274    fn legacy_import(legacy_message_id: i64, event_id: &str) -> LegacyDmMessageImport {
4275        LegacyDmMessageImport {
4276            legacy_message_id,
4277            event_id: event_id.to_string(),
4278            sender_user_id: 1,
4279            content: "{}".to_string(),
4280            origin_server_ts: 1000 + legacy_message_id,
4281        }
4282    }
4283
4284    fn adoption_of(legacy_dm_id: i64) -> DmAdoption<'static> {
4285        DmAdoption { room_id: ROOM, legacy_dm_id, key_events: &[], key_events_origin_server_ts: 500, now: T0 }
4286    }
4287
4288    #[test]
4289    fn adopt_dm_room_for_legacy_binds_the_room_and_imports_the_messages() {
4290        let mut conn = test_conn();
4291        ensure_legacy_dm_map_table(&conn);
4292        native_dm_room(&conn);
4293
4294        let counts = adopt_dm_room_for_legacy(&mut conn, adoption_of(7), &[legacy_import(1, "$l1"), legacy_import(2, "$l2")], &[])
4295            .expect("adopt")
4296            .expect("the room is adoptable");
4297        assert_eq!(counts.messages_imported, 2);
4298        assert_eq!(room_by_legacy_dm_id(&conn, 7).expect("query").expect("bound").id, ROOM);
4299        assert_eq!(highest_mapped_legacy_message_id(&conn, ROOM).expect("query"), Some(2));
4300        assert_eq!(get_event(&conn, "$l1").expect("query").expect("imported").room_id, ROOM);
4301    }
4302
4303    #[test]
4304    fn adopt_dm_room_for_legacy_writes_a_key_event_only_where_the_room_has_none() {
4305        let mut conn = test_conn();
4306        ensure_legacy_dm_map_table(&conn);
4307        native_dm_room(&conn);
4308        let existing = apply_state_event(&mut conn, &StateEventWrite { event_id: "$k-existing", room_id: ROOM, sender_user_id: 1, event_type: "org.example.legacy_dm_key", state_key: "@a:example.org", content: r#"{"public_key_b64":"AAAA"}"#, origin_server_ts: 900, now: T0 })
4309            .expect("existing key event");
4310        let key_events = [
4311            NewStateEvent {
4312                event_id: "$k-a".to_string(),
4313                sender_user_id: 1,
4314                event_type: "org.example.legacy_dm_key".to_string(),
4315                state_key: "@a:example.org".to_string(),
4316                content: r#"{"public_key_b64":"BBBB"}"#.to_string(),
4317            },
4318            NewStateEvent {
4319                event_id: "$k-b".to_string(),
4320                sender_user_id: 2,
4321                event_type: "org.example.legacy_dm_key".to_string(),
4322                state_key: "@b:example.org".to_string(),
4323                content: r#"{"public_key_b64":"CCCC"}"#.to_string(),
4324            },
4325        ];
4326        let adoption = DmAdoption { key_events: &key_events, ..adoption_of(7) };
4327        adopt_dm_room_for_legacy(&mut conn, adoption, &[], &[]).expect("adopt").expect("adoptable");
4328
4329        let a = current_state_event(&conn, ROOM, "org.example.legacy_dm_key", "@a:example.org").expect("query").expect("still present");
4330        assert_eq!(a.event_id, existing.event_id, "an existing key event is never overwritten");
4331        let b = current_state_event(&conn, ROOM, "org.example.legacy_dm_key", "@b:example.org").expect("query").expect("added");
4332        assert_eq!(b.event_id, "$k-b");
4333    }
4334
4335    #[test]
4336    fn adopt_dm_room_for_legacy_refuses_a_room_that_is_bound_or_not_a_dm() {
4337        let mut conn = test_conn();
4338        ensure_legacy_dm_map_table(&conn);
4339        make_room(&conn);
4340        assert!(adopt_dm_room_for_legacy(&mut conn, adoption_of(7), &[legacy_import(1, "$l1")], &[]).expect("adopt").is_none(), "a group room is never adopted");
4341        assert!(room_by_legacy_dm_id(&conn, 7).expect("query").is_none());
4342
4343        let mut bound = test_conn();
4344        create_room(&bound, ROOM, RoomKind::Dm, 1, T0, true, JoinRule::Invite, HistoryVisibility::Shared, Some("1:2"), Some(9)).expect("create bound room");
4345        assert!(adopt_dm_room_for_legacy(&mut bound, adoption_of(7), &[legacy_import(1, "$l1")], &[]).expect("adopt").is_none(), "an already-bound room is never re-bound");
4346        assert_eq!(room_by_legacy_dm_id(&bound, 9).expect("query").expect("still bound to 9").id, ROOM);
4347        assert!(get_event(&bound, "$l1").expect("query").is_none(), "a refused adoption writes nothing");
4348    }
4349
4350    #[test]
4351    fn adopt_dm_room_for_legacy_rolls_back_the_binding_when_an_import_fails() {
4352        let mut conn = test_conn();
4353        ensure_legacy_dm_map_table(&conn);
4354        native_dm_room(&conn);
4355
4356        // The second import reuses the first one's event id — refused by the
4357        // UNIQUE index, after the room was already bound inside the transaction.
4358        let result = adopt_dm_room_for_legacy(&mut conn, adoption_of(7), &[legacy_import(1, "$dup"), legacy_import(2, "$dup")], &[]);
4359        assert!(result.is_err());
4360
4361        assert!(room_by_legacy_dm_id(&conn, 7).expect("query").is_none(), "the binding must roll back with the failed import");
4362        assert!(get_event(&conn, "$dup").expect("query").is_none(), "no imported event survives");
4363        assert_eq!(highest_mapped_legacy_message_id(&conn, ROOM).expect("query"), None);
4364    }
4365}