Skip to main content

o_sfu_core/engine/room/
directory.rs

1//! Current-room indexes and lifecycle leases.
2//!
3//! [`RoomDirectory`] indexes one current room by UUID, issuer and instance ID.
4//! Each entry shares a [`RoomLifecycle`] gate. Accepted leases defer empty-room
5//! removal until the final mutation finishes. Reservation expiry removes only
6//! idle entries claimed by that gate.
7
8#[cfg(test)]
9use std::thread;
10use std::{
11    collections::BTreeMap,
12    sync::{Arc, Mutex},
13    time::Duration,
14};
15
16use time::{OffsetDateTime, format_description::well_known::Rfc3339};
17use tokio::time::Instant;
18use tracing::error;
19
20use super::Room;
21use crate::engine::{RoomInstanceId, sync::lock_unpoisoned};
22
23const UNKNOWN_REMOTE_ADDRESS: &str = "unknown";
24
25fn rfc3339_now() -> String {
26    match OffsetDateTime::now_utc().format(&Rfc3339) {
27        Ok(timestamp) => timestamp,
28        Err(_error) => String::from("1970-01-01T00:00:00Z"),
29    }
30}
31
32/// directory row for one current room instance
33///
34/// Cloned entries share the lifecycle gate. Manager admission holds the directory
35/// read guard until that gate accepts a lease, then releases it before room work.
36#[derive(Debug, Clone)]
37pub(crate) struct RoomDirectoryEntry {
38    pub room: Arc<Room>,
39    pub lifecycle: RoomLifecycle,
40    pub create_date: String,
41    pub remote_address: String,
42}
43
44impl RoomDirectoryEntry {
45    fn new(
46        room: Arc<Room>,
47        remote_address: Option<&str>,
48        reservation_ttl: Duration,
49        departure_grace: Duration,
50    ) -> Self {
51        Self {
52            room,
53            lifecycle: RoomLifecycle::new(reservation_ttl, departure_grace),
54            create_date: rfc3339_now(),
55            remote_address: remote_address.unwrap_or(UNKNOWN_REMOTE_ADDRESS).to_owned(),
56        }
57    }
58}
59
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61enum RoomLifecyclePhase {
62    Reservation { expires_at: Instant },
63    Alive,
64    Grace { expires_at: Instant },
65    Closing,
66}
67
68/// mutable state behind one directory entry's lifecycle lease gate
69///
70/// this lock is synchronous and short lived
71/// callers may hold a
72/// [`RoomLifecycleLease`] while awaiting, but this mutex is only held while a
73/// lease is accepted or released
74#[derive(Debug)]
75struct RoomLifecycleState {
76    /// accepted room work that has not finished or been dropped
77    active_mutations: usize,
78    /// empty-room removal request waiting for accepted work to drain
79    pending_removal: Option<RoomRemovalPolicy>,
80    /// lease length this reservation was published with and is renewed by
81    reservation_ttl: Duration,
82    /// grace duration this reservation was published with and is renewed by
83    departure_grace: Duration,
84    phase: RoomLifecyclePhase,
85}
86
87impl RoomLifecycleState {
88    fn new(reservation_ttl: Duration, departure_grace: Duration) -> Self {
89        Self {
90            active_mutations: 0,
91            pending_removal: None,
92            reservation_ttl,
93            departure_grace,
94            phase: RoomLifecyclePhase::Reservation {
95                expires_at: Instant::now() + reservation_ttl,
96            },
97        }
98    }
99}
100
101#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub(crate) enum ExpiryReason {
103    /// reservation lapsed before any join succeeded
104    ReservationLapsed,
105    /// room went empty and the departure grace ran out
106    GraceElapsed,
107}
108
109/// cloneable admission gate for the current room stored in one directory row
110///
111/// this type coordinates manager-level liveness only
112/// room membership ordering
113/// remains owned by [`Room`] and its state transition methods
114#[derive(Debug, Clone)]
115pub(crate) struct RoomLifecycle {
116    state: Arc<Mutex<RoomLifecycleState>>,
117}
118
119impl RoomLifecycle {
120    pub(crate) fn new(reservation_ttl: Duration, departure_grace: Duration) -> Self {
121        Self {
122            state: Arc::new(Mutex::new(RoomLifecycleState::new(
123                reservation_ttl,
124                departure_grace,
125            ))),
126        }
127    }
128
129    /// atomically claims cleanup responsibility for an expired grace/reservation period
130    pub(crate) fn claim_expired_room(&self) -> Option<ExpiryReason> {
131        let mut state = lock_unpoisoned(&self.state);
132        let expiry_reason = match state.phase {
133            RoomLifecyclePhase::Reservation { expires_at } if expires_at <= Instant::now() => {
134                Some(ExpiryReason::ReservationLapsed)
135            }
136            RoomLifecyclePhase::Grace { expires_at } if expires_at <= Instant::now() => {
137                Some(ExpiryReason::GraceElapsed)
138            }
139            _ => None,
140        };
141        if state.active_mutations == 0 && expiry_reason.is_some() {
142            state.phase = RoomLifecyclePhase::Closing;
143            drop(state);
144            return expiry_reason;
145        }
146        None
147    }
148
149    /// extends a room reservation; rearms it if coming from a grace period
150    pub(crate) fn renew_reservation(&self) {
151        let mut state = lock_unpoisoned(&self.state);
152        if matches!(
153            state.phase,
154            RoomLifecyclePhase::Reservation { .. } | RoomLifecyclePhase::Grace { .. }
155        ) {
156            state.phase = RoomLifecyclePhase::Reservation {
157                expires_at: Instant::now() + state.reservation_ttl,
158            }
159        }
160    }
161
162    #[cfg(any(test, feature = "testing-transport"))]
163    pub(crate) fn expire_reservation_now_for_test(&self) {
164        lock_unpoisoned(&self.state).phase = RoomLifecyclePhase::Reservation {
165            expires_at: Instant::now(),
166        };
167    }
168
169    #[cfg(any(test, feature = "testing-transport"))]
170    #[must_use]
171    pub(crate) fn has_reservation_deadline_for_test(&self) -> bool {
172        matches!(
173            lock_unpoisoned(&self.state).phase,
174            RoomLifecyclePhase::Reservation { .. }
175        )
176    }
177
178    #[cfg(any(test, feature = "testing-transport"))]
179    #[must_use]
180    pub(crate) fn has_departure_grace_for_test(&self) -> bool {
181        matches!(
182            lock_unpoisoned(&self.state).phase,
183            RoomLifecyclePhase::Grace { .. }
184        )
185    }
186
187    /// expires an armed departure grace, and reports whether one was armed
188    ///
189    /// this never arms a grace, so a test cannot expire a phase the room never
190    /// reached
191    #[cfg(any(test, feature = "testing-transport"))]
192    #[must_use]
193    pub(crate) fn expire_departure_grace_now_for_test(&self) -> bool {
194        let mut state = lock_unpoisoned(&self.state);
195        if !matches!(state.phase, RoomLifecyclePhase::Grace { .. }) {
196            return false;
197        }
198        state.phase = RoomLifecyclePhase::Grace {
199            expires_at: Instant::now(),
200        };
201        true
202    }
203
204    /// Accepts a lease unless immediate removal is pending or already claimed.
205    ///
206    /// Returns `None` when removal is pending or claimed or the lease count
207    /// cannot increase. Directory callers must retain their read guard through
208    /// admission to prove that the entry is current.
209    #[must_use]
210    pub(crate) fn begin(&self) -> Option<RoomLifecycleLease> {
211        let mut state = lock_unpoisoned(&self.state);
212        if matches!(state.phase, RoomLifecyclePhase::Closing)
213            || matches!(state.pending_removal, Some(RoomRemovalPolicy::Immediately))
214        {
215            return None;
216        }
217        state.active_mutations = state.active_mutations.checked_add(1)?;
218        let lease = RoomLifecycleLease {
219            state: Arc::clone(&self.state),
220            finished: false,
221        };
222        drop(state);
223        Some(lease)
224    }
225}
226
227#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
228pub(crate) enum RoomRemovalPolicy {
229    AfterGrace,
230    Immediately,
231}
232
233/// cancellation-safe permit for work accepted against a directory entry
234///
235/// dropping the lease releases admission without requesting removal
236/// manager
237/// teardown paths call [`Self::finish`] after checking whether the room is empty
238#[derive(Debug)]
239pub(crate) struct RoomLifecycleLease {
240    /// shared lease state for the directory entry that accepted this work
241    state: Arc<Mutex<RoomLifecycleState>>,
242    /// Prevents `Drop` from releasing a lease already completed by `finish`.
243    finished: bool,
244}
245
246impl RoomLifecycleLease {
247    /// Releases the lease and returns whether this caller claimed directory removal.
248    #[must_use]
249    pub(crate) fn finish(
250        mut self,
251        on_empty: Option<RoomRemovalPolicy>,
252        room_can_be_removed: bool,
253    ) -> bool {
254        self.release(on_empty, room_can_be_removed)
255    }
256
257    #[must_use]
258    fn release(&mut self, on_empty: Option<RoomRemovalPolicy>, room_can_be_removed: bool) -> bool {
259        if self.finished {
260            return false;
261        }
262        self.finished = true;
263        let mut state = lock_unpoisoned(&self.state);
264        // Keep the strongest policy; promote 0 seconds grace to immediate removal.
265        let removal_policy = match state.pending_removal.max(on_empty) {
266            Some(RoomRemovalPolicy::AfterGrace) if state.departure_grace.is_zero() => {
267                Some(RoomRemovalPolicy::Immediately)
268            }
269            policy => policy,
270        };
271
272        state.active_mutations = state.active_mutations.saturating_sub(1);
273        if state.active_mutations > 0 {
274            if on_empty.is_some() && room_can_be_removed {
275                state.pending_removal = removal_policy;
276            }
277            return false;
278        }
279        // The last lease consumes the policy below or voids it here, and a
280        // `None` policy means nothing was pending.
281        state.pending_removal = None;
282        if !room_can_be_removed {
283            return false;
284        }
285
286        match removal_policy {
287            Some(RoomRemovalPolicy::Immediately) => {
288                state.phase = RoomLifecyclePhase::Closing;
289                true
290            }
291            Some(RoomRemovalPolicy::AfterGrace) => {
292                if matches!(state.phase, RoomLifecyclePhase::Alive) {
293                    state.phase = RoomLifecyclePhase::Grace {
294                        expires_at: Instant::now() + state.departure_grace,
295                    };
296                }
297                false
298            }
299            None => false,
300        }
301    }
302
303    pub(crate) fn clear_expiration(&self) {
304        lock_unpoisoned(&self.state).phase = RoomLifecyclePhase::Alive;
305    }
306}
307
308impl Drop for RoomLifecycleLease {
309    fn drop(&mut self) {
310        if self.finished {
311            return;
312        }
313        error!("room lifecycle lease dropped without completing explicit cleanup");
314        #[cfg(test)]
315        assert!(
316            thread::panicking(),
317            "room lifecycle lease dropped without completing explicit cleanup"
318        );
319        let _ = self.release(None, false);
320    }
321}
322
323#[derive(Debug, Default)]
324pub(crate) struct RoomDirectory {
325    by_uuid: BTreeMap<String, RoomDirectoryEntry>,
326    uuid_by_instance: BTreeMap<RoomInstanceId, String>,
327    uuid_by_issuer: BTreeMap<String, String>,
328}
329
330impl RoomDirectory {
331    #[must_use]
332    pub(crate) fn get_by_uuid(&self, uuid: &str) -> Option<Arc<Room>> {
333        self.by_uuid.get(uuid).map(|entry| Arc::clone(&entry.room))
334    }
335
336    #[must_use]
337    pub(crate) fn entry(&self, uuid: &str) -> Option<&RoomDirectoryEntry> {
338        self.by_uuid.get(uuid)
339    }
340
341    #[must_use]
342    pub(crate) fn entry_by_issuer(&self, issuer: &str) -> Option<&RoomDirectoryEntry> {
343        let uuid = self.uuid_by_issuer.get(issuer)?;
344        self.by_uuid.get(uuid)
345    }
346
347    #[must_use]
348    pub(crate) fn get_by_instance_id(&self, room_instance_id: RoomInstanceId) -> Option<Arc<Room>> {
349        let uuid = self.uuid_by_instance.get(&room_instance_id)?;
350        self.get_by_uuid(uuid)
351    }
352
353    #[must_use]
354    pub(crate) fn entries(&self) -> Vec<RoomDirectoryEntry> {
355        self.by_uuid.values().cloned().collect()
356    }
357
358    #[must_use]
359    pub(crate) fn rooms(&self) -> Vec<Arc<Room>> {
360        self.by_uuid
361            .values()
362            .map(|entry| Arc::clone(&entry.room))
363            .collect()
364    }
365
366    pub(crate) fn insert(
367        &mut self,
368        room: Arc<Room>,
369        remote_address: Option<&str>,
370        reservation_ttl: Duration,
371        departure_grace: Duration,
372    ) {
373        let room_id = room.uuid().to_owned();
374        self.uuid_by_issuer
375            .insert(room.issuer().to_owned(), room_id.clone());
376        self.uuid_by_instance
377            .insert(room.instance_id(), room_id.clone());
378        self.by_uuid.insert(
379            room_id,
380            RoomDirectoryEntry::new(room, remote_address, reservation_ttl, departure_grace),
381        );
382    }
383
384    #[must_use]
385    pub(crate) fn contains_current(&self, uuid: &str, room: &Arc<Room>) -> bool {
386        self.by_uuid
387            .get(uuid)
388            .is_some_and(|entry| Arc::ptr_eq(&entry.room, room))
389    }
390
391    pub(crate) fn remove_if_current(&mut self, uuid: &str, room: &Arc<Room>) -> bool {
392        if self.contains_current(uuid, room) {
393            self.by_uuid.remove(uuid);
394            self.uuid_by_issuer.remove(room.issuer());
395            self.uuid_by_instance.remove(&room.instance_id());
396            return true;
397        }
398        false
399    }
400}