Skip to main content

mail4agent_server/
ephemeral.rs

1//! Typing, receipts, and read markers.
2//! [`crate::typing::TypingRegistry`] is in memory. This module does not wake.
3
4use std::collections::HashSet;
5use std::time::Instant;
6
7use rusqlite::Connection;
8
9use crate::error::MatrixError;
10use crate::store::{Membership, ReceiptType, Room};
11
12/// Applied when a `PUT .../typing/{userId}` body omits `timeout` and
13/// `typing: true` — matches
14/// [`crate::typing::TypingRegistry::set_typing`]'s own clamp
15/// ceiling (its private `TYPING_MAX_TIMEOUT_MS`), so an omitted timeout
16/// behaves the same as the largest timeout a client could ask for.
17pub const DEFAULT_TYPING_TIMEOUT_MS: u64 = 30_000;
18
19
20// ============================================================================
21// DB-only gate helper — own copy per this codebase's convention (mirrors
22// `routes::matrix::rooms::require_member`/`routes::matrix::messaging::
23// require_member` byte-for-byte).
24// ============================================================================
25
26pub fn require_member(membership: Option<Membership>) -> Result<(), MatrixError> {
27    match membership {
28        Some(Membership::Join) => Ok(()),
29        _ => Err(MatrixError::forbidden("not a member of this room")),
30    }
31}
32pub fn check_typing_target(target_mxid: &str, caller_mxid: &str) -> Result<(), MatrixError> {
33    if target_mxid == caller_mxid {
34        Ok(())
35    } else {
36        Err(MatrixError::forbidden("userId must be the caller's own mxid"))
37    }
38}
39pub fn apply_typing(
40    conn: &Connection,
41    typing_registry: &crate::typing::TypingRegistry,
42    room_id: &str,
43    caller_user_id: i64,
44    typing: bool,
45    timeout_ms: u64,
46    now: Instant,
47) -> Result<Option<HashSet<i64>>, MatrixError> {
48    crate::store::get_room(conn, room_id)?.ok_or_else(|| MatrixError::not_found("no such room"))?;
49    let membership = crate::store::room_member(conn, room_id, caller_user_id)?.map(|m| m.membership);
50    require_member(membership)?;
51
52    let changed = typing_registry.set_typing(room_id, caller_user_id, typing, timeout_ms, now);
53    if !changed {
54        return Ok(None);
55    }
56    crate::fed_edus::enqueue_typing(conn, room_id, caller_user_id, typing);
57    Ok(Some(crate::rooms::member_and_invited_ids(conn, room_id)?))
58}
59
60
61/// `POST /receipt/{receiptType}/{eventId}`'s whole DB-side decision (plan §4
62/// `receipt` row, and the `m.read`/`m.read.private` branches of `POST
63/// /read_markers`): Member, then the target event must be visible to the
64/// caller ([`crate::store::visible_upper_bound`] — for a `join`ed member
65/// (the only membership this gate accepts) this is always
66/// [`crate::store::HistoryWindow::All`] unless a `world_readable` room's
67/// event simply does not exist in THIS room at all, which the room-id filter
68/// below already catches), then [`crate::store::upsert_receipt`]'s own
69/// monotonic upsert. Returns `None` (no wake) when the upsert did not
70/// actually move the receipt forward — a backwards or repeat receipt is a
71/// silent no-op, detected here by comparing the receipt's `stream_id`
72/// before and after the upsert, never by re-deriving the monotonic
73/// comparison a second time.
74pub fn apply_receipt(
75    conn: &mut Connection,
76    room: &Room,
77    caller_user_id: i64,
78    receipt_type: ReceiptType,
79    event_id: &str,
80    now_ms: i64,
81) -> Result<Option<HashSet<i64>>, MatrixError> {
82    let caller_membership = crate::store::room_member(conn, &room.id, caller_user_id)?.map(|m| m.membership);
83    require_member(caller_membership)?;
84
85    let window = crate::store::visible_upper_bound(conn, room, caller_user_id)?;
86    let target = crate::store::get_event(conn, event_id)?
87        .filter(|e| e.room_id == room.id)
88        .ok_or_else(|| MatrixError::not_found("no such event"))?;
89    if !window.contains(target.stream_id) {
90        return Err(MatrixError::forbidden("event is outside your visible history"));
91    }
92
93    let before = crate::store::get_receipt(conn, &room.id, caller_user_id, receipt_type)?;
94    let new_stream_id = crate::store::upsert_receipt(conn, &room.id, caller_user_id, receipt_type, event_id, now_ms)?;
95    if before.as_ref().map(|r| r.stream_id) == Some(new_stream_id) {
96        return Ok(None);
97    }
98    if receipt_type == ReceiptType::Read {
99        crate::fed_edus::enqueue_receipt(conn, &room.id, caller_user_id, event_id, now_ms);
100    }
101
102    let wake_ids = match receipt_type {
103        // Visible to every joined+invited member by definition — filing a
104        // read receipt is itself proof the caller could read the event.
105        ReceiptType::Read => crate::rooms::member_and_invited_ids(conn, &room.id)?,
106        // Matrix's own privacy guarantee: a private receipt must never
107        // reach another user's `/sync` — only the caller's own devices.
108        ReceiptType::ReadPrivate => HashSet::from([caller_user_id]),
109    };
110    Ok(Some(wake_ids))
111}
112
113
114/// `POST /read_markers`'s whole DB-side decision (plan §4 `read_markers`
115/// row): Member, then `m.fully_read` (stored verbatim as per-room account
116/// data for the caller, waking only their own devices) and/or `m.read`/
117/// `m.read.private` (each delegated to [`apply_receipt`], which applies its
118/// own wake rule per type). Every field is independently optional — a call
119/// naming only one of the three touches only that one.
120pub fn apply_read_markers(
121    conn: &mut Connection,
122    room: &Room,
123    caller_user_id: i64,
124    fully_read_event_id: Option<&str>,
125    read_event_id: Option<&str>,
126    read_private_event_id: Option<&str>,
127    now_ms: i64,
128) -> Result<HashSet<i64>, MatrixError> {
129    let caller_membership = crate::store::room_member(conn, &room.id, caller_user_id)?.map(|m| m.membership);
130    require_member(caller_membership)?;
131
132    let mut wake_ids = HashSet::new();
133
134    if let Some(event_id) = fully_read_event_id {
135        crate::store::upsert_account_data(
136            conn,
137            caller_user_id,
138            &room.id,
139            "m.fully_read",
140            &serde_json::json!({ "event_id": event_id }).to_string(),
141        )?;
142        wake_ids.insert(caller_user_id);
143    }
144    if let Some(event_id) = read_event_id {
145        if let Some(ids) = apply_receipt(conn, room, caller_user_id, ReceiptType::Read, event_id, now_ms)? {
146            wake_ids.extend(ids);
147        }
148    }
149    if let Some(event_id) = read_private_event_id {
150        if let Some(ids) = apply_receipt(conn, room, caller_user_id, ReceiptType::ReadPrivate, event_id, now_ms)? {
151            wake_ids.extend(ids);
152        }
153    }
154
155    Ok(wake_ids)
156}
157
158
159// ============================================================================
160// PUT /_matrix/client/v3/rooms/{roomId}/typing/{userId}
161// ============================================================================
162
163#[derive(serde::Deserialize)]
164pub struct TypingBody {
165    pub typing: bool,
166    #[serde(default)]
167    pub timeout: Option<u64>,
168}
169
170
171// ============================================================================
172// POST /_matrix/client/v3/rooms/{roomId}/read_markers
173// ============================================================================
174
175#[derive(serde::Deserialize, Default)]
176pub struct ReadMarkersBody {
177    #[serde(rename = "m.fully_read", default)]
178    pub fully_read: Option<String>,
179    #[serde(rename = "m.read", default)]
180    pub read: Option<String>,
181    #[serde(rename = "m.read.private", default)]
182    pub read_private: Option<String>,
183}
184pub fn expired_typing_wake_ids(conn: &Connection, typing_registry: &crate::typing::TypingRegistry, now: Instant) -> (usize, HashSet<i64>) {
185    let expired_rooms = typing_registry.rooms_with_expired_typing(now);
186    let room_count = expired_rooms.len();
187    let mut ids = HashSet::new();
188    for room_id in expired_rooms {
189        match crate::rooms::member_and_invited_ids(conn, &room_id) {
190            Ok(members) => ids.extend(members),
191            Err(e) => tracing::error!("expired_typing_wake_ids: failed to list members of {}: {}", room_id, e),
192        }
193    }
194    (room_count, ids)
195}
196#[cfg(test)]
197mod tests { 
198    use super::*;
199    use std::time::Duration;
200
201    pub const T0: &str = "2026-09-24T00:00:00+00:00";
202    pub const ROOM: &str = "!testroom:example.org";
203
204    fn test_conn() -> Connection {
205        let conn = Connection::open_in_memory().expect("in-memory sqlite");
206        crate::store::create_matrix_schema(&conn).expect("matrix schema");
207        conn
208    }
209
210    /// A room with two joined members: user 1 ("alice") and user 2 ("bob").
211    fn make_room_with_two_members(conn: &mut Connection) -> (String, String) {
212        crate::store::create_room(
213            conn,
214            ROOM,
215            crate::store::RoomKind::Group,
216            1,
217            T0,
218            false,
219            crate::store::JoinRule::Invite,
220            crate::store::HistoryVisibility::Shared,
221            None,
222            None,
223        )
224        .expect("create room");
225        let alice = crate::store::ensure_matrix_user(conn, 1, "alice00000000000000000000000001", T0).expect("alice");
226        let bob = crate::store::ensure_matrix_user(conn, 2, "bob000000000000000000000000002", T0).expect("bob");
227        crate::store::apply_state_event(conn, &crate::store::StateEventWrite { event_id: "$m1", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &alice, content: r#"{"membership":"join"}"#, origin_server_ts: 900, now: T0 }).expect("alice joins");
228        crate::store::apply_state_event(conn, &crate::store::StateEventWrite { event_id: "$m2", room_id: ROOM, sender_user_id: 1, event_type: "m.room.member", state_key: &bob, content: r#"{"membership":"join"}"#, origin_server_ts: 1000, now: T0 }).expect("bob joins");
229        (alice, bob)
230    }
231
232    // ---- typing_for_another_user_is_forbidden ----
233
234    #[test]
235    fn typing_for_another_user_is_forbidden() {
236        let err = check_typing_target("@bob:example.org", "@alice:example.org").unwrap_err();
237        assert_eq!(err.errcode, "M_FORBIDDEN");
238        assert!(check_typing_target("@alice:example.org", "@alice:example.org").is_ok());
239    }
240
241    // ---- throttled_typing_does_not_wake ----
242
243    #[test]
244    fn throttled_typing_does_not_wake() {
245        let mut conn = test_conn();
246        make_room_with_two_members(&mut conn);
247        let typing = crate::typing::TypingRegistry::new();
248        let t0 = Instant::now();
249
250        let first = apply_typing(&conn, &typing, ROOM, 1, true, 5_000, t0).expect("first typing accepted");
251        assert!(first.is_some(), "the first typing:true must wake the room");
252
253        let second = apply_typing(&conn, &typing, ROOM, 1, true, 5_000, t0 + Duration::from_millis(500)).expect("throttled repeat");
254        assert!(second.is_none(), "a throttled repeat must not wake anyone");
255    }
256
257    #[test]
258    fn typing_by_a_non_member_is_forbidden() {
259        let mut conn = test_conn();
260        make_room_with_two_members(&mut conn);
261        let typing = crate::typing::TypingRegistry::new();
262
263        let err = apply_typing(&conn, &typing, ROOM, 999, true, 5_000, Instant::now()).unwrap_err();
264        assert_eq!(err.errcode, "M_FORBIDDEN");
265    }
266
267    // ---- private_receipt_wakes_only_own_devices ----
268
269    #[test]
270    fn private_receipt_wakes_only_own_devices() {
271        let mut conn = test_conn();
272        make_room_with_two_members(&mut conn);
273        let event = crate::store::insert_timeline_event(&mut conn, "$msg", ROOM, 2, "m.room.message", "{}", 1000).expect("message");
274
275        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
276        let wake_ids = apply_receipt(&mut conn, &room, 1, ReceiptType::ReadPrivate, &event.event_id, 1500)
277            .expect("apply private receipt")
278            .expect("a new receipt must report a change");
279        assert_eq!(wake_ids, HashSet::from([1]), "m.read.private must wake ONLY the caller's own devices, never other room members");
280    }
281
282    #[test]
283    fn read_receipt_wakes_every_joined_member() {
284        let mut conn = test_conn();
285        make_room_with_two_members(&mut conn);
286        let event = crate::store::insert_timeline_event(&mut conn, "$msg", ROOM, 2, "m.room.message", "{}", 1000).expect("message");
287
288        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
289        let wake_ids = apply_receipt(&mut conn, &room, 1, ReceiptType::Read, &event.event_id, 1500)
290            .expect("apply read receipt")
291            .expect("a new receipt must report a change");
292        assert_eq!(wake_ids, HashSet::from([1, 2]), "m.read must wake every joined member, including the caller's own other devices");
293    }
294
295    // ---- backwards_receipt_is_a_noop ----
296
297    #[test]
298    fn backwards_receipt_is_a_noop() {
299        let mut conn = test_conn();
300        make_room_with_two_members(&mut conn);
301        let e1 = crate::store::insert_timeline_event(&mut conn, "$e1", ROOM, 2, "m.room.message", "{}", 1000).expect("e1");
302        let e2 = crate::store::insert_timeline_event(&mut conn, "$e2", ROOM, 2, "m.room.message", "{}", 2000).expect("e2");
303        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
304
305        let forward = apply_receipt(&mut conn, &room, 1, ReceiptType::Read, &e2.event_id, 2500).expect("advance to e2");
306        assert!(forward.is_some());
307
308        let backward = apply_receipt(&mut conn, &room, 1, ReceiptType::Read, &e1.event_id, 3000).expect("attempted backward move");
309        assert!(backward.is_none(), "a backwards receipt must be a silent no-op, never a wake");
310
311        let stored = crate::store::get_receipt(&conn, ROOM, 1, ReceiptType::Read).expect("get receipt").expect("receipt exists");
312        assert_eq!(stored.event_id, e2.event_id, "the receipt must still point at e2, not have moved backward to e1");
313    }
314
315    // ---- fully_read_is_stored_as_room_account_data ----
316
317    #[test]
318    fn fully_read_is_stored_as_room_account_data() {
319        let mut conn = test_conn();
320        make_room_with_two_members(&mut conn);
321        let event = crate::store::insert_timeline_event(&mut conn, "$msg", ROOM, 2, "m.room.message", "{}", 1000).expect("message");
322        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
323
324        let wake_ids = apply_read_markers(&mut conn, &room, 1, Some(&event.event_id), None, None, 1500).expect("apply read markers");
325        assert_eq!(wake_ids, HashSet::from([1]), "m.fully_read must wake only the caller's own devices");
326
327        let stored = crate::store::get_account_data(&conn, 1, ROOM, "m.fully_read").expect("get account data").expect("row exists");
328        let content: serde_json::Value = serde_json::from_str(&stored.content).expect("json content");
329        assert_eq!(content, serde_json::json!({ "event_id": event.event_id }));
330    }
331
332    #[test]
333    fn read_markers_can_combine_fully_read_with_a_receipt() {
334        let mut conn = test_conn();
335        make_room_with_two_members(&mut conn);
336        let event = crate::store::insert_timeline_event(&mut conn, "$msg", ROOM, 2, "m.room.message", "{}", 1000).expect("message");
337        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
338
339        let wake_ids =
340            apply_read_markers(&mut conn, &room, 1, Some(&event.event_id), Some(&event.event_id), None, 1500).expect("apply read markers");
341        assert_eq!(wake_ids, HashSet::from([1, 2]), "the embedded m.read receipt must wake every joined member on top of the caller's own devices");
342        assert!(crate::store::get_account_data(&conn, 1, ROOM, "m.fully_read").expect("get account data").is_some());
343        assert!(crate::store::get_receipt(&conn, ROOM, 1, ReceiptType::Read).expect("get receipt").is_some());
344    }
345
346    // ---- receipt_on_invisible_event_is_refused ----
347
348    #[test]
349    fn receipt_on_invisible_event_is_refused() {
350        let mut conn = test_conn();
351        make_room_with_two_members(&mut conn);
352
353        pub const OTHER_ROOM: &str = "!otherroom:example.org";
354        crate::store::create_room(
355            &conn,
356            OTHER_ROOM,
357            crate::store::RoomKind::Group,
358            1,
359            T0,
360            false,
361            crate::store::JoinRule::Invite,
362            crate::store::HistoryVisibility::Shared,
363            None,
364            None,
365        )
366        .expect("create other room");
367        let foreign_event = crate::store::insert_timeline_event(&mut conn, "$foreign", OTHER_ROOM, 1, "m.room.message", "{}", 1000).expect("foreign event");
368
369        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
370        let err = apply_receipt(&mut conn, &room, 1, ReceiptType::Read, &foreign_event.event_id, 1500).unwrap_err();
371        assert_eq!(err.errcode, "M_NOT_FOUND");
372    }
373
374    #[test]
375    fn receipt_on_a_nonexistent_event_is_refused() {
376        let mut conn = test_conn();
377        make_room_with_two_members(&mut conn);
378        let room = crate::store::get_room(&conn, ROOM).expect("get room").expect("room exists");
379
380        let err = apply_receipt(&mut conn, &room, 1, ReceiptType::Read, "$never-existed", 1500).unwrap_err();
381        assert_eq!(err.errcode, "M_NOT_FOUND");
382    }
383
384    // ---- expired_typing_wake_ids (wake_rooms_with_expired_typing's DB-only half) ----
385
386    #[test]
387    pub fn expired_typing_wake_ids_covers_the_expired_rooms_joined_members() {
388        let mut conn = test_conn();
389        make_room_with_two_members(&mut conn);
390        let typing = crate::typing::TypingRegistry::new();
391        let t0 = Instant::now();
392        assert!(typing.set_typing(ROOM, 1, true, 100, t0));
393
394        let (room_count, ids) = expired_typing_wake_ids(&conn, &typing, t0 + Duration::from_millis(500));
395        assert_eq!(room_count, 1, "exactly one room had its typing state expire");
396        assert_eq!(ids, HashSet::from([1, 2]), "every joined member of the expired room must be woken");
397
398        // A second sweep with nothing newly expired reports nothing.
399        let (room_count_again, ids_again) = expired_typing_wake_ids(&conn, &typing, t0 + Duration::from_millis(600));
400        assert_eq!(room_count_again, 0, "a room with no NEWLY expired typing flag must not be reported again");
401        assert!(ids_again.is_empty(), "a room with no NEWLY expired typing flag must not be reported again");
402    }
403
404    #[test]
405    pub fn expired_typing_wake_ids_is_empty_when_nothing_is_typing() {
406        let conn = test_conn();
407        let typing = crate::typing::TypingRegistry::new();
408        let (room_count, ids) = expired_typing_wake_ids(&conn, &typing, Instant::now());
409        assert_eq!(room_count, 0);
410        assert!(ids.is_empty());
411    }
412
413}