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