Skip to main content

o_sfu_core/engine/media_transport/rtc/
bootstrap.rs

1//! cold-path rtc socket and session bootstrap
2//!
3//! this module creates worker sockets during construction and session-local
4//! str0m state during negotiation
5//!
6//! bootstrap stops at transport setup
7//! room policy, media registration, SDP
8//! staging and packet routing stay in the worker modules responsible for those
9//! contracts
10
11use std::{
12    io::ErrorKind,
13    net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, UdpSocket as StdUdpSocket},
14    sync::Arc,
15    time::{Duration, Instant},
16};
17
18use o_sfu_rfc::webrtc;
19use str0m::{Candidate, bwe::Bitrate as Str0mBitrate};
20use tracing::{info, warn};
21
22use super::{
23    RtpProfile,
24    consumer_egress::{ConsumerStreamStore, RTX_CACHE_MAX_PACKETS},
25    egress::RtcEgress,
26    packet_loop::{RtcUdpSocket, UdpIngress},
27    state::{
28        RtcSessionState, SessionSdpNegotiationState, SharedRtcSocket, bitrate::MediaBitrateCounter,
29        slots::SessionStore,
30    },
31};
32use crate::{
33    Bitrate, RtcPortRange, RtcUdpIoBackend,
34    engine::{
35        media_transport::{TransportAdapterError, TransportSessionKey},
36        metrics::RtcMetricsRecorder,
37    },
38};
39#[cfg(any(test, feature = "internal-benchmarks", fuzzing))]
40#[path = "TESTS/bootstrap.rs"]
41pub(super) mod test_support;
42
43/// bind the shared worker UDP socket and return the advertised candidate tuple
44///
45/// the bind address uses the unspecified address for the configured public IP
46/// family so one socket can receive traffic for every session assigned to the
47/// worker
48/// the advertised candidate keeps the configured public IP with the
49/// bound port because that is what browsers must see in SDP
50///
51/// tries configured ports in order and skips only [`ErrorKind::AddrInUse`]
52/// any other bind, nonblocking or backend initialization error aborts startup
53///
54/// # errors
55///
56/// returns [`TransportAdapterError::TransportUnavailable`] when every port is
57/// occupied or another socket operation fails
58pub(super) fn bind_shared_rtc_socket(
59    announced_ip: IpAddr,
60    rtc_port_range: RtcPortRange,
61    rtc_udp_io_backend: RtcUdpIoBackend,
62    rtc_metrics: &Arc<RtcMetricsRecorder>,
63) -> Result<SharedRtcSocket, TransportAdapterError> {
64    let bind_ip = bind_ip_for_announced_ip(announced_ip);
65    for port in rtc_port_range.ports() {
66        let bind_addr = SocketAddr::new(bind_ip, port);
67        let socket = match StdUdpSocket::bind(bind_addr) {
68            Ok(socket) => socket,
69            Err(error) if error.kind() == ErrorKind::AddrInUse => continue,
70            Err(error) => {
71                warn!(%bind_addr, ?error, "failed to bind shared rtc UDP socket");
72                return Err(TransportAdapterError::TransportUnavailable);
73            }
74        };
75        socket.set_nonblocking(true).map_err(|error| {
76            warn!(%bind_addr, ?error, "failed to configure shared rtc UDP socket");
77            TransportAdapterError::TransportUnavailable
78        })?;
79        let socket = RtcUdpSocket::from_std(socket, rtc_udp_io_backend).map_err(|error| {
80            warn!(
81                %bind_addr,
82                backend = rtc_udp_io_backend.wire_name(),
83                ?error,
84                "failed to initialize shared rtc UDP socket"
85            );
86            TransportAdapterError::TransportUnavailable
87        })?;
88        let candidate_addr = SocketAddr::new(announced_ip, port);
89        let ingress = UdpIngress::new(
90            socket.clone(),
91            bind_addr,
92            candidate_addr,
93            Arc::clone(rtc_metrics),
94        );
95        info!(
96            %bind_addr,
97            %candidate_addr,
98            "booted shared rtc UDP socket"
99        );
100        return Ok(SharedRtcSocket {
101            udp_candidate_addr: candidate_addr,
102            egress: RtcEgress::new(socket, candidate_addr, Arc::clone(rtc_metrics)),
103            ingress,
104        });
105    }
106    warn!(
107        %bind_ip,
108        port_min = rtc_port_range.min(),
109        port_max = rtc_port_range.max(),
110        "rtc UDP port range is unavailable"
111    );
112    Err(TransportAdapterError::TransportUnavailable)
113}
114
115/// choose the local bind address that matches the advertised IP family
116///
117/// workers bind to all local interfaces for that family while keeping SDP
118/// candidates anchored to the configured announced IP
119fn bind_ip_for_announced_ip(announced_ip: IpAddr) -> IpAddr {
120    match announced_ip {
121        IpAddr::V4(_) => IpAddr::V4(Ipv4Addr::UNSPECIFIED),
122        IpAddr::V6(_) => IpAddr::V6(Ipv6Addr::UNSPECIFIED),
123    }
124}
125
126/// ensure one worker-local [`RtcSessionState`] exists for a session
127///
128/// this is idempotent because negotiation may ask for readiness more than once
129/// while a session is still alive
130/// `Ok(true)` means a fresh [`str0m::Rtc`] was created and inserted
131/// `Ok(false)` means the existing session state still matches that key
132///
133/// new sessions start in ICE-lite mode, with RTP mode enabled, bandwidth
134/// estimation capped by `max_bitrate_out` and exactly one local host candidate
135/// attached to the shared worker socket's advertised address
136///
137/// # errors
138///
139/// returns `TransportUnavailable` if the local candidate cannot be represented
140/// by str0m or cannot be attached to the newly created rtc state
141pub(super) fn ensure_session_rtc_state(
142    users: &mut SessionStore,
143    room_id: Arc<str>,
144    session_key: &TransportSessionKey,
145    candidate_addr: SocketAddr,
146    max_bitrate_out: Bitrate,
147    profile: &RtpProfile,
148    stats_interval: Option<Duration>,
149) -> Result<bool, TransportAdapterError> {
150    if users.contains_key(session_key) {
151        return Ok(false);
152    }
153    let started_at = Instant::now();
154    let mut config = profile
155        .session_config()
156        .set_send_buffer_video(RTX_CACHE_MAX_PACKETS);
157    if let Some(stats_interval) = stats_interval {
158        config = config.set_stats_interval(Some(stats_interval));
159    }
160    let mut rtc = config
161        .enable_bwe(Some(Str0mBitrate::bps(max_bitrate_out.as_bps())))
162        .set_ice_lite(true)
163        .build(started_at);
164    let candidate = Candidate::host(candidate_addr, webrtc::IceTransport::Udp.as_str())
165        .map_err(|_error| TransportAdapterError::TransportUnavailable)?;
166    if rtc.add_local_candidate(candidate).is_none() {
167        return Err(TransportAdapterError::TransportUnavailable);
168    }
169    let local_ice_ufrag = rtc.direct_api().local_ice_credentials().ufrag;
170    users.insert(
171        session_key.clone(),
172        RtcSessionState {
173            room_id,
174            rtc,
175            started_at,
176            rtcp_ingress_budget: super::state::RtcpIngressBudget::new(started_at),
177            defer_rtx_expiry: false,
178            pending_rtp_input: None,
179            nack_totals: super::state::RtcNackTotals::default(),
180            egress_bitrate: Arc::new(MediaBitrateCounter::new(started_at)),
181            local_ice_ufrag,
182            #[cfg(test)]
183            max_bitrate_in: None,
184            #[cfg(test)]
185            max_bitrate_out: Some(max_bitrate_out),
186            receiver_bwe_target: None,
187            #[cfg(test)]
188            receiver_bwe_str0m_update_count: 0,
189            dtls_started: false,
190            packet_loop_dirty: false,
191            next_timeout: None,
192            sdp_negotiation: SessionSdpNegotiationState::default(),
193            consumer_streams: ConsumerStreamStore::default(),
194        },
195    );
196    Ok(true)
197}