File
Blob: firmware/vendor/str0m/src/lib.rs
| 1 | //! <image src="https://user-images.githubusercontent.com/227204/226143511-66fe5264-6ab7-47b9-9551-90ba7e155b96.svg" alt="str0m logo" ></image> |
| 2 | //! |
| 3 | //! A Sans I/O WebRTC implementation in Rust. |
| 4 | //! |
| 5 | //! This is a [Sans I/O][sansio] implementation meaning the `Rtc` instance itself is not doing any network |
| 6 | //! talking. Furthermore it has no internal threads or async tasks. All operations are happening from the |
| 7 | //! calls of the public API. |
| 8 | //! |
| 9 | //! This is deliberately not a standard `RTCPeerConnection` API since that isn't a great fit for Rust. |
| 10 | //! See more details in below section. |
| 11 | //! |
| 12 | //! # Join us |
| 13 | //! |
| 14 | //! We are discussing str0m things on Discord. Join us using this [invitation link][discord]. |
| 15 | //! |
| 16 | //! <image width="300px" src="https://user-images.githubusercontent.com/227204/209446544-f8a8d673-cb1b-4144-a0f2-42307b8d8869.gif" alt="silly clip showing video playing" ></image> |
| 17 | //! |
| 18 | //! # Usage |
| 19 | //! |
| 20 | //! The [`chat`][x-chat] example shows how to connect multiple browsers |
| 21 | //! together and act as an SFU (Selective Forwarding Unit). The example |
| 22 | //! multiplexes all traffic over one server UDP socket and uses two threads |
| 23 | //! (one for the web server, and one for the SFU loop). |
| 24 | //! |
| 25 | //! ## TLS |
| 26 | //! |
| 27 | //! For the browser to do WebRTC, all traffic must be under TLS. The |
| 28 | //! project ships with a self-signed certificate that is used for the |
| 29 | //! examples. The certificate is for hostname `str0m.test` since TLD .test |
| 30 | //! should never resolve to a real DNS name. |
| 31 | //! |
| 32 | //! ```text |
| 33 | //! cargo run --example chat |
| 34 | //! ``` |
| 35 | //! |
| 36 | //! The log should prompt you to connect a browser to https://10.0.0.103:3000 – this will |
| 37 | //! most likely cause a security warning that you must get the browser to accept. |
| 38 | //! |
| 39 | //! The [`http-post`][x-post] example roughly illustrates how to receive |
| 40 | //! media data from a browser client. The example is single threaded and |
| 41 | //! is a bit simpler than the chat. It is a good starting point to understand the API. |
| 42 | //! |
| 43 | //! ```text |
| 44 | //! cargo run --example http-post |
| 45 | //! ``` |
| 46 | //! |
| 47 | //! ### Real example |
| 48 | //! |
| 49 | //! To see how str0m is used in a real project, check out [BitWHIP][bitwhip] – |
| 50 | //! a CLI WebRTC Agent written in Rust. |
| 51 | //! |
| 52 | //! ## Passive |
| 53 | //! |
| 54 | //! For passive connections, i.e. where the media and initial OFFER is |
| 55 | //! made by a remote peer, we need these steps to open the connection. |
| 56 | //! |
| 57 | //! ```no_run |
| 58 | //! # use std::time::Instant; |
| 59 | //! # use str0m::{Rtc, Candidate}; |
| 60 | //! // Instantiate a new Rtc instance. |
| 61 | //! let mut rtc = Rtc::new(Instant::now()); |
| 62 | //! |
| 63 | //! // Add some ICE candidate such as a locally bound UDP port. |
| 64 | //! let addr = "1.2.3.4:5000".parse().unwrap(); |
| 65 | //! let candidate = Candidate::host(addr, "udp").unwrap(); |
| 66 | //! rtc.add_local_candidate(candidate); |
| 67 | //! |
| 68 | //! // Accept an incoming offer from the remote peer |
| 69 | //! // and get the corresponding answer. |
| 70 | //! let offer = todo!(); |
| 71 | //! let answer = rtc.sdp_api().accept_offer(offer).unwrap(); |
| 72 | //! |
| 73 | //! // Forward the answer to the remote peer. |
| 74 | //! |
| 75 | //! // Go to _run loop_ |
| 76 | //! ``` |
| 77 | //! |
| 78 | //! ## Active |
| 79 | //! |
| 80 | //! Active connections means we are making the inital OFFER and waiting for a |
| 81 | //! remote ANSWER to start the connection. |
| 82 | //! |
| 83 | //! ```no_run |
| 84 | //! # use std::time::Instant; |
| 85 | //! # use str0m::{Rtc, Candidate}; |
| 86 | //! # use str0m::media::{MediaKind, Direction}; |
| 87 | //! // Instantiate a new Rtc instance. |
| 88 | //! let mut rtc = Rtc::new(Instant::now()); |
| 89 | //! |
| 90 | //! // Add some ICE candidate such as a locally bound UDP port. |
| 91 | //! let addr = "1.2.3.4:5000".parse().unwrap(); |
| 92 | //! let candidate = Candidate::host(addr, "udp").unwrap(); |
| 93 | //! rtc.add_local_candidate(candidate); |
| 94 | //! |
| 95 | //! // Create a `SdpApi`. The change lets us make multiple changes |
| 96 | //! // before sending the offer. |
| 97 | //! let mut change = rtc.sdp_api(); |
| 98 | //! |
| 99 | //! // Do some change. A valid OFFER needs at least one "m-line" (media). |
| 100 | //! let mid = change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None); |
| 101 | //! |
| 102 | //! // Get the offer. |
| 103 | //! let (offer, pending) = change.apply().unwrap(); |
| 104 | //! |
| 105 | //! // Forward the offer to the remote peer and await the answer. |
| 106 | //! // How to transfer this is outside the scope for this library. |
| 107 | //! let answer = todo!(); |
| 108 | //! |
| 109 | //! // Apply answer. |
| 110 | //! rtc.sdp_api().accept_answer(pending, answer).unwrap(); |
| 111 | //! |
| 112 | //! // Go to _run loop_ |
| 113 | //! ``` |
| 114 | //! |
| 115 | //! ## Run loop |
| 116 | //! |
| 117 | //! ### The single-mutation invariant |
| 118 | //! |
| 119 | //! str0m's API has one strict contract that the run loop is built around: |
| 120 | //! |
| 121 | //! > **Every mutation of an `Rtc` instance must be followed by a complete |
| 122 | //! > drain of `poll_output` until it returns `Output::Timeout`, before the |
| 123 | //! > next mutation on the same `Rtc`.** |
| 124 | //! |
| 125 | //! A "mutation" is anything that takes `&mut Rtc` (directly or through a |
| 126 | //! handle obtained from it). The common ones are: |
| 127 | //! |
| 128 | //! - `Rtc::handle_input` — feeding a network packet or a timeout |
| 129 | //! - `Writer::write` / `Writer::request_keyframe` — sending media |
| 130 | //! - `Channel::write` — sending data channel data |
| 131 | //! - `SdpApi::apply` / `DirectApi::*` — negotiation |
| 132 | //! - `Rtc::add_local_candidate` / `Rtc::add_remote_candidate` |
| 133 | //! |
| 134 | //! Always: **mutate → drain to `Timeout` → mutate → drain to `Timeout` → …** |
| 135 | //! |
| 136 | //! Doing two mutations back-to-back without draining in between, or |
| 137 | //! waiting on I/O while the engine still has output queued, leaves the |
| 138 | //! engine in an inconsistent state and produces wrong behavior. Mutations |
| 139 | //! issued from inside the drain loop (e.g. calling `Writer::write` in |
| 140 | //! response to an `Output::Event`) are fine — the drain loop naturally |
| 141 | //! continues calling `poll_output` afterward and so the invariant holds. |
| 142 | //! |
| 143 | //! ### Canonical shape |
| 144 | //! |
| 145 | //! Driving an `Rtc` forward follows the same six-step shape, regardless |
| 146 | //! of sync or async: |
| 147 | //! |
| 148 | //! 1. **Wait** for one of: the next timeout firing, an incoming network |
| 149 | //! packet, or the application wanting to perform a mutation (e.g. |
| 150 | //! write media). |
| 151 | //! 2. **Perform that ONE mutation** — feed `Input` to `handle_input`, or |
| 152 | //! call into a writer / channel / SDP API. |
| 153 | //! 3. **Poll** `Rtc::poll_output`. |
| 154 | //! 4. **Handle** the output: `Output::Transmit` is sent on the socket, |
| 155 | //! `Output::Event` is dispatched to the application, `Output::Timeout` |
| 156 | //! records the next deadline. |
| 157 | //! 5. **Goto 3** until `poll_output` returns `Output::Timeout`. Only then |
| 158 | //! is the engine fully drained. |
| 159 | //! 6. The returned timeout is what we wait on next — **goto 1**. |
| 160 | //! |
| 161 | //! ```no_run |
| 162 | //! # use str0m::{Rtc, Output, IceConnectionState, Event, Input}; |
| 163 | //! # use str0m::net::{Receive, Protocol}; |
| 164 | //! # use std::io::ErrorKind; |
| 165 | //! # use std::net::UdpSocket; |
| 166 | //! # use std::time::{Duration, Instant}; |
| 167 | //! # let mut rtc = Rtc::new(Instant::now()); |
| 168 | //! // A UdpSocket we obtained _somehow_. |
| 169 | //! let socket: UdpSocket = todo!(); |
| 170 | //! |
| 171 | //! // Buffer for reading incoming UDP packets. |
| 172 | //! let mut buf = vec![0; 2000]; |
| 173 | //! |
| 174 | //! loop { |
| 175 | //! // === Steps 3-5: drain poll_output until we get the next timeout. === |
| 176 | //! let timeout = loop { |
| 177 | //! match rtc.poll_output().unwrap() { |
| 178 | //! // Step 5: a Timeout exits the drain loop. |
| 179 | //! Output::Timeout(t) => break t, |
| 180 | //! |
| 181 | //! // Step 4: transmit on the socket and keep draining. |
| 182 | //! // The destination IP comes from the ICE agent and may |
| 183 | //! // change during the session. |
| 184 | //! Output::Transmit(t) => { |
| 185 | //! socket.send_to(&t.contents, t.destination).unwrap(); |
| 186 | //! } |
| 187 | //! |
| 188 | //! // Step 4: hand the event to the application and keep draining. |
| 189 | //! // Events are mainly incoming media data from the remote peer, |
| 190 | //! // but also data channel data and statistics. |
| 191 | //! Output::Event(e) => { |
| 192 | //! if e == Event::IceConnectionStateChange(IceConnectionState::Disconnected) { |
| 193 | //! return; |
| 194 | //! } |
| 195 | //! // TODO: handle other events here, such as incoming media data. |
| 196 | //! } |
| 197 | //! } |
| 198 | //! }; |
| 199 | //! |
| 200 | //! // === Step 1: wait for ONE of: the timeout firing, an incoming |
| 201 | //! // packet, or application-side data. The example below uses a |
| 202 | //! // blocking UDP socket with a read timeout. With async you would |
| 203 | //! // `select!` over multiple futures; with application-side data you |
| 204 | //! // would also include a channel. |
| 205 | //! // |
| 206 | //! // set_read_timeout(Some(ZERO)) is not allowed, so clamp to >= 1ms. |
| 207 | //! let duration = (timeout - Instant::now()).max(Duration::from_millis(1)); |
| 208 | //! socket.set_read_timeout(Some(duration)).unwrap(); |
| 209 | //! buf.resize(2000, 0); |
| 210 | //! |
| 211 | //! // === Step 2: take ONE event and feed it as Input === |
| 212 | //! let input = match socket.recv_from(&mut buf) { |
| 213 | //! Ok((n, source)) => { |
| 214 | //! buf.truncate(n); |
| 215 | //! Input::Receive( |
| 216 | //! Instant::now(), |
| 217 | //! Receive { |
| 218 | //! proto: Protocol::Udp, |
| 219 | //! source, |
| 220 | //! destination: socket.local_addr().unwrap(), |
| 221 | //! contents: buf.as_slice().try_into().unwrap(), |
| 222 | //! }, |
| 223 | //! ) |
| 224 | //! } |
| 225 | //! |
| 226 | //! // The socket read timed out — feed Input::Timeout to advance |
| 227 | //! // the engine to the deadline. WouldBlock is the unix error, |
| 228 | //! // TimedOut is the windows error. |
| 229 | //! Err(e) if matches!(e.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut) => { |
| 230 | //! Input::Timeout(Instant::now()) |
| 231 | //! } |
| 232 | //! |
| 233 | //! Err(e) => { |
| 234 | //! eprintln!("Error: {:?}", e); |
| 235 | //! return; |
| 236 | //! } |
| 237 | //! }; |
| 238 | //! |
| 239 | //! rtc.handle_input(input).unwrap(); |
| 240 | //! |
| 241 | //! // === Step 6: back to the top of the outer loop (goto step 3). === |
| 242 | //! } |
| 243 | //! ``` |
| 244 | //! |
| 245 | //! ## Sending media data |
| 246 | //! |
| 247 | //! When creating the media, we can decide which codecs to support, and they |
| 248 | //! are negotiated with the remote side. Each codec corresponds to a |
| 249 | //! "payload type" (PT). To send media data we need to figure out which PT |
| 250 | //! to use when sending. |
| 251 | //! |
| 252 | //! ```no_run |
| 253 | //! # use str0m::Rtc; |
| 254 | //! # use str0m::media::Mid; |
| 255 | //! # let rtc: Rtc = todo!(); |
| 256 | //! // Obtain mid from Event::MediaAdded |
| 257 | //! let mid: Mid = todo!(); |
| 258 | //! |
| 259 | //! // Create a media writer for the mid. |
| 260 | //! let writer = rtc.writer(mid).unwrap(); |
| 261 | //! |
| 262 | //! // Get the payload type (pt) for the wanted codec. |
| 263 | //! let pt = writer.payload_params().nth(0).unwrap().pt(); |
| 264 | //! |
| 265 | //! // Write the data |
| 266 | //! let wallclock = todo!(); // Absolute time of the data |
| 267 | //! let media_time = todo!(); // Media time, in RTP time |
| 268 | //! let data: &[u8] = todo!(); // Actual data |
| 269 | //! writer.write(pt, wallclock, media_time, data).unwrap(); |
| 270 | //! ``` |
| 271 | //! |
| 272 | //! `Writer::write` is a mutation, so the |
| 273 | //! [single-mutation invariant](#the-single-mutation-invariant) applies: |
| 274 | //! after writing, drain `Rtc::poll_output` to `Output::Timeout` before |
| 275 | //! the next mutation on this `Rtc`. |
| 276 | //! |
| 277 | //! ## Media time, wallclock and local time |
| 278 | //! |
| 279 | //! str0m has three main concepts of time. "now", media time and wallclock. |
| 280 | //! |
| 281 | //! ### Now |
| 282 | //! |
| 283 | //! Some calls in str0m, such as `Rtc::handle_input` takes a `now` argument |
| 284 | //! that is a `std::time::Instant`. These calls "drive the time forward" in |
| 285 | //! the internal state. This is used for everything like deciding when |
| 286 | //! to produce various feedback reports (RTCP) to remote peers, to |
| 287 | //! bandwidth estimation (BWE) and statistics. |
| 288 | //! |
| 289 | //! Str0m has _no internal clock_ calls. I.e. str0m never calls |
| 290 | //! `Instant::now()` itself. All time is external input. That means it's |
| 291 | //! possible to construct test cases driving an `Rtc` instance faster |
| 292 | //! than realtime (see the [integration tests][intg]). |
| 293 | //! |
| 294 | //! ### Media time |
| 295 | //! |
| 296 | //! Each RTP header has a 32 bit number that str0m calls _media time_. |
| 297 | //! Media time is in some time base that is dependent on the codec, |
| 298 | //! however all codecs in str0m use 90_000Hz for video and 48_000Hz |
| 299 | //! for audio. |
| 300 | //! |
| 301 | //! For video the `MediaTime` type is `<timestamp>/90_000` str0m extends |
| 302 | //! the 32 bit number in the RTP header to 64 bit taking into account |
| 303 | //! "rollover". 64 bit is such a large number the user doesn't need to |
| 304 | //! think about rollovers. |
| 305 | //! |
| 306 | //! ### Wallclock |
| 307 | //! |
| 308 | //! With _wallclock_ str0m means the time a sample of media was produced |
| 309 | //! at an originating source. I.e. if we are talking into a microphone the |
| 310 | //! wallclock is the NTP time the sound is sampled. |
| 311 | //! |
| 312 | //! We can't know the exact wallclock for media from a remote peer since |
| 313 | //! not every device is synchronized with NTP. Every sender does |
| 314 | //! periodically produce a Sender Report (SR) that contains the peer's |
| 315 | //! idea of its wallclock, however this number can be very wrong compared to |
| 316 | //! "real" NTP time. |
| 317 | //! |
| 318 | //! Furthermore, not all remote devices will have a linear idea of |
| 319 | //! time passing that exactly matches the local time. A minute on the |
| 320 | //! remote peer might not be exactly one minute locally. |
| 321 | //! |
| 322 | //! These timestamps become important when handling simultaneous audio from |
| 323 | //! multiple peers. |
| 324 | //! |
| 325 | //! When writing media we need to provide str0m with an estimated wallclock. |
| 326 | //! The simplest strategy is to only trust local time and use arrival time |
| 327 | //! of the incoming UDP packet. Another simple strategy is to lock some |
| 328 | //! time T at the first UDP packet, and then offset each wallclock using |
| 329 | //! `MediaTime`, i.e. for video we could have `T + <media time>/90_000` |
| 330 | //! |
| 331 | //! A production worthy SFU probably needs an even more sophisticated |
| 332 | //! strategy weighing in all possible time sources to get a good estimate |
| 333 | //! of the remote wallclock for a packet. |
| 334 | //! |
| 335 | //! # Crypto backends |
| 336 | //! |
| 337 | //! str0m supports multiple crypto backends via feature flags. The default is `aws-lc-rs`. |
| 338 | //! |
| 339 | //! | Feature | Crate | DTLS | Platforms | |
| 340 | //! |-------------------|-----------------------|----------------------------------|-----------| |
| 341 | //! | `aws-lc-rs` | `str0m-aws-lc-rs` | dimpl + AWS-LC-RS | All | |
| 342 | //! | `rust-crypto` | `str0m-rust-crypto` | dimpl + RustCrypto | All | |
| 343 | //! | `openssl` | `str0m-openssl` | OpenSSL (DTLS 1.2 only) | All | |
| 344 | //! | `openssl-dimpl` | `str0m-openssl` | dimpl + OpenSSL crypto | All | |
| 345 | //! | `apple-crypto` | `str0m-apple-crypto` | dimpl + Apple CryptoKit | macOS/iOS | |
| 346 | //! | `wincrypto` | `str0m-wincrypto` | Windows SChannel (DTLS 1.2 only) | Windows | |
| 347 | //! | `wincrypto-dimpl` | `str0m-wincrypto` | dimpl + Windows CNG | Windows | |
| 348 | //! |
| 349 | //! If multiple backend features are enabled, str0m automatically selects the backend in this |
| 350 | //! priority order: `aws-lc-rs`, `rust-crypto`, `openssl-dimpl`, `openssl`, `apple-crypto` |
| 351 | //! (Apple platforms only), `wincrypt-dimpl` (Windows only), `wincrypto` (Windows only). |
| 352 | //! |
| 353 | //! If you disable the default features, you MUST explicitly configure an alternative |
| 354 | //! crypto backend either process-wide or per-instance. |
| 355 | //! |
| 356 | //! ## Process-wide default |
| 357 | //! |
| 358 | //! For applications, the easiest is to set a process-wide default at startup. |
| 359 | //! Note that you can use any backend crate directly without enabling its feature flag: |
| 360 | //! |
| 361 | //! ```no_run |
| 362 | //! // Set process default (will panic if called twice) |
| 363 | //! // No need to enable the "rust-crypto" feature flag |
| 364 | //! str0m_rust_crypto::default_provider().install_process_default(); |
| 365 | //! ``` |
| 366 | //! |
| 367 | //! ## Crypto provider per Rtc instance |
| 368 | //! |
| 369 | //! ```no_run |
| 370 | //! use std::sync::Arc; |
| 371 | //! use std::time::Instant; |
| 372 | //! use str0m::Rtc; |
| 373 | //! |
| 374 | //! let rtc = Rtc::builder() |
| 375 | //! .set_crypto_provider(Arc::new(str0m_rust_crypto::default_provider())) |
| 376 | //! .build(Instant::now()); |
| 377 | //! ``` |
| 378 | //! |
| 379 | //! # Project status |
| 380 | //! |
| 381 | //! Str0m was originally developed by Martin Algesten of |
| 382 | //! [Lookback][lookback]. We use str0m for a specific use case: str0m as a |
| 383 | //! server SFU (as opposed to peer-2-peer). That means we are heavily |
| 384 | //! testing and developing the parts needed for our use case. Str0m is |
| 385 | //! intended to be an all-purpose WebRTC library, which means it also |
| 386 | //! works for peer-2-peer, though that aspect has received less testing. |
| 387 | //! |
| 388 | //! Performance is very good, there have been some work the discover and |
| 389 | //! optimize bottlenecks. Such efforts are of course never ending with |
| 390 | //! diminishing returns. While there are no glaringly obvious performance |
| 391 | //! bottlenecks, more work is always welcome – both algorithmically and |
| 392 | //! allocation/cloning in hot paths etc. |
| 393 | //! |
| 394 | //! # Design |
| 395 | //! |
| 396 | //! Output from the `Rtc` instance can be grouped into three kinds. |
| 397 | //! |
| 398 | //! 1. Events (such as receiving media or data channel data). |
| 399 | //! 2. Network output. Data to be sent, typically from a UDP socket. |
| 400 | //! 3. Timeouts. Indicates when the instance next expects a time input. |
| 401 | //! |
| 402 | //! Input to the `Rtc` instance is: |
| 403 | //! |
| 404 | //! 1. User operations (such as sending media or data channel data). |
| 405 | //! 2. Network input. Typically read from a UDP socket. |
| 406 | //! 3. Timeouts. As obtained from the output above. |
| 407 | //! |
| 408 | //! The correct use can be seen in the above [Run loop](#run-loop) or in the |
| 409 | //! examples. |
| 410 | //! |
| 411 | //! Sans I/O is a pattern where we turn both network input/output as well |
| 412 | //! as time passing into external input to the API. This means str0m has |
| 413 | //! no internal threads, just an enormous state machine that is driven |
| 414 | //! forward by different kinds of input. |
| 415 | //! |
| 416 | //! ## Frame or RTP level? |
| 417 | //! |
| 418 | //! Str0m defaults to the "frame level" which treats the RTP as an internal detail. The user |
| 419 | //! will thus mainly interact with: |
| 420 | //! |
| 421 | //! 1. [`Event::MediaData`][evmed] to receive full frames (audio frames or video frames). |
| 422 | //! 2. [`Writer::write`][writer] to write full frames. |
| 423 | //! 3. [`Writer::request_keyframe`][reqkey] to request keyframes. |
| 424 | //! |
| 425 | //! ### Frame level |
| 426 | //! |
| 427 | //! All codecs such as h264, vp8, vp9 and opus outputs what we call |
| 428 | //! "Frames". A frame has a very specific meaning for video, but this |
| 429 | //! project uses it in a broader sense, where a frame is either a video |
| 430 | //! or audio time stamped chunk of encoded data that typically represents |
| 431 | //! a chunk of audio, or _one single frame for video_. |
| 432 | //! |
| 433 | //! Frames are not suitable to use directly in UDP (RTP) packets - for |
| 434 | //! one they are too big. Frames are therefore further chunked up by |
| 435 | //! codec specific payloaders into RTP packets. |
| 436 | //! |
| 437 | //! ### RTP mode |
| 438 | //! |
| 439 | //! Str0m also provides an RTP level API. This would be similar to many other |
| 440 | //! RTP libraries where the RTP packets themselves are the API surface |
| 441 | //! towards the user (when building an SFU one would often talk about "forwarding |
| 442 | //! RTP packets", while with str0m we can also "forward frames"). Using |
| 443 | //! this API requires a deeper knowledge of RTP and WebRTC. |
| 444 | //! |
| 445 | //! To enable RTP mode |
| 446 | //! |
| 447 | //! ``` |
| 448 | //! # use std::time::Instant; |
| 449 | //! # use str0m::Rtc; |
| 450 | //! let rtc = Rtc::builder() |
| 451 | //! // Enable RTP mode for this Rtc instance. |
| 452 | //! // This disables `MediaEvent` and the `Writer::write` API. |
| 453 | //! .set_rtp_mode(true) |
| 454 | //! .build(Instant::now()); |
| 455 | //! ``` |
| 456 | //! |
| 457 | //! RTP mode gives us some new API points. |
| 458 | //! |
| 459 | //! 1. [`Event::RtpPacket`][rtppak] emitted for every incoming RTP packet. Empty packets for bandwidth |
| 460 | //! estimation are silently discarded. |
| 461 | //! 2. [`StreamTx::write_rtp`][wrtrtp] to write outgoing RTP packets. |
| 462 | //! 3. [`StreamRx::request_keyframe`][reqkey2] to request keyframes from remote. |
| 463 | //! |
| 464 | //! ## NIC enumeration and TURN (and STUN) |
| 465 | //! |
| 466 | //! The [ICE RFC][ice] talks about "gathering ice candidates". This means |
| 467 | //! inspecting the local network interfaces and potentially binding UDP |
| 468 | //! sockets on each usable interface. Since str0m is Sans I/O, this part |
| 469 | //! is outside the scope of what str0m does. How the user figures out |
| 470 | //! local IP addresses, via config or via looking up local NICs is not |
| 471 | //! something str0m cares about. |
| 472 | //! |
| 473 | //! TURN is a way of obtaining IP addresses that can be used as fallback |
| 474 | //! in case direct connections fail. We consider TURN similar to |
| 475 | //! enumerating local network interfaces – it's a way of obtaining |
| 476 | //! sockets. |
| 477 | //! |
| 478 | //! All discovered candidates, be they local (NIC) or remote sockets |
| 479 | //! (TURN), are added to str0m and str0m will perform the task of ICE |
| 480 | //! agent, forming "candidate pairs" and figuring out the best connection |
| 481 | //! while the actual task of sending the network traffic is left to the |
| 482 | //! user. |
| 483 | //! |
| 484 | //! ## The importance of `&mut self` |
| 485 | //! |
| 486 | //! Rust shines when we can eschew locks and heavily rely `&mut` for data |
| 487 | //! write access. Since str0m has no internal threads, we never have to |
| 488 | //! deal with shared data. Furthermore the the internals of the library is |
| 489 | //! organized such that we don't need multiple references to the same |
| 490 | //! entities. In str0m there are no `Rc`, `Mutex`, `mpsc`, `Arc`(*), or |
| 491 | //! other locks. |
| 492 | //! |
| 493 | //! This means all input to the lib can be modelled as |
| 494 | //! `handle_something(&mut self, something)`. |
| 495 | //! |
| 496 | //! (*) Ok. There is one `Arc` if you use Windows where we also require openssl. |
| 497 | //! |
| 498 | //! ## Not a standard WebRTC "Peer Connection" API |
| 499 | //! |
| 500 | //! The library deliberately steps away from the "standard" WebRTC API as |
| 501 | //! seen in JavaScript and/or [webrtc-rs][webrtc-rs] (or [Pion][pion] in Go). |
| 502 | //! There are few reasons for this. |
| 503 | //! |
| 504 | //! First, in the standard API, events are callbacks, which are not a |
| 505 | //! great fit for Rust. Callbacks require some kind of reference |
| 506 | //! (ownership?) over the entity the callback is being dispatched |
| 507 | //! upon. I.e. if in Rust we want `pc.addEventListener(x)`, `x` needs |
| 508 | //! to be wholly owned by `pc`, or have some shared reference (like |
| 509 | //! `Arc`). Shared references means shared data, and to get mutable shared |
| 510 | //! data, we will need some kind of lock. i.e. `Arc<Mutex<EventListener>>` |
| 511 | //! or similar. |
| 512 | //! |
| 513 | //! As an alternative we could turn all events into `mpsc` channels, but |
| 514 | //! listening to multiple channels is awkward without async. |
| 515 | //! |
| 516 | //! Second, in the standard API, entities like `RTCPeerConnection` and |
| 517 | //! `RTCRtpTransceiver`, are easily clonable and/or long lived |
| 518 | //! references. I.e. `pc.getTranscievers()` returns objects that can be |
| 519 | //! retained and owned by the caller. This pattern is fine for garbage |
| 520 | //! collected or reference counted languages, but not great with Rust. |
| 521 | //! |
| 522 | //! ## Panics, Errors and unwraps |
| 523 | //! |
| 524 | //! Str0m adheres to [fail-fast][ff]. That means rather than brushing state |
| 525 | //! bugs under the carpet, it panics. We make a distinction between errors and |
| 526 | //! bugs. |
| 527 | //! |
| 528 | //! * Errors are as a result of incorrect or impossible to understand user input. |
| 529 | //! * Bugs are broken internal invariants (assumptions). |
| 530 | //! |
| 531 | //! If you scan the str0m code you find a few `unwrap()` (or `expect()`). These |
| 532 | //! will (should) always be accompanied by a code comment that explains why the |
| 533 | //! unwrap is okay. This is an internal invariant, a state assumption that |
| 534 | //! str0m is responsible for maintaining. |
| 535 | //! |
| 536 | //! We do not believe it's correct to change every `unwrap()`/`expect()` into |
| 537 | //! `unwrap_or_else()`, `if let Some(x) = x { ... }` etc, because doing so |
| 538 | //! brushes an actual problem (an incorrect assumption) under the carpet. Trying |
| 539 | //! to hobble along with an incorrect state would at best result in broken |
| 540 | //! behavior, at worst a security risk! |
| 541 | //! |
| 542 | //! Panics are our friends: *panic means bug* |
| 543 | //! |
| 544 | //! And also: str0m should *never* panic on any user input. If you encounter a panic, |
| 545 | //! please report it! |
| 546 | //! |
| 547 | //! ### Catching panics |
| 548 | //! |
| 549 | //! Panics should be incredibly rare, or we have a serious problem as a project. For an SFU, |
| 550 | //! it might not be ideal if str0m encounters a bug and brings the entire server down with it. |
| 551 | //! |
| 552 | //! For those who want an extra level of safety, we recommend looking at [`catch_unwind`][catch] |
| 553 | //! to safely discard a faulty `Rtc` instance. Since `Rtc` has no internal threads, locks or async |
| 554 | //! tasks, discarding the instance never risk poisoning locks or other issues that can happen |
| 555 | //! when catching a panic. |
| 556 | //! |
| 557 | //! ## FAQ |
| 558 | //! |
| 559 | //! ### Features |
| 560 | //! |
| 561 | //! Below is a brief comparison of features between libWebRTC and str0m to help you determine |
| 562 | //! if str0m is suitable for your project. |
| 563 | //! |
| 564 | //! | Feature | str0m | libWebRTC | |
| 565 | //! | ------------------------ | ------------------ | ------------------ | |
| 566 | //! | Peer Connection API | :x: | :white_check_mark: | |
| 567 | //! | SDP | :white_check_mark: | :white_check_mark: | |
| 568 | //! | ICE | :white_check_mark: | :white_check_mark: | |
| 569 | //! | Data Channels | :white_check_mark: | :white_check_mark: | |
| 570 | //! | Send/Recv Reports | :white_check_mark: | :white_check_mark: | |
| 571 | //! | Transport Wide CC | :white_check_mark: | :white_check_mark: | |
| 572 | //! | Bandwidth Estimation | :white_check_mark: | :white_check_mark: | |
| 573 | //! | Simulcast | :white_check_mark: | :white_check_mark: | |
| 574 | //! | NACK | :white_check_mark: | :white_check_mark: | |
| 575 | //! | Packetize | :white_check_mark: | :white_check_mark: | |
| 576 | //! | Fixed Depacketize Buffer | :white_check_mark: | :white_check_mark: | |
| 577 | //! | Adaptive Jitter Buffer | :x: | :white_check_mark: | |
| 578 | //! | Video/audio capture | :x: | :white_check_mark: | |
| 579 | //! | Video/audio encode | :x: | :white_check_mark: | |
| 580 | //! | Video/audio decode | :x: | :white_check_mark: | |
| 581 | //! | Audio render | :x: | :white_check_mark: | |
| 582 | //! | Turn | :x: | :white_check_mark: | |
| 583 | //! | Network interface enum | :x: | :white_check_mark: | |
| 584 | //! |
| 585 | //! ### Platform Support |
| 586 | //! |
| 587 | //! Platforms str0m is compiled and tested on: |
| 588 | //! |
| 589 | //! | Platform | Compiled | Tested | |
| 590 | //! | ------------------------------ | ----------------- | ----------------- | |
| 591 | //! | `x86_64-pc-windows-msvc` | :white_check_mark:| :white_check_mark:| |
| 592 | //! | `x86_64-unknown-linux-gnu` | :white_check_mark:| :white_check_mark:| |
| 593 | //! | `x86_64-apple-darwin` | :white_check_mark:| :white_check_mark:| |
| 594 | //! | `aarch64-apple-darwin` | :white_check_mark:| :white_check_mark:| |
| 595 | //! | `aarch64-unknown-linux-gnu` | :white_check_mark:| :white_check_mark:| |
| 596 | //! | `aarch64-pc-windows-msvc` | :white_check_mark:| :white_check_mark:| |
| 597 | //! | `aarch64-apple-ios` | :white_check_mark:| :x: | |
| 598 | //! | `aarch64-linux-android` | :white_check_mark:| :x: | |
| 599 | //! |
| 600 | //! If your platform isn't listed but is supported by Rust, we'd love for you to give str0m a try and |
| 601 | //! share your experience. We greatly appreciate your feedback! |
| 602 | //! |
| 603 | //! ### Does str0m support IPv4, IPv6, UDP and TCP? |
| 604 | //! |
| 605 | //! Certainly! str0m fully support IPv4, IPv6, UDP and TCP protocols. |
| 606 | //! |
| 607 | //! ### Can I utilize str0m with any Rust async runtime? |
| 608 | //! |
| 609 | //! Absolutely! str0m is fully sync, ensuring that it integrates seamlessly with any Rust async |
| 610 | //! runtime you opt for. |
| 611 | //! |
| 612 | //! ### Can I create a client with str0m? |
| 613 | //! |
| 614 | //! Of course! You have the freedom to create a client with str0m. However, please note that some |
| 615 | //! common client features like media encoding, decoding, and capture are not included in str0m. But |
| 616 | //! don't let that stop you from building amazing applications! |
| 617 | //! |
| 618 | //! ### Can I use str0m in a media server? |
| 619 | //! |
| 620 | //! Yes! str0m excels as a server component with support for both RTP API and Frame API. You can |
| 621 | //! easily build that recording server or SFU you dreamt of in Rust! |
| 622 | //! |
| 623 | //! ### Can I deploy the chat example into production? |
| 624 | //! |
| 625 | //! While the chat example showcases how to use str0m's API, it's not intended for production use or |
| 626 | //! heavy load. Writing a full-featured SFU or MCU (Multipoint Control Unit) is a significant |
| 627 | //! undertaking, involving various design decisions based on production requirements. |
| 628 | //! |
| 629 | //! ### Discovered a bug? Here's how to share it with us |
| 630 | //! |
| 631 | //! We'd love to hear about it! Please submit an issue and consider joining our Discord community |
| 632 | //! to discuss further. For a seamless reporting experience, refer to this exemplary |
| 633 | //! bug report: <https://github.com/algesten/str0m/issues/382>. We appreciate your contribution |
| 634 | //! to making str0m better! |
| 635 | //! |
| 636 | //! ### I am allergic to SDP can you help me? |
| 637 | //! |
| 638 | //! Yes use the direct API! |
| 639 | //! |
| 640 | //! [sansio]: https://sans-io.readthedocs.io |
| 641 | //! [quinn]: https://github.com/quinn-rs/quinn |
| 642 | //! [pion]: https://github.com/pion/webrtc |
| 643 | //! [webrtc-rs]: https://github.com/webrtc-rs/webrtc |
| 644 | //! [discord]: https://discord.gg/e2CC8UYebP |
| 645 | //! [zulip]: https://str0m.zulipchat.com |
| 646 | //! [ice]: https://www.rfc-editor.org/rfc/rfc8445 |
| 647 | //! [lookback]: https://www.lookback.com |
| 648 | //! [x-post]: https://github.com/algesten/str0m/blob/main/examples/http-post.rs |
| 649 | //! [x-chat]: https://github.com/algesten/str0m/blob/main/examples/chat.rs |
| 650 | //! [intg]: https://github.com/algesten/str0m/blob/main/tests/unidirectional.rs#L12 |
| 651 | //! [ff]: https://en.wikipedia.org/wiki/Fail-fast |
| 652 | //! [catch]: https://doc.rust-lang.org/std/panic/fn.catch_unwind.html |
| 653 | //! [evmed]: https://docs.rs/str0m/*/str0m/enum.Event.html#variant.MediaData |
| 654 | //! [writer]: https://docs.rs/str0m/*/str0m/media/struct.Writer.html#method.write |
| 655 | //! [reqkey]: https://docs.rs/str0m/*/str0m/media/struct.Writer.html#method.request_keyframe |
| 656 | //! [rtppak]: https://docs.rs/str0m/*/str0m/enum.Event.html#variant.RtpPacket |
| 657 | //! [wrtrtp]: https://docs.rs/str0m/*/str0m/rtp/struct.StreamTx.html#method.write_rtp |
| 658 | //! [reqkey2]: https://docs.rs/str0m/*/str0m/rtp/struct.StreamRx.html#method.request_keyframe |
| 659 | //! [bitwhip]: https://github.com/bitwhip/bitwhip |
| 660 | |
| 661 | #![forbid(unsafe_code)] |
| 662 | #![allow(clippy::new_without_default)] |
| 663 | #![allow(clippy::bool_to_int_with_if)] |
| 664 | #![allow(clippy::assertions_on_constants)] |
| 665 | #![allow(clippy::manual_range_contains)] |
| 666 | #![allow(clippy::get_first)] |
| 667 | #![allow(clippy::needless_lifetimes)] |
| 668 | #![allow(clippy::precedence)] |
| 669 | #![allow(clippy::doc_overindented_list_items)] |
| 670 | #![allow(clippy::uninlined_format_args)] |
| 671 | #![allow(unknown_lints, mismatched_lifetime_syntaxes)] |
| 672 | #![deny(clippy::needless_pass_by_ref_mut)] |
| 673 | #![deny(missing_docs)] |
| 674 | |
| 675 | #[macro_use] |
| 676 | extern crate tracing; |
| 677 | |
| 678 | use bwe::{Bwe, BweKind}; |
| 679 | use change::{DirectApi, SdpApi}; |
| 680 | use rtp::RawPacket; |
| 681 | use std::fmt; |
| 682 | use std::net::SocketAddr; |
| 683 | use std::sync::Arc; |
| 684 | use std::time::Instant; |
| 685 | use str0m_proto::Pii; |
| 686 | use streams::RtpPacket; |
| 687 | use streams::StreamPaused; |
| 688 | use util::InstantExt; |
| 689 | |
| 690 | // Identity `drv::ToStatic` (`Static = Self`) for the `Copy` identity types. |
| 691 | // Used instead of `#[derive(drv::Input)]` because the derive's generated |
| 692 | // shadow type isn't `Eq`/`Hash`/`Borrow<Self>`, so it can't be a `HashMap` |
| 693 | // key inside a `#[drv::memo]` projection; `Static = Self` can. `eq_static` |
| 694 | // defers to each type's own `PartialEq`. Each module invokes this for its |
| 695 | // own types. |
| 696 | #[cfg(feature = "drv")] |
| 697 | macro_rules! drv_identity_copy { |
| 698 | ($($t:ty),* $(,)?) => { |
| 699 | $( |
| 700 | impl drv::ToStatic for $t { |
| 701 | type Static = $t; |
| 702 | fn to_static(&self) -> $t { *self } |
| 703 | fn eq_static(&self, other: &$t) -> bool { self == other } |
| 704 | } |
| 705 | )* |
| 706 | }; |
| 707 | } |
| 708 | #[cfg(feature = "drv")] |
| 709 | pub(crate) use drv_identity_copy; |
| 710 | |
| 711 | /// Cryptographic provider traits and implementations. |
| 712 | /// |
| 713 | /// This module provides the traits for pluggable cryptographic operations |
| 714 | /// used in DTLS, SRTP, and STUN. |
| 715 | pub mod crypto; |
| 716 | use crypto::Fingerprint; |
| 717 | |
| 718 | mod dtls; |
| 719 | use crate::crypto::dtls::DtlsOutput; |
| 720 | use crate::crypto::{CryptoProvider, DtlsError, from_feature_flags}; |
| 721 | use crate::dtls::is_would_block; |
| 722 | use dtls::Dtls; |
| 723 | |
| 724 | use is::IceAgent; |
| 725 | use is::IceAgentEvent; |
| 726 | pub use is::{Candidate, CandidateBuilder, CandidateKind, IceConnectionState, IceCreds}; |
| 727 | |
| 728 | #[path = "config.rs"] |
| 729 | mod config_mod; |
| 730 | pub use config_mod::RtcConfig; |
| 731 | |
| 732 | /// Default target MTU used when none is configured. |
| 733 | pub use io::DATAGRAM_MTU_TARGET; |
| 734 | /// Upper bound that the target endpoint of [`RtcConfig::set_mtu`] must satisfy. |
| 735 | pub use io::DATAGRAM_MTU_TARGET_MAX; |
| 736 | /// Lower bound that the target endpoint of [`RtcConfig::set_mtu`] must satisfy. |
| 737 | pub use io::DATAGRAM_MTU_TARGET_MIN; |
| 738 | /// Default warning threshold for over-sized packets. |
| 739 | pub use io::DATAGRAM_MTU_WARN; |
| 740 | |
| 741 | /// Additional configuration. |
| 742 | pub mod config { |
| 743 | pub use super::crypto::dtls::{DtlsCert, DtlsVersion, KeyingMaterial}; |
| 744 | pub use super::crypto::{CryptoProvider, Fingerprint}; |
| 745 | } |
| 746 | |
| 747 | /// Low level ICE access. |
| 748 | // The ICE API is not necessary to interact with directly for "regular" |
| 749 | // use of str0m. This is exported for other libraries that want to |
| 750 | // reuse str0m's ICE implementation. This is now in the `is` crate. |
| 751 | #[doc(hidden)] |
| 752 | pub mod ice { |
| 753 | pub use is::IceCreds; |
| 754 | pub use is::stun::{StunMessage, StunMessageBuilder, StunPacket, TransId}; |
| 755 | pub use is::{IceAgent, IceAgentEvent}; |
| 756 | pub use is::{LocalPreference, default_local_preference}; |
| 757 | } |
| 758 | |
| 759 | mod io; |
| 760 | use io::DatagramRecvInner; |
| 761 | |
| 762 | mod packet; |
| 763 | |
| 764 | #[path = "rtp/mod.rs"] |
| 765 | mod rtp_; |
| 766 | use rtp_::{Bitrate, DataSize}; |
| 767 | |
| 768 | /// Low level RTP access. |
| 769 | pub mod rtp { |
| 770 | /// Feedback for RTP. |
| 771 | pub mod rtcp { |
| 772 | pub use crate::rtp_::AppSpecificFeedback; |
| 773 | pub use crate::rtp_::{Descriptions, ExtendedReport, Fir, Goodbye, Nack, Pli}; |
| 774 | pub use crate::rtp_::{Dlrr, NackEntry, ReceptionReport, ReportBlock}; |
| 775 | pub use crate::rtp_::{FirEntry, ReceiverReport, SenderInfo, SenderReport, Twcc}; |
| 776 | pub use crate::rtp_::{ReportList, Rrtr, Rtcp, Sdes, SdesType}; |
| 777 | } |
| 778 | use self::rtcp::Rtcp; |
| 779 | |
| 780 | /// Video Layers Allocation RTP Header Extension |
| 781 | pub mod vla; |
| 782 | pub use crate::rtp_::{AbsCaptureTime, ExtensionValues, UserExtensionValues}; |
| 783 | pub use crate::rtp_::{Extension, ExtensionMap, ExtensionSerializer}; |
| 784 | |
| 785 | pub use crate::packet::{ |
| 786 | Vp8Descriptor, Vp8DescriptorError, Vp8Patch, Vp8PatchBuilder, Vp8PatchError, |
| 787 | }; |
| 788 | pub use crate::rtp_::{RtpHeader, SeqNo, Ssrc, VideoOrientation}; |
| 789 | pub use crate::streams::{ |
| 790 | RtpPacket, RtpWrite, StreamPaused, StreamRx, StreamTx, StreamTxQueueInfo, |
| 791 | }; |
| 792 | |
| 793 | /// Debug output of the unencrypted RTP and RTCP packets. |
| 794 | /// |
| 795 | /// Enable using [`RtcConfig::enable_raw_packets()`][crate::RtcConfig::enable_raw_packets]. |
| 796 | /// This clones data, and is therefore expensive. |
| 797 | /// Should not be enabled outside of tests and troubleshooting. |
| 798 | #[derive(Debug)] |
| 799 | pub enum RawPacket { |
| 800 | /// Sent RTCP. |
| 801 | RtcpTx(Rtcp), |
| 802 | /// Incoming RTCP. |
| 803 | RtcpRx(Rtcp), |
| 804 | /// Sent RTP. |
| 805 | RtpTx(RtpHeader, Vec<u8>), |
| 806 | /// Incoming RTP. |
| 807 | RtpRx(RtpHeader, Vec<u8>), |
| 808 | } |
| 809 | } |
| 810 | |
| 811 | pub(crate) mod pacer; |
| 812 | |
| 813 | #[path = "bwe/mod.rs"] |
| 814 | pub(crate) mod bwe_; |
| 815 | |
| 816 | /// Bandwidth estimation. |
| 817 | pub mod bwe { |
| 818 | pub use crate::bwe_::api::*; |
| 819 | } |
| 820 | |
| 821 | mod sctp; |
| 822 | use sctp::{RtcSctp, SctpEvent, SctpInitData}; |
| 823 | |
| 824 | mod sdp; |
| 825 | |
| 826 | pub mod format; |
| 827 | use format::CodecConfig; |
| 828 | |
| 829 | pub mod channel; |
| 830 | use channel::{Channel, ChannelData, ChannelHandler, ChannelId}; |
| 831 | |
| 832 | pub mod media; |
| 833 | use media::AppSpecificFeedback; |
| 834 | use media::SenderFeedback; |
| 835 | use media::{Direction, Media, Mid, Pt, Rid, Writer}; |
| 836 | use media::{KeyframeRequest, KeyframeRequestKind}; |
| 837 | use media::{MediaAdded, MediaChanged, MediaData}; |
| 838 | |
| 839 | pub mod change; |
| 840 | |
| 841 | mod util; |
| 842 | use util::{Soonest, not_happening}; |
| 843 | |
| 844 | mod session; |
| 845 | use session::Session; |
| 846 | |
| 847 | pub mod stats; |
| 848 | |
| 849 | use stats::{CandidatePairStats, CandidateStats, MediaEgressStats, MediaIngressStats}; |
| 850 | use stats::{PeerStats, Stats, StatsEvent, StatsSnapshot}; |
| 851 | |
| 852 | mod streams; |
| 853 | |
| 854 | pub mod error; |
| 855 | |
| 856 | /// Network related types to get socket data in/out of [`Rtc`]. |
| 857 | pub mod net { |
| 858 | pub use crate::io::{DatagramRecv, DatagramSend, Protocol, Receive, TcpType, Transmit}; |
| 859 | } |
| 860 | |
| 861 | const VERSION: &str = env!("CARGO_PKG_VERSION"); |
| 862 | |
| 863 | pub use error::RtcError; |
| 864 | |
| 865 | /// Instance that does WebRTC. Main struct of the entire library. |
| 866 | /// |
| 867 | /// ## Usage |
| 868 | /// |
| 869 | /// The canonical run loop is: perform exactly one mutation (typically |
| 870 | /// `handle_input` with a network packet or timeout), then drain |
| 871 | /// `poll_output` until it returns `Output::Timeout`, then wait for the |
| 872 | /// next event. **Every mutation must be followed by a full drain before |
| 873 | /// the next mutation.** See the [crate docs](crate#run-loop) for the full |
| 874 | /// explanation of the single-mutation invariant. |
| 875 | /// |
| 876 | /// ```no_run |
| 877 | /// # use std::time::Instant; |
| 878 | /// # use str0m::{Rtc, Output, Input}; |
| 879 | /// let mut rtc = Rtc::new(Instant::now()); |
| 880 | /// |
| 881 | /// loop { |
| 882 | /// // Drain poll_output until it returns the next timeout. |
| 883 | /// let timeout = loop { |
| 884 | /// match rtc.poll_output().unwrap() { |
| 885 | /// Output::Timeout(t) => break t, |
| 886 | /// Output::Transmit(t) => { |
| 887 | /// // TODO: Send t.contents to t.destination on a UDP socket. |
| 888 | /// } |
| 889 | /// Output::Event(e) => { |
| 890 | /// // TODO: Handle event. |
| 891 | /// } |
| 892 | /// } |
| 893 | /// }; |
| 894 | /// |
| 895 | /// // TODO: Wait for ONE of: `timeout` firing, an incoming network |
| 896 | /// // packet, or application-side data to send. Build an |
| 897 | /// // `Input` from whichever fires first. |
| 898 | /// let input: Input = todo!(); |
| 899 | /// |
| 900 | /// rtc.handle_input(input).unwrap(); |
| 901 | /// } |
| 902 | /// ``` |
| 903 | pub struct Rtc { |
| 904 | state: RtcState, |
| 905 | ice: IceAgent, |
| 906 | dtls: Dtls, |
| 907 | dtls_connected: bool, |
| 908 | dtls_buf: Vec<u8>, |
| 909 | next_dtls_timeout: Option<Instant>, |
| 910 | sctp: RtcSctp, |
| 911 | chan: ChannelHandler, |
| 912 | stats: Option<Stats>, |
| 913 | session: Session, |
| 914 | remote_fingerprint: Option<Fingerprint>, |
| 915 | remote_addrs: Vec<SocketAddr>, |
| 916 | send_addr: Option<SendAddr>, |
| 917 | need_init_time: bool, |
| 918 | last_now: Instant, |
| 919 | peer_bytes_rx: u64, |
| 920 | peer_bytes_tx: u64, |
| 921 | change_counter: usize, |
| 922 | last_timeout_reason: Reason, |
| 923 | crypto_provider: Arc<crate::crypto::CryptoProvider>, |
| 924 | fingerprint_verification: bool, |
| 925 | close_dtls_started: bool, |
| 926 | } |
| 927 | |
| 928 | #[derive(Clone, Copy, PartialEq, Eq)] |
| 929 | enum RtcState { |
| 930 | Alive, |
| 931 | Closing, |
| 932 | Closed, |
| 933 | } |
| 934 | |
| 935 | struct SendAddr { |
| 936 | proto: net::Protocol, |
| 937 | source: SocketAddr, |
| 938 | destination: SocketAddr, |
| 939 | } |
| 940 | |
| 941 | /// Events produced by [`Rtc::poll_output()`]. |
| 942 | #[derive(Debug)] |
| 943 | #[non_exhaustive] |
| 944 | #[rustfmt::skip] |
| 945 | pub enum Event { |
| 946 | // =================== ICE related events =================== |
| 947 | |
| 948 | /// Emitted when we got ICE connection and established DTLS. |
| 949 | Connected, |
| 950 | |
| 951 | /// ICE connection state changes tells us whether the [`Rtc`] instance is |
| 952 | /// connected to the peer or not. |
| 953 | IceConnectionStateChange(IceConnectionState), |
| 954 | |
| 955 | // =================== Media related events ================== |
| 956 | |
| 957 | /// Upon detecting the remote side adding new media to the session. |
| 958 | /// |
| 959 | /// For locally added media, this event never fires. Thus it can be thought of as an |
| 960 | /// "SDP only" event. If the direct API is used on both sides, the declaration is local \ |
| 961 | /// to both sides and the event never fires. |
| 962 | /// |
| 963 | /// The [`Media`] instance is available via [`Rtc::media()`]. |
| 964 | MediaAdded(MediaAdded), |
| 965 | |
| 966 | /// Incoming media data sent by the remote peer. |
| 967 | MediaData(MediaData), |
| 968 | |
| 969 | /// Changes to the media may be emitted. |
| 970 | /// |
| 971 | ///. Currently only covers a change of direction. |
| 972 | MediaChanged(MediaChanged), |
| 973 | |
| 974 | // =================== Data channel related events =================== |
| 975 | |
| 976 | /// A data channel has opened. |
| 977 | /// |
| 978 | /// The string is the channel label which is set by the opening peer and can |
| 979 | /// be used to identify the purpose of the channel when there are more than one. |
| 980 | /// |
| 981 | /// The negotiation is to set up an SCTP association via DTLS. Subsequent data |
| 982 | /// channels reuse the same association. |
| 983 | /// |
| 984 | /// Upon this event, the [`Channel`] can be obtained via [`Rtc::channel()`]. |
| 985 | /// |
| 986 | /// For [`SdpApi`]: The first ever data channel results in an SDP |
| 987 | /// negotiation, and this events comes at the end of that. |
| 988 | ChannelOpen(ChannelId, String), |
| 989 | |
| 990 | /// Incoming data channel data from the remote peer. |
| 991 | ChannelData(ChannelData), |
| 992 | |
| 993 | /// A data channel has been closed. |
| 994 | ChannelClose(ChannelId), |
| 995 | |
| 996 | /// A data channel's buffered amount has dropped below the configured threshold. |
| 997 | ChannelBufferedAmountLow(ChannelId), |
| 998 | |
| 999 | // =================== Statistics and BWE related events =================== |
| 1000 | |
| 1001 | /// Statistics event for the Rtc instance |
| 1002 | /// |
| 1003 | /// Includes both media traffic (rtp payload) as well as all traffic |
| 1004 | PeerStats(PeerStats), |
| 1005 | |
| 1006 | /// Aggregated statistics for each media (mid, rid) in the ingress direction |
| 1007 | MediaIngressStats(MediaIngressStats), |
| 1008 | |
| 1009 | /// Aggregated statistics for each media (mid, rid) in the egress direction |
| 1010 | MediaEgressStats(MediaEgressStats), |
| 1011 | |
| 1012 | /// A new estimate from the bandwidth estimation subsystem. |
| 1013 | EgressBitrateEstimate(BweKind), |
| 1014 | |
| 1015 | // =================== RTP related events =================== |
| 1016 | |
| 1017 | /// Incoming keyframe request for media that we are sending to the remote peer. |
| 1018 | /// |
| 1019 | /// The request is either PLI (Picture Loss Indication) or FIR (Full Intra Request). |
| 1020 | KeyframeRequest(KeyframeRequest), |
| 1021 | |
| 1022 | /// Whether an incoming encoded stream is paused. |
| 1023 | /// |
| 1024 | /// This means the stream has not received any data for some time (default 1.5 seconds). |
| 1025 | StreamPaused(StreamPaused), |
| 1026 | |
| 1027 | /// Sender feedback for an incoming stream, derived from RTCP SR. |
| 1028 | SenderFeedback(SenderFeedback), |
| 1029 | |
| 1030 | /// Incoming RTP data. |
| 1031 | RtpPacket(RtpPacket), |
| 1032 | |
| 1033 | /// Incoming application-specific Payload-Specific Feedback (PSFB FMT=15, PT=206). |
| 1034 | /// |
| 1035 | /// Emitted when a non-REMB FMT=15 RTCP PSFB message is received. The payload |
| 1036 | /// is opaque and application-defined (RFC 4585 Section 6.4). |
| 1037 | AppSpecificFeedback(AppSpecificFeedback), |
| 1038 | |
| 1039 | /// The remote DTLS connection or SCTP association has closed. |
| 1040 | /// |
| 1041 | /// When this event is emitted, `Rtc` starts local shutdown signals and |
| 1042 | /// remains alive only while [`Rtc::poll_output()`] drains pending local |
| 1043 | /// close output such as RTCP BYE and the DTLS close_notify response. |
| 1044 | Closed, |
| 1045 | |
| 1046 | /// Debug output of incoming and outgoing RTCP/RTP packets. |
| 1047 | /// |
| 1048 | /// Enable using [`RtcConfig::enable_raw_packets()`]. |
| 1049 | /// This clones data, and is therefore expensive. |
| 1050 | /// Should not be enabled outside of tests and troubleshooting. |
| 1051 | RawPacket(Box<RawPacket>), |
| 1052 | |
| 1053 | /// For internal testing only. |
| 1054 | /// |
| 1055 | /// The probe cluster config when a probe fires. |
| 1056 | #[cfg(feature = "_internal_test_exports")] |
| 1057 | Probe(crate::bwe_::ProbeClusterConfig), |
| 1058 | } |
| 1059 | |
| 1060 | impl Event { |
| 1061 | /// Reference to the [`RawPacket`] if this is indeed an `Event::RawPacket`. |
| 1062 | pub fn as_raw_packet(&self) -> Option<&RawPacket> { |
| 1063 | if let Self::RawPacket(boxed) = &self { |
| 1064 | Some(&**boxed) |
| 1065 | } else { |
| 1066 | None |
| 1067 | } |
| 1068 | } |
| 1069 | } |
| 1070 | |
| 1071 | /// Input as expected by [`Rtc::handle_input()`]. Either network data or a timeout. |
| 1072 | #[derive(Debug)] |
| 1073 | #[allow(clippy::large_enum_variant)] // We purposely don't want to allocate. |
| 1074 | pub enum Input<'a> { |
| 1075 | /// A timeout without any network input. |
| 1076 | Timeout(Instant), |
| 1077 | /// Network input. The [`Instant`] is the time the packet was received from the socket - |
| 1078 | /// not necessarily "now" (e.g. when UDP demultiplexing runs on a separate thread). It |
| 1079 | /// drives time forward when more recent than the last "now" the instance has seen. |
| 1080 | Receive(Instant, net::Receive<'a>), |
| 1081 | } |
| 1082 | |
| 1083 | /// Output produced by [`Rtc::poll_output()`] |
| 1084 | #[allow(clippy::large_enum_variant)] |
| 1085 | #[derive(Debug)] |
| 1086 | pub enum Output { |
| 1087 | /// When the [`Rtc`] instance expects an [`Input::Timeout`]. |
| 1088 | Timeout(Instant), |
| 1089 | |
| 1090 | /// Network data that is to be sent. |
| 1091 | Transmit(net::Transmit), |
| 1092 | |
| 1093 | /// Some event such as media data arriving from the remote peer or connection events. |
| 1094 | Event(Event), |
| 1095 | } |
| 1096 | |
| 1097 | pub use crate::pacer::PacerReason; |
| 1098 | |
| 1099 | /// The reason for the next [`Output::Timeout`]. |
| 1100 | /// |
| 1101 | /// This enum is not considered stable API and may change in minor revisions. |
| 1102 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)] |
| 1103 | #[non_exhaustive] |
| 1104 | pub enum Reason { |
| 1105 | /// No timeout scheduled. |
| 1106 | /// |
| 1107 | /// The timeout value is in the distant future. |
| 1108 | #[default] |
| 1109 | NotHappening, |
| 1110 | |
| 1111 | /// The DTLS subsystem. |
| 1112 | /// |
| 1113 | /// Only relevant during handshaking. |
| 1114 | DTLS, |
| 1115 | |
| 1116 | /// The ICE agent. |
| 1117 | /// |
| 1118 | /// Includes checking candidate pairs and various cleanups. |
| 1119 | Ice, |
| 1120 | |
| 1121 | /// The SCTP subsystem. |
| 1122 | /// |
| 1123 | /// Things like handling retransmissions and keep-alive checks. |
| 1124 | Sctp, |
| 1125 | |
| 1126 | /// Data channels. |
| 1127 | /// |
| 1128 | /// Scheduled when we need to open allocations using SCTP. |
| 1129 | Channel, |
| 1130 | |
| 1131 | /// Stats gathering (if enabled). |
| 1132 | /// |
| 1133 | /// Periodic gathering of statistics. |
| 1134 | Stats, |
| 1135 | |
| 1136 | /// Regular RTP feedback. |
| 1137 | /// |
| 1138 | /// Receiver reports (RR) and sender reports (SR). |
| 1139 | Feedback, |
| 1140 | |
| 1141 | /// Sending of RTP NACK. |
| 1142 | /// |
| 1143 | /// When missing packets are discovered, a NACK is scheduled. |
| 1144 | Nack, |
| 1145 | |
| 1146 | /// Reporting of TWCC (if enabled). |
| 1147 | /// |
| 1148 | /// All incoming RTP packets are reported using TWCC. Enabled via SDP if both |
| 1149 | /// sides support it. |
| 1150 | Twcc, |
| 1151 | |
| 1152 | /// RTP streams not receiving data goes into a paused state. |
| 1153 | /// |
| 1154 | /// Whenever an RTP receive stream receives data, a new timeout is scheduled. |
| 1155 | PauseCheck, |
| 1156 | |
| 1157 | /// Preprocessing of RTP packets to be sent. |
| 1158 | /// |
| 1159 | /// Housekeeping task in RTP send streams. |
| 1160 | SendStream, |
| 1161 | |
| 1162 | /// Packetizing of media into RTP data (if used). |
| 1163 | /// |
| 1164 | /// Written media data needs packetizing. This is not used in RTP mode. |
| 1165 | Packetize, |
| 1166 | |
| 1167 | /// Pacer doing things. |
| 1168 | Pacer(PacerReason), |
| 1169 | |
| 1170 | /// The delay controller of the BWE subsystem. |
| 1171 | BweDelayControl, |
| 1172 | |
| 1173 | /// The probe controller of the BWE subsystem. |
| 1174 | BweProbeControl, |
| 1175 | |
| 1176 | /// The probe estimator of the BWE subsystem. |
| 1177 | BweProbeEstimator, |
| 1178 | } |
| 1179 | |
| 1180 | impl Rtc { |
| 1181 | /// Creates a new instance with default settings. |
| 1182 | /// |
| 1183 | /// To configure the instance, use [`RtcConfig`]. |
| 1184 | /// |
| 1185 | /// ``` |
| 1186 | /// use std::time::Instant; |
| 1187 | /// use str0m::Rtc; |
| 1188 | /// |
| 1189 | /// let rtc = Rtc::new(Instant::now()); |
| 1190 | /// ``` |
| 1191 | pub fn new(start: Instant) -> Self { |
| 1192 | let config = RtcConfig::default(); |
| 1193 | Self::new_from_config(config, start).expect("Failed to create Rtc from default config") |
| 1194 | } |
| 1195 | |
| 1196 | /// Creates a config builder that configures an [`Rtc`] instance. |
| 1197 | /// |
| 1198 | /// ``` |
| 1199 | /// # use std::time::Instant; |
| 1200 | /// # use str0m::Rtc; |
| 1201 | /// let rtc = Rtc::builder() |
| 1202 | /// .set_ice_lite(true) |
| 1203 | /// .build(Instant::now()); |
| 1204 | /// ``` |
| 1205 | pub fn builder() -> RtcConfig { |
| 1206 | RtcConfig::new() |
| 1207 | } |
| 1208 | |
| 1209 | pub(crate) fn new_from_config(config: RtcConfig, start: Instant) -> Result<Self, RtcError> { |
| 1210 | let crypto_provider = config |
| 1211 | .crypto_provider |
| 1212 | .clone() |
| 1213 | // If crypto_provider is not set in config, check process default |
| 1214 | .or_else(|| CryptoProvider::get_default().cloned().map(Arc::new)) |
| 1215 | // Or fall back on feature flags |
| 1216 | .or_else(|| Some(Arc::new(from_feature_flags()))) |
| 1217 | // from_feature_flags panics already, so we should never see |
| 1218 | // this expect message. |
| 1219 | .expect("a crash earlier if no crypto provider was set"); |
| 1220 | |
| 1221 | let session = Session::new(&config); |
| 1222 | |
| 1223 | // Capture before any partial moves of `config` below. |
| 1224 | let mtu = config.mtu.clone(); |
| 1225 | |
| 1226 | let local_creds = config.local_ice_credentials.unwrap_or_else(IceCreds::new); |
| 1227 | let mut ice = IceAgent::with_hmac(local_creds, crypto_provider.sha1_hmac_provider); |
| 1228 | if config.ice_lite { |
| 1229 | ice.set_ice_lite(config.ice_lite); |
| 1230 | } |
| 1231 | |
| 1232 | if let Some(initial_stun_rto) = config.initial_stun_rto { |
| 1233 | ice.set_initial_stun_rto(initial_stun_rto); |
| 1234 | } |
| 1235 | |
| 1236 | if let Some(max_stun_rto) = config.max_stun_rto { |
| 1237 | ice.set_max_stun_rto(max_stun_rto); |
| 1238 | } |
| 1239 | |
| 1240 | if let Some(max_stun_retransmits) = config.max_stun_retransmits { |
| 1241 | ice.set_max_stun_retransmits(max_stun_retransmits); |
| 1242 | } |
| 1243 | |
| 1244 | ice.set_mtu(mtu.clone()); |
| 1245 | |
| 1246 | let dtls_cert = config |
| 1247 | .dtls_cert |
| 1248 | .or_else(|| crypto_provider.dtls_provider.generate_certificate()) |
| 1249 | .expect( |
| 1250 | "No DTLS certificate provided and the crypto provider cannot generate one. \ |
| 1251 | Either provide a certificate via RtcConfig::set_dtls_cert or use a \ |
| 1252 | crypto provider that supports certificate generation.", |
| 1253 | ); |
| 1254 | |
| 1255 | let mut sctp = RtcSctp::with_receive_limits(*mtu.start(), config.sctp_receive_limits); |
| 1256 | if config.snap_enabled { |
| 1257 | sctp.enable_snap(); |
| 1258 | } |
| 1259 | |
| 1260 | Ok(Rtc { |
| 1261 | state: RtcState::Alive, |
| 1262 | ice, |
| 1263 | dtls: Dtls::new( |
| 1264 | &dtls_cert, |
| 1265 | crypto_provider.dtls_provider, |
| 1266 | crypto_provider.sha256_provider, |
| 1267 | start, |
| 1268 | config.dtls_version, |
| 1269 | mtu, |
| 1270 | ) |
| 1271 | .expect("DTLS to init without problem"), |
| 1272 | dtls_connected: false, |
| 1273 | dtls_buf: vec![0; 2000], |
| 1274 | next_dtls_timeout: None, |
| 1275 | session, |
| 1276 | sctp, |
| 1277 | chan: ChannelHandler::default(), |
| 1278 | stats: config.stats_interval.map(Stats::new), |
| 1279 | remote_fingerprint: None, |
| 1280 | remote_addrs: vec![], |
| 1281 | send_addr: None, |
| 1282 | need_init_time: true, |
| 1283 | last_now: start, |
| 1284 | peer_bytes_rx: 0, |
| 1285 | peer_bytes_tx: 0, |
| 1286 | change_counter: 0, |
| 1287 | last_timeout_reason: Reason::NotHappening, |
| 1288 | crypto_provider, |
| 1289 | fingerprint_verification: config.fingerprint_verification, |
| 1290 | close_dtls_started: false, |
| 1291 | }) |
| 1292 | } |
| 1293 | |
| 1294 | /// Tests if this instance is still working. |
| 1295 | /// |
| 1296 | /// Certain events will straight away disconnect the `Rtc` instance, such as |
| 1297 | /// the DTLS fingerprint from the setup not matching that of the TLS negotiation |
| 1298 | /// (since that would potentially indicate a MITM attack!). |
| 1299 | /// |
| 1300 | /// The instance can be manually disconnected using [`Rtc::disconnect()`]. |
| 1301 | /// |
| 1302 | /// During shutdown, this remains `true` until [`Rtc::poll_output()`] |
| 1303 | /// has drained pending local close output. App-facing operations such as |
| 1304 | /// [`Rtc::writer()`] and [`Rtc::channel()`] are unavailable during that |
| 1305 | /// closing drain. |
| 1306 | /// |
| 1307 | /// ``` |
| 1308 | /// # use std::time::Instant; |
| 1309 | /// # use str0m::Rtc; |
| 1310 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1311 | /// |
| 1312 | /// assert!(rtc.is_alive()); |
| 1313 | /// |
| 1314 | /// rtc.disconnect(); |
| 1315 | /// assert!(!rtc.is_alive()); |
| 1316 | /// ``` |
| 1317 | pub fn is_alive(&self) -> bool { |
| 1318 | self.state != RtcState::Closed |
| 1319 | } |
| 1320 | |
| 1321 | /// Force disconnects the instance making [`Rtc::is_alive()`] return `false`. |
| 1322 | /// |
| 1323 | /// This makes [`Rtc::poll_output`] and [`Rtc::handle_input`] go inert and not |
| 1324 | /// produce anymore network output or events. |
| 1325 | /// |
| 1326 | /// ``` |
| 1327 | /// # use std::time::Instant; |
| 1328 | /// # use str0m::Rtc; |
| 1329 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1330 | /// |
| 1331 | /// rtc.disconnect(); |
| 1332 | /// assert!(!rtc.is_alive()); |
| 1333 | /// ``` |
| 1334 | pub fn disconnect(&mut self) { |
| 1335 | self.state = RtcState::Closed; |
| 1336 | } |
| 1337 | |
| 1338 | /// Add a local ICE candidate. Local candidates are socket addresses the `Rtc` instance |
| 1339 | /// use for communicating with the peer. |
| 1340 | /// |
| 1341 | /// If the candidate is accepted by the `Rtc` instance, it will return `Some` with a reference |
| 1342 | /// to it. You should then signal this candidate to the remote peer. |
| 1343 | /// |
| 1344 | /// This library has no built-in discovery of local network addresses on the host |
| 1345 | /// or NATed addresses via a STUN server or TURN server. The user of the library |
| 1346 | /// is expected to add new local candidates as they are discovered. |
| 1347 | /// |
| 1348 | /// In WebRTC lingo, the `Rtc` instance is permanently in a mode of [Trickle Ice][1]. It's |
| 1349 | /// however advisable to add at least one local candidate before starting the instance. |
| 1350 | /// |
| 1351 | /// ``` |
| 1352 | /// # use std::time::Instant; |
| 1353 | /// # use str0m::{Rtc, Candidate}; |
| 1354 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1355 | /// |
| 1356 | /// let a = "127.0.0.1:5000".parse().unwrap(); |
| 1357 | /// let c = Candidate::host(a, "udp").unwrap(); |
| 1358 | /// |
| 1359 | /// rtc.add_local_candidate(c); |
| 1360 | /// ``` |
| 1361 | /// |
| 1362 | /// [1]: https://www.rfc-editor.org/rfc/rfc8838.txt |
| 1363 | pub fn add_local_candidate(&mut self, c: Candidate) -> Option<&Candidate> { |
| 1364 | self.ice.add_local_candidate(c) |
| 1365 | } |
| 1366 | |
| 1367 | /// Add a remote ICE candidate. Remote candidates are addresses of the peer. |
| 1368 | /// |
| 1369 | /// For [`SdpApi`]: Remote candidates are typically added via |
| 1370 | /// receiving a remote [`SdpOffer`][change::SdpOffer] or [`SdpAnswer`][change::SdpAnswer]. |
| 1371 | /// |
| 1372 | /// However for the case of [Trickle Ice][1], this is the way to add remote candidates |
| 1373 | /// that are "trickled" from the other side. |
| 1374 | /// |
| 1375 | /// ``` |
| 1376 | /// # use std::time::Instant; |
| 1377 | /// # use str0m::{Rtc, Candidate}; |
| 1378 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1379 | /// |
| 1380 | /// let a = "1.2.3.4:5000".parse().unwrap(); |
| 1381 | /// let c = Candidate::host(a, "udp").unwrap(); |
| 1382 | /// |
| 1383 | /// rtc.add_remote_candidate(c); |
| 1384 | /// ``` |
| 1385 | /// |
| 1386 | /// [1]: https://www.rfc-editor.org/rfc/rfc8838.txt |
| 1387 | pub fn add_remote_candidate(&mut self, c: Candidate) { |
| 1388 | self.ice.add_remote_candidate(c); |
| 1389 | } |
| 1390 | |
| 1391 | /// Checks if we are connected. |
| 1392 | /// |
| 1393 | /// This tests if we have ICE connection, DTLS and the SRTP crypto derived contexts are up. |
| 1394 | pub fn is_connected(&self) -> bool { |
| 1395 | self.ice.state().is_connected() && self.dtls_connected && self.session.is_connected() |
| 1396 | } |
| 1397 | |
| 1398 | /// Make changes to the Rtc session via SDP. |
| 1399 | /// |
| 1400 | /// ```no_run |
| 1401 | /// # use std::time::Instant; |
| 1402 | /// # use str0m::Rtc; |
| 1403 | /// # use str0m::media::{MediaKind, Direction}; |
| 1404 | /// # use str0m::change::SdpAnswer; |
| 1405 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1406 | /// |
| 1407 | /// let mut changes = rtc.sdp_api(); |
| 1408 | /// let mid_audio = changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 1409 | /// let mid_video = changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 1410 | /// |
| 1411 | /// let (offer, pending) = changes.apply().unwrap(); |
| 1412 | /// let json = serde_json::to_vec(&offer).unwrap(); |
| 1413 | /// |
| 1414 | /// // Send json OFFER to remote peer. Receive an answer back. |
| 1415 | /// let answer: SdpAnswer = todo!(); |
| 1416 | /// |
| 1417 | /// rtc.sdp_api().accept_answer(pending, answer).unwrap(); |
| 1418 | /// ``` |
| 1419 | pub fn sdp_api(&mut self) -> SdpApi { |
| 1420 | SdpApi::new(self) |
| 1421 | } |
| 1422 | |
| 1423 | /// Makes direct changes to the Rtc session. |
| 1424 | /// |
| 1425 | /// This is a low level API. For "normal" use via SDP, see [`Rtc::sdp_api()`]. |
| 1426 | pub fn direct_api(&mut self) -> DirectApi { |
| 1427 | DirectApi::new(self) |
| 1428 | } |
| 1429 | |
| 1430 | /// Send outgoing media data (frames) or request keyframes. |
| 1431 | /// |
| 1432 | /// Returns `None` if the direction isn't sending (`sendrecv` or `sendonly`). |
| 1433 | /// |
| 1434 | /// ```no_run |
| 1435 | /// # use std::time::Instant; |
| 1436 | /// # use str0m::Rtc; |
| 1437 | /// # use str0m::media::{MediaData, Mid}; |
| 1438 | /// # use str0m::format::PayloadParams; |
| 1439 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1440 | /// |
| 1441 | /// // add candidates, do SDP negotiation |
| 1442 | /// let mid: Mid = todo!(); // obtain mid from Event::MediaAdded. |
| 1443 | /// |
| 1444 | /// // Writer for this mid. |
| 1445 | /// let writer = rtc.writer(mid).unwrap(); |
| 1446 | /// |
| 1447 | /// // Get incoming media data from another peer |
| 1448 | /// let data: MediaData = todo!(); |
| 1449 | /// |
| 1450 | /// // Match incoming PT to an outgoing PT. |
| 1451 | /// let pt = writer.match_params(data.params).unwrap(); |
| 1452 | /// |
| 1453 | /// writer.write(pt, data.network_time, data.time, data.data).unwrap(); |
| 1454 | /// ``` |
| 1455 | /// |
| 1456 | /// This is a frame level API: For RTP level see [`DirectApi::stream_tx()`] |
| 1457 | /// and [`DirectApi::stream_rx()`]. |
| 1458 | /// |
| 1459 | pub fn writer(&mut self, mid: Mid) -> Option<Writer> { |
| 1460 | if self.session.rtp_mode { |
| 1461 | panic!("In rtp_mode use direct_api().stream_tx().write_rtp()"); |
| 1462 | } |
| 1463 | |
| 1464 | if self.state != RtcState::Alive { |
| 1465 | return None; |
| 1466 | } |
| 1467 | |
| 1468 | // This does not catch potential RIDs required to send simulcast, but |
| 1469 | // it's a good start. An error might arise later on RID mismatch. |
| 1470 | self.session.media_by_mid_mut(mid)?; |
| 1471 | |
| 1472 | Some(Writer::new(&mut self.session, mid)) |
| 1473 | } |
| 1474 | |
| 1475 | /// Currently configured media. |
| 1476 | /// |
| 1477 | /// Read only access. Changes are made via [`Rtc::sdp_api()`] or [`Rtc::direct_api()`]. |
| 1478 | pub fn media(&self, mid: Mid) -> Option<&Media> { |
| 1479 | self.session.media_by_mid(mid) |
| 1480 | } |
| 1481 | |
| 1482 | fn init_dtls(&mut self, active: bool) -> Result<(), RtcError> { |
| 1483 | if self.dtls.is_inited() { |
| 1484 | return Ok(()); |
| 1485 | } |
| 1486 | |
| 1487 | debug!("DTLS setup is: {:?}", active); |
| 1488 | self.dtls.set_active(active); |
| 1489 | |
| 1490 | // Initialize the DTLS state (client or server) before any operations |
| 1491 | // This ensures internal state like random (client) or last_now (server) is initialized |
| 1492 | self.dtls.handle_timeout(self.last_now)?; |
| 1493 | |
| 1494 | if active { |
| 1495 | // Drive handshake by sending an empty packet to trigger ClientHello |
| 1496 | let _ = self.dtls.handle_receive(&[]); |
| 1497 | } |
| 1498 | |
| 1499 | Ok(()) |
| 1500 | } |
| 1501 | |
| 1502 | fn try_init_sctp( |
| 1503 | &mut self, |
| 1504 | client: bool, |
| 1505 | sctp_init_data: Option<SctpInitData>, |
| 1506 | remote_max_message_size: Option<u32>, |
| 1507 | ) -> Result<(), RtcError> { |
| 1508 | // If we got an m=application line, ensure we have negotiated the |
| 1509 | // SCTP association with the other side. |
| 1510 | if self.sctp.is_inited() { |
| 1511 | return Ok(()); |
| 1512 | } |
| 1513 | |
| 1514 | self.sctp.init( |
| 1515 | client, |
| 1516 | self.last_now, |
| 1517 | sctp_init_data, |
| 1518 | remote_max_message_size, |
| 1519 | )?; |
| 1520 | Ok(()) |
| 1521 | } |
| 1522 | |
| 1523 | /// Creates a new Mid that is not in the session already. |
| 1524 | pub(crate) fn new_mid(&self) -> Mid { |
| 1525 | loop { |
| 1526 | let mid = Mid::new(); |
| 1527 | if !self.session.has_mid(mid) { |
| 1528 | break mid; |
| 1529 | } |
| 1530 | } |
| 1531 | } |
| 1532 | |
| 1533 | /// Poll the `Rtc` instance for output. Output can be three things, something to _Transmit_ |
| 1534 | /// via a UDP socket (maybe via a TURN server). An _Event_, such as receiving media data, |
| 1535 | /// or a _Timeout_. |
| 1536 | /// |
| 1537 | /// The user of the library is expected to continuously call this function and deal with |
| 1538 | /// the output until it encounters an [`Output::Timeout`] at which point no further output |
| 1539 | /// is produced (if polled again, it will result in just another timeout). |
| 1540 | /// |
| 1541 | /// After exhausting the `poll_output`, the function will only produce more output again |
| 1542 | /// when one of two things happen: |
| 1543 | /// |
| 1544 | /// 1. The polled timeout is reached. |
| 1545 | /// 2. New network input. |
| 1546 | /// |
| 1547 | /// See [`Rtc`] instance documentation for how this is expected to be used in a loop. |
| 1548 | pub fn poll_output(&mut self) -> Result<Output, RtcError> { |
| 1549 | let o = self.do_poll_output()?; |
| 1550 | |
| 1551 | match &o { |
| 1552 | Output::Event(e) => match e { |
| 1553 | Event::ChannelData(_) |
| 1554 | | Event::MediaData(_) |
| 1555 | | Event::RtpPacket(_) |
| 1556 | | Event::SenderFeedback(_) |
| 1557 | | Event::MediaEgressStats(_) |
| 1558 | | Event::MediaIngressStats(_) |
| 1559 | | Event::PeerStats(_) |
| 1560 | | Event::ChannelBufferedAmountLow(_) |
| 1561 | | Event::EgressBitrateEstimate(_) |
| 1562 | | Event::KeyframeRequest(_) => { |
| 1563 | trace!("{:?}", e) |
| 1564 | } |
| 1565 | _ => debug!("{:?}", e), |
| 1566 | }, |
| 1567 | Output::Transmit(t) => { |
| 1568 | self.peer_bytes_tx += t.contents.len() as u64; |
| 1569 | trace!("OUT {:?}", t) |
| 1570 | } |
| 1571 | Output::Timeout(_t) => {} |
| 1572 | } |
| 1573 | |
| 1574 | Ok(o) |
| 1575 | } |
| 1576 | |
| 1577 | fn do_poll_output(&mut self) -> Result<Output, RtcError> { |
| 1578 | if self.state == RtcState::Closed { |
| 1579 | self.last_timeout_reason = Reason::NotHappening; |
| 1580 | return Ok(Output::Timeout(not_happening())); |
| 1581 | } |
| 1582 | |
| 1583 | while let Some(e) = self.ice.poll_event() { |
| 1584 | match e { |
| 1585 | IceAgentEvent::IceRestart(_) => { |
| 1586 | // |
| 1587 | } |
| 1588 | IceAgentEvent::IceConnectionStateChange(v) => { |
| 1589 | return Ok(Output::Event(Event::IceConnectionStateChange(v))); |
| 1590 | } |
| 1591 | IceAgentEvent::DiscoveredRecv { proto, source } => { |
| 1592 | debug!("ICE remote address: {:?}/{:?}", Pii(source), proto); |
| 1593 | self.remote_addrs.push(source); |
| 1594 | while self.remote_addrs.len() > 20 { |
| 1595 | self.remote_addrs.remove(0); |
| 1596 | } |
| 1597 | } |
| 1598 | IceAgentEvent::NominatedSend { |
| 1599 | proto, |
| 1600 | source, |
| 1601 | destination, |
| 1602 | } => { |
| 1603 | debug!( |
| 1604 | "ICE nominated send from: {:?} to: {:?} with protocol {:?}", |
| 1605 | Pii(source), |
| 1606 | Pii(destination), |
| 1607 | proto, |
| 1608 | ); |
| 1609 | self.send_addr = Some(SendAddr { |
| 1610 | proto, |
| 1611 | source, |
| 1612 | destination, |
| 1613 | }); |
| 1614 | } |
| 1615 | } |
| 1616 | } |
| 1617 | |
| 1618 | // Handle DTLS timeout before polling output so any retransmit packets |
| 1619 | // queued by dimpl and the re-armed flight timer are picked up by the |
| 1620 | // poll loop below in the same iteration. |
| 1621 | if let Some(timeout) = self.next_dtls_timeout { |
| 1622 | if timeout <= self.last_now { |
| 1623 | self.next_dtls_timeout = None; |
| 1624 | |
| 1625 | if let Err(error) = self.dtls.handle_timeout(self.last_now) { |
| 1626 | // A failed timer transition leaves DTLS unrecoverable. Close the |
| 1627 | // RTC so continued polling cannot expose the same due timeout again. |
| 1628 | self.disconnect(); |
| 1629 | return Err(error.into()); |
| 1630 | } |
| 1631 | } |
| 1632 | } |
| 1633 | |
| 1634 | // Poll DTLS output - collect packets, handle events |
| 1635 | let mut just_connected = false; |
| 1636 | loop { |
| 1637 | match self.dtls.poll_output(&mut self.dtls_buf) { |
| 1638 | DtlsOutput::Packet(_) => { |
| 1639 | unreachable!("We don't expect DTLS packets here since we use poll_packet"); |
| 1640 | } |
| 1641 | DtlsOutput::Connected => { |
| 1642 | if !self.dtls_connected { |
| 1643 | debug!("DTLS connected"); |
| 1644 | self.dtls_connected = true; |
| 1645 | just_connected = true; |
| 1646 | } |
| 1647 | } |
| 1648 | DtlsOutput::KeyingMaterial(km, profile) => { |
| 1649 | use config::KeyingMaterial; |
| 1650 | let km_bytes = km.as_ref().to_vec(); |
| 1651 | debug!("DTLS set SRTP keying material and profile: {}", profile); |
| 1652 | let active = self.dtls.is_active().expect("DTLS must be inited by now"); |
| 1653 | self.session.set_keying_material( |
| 1654 | KeyingMaterial::new(&km_bytes), |
| 1655 | &self.crypto_provider, |
| 1656 | profile, |
| 1657 | active, |
| 1658 | ); |
| 1659 | } |
| 1660 | DtlsOutput::PeerCert(der) => { |
| 1661 | debug!("DTLS verify remote fingerprint"); |
| 1662 | // Compute fingerprint from peer's DER certificate |
| 1663 | let fingerprint = crate::crypto::Fingerprint { |
| 1664 | hash_func: "sha-256".to_string(), |
| 1665 | bytes: self.crypto_provider.sha256_provider.sha256(der).to_vec(), |
| 1666 | }; |
| 1667 | self.dtls.set_remote_fingerprint(fingerprint.clone()); |
| 1668 | if let Some(expected) = &self.remote_fingerprint { |
| 1669 | if !self.fingerprint_verification { |
| 1670 | debug!("DTLS fingerprint verification disabled"); |
| 1671 | } else if fingerprint != *expected { |
| 1672 | self.disconnect(); |
| 1673 | return Err(RtcError::RemoteSdp("remote fingerprint no match".into())); |
| 1674 | } |
| 1675 | } else { |
| 1676 | self.disconnect(); |
| 1677 | return Err(RtcError::RemoteSdp("no a=fingerprint before dtls".into())); |
| 1678 | } |
| 1679 | } |
| 1680 | DtlsOutput::ApplicationData(data) => { |
| 1681 | self.sctp.handle_input(self.last_now, data); |
| 1682 | } |
| 1683 | DtlsOutput::Timeout(t) => { |
| 1684 | self.next_dtls_timeout = Some(t); |
| 1685 | break; |
| 1686 | } |
| 1687 | DtlsOutput::CloseNotify => { |
| 1688 | self.start_close()?; |
| 1689 | return Ok(Output::Event(Event::Closed)); |
| 1690 | } |
| 1691 | other => { |
| 1692 | return Err(RtcError::Dtls(DtlsError::Io(std::io::Error::other( |
| 1693 | format!("Unexpected DTLS output: {other:?}"), |
| 1694 | )))); |
| 1695 | } |
| 1696 | } |
| 1697 | } |
| 1698 | |
| 1699 | if just_connected { |
| 1700 | return Ok(Output::Event(Event::Connected)); |
| 1701 | } |
| 1702 | |
| 1703 | while let Some(e) = self.sctp.poll() { |
| 1704 | match e { |
| 1705 | SctpEvent::Transmit { mut packets } => { |
| 1706 | if let Some(v) = packets.front() { |
| 1707 | if let Err(e) = self.dtls.handle_input(v) { |
| 1708 | if is_would_block(&e) { |
| 1709 | self.sctp.push_back_transmit(packets); |
| 1710 | break; |
| 1711 | } else if self.state == RtcState::Closing { |
| 1712 | debug!( |
| 1713 | "Dropping SCTP transmit while closing after DTLS error: {e}" |
| 1714 | ); |
| 1715 | packets.pop_front(); |
| 1716 | if !packets.is_empty() { |
| 1717 | self.sctp.push_back_transmit(packets); |
| 1718 | } |
| 1719 | return self.do_poll_output(); |
| 1720 | } else { |
| 1721 | return Err(e.into()); |
| 1722 | } |
| 1723 | } |
| 1724 | |
| 1725 | packets.pop_front(); |
| 1726 | // If there are still packets, they are sent on next |
| 1727 | // poll_output() |
| 1728 | if !packets.is_empty() { |
| 1729 | self.sctp.push_back_transmit(packets); |
| 1730 | } |
| 1731 | |
| 1732 | // Run again since this would feed the DTLS subsystem |
| 1733 | // to produce a packet now. |
| 1734 | return self.do_poll_output(); |
| 1735 | } |
| 1736 | } |
| 1737 | SctpEvent::Open { id, label } => { |
| 1738 | self.chan.ensure_channel_id_for(id); |
| 1739 | let id = self.chan.channel_id_by_stream_id(id).unwrap(); |
| 1740 | return Ok(Output::Event(Event::ChannelOpen(id, label))); |
| 1741 | } |
| 1742 | SctpEvent::Close { |
| 1743 | id: stream_id, |
| 1744 | reset_pending, |
| 1745 | } => { |
| 1746 | let Some(channel_id) = self.chan.channel_id_by_stream_id(stream_id) else { |
| 1747 | warn!("Drop ChannelClose event for id: {:?}", stream_id); |
| 1748 | continue; |
| 1749 | }; |
| 1750 | // When a reset is outstanding, remove_channel holds the stream id |
| 1751 | // back from reallocation until the handshake completes. |
| 1752 | self.chan.remove_channel(channel_id, reset_pending); |
| 1753 | return Ok(Output::Event(Event::ChannelClose(channel_id))); |
| 1754 | } |
| 1755 | SctpEvent::StreamResetComplete { id } => { |
| 1756 | // The reset handshake completed, the stream id can be used again. |
| 1757 | self.chan.stream_reset_complete(id); |
| 1758 | continue; |
| 1759 | } |
| 1760 | SctpEvent::AssociationLost => { |
| 1761 | self.chan.association_lost(); |
| 1762 | self.start_close()?; |
| 1763 | return Ok(Output::Event(Event::Closed)); |
| 1764 | } |
| 1765 | SctpEvent::Data { id, binary, data } => { |
| 1766 | let Some(id) = self.chan.channel_id_by_stream_id(id) else { |
| 1767 | warn!("Drop ChannelData event for id: {:?}", id); |
| 1768 | continue; |
| 1769 | }; |
| 1770 | let cd = ChannelData { id, binary, data }; |
| 1771 | return Ok(Output::Event(Event::ChannelData(cd))); |
| 1772 | } |
| 1773 | SctpEvent::BufferedAmountLow { id } => { |
| 1774 | let Some(id) = self.chan.channel_id_by_stream_id(id) else { |
| 1775 | warn!("Drop BufferedAmountLow for id: {:?}", id); |
| 1776 | continue; |
| 1777 | }; |
| 1778 | return Ok(Output::Event(Event::ChannelBufferedAmountLow(id))); |
| 1779 | } |
| 1780 | } |
| 1781 | } |
| 1782 | |
| 1783 | if let Some(ev) = self.session.poll_event() { |
| 1784 | return Ok(Output::Event(ev)); |
| 1785 | } |
| 1786 | |
| 1787 | // Some polling needs to bubble up errors. |
| 1788 | if let Some(ev) = self.session.poll_event_fallible()? { |
| 1789 | return Ok(Output::Event(ev)); |
| 1790 | } |
| 1791 | |
| 1792 | if let Some(e) = self.stats.as_mut().and_then(|s| s.poll_output()) { |
| 1793 | return Ok(match e { |
| 1794 | StatsEvent::Peer(s) => Output::Event(Event::PeerStats(s)), |
| 1795 | StatsEvent::MediaIngress(s) => Output::Event(Event::MediaIngressStats(s)), |
| 1796 | StatsEvent::MediaEgress(s) => Output::Event(Event::MediaEgressStats(s)), |
| 1797 | }); |
| 1798 | } |
| 1799 | |
| 1800 | if let Some(v) = self.ice.poll_transmit() { |
| 1801 | return Ok(Output::Transmit(v)); |
| 1802 | } |
| 1803 | |
| 1804 | if let Some(send) = &self.send_addr { |
| 1805 | // These can only be sent after we got an ICE connection. |
| 1806 | let datagram = None |
| 1807 | .or_else(|| self.dtls.poll_packet()) |
| 1808 | .or_else(|| self.session.poll_datagram(self.last_now)); |
| 1809 | |
| 1810 | if let Some(contents) = datagram { |
| 1811 | let t = net::Transmit { |
| 1812 | proto: send.proto, |
| 1813 | source: send.source, |
| 1814 | destination: send.destination, |
| 1815 | contents, |
| 1816 | }; |
| 1817 | return Ok(Output::Transmit(t)); |
| 1818 | } |
| 1819 | } else { |
| 1820 | // Don't allow accumulated feedback to build up indefinitely |
| 1821 | self.session.clear_feedback(); |
| 1822 | } |
| 1823 | |
| 1824 | let stats_timeout = self.stats.as_mut().and_then(|s| s.poll_timeout()); |
| 1825 | |
| 1826 | let time_and_reason = (None, Reason::NotHappening) |
| 1827 | .soonest((self.next_dtls_timeout, Reason::DTLS)) |
| 1828 | .soonest((self.ice.poll_timeout(), Reason::Ice)) |
| 1829 | .soonest(self.session.poll_timeout()) |
| 1830 | .soonest((self.sctp.poll_timeout(), Reason::Sctp)) |
| 1831 | .soonest((self.chan.poll_timeout(&self.sctp), Reason::Channel)) |
| 1832 | .soonest((stats_timeout, Reason::Stats)); |
| 1833 | |
| 1834 | // trace!("poll_output timeout reason: {}", time_and_reason.1); |
| 1835 | |
| 1836 | let time = time_and_reason.0.unwrap_or_else(not_happening); |
| 1837 | let reason = time_and_reason.1; |
| 1838 | |
| 1839 | // We want to guarantee time doesn't go backwards. |
| 1840 | let next = if time < self.last_now { |
| 1841 | self.last_now |
| 1842 | } else { |
| 1843 | time |
| 1844 | }; |
| 1845 | |
| 1846 | if self.state == RtcState::Closing && self.close_drain_complete() { |
| 1847 | self.state = RtcState::Closed; |
| 1848 | self.last_timeout_reason = Reason::NotHappening; |
| 1849 | return Ok(Output::Timeout(not_happening())); |
| 1850 | } |
| 1851 | |
| 1852 | self.last_timeout_reason = reason; |
| 1853 | Ok(Output::Timeout(next)) |
| 1854 | } |
| 1855 | |
| 1856 | /// The reason for the last [`Output::Timeout`] |
| 1857 | /// |
| 1858 | /// This is updated when calling [`Rtc::poll_output()`] and the next output |
| 1859 | /// is a timeout. |
| 1860 | /// |
| 1861 | /// ``` |
| 1862 | /// # use str0m::{Rtc, Reason}; |
| 1863 | /// # use std::time::Instant; |
| 1864 | /// let rtc = Rtc::new(Instant::now()); |
| 1865 | /// |
| 1866 | /// // Before any call to poll_output(), the reason is the default. |
| 1867 | /// assert_eq!(rtc.last_timeout_reason(), Reason::NotHappening); |
| 1868 | /// ``` |
| 1869 | pub fn last_timeout_reason(&self) -> Reason { |
| 1870 | self.last_timeout_reason |
| 1871 | } |
| 1872 | |
| 1873 | /// Check if this `Rtc` instance accepts the given input. This is used for demultiplexing |
| 1874 | /// several `Rtc` instances over the same UDP server socket. |
| 1875 | /// |
| 1876 | /// [`Input::Timeout`] is always accepted. [`Input::Receive`] is tested against the nominated |
| 1877 | /// ICE candidate. If that doesn't match and the incoming data is a STUN packet, the accept call |
| 1878 | /// is delegated to the ICE agent which recognizes the remote peer from `a=ufrag`/`a=password` |
| 1879 | /// credentials negotiated in the SDP. If that also doesn't match, all remote ICE candidates are |
| 1880 | /// checked for a match. |
| 1881 | /// |
| 1882 | /// In a server setup, the server would try to find an `Rtc` instances using [`Rtc::accepts()`]. |
| 1883 | /// The first found instance would be given the input via [`Rtc::handle_input()`]. |
| 1884 | /// |
| 1885 | /// ```no_run |
| 1886 | /// # use std::time::Instant; |
| 1887 | /// # use str0m::{Rtc, Input}; |
| 1888 | /// // A vec holding the managed rtc instances. One instance per remote peer. |
| 1889 | /// let now = Instant::now(); |
| 1890 | /// let mut rtcs = vec![Rtc::new(now), Rtc::new(now), Rtc::new(now)]; |
| 1891 | /// |
| 1892 | /// // Configure instances with local ice candidates etc. |
| 1893 | /// |
| 1894 | /// loop { |
| 1895 | /// // TODO poll_timeout() and handle the output. |
| 1896 | /// |
| 1897 | /// let input: Input = todo!(); // read network data from socket. |
| 1898 | /// for rtc in &mut rtcs { |
| 1899 | /// if rtc.accepts(&input) { |
| 1900 | /// rtc.handle_input(input).unwrap(); |
| 1901 | /// } |
| 1902 | /// } |
| 1903 | /// } |
| 1904 | /// ``` |
| 1905 | pub fn accepts(&self, input: &Input) -> bool { |
| 1906 | let Input::Receive(_, r) = input else { |
| 1907 | // always accept the Input::Timeout. |
| 1908 | return true; |
| 1909 | }; |
| 1910 | |
| 1911 | // Fast path: DTLS, RTP, and RTCP traffic coming in from the same socket address |
| 1912 | // we've nominated for sending via the ICE agent. This is the typical case |
| 1913 | if let Some(send_addr) = &self.send_addr { |
| 1914 | if r.source == send_addr.destination { |
| 1915 | return true; |
| 1916 | } |
| 1917 | } |
| 1918 | |
| 1919 | // STUN can use the ufrag/password to identify that a message belongs |
| 1920 | // to this Rtc instance. |
| 1921 | if let DatagramRecvInner::Stun(v) = &r.contents.inner { |
| 1922 | return self.ice.accepts_message(v); |
| 1923 | } |
| 1924 | |
| 1925 | // Slow path: Occasionally, traffic comes in on a socket address corresponding |
| 1926 | // to a successful candidate pair other than the one we've currently nominated. |
| 1927 | // This typically happens at the beginning of the connection |
| 1928 | if self.ice.has_viable_remote_candidate(r.source) { |
| 1929 | return true; |
| 1930 | } |
| 1931 | |
| 1932 | false |
| 1933 | } |
| 1934 | |
| 1935 | /// Provide input to this `Rtc` instance. Input is either a [`Input::Timeout`] for some |
| 1936 | /// time that was previously obtained from [`Rtc::poll_output()`], or [`Input::Receive`] |
| 1937 | /// for network data. |
| 1938 | /// |
| 1939 | /// Both the timeout and the network data contains a [`std::time::Instant`] which drives |
| 1940 | /// time forward in the instance. For network data, the intention is to record the time |
| 1941 | /// of receiving the network data as precise as possible. This time is used to calculate |
| 1942 | /// things like jitter and bandwidth. |
| 1943 | /// |
| 1944 | /// It's always okay to call [`Rtc::handle_input()`] with a timeout, also before the |
| 1945 | /// time obtained via [`Rtc::poll_output()`]. |
| 1946 | /// |
| 1947 | /// ```no_run |
| 1948 | /// # use str0m::{Rtc, Input}; |
| 1949 | /// # use std::time::Instant; |
| 1950 | /// let mut rtc = Rtc::new(Instant::now()); |
| 1951 | /// |
| 1952 | /// loop { |
| 1953 | /// let timeout: Instant = todo!(); // rtc.poll_output() until we get a timeout. |
| 1954 | /// |
| 1955 | /// let input: Input = todo!(); // wait for network data or timeout. |
| 1956 | /// rtc.handle_input(input); |
| 1957 | /// } |
| 1958 | /// ``` |
| 1959 | pub fn handle_input(&mut self, input: Input) -> Result<(), RtcError> { |
| 1960 | if self.state == RtcState::Closed { |
| 1961 | return Ok(()); |
| 1962 | } |
| 1963 | |
| 1964 | match input { |
| 1965 | Input::Timeout(now) => self.do_handle_timeout(now)?, |
| 1966 | Input::Receive(recv_time, r) => { |
| 1967 | self.do_handle_receive(recv_time, r)?; |
| 1968 | self.do_handle_timeout(recv_time)?; |
| 1969 | } |
| 1970 | } |
| 1971 | Ok(()) |
| 1972 | } |
| 1973 | |
| 1974 | /// Starts local shutdown. |
| 1975 | /// |
| 1976 | /// The `Rtc` instance schedules RTCP BYE, requests SCTP shutdown, and sends |
| 1977 | /// DTLS close_notify without waiting for remote replies. It enters a closing |
| 1978 | /// drain state where app-facing operations stop, but [`Rtc::poll_output()`] |
| 1979 | /// continues polling until pending local close output has drained. Once the |
| 1980 | /// drain is complete, [`Rtc::is_alive()`] returns `false`. |
| 1981 | pub fn close(&mut self) -> Result<(), RtcError> { |
| 1982 | self.start_close() |
| 1983 | } |
| 1984 | |
| 1985 | fn start_close(&mut self) -> Result<(), RtcError> { |
| 1986 | if self.state == RtcState::Closed { |
| 1987 | return Ok(()); |
| 1988 | } |
| 1989 | |
| 1990 | if self.state != RtcState::Closing { |
| 1991 | self.session.close_rtp(); |
| 1992 | self.sctp.close()?; |
| 1993 | self.state = RtcState::Closing; |
| 1994 | } |
| 1995 | |
| 1996 | self.start_dtls_close() |
| 1997 | } |
| 1998 | |
| 1999 | fn start_dtls_close(&mut self) -> Result<(), RtcError> { |
| 2000 | if !self.close_dtls_started { |
| 2001 | self.dtls.close()?; |
| 2002 | self.close_dtls_started = true; |
| 2003 | } |
| 2004 | |
| 2005 | Ok(()) |
| 2006 | } |
| 2007 | |
| 2008 | fn close_drain_complete(&self) -> bool { |
| 2009 | self.close_dtls_started && self.dtls.is_closed() && self.session.is_rtp_closed() |
| 2010 | } |
| 2011 | |
| 2012 | fn init_time(&mut self, now: Instant) { |
| 2013 | // The operation is somewhat expensive, hence we only do it once. |
| 2014 | if !self.need_init_time { |
| 2015 | return; |
| 2016 | } |
| 2017 | |
| 2018 | // We assume this first "now" is a time 0 start point for calculating ntp/unix time offsets. |
| 2019 | // This initializes the conversion of Instant -> NTP/Unix time. |
| 2020 | let _ = now.to_unix_duration(); |
| 2021 | |
| 2022 | self.need_init_time = false; |
| 2023 | } |
| 2024 | |
| 2025 | fn do_handle_timeout(&mut self, now: Instant) -> Result<(), RtcError> { |
| 2026 | self.init_time(now); |
| 2027 | |
| 2028 | // Prevent time from going backwards. |
| 2029 | if now < self.last_now { |
| 2030 | return Ok(()); |
| 2031 | } |
| 2032 | |
| 2033 | self.last_now = now; |
| 2034 | self.ice.handle_timeout(now); |
| 2035 | self.sctp.handle_timeout(now); |
| 2036 | self.chan.handle_timeout(now, &mut self.sctp); |
| 2037 | self.session.handle_timeout(now)?; |
| 2038 | |
| 2039 | if let Some(stats) = &mut self.stats { |
| 2040 | if stats.wants_timeout(now) { |
| 2041 | let mut snapshot = StatsSnapshot::new(now); |
| 2042 | snapshot.peer_rx = self.peer_bytes_rx; |
| 2043 | snapshot.peer_tx = self.peer_bytes_tx; |
| 2044 | let current_round_trip_time = self.ice.nominated_pair_rtt(); |
| 2045 | let total_round_trip_time = self.ice.nominated_pair_total_rtt().unwrap_or_default(); |
| 2046 | let responses_received = self.ice.nominated_pair_responses_received().unwrap_or(0); |
| 2047 | snapshot.selected_candidate_pair = |
| 2048 | self.send_addr.as_ref().map(|s| CandidatePairStats { |
| 2049 | protocol: s.proto, |
| 2050 | local: CandidateStats { addr: s.source }, |
| 2051 | remote: CandidateStats { |
| 2052 | addr: s.destination, |
| 2053 | }, |
| 2054 | current_round_trip_time, |
| 2055 | total_round_trip_time, |
| 2056 | responses_received, |
| 2057 | }); |
| 2058 | self.session.visit_stats(now, &mut snapshot); |
| 2059 | stats.do_handle_timeout(&mut snapshot); |
| 2060 | } |
| 2061 | } |
| 2062 | |
| 2063 | Ok(()) |
| 2064 | } |
| 2065 | |
| 2066 | fn do_handle_receive(&mut self, recv_time: Instant, r: net::Receive) -> Result<(), RtcError> { |
| 2067 | trace!("IN {:?}", r); |
| 2068 | use DatagramRecvInner::*; |
| 2069 | |
| 2070 | let bytes_rx = match r.contents.inner { |
| 2071 | // TODO: stun is already parsed (depacketized) here |
| 2072 | Stun(_) => 0, |
| 2073 | Dtls(v) | Rtp(v) | Rtcp(v) => v.len(), |
| 2074 | }; |
| 2075 | |
| 2076 | self.peer_bytes_rx += bytes_rx as u64; |
| 2077 | |
| 2078 | match r.contents.inner { |
| 2079 | Stun(stun) => { |
| 2080 | let packet = is::stun::StunPacket { |
| 2081 | proto: r.proto, |
| 2082 | source: r.source, |
| 2083 | destination: r.destination, |
| 2084 | message: stun, |
| 2085 | }; |
| 2086 | self.ice.handle_packet(recv_time, packet); |
| 2087 | } |
| 2088 | Dtls(dtls) => self.dtls.handle_receive(dtls)?, |
| 2089 | Rtp(rtp) => self.session.handle_rtp_receive(recv_time, rtp), |
| 2090 | Rtcp(rtcp) => self.session.handle_rtcp_receive(recv_time, rtcp), |
| 2091 | } |
| 2092 | |
| 2093 | Ok(()) |
| 2094 | } |
| 2095 | |
| 2096 | /// Obtain handle for writing to a data channel. |
| 2097 | /// |
| 2098 | /// This is first available when a [`ChannelId`] is advertised via [`Event::ChannelOpen`]. |
| 2099 | /// The function returns `None` also for IDs from [`SdpApi::add_channel()`]. |
| 2100 | /// |
| 2101 | /// Incoming channel data is via the [`Event::ChannelData`] event. |
| 2102 | /// |
| 2103 | /// ```no_run |
| 2104 | /// # use std::time::Instant; |
| 2105 | /// # use str0m::{Rtc, channel::ChannelId}; |
| 2106 | /// let mut rtc = Rtc::new(Instant::now()); |
| 2107 | /// |
| 2108 | /// let cid: ChannelId = todo!(); // obtain channel id from Event::ChannelOpen |
| 2109 | /// let channel = rtc.channel(cid).unwrap(); |
| 2110 | /// // TODO write data channel data. |
| 2111 | /// ``` |
| 2112 | pub fn channel(&mut self, id: ChannelId) -> Option<Channel<'_>> { |
| 2113 | if self.state != RtcState::Alive { |
| 2114 | return None; |
| 2115 | } |
| 2116 | |
| 2117 | let sctp_stream_id = self.chan.stream_id_by_channel_id(id)?; |
| 2118 | |
| 2119 | if !self.sctp.is_open(sctp_stream_id) { |
| 2120 | return None; |
| 2121 | } |
| 2122 | |
| 2123 | Some(Channel::new(sctp_stream_id, self)) |
| 2124 | } |
| 2125 | |
| 2126 | /// Configure the Bandwidth Estimate (BWE) subsystem. |
| 2127 | /// |
| 2128 | /// Only relevant if BWE was enabled in the [`RtcConfig::enable_bwe()`] |
| 2129 | pub fn bwe(&mut self) -> Bwe { |
| 2130 | Bwe(self) |
| 2131 | } |
| 2132 | |
| 2133 | fn is_correct_change_id(&self, change_id: usize) -> bool { |
| 2134 | self.change_counter == change_id + 1 |
| 2135 | } |
| 2136 | |
| 2137 | fn next_change_id(&mut self) -> usize { |
| 2138 | let n = self.change_counter; |
| 2139 | self.change_counter += 1; |
| 2140 | n |
| 2141 | } |
| 2142 | |
| 2143 | /// The codec configs for sending/receiving data. |
| 2144 | /// |
| 2145 | /// The configurations can be set with [`RtcConfig`] before setting up the session, and they |
| 2146 | /// might be further updated by SDP negotiation. |
| 2147 | pub fn codec_config(&self) -> &CodecConfig { |
| 2148 | &self.session.codec_config |
| 2149 | } |
| 2150 | } |
| 2151 | |
| 2152 | impl PartialEq for Event { |
| 2153 | fn eq(&self, other: &Self) -> bool { |
| 2154 | match (self, other) { |
| 2155 | (Self::IceConnectionStateChange(l0), Self::IceConnectionStateChange(r0)) => l0 == r0, |
| 2156 | (Self::MediaAdded(m0), Self::MediaAdded(m1)) => m0 == m1, |
| 2157 | (Self::MediaData(m1), Self::MediaData(m2)) => m1 == m2, |
| 2158 | (Self::ChannelOpen(l0, l1), Self::ChannelOpen(r0, r1)) => l0 == r0 && l1 == r1, |
| 2159 | (Self::ChannelData(l0), Self::ChannelData(r0)) => l0 == r0, |
| 2160 | (Self::ChannelClose(l0), Self::ChannelClose(r0)) => l0 == r0, |
| 2161 | _ => false, |
| 2162 | } |
| 2163 | } |
| 2164 | } |
| 2165 | |
| 2166 | impl Eq for Event {} |
| 2167 | |
| 2168 | impl fmt::Debug for Rtc { |
| 2169 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 2170 | f.debug_struct("Rtc").finish() |
| 2171 | } |
| 2172 | } |
| 2173 | |
| 2174 | /// Log a CSV like stat to stdout. |
| 2175 | /// |
| 2176 | /// ```ignore |
| 2177 | /// log_stat!("MY_STAT", 1, "hello", 3); |
| 2178 | /// ``` |
| 2179 | /// |
| 2180 | /// will result in the following being printed |
| 2181 | /// |
| 2182 | /// ```text |
| 2183 | /// MY_STAT 1, hello, 3, {unix_timestamp_ms} |
| 2184 | /// ```` |
| 2185 | /// |
| 2186 | /// These logs can be easily grepped for, parsed and graphed, or otherwise analyzed. |
| 2187 | /// |
| 2188 | /// This macro turns into a NO-OP if the `_internal_dont_use_log_stats` feature is not enabled |
| 2189 | macro_rules! log_stat { |
| 2190 | ($name:expr, $($arg:expr),+) => { |
| 2191 | #[cfg(feature = "_internal_dont_use_log_stats")] |
| 2192 | { |
| 2193 | use std::time::SystemTime; |
| 2194 | use std::io::{self, Write}; |
| 2195 | |
| 2196 | let now = SystemTime::now(); |
| 2197 | let since_epoch = now.duration_since(SystemTime::UNIX_EPOCH).unwrap(); |
| 2198 | let unix_time_ms = since_epoch.as_millis(); |
| 2199 | let mut lock = io::stdout().lock(); |
| 2200 | write!(lock, "{} ", $name).expect("Failed to write to stdout"); |
| 2201 | |
| 2202 | $( |
| 2203 | write!(lock, "{},", $arg).expect("Failed to write to stdout"); |
| 2204 | )+ |
| 2205 | writeln!(lock, "{}", unix_time_ms).expect("Failed to write to stdout"); |
| 2206 | } |
| 2207 | }; |
| 2208 | } |
| 2209 | pub(crate) use log_stat; |
| 2210 | |
| 2211 | #[cfg(test)] |
| 2212 | #[doc(hidden)] |
| 2213 | pub fn init_crypto_default() { |
| 2214 | crate::crypto::from_feature_flags().install_process_default(); |
| 2215 | } |
| 2216 | |
| 2217 | #[cfg(test)] |
| 2218 | mod test { |
| 2219 | use std::panic::UnwindSafe; |
| 2220 | |
| 2221 | use super::*; |
| 2222 | |
| 2223 | #[test] |
| 2224 | fn rtc_is_send() { |
| 2225 | fn is_send<T: Send>(_t: T) {} |
| 2226 | fn is_sync<T: Sync>(_t: T) {} |
| 2227 | is_send(Rtc::new(Instant::now())); |
| 2228 | is_sync(Rtc::new(Instant::now())); |
| 2229 | } |
| 2230 | |
| 2231 | #[test] |
| 2232 | fn rtc_is_unwind_safe() { |
| 2233 | fn is_unwind_safe<T: UnwindSafe>(_t: T) {} |
| 2234 | is_unwind_safe(Rtc::new(Instant::now())); |
| 2235 | } |
| 2236 | |
| 2237 | #[test] |
| 2238 | fn event_is_reasonably_sized() { |
| 2239 | let n = std::mem::size_of::<Event>(); |
| 2240 | assert!(n < 490); // Increased to accommodate abs-capture-time fields in ExtensionValues |
| 2241 | } |
| 2242 | } |
| 2243 | |
| 2244 | #[cfg(feature = "_internal_test_exports")] |
| 2245 | #[allow(missing_docs)] |
| 2246 | pub mod _internal_test_exports; |
| 2247 | |
| 2248 | #[cfg(feature = "unversioned")] |
| 2249 | pub mod unversioned { |
| 2250 | //! This module provides functionality that is not versioned according to semver. |
| 2251 | //! It may change in breaking ways between minor/patch releases, there are no guarantees. |
| 2252 | //! USE AT YOUR OWN RISK. |
| 2253 | //! |
| 2254 | //! To use this module, enable the `unversioned` feature flag in your Cargo.toml. |
| 2255 | |
| 2256 | pub use super::packet::{ |
| 2257 | Depacketizer, H264Depacketizer, H264Packetizer, OpusPacketizer, Packetizer, |
| 2258 | Vp8Depacketizer, Vp8Packetizer, |
| 2259 | }; |
| 2260 | } |