Skip to main content

o_sfu_core/
server.rs

1//! Server-runtime integration surface.
2//!
3//! The top-level runtime uses these facades for diagnostics, metrics, room
4//! management, transport construction and packet sinks. The caller-facing
5//! session API remains under [`crate::prelude`].
6
7/// diagnostics response types used by runtime inspection endpoints
8pub mod diagnostics {
9    pub use o_sfu_telemetry::diagnostics::*;
10}
11
12/// process-local metric catalog and typed recorders
13///
14/// runtime edges record through this facade instead of assembling metric names
15/// or label sets manually
16/// the Prometheus renderer reads the catalog declared by these types
17pub mod metrics {
18    pub use crate::engine::metrics::*;
19}
20
21/// recording sink trait used by media routing code
22///
23/// recording integrations register packet sinks through the packet-sink
24/// registry while transport code only depends on this narrow sink contract
25pub mod recording {
26    pub use crate::engine::packet_sink_registry::PacketSink as MediaPacketSink;
27}
28
29/// room packet-sink registry shared by transport workers
30///
31/// packet sinks are looked up by room when media has to fan out to recording or
32/// other non-local destinations
33/// the registry keeps that routing concern out of packet-loop callers
34pub mod packet_sinks {
35    pub use crate::engine::packet_sink_registry::{
36        PacketSink, RegisteredPacketSink, RoomPacketSinkRegistry,
37    };
38}
39
40/// room manager, room runtime and user-session integration types
41///
42/// this facade is the server crate's entry point for admitting users, sending
43/// outbound messages and reading room statistics
44/// pure routing and transport details remain behind the room API
45///
46/// # Slow-consumer shutdown
47///
48/// [`room::UserOutboundSender::send`] queues output without waiting for
49/// capacity
50/// when the queue is full it also records an overflow sentinel
51/// [`room::UserOutboundReceiver::recv_event`] returns that sentinel before
52/// output that was already queued
53/// the caller should stop normal draining after
54/// [`room::UserOutboundEvent::Overflow`]
55///
56/// ```
57/// # use std::sync::Arc;
58/// # use o_sfu_core::server::{
59/// #     metrics::RuntimeMetrics,
60/// #     room::{
61/// #         RoomEventMessage, UserOutbound, UserOutboundEvent, UserOutboundSendError,
62/// #         UserOutboundSender,
63/// #     },
64/// #     session::UserId,
65/// # };
66/// # #[tokio::main(flavor = "current_thread")]
67/// # async fn main() {
68/// let (sender, mut receiver) =
69///     UserOutboundSender::channel(1, Arc::new(RuntimeMetrics::default()));
70/// let departed = |id| {
71///     UserOutbound::Message(RoomEventMessage::UserDeparted {
72///         user_id: UserId::Integer(id),
73///     })
74/// };
75/// let queued_output = departed(1);
76/// assert!(sender.send(queued_output).is_ok());
77///
78/// assert!(matches!(
79///     sender.send(departed(2)),
80///     Err(UserOutboundSendError::Full(_))
81/// ));
82/// assert!(matches!(
83///     receiver.recv_event().await,
84///     UserOutboundEvent::Overflow(_)
85/// ));
86/// # }
87/// ```
88pub mod room {
89    #[cfg(any(test, feature = "testing-transport"))]
90    pub mod test_support {
91        pub use crate::engine::{
92            room::{NegotiatedPublish, RoomManagerTestApi, RoomTestApi},
93            source_model::test_support::{
94                TestSourceKind, TestSubscriptionStates, source_kind_for_stream_id,
95                source_publish_intent_for_source, stream_id_for_source,
96                subscription_intents_from_test_states,
97            },
98        };
99    }
100
101    #[cfg(feature = "internal-benchmarks")]
102    pub mod benchmark_support {
103        pub use crate::engine::room::run_source_policy_turn_for_benchmark;
104    }
105
106    #[cfg(any(test, feature = "testing-transport"))]
107    pub use crate::engine::room::{ConsumerRouteState, JoinPlacementTestGate};
108    pub use crate::{
109        MediaWorkerId,
110        engine::room::{
111            BroadcastPayload, BroadcastPayloadError, DEFAULT_USER_OUTBOUND_QUEUE_BYTE_CAPACITY,
112            DEFAULT_USER_OUTBOUND_QUEUE_CAPACITY, IncomingBitrateSnapshot, JoinUserRequest,
113            MAX_BROADCAST_PAYLOAD_BYTES, RemoteTrackProjection, RemoteTrackSnapshot, Room,
114            RoomAdmissionPolicy, RoomConfig, RoomDetailCapture, RoomEventMessage, RoomJoinError,
115            RoomManager, RoomManagerJoinError, RoomManagerServeError, RoomMediaCounts,
116            RoomOverviewCapture, RoomRuntimeContext, RoomRuntimePolicy, RoomUserAdmission,
117            RoomUserCapture, RoomUserPermissions, RoomUserStatsSnapshot, RoomUsersCapture,
118            RouterPlacement, RouterPlacements, RouterPlacementsError, RuntimeRoomDirectorySnapshot,
119            RuntimeRoomStatsSnapshot, UserCloseReason, UserOutbound, UserOutboundEvent,
120            UserOutboundOverflow, UserOutboundOverflowKind, UserOutboundQueueLimits,
121            UserOutboundReceiver, UserOutboundSendError, UserOutboundSender,
122        },
123    };
124}
125
126/// signaling-domain payloads shared by rooms and WebSocket sessions
127///
128/// these types describe users, features, permissionsm recording state and
129/// host-visible close codes
130/// they are re-exported here so the runtime does not import private room or
131/// engine modules for protocol payload construction
132pub mod session {
133    pub use crate::engine::{
134        AvailableFeatures, JsonPayload, PeerSnapshot, RecordingOptions, RecordingState,
135        RecordingStateUpdate, StopCode, UserId, UserInfo, UserPermissions, VideoLayoutIntent,
136        WebSocketCloseCode,
137    };
138}
139
140/// source descriptor model accepted by publication and subscription flows
141///
142/// source-model types make published media identity explicit before it reaches
143/// room routing or transport negotiation
144pub mod source_model {
145    pub use crate::engine::source_model::{
146        PublishedSourceDescriptor, PublishedSourceDescriptorParts, PublishedSourceId,
147        PublishedSourceOwner, SourceEncodingDescriptor, SourceEncodingDescriptorParts,
148        SourceEncodingId, SourceModelError,
149    };
150}
151
152/// media transport construction and extension boundary
153///
154/// the runtime builds one `MediaTransport` from owner configuration and process
155/// services then room code uses the opaque handle for media operations
156///
157/// code above this module should not branch on RTC worker internals
158pub mod transport {
159
160    #[cfg(any(test, feature = "testing-transport"))]
161    pub mod test_support {
162        //! non-production media transport route inspectors
163
164        pub use crate::engine::media_transport::test_support::*;
165    }
166
167    #[cfg(feature = "internal-benchmarks")]
168    pub mod benchmark_support {
169        //! feature-gated packet-loop benchmark fixtures
170
171        pub use crate::engine::media_transport::benchmark_support::*;
172    }
173
174    #[cfg(any(test, fuzzing))]
175    pub mod fuzz_support {
176        pub use crate::engine::media_transport::fuzz_support::{
177            client_rtp_capabilities_from_answer, route_packet_loop_ingress_demux,
178        };
179    }
180
181    pub use crate::{
182        MediaWorkerId,
183        engine::media_transport::{
184            ActiveSpeakerActivityReason, ActiveSpeakerActivityState, ActiveSpeakerSource,
185            ActiveSpeakerSourceDiagnostic, AppliedProducer, AppliedSessionAnswer, ConsumerActivity,
186            MediaTransport, MediaTransportBuildError, MediaTransportConfig, MediaTransportDeps,
187            ProducerActivity, ReceiverBandwidthSnapshot, RelayRouteActivity, SessionOffer,
188            SessionUploadEncoding, SessionUploadSlot, SourcePacketGate, SourcePolicySignal,
189            SourcePolicyUpdateSubscription, TransportAdapterError, TransportBitrateSnapshot,
190            TransportConsumerRoute, TransportHealthSnapshot, TransportMediaId,
191            TransportQualitySample, TransportQualitySnapshot, TransportRelayRouteAction,
192            TransportRelayRouteEffect, TransportResult, TransportSessionHealth,
193            TransportSessionKey, TransportSourceDiagnosticsSnapshot, TransportSourceKey,
194            TransportWorkerPressureSnapshot,
195        },
196        prelude::SessionBitrateLimits,
197    };
198}