1use std::sync::Arc;
9
10use axum::extract::{Path, State};
11use axum::http::HeaderMap;
12use axum::routing::get;
13use axum::{Json, Router};
14use rusqlite::{params, Connection, OptionalExtension};
15use serde_json::{json, Value};
16
17use super::{resolve_caller, with_conn_pub, wake_users, Homeserver};
18use crate::error::MatrixError;
19use crate::store;
20
21const IDLE_MS: i64 = 5 * 60 * 1000;
23
24pub fn enabled() -> bool {
25 cfg!(feature = "presence") && std::env::var("M4A_PRESENCE").map(|v| v != "off").unwrap_or(true)
26}
27
28pub(super) fn routes() -> Router<Arc<Homeserver>> {
29 Router::new().route("/client/v3/presence/{user_id}/status", get(get_status).put(put_status))
30}
31
32pub fn create_schema(conn: &Connection) -> rusqlite::Result<()> {
33 conn.execute_batch(
34 "CREATE TABLE IF NOT EXISTS presence (
35 user_id INTEGER PRIMARY KEY,
36 state TEXT NOT NULL,
37 status_msg TEXT,
38 last_active_ms INTEGER NOT NULL,
39 stream_id INTEGER NOT NULL
40 );
41 CREATE INDEX IF NOT EXISTS idx_presence_stream ON presence(stream_id);",
42 )
43}
44
45fn now_ms() -> i64 {
46 chrono::Utc::now().timestamp_millis()
47}
48
49pub fn content_of(conn: &Connection, user_id: i64) -> Option<Value> {
51 let (state, msg, last): (String, Option<String>, i64) = conn
52 .query_row("SELECT state, status_msg, last_active_ms FROM presence WHERE user_id = ?1", [user_id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))
53 .optional()
54 .ok()??;
55 let ago = (now_ms() - last).max(0);
56 let shown = if state == "online" && ago > IDLE_MS { "unavailable".to_string() } else { state };
57 let mut v = json!({ "presence": shown, "last_active_ago": ago, "currently_active": shown == "online" });
58 if let Some(m) = msg {
59 v["status_msg"] = json!(m);
60 }
61 Some(v)
62}
63
64pub fn set(conn: &mut Connection, user_id: i64, state: &str, msg: Option<&str>, last_active_ms: i64) -> rusqlite::Result<std::collections::HashSet<i64>> {
66 let tx = conn.transaction()?;
67 let stream = store::next_stream_id(&tx)?;
68 tx.execute(
69 "INSERT INTO presence (user_id, state, status_msg, last_active_ms, stream_id) VALUES (?1, ?2, ?3, ?4, ?5)
70 ON CONFLICT(user_id) DO UPDATE SET state = excluded.state, status_msg = excluded.status_msg, last_active_ms = excluded.last_active_ms, stream_id = excluded.stream_id",
71 params![user_id, state, msg, last_active_ms, stream],
72 )?;
73 tx.commit()?;
74 let mut ids = crate::key_ops::peers_sharing_a_room_with(conn, user_id)?;
75 ids.retain(|u| *u > 0);
76 Ok(ids)
77}
78
79pub fn sync_events(conn: &Connection, caller: i64, since_stream: i64, upto: i64) -> rusqlite::Result<Vec<Value>> {
81 if !enabled() {
82 return Ok(vec![]);
83 }
84 let peers = crate::key_ops::peers_sharing_a_room_with(conn, caller)?;
85 let mut st = conn.prepare("SELECT user_id FROM presence WHERE stream_id > ?1 AND stream_id <= ?2")?;
86 let changed: Vec<i64> = st.query_map(params![since_stream, upto], |r| r.get(0))?.flatten().collect();
87 let mut out = Vec::new();
88 for u in changed.into_iter().filter(|u| peers.contains(u) && *u != caller) {
89 if let (Some(mxid), Some(content)) = (store::mxid_of(conn, u)?, content_of(conn, u)) {
90 out.push(json!({ "type": "m.presence", "sender": mxid, "content": content }));
91 }
92 }
93 Ok(out)
94}
95
96async fn put_status(State(state): State<Arc<Homeserver>>, headers: HeaderMap, Path(user_id): Path<String>, Json(body): Json<Value>) -> Result<Json<Value>, MatrixError> {
97 let caller = resolve_caller(&state, &headers, None).await?;
98 crate::account::check_caller_owns_user_id(&user_id, &caller.mxid)?;
99 let st = body.get("presence").and_then(Value::as_str).ok_or_else(|| MatrixError::bad_json("presence"))?.to_string();
100 if !["online", "offline", "unavailable"].contains(&st.as_str()) {
101 return Err(MatrixError::invalid_param("presence must be online, offline or unavailable"));
102 }
103 let msg = body.get("status_msg").and_then(Value::as_str).map(|s| s.chars().take(256).collect::<String>());
104 if !enabled() {
105 return Ok(Json(json!({})));
106 }
107 let uid = caller.user_id;
108 let ids = with_conn_pub(&state, move |c| {
109 let ids = set(c, uid, &st, msg.as_deref(), now_ms()).map_err(|_| MatrixError::internal())?;
110 crate::fed_edus::enqueue_presence(c, uid);
111 Ok(ids)
112 })
113 .await?;
114 wake_users(&state, ids);
115 Ok(Json(json!({})))
116}
117
118async fn get_status(State(state): State<Arc<Homeserver>>, headers: HeaderMap, Path(user_id): Path<String>) -> Result<Json<Value>, MatrixError> {
119 let caller = resolve_caller(&state, &headers, None).await?;
120 if !enabled() {
121 return Ok(Json(json!({ "presence": "offline" })));
122 }
123 let me = caller.user_id;
124 with_conn_pub(&state, move |c| {
125 let target = store::user_id_of(c, &user_id)?.ok_or_else(|| MatrixError::not_found("unknown user"))?;
126 if target != me && !crate::key_ops::peers_sharing_a_room_with(c, me)?.contains(&target) {
127 return Err(MatrixError::forbidden("you share no room with that user"));
128 }
129 Ok(Json(content_of(c, target).unwrap_or_else(|| json!({ "presence": "offline" }))))
130 })
131 .await
132}