1use 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
33pub 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
42pub 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
52pub 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
62pub 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
75pub 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
84pub 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}