Skip to main content

mail4agent_server/
fed_rooms.rs

1//! Federation F1/F2 room layer (storage side, no network).
2//!
3//! Profile (documented in the project plan as `m4a-fed-1`):
4//! * every server that has members in a room keeps a full replica of it and
5//!   pushes its own users' events to every other participating server
6//!   (full mesh), so there is no home-server sequencer and no event DAG;
7//! * a PDU carries an explicit `event_id`, a content hash and the origin's
8//!   signature over the *redacted* form, so the stored PDU stays verifiable
9//!   after the ciphertext is dropped (skeleton);
10//! * remote users are rows in `matrix_users` with negative ids and their real
11//!   mxids, so membership, sync, receipts and key visibility work unchanged;
12//! * authorization is a minimal membership/power-level check against the
13//!   local replica; there is no state resolution (concurrent conflicting
14//!   state changes are last-write-wins per server).
15
16use std::collections::{BTreeSet, HashSet};
17
18use base64::engine::general_purpose::STANDARD_NO_PAD;
19use base64::Engine;
20use rusqlite::{params, Connection, OptionalExtension};
21use serde_json::{json, Map, Value};
22use sha2::{Digest, Sha256};
23
24use crate::error::MatrixError;
25use crate::federation::{active_signing_key, canonical_json, parse_server_name, sign_json, FedError};
26use crate::store::{self, JoinRule, MatrixEvent, Membership, PowerAction};
27
28/// Tables owned by this layer.
29pub fn create_fed_schema(conn: &Connection) -> rusqlite::Result<()> {
30    conn.execute_batch(
31        r#"
32        CREATE TABLE IF NOT EXISTS fed_pdus (
33            event_id TEXT PRIMARY KEY,
34            room_id  TEXT NOT NULL,
35            pdu      TEXT NOT NULL
36        );
37        CREATE TABLE IF NOT EXISTS fed_skeleton (
38            event_id TEXT PRIMARY KEY
39        );
40        CREATE TABLE IF NOT EXISTS fed_outbox (
41            id          INTEGER PRIMARY KEY AUTOINCREMENT,
42            destination TEXT NOT NULL,
43            kind        TEXT NOT NULL,
44            room_id     TEXT NOT NULL DEFAULT '',
45            event_id    TEXT NOT NULL DEFAULT '',
46            payload     TEXT NOT NULL,
47            created_ms  INTEGER NOT NULL,
48            attempts    INTEGER NOT NULL DEFAULT 0,
49            next_try_ms INTEGER NOT NULL DEFAULT 0
50        );
51        CREATE INDEX IF NOT EXISTS idx_fed_outbox_dest ON fed_outbox(destination, id);
52        CREATE TABLE IF NOT EXISTS fed_export_cursor (
53            id        INTEGER PRIMARY KEY CHECK (id = 1),
54            stream_id INTEGER NOT NULL
55        );
56        INSERT OR IGNORE INTO fed_export_cursor (id, stream_id) VALUES (1, 0);
57        "#,
58    )
59}
60
61fn db_err(e: rusqlite::Error) -> FedError {
62    FedError::Db(e.to_string())
63}
64
65// ------------------------------------------------------------ remote users
66
67/// Domain part of an mxid (`@u:d` gives `d`, ports included).
68pub fn domain_of(mxid: &str) -> Option<&str> {
69    mxid.strip_prefix('@')?.split_once(':').map(|(_, d)| d)
70}
71
72/// An mxid addressed to another server.
73pub fn is_remote_mxid(mxid: &str) -> bool {
74    domain_of(mxid).is_some_and(|d| parse_server_name(d).is_some() && !store::is_local_server_name(d))
75}
76
77/// Proxy row for a remote user (negative id, real mxid). Idempotent.
78pub fn ensure_remote_user(conn: &Connection, mxid: &str, now: &str) -> Result<i64, FedError> {
79    if !is_remote_mxid(mxid) {
80        return Err(FedError::Malformed(format!("not a remote mxid: {mxid}")));
81    }
82    if let Some(id) = store::user_id_of(conn, mxid).map_err(db_err)? {
83        return Ok(id);
84    }
85    let next: i64 = conn
86        .query_row("SELECT COALESCE(MIN(user_id), 0) FROM matrix_users WHERE user_id < 0", [], |r| r.get::<_, i64>(0))
87        .map_err(db_err)?
88        - 1;
89    conn.execute("INSERT INTO matrix_users (user_id, mxid, created_at) VALUES (?1, ?2, ?3)", params![next, mxid, now]).map_err(db_err)?;
90    Ok(next)
91}
92
93/// Domains of remote users that have any membership row in the room.
94pub fn remote_domains(conn: &Connection, room_id: &str) -> rusqlite::Result<BTreeSet<String>> {
95    let mut stmt = conn.prepare(
96        "SELECT u.mxid FROM room_members m JOIN matrix_users u ON u.user_id = m.user_id WHERE m.room_id = ?1 AND m.user_id < 0",
97    )?;
98    let rows = stmt.query_map(params![room_id], |r| r.get::<_, String>(0))?;
99    let mut out = BTreeSet::new();
100    for r in rows {
101        if let Some(d) = domain_of(&r?) {
102            out.insert(d.to_string());
103        }
104    }
105    Ok(out)
106}
107
108/// Whether the room has remote members (federated closed rooms get skeleton redaction).
109pub fn room_is_federated(conn: &Connection, room_id: &str) -> rusqlite::Result<bool> {
110    Ok(!remote_domains(conn, room_id)?.is_empty())
111}
112
113// ---------------------------------------------------------------- PDU form
114
115fn b64(bytes: &[u8]) -> String {
116    STANDARD_NO_PAD.encode(bytes)
117}
118
119/// The signed/hash-stable subset of a PDU (Matrix redaction rules, room v11 shape).
120pub fn redact_pdu(pdu: &Value) -> Value {
121    const KEEP: [&str; 11] = ["event_id", "type", "room_id", "sender", "state_key", "hashes", "signatures", "depth", "origin", "origin_server_ts", "redacts"];
122    let mut out = Map::new();
123    for k in KEEP {
124        if let Some(v) = pdu.get(k) {
125            out.insert(k.to_string(), v.clone());
126        }
127    }
128    let keys: &[&str] = match pdu.get("type").and_then(Value::as_str).unwrap_or("") {
129        "m.room.member" => &["membership", "join_authorised_via_users_server"],
130        "m.room.join_rules" => &["join_rule", "allow"],
131        "m.room.power_levels" => &["ban", "events", "events_default", "invite", "kick", "redact", "state_default", "users", "users_default"],
132        "m.room.history_visibility" => &["history_visibility"],
133        "m.room.create" => &["creator", "room_version", "type", "m.federate"],
134        _ => &[],
135    };
136    let mut content = Map::new();
137    if let Some(c) = pdu.get("content").and_then(Value::as_object) {
138        for k in keys {
139            if let Some(v) = c.get(*k) {
140                content.insert((*k).to_string(), v.clone());
141            }
142        }
143    }
144    out.insert("content".into(), Value::Object(content));
145    Value::Object(out)
146}
147
148fn content_hash(pdu: &Value) -> String {
149    let mut bare = pdu.as_object().cloned().unwrap_or_default();
150    for k in ["hashes", "signatures", "unsigned"] {
151        bare.remove(k);
152    }
153    b64(&Sha256::digest(canonical_json(&Value::Object(bare))))
154}
155
156/// Add `hashes` and the origin's signature (over the redacted form).
157pub fn finalize_pdu(conn: &Connection, mut pdu: Map<String, Value>, origin: &str, now_ms: i64) -> Result<Value, FedError> {
158    for k in ["hashes", "signatures", "unsigned"] {
159        pdu.remove(k);
160    }
161    let h = content_hash(&Value::Object(pdu.clone()));
162    pdu.insert("hashes".into(), json!({ "sha256": h }));
163    let (key_id, key) = active_signing_key(conn, now_ms)?;
164    let mut red = redact_pdu(&Value::Object(pdu.clone())).as_object().cloned().unwrap_or_default();
165    sign_json(&mut red, origin, &key_id, &key);
166    if let Some(sigs) = red.remove("signatures") {
167        pdu.insert("signatures".into(), sigs);
168    }
169    Ok(Value::Object(pdu))
170}
171
172/// `(signing server, key id)` named by a PDU: the sender's server, first key listed.
173pub fn pdu_signer(pdu: &Value) -> Option<(String, String)> {
174    let server = domain_of(pdu.get("sender")?.as_str()?)?.to_string();
175    let key = pdu.get("signatures")?.get(&server)?.as_object()?.keys().next()?.clone();
176    Some((server, key))
177}
178
179/// Check a PDU against the signer's public key. `Ok(true)`: signature and
180/// content hash good; `Ok(false)`: signature good but content does not match
181/// the hash (caller must drop the content, i.e. keep only the redacted form).
182pub fn verify_pdu_with_key(pdu: &Value, signer: &str, key_id: &str, public_key_b64: &str) -> Result<bool, FedError> {
183    crate::federation::verify_json(&redact_pdu(pdu), signer, key_id, public_key_b64)?;
184    let want = pdu.get("hashes").and_then(|h| h.get("sha256")).and_then(Value::as_str).ok_or_else(|| FedError::Malformed("hashes".into()))?;
185    Ok(want == content_hash(pdu))
186}
187
188/// Signed PDU for a locally created event; stored so it is served unchanged later.
189pub fn pdu_for_event(conn: &Connection, ev: &MatrixEvent, local: &str, now_ms: i64) -> Result<Value, FedError> {
190    if let Some(s) = conn.query_row("SELECT pdu FROM fed_pdus WHERE event_id = ?1", params![ev.event_id], |r| r.get::<_, String>(0)).optional().map_err(db_err)? {
191        return serde_json::from_str(&s).map_err(|_| FedError::Malformed("stored pdu".into()));
192    }
193    let sender = store::mxid_of(conn, ev.sender_user_id).map_err(db_err)?.ok_or_else(|| FedError::Malformed("sender".into()))?;
194    if is_remote_mxid(&sender) {
195        return Err(FedError::Malformed("event of a remote sender has no stored pdu".into()));
196    }
197    let content: Value = serde_json::from_str(&ev.content).map_err(|_| FedError::Malformed("content".into()))?;
198    let mut o = Map::new();
199    o.insert("event_id".into(), json!(ev.event_id));
200    o.insert("room_id".into(), json!(ev.room_id));
201    o.insert("sender".into(), json!(sender));
202    o.insert("type".into(), json!(ev.event_type));
203    if let Some(sk) = &ev.state_key {
204        o.insert("state_key".into(), json!(sk));
205    }
206    o.insert("content".into(), content);
207    o.insert("origin_server_ts".into(), json!(ev.origin_server_ts));
208    o.insert("origin".into(), json!(local));
209    o.insert("depth".into(), json!(ev.stream_id));
210    if let Some(r) = &ev.redacts {
211        o.insert("redacts".into(), json!(r));
212    }
213    let pdu = finalize_pdu(conn, o, local, now_ms)?;
214    conn.execute("INSERT OR IGNORE INTO fed_pdus (event_id, room_id, pdu) VALUES (?1, ?2, ?3)", params![ev.event_id, ev.room_id, pdu.to_string()]).map_err(db_err)?;
215    Ok(pdu)
216}
217
218// ------------------------------------------------------------- room info
219
220/// Room metadata a remote replica needs (the part that is not event content).
221pub fn room_info(conn: &Connection, room_id: &str) -> Result<Value, FedError> {
222    let room = store::get_room(conn, room_id).map_err(db_err)?.ok_or_else(|| FedError::Malformed("unknown room".into()))?;
223    let creator = store::mxid_of(conn, room.creator_user_id).map_err(db_err)?.unwrap_or_default();
224    Ok(json!({
225        "kind": room.kind.as_str(),
226        "is_encrypted": room.is_encrypted,
227        "join_rule": room.join_rule.as_str(),
228        "history_visibility": room.history_visibility.as_str(),
229        "room_version": room.room_version,
230        "creator": creator,
231    }))
232}
233
234/// Create the local replica row for a remote room. No-op if present.
235pub fn create_replica_room(conn: &Connection, room_id: &str, info: &Value, now: &str) -> Result<(), MatrixError> {
236    if store::get_room(conn, room_id)?.is_some() {
237        return Ok(());
238    }
239    let kind = info.get("kind").and_then(Value::as_str).and_then(store::RoomKind::from_wire_name).ok_or_else(|| MatrixError::bad_json("room kind"))?;
240    let join_rule = info.get("join_rule").and_then(Value::as_str).and_then(JoinRule::from_wire_name).unwrap_or(JoinRule::Invite);
241    let hv = info.get("history_visibility").and_then(Value::as_str).and_then(store::HistoryVisibility::from_wire_name).unwrap_or(store::HistoryVisibility::Shared);
242    let creator = info.get("creator").and_then(Value::as_str).unwrap_or_default();
243    let creator_id = if is_remote_mxid(creator) { ensure_remote_user(conn, creator, now).map_err(|_| MatrixError::bad_json("creator"))? } else { store::user_id_of(conn, creator)?.unwrap_or(-1) };
244    let encrypted = info.get("is_encrypted").and_then(Value::as_bool).unwrap_or(false);
245    store::create_room(conn, room_id, kind, creator_id, now, encrypted, join_rule, hv, None, None)?;
246    Ok(())
247}
248
249// ---------------------------------------------------------------- ingest
250
251/// How strictly an incoming PDU is authorized.
252#[derive(Debug, Clone, Copy, PartialEq, Eq)]
253pub enum Mode {
254    /// Live delivery: membership and power checks against the local replica.
255    Live,
256    /// State/history inside a signed join or invite response: signature only.
257    Trusted,
258}
259
260/// Result of one ingest.
261#[derive(Debug, Clone, Default)]
262pub struct Ingested {
263    /// Event id.
264    pub event_id: String,
265    /// Already stored, nothing done.
266    pub duplicate: bool,
267    /// Local users to wake.
268    pub wake: HashSet<i64>,
269}
270
271fn auth_live(conn: &Connection, room: &store::Room, sender_uid: i64, sender: &str, pdu: &Value, content: &Value) -> Result<(), MatrixError> {
272    let ty = pdu.get("type").and_then(Value::as_str).unwrap_or("");
273    let current = store::room_member(conn, &room.id, sender_uid)?.map(|m| m.membership);
274    let pl = crate::rooms::power_levels_of(conn, &room.id)?;
275    if ty == "m.room.member" {
276        let sk = pdu.get("state_key").and_then(Value::as_str).ok_or_else(|| MatrixError::bad_json("state_key"))?;
277        let target_uid = store::user_id_of(conn, sk)?;
278        let target_cur = match target_uid {
279            Some(u) => store::room_member(conn, &room.id, u)?.map(|m| m.membership),
280            None => None,
281        };
282        return match content.get("membership").and_then(Value::as_str).unwrap_or("") {
283            "join" => {
284                if sk != sender {
285                    return Err(MatrixError::forbidden("join for another user"));
286                }
287                match current {
288                    Some(Membership::Ban) => Err(MatrixError::forbidden("banned")),
289                    Some(Membership::Join) | Some(Membership::Invite) => Ok(()),
290                    _ if room.join_rule == JoinRule::Public => Ok(()),
291                    _ => Err(MatrixError::forbidden("no invitation")),
292                }
293            }
294            "invite" => {
295                crate::rooms::require_member(current)?;
296                crate::rooms::require_power(&pl, sender, PowerAction::Invite)?;
297                match target_cur {
298                    Some(Membership::Join) | Some(Membership::Ban) => Err(MatrixError::forbidden("target cannot be invited")),
299                    _ => Ok(()),
300                }
301            }
302            "leave" if sk == sender => match current {
303                Some(Membership::Join) | Some(Membership::Invite) => Ok(()),
304                _ => Err(MatrixError::forbidden("not a member")),
305            },
306            "leave" => {
307                crate::rooms::require_member(current)?;
308                crate::rooms::require_power(&pl, sender, PowerAction::Kick)
309            }
310            "ban" => {
311                crate::rooms::require_member(current)?;
312                crate::rooms::require_power(&pl, sender, PowerAction::Ban)
313            }
314            _ => Err(MatrixError::forbidden("unsupported membership")),
315        };
316    }
317    crate::rooms::require_member(current)?;
318    let is_state = pdu.get("state_key").is_some();
319    if store::user_level(&pl, sender) < store::event_level(&pl, ty, is_state) {
320        return Err(MatrixError::forbidden("insufficient power level"));
321    }
322    Ok(())
323}
324
325/// Store one verified PDU into the local replica.
326pub fn ingest_pdu(conn: &mut Connection, pdu: &Value, content_ok: bool, mode: Mode, now: &str) -> Result<Ingested, MatrixError> {
327    let get = |k: &str| pdu.get(k).and_then(Value::as_str).map(str::to_string);
328    let event_id = get("event_id").ok_or_else(|| MatrixError::bad_json("event_id"))?;
329    let room_id = get("room_id").ok_or_else(|| MatrixError::bad_json("room_id"))?;
330    let sender = get("sender").ok_or_else(|| MatrixError::bad_json("sender"))?;
331    let ty = get("type").ok_or_else(|| MatrixError::bad_json("type"))?;
332    let ts = pdu.get("origin_server_ts").and_then(Value::as_i64).ok_or_else(|| MatrixError::bad_json("origin_server_ts"))?;
333    if store::get_event(conn, &event_id)?.is_some() {
334        return Ok(Ingested { event_id, duplicate: true, wake: HashSet::new() });
335    }
336    if !is_remote_mxid(&sender) {
337        return Err(MatrixError::forbidden("sender is not a remote user"));
338    }
339    let room = store::get_room(conn, &room_id)?.ok_or_else(|| MatrixError::not_found("unknown room"))?;
340    let sender_uid = ensure_remote_user(conn, &sender, now).map_err(|_| MatrixError::bad_json("sender"))?;
341    let content = if content_ok { pdu.get("content").cloned().unwrap_or_else(|| json!({})) } else { redact_pdu(pdu)["content"].clone() };
342    let state_key = get("state_key");
343    if ty == "m.room.member" {
344        let sk = state_key.as_deref().ok_or_else(|| MatrixError::bad_json("state_key"))?;
345        if is_remote_mxid(sk) {
346            ensure_remote_user(conn, sk, now).map_err(|_| MatrixError::bad_json("state_key"))?;
347        } else if store::user_id_of(conn, sk)?.is_none() {
348            return Err(MatrixError::not_found("unknown local user"));
349        }
350    }
351    if mode == Mode::Live {
352        auth_live(conn, &room, sender_uid, &sender, pdu, &content)?;
353    }
354    let text = content.to_string();
355    match &state_key {
356        Some(sk) => {
357            store::apply_state_event(conn, &store::StateEventWrite { event_id: &event_id, room_id: &room_id, sender_user_id: sender_uid, event_type: &ty, state_key: sk, content: &text, origin_server_ts: ts, now })?;
358        }
359        None if crate::public_channels::is_public_room(conn, &room_id)? => {
360            crate::public_channels::insert_event_deduped(conn, "fed", &event_id, &event_id, &room_id, sender_uid, &ty, &text, ts).map_err(|_| MatrixError::internal())?;
361        }
362        None => {
363            store::insert_timeline_event(conn, &event_id, &room_id, sender_uid, &ty, &text, ts)?;
364        }
365    }
366    let stored = if content_ok { pdu.clone() } else { redact_pdu(pdu) };
367    conn.execute("INSERT OR REPLACE INTO fed_pdus (event_id, room_id, pdu) VALUES (?1, ?2, ?3)", params![event_id, room_id, stored.to_string()])?;
368    let wake = crate::rooms::member_and_invited_ids(conn, &room_id)?;
369    Ok(Ingested { event_id, duplicate: false, wake })
370}
371
372/// Store a PDU that a local user created (their own join in a remote room), keeping its id.
373pub fn store_own_pdu(conn: &mut Connection, pdu: &Value, user_id: i64, mxid: &str, now: &str) -> Result<(), MatrixError> {
374    let get = |k: &str| pdu.get(k).and_then(Value::as_str).map(str::to_string);
375    let event_id = get("event_id").ok_or_else(|| MatrixError::bad_json("event_id"))?;
376    let room_id = get("room_id").ok_or_else(|| MatrixError::bad_json("room_id"))?;
377    let ts = pdu.get("origin_server_ts").and_then(Value::as_i64).unwrap_or(0);
378    if pdu.get("sender").and_then(Value::as_str) != Some(mxid) || pdu.get("state_key").and_then(Value::as_str) != Some(mxid) {
379        return Err(MatrixError::bad_json("own pdu mismatch"));
380    }
381    let content = pdu.get("content").cloned().unwrap_or_else(|| json!({}));
382    store::apply_state_event(conn, &store::StateEventWrite { event_id: &event_id, room_id: &room_id, sender_user_id: user_id, event_type: "m.room.member", state_key: mxid, content: &content.to_string(), origin_server_ts: ts, now })?;
383    conn.execute("INSERT OR REPLACE INTO fed_pdus (event_id, room_id, pdu) VALUES (?1, ?2, ?3)", params![event_id, room_id, pdu.to_string()])?;
384    Ok(())
385}
386
387// ---------------------------------------------------------------- outbox
388
389/// Queue an item for a destination.
390pub fn enqueue(conn: &Connection, destination: &str, kind: &str, room_id: &str, event_id: &str, payload: &Value, now_ms: i64) -> Result<(), FedError> {
391    conn.execute(
392        "INSERT INTO fed_outbox (destination, kind, room_id, event_id, payload, created_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
393        params![destination, kind, room_id, event_id, payload.to_string(), now_ms],
394    )
395    .map_err(db_err)?;
396    Ok(())
397}
398
399/// Queue `pdu` as a normal transaction item for every remote server of the room except `exclude`.
400pub fn relay_pdu(conn: &Connection, room_id: &str, pdu: &Value, exclude: &[&str], now_ms: i64) -> Result<(), FedError> {
401    let event_id = pdu.get("event_id").and_then(Value::as_str).unwrap_or_default();
402    for d in remote_domains(conn, room_id).map_err(db_err)? {
403        if !exclude.contains(&d.as_str()) {
404            enqueue(conn, &d, "send", room_id, event_id, pdu, now_ms)?;
405        }
406    }
407    Ok(())
408}
409
410/// State PDUs of the room that this server can vouch for (its own senders').
411pub fn local_state_pdus(conn: &Connection, room_id: &str, local: &str, now_ms: i64) -> Result<Vec<Value>, FedError> {
412    let mut out = Vec::new();
413    for ev in store::current_state_all(conn, room_id).map_err(db_err)? {
414        if let Ok(p) = pdu_for_event(conn, &ev, local, now_ms) {
415            out.push(p);
416        }
417    }
418    Ok(out)
419}
420
421/// Turn new local-sender events in federated rooms into outbox items. Returns how many events were exported.
422pub fn export_local_events(conn: &Connection, local: &str, now_ms: i64) -> Result<usize, FedError> {
423    let cursor: i64 = conn.query_row("SELECT stream_id FROM fed_export_cursor WHERE id = 1", [], |r| r.get(0)).map_err(db_err)?;
424    let ids: Vec<(i64, String)> = {
425        let mut stmt = conn
426            .prepare("SELECT stream_id, event_id FROM (SELECT stream_id, event_id, sender_user_id FROM events UNION ALL SELECT stream_id, event_id, sender_user_id FROM pub_events) WHERE stream_id > ?1 AND sender_user_id > 0 ORDER BY stream_id LIMIT 500")
427            .map_err(db_err)?;
428        let rows = stmt.query_map(params![cursor], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?))).map_err(db_err)?;
429        rows.collect::<Result<_, _>>().map_err(db_err)?
430    };
431    let mut exported = 0;
432    let mut max_seen = cursor;
433    for (stream, event_id) in ids {
434        max_seen = max_seen.max(stream);
435        let Some(ev) = store::get_event(conn, &event_id).map_err(db_err)? else { continue };
436        let domains = remote_domains(conn, &ev.room_id).map_err(db_err)?;
437        if domains.is_empty() {
438            continue;
439        }
440        let pdu = pdu_for_event(conn, &ev, local, now_ms)?;
441        let invitee_domain = if ev.event_type == "m.room.member" && serde_json::from_str::<Value>(&ev.content).ok().and_then(|c| c.get("membership").and_then(Value::as_str).map(str::to_string)).as_deref() == Some("invite") {
442            ev.state_key.as_deref().filter(|s| is_remote_mxid(s)).and_then(domain_of).map(str::to_string)
443        } else {
444            None
445        };
446        for d in &domains {
447            if invitee_domain.as_deref() == Some(d.as_str()) {
448                let payload = json!({ "event": pdu, "room_info": room_info(conn, &ev.room_id)?, "state": local_state_pdus(conn, &ev.room_id, local, now_ms)? });
449                enqueue(conn, d, "invite", &ev.room_id, &ev.event_id, &payload, now_ms)?;
450            } else {
451                enqueue(conn, d, "send", &ev.room_id, &ev.event_id, &pdu, now_ms)?;
452            }
453        }
454        exported += 1;
455    }
456    conn.execute("UPDATE fed_export_cursor SET stream_id = ?1 WHERE id = 1", params![max_seen]).map_err(db_err)?;
457    Ok(exported)
458}
459
460// -------------------------------------------------------------- skeleton
461
462/// Drop the content of an event but keep its place in the room: the row stays
463/// with empty content and the stored PDU becomes its redacted form (which
464/// still verifies). Used by retention for federated closed rooms.
465pub fn skeletonize_event(conn: &Connection, event_id: &str) -> rusqlite::Result<()> {
466    conn.execute("UPDATE events SET content = '{}' WHERE event_id = ?1", params![event_id])?;
467    if let Some(s) = conn.query_row("SELECT pdu FROM fed_pdus WHERE event_id = ?1", params![event_id], |r| r.get::<_, String>(0)).optional()? {
468        if let Ok(v) = serde_json::from_str::<Value>(&s) {
469            conn.execute("UPDATE fed_pdus SET pdu = ?2 WHERE event_id = ?1", params![event_id, redact_pdu(&v).to_string()])?;
470        }
471    }
472    conn.execute("INSERT OR IGNORE INTO fed_skeleton (event_id) VALUES (?1)", params![event_id])?;
473    Ok(())
474}
475
476// ------------------------------------------------------------ key serving
477
478/// Local users whose keys `origin` may see: members/invitees of rooms that
479/// also have a member on `origin`.
480pub fn users_visible_to_origin(conn: &Connection, origin: &str) -> rusqlite::Result<HashSet<i64>> {
481    let mut out = HashSet::new();
482    let mut stmt = conn.prepare(
483        "SELECT DISTINCT t.user_id FROM room_members r
484           JOIN matrix_users ru ON ru.user_id = r.user_id
485           JOIN room_members t ON t.room_id = r.room_id
486          WHERE r.user_id < 0 AND r.membership IN ('join','invite') AND t.user_id > 0 AND t.membership IN ('join','invite')
487            AND (ru.mxid LIKE ?1)",
488    )?;
489    let rows = stmt.query_map(params![format!("@%:{origin}")], |r| r.get::<_, i64>(0))?;
490    for r in rows {
491        out.insert(r?);
492    }
493    Ok(out)
494}
495
496#[cfg(test)]
497mod tests {
498    use super::*;
499
500    fn db() -> Connection {
501        let c = Connection::open_in_memory().unwrap();
502        crate::store::create_matrix_schema(&c).unwrap();
503        crate::keys::create_matrix_keys_schema(&c).unwrap();
504        c
505    }
506
507    #[test]
508    fn redaction_keeps_membership_and_drops_ciphertext() {
509        let pdu = json!({"type":"m.room.encrypted","content":{"ciphertext":"x"},"sender":"@a:b.example","event_id":"$1","room_id":"!r:b.example","origin_server_ts":1,"hashes":{"sha256":"h"},"unsigned":{"a":1}});
510        let r = redact_pdu(&pdu);
511        assert_eq!(r["content"], json!({}));
512        assert!(r.get("unsigned").is_none());
513        let m = redact_pdu(&json!({"type":"m.room.member","content":{"membership":"join","displayname":"x"}}));
514        assert_eq!(m["content"], json!({"membership":"join"}));
515    }
516
517    #[test]
518    fn signed_pdu_verifies_even_after_content_is_dropped() {
519        let c = db();
520        let mut o = Map::new();
521        for (k, v) in [("event_id", json!("$1")), ("room_id", json!("!r:a.example")), ("sender", json!("@u:a.example")), ("type", json!("m.room.encrypted")), ("content", json!({"ciphertext":"secret"})), ("origin_server_ts", json!(5)), ("origin", json!("a.example"))] {
522            o.insert(k.into(), v);
523        }
524        let pdu = finalize_pdu(&c, o, "a.example", 1).unwrap();
525        let (server, key_id) = pdu_signer(&pdu).unwrap();
526        let (_, key) = active_signing_key(&c, 1).unwrap();
527        let pk = b64(key.verifying_key().as_bytes());
528        assert_eq!(verify_pdu_with_key(&pdu, &server, &key_id, &pk), Ok(true));
529        let mut tampered = pdu.clone();
530        tampered["content"]["ciphertext"] = json!("other");
531        assert_eq!(verify_pdu_with_key(&tampered, &server, &key_id, &pk), Ok(false), "signature holds, hash flags the content");
532        let skeleton = redact_pdu(&pdu);
533        assert_eq!(skeleton["content"], json!({}));
534        assert!(verify_pdu_with_key(&skeleton, &server, &key_id, &pk).is_ok(), "skeleton still verifies the signature");
535        let mut forged = pdu.clone();
536        forged["sender"] = json!("@v:a.example");
537        assert!(verify_pdu_with_key(&forged, &server, &key_id, &pk).is_err());
538    }
539
540    #[test]
541    fn remote_users_get_negative_ids_and_domains_are_tracked() {
542        let c = db();
543        crate::store::ensure_matrix_user(&c, 1, "alice000000000000000000000000a1", "t").unwrap();
544        let id = ensure_remote_user(&c, "@bob:b.example", "t").unwrap();
545        assert!(id < 0);
546        assert_eq!(ensure_remote_user(&c, "@bob:b.example", "t").unwrap(), id);
547        assert_ne!(ensure_remote_user(&c, "@eve:b.example", "t").unwrap(), id);
548        assert!(ensure_remote_user(&c, "@x:example.org", "t").is_err(), "own name is not remote");
549        assert!(is_remote_mxid("@x:b.example:8448"));
550        assert_eq!(domain_of("@x:b.example:8448"), Some("b.example:8448"));
551    }
552
553    #[test]
554    fn retention_leaves_a_verifiable_skeleton_in_federated_rooms_and_deletes_in_local_ones() {
555        let mut c = db();
556        crate::retention::create_retention_schema(&c).unwrap();
557        crate::store::ensure_matrix_user(&c, 1, "alice000000000000000000000000a1", "t").unwrap();
558        let bob = ensure_remote_user(&c, "@bob:b.example", "t").unwrap();
559        for room in ["!fed:example.org", "!loc:example.org"] {
560            c.execute("INSERT INTO rooms (id, kind, creator_user_id, created_at, is_encrypted) VALUES (?1, 'group', 1, 't', 1)", params![room]).unwrap();
561        }
562        c.execute("INSERT INTO room_members (room_id, user_id, membership, updated_at) VALUES ('!fed:example.org', ?1, 'join', 't')", params![bob]).unwrap();
563        let mut ids = Vec::new();
564        for (i, room) in ["!fed:example.org", "!loc:example.org"].iter().enumerate() {
565            let id = format!("$e{i}");
566            store::insert_timeline_event(&mut c, &id, room, 1, "m.room.encrypted", r#"{"ciphertext":"secret"}"#, 1_000).unwrap();
567            ids.push(id);
568        }
569        let ev = store::get_event(&c, &ids[0]).unwrap().unwrap();
570        // the local server name is the process default here; sign as a stand-in origin
571        let pdu = pdu_for_event(&c, &ev, "example.org", 1).unwrap();
572        let policy = crate::retention::RetentionPolicy { ttl_ms: 5_000, ack_grace_ms: 0, keep_last: 0, stale_device_ms: 1 };
573        assert_eq!(crate::retention::purge_delivered_events(&mut c, 1_000_000, &policy).unwrap(), 2);
574        let fed: String = c.query_row("SELECT content FROM events WHERE event_id = '$e0'", [], |r| r.get(0)).unwrap();
575        assert_eq!(fed, "{}", "federated closed room keeps a skeleton row, ciphertext gone");
576        let gone: i64 = c.query_row("SELECT COUNT(*) FROM events WHERE event_id = '$e1'", [], |r| r.get(0)).unwrap();
577        assert_eq!(gone, 0, "local-only room: event deleted as before");
578        let stored: Value = serde_json::from_str(&c.query_row("SELECT pdu FROM fed_pdus WHERE event_id = '$e0'", [], |r| r.get::<_, String>(0)).unwrap()).unwrap();
579        assert_eq!(stored["content"], json!({}));
580        assert_eq!(stored["signatures"], pdu["signatures"], "signature and hash survive the purge");
581        assert_eq!(stored["hashes"], pdu["hashes"]);
582        let (server, key_id) = pdu_signer(&stored).unwrap();
583        let (_, key) = active_signing_key(&c, 1).unwrap();
584        assert!(verify_pdu_with_key(&stored, &server, &key_id, &b64(key.verifying_key().as_bytes())).is_ok());
585        assert_eq!(crate::retention::purge_delivered_events(&mut c, 1_000_000, &policy).unwrap(), 0, "skeleton is not purged again");
586    }
587}