Skip to main content

fraiseql_server/subscriptions/
presence.rs

1//! Room-based presence tracking for realtime member awareness.
2//!
3//! Clients join a "room" with an initial state payload.  The server tracks
4//! membership and emits `PRESENCE_STATE` (full roster on join) and
5//! `PRESENCE_DIFF` (join/leave/update deltas) events.
6//!
7//! All state is in-memory — lost on server restart (acceptable for v1).
8
9use 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// ──────────────────────── Configuration ────────────────────────
20
21/// Configuration for the presence subsystem.
22#[derive(Debug, Clone)]
23pub struct PresenceConfig {
24    /// Maximum members per room (prevents memory abuse).
25    pub max_members_per_room: usize,
26
27    /// Maximum number of rooms that can exist simultaneously.
28    pub max_rooms: usize,
29
30    /// Heartbeat timeout — members are evicted after this duration without a ping.
31    pub heartbeat_timeout: Duration,
32}
33
34impl PresenceConfig {
35    /// Create config with production defaults.
36    #[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// ──────────────────────── Types ────────────────────────
53
54/// A member's presence in a room.
55#[derive(Debug, Clone, Serialize, Deserialize)]
56pub struct PresenceMember {
57    /// Unique identifier for this member (typically connection ID or user ID).
58    pub id: String,
59
60    /// Arbitrary JSON state (e.g., cursor position, status, avatar).
61    pub state: serde_json::Value,
62
63    /// When this member last sent a heartbeat.
64    #[serde(skip, default = "Instant::now")]
65    pub last_seen: Instant,
66}
67
68/// Full room state — sent to a client on join (`PRESENCE_STATE`).
69#[derive(Debug, Clone, Serialize)]
70pub struct PresenceState {
71    /// Room name.
72    pub room: String,
73
74    /// All current members.
75    pub members: Vec<PresenceMember>,
76}
77
78/// Delta event — sent when members join, leave, or update (`PRESENCE_DIFF`).
79#[derive(Debug, Clone, Serialize)]
80pub struct PresenceDiff {
81    /// Room name.
82    pub room: String,
83
84    /// Members who joined.
85    pub joins: Vec<PresenceMember>,
86
87    /// Member IDs who left (or were evicted).
88    pub leaves: Vec<String>,
89}
90
91// ──────────────────────── Room ────────────────────────
92
93/// A single presence room.
94#[derive(Debug)]
95struct PresenceRoom {
96    /// Members indexed by their ID.
97    members: HashMap<String, PresenceMember>,
98}
99
100impl PresenceRoom {
101    fn new() -> Self {
102        Self {
103            members: HashMap::new(),
104        }
105    }
106}
107
108// ──────────────────────── Manager ────────────────────────
109
110/// Statistics for the presence subsystem.
111#[derive(Debug, Clone)]
112pub struct PresenceStats {
113    /// Total rooms currently tracked.
114    pub active_rooms: usize,
115
116    /// Total members across all rooms.
117    pub total_members: usize,
118
119    /// Total join events processed.
120    pub joins_total: u64,
121
122    /// Total leave events processed.
123    pub leaves_total: u64,
124
125    /// Total heartbeat evictions.
126    pub evictions_total: u64,
127}
128
129/// Manages room-based presence state.
130///
131/// Thread-safe via `RwLock` for the room map and atomics for counters.
132#[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    /// Create a new presence manager.
143    #[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    /// Join a room with initial state.
155    ///
156    /// Returns `PresenceState` (current members including the new one) and
157    /// a `PresenceDiff` (announcing the join) for broadcasting.
158    ///
159    /// # Errors
160    ///
161    /// Returns error if the room is full or room limit is exceeded.
162    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        // Create room if needed
171        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            // Unreachable: we just inserted the room above if it was missing.
182            return Err(PresenceError::TooManyRooms {
183                max: self.config.max_rooms,
184            });
185        };
186
187        // Check room capacity (only if this is a new member, not a rejoin)
188        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    /// Leave a room.
227    ///
228    /// Returns a `PresenceDiff` for broadcasting, or `None` if the member
229    /// wasn't in the room.
230    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        // Clean up empty rooms
246        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    /// Record a heartbeat for a member, resetting their eviction timer.
255    ///
256    /// Returns `true` if the heartbeat was accepted (member exists in room).
257    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    /// Update a member's state payload.
271    ///
272    /// Returns a `PresenceDiff` with the updated member in `joins` (same
273    /// semantics as Supabase Realtime — updates appear as joins).
274    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    /// Evict members whose heartbeat has expired.
295    ///
296    /// Returns `PresenceDiff` events for each room that had evictions.
297    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    /// Get current members of a room.
340    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    /// Get statistics.
351    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/// Errors from presence operations.
366#[derive(Debug, thiserror::Error)]
367pub enum PresenceError {
368    /// Room has reached its member cap.
369    #[error("room '{room}' is full: max {max} members")]
370    RoomFull {
371        /// Room name.
372        room: String,
373        /// Maximum members.
374        max:  usize,
375    },
376
377    /// Too many rooms exist.
378    #[error("room limit exceeded: max {max} rooms")]
379    TooManyRooms {
380        /// Maximum rooms.
381        max: usize,
382    },
383}