1use std::{collections::BTreeMap, mem::replace, sync::Arc};
20
21pub use crate::engine::media_transport::{
22 SessionOffer as NegotiationOffer, SessionUploadEncoding as UploadEncoding,
23 SessionUploadSlot as UploadSlot,
24};
25use crate::{
26 ConnectionId,
27 engine::{
28 AvailableFeatures, JsonPayload, PeerSnapshot, RecordingOptions, RecordingState, UserId,
29 UserInfo,
30 media_transport::{
31 MediaTransport, TransportAdapterError, TransportSessionHealth, TransportSessionKey,
32 },
33 room::{
34 BroadcastPayloadError, DeactivateIntentOutcome, JoinUserRequest, PublishIntentOutcome,
35 Room, RoomManager, RoomManagerJoinError, RoomUserOperation,
36 },
37 source_model::{
38 SourceDeactivateIntent, SourcePublishIntent, SourceSubscriptionIntent, UserStreamId,
39 },
40 },
41};
42
43#[derive(Debug, Default)]
49enum SessionPhase {
50 #[default]
52 BeforeInitialOffer,
53 Stable,
55 WaitingForAnswer(InFlightOffer),
57}
58
59#[derive(Debug)]
61struct InFlightOffer {
62 purpose: SessionOfferPurpose,
63 queued_publishes: BTreeMap<UserStreamId, SourcePublishIntent>,
64 follow_up_renegotiation: bool,
65}
66
67impl SessionPhase {
68 fn can_stage_publish(&self) -> bool {
69 !matches!(self, Self::WaitingForAnswer(_))
70 }
71
72 fn has_queued_publish(&self, stream_id: &UserStreamId) -> bool {
73 matches!(
74 self,
75 Self::WaitingForAnswer(pending)
76 if pending.queued_publishes.contains_key(stream_id)
77 )
78 }
79
80 fn queue_publish(&mut self, intent: SourcePublishIntent) {
81 if let Self::WaitingForAnswer(pending) = self {
82 let stream_id = intent.stream_id().clone();
83 pending.queued_publishes.insert(stream_id, intent);
84 }
85 }
86
87 fn remove_queued_publish(&mut self, stream_id: &UserStreamId) -> bool {
88 let Self::WaitingForAnswer(pending) = self else {
89 return false;
90 };
91 pending.queued_publishes.remove(stream_id).is_some()
92 }
93
94 fn clear_queued_publishes(&mut self) {
95 if let Self::WaitingForAnswer(pending) = self {
96 pending.queued_publishes.clear();
97 }
98 }
99
100 fn request_renegotiation(&mut self) -> bool {
101 match self {
102 Self::BeforeInitialOffer => false,
103 Self::Stable => true,
104 Self::WaitingForAnswer(pending) => {
105 pending.follow_up_renegotiation = true;
106 false
107 }
108 }
109 }
110
111 fn mark_follow_up_renegotiation(&mut self) {
112 if let Self::WaitingForAnswer(pending) = self {
113 pending.follow_up_renegotiation = true;
114 }
115 }
116
117 fn wait_for_answer(&mut self, purpose: SessionOfferPurpose) {
118 *self = Self::WaitingForAnswer(InFlightOffer {
119 purpose,
120 queued_publishes: BTreeMap::new(),
121 follow_up_renegotiation: false,
122 });
123 }
124
125 #[expect(
126 clippy::unreachable,
127 reason = "answer validates the phase before awaiting with exclusive session access"
128 )]
129 fn complete_answer(&mut self) -> InFlightOffer {
130 match replace(self, Self::Stable) {
131 Self::WaitingForAnswer(pending) => pending,
132 _ => unreachable!("answer completion requires an in-flight offer"),
133 }
134 }
135}
136
137#[derive(Debug)]
139enum SessionOfferPurpose {
140 EstablishSession,
141 RefreshSession,
142}
143
144#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
146pub enum SessionError {
147 #[error("no pending media request")]
149 NoPendingRequest,
150 #[error(transparent)]
152 Core(#[from] SfuCoreError),
153}
154
155impl SessionError {
156 #[must_use]
161 pub const fn is_client_error(self) -> bool {
162 match self {
163 Self::NoPendingRequest => true,
164 Self::Core(error) => error.is_client_error(),
165 }
166 }
167}
168
169#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
170pub enum SfuCoreError {
171 #[error("transport operation failed")]
173 Transport(#[source] TransportAdapterError),
174 #[error("capability projection failed")]
176 CapabilityProjection(#[source] TransportAdapterError),
177 #[error("session negotiation rejected")]
179 SessionNegotiationRejected,
180 #[error("session refresh rejected")]
182 SessionRefreshRejected,
183 #[error("subscription update rejected")]
185 SubscriptionUpdateRejected,
186}
187
188impl SfuCoreError {
189 #[must_use]
191 pub const fn is_client_error(self) -> bool {
192 matches!(
193 self,
194 Self::Transport(TransportAdapterError::InvalidInput)
195 | Self::CapabilityProjection(_)
196 | Self::SessionNegotiationRejected
197 | Self::SessionRefreshRejected
198 | Self::SubscriptionUpdateRejected
199 )
200 }
201}
202
203#[derive(Debug, Clone)]
205pub struct SfuCore {
206 media_transport: MediaTransport,
207 rooms: Arc<RoomManager>,
208}
209
210#[derive(Debug)]
229pub struct MediaSession {
230 core: SfuCore,
231 room: Arc<Room>,
232 transport_user_key: TransportSessionKey,
233 phase: SessionPhase,
234 closed: bool,
235}
236
237impl SfuCore {
238 #[must_use]
239 pub fn new(media_transport: MediaTransport, rooms: Arc<RoomManager>) -> Self {
240 Self {
241 media_transport,
242 rooms,
243 }
244 }
245
246 pub async fn admit_user(
267 &self,
268 room_id: &str,
269 request: JoinUserRequest,
270 ) -> Result<MediaSession, RoomManagerJoinError> {
271 let admission = self
272 .rooms
273 .join_user(room_id, request, &self.media_transport)
274 .await?;
275 Ok(MediaSession {
276 core: self.clone(),
277 room: admission.room,
278 transport_user_key: admission.transport_session_key,
279 phase: SessionPhase::default(),
280 closed: false,
281 })
282 }
283}
284
285impl MediaSession {
286 pub async fn establish(&mut self) -> Result<Option<NegotiationOffer>, SessionError> {
297 if !matches!(self.phase, SessionPhase::BeforeInitialOffer) {
298 return Ok(None);
299 }
300 let offer = self
301 .core
302 .media_transport
303 .create_initial_session_offer(self.room.uuid(), &self.transport_user_key)
304 .await
305 .map_err(SfuCoreError::Transport)?;
306 self.phase
307 .wait_for_answer(SessionOfferPurpose::EstablishSession);
308 Ok(Some(offer))
309 }
310
311 pub async fn answer(&mut self, sdp: &str) -> Result<Option<NegotiationOffer>, SessionError> {
326 if !matches!(self.phase, SessionPhase::WaitingForAnswer(_)) {
327 return Err(SessionError::NoPendingRequest);
328 }
329 let applied_answer = self
330 .core
331 .media_transport
332 .apply_session_answer(&self.transport_user_key, sdp)
333 .await
334 .map_err(SfuCoreError::Transport)?;
335 let InFlightOffer {
336 purpose,
337 queued_publishes,
338 follow_up_renegotiation,
339 } = self.phase.complete_answer();
340 match purpose {
341 SessionOfferPurpose::EstablishSession => {
342 let client_capabilities = applied_answer.client_capabilities().cloned().ok_or(
343 SfuCoreError::CapabilityProjection(TransportAdapterError::InvalidInput),
344 )?;
345 self.room_operation()
346 .apply_session_negotiated(
347 client_capabilities,
348 applied_answer.declined_consumers(),
349 )
350 .await
351 .ok_or(SfuCoreError::SessionNegotiationRejected)?;
352 }
353 SessionOfferPurpose::RefreshSession => {
354 self.room_operation()
355 .apply_session_refreshed(applied_answer.declined_consumers())
356 .await
357 .ok_or(SfuCoreError::SessionRefreshRejected)?;
358 }
359 }
360 self.room_operation()
361 .commit_staged_publishes(&applied_answer)
362 .await;
363 let staged = self.stage_queued_publishes(queued_publishes).await?;
364 if staged || follow_up_renegotiation {
365 return self.renegotiate().await;
366 }
367 Ok(None)
368 }
369
370 pub async fn publish(
382 &mut self,
383 intent: SourcePublishIntent,
384 ) -> Result<Option<NegotiationOffer>, SessionError> {
385 if self.phase.has_queued_publish(intent.stream_id()) {
386 return Ok(None);
387 }
388 match self
389 .start_publish(&intent, self.phase.can_stage_publish())
390 .await?
391 {
392 PublishIntentOutcome::Noop | PublishIntentOutcome::Activated => Ok(None),
393 PublishIntentOutcome::Queue => {
394 self.phase.queue_publish(intent);
395 Ok(None)
396 }
397 PublishIntentOutcome::Staged => self.renegotiate().await,
398 }
399 }
400
401 pub async fn deactivate_publication(&mut self, intent: SourceDeactivateIntent) {
409 if self.phase.remove_queued_publish(intent.stream_id()) {
410 return;
411 }
412 match self.room_operation().deactivate_publication(&intent).await {
413 DeactivateIntentOutcome::RolledBack => {
414 self.phase.mark_follow_up_renegotiation();
415 }
416 DeactivateIntentOutcome::Deactivated | DeactivateIntentOutcome::Noop => {}
417 }
418 }
419
420 pub async fn close(&mut self) -> bool {
429 if self.closed {
430 return false;
431 }
432 self.phase.clear_queued_publishes();
433 let did_close = self
434 .core
435 .rooms
436 .close_session(
437 self.room_id(),
438 self.user_id(),
439 self.connection_id(),
440 &self.core.media_transport,
441 )
442 .await;
443 self.closed = true;
444 did_close
445 }
446
447 pub async fn renegotiate(&mut self) -> Result<Option<NegotiationOffer>, SessionError> {
459 if !self.phase.request_renegotiation() {
460 return Ok(None);
461 }
462 let offer = match self
463 .core
464 .media_transport
465 .create_session_renegotiation_offer(&self.transport_user_key)
466 .await
467 {
468 Ok(offer) => offer,
469 Err(TransportAdapterError::UnsupportedFeature) => return Ok(None),
470 Err(error) => return Err(SfuCoreError::Transport(error).into()),
471 };
472 self.phase
473 .wait_for_answer(SessionOfferPurpose::RefreshSession);
474 Ok(Some(offer))
475 }
476
477 pub async fn subscribe(
489 &self,
490 target_user_id: &UserId,
491 intents: &BTreeMap<UserStreamId, SourceSubscriptionIntent>,
492 ) -> Result<(), SessionError> {
493 self.room_operation()
494 .apply_receiver_intent(target_user_id, intents)
495 .await
496 .ok_or(SfuCoreError::SubscriptionUpdateRejected)?;
497 Ok(())
498 }
499
500 #[must_use]
502 pub fn endpoint_health(&self) -> Option<TransportSessionHealth> {
503 self.core
504 .media_transport
505 .session_transport_health(&self.transport_user_key)
506 }
507
508 #[must_use]
509 pub fn user_id(&self) -> &UserId {
510 self.transport_user_key.user_id()
511 }
512
513 #[must_use]
514 pub const fn connection_id(&self) -> ConnectionId {
515 self.transport_user_key.connection_id()
516 }
517
518 #[must_use]
519 pub fn room_id(&self) -> &str {
520 self.room.uuid()
521 }
522
523 pub async fn is_current_connection(&self) -> bool {
524 self.room
525 .has_connection(self.user_id(), self.connection_id())
526 .await
527 }
528
529 #[must_use]
530 pub fn available_features(&self) -> AvailableFeatures {
531 self.room.available_features()
532 }
533
534 pub async fn recording_state(&self) -> RecordingState {
535 self.room.recording_state().await
536 }
537
538 pub async fn peer_snapshots(&self) -> Vec<PeerSnapshot> {
540 self.room.user_snapshots_except(self.user_id()).await
541 }
542
543 fn room_operation(&self) -> RoomUserOperation<'_> {
544 self.room.user_operation(
545 self.user_id(),
546 self.connection_id(),
547 &self.core.media_transport,
548 )
549 }
550
551 async fn start_publish(
552 &self,
553 intent: &SourcePublishIntent,
554 can_stage: bool,
555 ) -> Result<PublishIntentOutcome, SfuCoreError> {
556 self.room_operation()
557 .start_publish(intent, can_stage)
558 .await
559 .map_err(SfuCoreError::Transport)
560 }
561
562 pub async fn update_info(&self, info: UserInfo) {
568 self.room
569 .update_user_info(
570 self.user_id(),
571 self.connection_id(),
572 &self.core.media_transport,
573 info,
574 )
575 .await;
576 }
577
578 pub async fn broadcast(&self, message: JsonPayload) -> Result<(), BroadcastPayloadError> {
589 self.room
590 .broadcast(self.user_id(), self.connection_id(), message)
591 .await
592 }
593
594 #[must_use]
599 #[expect(
600 clippy::unused_async,
601 reason = "keeps the public MediaSession recording facade async while disabled recording is synchronous"
602 )]
603 pub async fn start_recording(&self, options: RecordingOptions) -> bool {
604 self.room
605 .apply_recording_start(self.user_id(), self.connection_id(), options)
606 }
607
608 #[must_use]
612 #[expect(
613 clippy::unused_async,
614 reason = "keeps the public MediaSession recording facade async while disabled recording is synchronous"
615 )]
616 pub async fn stop_recording(&self) -> bool {
617 self.room
618 .apply_recording_stop(self.user_id(), self.connection_id())
619 }
620
621 async fn stage_queued_publishes(
622 &self,
623 queued: BTreeMap<UserStreamId, SourcePublishIntent>,
624 ) -> Result<bool, SessionError> {
625 let mut staged = false;
626 for intent in queued.into_values() {
627 if self.start_publish(&intent, true).await? == PublishIntentOutcome::Staged {
628 staged = true;
629 }
630 }
631 Ok(staged)
632 }
633}