Skip to content
File

Blob: firmware/vendor/str0m/src/lib.rs

rust2261 lines
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]
676extern crate tracing;
677 
678use bwe::{Bwe, BweKind};
679use change::{DirectApi, SdpApi};
680use rtp::RawPacket;
681use std::fmt;
682use std::net::SocketAddr;
683use std::sync::Arc;
684use std::time::Instant;
685use str0m_proto::Pii;
686use streams::RtpPacket;
687use streams::StreamPaused;
688use 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")]
697macro_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")]
709pub(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.
715pub mod crypto;
716use crypto::Fingerprint;
717 
718mod dtls;
719use crate::crypto::dtls::DtlsOutput;
720use crate::crypto::{CryptoProvider, DtlsError, from_feature_flags};
721use crate::dtls::is_would_block;
722use dtls::Dtls;
723 
724use is::IceAgent;
725use is::IceAgentEvent;
726pub use is::{Candidate, CandidateBuilder, CandidateKind, IceConnectionState, IceCreds};
727 
728#[path = "config.rs"]
729mod config_mod;
730pub use config_mod::RtcConfig;
731 
732/// Default target MTU used when none is configured.
733pub use io::DATAGRAM_MTU_TARGET;
734/// Upper bound that the target endpoint of [`RtcConfig::set_mtu`] must satisfy.
735pub use io::DATAGRAM_MTU_TARGET_MAX;
736/// Lower bound that the target endpoint of [`RtcConfig::set_mtu`] must satisfy.
737pub use io::DATAGRAM_MTU_TARGET_MIN;
738/// Default warning threshold for over-sized packets.
739pub use io::DATAGRAM_MTU_WARN;
740 
741/// Additional configuration.
742pub 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)]
752pub 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 
759mod io;
760use io::DatagramRecvInner;
761 
762mod packet;
763 
764#[path = "rtp/mod.rs"]
765mod rtp_;
766use rtp_::{Bitrate, DataSize};
767 
768/// Low level RTP access.
769pub 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 
811pub(crate) mod pacer;
812 
813#[path = "bwe/mod.rs"]
814pub(crate) mod bwe_;
815 
816/// Bandwidth estimation.
817pub mod bwe {
818 pub use crate::bwe_::api::*;
819}
820 
821mod sctp;
822use sctp::{RtcSctp, SctpEvent, SctpInitData};
823 
824mod sdp;
825 
826pub mod format;
827use format::CodecConfig;
828 
829pub mod channel;
830use channel::{Channel, ChannelData, ChannelHandler, ChannelId};
831 
832pub mod media;
833use media::AppSpecificFeedback;
834use media::SenderFeedback;
835use media::{Direction, Media, Mid, Pt, Rid, Writer};
836use media::{KeyframeRequest, KeyframeRequestKind};
837use media::{MediaAdded, MediaChanged, MediaData};
838 
839pub mod change;
840 
841mod util;
842use util::{Soonest, not_happening};
843 
844mod session;
845use session::Session;
846 
847pub mod stats;
848 
849use stats::{CandidatePairStats, CandidateStats, MediaEgressStats, MediaIngressStats};
850use stats::{PeerStats, Stats, StatsEvent, StatsSnapshot};
851 
852mod streams;
853 
854pub mod error;
855 
856/// Network related types to get socket data in/out of [`Rtc`].
857pub mod net {
858 pub use crate::io::{DatagramRecv, DatagramSend, Protocol, Receive, TcpType, Transmit};
859}
860 
861const VERSION: &str = env!("CARGO_PKG_VERSION");
862 
863pub 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/// ```
903pub 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)]
929enum RtcState {
930 Alive,
931 Closing,
932 Closed,
933}
934 
935struct 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]
945pub 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 
1060impl 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.
1074pub 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)]
1086pub 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 
1097pub 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]
1104pub 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 
1180impl 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 
2152impl 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 
2166impl Eq for Event {}
2167 
2168impl 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
2189macro_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}
2209pub(crate) use log_stat;
2210 
2211#[cfg(test)]
2212#[doc(hidden)]
2213pub fn init_crypto_default() {
2214 crate::crypto::from_feature_flags().install_process_default();
2215}
2216 
2217#[cfg(test)]
2218mod 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)]
2246pub mod _internal_test_exports;
2247 
2248#[cfg(feature = "unversioned")]
2249pub 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}