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 #[cfg(feature = "f3-hash-ids")]
331 if crate::f3::is_f3_room(conn, &room_id) {
332 return Err(MatrixError::forbidden("this room is a DAG room: events arrive as hashed events"));
333 }
334 let sender = get("sender").ok_or_else(|| MatrixError::bad_json("sender"))?;
335 let ty = get("type").ok_or_else(|| MatrixError::bad_json("type"))?;
336 let ts = pdu.get("origin_server_ts").and_then(Value::as_i64).ok_or_else(|| MatrixError::bad_json("origin_server_ts"))?;
337 if store::get_event(conn, &event_id)?.is_some() {
338 return Ok(Ingested { event_id, duplicate: true, wake: HashSet::new() });
339 }
340 if !is_remote_mxid(&sender) {
341 return Err(MatrixError::forbidden("sender is not a remote user"));
342 }
343 let room = store::get_room(conn, &room_id)?.ok_or_else(|| MatrixError::not_found("unknown room"))?;
344 let sender_uid = ensure_remote_user(conn, &sender, now).map_err(|_| MatrixError::bad_json("sender"))?;
345 let content = if content_ok { pdu.get("content").cloned().unwrap_or_else(|| json!({})) } else { redact_pdu(pdu)["content"].clone() };
346 let state_key = get("state_key");
347 if ty == "m.room.member" {
348 let sk = state_key.as_deref().ok_or_else(|| MatrixError::bad_json("state_key"))?;
349 if is_remote_mxid(sk) {
350 ensure_remote_user(conn, sk, now).map_err(|_| MatrixError::bad_json("state_key"))?;
351 } else if store::user_id_of(conn, sk)?.is_none() {
352 return Err(MatrixError::not_found("unknown local user"));
353 }
354 }
355 if mode == Mode::Live {
356 auth_live(conn, &room, sender_uid, &sender, pdu, &content)?;
357 }
358 let text = content.to_string();
359 match &state_key {
360 Some(sk) => {
361 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 })?;
362 }
363 None if crate::public_channels::is_public_room(conn, &room_id)? => {
364 crate::public_channels::insert_event_deduped(conn, "fed", &event_id, &event_id, &room_id, sender_uid, &ty, &text, ts).map_err(|_| MatrixError::internal())?;
365 }
366 None => {
367 store::insert_timeline_event(conn, &event_id, &room_id, sender_uid, &ty, &text, ts)?;
368 }
369 }
370 let stored = if content_ok { pdu.clone() } else { redact_pdu(pdu) };
371 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()])?;
372 let wake = crate::rooms::member_and_invited_ids(conn, &room_id)?;
373 Ok(Ingested { event_id, duplicate: false, wake })
374}
375
376pub fn store_own_pdu(conn: &mut Connection, pdu: &Value, user_id: i64, mxid: &str, now: &str) -> Result<(), MatrixError> {
378 let get = |k: &str| pdu.get(k).and_then(Value::as_str).map(str::to_string);
379 let event_id = get("event_id").ok_or_else(|| MatrixError::bad_json("event_id"))?;
380 let room_id = get("room_id").ok_or_else(|| MatrixError::bad_json("room_id"))?;
381 let ts = pdu.get("origin_server_ts").and_then(Value::as_i64).unwrap_or(0);
382 if pdu.get("sender").and_then(Value::as_str) != Some(mxid) || pdu.get("state_key").and_then(Value::as_str) != Some(mxid) {
383 return Err(MatrixError::bad_json("own pdu mismatch"));
384 }
385 let content = pdu.get("content").cloned().unwrap_or_else(|| json!({}));
386 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 })?;
387 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()])?;
388 Ok(())
389}
390
391pub fn enqueue(conn: &Connection, destination: &str, kind: &str, room_id: &str, event_id: &str, payload: &Value, now_ms: i64) -> Result<(), FedError> {
395 conn.execute(
396 "INSERT INTO fed_outbox (destination, kind, room_id, event_id, payload, created_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
397 params![destination, kind, room_id, event_id, payload.to_string(), now_ms],
398 )
399 .map_err(db_err)?;
400 Ok(())
401}
402
403pub fn relay_pdu(conn: &Connection, room_id: &str, pdu: &Value, exclude: &[&str], now_ms: i64) -> Result<(), FedError> {
405 let event_id = pdu.get("event_id").and_then(Value::as_str).unwrap_or_default();
406 for d in remote_domains(conn, room_id).map_err(db_err)? {
407 if !exclude.contains(&d.as_str()) {
408 enqueue(conn, &d, "send", room_id, event_id, pdu, now_ms)?;
409 }
410 }
411 Ok(())
412}
413
414pub fn local_state_pdus(conn: &Connection, room_id: &str, local: &str, now_ms: i64) -> Result<Vec<Value>, FedError> {
416 let mut out = Vec::new();
417 for ev in store::current_state_all(conn, room_id).map_err(db_err)? {
418 if let Ok(p) = pdu_for_event(conn, &ev, local, now_ms) {
419 out.push(p);
420 }
421 }
422 Ok(out)
423}
424
425pub fn export_local_events(conn: &Connection, local: &str, now_ms: i64) -> Result<usize, FedError> {
427 let cursor: i64 = conn.query_row("SELECT stream_id FROM fed_export_cursor WHERE id = 1", [], |r| r.get(0)).map_err(db_err)?;
428 let ids: Vec<(i64, String)> = {
429 let mut stmt = conn
430 .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")
431 .map_err(db_err)?;
432 let rows = stmt.query_map(params![cursor], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?))).map_err(db_err)?;
433 rows.collect::<Result<_, _>>().map_err(db_err)?
434 };
435 let mut exported = 0;
436 let mut max_seen = cursor;
437 for (stream, event_id) in ids {
438 max_seen = max_seen.max(stream);
439 let Some(ev) = store::get_event(conn, &event_id).map_err(db_err)? else { continue };
440 let domains = remote_domains(conn, &ev.room_id).map_err(db_err)?;
441 if domains.is_empty() {
442 continue;
443 }
444 let pdu = pdu_for_event(conn, &ev, local, now_ms)?;
445 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") {
446 ev.state_key.as_deref().filter(|s| is_remote_mxid(s)).and_then(domain_of).map(str::to_string)
447 } else {
448 None
449 };
450 for d in &domains {
451 if invitee_domain.as_deref() == Some(d.as_str()) {
452 #[allow(unused_mut)]
453 let mut payload = json!({ "event": pdu, "room_info": room_info(conn, &ev.room_id)?, "state": local_state_pdus(conn, &ev.room_id, local, now_ms)? });
454 #[cfg(feature = "f3-hash-ids")]
455 if crate::f3::is_f3_room(conn, &ev.room_id) {
456 payload["state"] = json!([]);
458 payload["f3"] = crate::f3::snapshot_json(conn, &ev.room_id, &ev.event_id).map_err(|e| FedError::Malformed(e.error))?;
459 }
460 enqueue(conn, d, "invite", &ev.room_id, &ev.event_id, &payload, now_ms)?;
461 } else {
462 enqueue(conn, d, "send", &ev.room_id, &ev.event_id, &pdu, now_ms)?;
463 }
464 }
465 exported += 1;
466 }
467 conn.execute("UPDATE fed_export_cursor SET stream_id = ?1 WHERE id = 1", params![max_seen]).map_err(db_err)?;
468 Ok(exported)
469}
470
471pub fn skeletonize_event(conn: &Connection, event_id: &str) -> rusqlite::Result<()> {
477 conn.execute("UPDATE events SET content = '{}' WHERE event_id = ?1", params![event_id])?;
478 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()? {
479 if let Ok(v) = serde_json::from_str::<Value>(&s) {
480 conn.execute("UPDATE fed_pdus SET pdu = ?2 WHERE event_id = ?1", params![event_id, redact_pdu(&v).to_string()])?;
481 }
482 }
483 conn.execute("INSERT OR IGNORE INTO fed_skeleton (event_id) VALUES (?1)", params![event_id])?;
484 Ok(())
485}
486
487pub fn users_visible_to_origin(conn: &Connection, origin: &str) -> rusqlite::Result<HashSet<i64>> {
492 let mut out = HashSet::new();
493 let mut stmt = conn.prepare(
494 "SELECT DISTINCT t.user_id FROM room_members r
495 JOIN matrix_users ru ON ru.user_id = r.user_id
496 JOIN room_members t ON t.room_id = r.room_id
497 WHERE r.user_id < 0 AND r.membership IN ('join','invite') AND t.user_id > 0 AND t.membership IN ('join','invite')
498 AND (ru.mxid LIKE ?1)",
499 )?;
500 let rows = stmt.query_map(params![format!("@%:{origin}")], |r| r.get::<_, i64>(0))?;
501 for r in rows {
502 out.insert(r?);
503 }
504 Ok(out)
505}
506
507#[cfg(test)]
508mod tests {
509 use super::*;
510
511 fn db() -> Connection {
512 let c = Connection::open_in_memory().unwrap();
513 crate::store::create_matrix_schema(&c).unwrap();
514 crate::keys::create_matrix_keys_schema(&c).unwrap();
515 c
516 }
517
518 #[test]
519 fn redaction_keeps_membership_and_drops_ciphertext() {
520 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}});
521 let r = redact_pdu(&pdu);
522 assert_eq!(r["content"], json!({}));
523 assert!(r.get("unsigned").is_none());
524 let m = redact_pdu(&json!({"type":"m.room.member","content":{"membership":"join","displayname":"x"}}));
525 assert_eq!(m["content"], json!({"membership":"join"}));
526 }
527
528 #[test]
529 fn signed_pdu_verifies_even_after_content_is_dropped() {
530 let c = db();
531 let mut o = Map::new();
532 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"))] {
533 o.insert(k.into(), v);
534 }
535 let pdu = finalize_pdu(&c, o, "a.example", 1).unwrap();
536 let (server, key_id) = pdu_signer(&pdu).unwrap();
537 let (_, key) = active_signing_key(&c, 1).unwrap();
538 let pk = b64(key.verifying_key().as_bytes());
539 assert_eq!(verify_pdu_with_key(&pdu, &server, &key_id, &pk), Ok(true));
540 let mut tampered = pdu.clone();
541 tampered["content"]["ciphertext"] = json!("other");
542 assert_eq!(verify_pdu_with_key(&tampered, &server, &key_id, &pk), Ok(false), "signature holds, hash flags the content");
543 let skeleton = redact_pdu(&pdu);
544 assert_eq!(skeleton["content"], json!({}));
545 assert!(verify_pdu_with_key(&skeleton, &server, &key_id, &pk).is_ok(), "skeleton still verifies the signature");
546 let mut forged = pdu.clone();
547 forged["sender"] = json!("@v:a.example");
548 assert!(verify_pdu_with_key(&forged, &server, &key_id, &pk).is_err());
549 }
550
551 #[test]
552 fn remote_users_get_negative_ids_and_domains_are_tracked() {
553 let c = db();
554 crate::store::ensure_matrix_user(&c, 1, "alice000000000000000000000000a1", "t").unwrap();
555 let id = ensure_remote_user(&c, "@bob:b.example", "t").unwrap();
556 assert!(id < 0);
557 assert_eq!(ensure_remote_user(&c, "@bob:b.example", "t").unwrap(), id);
558 assert_ne!(ensure_remote_user(&c, "@eve:b.example", "t").unwrap(), id);
559 assert!(ensure_remote_user(&c, "@x:example.org", "t").is_err(), "own name is not remote");
560 assert!(is_remote_mxid("@x:b.example:8448"));
561 assert_eq!(domain_of("@x:b.example:8448"), Some("b.example:8448"));
562 }
563
564 #[test]
565 fn retention_leaves_a_verifiable_skeleton_in_federated_rooms_and_deletes_in_local_ones() {
566 let mut c = db();
567 crate::retention::create_retention_schema(&c).unwrap();
568 crate::store::ensure_matrix_user(&c, 1, "alice000000000000000000000000a1", "t").unwrap();
569 let bob = ensure_remote_user(&c, "@bob:b.example", "t").unwrap();
570 for room in ["!fed:example.org", "!loc:example.org"] {
571 c.execute("INSERT INTO rooms (id, kind, creator_user_id, created_at, is_encrypted) VALUES (?1, 'group', 1, 't', 1)", params![room]).unwrap();
572 }
573 c.execute("INSERT INTO room_members (room_id, user_id, membership, updated_at) VALUES ('!fed:example.org', ?1, 'join', 't')", params![bob]).unwrap();
574 let mut ids = Vec::new();
575 for (i, room) in ["!fed:example.org", "!loc:example.org"].iter().enumerate() {
576 let id = format!("$e{i}");
577 store::insert_timeline_event(&mut c, &id, room, 1, "m.room.encrypted", r#"{"ciphertext":"secret"}"#, 1_000).unwrap();
578 ids.push(id);
579 }
580 let ev = store::get_event(&c, &ids[0]).unwrap().unwrap();
581 let pdu = pdu_for_event(&c, &ev, "example.org", 1).unwrap();
583 let policy = crate::retention::RetentionPolicy { ttl_ms: 5_000, ack_grace_ms: 0, keep_last: 0, stale_device_ms: 1 };
584 assert_eq!(crate::retention::purge_delivered_events(&mut c, 1_000_000, &policy).unwrap(), 2);
585 let fed: String = c.query_row("SELECT content FROM events WHERE event_id = '$e0'", [], |r| r.get(0)).unwrap();
586 assert_eq!(fed, "{}", "federated closed room keeps a skeleton row, ciphertext gone");
587 let gone: i64 = c.query_row("SELECT COUNT(*) FROM events WHERE event_id = '$e1'", [], |r| r.get(0)).unwrap();
588 assert_eq!(gone, 0, "local-only room: event deleted as before");
589 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();
590 assert_eq!(stored["content"], json!({}));
591 assert_eq!(stored["signatures"], pdu["signatures"], "signature and hash survive the purge");
592 assert_eq!(stored["hashes"], pdu["hashes"]);
593 let (server, key_id) = pdu_signer(&stored).unwrap();
594 let (_, key) = active_signing_key(&c, 1).unwrap();
595 assert!(verify_pdu_with_key(&stored, &server, &key_id, &b64(key.verifying_key().as_bytes())).is_ok());
596 assert_eq!(crate::retention::purge_delivered_events(&mut c, 1_000_000, &policy).unwrap(), 0, "skeleton is not purged again");
597 }
598}