Skip to main content

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}