Skip to main content

mail4agent_server/http/
presence.rs

1//! Presence: `PUT/GET /presence/{userId}/status`, the sync `presence` section, and `m.presence`
2//! EDUs between servers.
3//!
4//! Optional: the `presence` cargo feature (default on) and the environment switch
5//! `M4A_PRESENCE=off`. Off, a PUT is accepted and forgotten, a GET says `offline`, sync carries no
6//! presence and no EDU is sent or applied.
7
8use 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
21/// A user who has been silent this long while "online" is reported "unavailable".
22const 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
49/// The content of one user's presence (as in `m.presence` events), if any is known.
50pub 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
64/// Record a presence change (own user or a remote one learned by EDU); returns the users to wake.
65pub 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
79/// Sync's `presence.events` for `caller`: peers whose presence changed after `since_stream`.
80pub 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}