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#[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 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 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 #[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}