Skip to main content

o_sfu_core/engine/room/
instance.rs

1use std::{fmt, sync::Arc};
2
3use tokio::sync::{Mutex as AsyncMutex, MutexGuard, RwLock};
4
5use super::{definition::RoomDefinition, factory::RoomInit, state::RoomState};
6use crate::{
7    RoomWorkerPolicy,
8    engine::{
9        AvailableFeatures, ConnectionId, MediaWorkerId, PeerSnapshot, RecordingState,
10        RoomInstanceId, UserId,
11        media_transport::{MediaTransport, TransportSessionKey},
12        metrics::RuntimeMetrics,
13    },
14};
15
16#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
17pub enum RoomJoinError {
18    #[error("room is full")]
19    RoomFull,
20    #[error("no usable media worker")]
21    NoUsableWorker,
22    #[error("router state error")]
23    RouterState,
24}
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
27pub enum RoomManagerJoinError {
28    #[error("room not found")]
29    MissingRoom,
30    #[error("room is full")]
31    RoomFull,
32    #[error("no usable media worker")]
33    NoUsableWorker,
34    #[error("router state error")]
35    RouterState,
36}
37
38#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
39pub enum RoomManagerServeError {
40    #[error("conflicting room reservation")]
41    ConflictingReservation,
42}
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
45pub struct RoomMediaCounts {
46    pub publications: usize,
47    pub subscriptions: usize,
48}
49
50#[derive(Clone, Copy)]
51pub(crate) struct RoomUserOperation<'a> {
52    pub room: &'a Room,
53    pub user_id: &'a UserId,
54    pub connection_id: ConnectionId,
55    pub media_transport: &'a MediaTransport,
56}
57
58/// Holds publication and source-policy ordering for one room.
59///
60/// Guarded effects derive their room from this guard so another room's lock
61/// cannot authorize them. Keep the guard across each state commit and its
62/// effects, including every publication resolved by one answer.
63#[must_use = "dropping this guard releases source-policy ordering"]
64pub(super) struct SourcePolicyGuard<'a> {
65    room: &'a Room,
66    _lock: MutexGuard<'a, ()>,
67}
68
69impl SourcePolicyGuard<'_> {
70    pub const fn room(&self) -> &Room {
71        self.room
72    }
73}
74
75pub struct Room {
76    pub(super) definition: RoomDefinition,
77    pub(super) metrics: Arc<RuntimeMetrics>,
78    /// Serializes publication commits and activity changes with source-policy
79    /// turns so their transport effects cannot overtake each other.
80    source_policy_turn: AsyncMutex<()>,
81    pub(super) state: RwLock<RoomState>,
82}
83
84impl Room {
85    pub(super) fn new(init: RoomInit) -> Self {
86        let RoomInit {
87            runtime_context,
88            runtime_policy,
89            issuer,
90            key,
91            config,
92            metrics,
93        } = init;
94        let definition =
95            RoomDefinition::new(&runtime_context, &runtime_policy, issuer, key, config);
96        Self {
97            definition,
98            metrics,
99            source_policy_turn: AsyncMutex::new(()),
100            state: RwLock::new(RoomState::new(
101                &runtime_context,
102                runtime_policy.admission_policy,
103                runtime_policy.media_limits,
104                runtime_policy.video_adaptation_tuning,
105                runtime_policy.router_rtp_capabilities,
106            )),
107        }
108    }
109
110    pub(crate) fn user_operation<'a>(
111        &'a self,
112        user_id: &'a UserId,
113        connection_id: ConnectionId,
114        media_transport: &'a MediaTransport,
115    ) -> RoomUserOperation<'a> {
116        RoomUserOperation {
117            room: self,
118            user_id,
119            connection_id,
120            media_transport,
121        }
122    }
123
124    /// Serializes publication commits and effects with source-policy turns.
125    ///
126    /// Call guarded operations while this guard is held. Acquiring another
127    /// source-policy guard for this room before releasing it would deadlock.
128    pub(super) async fn lock_source_policy(&self) -> SourcePolicyGuard<'_> {
129        SourcePolicyGuard {
130            room: self,
131            _lock: self.source_policy_turn.lock().await,
132        }
133    }
134
135    #[must_use]
136    pub fn uuid(&self) -> &str {
137        self.definition.uuid()
138    }
139
140    #[must_use]
141    pub fn issuer(&self) -> &str {
142        self.definition.issuer()
143    }
144
145    #[must_use]
146    pub fn key(&self) -> &secrecy::SecretSlice<u8> {
147        self.definition.key()
148    }
149
150    #[must_use]
151    pub(crate) fn available_features(&self) -> AvailableFeatures {
152        self.definition.available_features()
153    }
154
155    pub async fn recording_state(&self) -> RecordingState {
156        self.state.read().await.recording_state()
157    }
158
159    pub(crate) async fn user_snapshots_except(
160        &self,
161        excluded_user_id: &UserId,
162    ) -> Vec<PeerSnapshot> {
163        self.state
164            .read()
165            .await
166            .user_snapshots_except(excluded_user_id)
167    }
168
169    #[must_use]
170    pub fn web_rtc_enabled(&self) -> bool {
171        self.definition.web_rtc_enabled()
172    }
173
174    /// Returns the transport key for an exact committed router placement.
175    ///
176    /// # Panics
177    ///
178    /// Panics when `user_id` and `connection_id` have no committed router
179    /// placement.
180    #[must_use]
181    pub async fn transport_user_key(
182        &self,
183        user_id: &UserId,
184        connection_id: ConnectionId,
185    ) -> TransportSessionKey {
186        self.state
187            .read()
188            .await
189            .transport_user_key(user_id, connection_id)
190    }
191
192    pub fn room_worker_policy(&self) -> RoomWorkerPolicy {
193        self.definition.room_worker_policy()
194    }
195
196    #[must_use]
197    pub(crate) fn instance_id(&self) -> RoomInstanceId {
198        self.definition.instance_id()
199    }
200}
201
202impl fmt::Debug for Room {
203    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
204        let media_worker_id = self
205            .state
206            .try_read()
207            .ok()
208            .and_then(|state| state.assigned_primary_media_worker_id())
209            .map(MediaWorkerId::as_usize);
210        formatter
211            .debug_struct("Room")
212            .field("instance_id", &self.definition.instance_id())
213            .field("media_worker_id", &media_worker_id)
214            .field("uuid", &self.definition.uuid())
215            .field("issuer", &self.definition.issuer())
216            .field("web_rtc_enabled", &self.definition.web_rtc_enabled())
217            .finish_non_exhaustive()
218    }
219}