o_sfu_core/engine/media_transport/rtc/
bootstrap.rs1use 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
43pub(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
115fn 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
126pub(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}