1use 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
28pub 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
65pub fn domain_of(mxid: &str) -> Option<&str> {
69 mxid.strip_prefix('@')?.split_once(':').map(|(_, d)| d)
70}
71
72pub 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
77pub 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
93pub 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
108pub fn room_is_federated(conn: &Connection, room_id: &str) -> rusqlite::Result<bool> {
110 Ok(!remote_domains(conn, room_id)?.is_empty())
111}
112
113fn b64(bytes: &[u8]) -> String {
116 STANDARD_NO_PAD.encode(bytes)
117}
118
119pub 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
156pub 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
172pub 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
179pub 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
188pub 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
218pub 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
234pub 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
253pub enum Mode {
254 Live,
256 Trusted,
258}
259
260#[derive(Debug, Clone, Default)]
262pub struct Ingested {
263 pub event_id: String,
265 pub duplicate: bool,
267 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
325pub 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
372pub 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
387pub 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
399pub 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
410pub 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
421pub 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
460pub 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
476pub 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 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}