Skip to main content

o_sfu_core/
sfu.rs

1//! [`SfuCore::admit_user`] returns one [`MediaSession`] per admitted room connection.
2//!
3//! The session sequences offer/answer, publication, subscription, recording and
4//! cleanup without exposing room or transport internals.
5//!
6//! ```text
7//! SfuCore::admit_user -> MediaSession
8//!
9//! establish             -> offer -> answer
10//! publish               -> [offer -> answer]?
11//! deactivate_publication -> cancel pending or suppress committed source
12//! subscribe             -> persist intent and reconcile eligible routes
13//! close                 -> remove only the current connection
14//! ```
15//!
16//! `&mut MediaSession` serializes negotiation. [`MediaSession::publish`] queues
17//! the first intent for each stream while an offer awaits its answer. A successful
18//! [`MediaSession::answer`] may return the follow-up offer.
19use 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/// media-session negotiation phase
44///
45/// queued publishes live beside the offer state so a user action that arrives
46/// while a browser answer is pending cannot interleave with the in-flight SDP
47/// exchange
48#[derive(Debug, Default)]
49enum SessionPhase {
50    /// no initial offer has been created yet
51    #[default]
52    BeforeInitialOffer,
53    /// no offer is awaiting an answer
54    Stable,
55    /// one offer has been sent and must be answered before the next offer
56    WaitingForAnswer(InFlightOffer),
57}
58
59/// in-flight offer state plus mutations deferred until its answer is accepted
60#[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/// reason an offer is waiting for an answer
138#[derive(Debug)]
139enum SessionOfferPurpose {
140    EstablishSession,
141    RefreshSession,
142}
143
144/// error returned by [`MediaSession`] operations
145#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
146pub enum SessionError {
147    /// [`MediaSession::answer`] was called without an in-flight offer
148    #[error("no pending media request")]
149    NoPendingRequest,
150    /// lower core operation failed or rejected the request
151    #[error(transparent)]
152    Core(#[from] SfuCoreError),
153}
154
155impl SessionError {
156    /// whether the runtime should report the error as a client or protocol fault
157    ///
158    /// transport failures that indicate malformed client input are client
159    /// errors, while infrastructure failures are internal errors
160    #[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    /// media transport command failed
172    #[error("transport operation failed")]
173    Transport(#[source] TransportAdapterError),
174    /// accepted answer did not yield client capabilities needed by room state
175    #[error("capability projection failed")]
176    CapabilityProjection(#[source] TransportAdapterError),
177    /// initial answer was valid transport input but stale for room state
178    #[error("session negotiation rejected")]
179    SessionNegotiationRejected,
180    /// refresh answer was valid transport input but stale for room state
181    #[error("session refresh rejected")]
182    SessionRefreshRejected,
183    /// subscription intent targeted a stale or invalid room connection
184    #[error("subscription update rejected")]
185    SubscriptionUpdateRejected,
186}
187
188impl SfuCoreError {
189    /// whether this error should close the client as a protocol fault
190    #[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/// A cloneable core handle that admits users into room-bound [`MediaSession`]s.
204#[derive(Debug, Clone)]
205pub struct SfuCore {
206    media_transport: MediaTransport,
207    rooms: Arc<RoomManager>,
208}
209
210/// One admitted user connection in one room.
211///
212/// Room mutations revalidate the connection before committing room state.
213/// [`close`](Self::close) cannot remove a replacement connection and drains
214/// connection-scoped staged media when this session is current.
215///
216/// Futures returned by [`establish`](Self::establish), [`answer`](Self::answer),
217/// [`publish`](Self::publish),
218/// [`deactivate_publication`](Self::deactivate_publication),
219/// [`renegotiate`](Self::renegotiate), [`subscribe`](Self::subscribe),
220/// [`update_info`](Self::update_info) and [`close`](Self::close) are not
221/// cancellation safe. Once polled, await them to completion. Negotiation can
222/// otherwise leave room and transport state out of step with the local phase.
223/// Other mutations can leave committed room state without its transport or
224/// output effects.
225///
226/// Call [`close`](Self::close) before dropping the session. `Drop` performs no
227/// room or transport cleanup.
228#[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    /// Admits one room user and returns its connection-owning media session.
247    ///
248    /// An existing `request.user_id` replaces its current connection.
249    /// Replacement room and transport effects finish before return.
250    /// This future is not cancellation safe. Once polled, await it to completion.
251    /// Dropping it after membership commits can leave room or replacement
252    /// transport effects unfinished.
253    ///
254    /// # Errors
255    ///
256    /// Returns [`RoomManagerJoinError::MissingRoom`] when `room_id` is not
257    /// current. Returns [`RoomManagerJoinError::RoomFull`] when a new user
258    /// exceeds room capacity. Returns [`RoomManagerJoinError::NoUsableWorker`]
259    /// when no configured worker is usable. Returns [`RoomManagerJoinError::RouterState`]
260    /// when router placement cannot commit.
261    ///
262    /// # Panics
263    ///
264    /// Panics when existing relay state refers to an uncommitted source
265    /// placement.
266    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    /// creates the first browser offer for this connection
287    ///
288    /// returns `Ok(None)` after the initial offer has already been requested
289    /// this lets reconnect or duplicate-start paths retry safely without
290    /// creating a second initial offer
291    ///
292    /// # Errors
293    ///
294    /// returns [`SessionError::Core`] when the transport cannot create the
295    /// initial offer
296    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    /// accepts the answer for the pending offer and commits any ready room work
312    ///
313    /// a rejection before the worker consumes its pending offer leaves the
314    /// application round in place so a direct caller may retry it
315    /// failures after the worker consumes that offer do not promise retry even
316    /// when the RTC backend rejects the answer
317    /// when queued publish intent needs another SDP round the returned offer
318    /// must be sent to the client before the next answer
319    ///
320    /// # Errors
321    ///
322    /// returns [`SessionError::NoPendingRequest`] when no offer is pending
323    /// returns [`SessionError::Core`] when answer application fails, capability
324    /// projection fails or room state rejects the accepted answer as stale
325    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    /// applies publish intent for one user stream
371    ///
372    /// returns no offer when the intent is already queued, already active or
373    /// must wait for an in-flight answer
374    /// returns an offer when the browser must answer a new offer before the
375    /// publication can commit
376    ///
377    /// # Errors
378    ///
379    /// returns [`SessionError::Core`] when the media backend cannot stage a
380    /// publish that needs negotiation
381    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    /// deactivates one publication without changing negotiated media
402    ///
403    /// a queued first publication is cancelled
404    /// a staged first publication is rolled back and its pending answer creates
405    /// the cleanup offer
406    /// a committed publication keeps its source identity, routes and negotiated
407    /// MID until session teardown
408    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    /// closes this media session and removes its room connection if still current
421    ///
422    /// the call is idempotent
423    /// it returns `true` only when the room manager removed the current
424    /// connection
425    /// current-session cleanup drains connection-scoped staged media through
426    /// room state
427    /// stale sessions do not remove a replacement connection for the same user
428    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    /// creates a refresh offer when the stable session needs renegotiation
448    ///
449    /// returns `Ok(None)` before the initial offer, while an answer is pending
450    /// or when the transport reports that the requested refresh is unsupported
451    /// a call made while an answer is pending records that another offer should
452    /// be created after the answer commits
453    ///
454    /// # Errors
455    ///
456    /// returns [`SessionError::Core`] when the transport rejects
457    /// renegotiation
458    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    /// applies receiver intent for sources published by another user
478    ///
479    /// subscription intent is remembered even when no producer is currently
480    /// routable
481    /// once negotiation makes the receiver consumable, room effects create the
482    /// missing consumer routes
483    ///
484    /// # Errors
485    ///
486    /// returns [`SessionError::Core`] when room state rejects this connection
487    /// as stale
488    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    /// returns `None` when the transport has no endpoint for this session key
501    #[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    /// Returns snapshots for current room users other than this session's user.
539    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    /// Updates user information for this room connection.
563    ///
564    /// Missing or stale connections are ignored. `is_camera_on` and
565    /// `is_screen_sharing_on` are derived from publication state and discarded
566    /// from `info`.
567    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    /// Broadcasts `message` to every other current room user.
579    ///
580    /// Missing or stale sender connections return `Ok(())` without delivery.
581    ///
582    /// # Errors
583    ///
584    /// Returns [`BroadcastPayloadError::TooLarge`] when the serialized payload
585    /// exceeds the room broadcast limit. Returns
586    /// [`BroadcastPayloadError::JsonSerialization`] when JSON serialization
587    /// fails.
588    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    /// Rejects recording start because no persistent recording backend is enabled.
595    ///
596    /// `options` has no effect. The request records a rejection and returns
597    /// `false`.
598    #[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    /// Rejects recording stop because no persistent recording backend is enabled.
609    ///
610    /// The request records a rejection and returns `false`.
611    #[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}