Skip to main content

mail4agent_server/
fed_edus.rs

1//! Ephemeral data units between servers: typing, read receipts, device-list updates. (Presence is
2//! not carried: this server has none, and an incoming `m.presence` is accepted and dropped.)
3//! Outgoing ones go through the federation outbox like any other item; incoming ones are applied
4//! only for users of the sending server who are members of a room here.
5
6use std::collections::HashSet;
7use std::time::Instant;
8
9use rusqlite::Connection;
10use serde_json::{json, Value};
11
12use crate::fed_rooms as fr;
13use crate::store::{self, Membership};
14
15fn now_ms() -> i64 {
16    crate::federation::now_ms()
17}
18
19fn remote_servers_of_user(conn: &Connection, user_id: i64) -> Vec<String> {
20    let mut out = std::collections::BTreeSet::new();
21    for room in store::rooms_for_user(conn, user_id, Some(Membership::Join)).unwrap_or_default() {
22        out.extend(fr::remote_domains(conn, &room).unwrap_or_default());
23    }
24    out.into_iter().collect()
25}
26
27fn queue(conn: &Connection, servers: impl IntoIterator<Item = String>, edu: &Value) {
28    for d in servers {
29        let _ = fr::enqueue(conn, &d, "edu", "", "", edu, now_ms());
30    }
31}
32
33/// A local user's typing state in a room that has remote members.
34pub fn enqueue_typing(conn: &Connection, room_id: &str, user_id: i64, typing: bool) {
35    if user_id <= 0 {
36        return;
37    }
38    let (Ok(servers), Ok(Some(mxid))) = (fr::remote_domains(conn, room_id), store::mxid_of(conn, user_id)) else { return };
39    queue(conn, servers, &json!({ "edu_type": "m.typing", "content": { "room_id": room_id, "user_id": mxid, "typing": typing } }));
40}
41
42/// A local user's read receipt in a room that has remote members.
43pub fn enqueue_receipt(conn: &Connection, room_id: &str, user_id: i64, event_id: &str, ts: i64) {
44    if user_id <= 0 {
45        return;
46    }
47    let (Ok(servers), Ok(Some(mxid))) = (fr::remote_domains(conn, room_id), store::mxid_of(conn, user_id)) else { return };
48    let content = json!({ room_id: { "m.read": { mxid: { "event_ids": [event_id], "data": { "ts": ts } } } } });
49    queue(conn, servers, &json!({ "edu_type": "m.receipt", "content": content }));
50}
51
52/// A local user's devices or keys changed: tell every server that shares a room with them.
53pub fn enqueue_device_list(conn: &Connection, user_id: i64, stream_id: i64) {
54    if user_id <= 0 {
55        return;
56    }
57    let Ok(Some(mxid)) = store::mxid_of(conn, user_id) else { return };
58    let servers = remote_servers_of_user(conn, user_id);
59    queue(conn, servers, &json!({ "edu_type": "m.device_list_update", "content": { "user_id": mxid, "device_id": "*", "stream_id": stream_id, "prev_id": [] } }));
60}
61
62/// A local user's presence, to every server that shares a room with them.
63pub fn enqueue_presence(conn: &Connection, user_id: i64) {
64    if user_id <= 0 || !crate::http::presence::enabled() {
65        return;
66    }
67    let (Ok(Some(mxid)), Some(c)) = (store::mxid_of(conn, user_id), crate::http::presence::content_of(conn, user_id)) else { return };
68    let mut entry = json!({ "user_id": mxid, "presence": c["presence"], "last_active_ago": c["last_active_ago"], "currently_active": c["currently_active"] });
69    if let Some(m) = c.get("status_msg") {
70        entry["status_msg"] = m.clone();
71    }
72    queue(conn, remote_servers_of_user(conn, user_id), &json!({ "edu_type": "m.presence", "content": { "push": [entry] } }));
73}
74
75/// The newest device-list change of a user (0 when there is none).
76pub fn device_list_stream(conn: &Connection, user_id: i64) -> i64 {
77    conn.query_row("SELECT COALESCE(MAX(stream_id), 0) FROM device_list_changes WHERE user_id = ?1", [user_id], |r| r.get(0)).unwrap_or(0)
78}
79
80fn from_origin(mxid: &str, origin: &str) -> bool {
81    fr::domain_of(mxid) == Some(origin)
82}
83
84/// Apply one incoming EDU from `origin`. Returns the local users to wake.
85pub fn apply_inbound(conn: &mut Connection, typing: &crate::typing::TypingRegistry, origin: &str, edu: &Value) -> HashSet<i64> {
86    let mut wake = HashSet::new();
87    let content = edu.get("content").cloned().unwrap_or(Value::Null);
88    match edu.get("edu_type").and_then(Value::as_str) {
89        Some("m.typing") => {
90            let (Some(room), Some(user), Some(on)) = (content.get("room_id").and_then(Value::as_str), content.get("user_id").and_then(Value::as_str), content.get("typing").and_then(Value::as_bool)) else { return wake };
91            if !from_origin(user, origin) {
92                return wake;
93            }
94            let Ok(Some(uid)) = store::user_id_of(conn, user) else { return wake };
95            if store::room_member(conn, room, uid).ok().flatten().map(|m| m.membership) != Some(Membership::Join) {
96                return wake;
97            }
98            if typing.set_typing(room, uid, on, 30_000, Instant::now()) {
99                wake = crate::rooms::member_and_invited_ids(conn, room).unwrap_or_default();
100            }
101        }
102        Some("m.receipt") => {
103            for (room, by_type) in content.as_object().into_iter().flatten() {
104                let Some(users) = by_type.get("m.read").and_then(Value::as_object) else { continue };
105                for (user, data) in users {
106                    if !from_origin(user, origin) {
107                        continue;
108                    }
109                    let Ok(Some(uid)) = store::user_id_of(conn, user) else { continue };
110                    if store::room_member(conn, room, uid).ok().flatten().map(|m| m.membership) != Some(Membership::Join) {
111                        continue;
112                    }
113                    let ts = data.pointer("/data/ts").and_then(Value::as_i64).unwrap_or_else(now_ms);
114                    for ev in data.get("event_ids").and_then(Value::as_array).into_iter().flatten().filter_map(Value::as_str) {
115                        if store::get_event(conn, ev).ok().flatten().is_some_and(|e| e.room_id == *room) {
116                            if let Ok(w) = store::upsert_receipt(conn, room, uid, store::ReceiptType::Read, ev, ts).map(|_| crate::rooms::member_and_invited_ids(conn, room).unwrap_or_default()) {
117                                wake.extend(w);
118                            }
119                        }
120                    }
121                }
122            }
123        }
124        Some("m.device_list_update") => {
125            let Some(user) = content.get("user_id").and_then(Value::as_str) else { return wake };
126            if !from_origin(user, origin) {
127                return wake;
128            }
129            if let Ok(Some(uid)) = store::user_id_of(conn, user) {
130                if uid < 0 {
131                    let _ = crate::keys::log_device_list_change(conn, uid, &chrono::Utc::now().to_rfc3339());
132                    for room in store::rooms_for_user(conn, uid, Some(Membership::Join)).unwrap_or_default() {
133                        wake.extend(store::room_members(conn, &room, Some(Membership::Join)).unwrap_or_default().into_iter().map(|m| m.user_id).filter(|u| *u > 0));
134                    }
135                }
136            }
137        }
138        Some("m.presence") if crate::http::presence::enabled() => {
139            for p in content.get("push").and_then(Value::as_array).into_iter().flatten() {
140                let (Some(user), Some(state)) = (p.get("user_id").and_then(Value::as_str), p.get("presence").and_then(Value::as_str)) else { continue };
141                if !from_origin(user, origin) || !["online", "offline", "unavailable"].contains(&state) {
142                    continue;
143                }
144                let Ok(Some(uid)) = store::user_id_of(conn, user) else { continue };
145                if uid >= 0 {
146                    continue;
147                }
148                let ago = p.get("last_active_ago").and_then(Value::as_i64).unwrap_or(0).max(0);
149                let msg = p.get("status_msg").and_then(Value::as_str).map(|s| s.chars().take(256).collect::<String>());
150                if let Ok(ids) = crate::http::presence::set(conn, uid, state, msg.as_deref(), now_ms() - ago) {
151                    wake.extend(ids);
152                }
153            }
154        }
155        _ => {}
156    }
157    wake
158}