Skip to main content

o_sfu/
lib.rs

1//! Odoo's Selective Forwarding Unit (SFU) for audio/video calls.
2//!
3//! `o-sfu` is a multi-tenant SFU designed to be compatible with the Odoo saas/.sh
4//! architecture, which means that a single o-sfu server can serve
5//! independent (individually authenticated) "rooms" which can be used by many
6//! independent tenants (like odoo saas "databases").
7//!
8//! # Reading Map
9//!
10//! - **Start or embed the server**: [`run`], [`Runtime`] and [`config`].
11//! - **Understand admission**: [`auth`], [`http`] and [`websocket`].
12//! - **Change room and media behavior**: [`core::server::room::Room`],
13//!   [`o_sfu_router::Router`] and [`core::server::transport::MediaTransport`].
14//! - **Integrate a browser client**: `SfuClient` exposes the Odoo-facing API.
15//!   [`o_sfu_protocol::host::ProtocolCore`] supplies its signaling state machine.
16//!
17//!
18//! # Core Concepts
19//!
20//! `o-sfu` separates the **control plane** for admission and room policy, the
21//! **routing plane** for user-to-connection placement and the **packet plane**
22//! for RTP forwarding on worker loops.
23//!
24//! - **[`Runtime`]**: Owns the process lifecycle, the HTTP/WebSocket servers and graceful shutdown.
25//! - **[`core::server::room::Room`]**: The control plane boundary for a set of participants. It commits membership and media relationships.
26//! - **[`o_sfu_router::Router`]**: The routing plane, a pure, sans-I/O engine owning the placement graph that maps users to connections.
27//! - **[`core::prelude::MediaSession`]**: Orchestrates a user's connection, bridging room intent to transport effects.
28//! - **[`core::server::transport::MediaTransport`]**: Owns the media workers and abstracts their threading model. Each worker holds a packet loop and applies projected routes to incoming datagrams.
29//!
30//! # Architecture
31//!
32//! Control and routing decisions happen above the packet loops, which apply the
33//! resulting transport state to UDP datagrams.
34//!
35//! ```text
36//! HTTP control API ------------------------------> RoomManager
37//!                                                       |
38//! WebSocket session -> SfuCore -> MediaSession          |
39//!                                    \                  /
40//!                                     +----> Room <----+
41//!                                             |
42//!                                             +<----> Router
43//!                                             |
44//!                                             v
45//!                                       MediaTransport
46//!                                        |    |    |
47//!                                        v    v    v
48//!                          RTC workers / packet loops
49//!                           (one thread per worker)
50//!                         |            |             |
51//!                      UDP socket   UDP socket   UDP socket
52//!                         |            |             |
53//!                         v            v             v
54//!                      fanout       fanout        fanout
55//!                      UDP OUT      UDP OUT       UDP OUT
56//! ```
57//!
58//! # Admission Edge
59//!
60//! Applications provision rooms through
61//! [`RoomManager::serve_room`](core::server::room::RoomManager::serve_room).
62//! WebSocket clients join through
63//! [`SfuCore::admit_user`](core::prelude::SfuCore::admit_user). Both paths require
64//! JWT authentication before admission.
65//!
66//! Upgraded sockets retain their connection permits. Before authentication,
67//! they have a first-frame deadline and global plus per-origin caps
68//! (one bucket per IPv4 address or IPv6 /64).
69//!
70//! ```text
71//! Incoming TCP connection
72//!     |
73//!     v
74//! connection cap + HTTP header deadline
75//!     |
76//!     v
77//! Axum Router
78//!     |
79//!     +-> GET /v1/noop
80//!     |     public liveness -> noop
81//!     |
82//!     +-> GET /v1/channel
83//!     |     Authorization JWT -> VerifiedRoomRequest -> room
84//!     |
85//!     +-> POST /v1/disconnect
86//!     |     request-body JWT -> VerifiedDisconnectClaims -> disconnect
87//!     |
88//!     +-> WebSocket /
89//!     |     upgrade -> first-frame JWT -> admit_user -> session
90//!     |
91//!     +-> GET /v1/stats, /metrics and /internal/diagnostics/...
92//!           authorize_operator -> stats, metrics or diagnostics
93//! ```
94//!
95//! - **HTTP**: Parses server-to-server requests using [`http::CreateRoomQuery`]. Verifies [`auth::HttpRoomClaims`]. The request that creates the current room fixes its signing key.
96//! - **WebSocket**: Client connection frames are decoded by [`websocket::decode_auth_payload_text`]. A hint selects a candidate key. Claims are verified, normalized into [`auth::WebSocketConnectClaims`] and trusted for access.
97//!
98//! # Security Model
99//!
100//! `o-sfu` secures two planes independently. Application-layer JWTs gate room
101//! admission on the control plane. `str0m` encrypts media on the packet plane
102//! with DTLS-SRTP. An external reverse proxy provides TLS for HTTP and
103//! WebSocket signaling.
104//!
105//! ```text
106//! control plane    JWT HS256           admission trust
107//! packet plane     DTLS-SRTP (str0m)   media confidentiality
108//! signaling wire   TLS (via proxy)     transport confidentiality
109//! ```
110//!
111//! ## JWT Admission
112//!
113//! Tokens are `HS256` only. [`auth::verify`] rejects any other `alg`, checks the
114//! HMAC in constant time and requires an unexpired `exp`. It validates `nbf`
115//! and the `iat` future-skew bound when present. Current Odoo does not issue
116//! `iat`. Tokens without `exp` fail with [`auth::AuthenticationError::MissingExpiry`].
117//! It caps token size at
118//! [`auth::MAX_JWT_TOKEN_BYTES`].
119//!
120//! There are two different keys:
121//!
122//! - **Server-to-server key**: `AUTH_KEY` (base64, at least 32 decoded bytes) verifies
123//!   the HTTP [`http::CreateRoomQuery`] path through [`auth::HttpRoomClaims`] and
124//!   [`auth::HttpDisconnectClaims`]. See [`config`].
125//! - **Per-room key**: the request that creates the current room pins the signing
126//!   key from the `key` or `keySeed` claim in [`auth::HttpRoomClaims`]. A direct
127//!   key must decode to at least 32 bytes or room creation returns `400`.
128//!   `keySeed` instead derives the room's signing bytes from `AUTH_KEY`:
129//!   ```text
130//!   room_key = HMAC-SHA256(
131//!       key = Base64Decode(AUTH_KEY),
132//!       message = Base64Decode(keySeed)
133//!   )
134//!   ```
135//!   WebSocket [`auth::WebSocketConnectClaims`] verify against that room key,
136//!   never against `AUTH_KEY`. Runtime construction decodes the global key once.
137//!   Room creation retains decoded secret bytes, so verification does not decode
138//!   stored keys and equivalent base64 encodings match the same reservation.
139//!
140//! HTTP room creation uses the
141//! `Authorization` header, HTTP disconnect uses the request body and the
142//! WebSocket client sends a first-frame auth envelope decoded by
143//! [`websocket::decode_auth_payload_text`]. An unverified room id selects only a
144//! candidate key, then the same token is verified once against it. Modern
145//! [`auth::WebSocketConnectClaims`] must name the selected room. Legacy Odoo
146//! tokens select it through the auth envelope's `channel` and are normalized
147//! only after verification with that room's key. Odoo may send both
148//! `session_id` and account `user_id`. The session takes precedence as the
149//! participant identity. Tokens without either identity are rejected.
150//!
151//! Admission establishes identity and room scope. It does not enforce the
152//! per-user `permissions` claim provided by each tenant.
153//!
154//! ## Media Transport
155//!
156//! `str0m` handles ICE, DTLS and SRTP over UDP. `o-sfu` builds and drives the
157//! `str0m` sessions to forward RTP between participants.
158//!
159//! - **Keying**: the DTLS handshake derives SRTP keys per RFC 5764 DTLS-SRTP.
160//! - **Certificate**: `str0m` generates a self-signed certificate when each RTC
161//!   session is built and advertises its SHA-256 fingerprint in the SDP offer.
162//!   Accepting the answer stores the expected remote fingerprint. The DTLS
163//!   handshake verifies it against the peer certificate.
164//! - **ICE**: `o-sfu` runs ICE-lite with `a=setup:actpass` and advertises
165//!   `ANNOUNCED_IP`, so media UDP must reach the host directly.
166//!
167//! ## Signaling Transport
168//!
169//! HTTPS and WSS are expected to be terminated by an external reverse proxy.
170//! Setting `PROXY=true` in [`config`] requires `TRUSTED_PROXIES` CIDRs.
171//! Forwarded headers are honored only when the TCP peer belongs to that set.
172//! The public edge must overwrite client-supplied forwarded host and protocol
173//! headers. See [`http::resolve_request_origin`] for origin
174//! resolution and [`http`] for operator route access.
175//!
176//! # Room and Router Ownership
177//!
178//! Room transitions produce typed commits while holding short exclusive state
179//! locks. `RoomEffects` consumes deferred transport, source-policy and WebSocket
180//! output work after the lock is released.
181//!
182//! ```text
183//! room state lock held                  lock released
184//! +--------------------------------+    +------------------------------+
185//! | validate user and connection   |    | MediaTransport commands      |
186//! | commit room topology           |    | source-policy turn           |
187//! | capture transition commit      |--->| websocket output             |
188//! |                                |    | idempotent teardown          |
189//! +--------------------------------+    +------------------------------+
190//! ```
191//!
192//! [`o_sfu_router::Router`] owns exact user-to-connection placement. Receiver shadows are foreign local sessions derived from active consumer dependencies, disappearing with their final consumer.
193//!
194//! Absent subscription targets are capped per receiver. Eviction drops the
195//! oldest absent target's intent. Present members do not consume that allowance.
196//!
197//! # Signaling and Client Bundle
198//!
199//! Browsers use `SfuClient` for connection, publication, subscription and room
200//! control. Signaling state stays in [`o_sfu_protocol::host::ProtocolCore`] and
201//! yields ordered [`o_sfu_protocol::host::Command`] values. The WASM bridge
202//! serializes commands for `BrowserRuntime`, which executes effects through
203//! browser `WebSocket`, `RTCPeerConnection` and timer APIs. Protocol events are
204//! mapped to Odoo bundle updates during command serialization. `BrowserRuntime`
205//! applies those updates to client state and notifies the application.
206//!
207//! Outbound overflow closes with `4110` (`Overloaded`), allowing reconnect and
208//! intent replay. Replacement or explicit removal uses terminal `4108` (`Kicked`).
209//!
210//! ```text
211//! SfuClient (public API)
212//!        |
213//!        v
214//! BrowserRuntime
215//!        |
216//!        v
217//! ProtocolCore (sans-I/O) -> Vec<Command>
218//!        |
219//!        v
220//! WASM serialization -> TypeScript command union
221//!        |
222//!        v
223//! BrowserRuntime -> WebSocket, RTCPeerConnection, timers
224//!        ^                                        |
225//!        +--------------- browser events ---------+
226//! ```
227//!
228//! # Packet Path
229//!
230//! [`core::server::transport::MediaTransport`] owns the media workers, which hold the packet loops. These loops receive UDP datagrams, drive WebRTC state (`str0m`), apply route tables and forward RTP.
231//!
232//! ```text
233//! UDP datagram
234//!     |
235//!     v
236//! worker ingress and `str0m` drain
237//!     |
238//!     v
239//! packet facts (source, RID and codec)
240//!     |
241//!     +-> origin packet sinks
242//!     |
243//!     +-> source and destination packet gates
244//!              |
245//!              +-> relay fanout
246//!              +-> local RTC -> RTP identity and codec rewrite
247//! ```
248//!
249//! Registered origin packet sinks observe publisher packets before route gates,
250//! including publishers without active receivers. Source, relay and receiver
251//! gates then narrow routed fanout. Same-process relays share payload data with
252//! another worker for local delivery.
253//!
254//! Worker BWE and audio observations feed into room source policy, which updates route gates for later packets.
255//!
256//! ```text
257//! worker BWE and audio observations
258//!                |
259//!                v
260//!        room source policy
261//!                |
262//!                v
263//!   route gates for later packets
264//! ```
265//!
266//! # Observability
267//!
268//! Monitored through the [`o_sfu_telemetry`] sub-crate. See [`http::telemetry`] for the HTTP contracts.
269//!
270//! - **Metrics**: [`http::telemetry::metrics`] exposes Prometheus text exposition.
271//! - **Diagnostics**: [`http::telemetry::diagnostics`] exposes JSON state summaries.
272//!
273//! # Scaling
274//!
275//! Rooms use one [`o_sfu_router::Router`] facade and default to one local router.
276//! [`config::RoomWorkerPolicy`] can enable additional same-process local routers.
277//! Only running workers are eligible. Joins prefer assigned workers with a known
278//! delay below the configured threshold. When none qualifies, a join may attach
279//! an unused healthy worker within the router cap and worker count. Otherwise,
280//! joins reuse a running assigned worker even when it exceeds the delay threshold.
281//! Placements on failed workers do not count toward the cap. If no assigned worker
282//! is running, a later join can attach a fresh placement on another running worker
283//! even under the single-router policy. Admission fails when no worker is running.
284//!
285//! # Feature Flags
286//!
287//! Core media behavior is configured at runtime through [`config`], not Cargo features.
288//! The default feature `otel-tracing` enables OpenTelemetry tracing support through [`o_sfu_telemetry::TraceExportConfig`].
289//! Other features are only used for tests and benchmarking.
290//!
291//! # Sub-crates
292//!
293//! | Crate | Role |
294//! | --- | --- |
295//! | [`o_sfu_rfc`] | RFC-backed JWT, RTP, RTCP, SDP and WebRTC consts/types |
296//! | `o-sfu-model` | Shared call data ([`o_sfu_protocol::wire::UserId`], etc.) |
297//! | [`o_sfu_router`] | Sans-I/O [`o_sfu_router::Router`] facade for room placement and routed media lifetimes |
298//! | [`o_sfu_core`] | Room engine, [`core::prelude::SourcePolicy`], recording taps and [`core::server::transport::MediaTransport`] projection |
299//! | [`o_sfu_protocol`] | Sans-I/O [`o_sfu_protocol::host::ProtocolCore`] and typed commands |
300//! | [`o_sfu_telemetry`] | Tracing setup, metrics and diagnostics response types |
301pub mod config;
302pub mod core {
303    pub use o_sfu_core::{prelude, server};
304}
305pub(crate) mod application;
306mod runtime;
307
308pub mod auth {
309    pub use crate::runtime::auth::{
310        AuthenticationError, HttpDisconnectClaims, HttpRoomClaims, MAX_JWT_TOKEN_BYTES,
311        RegisteredJwtClaims, WebSocketConnectClaims, sign, verify,
312    };
313}
314
315#[cfg(any(test, feature = "testing"))]
316pub mod test_support {
317    pub use crate::runtime::auth::test_support::TestHttpRoomClaims;
318}
319
320/// HTTP route and payload contracts.
321///
322/// `/v1/stats`, `/metrics` and diagnostics require the configured
323/// [`crate::config::DiagnosticsConfig::auth_token`] on every listener. Without
324/// one, the actual listener must be loopback. Missing or invalid tokens return
325/// `401 Unauthorized` with `WWW-Authenticate: Bearer realm="o-sfu"`.
326/// Tokenless non-loopback access returns `403 Forbidden`.
327pub mod http {
328    pub use crate::runtime::{
329        http_server::contract::{
330            CreateRoomQuery, IncomingBitRateStatsResponse, NoopResponse, RoomResponse,
331            StatsResponse, route,
332        },
333        request_origin::{RequestOrigin, resolve_request_origin},
334    };
335
336    /// Operator-facing metrics and diagnostics contracts.
337    pub mod telemetry {
338        /// Prometheus metric scrape contract.
339        ///
340        /// `GET` [`metrics::PATH`] returns `200 OK` Prometheus text exposition
341        /// with [`metrics::CONTENT_TYPE`].
342        ///
343        /// This is a scrape endpoint.
344        /// Configure Prometheus to scrape [`metrics::PATH`], then issue `PromQL`
345        /// queries to Prometheus.
346        /// Histogram families render `<name>_bucket` with the additional `le`
347        /// label plus `<name>_sum` and `<name>_count`.
348        ///
349        /// # Scrape and Query
350        ///
351        /// ```yaml
352        /// scrape_configs:
353        ///   - job_name: o-sfu
354        ///     scheme: https
355        ///     metrics_path: /metrics
356        ///     authorization:
357        ///       type: Bearer
358        ///       credentials_file: /run/secrets/o_sfu_diagnostics_token
359        ///     tls_config:
360        ///       ca_file: /run/secrets/o_sfu_observability_ca
361        ///       server_name: o-sfu-observability.internal
362        ///     static_configs:
363        ///       - targets: ["o-sfu-observability.internal:443"]
364        /// ```
365        ///
366        /// The endpoint returns Prometheus text exposition.
367        ///
368        /// ```text
369        /// # HELP osfu_rooms_active Current number of active rooms owned by this runtime.
370        /// # TYPE osfu_rooms_active gauge
371        /// osfu_rooms_active 3
372        /// # TYPE osfu_worker_rtp_packets_total counter
373        /// osfu_worker_rtp_packets_total{media_worker_id="0",direction="ingress"} 1240
374        /// ```
375        ///
376        /// Query clients send `PromQL` to the Prometheus-compatible backend
377        /// rather than to o-sfu.
378        /// These examples cover a gauge, counter and histogram.
379        ///
380        /// ```promql
381        /// sum(osfu_users_active)
382        /// sum by (stage) (rate(osfu_ws_connections_total[5m]))
383        /// histogram_quantile(
384        ///   0.95,
385        ///   sum by (le, route) (rate(osfu_http_request_duration_seconds_bucket[10m]))
386        /// )
387        /// ```
388        ///
389        /// Read the [`metrics::MetricName`] variants below to find every
390        /// exported name and its meaning.
391        /// Query [`metrics::PATH`] to see each family's `HELP`, `TYPE`, label
392        /// keys and current label values before building selectors.
393        pub mod metrics {
394            pub use o_sfu_telemetry::{
395                metrics::MetricName, prometheus::PROMETHEUS_CONTENT_TYPE as CONTENT_TYPE,
396            };
397
398            pub use crate::http::route::METRICS as PATH;
399        }
400
401        /// JSON diagnostics contract.
402        ///
403        /// Every constant in [`diagnostics::route`] is a `GET` endpoint returning
404        /// `200 OK` JSON on success.
405        ///
406        /// # Routes and Parameters
407        ///
408        /// | request | JSON response | parameter source |
409        /// | --- | --- | --- |
410        /// | `GET /internal/diagnostics/summary` | one [`diagnostics::DiagnosticsSummaryResponse`] | none |
411        /// | `GET /internal/diagnostics/rooms` | array of [`diagnostics::DiagnosticsRoomSummary`] | none |
412        /// | `GET /internal/diagnostics/workers` | array of [`diagnostics::DiagnosticsWorkerSummary`] | none |
413        /// | `GET /internal/diagnostics/rooms/{uuid}` | one [`diagnostics::DiagnosticsRoomDetail`] | `uuid` from the rooms response |
414        /// | `GET /internal/diagnostics/rooms/{uuid}/users` | array of [`diagnostics::DiagnosticsUserSummary`] | `uuid` from the rooms response |
415        /// | `GET /internal/diagnostics/rooms/{uuid}/users/{id}` | one [`diagnostics::DiagnosticsUserDetail`] | `uuid` from rooms and `userKey` from room users |
416        ///
417        /// `userId` may be a JSON number or string.
418        /// `userKey` is always the string to put into `{id}`.
419        /// URL-encode both path values before substitution.
420        ///
421        /// # Summary Request and Response over HTTPS
422        ///
423        /// ```text
424        /// GET /internal/diagnostics/summary HTTP/1.1
425        /// Host: o-sfu-observability.internal
426        /// Authorization: Bearer <diagnostics-token>
427        /// Accept: application/json
428        ///
429        /// HTTP/1.1 200 OK
430        /// Content-Type: application/json
431        ///
432        /// {
433        ///   "roomsActive": 1,
434        ///   "publicationsActive": 1,
435        ///   "recordingRoomsActive": 0,
436        ///   "usersActive": 2,
437        ///   "subscriptionsActive": 1,
438        ///   "transport": {
439        ///     "connectedUsers": 2,
440        ///     "disconnectedUsers": 0,
441        ///     "totalUsers": 2,
442        ///     "unknownUsers": 0
443        ///   }
444        /// }
445        /// ```
446        ///
447        /// # JavaScript Fetch Example
448        ///
449        /// ```javascript
450        /// const origin = "https://o-sfu-observability.internal";
451        /// const headers = {
452        ///   Authorization: `Bearer ${process.env.DIAGNOSTICS_AUTH_TOKEN}`,
453        /// };
454        ///
455        /// async function getJson(path) {
456        ///   const response = await fetch(`${origin}${path}`, { headers });
457        ///   if (!response.ok) {
458        ///     throw new Error(`${response.status} ${await response.text()}`);
459        ///   }
460        ///   return response.json();
461        /// }
462        ///
463        /// async function main() {
464        ///   const rooms = await getJson("/internal/diagnostics/rooms");
465        ///   const roomUuid = encodeURIComponent(rooms[0].uuid);
466        ///   const room = await getJson(`/internal/diagnostics/rooms/${roomUuid}`);
467        ///   const users = await getJson(`/internal/diagnostics/rooms/${roomUuid}/users`);
468        ///   const userKey = encodeURIComponent(users[0].userKey);
469        ///   const user = await getJson(
470        ///     `/internal/diagnostics/rooms/${roomUuid}/users/${userKey}`,
471        ///   );
472        ///
473        ///   console.log(room.summary, room.users, room.sources);
474        ///   console.log(user.user, user.recordingState);
475        /// }
476        ///
477        /// main().catch((error) => {
478        ///   console.error(error);
479        ///   process.exitCode = 1;
480        /// });
481        /// ```
482        ///
483        /// The rooms response has this shape.
484        ///
485        /// ```json
486        /// [
487        ///   {
488        ///     "createDate": "2026-07-15T10:20:30.000Z",
489        ///     "mediaWorkerId": 0,
490        ///     "publicationCount": 1,
491        ///     "recordingState": {
492        ///       "recording": false,
493        ///       "audio": false,
494        ///       "transcription": false,
495        ///       "video": false
496        ///     },
497        ///     "remoteAddress": "203.0.113.10",
498        ///     "sourceCount": 1,
499        ///     "userCount": 2,
500        ///     "subscriptionCount": 1,
501        ///     "transport": {
502        ///       "connectedUsers": 2,
503        ///       "disconnectedUsers": 0,
504        ///       "totalUsers": 2,
505        ///       "unknownUsers": 0
506        ///     },
507        ///     "uuid": "550e8400-e29b-41d4-a716-446655440000",
508        ///     "webRtcEnabled": true
509        ///   }
510        /// ]
511        /// ```
512        ///
513        /// The room users response has this shape.
514        ///
515        /// ```json
516        /// [
517        ///   {
518        ///     "audioIncomingBitrateBps": 32000,
519        ///     "cameraIncomingBitrateBps": 600000,
520        ///     "connectionId": 91,
521        ///     "health": "connected",
522        ///     "incomingBitrateBps": 632000,
523        ///     "mediaWorkerId": 0,
524        ///     "publicationCount": 2,
525        ///     "roomId": "550e8400-e29b-41d4-a716-446655440000",
526        ///     "screenIncomingBitrateBps": 0,
527        ///     "subscriptionCount": 1,
528        ///     "userId": 42,
529        ///     "userKey": "42"
530        ///   }
531        /// ]
532        /// ```
533        ///
534        /// The response structs below list every field in each payload.
535        /// Wire names are `camelCase` unless a field documents an exception.
536        ///
537        /// User detail is room-scoped because the same user key can be active
538        /// in several rooms.
539        pub mod diagnostics {
540            pub use o_sfu_telemetry::diagnostics::{
541                DiagnosticsRoomDetail, DiagnosticsRoomSummary, DiagnosticsSummaryResponse,
542                DiagnosticsUserDetail, DiagnosticsUserSummary, DiagnosticsWorkerSummary,
543            };
544
545            pub use crate::http::route::diagnostics as route;
546        }
547    }
548}
549
550pub mod websocket {
551    pub use crate::runtime::websocket_server::{
552        ClientBatchDecodeError, ClientBatchDecodeFailureKind, MAX_CLIENT_BATCH_ENVELOPES,
553        MAX_CLIENT_FRAME_BYTES, decode_auth_payload_text, decode_client_batch,
554    };
555}
556
557pub use self::runtime::{Runtime, ServeError, run};