fraiseql_server/subscriptions/
presence.rs1use std::{
10 collections::HashMap,
11 sync::atomic::{AtomicU64, Ordering},
12 time::Duration,
13};
14
15use serde::{Deserialize, Serialize};
16use tokio::{sync::RwLock, time::Instant};
17use tracing::debug;
18
19#[derive(Debug, Clone)]
23pub struct PresenceConfig {
24 pub max_members_per_room: usize,
26
27 pub max_rooms: usize,
29
30 pub heartbeat_timeout: Duration,
32}
33
34impl PresenceConfig {
35 #[must_use]
37 pub const fn new() -> Self {
38 Self {
39 max_members_per_room: 500,
40 max_rooms: 10_000,
41 heartbeat_timeout: Duration::from_secs(30),
42 }
43 }
44}
45
46impl Default for PresenceConfig {
47 fn default() -> Self {
48 Self::new()
49 }
50}
51
52#[derive(Debug, Clone, Serialize, Deserialize)]
56pub struct PresenceMember {
57 pub id: String,
59
60 pub state: serde_json::Value,
62
63 #[serde(skip, default = "Instant::now")]
65 pub last_seen: Instant,
66}
67
68#[derive(Debug, Clone, Serialize)]
70pub struct PresenceState {
71 pub room: String,
73
74 pub members: Vec<PresenceMember>,
76}
77
78#[derive(Debug, Clone, Serialize)]
80pub struct PresenceDiff {
81 pub room: String,
83
84 pub joins: Vec<PresenceMember>,
86
87 pub leaves: Vec<String>,
89}
90
91#[derive(Debug)]
95struct PresenceRoom {
96 members: HashMap<String, PresenceMember>,
98}
99
100impl PresenceRoom {
101 fn new() -> Self {
102 Self {
103 members: HashMap::new(),
104 }
105 }
106}
107
108#[derive(Debug, Clone)]
112pub struct PresenceStats {
113 pub active_rooms: usize,
115
116 pub total_members: usize,
118
119 pub joins_total: u64,
121
122 pub leaves_total: u64,
124
125 pub evictions_total: u64,
127}
128
129#[derive(Debug)]
133pub struct PresenceManager {
134 rooms: RwLock<HashMap<String, PresenceRoom>>,
135 config: PresenceConfig,
136 joins_total: AtomicU64,
137 leaves_total: AtomicU64,
138 evictions_total: AtomicU64,
139}
140
141impl PresenceManager {
142 #[must_use]
144 pub fn new(config: PresenceConfig) -> Self {
145 Self {
146 rooms: RwLock::new(HashMap::new()),
147 config,
148 joins_total: AtomicU64::new(0),
149 leaves_total: AtomicU64::new(0),
150 evictions_total: AtomicU64::new(0),
151 }
152 }
153
154 pub async fn join(
163 &self,
164 room: &str,
165 member_id: &str,
166 state: serde_json::Value,
167 ) -> Result<(PresenceState, PresenceDiff), PresenceError> {
168 let mut rooms = self.rooms.write().await;
169
170 if !rooms.contains_key(room) {
172 if rooms.len() >= self.config.max_rooms {
173 return Err(PresenceError::TooManyRooms {
174 max: self.config.max_rooms,
175 });
176 }
177 rooms.insert(room.to_string(), PresenceRoom::new());
178 }
179
180 let Some(presence_room) = rooms.get_mut(room) else {
181 return Err(PresenceError::TooManyRooms {
183 max: self.config.max_rooms,
184 });
185 };
186
187 if !presence_room.members.contains_key(member_id)
189 && presence_room.members.len() >= self.config.max_members_per_room
190 {
191 return Err(PresenceError::RoomFull {
192 room: room.to_string(),
193 max: self.config.max_members_per_room,
194 });
195 }
196
197 let member = PresenceMember {
198 id: member_id.to_string(),
199 state,
200 last_seen: Instant::now(),
201 };
202
203 presence_room.members.insert(member_id.to_string(), member.clone());
204 self.joins_total.fetch_add(1, Ordering::Relaxed);
205
206 let presence_state = PresenceState {
207 room: room.to_string(),
208 members: presence_room.members.values().cloned().collect(),
209 };
210
211 let diff = PresenceDiff {
212 room: room.to_string(),
213 joins: vec![member],
214 leaves: vec![],
215 };
216
217 debug!(
218 room,
219 member_id,
220 members = presence_room.members.len(),
221 "presence: member joined"
222 );
223 Ok((presence_state, diff))
224 }
225
226 pub async fn leave(&self, room: &str, member_id: &str) -> Option<PresenceDiff> {
231 let mut rooms = self.rooms.write().await;
232
233 let presence_room = rooms.get_mut(room)?;
234 presence_room.members.remove(member_id)?;
235 self.leaves_total.fetch_add(1, Ordering::Relaxed);
236
237 debug!(room, member_id, members = presence_room.members.len(), "presence: member left");
238
239 let diff = PresenceDiff {
240 room: room.to_string(),
241 joins: vec![],
242 leaves: vec![member_id.to_string()],
243 };
244
245 if presence_room.members.is_empty() {
247 rooms.remove(room);
248 debug!(room, "presence: room removed (empty)");
249 }
250
251 Some(diff)
252 }
253
254 pub async fn heartbeat(&self, room: &str, member_id: &str) -> bool {
258 let mut rooms = self.rooms.write().await;
259
260 if let Some(presence_room) = rooms.get_mut(room) {
261 if let Some(member) = presence_room.members.get_mut(member_id) {
262 member.last_seen = Instant::now();
263 return true;
264 }
265 }
266
267 false
268 }
269
270 pub async fn update_state(
275 &self,
276 room: &str,
277 member_id: &str,
278 new_state: serde_json::Value,
279 ) -> Option<PresenceDiff> {
280 let mut rooms = self.rooms.write().await;
281 let presence_room = rooms.get_mut(room)?;
282 let member = presence_room.members.get_mut(member_id)?;
283
284 member.state = new_state;
285 member.last_seen = Instant::now();
286
287 Some(PresenceDiff {
288 room: room.to_string(),
289 joins: vec![member.clone()],
290 leaves: vec![],
291 })
292 }
293
294 pub async fn evict_stale(&self) -> Vec<PresenceDiff> {
298 let timeout = self.config.heartbeat_timeout;
299 let mut rooms = self.rooms.write().await;
300 let mut diffs = Vec::new();
301 let mut empty_rooms = Vec::new();
302
303 for (room_name, room) in rooms.iter_mut() {
304 let mut evicted = Vec::new();
305
306 room.members.retain(|id, member| {
307 if member.last_seen.elapsed() > timeout {
308 evicted.push(id.clone());
309 false
310 } else {
311 true
312 }
313 });
314
315 if !evicted.is_empty() {
316 let count = evicted.len();
317 self.evictions_total.fetch_add(count as u64, Ordering::Relaxed);
318 debug!(room = %room_name, evicted = count, "presence: evicted stale members");
319
320 diffs.push(PresenceDiff {
321 room: room_name.clone(),
322 joins: vec![],
323 leaves: evicted,
324 });
325 }
326
327 if room.members.is_empty() {
328 empty_rooms.push(room_name.clone());
329 }
330 }
331
332 for room_name in empty_rooms {
333 rooms.remove(&room_name);
334 }
335
336 diffs
337 }
338
339 pub async fn get_room(&self, room: &str) -> Option<PresenceState> {
341 let rooms = self.rooms.read().await;
342 let presence_room = rooms.get(room)?;
343
344 Some(PresenceState {
345 room: room.to_string(),
346 members: presence_room.members.values().cloned().collect(),
347 })
348 }
349
350 pub async fn stats(&self) -> PresenceStats {
352 let rooms = self.rooms.read().await;
353 let total_members: usize = rooms.values().map(|r| r.members.len()).sum();
354
355 PresenceStats {
356 active_rooms: rooms.len(),
357 total_members,
358 joins_total: self.joins_total.load(Ordering::Relaxed),
359 leaves_total: self.leaves_total.load(Ordering::Relaxed),
360 evictions_total: self.evictions_total.load(Ordering::Relaxed),
361 }
362 }
363}
364
365#[derive(Debug, thiserror::Error)]
367pub enum PresenceError {
368 #[error("room '{room}' is full: max {max} members")]
370 RoomFull {
371 room: String,
373 max: usize,
375 },
376
377 #[error("room limit exceeded: max {max} rooms")]
379 TooManyRooms {
380 max: usize,
382 },
383}