use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Mutex;
use std::time::{Duration, Instant};
const TYPING_THROTTLE_WINDOW: Duration = Duration::from_secs(3);
const TYPING_MAX_TIMEOUT_MS: u64 = 30_000;
struct TypingUserState {
expires_at: Instant,
is_active: bool,
last_true_accepted_at: Option<Instant>,
}
#[derive(Default)]
struct RoomTyping {
users: HashMap<i64, TypingUserState>,
serial: u64,
}
fn reap_room(room: &mut RoomTyping, now: Instant, next_gen: &AtomicU64) -> bool {
let mut bumped = false;
for state in room.users.values_mut() {
if state.is_active && state.expires_at <= now {
state.is_active = false;
bumped = true;
}
}
if bumped {
room.serial = next_gen.fetch_add(1, Ordering::AcqRel) + 1;
}
room.users.retain(|_, state| {
state.is_active
|| state
.last_true_accepted_at
.is_some_and(|last| now.saturating_duration_since(last) <= TYPING_THROTTLE_WINDOW)
});
bumped
}
pub struct TypingRegistry {
rooms: Mutex<HashMap<String, RoomTyping>>,
next_gen: AtomicU64,
}
impl Default for TypingRegistry {
fn default() -> Self {
Self::new()
}
}
impl TypingRegistry {
pub fn new() -> Self {
Self::with_seed(0)
}
pub fn with_seed(seed: u64) -> Self {
Self { rooms: Mutex::new(HashMap::new()), next_gen: AtomicU64::new(seed) }
}
pub fn set_typing(
&self,
room_id: &str,
user_id: i64,
typing: bool,
timeout_ms: u64,
now: Instant,
) -> bool {
let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
let room = rooms.entry(room_id.to_string()).or_default();
let changed = if typing {
let timeout = Duration::from_millis(timeout_ms.min(TYPING_MAX_TIMEOUT_MS));
let state = room.users.entry(user_id).or_insert_with(|| TypingUserState {
expires_at: now,
is_active: false,
last_true_accepted_at: None,
});
let throttled = state
.last_true_accepted_at
.is_some_and(|last| now.saturating_duration_since(last) < TYPING_THROTTLE_WINDOW);
if throttled {
false
} else {
state.last_true_accepted_at = Some(now);
state.expires_at = now + timeout;
let was_active = state.is_active;
state.is_active = true;
!was_active
}
} else {
match room.users.get_mut(&user_id) {
Some(state) if state.is_active => {
state.is_active = false;
state.expires_at = now;
true
}
_ => false,
}
};
if changed {
room.serial = self.next_gen.fetch_add(1, Ordering::AcqRel) + 1;
}
changed
}
pub fn current_typing_gen(&self) -> u64 {
self.next_gen.load(Ordering::Acquire)
}
pub fn typing_users(&self, room_id: &str, now: Instant) -> Vec<i64> {
let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
let mut users = Vec::new();
let mut room_now_empty = false;
if let Some(room) = rooms.get_mut(room_id) {
reap_room(room, now, &self.next_gen);
users.extend(room.users.iter().filter(|(_, state)| state.is_active).map(|(id, _)| *id));
room_now_empty = room.users.is_empty();
}
if room_now_empty {
rooms.remove(room_id);
}
users.sort_unstable();
users
}
pub fn typing_serial(&self, room_id: &str, now: Instant) -> u64 {
let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
let mut serial = 0;
let mut room_now_empty = false;
if let Some(room) = rooms.get_mut(room_id) {
reap_room(room, now, &self.next_gen);
serial = room.serial;
room_now_empty = room.users.is_empty();
}
if room_now_empty {
rooms.remove(room_id);
}
serial
}
pub fn rooms_with_expired_typing(&self, now: Instant) -> Vec<String> {
let mut rooms = self.rooms.lock().unwrap_or_else(|e| e.into_inner());
let mut changed = Vec::new();
rooms.retain(|room_id, room| {
if reap_room(room, now, &self.next_gen) {
changed.push(room_id.clone());
}
!room.users.is_empty()
});
changed
}
}