1use std::collections::HashSet;
5use std::time::Instant;
6
7use rusqlite::Connection;
8
9use crate::error::MatrixError;
10use crate::store::{Membership, ReceiptType, Room};
11
12pub const DEFAULT_TYPING_TIMEOUT_MS: u64 = 30_000;
18
19
20pub 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
60pub 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 ReceiptType::Read => crate::rooms::member_and_invited_ids(conn, &room.id)?,
102 ReceiptType::ReadPrivate => HashSet::from([caller_user_id]),
105 };
106 Ok(Some(wake_ids))
107}
108
109
110pub 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#[derive(serde::Deserialize)]
160pub struct TypingBody {
161 pub typing: bool,
162 #[serde(default)]
163 pub timeout: Option<u64>,
164}
165
166
167#[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 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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 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}