mail4agent_server/typing.rs
1//! In-memory typing set. Never written to the messenger database.
2//! The builder seeds [`TypingRegistry::with_seed`] from a clock and notifies
3//! clients itself when `set_typing` or `rooms_with_expired_typing` report a change.
4
5use std::collections::HashMap;
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::sync::Mutex;
8use std::time::{Duration, Instant};
9
10/// A repeat `typing: true` PUT for the same `(room, user)` inside this
11/// window is throttled — accepted independently of the IP rate limiter
12/// (§3.6/§3.8), so a chatty client cannot force a wake on every keystroke.
13const TYPING_THROTTLE_WINDOW: Duration = Duration::from_secs(3);
14
15/// A client-supplied typing timeout above this is clamped down to it.
16const TYPING_MAX_TIMEOUT_MS: u64 = 30_000;
17
18/// One `(room, user)`'s typing state: `expires_at` is when a currently-true
19/// typing flag auto-clears; `is_active` caches whether this user currently
20/// counts as typing (kept in sync by every touch and by the lazy reap, so a
21/// membership transition is detected exactly once); `last_true_accepted_at`
22/// is `None` until the first accepted `typing: true`, purely for the
23/// throttle above — a `typing: false` never touches it.
24struct TypingUserState {
25 expires_at: Instant,
26 is_active: bool,
27 last_true_accepted_at: Option<Instant>,
28}
29
30/// One room's typing state plus a serial stamped from [`TypingRegistry`]'s
31/// GLOBAL `next_gen` counter on every EFFECTIVE change (a user starting or
32/// stopping counting as typing) — a refreshed-but-still-active `typing:
33/// true`, or a throttled repeat, bumps nothing, since nothing about what a
34/// peer would see has changed. See the module doc's "The typing generation
35/// is GLOBAL, not per-room" section for why this is not a local counter.
36#[derive(Default)]
37struct RoomTyping {
38 users: HashMap<i64, TypingUserState>,
39 serial: u64,
40}
41
42/// Marks every entry in `room` whose typing flag has timed out as no
43/// longer active (stamping `room.serial` from `next_gen` once if at least
44/// one did), then drops entries that are both inactive and outside the
45/// throttle memory window (nothing left worth remembering). Returns whether
46/// the serial bumped, so callers can report "this room's typing set just
47/// changed" without the caller re-deriving it.
48fn reap_room(room: &mut RoomTyping, now: Instant, next_gen: &AtomicU64) -> bool {
49 let mut bumped = false;
50 for state in room.users.values_mut() {
51 if state.is_active && state.expires_at <= now {
52 state.is_active = false;
53 bumped = true;
54 }
55 }
56 if bumped {
57 room.serial = next_gen.fetch_add(1, Ordering::AcqRel) + 1;
58 }
59 room.users.retain(|_, state| {
60 state.is_active
61 || state
62 .last_true_accepted_at
63 .is_some_and(|last| now.saturating_duration_since(last) <= TYPING_THROTTLE_WINDOW)
64 });
65 bumped
66}
67
68/// In-memory, per-room typing indicator state (§3.6) — deliberately never
69/// written to `messenger.db` (see the module doc). Not related to
70/// [`LiveRegistry`]'s own keys; a caller that gets `changed == true` back
71/// from [`Self::set_typing`] (or a non-empty room list back from
72/// [`Self::rooms_with_expired_typing`]) is the one that decides which
73/// `LiveRegistry` keys to wake for it.
74pub struct TypingRegistry {
75 rooms: Mutex<HashMap<String, RoomTyping>>,
76 /// The one counter every room's [`RoomTyping::serial`] is stamped from —
77 /// see the module doc's "The typing generation is GLOBAL, not per-room"
78 /// section.
79 next_gen: AtomicU64,
80}
81
82impl Default for TypingRegistry {
83 fn default() -> Self {
84 Self::new()
85 }
86}
87
88impl TypingRegistry {
89 /// A fresh registry whose typing generation starts at `0` — this
90 /// codebase's own tests use this (deterministic small values are easy
91 /// to assert against); production boot uses [`Self::with_seed`] instead
92 /// (see that constructor's own doc for why).
93 pub fn new() -> Self {
94 Self::with_seed(0)
95 }
96
97 /// A fresh registry whose typing generation starts at `seed` (the FIRST
98 /// effective change stamps `seed + 1`, matching [`Self::new`]'s own
99 /// "first change stamps 1" behavior when `seed == 0`). Production boot
100 /// (`main.rs`) seeds this from the current Unix time in milliseconds —
101 /// see the module doc's "The typing generation is GLOBAL, not per-room"
102 /// section: an in-memory counter starting at `0` on every restart can
103 /// under-report a typing change to a client still holding a
104 /// higher-`typing_gen` token from before that restart. Seeding from a
105 /// wall-clock reading that only ever increases makes a POST-restart
106 /// seed reliably higher than anything a PRE-restart process could ever
107 /// have reached (an `AtomicU64` counter incrementing once per typing
108 /// change, even at an implausible sustained rate, cannot climb into the
109 /// same range as milliseconds-since-1970 within any real process
110 /// lifetime) — closing the under-reporting gap [`Self::current_typing_gen`]'s
111 /// own doc used to accept as a known limitation.
112 pub fn with_seed(seed: u64) -> Self {
113 Self { rooms: Mutex::new(HashMap::new()), next_gen: AtomicU64::new(seed) }
114 }
115
116 /// Apply one `PUT /rooms/{roomId}/typing/{userId}` call. `timeout_ms`
117 /// (only meaningful when `typing == true`) is clamped to
118 /// [`TYPING_MAX_TIMEOUT_MS`]. Returns whether this call actually
119 /// changed who counts as typing in `room_id` — `false` for a throttled
120 /// repeat `true` (§3.6), a refresh of an already-active `true`, or a
121 /// `false` for a user who was not counted as typing — callers should
122 /// skip waking anyone when this returns `false`.
123 pub fn set_typing(
124 &self,
125 room_id: &str,
126 user_id: i64,
127 typing: bool,
128 timeout_ms: u64,
129 now: Instant,
130 ) -> bool {
131 let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
132 let room = rooms.entry(room_id.to_string()).or_default();
133
134 let changed = if typing {
135 let timeout = Duration::from_millis(timeout_ms.min(TYPING_MAX_TIMEOUT_MS));
136 let state = room.users.entry(user_id).or_insert_with(|| TypingUserState {
137 expires_at: now,
138 is_active: false,
139 last_true_accepted_at: None,
140 });
141 let throttled = state
142 .last_true_accepted_at
143 .is_some_and(|last| now.saturating_duration_since(last) < TYPING_THROTTLE_WINDOW);
144 if throttled {
145 false
146 } else {
147 state.last_true_accepted_at = Some(now);
148 state.expires_at = now + timeout;
149 let was_active = state.is_active;
150 state.is_active = true;
151 !was_active
152 }
153 } else {
154 match room.users.get_mut(&user_id) {
155 Some(state) if state.is_active => {
156 state.is_active = false;
157 state.expires_at = now;
158 true
159 }
160 _ => false,
161 }
162 };
163
164 if changed {
165 room.serial = self.next_gen.fetch_add(1, Ordering::AcqRel) + 1;
166 }
167 changed
168 }
169
170 /// The current global typing generation — every effective typing change
171 /// across every room takes its own room's [`RoomTyping::serial`] from
172 /// this SAME counter (see the module doc), so "room serial > token's
173 /// typing_gen" identifies exactly the rooms whose typing changed since
174 /// that token was issued. `0` means no typing change has EVER been
175 /// accepted by this registry instance — a fresh boot, or simply an idle
176 /// server; a token carrying `typing_gen: 0` (the bare `s{stream_id}`
177 /// legacy form `routes::matrix::sync_token` still accepts, or a token
178 /// issued before this server ever saw a single typing PUT) never
179 /// wrongly suppresses a real future change, since every real change
180 /// stamps a value `>= 1`.
181 ///
182 /// In-memory only, like every other fact this registry holds — a plain
183 /// `TypingRegistry::new()` restarts this counter at `0` on every
184 /// process restart. Production boot avoids the resulting under-report
185 /// risk (a token issued before a restart carrying a `typing_gen` LARGER
186 /// than anything a freshly-zeroed registry has assigned yet) by
187 /// constructing via [`Self::with_seed`] instead, seeded from a
188 /// wall-clock reading that is always higher than any pre-restart value
189 /// — see that constructor's own doc.
190 pub fn current_typing_gen(&self) -> u64 {
191 self.next_gen.load(Ordering::Acquire)
192 }
193
194 /// Current (non-expired) typers in `room_id`, sorted for a deterministic
195 /// response — expired entries are dropped lazily as a side effect.
196 /// Empty (including for a room this registry has never heard of) rather
197 /// than an error — matches `m.typing`'s own "no event at all if nobody
198 /// is typing" convention (§3.6).
199 pub fn typing_users(&self, room_id: &str, now: Instant) -> Vec<i64> {
200 let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
201 let mut users = Vec::new();
202 let mut room_now_empty = false;
203 if let Some(room) = rooms.get_mut(room_id) {
204 reap_room(room, now, &self.next_gen);
205 users.extend(room.users.iter().filter(|(_, state)| state.is_active).map(|(id, _)| *id));
206 room_now_empty = room.users.is_empty();
207 }
208 if room_now_empty {
209 rooms.remove(room_id);
210 }
211 users.sort_unstable();
212 users
213 }
214
215 /// `room_id`'s current typing serial — stamped from the SAME global
216 /// counter [`Self::current_typing_gen`] reads (see the module doc) on
217 /// every effective change; `0` for a room this registry has never heard
218 /// of, or one whose typing set has never actually changed since boot.
219 /// Reaps expired entries first, so this is never stale relative to
220 /// `now`. `routes::matrix::sync`'s own per-room condition is `is_initial
221 /// || typing_serial(room, now) > token.typing_gen`.
222 pub fn typing_serial(&self, room_id: &str, now: Instant) -> u64 {
223 let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
224 let mut serial = 0;
225 let mut room_now_empty = false;
226 if let Some(room) = rooms.get_mut(room_id) {
227 reap_room(room, now, &self.next_gen);
228 serial = room.serial;
229 room_now_empty = room.users.is_empty();
230 }
231 if room_now_empty {
232 rooms.remove(room_id);
233 }
234 serial
235 }
236
237 /// Sweep every room, reaping any typing flag that has timed out since
238 /// it was last touched, and return the `room_id`s whose typing set
239 /// changed as a result. A periodic tick (or the `/sync` path itself)
240 /// calls this and wakes the returned rooms' members — without it, a
241 /// member who stops typing and never sends another `/sync`-triggering
242 /// action would never have their typing flag cleared for anyone
243 /// currently blocked in a long poll.
244 pub fn rooms_with_expired_typing(&self, now: Instant) -> Vec<String> {
245 let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
246 let mut changed = Vec::new();
247 rooms.retain(|room_id, room| {
248 if reap_room(room, now, &self.next_gen) {
249 changed.push(room_id.clone());
250 }
251 !room.users.is_empty()
252 });
253 changed
254 }
255}