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 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
61pub 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 ReceiptType::Read => crate::rooms::member_and_invited_ids(conn, &room.id)?,
106 ReceiptType::ReadPrivate => HashSet::from([caller_user_id]),
109 };
110 Ok(Some(wake_ids))
111}
112
113
114pub 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#[derive(serde::Deserialize)]
164pub struct TypingBody {
165 pub typing: bool,
166 #[serde(default)]
167 pub timeout: Option<u64>,
168}
169
170
171#[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 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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 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}