From 6de554b687b87b6a3f58cf2c02007ef1321e55f6 Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Tue, 22 Sep 2026 20:14:18 +0200 Subject: [PATCH 1/7] refactor: Consolidate h264 and vpx codecs inside `codecs/` module --- .github/workflows/ci.yml | 2 +- Cargo.toml | 2 +- benches/media_baseline.rs | 7 +- benches/support/rtc_sources.rs | 12 +- src/rtc/{ => codecs}/h264.rs | 4 +- src/rtc/codecs/mod.rs | 6 + src/rtc/{ => codecs}/rtp_h264.rs | 2 +- src/rtc/{ => codecs}/rtp_vpx.rs | 0 src/rtc/{ => codecs}/vpx.rs | 321 +++++++++++++++++++++++++++++- src/rtc/local_track.rs | 8 +- src/rtc/mod.rs | 6 +- src/rtc/remote_track.rs | 7 +- src/rtc/vpx_decode.rs | 329 ------------------------------- 13 files changed, 338 insertions(+), 368 deletions(-) rename src/rtc/{ => codecs}/h264.rs (99%) create mode 100644 src/rtc/codecs/mod.rs rename src/rtc/{ => codecs}/rtp_h264.rs (99%) rename src/rtc/{ => codecs}/rtp_vpx.rs (100%) rename src/rtc/{ => codecs}/vpx.rs (71%) delete mode 100644 src/rtc/vpx_decode.rs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b8502a3..359328b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -93,6 +93,6 @@ jobs: cargo test --locked --lib -Zbuild-std --target x86_64-unknown-linux-gnu - rtc::vpx + rtc::codecs::vpx env: RUSTFLAGS: -Zsanitizer=${{ matrix.sanitizer }} diff --git a/Cargo.toml b/Cargo.toml index 165cbec..66c9285 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -69,7 +69,7 @@ tokio-tungstenite = { version = "0.30.0", default-features = false, features = [ tracing = "0.1.44" url = "2.5.8" uuid = { version = "1.24.0", features = ["v4"] } -# VP8/VP9 encoder FFI for the video publish path (`rtc::vpx`). The `generate` +# VP8/VP9 encoder FFI for the video publish path (`rtc::codecs::vpx`). The `generate` # feature runs bindgen at build time so any installed libvpx version works (no # version-locked pre-generated bindings). This links a **system** libvpx # (`brew install libvpx` / `apt install libvpx-dev`) discovered via pkg-config — diff --git a/benches/media_baseline.rs b/benches/media_baseline.rs index 0a69426..180fab6 100644 --- a/benches/media_baseline.rs +++ b/benches/media_baseline.rs @@ -16,10 +16,9 @@ use getstream::rtc::{PcmFrame, StreamResampler}; #[path = "support/rtc_sources.rs"] mod rtc; -use rtc::rtp_vpx::VpxRtpPacketizer; -use rtc::vpx::{VpxCodec, VpxEncoder}; -use rtc::vpx_decode::VpxDecoder; -use rtc::{ +use rtc::codecs::rtp_vpx::VpxRtpPacketizer; +use rtc::codecs::vpx::{VpxCodec, VpxDecoder, VpxEncoder}; +use rtc::codecs::{ h264::{H264Decoder, H264Encoder}, rtp_h264::H264RtpPacketizer, }; diff --git a/benches/support/rtc_sources.rs b/benches/support/rtc_sources.rs index 421d612..7f3098a 100644 --- a/benches/support/rtc_sources.rs +++ b/benches/support/rtc_sources.rs @@ -10,15 +10,7 @@ pub(crate) mod error { pub(crate) use getstream::rtc::{RtcError, RtcResult as Result}; } -#[path = "../../src/rtc/h264.rs"] -pub(crate) mod h264; -#[path = "../../src/rtc/rtp_h264.rs"] -pub(crate) mod rtp_h264; -#[path = "../../src/rtc/rtp_vpx.rs"] -pub(crate) mod rtp_vpx; +#[path = "../../src/rtc/codecs/mod.rs"] +pub(crate) mod codecs; #[path = "../../src/rtc/video_frame.rs"] pub(crate) mod video_frame; -#[path = "../../src/rtc/vpx.rs"] -pub(crate) mod vpx; -#[path = "../../src/rtc/vpx_decode.rs"] -pub(crate) mod vpx_decode; diff --git a/src/rtc/h264.rs b/src/rtc/codecs/h264.rs similarity index 99% rename from src/rtc/h264.rs rename to src/rtc/codecs/h264.rs index 4cacce5..874b428 100644 --- a/src/rtc/h264.rs +++ b/src/rtc/codecs/h264.rs @@ -12,8 +12,8 @@ use openh264::encoder::{ }; use openh264::formats::YUVSource; -use super::error::{Result, RtcError}; -use super::video_frame::{VideoFrame, i420_len}; +use crate::rtc::error::{Result, RtcError}; +use crate::rtc::video_frame::{VideoFrame, i420_len}; const MAX_H264_DECODE_DIMENSION: usize = 3_840; const MAX_H264_DECODE_PIXELS: usize = 3_840 * 2_160; diff --git a/src/rtc/codecs/mod.rs b/src/rtc/codecs/mod.rs new file mode 100644 index 0000000..29b8304 --- /dev/null +++ b/src/rtc/codecs/mod.rs @@ -0,0 +1,6 @@ +//! Video codecs and their RTP payload formats. + +pub(crate) mod h264; +pub(crate) mod rtp_h264; +pub(crate) mod rtp_vpx; +pub(crate) mod vpx; diff --git a/src/rtc/rtp_h264.rs b/src/rtc/codecs/rtp_h264.rs similarity index 99% rename from src/rtc/rtp_h264.rs rename to src/rtc/codecs/rtp_h264.rs index 983d16b..68c188d 100644 --- a/src/rtc/rtp_h264.rs +++ b/src/rtc/codecs/rtp_h264.rs @@ -5,7 +5,7 @@ use webrtc::rtp::Error as RtpError; use webrtc::rtp::codecs::h264::H264Payloader; use webrtc::rtp::packetizer::{Depacketizer, Payloader}; -use super::error::{Result, RtcError}; +use crate::rtc::error::{Result, RtcError}; const ANNEX_B_START_CODE: &[u8] = &[0, 0, 0, 1]; const NAL_TYPE_MASK: u8 = 0x1f; diff --git a/src/rtc/rtp_vpx.rs b/src/rtc/codecs/rtp_vpx.rs similarity index 100% rename from src/rtc/rtp_vpx.rs rename to src/rtc/codecs/rtp_vpx.rs diff --git a/src/rtc/vpx.rs b/src/rtc/codecs/vpx.rs similarity index 71% rename from src/rtc/vpx.rs rename to src/rtc/codecs/vpx.rs index 4b6be40..c9b52c1 100644 --- a/src/rtc/vpx.rs +++ b/src/rtc/codecs/vpx.rs @@ -1,6 +1,7 @@ -//! Minimal libvpx VP8/VP9 encoder for the outbound video path. +//! Minimal libvpx VP8/VP9 encoder and decoder for the outbound and inbound +//! video paths. //! -//! [`LocalVideoTrack`](super::local_track::LocalVideoTrack) needs to turn raw +//! [`LocalVideoTrack`](crate::rtc::local_track::LocalVideoTrack) needs to turn raw //! I420 frames into encoded VP8/VP9 for the SFU. The `vpx-encode` crate is too //! restrictive for realtime streaming — it hardcodes the encoder config, so it //! cannot set `g_lag_in_frames = 0` (VP9 otherwise buffers frames and emits @@ -9,10 +10,16 @@ //! that late subscribers miss). This module binds `libvpx` directly (via //! `env-libvpx-sys`, exposed as `vpx_sys`) with the correct realtime config. //! +//! [`RemoteTrack::next_video_frame`](crate::rtc::remote_track::RemoteTrack::next_video_frame) +//! needs raw frames, not RTP: an agent that wants to *see* the call has to turn +//! reassembled VP8/VP9 samples into pixels. The decoder symbols come from the +//! same `env-libvpx-sys` bindings, so inbound video adds no new native +//! dependency. +//! //! The libvpx C API is inherently unsafe; every FFI call is wrapped here and the -//! module surface is safe. libvpx encoder contexts are single-threaded but not -//! thread-*affine*, so [`VpxEncoder`] is `Send` (accessed under a mutex by the -//! caller) but not `Sync`. +//! module surface is safe. libvpx codec contexts are single-threaded but not +//! thread-*affine*, so [`VpxEncoder`], [`VpxSvcEncoder`] and [`VpxDecoder`] are +//! `Send` (accessed under a mutex by the caller) but not `Sync`. use std::mem::MaybeUninit; use std::os::raw::{c_int, c_uint, c_ulong, c_void}; @@ -31,7 +38,12 @@ use vpx_sys::vpx_img_fmt::VPX_IMG_FMT_I420; use vpx_sys::vpx_rc_mode::VPX_CBR; use vpx_sys::*; -use super::error::{Result, RtcError}; +use crate::rtc::error::{Result, RtcError}; +use crate::rtc::video_frame::{VideoFrame, i420_len}; + +/// Decoder worker threads. Two is enough to keep up with the 360p–720p an agent +/// subscribes to without competing with the rest of the runtime for cores. +const DECODE_THREADS: c_uint = 2; /// Which VPx codec a [`VpxEncoder`] produces. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -173,7 +185,7 @@ unsafe extern "C" fn collect_svc_packet(pkt: *mut vpx_codec_cx_pkt_t, user_data: /// Turn a libvpx `vpx_codec_err_t` into a `Result`. The enum is `repr(i32)`, so /// a non-zero discriminant is an error (`VPX_CODEC_OK == 0`). -pub(super) fn check(res: vpx_codec_err_t, what: &str) -> Result<()> { +fn check(res: vpx_codec_err_t, what: &str) -> Result<()> { let code: i32 = res as i32; if code == 0 { Ok(()) @@ -650,6 +662,170 @@ impl Drop for VpxSvcEncoder { } } +/// A libvpx VP8/VP9 decoder producing packed I420 [`VideoFrame`]s. +pub(crate) struct VpxDecoder { + ctx: vpx_codec_ctx_t, +} + +// SAFETY: `vpx_codec_ctx_t` holds raw pointers into libvpx internal state, which +// is why the auto `Send` impl is withheld. libvpx does not pin a context to the +// thread that created it — it only forbids *concurrent* use. Callers always hold +// the decoder behind a `std::sync::Mutex`, which serializes access, so moving +// the value across threads is sound. +unsafe impl Send for VpxDecoder {} + +impl VpxDecoder { + /// Create a decoder for `codec`. Frame size is learned from the bitstream, + /// so no dimensions are needed up front (`w`/`h` of 0). + pub(crate) fn new(codec: VpxCodec) -> Result { + // SAFETY: `iface` is a static libvpx interface pointer. `cfg` is fully + // initialized by us before the call, and `ctx` is left uninitialized and + // populated by `vpx_codec_dec_init_ver` before `assume_init` reads it. + // All pointers are valid for the call's scope. + unsafe { + let iface = match codec { + VpxCodec::Vp8 => vpx_codec_vp8_dx(), + VpxCodec::Vp9 => vpx_codec_vp9_dx(), + }; + if iface.is_null() { + return Err(RtcError::Media( + "libvpx decoder interface unavailable".to_owned(), + )); + } + + let cfg = vpx_codec_dec_cfg_t { + threads: DECODE_THREADS, + w: 0, + h: 0, + }; + + let mut ctx = MaybeUninit::::uninit(); + check( + vpx_codec_dec_init_ver( + ctx.as_mut_ptr(), + iface, + &cfg, + VPX_CODEC_USE_FRAME_THREADING as vpx_codec_flags_t, + VPX_DECODER_ABI_VERSION as c_int, + ), + "dec_init", + )?; + Ok(Self { + ctx: ctx.assume_init(), + }) + } + } + + /// Decode one reassembled VP8/VP9 frame, returning every picture libvpx + /// produced. `rtp_timestamp` is copied onto each frame. + /// + /// An empty result is normal, not an error: with frame threading libvpx may + /// buffer briefly at the start of a stream, and a frame that only updates + /// reference state yields no picture. + pub(crate) fn decode(&mut self, data: &[u8], rtp_timestamp: u32) -> Result> { + if data.is_empty() { + return Ok(Vec::new()); + } + + // SAFETY: `data` is a valid slice libvpx only reads for the duration of + // the call. The iterator-driven drain below operates on our owned + // context, and each `vpx_image_t` is copied into an owned buffer before + // the next call can recycle it. + unsafe { + check( + vpx_codec_decode( + &mut self.ctx, + data.as_ptr(), + data.len() as c_uint, + ptr::null_mut(), + 0, + ), + "decode", + )?; + + let mut frames = Vec::new(); + let mut iter: vpx_codec_iter_t = ptr::null(); + loop { + let img = vpx_codec_get_frame(&mut self.ctx, &mut iter); + if img.is_null() { + break; + } + match copy_i420(&*img, rtp_timestamp) { + Some(frame) => frames.push(frame), + None => tracing::debug!( + fmt = ?(*img).fmt, + "stream.rtc.vpx.unsupported_decoded_format" + ), + } + } + Ok(frames) + } + } +} + +impl Drop for VpxDecoder { + fn drop(&mut self) { + // SAFETY: `ctx` was successfully initialized in `new` (we only construct + // `Self` on success) and is destroyed exactly once here. + unsafe { + let _ = vpx_codec_destroy(&mut self.ctx); + } + } +} + +/// Copy a libvpx image into a packed I420 [`VideoFrame`]. +/// +/// libvpx hands back plane pointers with their own row strides, which are +/// padded well past the visible width — copying `w * h` bytes straight off +/// `planes[0]` produces the classic sheared/garbled frame. Every row is copied +/// individually at `stride[plane]`. +/// +/// # Safety +/// +/// `img` must be a live `vpx_image_t` returned by `vpx_codec_get_frame`, whose +/// planes stay valid until the next `vpx_codec_decode` call. +unsafe fn copy_i420(img: &vpx_image_t, rtp_timestamp: u32) -> Option { + if img.fmt != VPX_IMG_FMT_I420 { + return None; + } + let (width, height) = (img.d_w, img.d_h); + if width == 0 || height == 0 { + return None; + } + + let mut data = Vec::with_capacity(i420_len(width, height)); + // I420: chroma planes are half resolution in both axes (the shifts are 1, + // but read them from the image rather than assuming). + let cw = (width as usize).div_ceil(1 << img.x_chroma_shift); + let ch = (height as usize).div_ceil(1 << img.y_chroma_shift); + let planes = [ + (VPX_PLANE_Y as usize, width as usize, height as usize), + (VPX_PLANE_U as usize, cw, ch), + (VPX_PLANE_V as usize, cw, ch), + ]; + + for (plane, row_len, rows) in planes { + let base = img.planes[plane]; + let stride = img.stride[plane]; + if base.is_null() || stride < 0 || (stride as usize) < row_len { + return None; + } + for row in 0..rows { + // SAFETY: libvpx guarantees `rows` rows of at least `stride` bytes + // from `base`, and `row_len <= stride` is checked above. + let start = unsafe { base.add(row * stride as usize) }; + data.extend_from_slice(unsafe { std::slice::from_raw_parts(start, row_len) }); + } + } + + Some(VideoFrame { + width, + height, + data, + rtp_timestamp, + }) +} + #[cfg(test)] mod tests { use super::*; @@ -802,4 +978,135 @@ mod tests { assert_eq!(temporal_ids, [0, 2, 1, 2, 0, 0]); } + + /// A frame with a horizontal luma ramp: a stride bug shears or mangles the + /// gradient, so this catches the classic packed-vs-strided copy mistake in a + /// way a solid color cannot. + fn ramp_i420(width: u32, height: u32) -> Vec { + let (w, h) = (width as usize, height as usize); + let mut buf = Vec::with_capacity(i420_len(width, height)); + for _ in 0..h { + for x in 0..w { + buf.push((x * 255 / w.max(1)) as u8); + } + } + buf.extend(std::iter::repeat_n(128u8, (w / 2) * (h / 2))); + buf.extend(std::iter::repeat_n(128u8, (w / 2) * (h / 2))); + buf + } + + /// Round-trip through the real encoder and decoder: a keyframe in must come + /// back out as a frame of the same size with recognisable content. + fn round_trip(codec: VpxCodec, width: u32, height: u32) { + let mut enc = VpxEncoder::new(codec, width, height, 1_000).expect("encoder"); + let mut dec = VpxDecoder::new(codec).expect("decoder"); + let source = ramp_i420(width, height); + + let packets = enc.encode(&source, 0, 33, true).expect("encode"); + assert!(!packets.is_empty(), "encoder produced no packet"); + + let mut decoded = Vec::new(); + for packet in &packets { + decoded.extend(dec.decode(&packet.data, 12_345).expect("decode")); + } + let frame = decoded.first().expect("decoder produced no frame"); + + assert_eq!((frame.width, frame.height), (width, height)); + assert_eq!(frame.data.len(), i420_len(width, height)); + assert_eq!(frame.rtp_timestamp, 12_345); + + // The luma ramp must survive: left edge dark, right edge bright, and + // every row identical (a stride bug breaks all three). + let w = width as usize; + let y = &frame.data[..w * height as usize]; + assert!(y[0] < 40, "left edge should be dark, got {}", y[0]); + assert!( + y[w - 1] > 215, + "right edge should be bright, got {}", + y[w - 1] + ); + let mid_row = &y[(height as usize / 2) * w..][..w]; + assert!( + mid_row[0] < 40 && mid_row[w - 1] > 215, + "the ramp is not consistent across rows (stride handling)" + ); + } + + #[test] + fn vp8_round_trips_a_luma_ramp() { + round_trip(VpxCodec::Vp8, 320, 240); + } + + #[test] + fn vp9_round_trips_a_luma_ramp() { + round_trip(VpxCodec::Vp9, 320, 240); + } + + /// Odd widths make libvpx's stride padding differ most from the visible + /// width, so this is the strongest stride regression guard. + #[test] + fn vp9_round_trips_a_non_multiple_of_16_size() { + round_trip(VpxCodec::Vp9, 322, 178); + } + + /// `copy_i420` documents that libvpx planes die at the next `decode` call, + /// so a frame handed out earlier must own its pixels. + #[test] + fn decoded_frames_survive_the_next_decode_call() { + let mut enc = VpxEncoder::new(VpxCodec::Vp9, 320, 240, 1_000).expect("encoder"); + let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); + let source = ramp_i420(320, 240); + + let mut held = Vec::new(); + for packet in &enc.encode(&source, 0, 33, true).expect("encode key") { + held.extend(dec.decode(&packet.data, 1).expect("decode key")); + } + assert!(!held.is_empty(), "keyframe produced no frame"); + + for packet in &enc.encode(&source, 33, 33, false).expect("encode delta") { + dec.decode(&packet.data, 2).expect("decode delta"); + } + + let frame = &held[0]; + let w = 320; + let y = &frame.data[..w * 240]; + assert_eq!(frame.rtp_timestamp, 1); + assert!( + y[0] < 40 && y[w - 1] > 215, + "the held frame lost its ramp after a later decode" + ); + } + + /// `copy_i420` reads the dimensions off each image, so one decoder must + /// follow a stream whose resolution changes. + #[test] + fn decoder_follows_a_resolution_change() { + let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); + let mut sizes = Vec::new(); + + for (width, height) in [(320u32, 240u32), (160, 120)] { + let mut enc = VpxEncoder::new(VpxCodec::Vp9, width, height, 800).expect("encoder"); + let source = ramp_i420(width, height); + for packet in &enc.encode(&source, 0, 33, true).expect("encode") { + for frame in dec.decode(&packet.data, 0).expect("decode") { + assert_eq!(frame.data.len(), i420_len(frame.width, frame.height)); + sizes.push((frame.width, frame.height)); + } + } + } + + assert_eq!(sizes, vec![(320, 240), (160, 120)]); + } + + #[test] + fn empty_input_yields_no_frames() { + let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); + assert!(dec.decode(&[], 0).expect("decode empty").is_empty()); + } + + #[test] + fn garbage_input_is_an_error_not_a_panic() { + let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); + assert!(dec.decode(&[0xff; 32], 0).is_err()); + } } diff --git a/src/rtc/local_track.rs b/src/rtc/local_track.rs index a01ac7e..8d75233 100644 --- a/src/rtc/local_track.rs +++ b/src/rtc/local_track.rs @@ -43,15 +43,15 @@ use webrtc::track::track_local::TrackLocal; use webrtc::track::track_local::TrackLocalWriter; use webrtc::track::track_local::track_local_static_rtp::TrackLocalStaticRTP; +use super::codecs::h264::{H264Encoder, validate_h264_encode_request}; +use super::codecs::rtp_h264::H264RtpPacketizer; +use super::codecs::rtp_vpx::VpxRtpPacketizer; +use super::codecs::vpx::{Vp9SvcMode, VpxCodec, VpxEncoder, VpxSvcEncoder}; use super::error::{Result, RtcError}; -use super::h264::{H264Encoder, validate_h264_encode_request}; use super::layers::{PlannedVideoLayer, simulcast_layers, single_layer}; use super::pcm::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame, StreamResampler, rms_i16}; use super::proto::event::VideoLayerSetting; use super::proto::models::{PublishOption, TrackType, VideoLayer}; -use super::rtp_h264::H264RtpPacketizer; -use super::rtp_vpx::VpxRtpPacketizer; -use super::vpx::{Vp9SvcMode, VpxCodec, VpxEncoder, VpxSvcEncoder}; /// A raw RTP packet (`webrtc::rtp::packet::Packet`), used by the RTP-forward /// republish path. diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 4a576d3..10e835f 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -21,10 +21,10 @@ //! the re-exports below, which are covered by the crate's semver policy. pub mod client; +mod codecs; pub mod coordinator; pub mod coordinator_ws; pub mod error; -mod h264; pub mod identity; pub mod join; mod layers; @@ -36,16 +36,12 @@ mod publish_options; pub mod publisher; pub mod reconnect; pub mod remote_track; -mod rtp_h264; -mod rtp_vpx; pub mod sfu_ws; pub mod signal; pub mod stats; pub mod subscriptions; pub mod tracer; pub mod video_frame; -mod vpx; -mod vpx_decode; pub use client::{RtcCall, RtcClient, TokenFuture, TokenProvider}; pub use coordinator::{ diff --git a/src/rtc/remote_track.rs b/src/rtc/remote_track.rs index 4a1f783..fb0c643 100644 --- a/src/rtc/remote_track.rs +++ b/src/rtc/remote_track.rs @@ -32,15 +32,14 @@ use webrtc::rtp::codecs::vp8::Vp8Packet; use webrtc::rtp::codecs::vp9::Vp9Packet; use webrtc::track::track_remote::TrackRemote; +use super::codecs::h264::{H264Decoder, access_unit_has_idr}; +use super::codecs::rtp_h264::H264Depacketizer; +use super::codecs::vpx::{VpxCodec, VpxDecoder}; use super::error::{Result, RtcError}; -use super::h264::{H264Decoder, access_unit_has_idr}; use super::local_track::RtpPacket; use super::pcm::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame}; use super::proto::models::{self, TrackType}; -use super::rtp_h264::H264Depacketizer; use super::video_frame::VideoFrame; -use super::vpx::VpxCodec; -use super::vpx_decode::VpxDecoder; /// How many packets the video [`SampleBuilder`] buffers while waiting for gaps /// to be filled (by NACK/RTX or a reordered arrival) before giving up on a diff --git a/src/rtc/vpx_decode.rs b/src/rtc/vpx_decode.rs deleted file mode 100644 index 64c535b..0000000 --- a/src/rtc/vpx_decode.rs +++ /dev/null @@ -1,329 +0,0 @@ -//! Minimal libvpx VP8/VP9 decoder for the inbound video path. -//! -//! [`RemoteTrack::next_video_frame`](super::remote_track::RemoteTrack::next_video_frame) -//! needs raw frames, not RTP: a bot that wants to *see* the call has to turn -//! reassembled VP8/VP9 samples into pixels. This is the mirror of -//! [`vpx`](super::vpx), which binds the encoder for the outbound path — the -//! decoder symbols come from the same `env-libvpx-sys` bindings, so supporting -//! inbound video adds no new native dependency. -//! -//! The libvpx C API is inherently unsafe; every FFI call is wrapped here and the -//! module surface is safe. libvpx decoder contexts are not thread-*affine*, so -//! [`VpxDecoder`] is `Send` (accessed under a mutex by the caller) but not -//! `Sync`. - -use std::mem::MaybeUninit; -use std::os::raw::{c_int, c_uint}; -use std::ptr; - -use vpx_sys::vpx_img_fmt::VPX_IMG_FMT_I420; -use vpx_sys::*; - -use super::error::{Result, RtcError}; -use super::video_frame::{VideoFrame, i420_len}; -use super::vpx::{VpxCodec, check}; - -/// Decoder worker threads. Two is enough to keep up with the 360p–720p a bot -/// subscribes to without competing with the rest of the runtime for cores. -const DECODE_THREADS: c_uint = 2; - -/// A libvpx VP8/VP9 decoder producing packed I420 [`VideoFrame`]s. -pub(crate) struct VpxDecoder { - ctx: vpx_codec_ctx_t, -} - -// SAFETY: `vpx_codec_ctx_t` holds raw pointers into libvpx internal state, which -// is why the auto `Send` impl is withheld. libvpx does not pin a context to the -// thread that created it — it only forbids *concurrent* use. Callers always hold -// the decoder behind a `std::sync::Mutex`, which serializes access, so moving -// the value across threads is sound. -unsafe impl Send for VpxDecoder {} - -impl VpxDecoder { - /// Create a decoder for `codec`. Frame size is learned from the bitstream, - /// so no dimensions are needed up front (`w`/`h` of 0). - pub(crate) fn new(codec: VpxCodec) -> Result { - // SAFETY: `iface` is a static libvpx interface pointer. `cfg` is fully - // initialized by us before the call, and `ctx` is left uninitialized and - // populated by `vpx_codec_dec_init_ver` before `assume_init` reads it. - // All pointers are valid for the call's scope. - unsafe { - let iface = match codec { - VpxCodec::Vp8 => vpx_codec_vp8_dx(), - VpxCodec::Vp9 => vpx_codec_vp9_dx(), - }; - if iface.is_null() { - return Err(RtcError::Media( - "libvpx decoder interface unavailable".to_owned(), - )); - } - - let cfg = vpx_codec_dec_cfg_t { - threads: DECODE_THREADS, - w: 0, - h: 0, - }; - - let mut ctx = MaybeUninit::::uninit(); - check( - vpx_codec_dec_init_ver( - ctx.as_mut_ptr(), - iface, - &cfg, - VPX_CODEC_USE_FRAME_THREADING as vpx_codec_flags_t, - VPX_DECODER_ABI_VERSION as c_int, - ), - "dec_init", - )?; - Ok(Self { - ctx: ctx.assume_init(), - }) - } - } - - /// Decode one reassembled VP8/VP9 frame, returning every picture libvpx - /// produced. `rtp_timestamp` is copied onto each frame. - /// - /// An empty result is normal, not an error: with frame threading libvpx may - /// buffer briefly at the start of a stream, and a frame that only updates - /// reference state yields no picture. - pub(crate) fn decode(&mut self, data: &[u8], rtp_timestamp: u32) -> Result> { - if data.is_empty() { - return Ok(Vec::new()); - } - - // SAFETY: `data` is a valid slice libvpx only reads for the duration of - // the call. The iterator-driven drain below operates on our owned - // context, and each `vpx_image_t` is copied into an owned buffer before - // the next call can recycle it. - unsafe { - check( - vpx_codec_decode( - &mut self.ctx, - data.as_ptr(), - data.len() as c_uint, - ptr::null_mut(), - 0, - ), - "decode", - )?; - - let mut frames = Vec::new(); - let mut iter: vpx_codec_iter_t = ptr::null(); - loop { - let img = vpx_codec_get_frame(&mut self.ctx, &mut iter); - if img.is_null() { - break; - } - match copy_i420(&*img, rtp_timestamp) { - Some(frame) => frames.push(frame), - None => tracing::debug!( - fmt = ?(*img).fmt, - "stream.rtc.vpx.unsupported_decoded_format" - ), - } - } - Ok(frames) - } - } -} - -impl Drop for VpxDecoder { - fn drop(&mut self) { - // SAFETY: `ctx` was successfully initialized in `new` (we only construct - // `Self` on success) and is destroyed exactly once here. - unsafe { - let _ = vpx_codec_destroy(&mut self.ctx); - } - } -} - -/// Copy a libvpx image into a packed I420 [`VideoFrame`]. -/// -/// libvpx hands back plane pointers with their own row strides, which are -/// padded well past the visible width — copying `w * h` bytes straight off -/// `planes[0]` produces the classic sheared/garbled frame. Every row is copied -/// individually at `stride[plane]`. -/// -/// # Safety -/// -/// `img` must be a live `vpx_image_t` returned by `vpx_codec_get_frame`, whose -/// planes stay valid until the next `vpx_codec_decode` call. -unsafe fn copy_i420(img: &vpx_image_t, rtp_timestamp: u32) -> Option { - if img.fmt != VPX_IMG_FMT_I420 { - return None; - } - let (width, height) = (img.d_w, img.d_h); - if width == 0 || height == 0 { - return None; - } - - let mut data = Vec::with_capacity(i420_len(width, height)); - // I420: chroma planes are half resolution in both axes (the shifts are 1, - // but read them from the image rather than assuming). - let cw = (width as usize).div_ceil(1 << img.x_chroma_shift); - let ch = (height as usize).div_ceil(1 << img.y_chroma_shift); - let planes = [ - (VPX_PLANE_Y as usize, width as usize, height as usize), - (VPX_PLANE_U as usize, cw, ch), - (VPX_PLANE_V as usize, cw, ch), - ]; - - for (plane, row_len, rows) in planes { - let base = img.planes[plane]; - let stride = img.stride[plane]; - if base.is_null() || stride < 0 || (stride as usize) < row_len { - return None; - } - for row in 0..rows { - // SAFETY: libvpx guarantees `rows` rows of at least `stride` bytes - // from `base`, and `row_len <= stride` is checked above. - let start = unsafe { base.add(row * stride as usize) }; - data.extend_from_slice(unsafe { std::slice::from_raw_parts(start, row_len) }); - } - } - - Some(VideoFrame { - width, - height, - data, - rtp_timestamp, - }) -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::rtc::vpx::VpxEncoder; - - /// A frame with a horizontal luma ramp: a stride bug shears or mangles the - /// gradient, so this catches the classic packed-vs-strided copy mistake in a - /// way a solid color cannot. - fn ramp_i420(width: u32, height: u32) -> Vec { - let (w, h) = (width as usize, height as usize); - let mut buf = Vec::with_capacity(i420_len(width, height)); - for _ in 0..h { - for x in 0..w { - buf.push((x * 255 / w.max(1)) as u8); - } - } - buf.extend(std::iter::repeat_n(128u8, (w / 2) * (h / 2))); - buf.extend(std::iter::repeat_n(128u8, (w / 2) * (h / 2))); - buf - } - - /// Round-trip through the real encoder and decoder: a keyframe in must come - /// back out as a frame of the same size with recognisable content. - fn round_trip(codec: VpxCodec, width: u32, height: u32) { - let mut enc = VpxEncoder::new(codec, width, height, 1_000).expect("encoder"); - let mut dec = VpxDecoder::new(codec).expect("decoder"); - let source = ramp_i420(width, height); - - let packets = enc.encode(&source, 0, 33, true).expect("encode"); - assert!(!packets.is_empty(), "encoder produced no packet"); - - let mut decoded = Vec::new(); - for packet in &packets { - decoded.extend(dec.decode(&packet.data, 12_345).expect("decode")); - } - let frame = decoded.first().expect("decoder produced no frame"); - - assert_eq!((frame.width, frame.height), (width, height)); - assert_eq!(frame.data.len(), i420_len(width, height)); - assert_eq!(frame.rtp_timestamp, 12_345); - - // The luma ramp must survive: left edge dark, right edge bright, and - // every row identical (a stride bug breaks all three). - let w = width as usize; - let y = &frame.data[..w * height as usize]; - assert!(y[0] < 40, "left edge should be dark, got {}", y[0]); - assert!( - y[w - 1] > 215, - "right edge should be bright, got {}", - y[w - 1] - ); - let mid_row = &y[(height as usize / 2) * w..][..w]; - assert!( - mid_row[0] < 40 && mid_row[w - 1] > 215, - "the ramp is not consistent across rows (stride handling)" - ); - } - - #[test] - fn vp8_round_trips_a_luma_ramp() { - round_trip(VpxCodec::Vp8, 320, 240); - } - - #[test] - fn vp9_round_trips_a_luma_ramp() { - round_trip(VpxCodec::Vp9, 320, 240); - } - - /// Odd widths make libvpx's stride padding differ most from the visible - /// width, so this is the strongest stride regression guard. - #[test] - fn vp9_round_trips_a_non_multiple_of_16_size() { - round_trip(VpxCodec::Vp9, 322, 178); - } - - /// `copy_i420` documents that libvpx planes die at the next `decode` call, - /// so a frame handed out earlier must own its pixels. - #[test] - fn decoded_frames_survive_the_next_decode_call() { - let mut enc = VpxEncoder::new(VpxCodec::Vp9, 320, 240, 1_000).expect("encoder"); - let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); - let source = ramp_i420(320, 240); - - let mut held = Vec::new(); - for packet in &enc.encode(&source, 0, 33, true).expect("encode key") { - held.extend(dec.decode(&packet.data, 1).expect("decode key")); - } - assert!(!held.is_empty(), "keyframe produced no frame"); - - for packet in &enc.encode(&source, 33, 33, false).expect("encode delta") { - dec.decode(&packet.data, 2).expect("decode delta"); - } - - let frame = &held[0]; - let w = 320; - let y = &frame.data[..w * 240]; - assert_eq!(frame.rtp_timestamp, 1); - assert!( - y[0] < 40 && y[w - 1] > 215, - "the held frame lost its ramp after a later decode" - ); - } - - /// `copy_i420` reads the dimensions off each image, so one decoder must - /// follow a stream whose resolution changes. - #[test] - fn decoder_follows_a_resolution_change() { - let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); - let mut sizes = Vec::new(); - - for (width, height) in [(320u32, 240u32), (160, 120)] { - let mut enc = VpxEncoder::new(VpxCodec::Vp9, width, height, 800).expect("encoder"); - let source = ramp_i420(width, height); - for packet in &enc.encode(&source, 0, 33, true).expect("encode") { - for frame in dec.decode(&packet.data, 0).expect("decode") { - assert_eq!(frame.data.len(), i420_len(frame.width, frame.height)); - sizes.push((frame.width, frame.height)); - } - } - } - - assert_eq!(sizes, vec![(320, 240), (160, 120)]); - } - - #[test] - fn empty_input_yields_no_frames() { - let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); - assert!(dec.decode(&[], 0).expect("decode empty").is_empty()); - } - - #[test] - fn garbage_input_is_an_error_not_a_panic() { - let mut dec = VpxDecoder::new(VpxCodec::Vp9).expect("decoder"); - assert!(dec.decode(&[0xff; 32], 0).is_err()); - } -} From 64c49d6e5ee2db62a6f55c9d5bec1f0b2eadf60d Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Tue, 22 Sep 2026 20:27:21 +0200 Subject: [PATCH 2/7] refactor: consolidate peer connections & ICE code under peer/ --- src/rtc/join/connection.rs | 125 +-------------------------- src/rtc/join/mod.rs | 58 +------------ src/rtc/join/publication.rs | 2 +- src/rtc/join/publish.rs | 4 +- src/rtc/mod.rs | 3 +- src/rtc/peer/ice.rs | 144 +++++++++++++++++++++++++++++++ src/rtc/{peer.rs => peer/mod.rs} | 7 ++ src/rtc/{ => peer}/publisher.rs | 21 ++--- src/rtc/peer/subscriber.rs | 51 +++++++++++ 9 files changed, 225 insertions(+), 190 deletions(-) create mode 100644 src/rtc/peer/ice.rs rename src/rtc/{peer.rs => peer/mod.rs} (98%) rename src/rtc/{ => peer}/publisher.rs (96%) create mode 100644 src/rtc/peer/subscriber.rs diff --git a/src/rtc/join/connection.rs b/src/rtc/join/connection.rs index 0ea679f..fc237fe 100644 --- a/src/rtc/join/connection.rs +++ b/src/rtc/join/connection.rs @@ -63,53 +63,6 @@ pub(super) async fn await_join_response( } } -/// Register a subscriber/publisher `on_ice_candidate` handler that trickles -/// gathered candidates to the SFU over Twirp. -pub(super) fn register_ice_trickle( - pc: &Arc, - signal: SignalClient, - session_id: String, - peer_type: PeerType, - tracer: Arc, -) { - pc.on_ice_candidate(Box::new(move |candidate: Option| { - let signal = signal.clone(); - let session_id = session_id.clone(); - let tracer = tracer.clone(); - Box::pin(async move { - let Some(candidate) = candidate else { return }; - let init = match candidate.to_json() { - Ok(init) => init, - Err(e) => { - tracing::debug!(error = %e, "stream.rtc.ice.to_json_failed"); - return; - } - }; - // Match JS `onicecandidate`: trace the candidate init object. - tracer.trace( - "onicecandidate", - serde_json::to_value(&init).unwrap_or(serde_json::Value::Null), - ); - let ice_candidate = match serde_json::to_string(&init) { - Ok(s) => s, - Err(e) => { - tracing::debug!(error = %e, "stream.rtc.ice.serialize_failed"); - return; - } - }; - let trickle = models::IceTrickle { - peer_type: peer_type as i32, - ice_candidate, - session_id, - }; - match signal.ice_trickle(trickle).await { - Ok(_) => tracing::debug!(?peer_type, "stream.rtc.ice.trickle_sent"), - Err(e) => tracing::debug!(error = %e, "stream.rtc.ice.trickle_failed"), - } - }) - })); -} - /// Register the subscriber `on_track` handler, delivering each inbound track to /// the core's correlation + `on_track` callback path. Also traces `ontrack` /// (`: [stream:]`) to match JS. @@ -279,13 +232,10 @@ pub(super) async fn handle_event( .await?; } E::IceTrickle(trickle) => { - add_remote_candidate( - &context.subscriber, - &context.publisher, - trickle, - &context.pending_ice, - ) - .await?; + context + .pending_ice + .add_remote(&context.subscriber, &context.publisher, trickle) + .await?; } E::ConnectionQualityChanged(event) => { core.update_connection_quality(&event.connection_quality_updates); @@ -437,73 +387,6 @@ pub(super) async fn handle_event( Ok(()) } -/// Answer an SFU subscriber offer and post the answer over Twirp. -pub(super) async fn negotiate_subscriber( - subscriber: &Arc, - signal: &SignalClient, - session_id: &str, - offer: super::super::proto::event::SubscriberOffer, - pending_ice: &Arc, -) -> Result<()> { - let remote = RTCSessionDescription::offer(offer.sdp) - .map_err(|e| RtcError::Negotiation(super::super::error::NegotiationError(e.to_string())))?; - subscriber - .set_remote_description(remote) - .await - .map_err(|e| RtcError::Negotiation(super::super::error::NegotiationError(e.to_string())))?; - // The remote description now exists: release any candidates the SFU trickled - // before this offer arrived. - flush_candidates(subscriber, &pending_ice.subscriber).await; - let answer = subscriber - .create_answer(None) - .await - .map_err(|e| RtcError::Negotiation(super::super::error::NegotiationError(e.to_string())))?; - subscriber - .set_local_description(answer.clone()) - .await - .map_err(|e| RtcError::Negotiation(super::super::error::NegotiationError(e.to_string())))?; - - signal - .send_answer(signal::SendAnswerRequest { - peer_type: PeerType::Subscriber as i32, - sdp: answer.sdp, - session_id: session_id.to_owned(), - negotiation_id: offer.negotiation_id, - }) - .await?; - tracing::debug!(session_id, "stream.rtc.subscriber.answer_sent"); - Ok(()) -} - -/// Add a remote ICE candidate to the publisher or subscriber PeerConnection, -/// buffering it if the remote description is not set yet. -pub(super) async fn add_remote_candidate( - subscriber: &Arc, - publisher: &Arc, - trickle: models::IceTrickle, - pending_ice: &Arc, -) -> Result<()> { - let init: RTCIceCandidateInit = serde_json::from_str(&trickle.ice_candidate)?; - let (target, queue) = if trickle.peer_type == PeerType::Subscriber as i32 { - (subscriber, &pending_ice.subscriber) - } else { - (publisher, &pending_ice.publisher) - }; - for candidate in queue.offer(init) { - target.add_ice_candidate(candidate).await?; - } - Ok(()) -} - -/// Release every buffered candidate now that `pc`'s remote description is set. -pub(super) async fn flush_candidates(pc: &Arc, queue: &CandidateQueue) { - for candidate in queue.mark_ready() { - if let Err(e) = pc.add_ice_candidate(candidate).await { - tracing::debug!(error = %e, "stream.rtc.ice.flush_add_failed"); - } - } -} - /// Health-check ping loop (JS 5s cadence) keeping the SFU session alive. pub(super) async fn ping_loop( core: Arc, diff --git a/src/rtc/join/mod.rs b/src/rtc/join/mod.rs index ee9a285..da65942 100644 --- a/src/rtc/join/mod.rs +++ b/src/rtc/join/mod.rs @@ -38,10 +38,8 @@ use std::time::{Duration, Instant}; use tokio::sync::{Mutex as TokioMutex, Notify, broadcast}; use tokio::task::JoinHandle; use url::Url; -use webrtc::ice_transport::ice_candidate::{RTCIceCandidate, RTCIceCandidateInit}; use webrtc::peer_connection::RTCPeerConnection; use webrtc::peer_connection::peer_connection_state::RTCPeerConnectionState; -use webrtc::peer_connection::sdp::session_description::RTCSessionDescription; use webrtc::rtp_transceiver::RTCRtpTransceiverInit; use webrtc::rtp_transceiver::rtp_codec::RTPCodecType; use webrtc::rtp_transceiver::rtp_transceiver_direction::RTCRtpTransceiverDirection; @@ -56,12 +54,11 @@ use super::coordinator_ws::{self, ConnectUserDetails, CoordinatorEvent, WsAuthMe use super::error::{Result, RtcError, SfuJoinError, SfuTimeoutError}; use super::identity; use super::local_track::LocalTrack; -use super::peer; +use super::peer::{self, PendingIce, negotiate_subscriber, publisher, register_ice_trickle}; use super::proto::event::{self, JoinRequest, JoinResponse, ReconnectDetails, SfuEvent, sfu_event}; use super::proto::models::{self, PeerType, TrackType}; use super::proto::signal; use super::publish_options::ClientPublishOptions; -use super::publisher; use super::reconnect::{ self, FailureCaps, ReconnectStrategy, SlidingWindowRateLimiter, escalate_strategy, strategy_after_signal_close, @@ -84,8 +81,8 @@ mod roster; mod subscriptions_runtime; use connection::{ - await_join_response, build_sfu_ws_url, event_loop, flush_candidates, ping_loop, - register_connection_state, register_ice_trickle, register_on_track, + await_join_response, build_sfu_ws_url, event_loop, ping_loop, register_connection_state, + register_on_track, }; use publication::{MediaState, PublicationStatus}; use roster::{CallStateCache, RosterEntry}; @@ -272,26 +269,6 @@ pub struct CallStateSnapshot { pub current_grants: Option, } -/// Buffers remote ICE candidates that arrive before a PeerConnection's remote -/// description is set, then releases them once it is. -/// -/// webrtc-rs rejects `add_ice_candidate` before the remote description exists, -/// and the SFU trickles its candidates as soon as it receives our offer — often -/// before our `set_remote_description` runs. Dropping those candidates leaves -/// the agent with no pairs and ICE fails. The queue serializes "buffer vs add" -/// under one lock so no candidate is lost to the race (JS `SfuClient` pending -/// candidate handling). -#[derive(Default)] -struct CandidateQueue { - inner: StdMutex, -} - -#[derive(Default)] -struct CandidateQueueInner { - remote_set: bool, - pending: Vec, -} - #[derive(Debug)] struct Lifecycle { state: CallingState, @@ -300,35 +277,6 @@ struct Lifecycle { generation_publish_options: ClientPublishOptions, } -impl CandidateQueue { - /// Offer a freshly-trickled candidate: returns the candidates to add now - /// (the new one if the remote description is set, else none — it is buffered). - fn offer(&self, init: RTCIceCandidateInit) -> Vec { - let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); - if g.remote_set { - vec![init] - } else { - g.pending.push(init); - Vec::new() - } - } - - /// Mark the remote description as set and return every buffered candidate to - /// be added now. Idempotent across renegotiations. - fn mark_ready(&self) -> Vec { - let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); - g.remote_set = true; - std::mem::take(&mut g.pending) - } -} - -/// Per-connection ICE candidate buffers for both PeerConnections. -#[derive(Default)] -struct PendingIce { - publisher: CandidateQueue, - subscriber: CandidateQueue, -} - /// A live SFU connection bundle. Swapped out wholesale on REJOIN/MIGRATE. struct Connection { generation: u64, diff --git a/src/rtc/join/publication.rs b/src/rtc/join/publication.rs index 8b373e1..c6a4914 100644 --- a/src/rtc/join/publication.rs +++ b/src/rtc/join/publication.rs @@ -6,7 +6,7 @@ use super::super::error::Result; use super::super::local_track::LocalTrack; use super::super::proto::event; use super::super::proto::models::{self, TrackType}; -use super::super::publisher; +use crate::rtc::peer::publisher; #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(super) enum PublicationStatus { diff --git a/src/rtc/join/publish.rs b/src/rtc/join/publish.rs index 7335881..ffb1a8a 100644 --- a/src/rtc/join/publish.rs +++ b/src/rtc/join/publish.rs @@ -320,7 +320,9 @@ impl RtcCore { .as_ref() .map(|connection| (connection.publisher.clone(), connection.pending_ice.clone())); if let Some((publisher, pending_ice)) = handles { - flush_candidates(&publisher, &pending_ice.publisher).await; + pending_ice + .flush(&publisher, PeerType::PublisherUnspecified) + .await; } } diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 10e835f..442c0c4 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -15,7 +15,7 @@ //! # Stability //! //! The wire-layer modules — [`proto`], [`peer`], [`sfu_ws`], [`signal`], -//! [`publisher`], [`tracer`], and [`coordinator_ws`] — mirror Stream's SFU +//! [`tracer`], and [`coordinator_ws`] — mirror Stream's SFU //! protocol and change with it. They are exempt from this crate's compatibility //! guarantees at any version bump. Prefer [`crate::Call`], [`RtcClient`], and //! the re-exports below, which are covered by the crate's semver policy. @@ -33,7 +33,6 @@ pub mod pcm; pub mod peer; pub mod proto; mod publish_options; -pub mod publisher; pub mod reconnect; pub mod remote_track; pub mod sfu_ws; diff --git a/src/rtc/peer/ice.rs b/src/rtc/peer/ice.rs new file mode 100644 index 0000000..679827f --- /dev/null +++ b/src/rtc/peer/ice.rs @@ -0,0 +1,144 @@ +//! ICE candidate exchange with the SFU: trickling local candidates and +//! buffering remote ones until a PeerConnection can accept them. + +use std::sync::{Arc, Mutex as StdMutex}; + +use webrtc::ice_transport::ice_candidate::{RTCIceCandidate, RTCIceCandidateInit}; +use webrtc::peer_connection::RTCPeerConnection; + +use crate::rtc::error::Result; +use crate::rtc::proto::models::{self, PeerType}; +use crate::rtc::signal::SignalClient; +use crate::rtc::tracer::Tracer; + +/// Buffers remote ICE candidates that arrive before a PeerConnection's remote +/// description is set, then releases them once it is. +/// +/// webrtc-rs rejects `add_ice_candidate` before the remote description exists, +/// and the SFU trickles its candidates as soon as it receives our offer — often +/// before our `set_remote_description` runs. Dropping those candidates leaves +/// the agent with no pairs and ICE fails. The queue serializes "buffer vs add" +/// under one lock so no candidate is lost to the race (JS `SfuClient` pending +/// candidate handling). +#[derive(Default)] +struct CandidateQueue { + inner: StdMutex, +} + +#[derive(Default)] +struct CandidateQueueInner { + remote_set: bool, + pending: Vec, +} + +impl CandidateQueue { + /// Offer a freshly-trickled candidate: returns the candidates to add now + /// (the new one if the remote description is set, else none — it is buffered). + fn offer(&self, init: RTCIceCandidateInit) -> Vec { + let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); + if g.remote_set { + vec![init] + } else { + g.pending.push(init); + Vec::new() + } + } + + /// Mark the remote description as set and return every buffered candidate to + /// be added now. Idempotent across renegotiations. + fn mark_ready(&self) -> Vec { + let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); + g.remote_set = true; + std::mem::take(&mut g.pending) + } +} + +/// Per-connection ICE candidate buffers for both PeerConnections. +#[derive(Default)] +pub struct PendingIce { + publisher: CandidateQueue, + subscriber: CandidateQueue, +} + +impl PendingIce { + /// Add a remote ICE candidate to the publisher or subscriber PeerConnection, + /// buffering it if the remote description is not set yet. + pub async fn add_remote( + &self, + subscriber: &Arc, + publisher: &Arc, + trickle: models::IceTrickle, + ) -> Result<()> { + let init: RTCIceCandidateInit = serde_json::from_str(&trickle.ice_candidate)?; + let (target, queue) = if trickle.peer_type == PeerType::Subscriber as i32 { + (subscriber, &self.subscriber) + } else { + (publisher, &self.publisher) + }; + for candidate in queue.offer(init) { + target.add_ice_candidate(candidate).await?; + } + Ok(()) + } + + /// Release every buffered candidate now that `pc`'s remote description is set. + pub async fn flush(&self, pc: &Arc, peer_type: PeerType) { + let queue = if peer_type == PeerType::Subscriber { + &self.subscriber + } else { + &self.publisher + }; + for candidate in queue.mark_ready() { + if let Err(e) = pc.add_ice_candidate(candidate).await { + tracing::debug!(error = %e, "stream.rtc.ice.flush_add_failed"); + } + } + } +} + +/// Register a subscriber/publisher `on_ice_candidate` handler that trickles +/// gathered candidates to the SFU over Twirp. +pub fn register_ice_trickle( + pc: &Arc, + signal: SignalClient, + session_id: String, + peer_type: PeerType, + tracer: Arc, +) { + pc.on_ice_candidate(Box::new(move |candidate: Option| { + let signal = signal.clone(); + let session_id = session_id.clone(); + let tracer = tracer.clone(); + Box::pin(async move { + let Some(candidate) = candidate else { return }; + let init = match candidate.to_json() { + Ok(init) => init, + Err(e) => { + tracing::debug!(error = %e, "stream.rtc.ice.to_json_failed"); + return; + } + }; + // Match JS `onicecandidate`: trace the candidate init object. + tracer.trace( + "onicecandidate", + serde_json::to_value(&init).unwrap_or(serde_json::Value::Null), + ); + let ice_candidate = match serde_json::to_string(&init) { + Ok(s) => s, + Err(e) => { + tracing::debug!(error = %e, "stream.rtc.ice.serialize_failed"); + return; + } + }; + let trickle = models::IceTrickle { + peer_type: peer_type as i32, + ice_candidate, + session_id, + }; + match signal.ice_trickle(trickle).await { + Ok(_) => tracing::debug!(?peer_type, "stream.rtc.ice.trickle_sent"), + Err(e) => tracing::debug!(error = %e, "stream.rtc.ice.trickle_failed"), + } + }) + })); +} diff --git a/src/rtc/peer.rs b/src/rtc/peer/mod.rs similarity index 98% rename from src/rtc/peer.rs rename to src/rtc/peer/mod.rs index 3899e0a..8e8ad6e 100644 --- a/src/rtc/peer.rs +++ b/src/rtc/peer/mod.rs @@ -7,6 +7,10 @@ //! SFU inspects to learn our codec capabilities on the `JoinRequest` //! (JS `getGenericSdp`, stream-py `create_join_request`). +mod ice; +pub mod publisher; +mod subscriber; + use std::sync::Arc; use serde_json::json; @@ -33,6 +37,9 @@ use super::error::Result; use super::publish_options::H264_FMTP; use super::tracer::Tracer; +pub(super) use ice::{PendingIce, register_ice_trickle}; +pub(super) use subscriber::negotiate_subscriber; + const OPUS_PAYLOAD_TYPE: u8 = 111; const VP8_PAYLOAD_TYPE: u8 = 96; const VP9_PAYLOAD_TYPE: u8 = 98; diff --git a/src/rtc/publisher.rs b/src/rtc/peer/publisher.rs similarity index 96% rename from src/rtc/publisher.rs rename to src/rtc/peer/publisher.rs index c7ad679..b5b3e72 100644 --- a/src/rtc/publisher.rs +++ b/src/rtc/peer/publisher.rs @@ -15,11 +15,11 @@ use webrtc::peer_connection::sdp::sdp_type::RTCSdpType; use webrtc::peer_connection::sdp::session_description::RTCSessionDescription; use webrtc::peer_connection::signaling_state::RTCSignalingState; -use super::error::{NegotiationError, Result, RtcError}; -use super::local_track::LocalTrack; -use super::proto::models::{PublishOption, TrackInfo, TrackType}; -use super::proto::signal::SetPublisherRequest; -use super::signal::SignalClient; +use crate::rtc::error::{NegotiationError, Result, RtcError}; +use crate::rtc::local_track::LocalTrack; +use crate::rtc::proto::models::{PublishOption, TrackInfo, TrackType}; +use crate::rtc::proto::signal::SetPublisherRequest; +use crate::rtc::signal::SignalClient; /// Renegotiate the publisher PeerConnection with the SFU for `tracks`. /// @@ -184,7 +184,7 @@ pub(crate) async fn add_transceiver_for_track( .next() .ok_or_else(|| RtcError::Media("local publication has no physical encodings".to_owned()))?; let transceiver = publisher - .add_transceiver_from_track(first, Some(super::join::send_only())) + .add_transceiver_from_track(first, Some(crate::rtc::join::send_only())) .await .map_err(RtcError::from)?; let sender = transceiver.sender().await; @@ -296,10 +296,11 @@ pub(crate) async fn build_track_infos( #[cfg(test)] mod tests { - use super::super::local_track::{LocalVideoTrack, LocalVideoTrackConfig}; - use super::super::peer; - use super::super::proto::models::{Codec, VideoDimension}; use super::*; + use crate::rtc::local_track::{LocalVideoTrack, LocalVideoTrackConfig}; + use crate::rtc::peer; + use crate::rtc::proto::event::VideoLayerSetting; + use crate::rtc::proto::models::{Codec, VideoDimension}; fn video_option(name: &str) -> PublishOption { PublishOption { @@ -396,7 +397,7 @@ mod tests { assert!(offer.sdp.contains("a=rid:f send")); assert!(offer.sdp.contains("a=simulcast:send q;h;f")); pc.set_local_description(offer).await.expect("set offer"); - local.apply_video_layer_settings(&[super::super::proto::event::VideoLayerSetting { + local.apply_video_layer_settings(&[VideoLayerSetting { name: "h".to_owned(), active: false, max_bitrate: 450_000, diff --git a/src/rtc/peer/subscriber.rs b/src/rtc/peer/subscriber.rs new file mode 100644 index 0000000..2c58b6b --- /dev/null +++ b/src/rtc/peer/subscriber.rs @@ -0,0 +1,51 @@ +//! Subscriber-side SDP negotiation: answering the SFU's `SubscriberOffer`. + +use std::sync::Arc; + +use webrtc::peer_connection::RTCPeerConnection; +use webrtc::peer_connection::sdp::session_description::RTCSessionDescription; + +use super::ice::PendingIce; +use crate::rtc::error::{NegotiationError, Result, RtcError}; +use crate::rtc::proto::event::SubscriberOffer; +use crate::rtc::proto::models::PeerType; +use crate::rtc::proto::signal; +use crate::rtc::signal::SignalClient; + +/// Answer an SFU subscriber offer and post the answer over Twirp. +pub async fn negotiate_subscriber( + subscriber: &Arc, + signal: &SignalClient, + session_id: &str, + offer: SubscriberOffer, + pending_ice: &Arc, +) -> Result<()> { + let remote = RTCSessionDescription::offer(offer.sdp) + .map_err(|e| RtcError::Negotiation(NegotiationError(e.to_string())))?; + subscriber + .set_remote_description(remote) + .await + .map_err(|e| RtcError::Negotiation(NegotiationError(e.to_string())))?; + // The remote description now exists: release any candidates the SFU trickled + // before this offer arrived. + pending_ice.flush(subscriber, PeerType::Subscriber).await; + let answer = subscriber + .create_answer(None) + .await + .map_err(|e| RtcError::Negotiation(NegotiationError(e.to_string())))?; + subscriber + .set_local_description(answer.clone()) + .await + .map_err(|e| RtcError::Negotiation(NegotiationError(e.to_string())))?; + + signal + .send_answer(signal::SendAnswerRequest { + peer_type: PeerType::Subscriber as i32, + sdp: answer.sdp, + session_id: session_id.to_owned(), + negotiation_id: offer.negotiation_id, + }) + .await?; + tracing::debug!(session_id, "stream.rtc.subscriber.answer_sent"); + Ok(()) +} From c6e9a3373efc6405d2a8a06978e38593dfbc588c Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Wed, 23 Sep 2026 14:17:58 +0200 Subject: [PATCH 3/7] refactor: consolidate SFU code under sfu/ --- README.md | 12 ++++++------ src/rtc/join/lifecycle.rs | 26 ++++++++++++++------------ src/rtc/join/mod.rs | 4 ++-- src/rtc/join/reconnect_runtime.rs | 2 +- src/rtc/mod.rs | 15 +++++++-------- src/rtc/peer/ice.rs | 2 +- src/rtc/peer/publisher.rs | 2 +- src/rtc/peer/subscriber.rs | 2 +- src/rtc/sfu/mod.rs | 5 +++++ src/rtc/{ => sfu}/signal.rs | 8 ++++---- src/rtc/{sfu_ws.rs => sfu/ws.rs} | 8 ++++---- src/rtc/stats.rs | 2 +- 12 files changed, 47 insertions(+), 41 deletions(-) create mode 100644 src/rtc/sfu/mod.rs rename src/rtc/{ => sfu}/signal.rs (98%) rename src/rtc/{sfu_ws.rs => sfu/ws.rs} (97%) diff --git a/README.md b/README.md index 5e67747..9dde746 100644 --- a/README.md +++ b/README.md @@ -58,9 +58,9 @@ Rust crate, published as [`getstream`](https://crates.io/crates/getstream). | [`ClientConfig`](https://docs.rs/getstream/latest/getstream/struct.ClientConfig.html) | HTTP timeouts, retries, and payload limits | | [`webhook`](https://docs.rs/getstream/latest/getstream/webhook/index.html) | Signature verification and typed events | -The wire-level `rtc` transport modules (`proto`, `peer`, `sfu_ws`, `signal`, -`publisher`, `tracer`, `coordinator_ws`) are public because they track Stream's -SFU protocol, but they are exempt from compatibility guarantees. +The wire-level `rtc` transport modules (`proto`, `peer`, `sfu`, `tracer`, +`coordinator_ws`) are public because they track Stream's SFU protocol, but they +are exempt from compatibility guarantees. ## Requirements @@ -100,9 +100,9 @@ The API reference is published at [docs.rs/getstream](https://docs.rs/getstream) from a checkout, generate it locally with `cargo doc --open`. This is a `0.x` preview, so minor releases may contain breaking changes. The -wire-level `rtc` transport modules (`proto`, `peer`, `sfu_ws`, `signal`, -`publisher`, `tracer`, `coordinator_ws`) track Stream's SFU protocol directly and -are exempt from compatibility guarantees at any version bump. +wire-level `rtc` transport modules (`proto`, `peer`, `sfu`, `tracer`, +`coordinator_ws`) track Stream's SFU protocol directly and are exempt from +compatibility guarantees at any version bump. ## Getting started diff --git a/src/rtc/join/lifecycle.rs b/src/rtc/join/lifecycle.rs index f672bfa..9ef62d2 100644 --- a/src/rtc/join/lifecycle.rs +++ b/src/rtc/join/lifecycle.rs @@ -453,18 +453,20 @@ impl RtcCore { &self.cid(), attempt, )?; - let (mut sender, mut receiver) = - match sfu_ws::connect_with_limit(&ws_url, self.client.max_websocket_message_bytes()) - .await - { - Ok(pair) => pair, - Err(e) => { - signal.trace("signal.close", json!(e.to_string())); - return Err(RtcError::WsConnection( - super::super::error::WsConnectionError::transport(e.to_string()), - )); - } - }; + let (mut sender, mut receiver) = match ws::connect_with_limit( + &ws_url, + self.client.max_websocket_message_bytes(), + ) + .await + { + Ok(pair) => pair, + Err(e) => { + signal.trace("signal.close", json!(e.to_string())); + return Err(RtcError::WsConnection( + super::super::error::WsConnectionError::transport(e.to_string()), + )); + } + }; signal.trace("signal.ws.open", json!(credentials.server.edge_name)); // Build + send the JoinRequest. `fast_reconnect` is deprecated upstream; diff --git a/src/rtc/join/mod.rs b/src/rtc/join/mod.rs index da65942..fe5eafa 100644 --- a/src/rtc/join/mod.rs +++ b/src/rtc/join/mod.rs @@ -64,8 +64,8 @@ use super::reconnect::{ strategy_after_signal_close, }; use super::remote_track::{RemoteParticipant, RemoteTrack}; -use super::sfu_ws::{self, SfuReceiver, SfuSender}; -use super::signal::SignalClient; +use super::sfu::signal::SignalClient; +use super::sfu::ws::{self, SfuReceiver, SfuSender}; use super::stats::{self, StatsReporter, StatsReporterParts}; use super::subscriptions::{SubscriptionConfig, SubscriptionTarget, TrackKey}; use super::tracer::Tracer; diff --git a/src/rtc/join/reconnect_runtime.rs b/src/rtc/join/reconnect_runtime.rs index 9f52a76..59c187e 100644 --- a/src/rtc/join/reconnect_runtime.rs +++ b/src/rtc/join/reconnect_runtime.rs @@ -804,7 +804,7 @@ impl RtcCore { self.reconnect_attempts.load(Ordering::SeqCst), )?; let (mut next_sender, mut receiver) = - sfu_ws::connect_with_limit(&ws_url, self.client.max_websocket_message_bytes()).await?; + ws::connect_with_limit(&ws_url, self.client.max_websocket_message_bytes()).await?; let mut join_request = JoinRequest { token: credentials.token, session_id: session_id.clone(), diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 442c0c4..28fec08 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -2,8 +2,8 @@ //! Cargo feature. //! //! The wire layer holds the generated protobuf types ([`proto`]), the Twirp -//! signal client ([`signal`]), the SFU protobuf WebSocket ([`sfu_ws`]), and the -//! coordinator auth WebSocket ([`coordinator_ws`]). +//! signal client ([`sfu::signal`]), the SFU protobuf WebSocket ([`sfu::ws`]), and +//! the coordinator auth WebSocket ([`coordinator_ws`]). //! //! The participant layer sits on top: the [`coordinator`] join REST, dual //! publisher/subscriber PeerConnections ([`peer`]), and the [`join`] state @@ -14,8 +14,8 @@ //! //! # Stability //! -//! The wire-layer modules — [`proto`], [`peer`], [`sfu_ws`], [`signal`], -//! [`tracer`], and [`coordinator_ws`] — mirror Stream's SFU +//! The wire-layer modules — [`proto`], [`peer`], [`sfu`], [`tracer`], and +//! [`coordinator_ws`] — mirror Stream's SFU //! protocol and change with it. They are exempt from this crate's compatibility //! guarantees at any version bump. Prefer [`crate::Call`], [`RtcClient`], and //! the re-exports below, which are covered by the crate's semver policy. @@ -35,8 +35,7 @@ pub mod proto; mod publish_options; pub mod reconnect; pub mod remote_track; -pub mod sfu_ws; -pub mod signal; +pub mod sfu; pub mod stats; pub mod subscriptions; pub mod tracer; @@ -69,8 +68,8 @@ pub use reconnect::{ DEFAULT_MAX_JOIN_RETRIES, JoinAttemptOutcome, ReconnectStrategy, retry_interval, }; pub use remote_track::{Codec, RemoteParticipant, RemoteTrack}; -pub use sfu_ws::{SfuReceiver, SfuSender}; -pub use signal::SignalClient; +pub use sfu::signal::SignalClient; +pub use sfu::ws::{SfuReceiver, SfuSender}; pub use stats::{DEFAULT_REPORTING_INTERVAL_MS, reporting_interval}; pub use subscriptions::{SubscriptionConfig, SubscriptionTarget}; pub use tracer::{TraceRecord, Tracer}; diff --git a/src/rtc/peer/ice.rs b/src/rtc/peer/ice.rs index 679827f..1029f39 100644 --- a/src/rtc/peer/ice.rs +++ b/src/rtc/peer/ice.rs @@ -8,7 +8,7 @@ use webrtc::peer_connection::RTCPeerConnection; use crate::rtc::error::Result; use crate::rtc::proto::models::{self, PeerType}; -use crate::rtc::signal::SignalClient; +use crate::rtc::sfu::signal::SignalClient; use crate::rtc::tracer::Tracer; /// Buffers remote ICE candidates that arrive before a PeerConnection's remote diff --git a/src/rtc/peer/publisher.rs b/src/rtc/peer/publisher.rs index b5b3e72..fc21705 100644 --- a/src/rtc/peer/publisher.rs +++ b/src/rtc/peer/publisher.rs @@ -19,7 +19,7 @@ use crate::rtc::error::{NegotiationError, Result, RtcError}; use crate::rtc::local_track::LocalTrack; use crate::rtc::proto::models::{PublishOption, TrackInfo, TrackType}; use crate::rtc::proto::signal::SetPublisherRequest; -use crate::rtc::signal::SignalClient; +use crate::rtc::sfu::signal::SignalClient; /// Renegotiate the publisher PeerConnection with the SFU for `tracks`. /// diff --git a/src/rtc/peer/subscriber.rs b/src/rtc/peer/subscriber.rs index 2c58b6b..a58b5b3 100644 --- a/src/rtc/peer/subscriber.rs +++ b/src/rtc/peer/subscriber.rs @@ -10,7 +10,7 @@ use crate::rtc::error::{NegotiationError, Result, RtcError}; use crate::rtc::proto::event::SubscriberOffer; use crate::rtc::proto::models::PeerType; use crate::rtc::proto::signal; -use crate::rtc::signal::SignalClient; +use crate::rtc::sfu::signal::SignalClient; /// Answer an SFU subscriber offer and post the answer over Twirp. pub async fn negotiate_subscriber( diff --git a/src/rtc/sfu/mod.rs b/src/rtc/sfu/mod.rs new file mode 100644 index 0000000..e749b3f --- /dev/null +++ b/src/rtc/sfu/mod.rs @@ -0,0 +1,5 @@ +//! SFU transports: the Twirp signal client ([`signal`]) and the protobuf +//! WebSocket ([`ws`]). + +pub mod signal; +pub mod ws; diff --git a/src/rtc/signal.rs b/src/rtc/sfu/signal.rs similarity index 98% rename from src/rtc/signal.rs rename to src/rtc/sfu/signal.rs index 8752545..f826f99 100644 --- a/src/rtc/signal.rs +++ b/src/rtc/sfu/signal.rs @@ -24,10 +24,10 @@ use url::Url; use crate::client::DEFAULT_MAX_RESPONSE_BODY_BYTES; -use super::error::{Result, RtcError, TwirpError}; -use super::identity; -use super::proto::{models, signal}; -use super::tracer::Tracer; +use crate::rtc::error::{Result, RtcError, TwirpError}; +use crate::rtc::identity; +use crate::rtc::proto::{models, signal}; +use crate::rtc::tracer::Tracer; /// Fully-qualified Twirp service name (`package.Service`). const SERVICE: &str = "stream.video.sfu.signal.SignalServer"; diff --git a/src/rtc/sfu_ws.rs b/src/rtc/sfu/ws.rs similarity index 97% rename from src/rtc/sfu_ws.rs rename to src/rtc/sfu/ws.rs index 29a8120..f0541bf 100644 --- a/src/rtc/sfu_ws.rs +++ b/src/rtc/sfu/ws.rs @@ -13,7 +13,7 @@ //! This module intentionally provides **only** the framed transport: connect, //! a typed send path, and an async event receiver. The join handshake and //! reconnect logic (waiting for `JoinResponse`, ping cadence, health watchdog) -//! belong to the join state machine in [`super::join`]. +//! belong to the join state machine in [`crate::rtc::join`]. use bytes::Bytes; use futures_util::stream::{SplitSink, SplitStream}; @@ -26,8 +26,8 @@ use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async_with_conf use crate::client::DEFAULT_MAX_WEBSOCKET_MESSAGE_BYTES; -use super::error::{Result, RtcError}; -use super::proto::event::{ +use crate::rtc::error::{Result, RtcError}; +use crate::rtc::proto::event::{ HealthCheckRequest, JoinRequest, LeaveCallRequest, SfuEvent, SfuRequest, sfu_request, }; @@ -38,7 +38,7 @@ type WsStream = WebSocketStream>; /// /// `endpoint` is `credentials.server.ws_endpoint` with the query string the /// SFU expects (`?attempt=…&user_id=…&api_key=…&user_session_id=…&cid=…`), -/// built by the join code in [`super::join`]. +/// built by the join code in [`crate::rtc::join`]. pub async fn connect(endpoint: &str) -> Result<(SfuSender, SfuReceiver)> { connect_with_limit(endpoint, DEFAULT_MAX_WEBSOCKET_MESSAGE_BYTES).await } diff --git a/src/rtc/stats.rs b/src/rtc/stats.rs index 0d51b05..2421267 100644 --- a/src/rtc/stats.rs +++ b/src/rtc/stats.rs @@ -27,7 +27,7 @@ use super::coordinator::StatsOptions; use super::error::Result; use super::identity; use super::proto::signal::SendStatsRequest; -use super::signal::SignalClient; +use super::sfu::signal::SignalClient; use super::tracer::{TraceRecord, Tracer, now_ms}; /// Default stats reporting cadence (ms) used when the coordinator does not From 3aca3c8b95d32b6fa37627ceb042766714761556 Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Wed, 23 Sep 2026 14:36:33 +0200 Subject: [PATCH 4/7] refactor: consolidate coordinator code under coordinator/ --- README.md | 4 ++-- src/rtc/coordinator/mod.rs | 19 +++++++++++++++++++ .../{coordinator.rs => coordinator/rest.rs} | 9 +-------- .../{coordinator_ws.rs => coordinator/ws.rs} | 6 +++--- src/rtc/join/lifecycle.rs | 4 ++-- src/rtc/join/mod.rs | 2 +- src/rtc/mod.rs | 11 +++++------ 7 files changed, 33 insertions(+), 22 deletions(-) create mode 100644 src/rtc/coordinator/mod.rs rename src/rtc/{coordinator.rs => coordinator/rest.rs} (94%) rename src/rtc/{coordinator_ws.rs => coordinator/ws.rs} (99%) diff --git a/README.md b/README.md index 9dde746..88eae92 100644 --- a/README.md +++ b/README.md @@ -59,7 +59,7 @@ Rust crate, published as [`getstream`](https://crates.io/crates/getstream). | [`webhook`](https://docs.rs/getstream/latest/getstream/webhook/index.html) | Signature verification and typed events | The wire-level `rtc` transport modules (`proto`, `peer`, `sfu`, `tracer`, -`coordinator_ws`) are public because they track Stream's SFU protocol, but they +`coordinator::ws`) are public because they track Stream's SFU protocol, but they are exempt from compatibility guarantees. ## Requirements @@ -101,7 +101,7 @@ from a checkout, generate it locally with `cargo doc --open`. This is a `0.x` preview, so minor releases may contain breaking changes. The wire-level `rtc` transport modules (`proto`, `peer`, `sfu`, `tracer`, -`coordinator_ws`) track Stream's SFU protocol directly and are exempt from +`coordinator::ws`) track Stream's SFU protocol directly and are exempt from compatibility guarantees at any version bump. ## Getting started diff --git a/src/rtc/coordinator/mod.rs b/src/rtc/coordinator/mod.rs new file mode 100644 index 0000000..dcbb636 --- /dev/null +++ b/src/rtc/coordinator/mod.rs @@ -0,0 +1,19 @@ +//! Coordinator `JoinCall` REST + location discovery. +//! +//! `POST /api/v2/video/call/{type}/{id}/join` is a **user-token** operation +//! (unlike the server-token REST in [`crate::video`]): the coordinator returns the SFU +//! credentials the participant needs — the Twirp base URL, the SFU token, the +//! signaling WebSocket endpoint, the ICE servers, and the stats options to +//! cache. Ported from stream-py `connection_utils.join_call_coordinator_request` +//! and the OpenAPI coordinator models shared by all SDKs. +//! +//! The coordinator auth WebSocket lives in [`ws`]. + +mod rest; +pub mod ws; + +pub(crate) use rest::join_call; +pub use rest::{ + Credentials, FALLBACK_LOCATION, IceServer, JoinCallRequest, JoinCallResponse, SfuServer, + StatsOptions, discover_location, +}; diff --git a/src/rtc/coordinator.rs b/src/rtc/coordinator/rest.rs similarity index 94% rename from src/rtc/coordinator.rs rename to src/rtc/coordinator/rest.rs index a566013..3bdf72b 100644 --- a/src/rtc/coordinator.rs +++ b/src/rtc/coordinator/rest.rs @@ -1,11 +1,4 @@ //! Coordinator `JoinCall` REST + location discovery. -//! -//! `POST /api/v2/video/call/{type}/{id}/join` is a **user-token** operation -//! (unlike the server-token REST in [`crate::video`]): the coordinator returns the SFU -//! credentials the participant needs — the Twirp base URL, the SFU token, the -//! signaling WebSocket endpoint, the ICE servers, and the stats options to -//! cache. Ported from stream-py `connection_utils.join_call_coordinator_request` -//! and the OpenAPI coordinator models shared by all SDKs. use std::sync::Arc; use std::time::Duration; @@ -17,7 +10,7 @@ use serde_json::Value; use crate::client::Client; use crate::models::{CallRequest, CallResponse, MemberResponse}; -use super::error::{Result, RtcError}; +use crate::rtc::error::{Result, RtcError}; /// The CloudFront hint endpoint used to discover the caller's edge location /// (JS `getLocationHint` / stream-py `location_discovery`). diff --git a/src/rtc/coordinator_ws.rs b/src/rtc/coordinator/ws.rs similarity index 99% rename from src/rtc/coordinator_ws.rs rename to src/rtc/coordinator/ws.rs index c24a048..30ba4a2 100644 --- a/src/rtc/coordinator_ws.rs +++ b/src/rtc/coordinator/ws.rs @@ -10,7 +10,7 @@ //! - the client pings with a `health.check` event (~20s in videosdk) //! //! This module covers connect, auth, and typed event decode only. The join and -//! call-watch flow that builds on it lives in [`super::join`]. +//! call-watch flow that builds on it lives in [`crate::rtc::join`]. use std::time::Duration; @@ -28,8 +28,8 @@ use url::Url; use crate::client::{DEFAULT_BASE_URL, DEFAULT_MAX_WEBSOCKET_MESSAGE_BYTES}; -use super::error::{Result, RtcError, SfuTimeoutError}; -use super::identity; +use crate::rtc::error::{Result, RtcError, SfuTimeoutError}; +use crate::rtc::identity; type WsStream = WebSocketStream>; diff --git a/src/rtc/join/lifecycle.rs b/src/rtc/join/lifecycle.rs index 9ef62d2..8f5f568 100644 --- a/src/rtc/join/lifecycle.rs +++ b/src/rtc/join/lifecycle.rs @@ -696,8 +696,8 @@ impl RtcCore { return Err(join_cancelled()); } let auth = WsAuthMessage::video(user_token, ConnectUserDetails::new(user_id)); - let coordinator_url = coordinator_ws::coordinator_ws_url(self.client.base_url())?; - let (mut coordinator, mut events, connected) = coordinator_ws::connect_with_limit( + let coordinator_url = coordinator::ws::coordinator_ws_url(self.client.base_url())?; + let (mut coordinator, mut events, connected) = coordinator::ws::connect_with_limit( coordinator_url.as_str(), &self.api_key, user_id, diff --git a/src/rtc/join/mod.rs b/src/rtc/join/mod.rs index fe5eafa..71d1a86 100644 --- a/src/rtc/join/mod.rs +++ b/src/rtc/join/mod.rs @@ -49,8 +49,8 @@ use crate::client::Client; use crate::models::CallRequest; use super::client::UserTokenSource; +use super::coordinator::ws::{ConnectUserDetails, CoordinatorEvent, WsAuthMessage}; use super::coordinator::{self, Credentials, JoinCallRequest, StatsOptions}; -use super::coordinator_ws::{self, ConnectUserDetails, CoordinatorEvent, WsAuthMessage}; use super::error::{Result, RtcError, SfuJoinError, SfuTimeoutError}; use super::identity; use super::local_track::LocalTrack; diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 28fec08..290446c 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -3,7 +3,7 @@ //! //! The wire layer holds the generated protobuf types ([`proto`]), the Twirp //! signal client ([`sfu::signal`]), the SFU protobuf WebSocket ([`sfu::ws`]), and -//! the coordinator auth WebSocket ([`coordinator_ws`]). +//! the coordinator auth WebSocket ([`coordinator::ws`]). //! //! The participant layer sits on top: the [`coordinator`] join REST, dual //! publisher/subscriber PeerConnections ([`peer`]), and the [`join`] state @@ -15,7 +15,7 @@ //! # Stability //! //! The wire-layer modules — [`proto`], [`peer`], [`sfu`], [`tracer`], and -//! [`coordinator_ws`] — mirror Stream's SFU +//! [`coordinator::ws`] — mirror Stream's SFU //! protocol and change with it. They are exempt from this crate's compatibility //! guarantees at any version bump. Prefer [`crate::Call`], [`RtcClient`], and //! the re-exports below, which are covered by the crate's semver policy. @@ -23,7 +23,6 @@ pub mod client; mod codecs; pub mod coordinator; -pub mod coordinator_ws; pub mod error; pub mod identity; pub mod join; @@ -42,12 +41,12 @@ pub mod tracer; pub mod video_frame; pub use client::{RtcCall, RtcClient, TokenFuture, TokenProvider}; +pub use coordinator::ws::{ + ConnectUserDetails, CoordinatorEvent, CoordinatorEvents, CoordinatorWs, WsAuthMessage, +}; pub use coordinator::{ Credentials, IceServer, JoinCallRequest, JoinCallResponse, SfuServer, StatsOptions, }; -pub use coordinator_ws::{ - ConnectUserDetails, CoordinatorEvent, CoordinatorEvents, CoordinatorWs, WsAuthMessage, -}; pub use error::{ ErrorFromResponse, NegotiationError, Result as RtcResult, RtcError, SfuJoinError, SfuTimeoutError, TwirpError, WsConnectionError, is_join_error_code, From 3b7e7461557320cd9edb9acccff0b43f2747de40 Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Wed, 23 Sep 2026 15:06:07 +0200 Subject: [PATCH 5/7] refactor: consolidate tracks code under tracks/ --- src/rtc/client.rs | 3 +-- src/rtc/codecs/vpx.rs | 4 +-- src/rtc/join/mod.rs | 3 +-- src/rtc/join/publication.rs | 2 +- src/rtc/mod.rs | 14 +++++------ src/rtc/peer/publisher.rs | 4 +-- src/rtc/{ => tracks}/layers.rs | 2 +- src/rtc/{local_track.rs => tracks/local.rs} | 25 ++++++++++--------- src/rtc/tracks/mod.rs | 12 +++++++++ src/rtc/{remote_track.rs => tracks/remote.rs} | 16 ++++++------ src/rtc/video_frame.rs | 4 +-- 11 files changed, 49 insertions(+), 40 deletions(-) rename src/rtc/{ => tracks}/layers.rs (98%) rename src/rtc/{local_track.rs => tracks/local.rs} (99%) create mode 100644 src/rtc/tracks/mod.rs rename src/rtc/{remote_track.rs => tracks/remote.rs} (98%) diff --git a/src/rtc/client.rs b/src/rtc/client.rs index 1235efa..1d73aab 100644 --- a/src/rtc/client.rs +++ b/src/rtc/client.rs @@ -16,10 +16,9 @@ use crate::token::{self, TokenOptions}; use super::error::{Result, RtcError}; use super::join::{CallEvent, CallStateSnapshot, CallingState, JoinCallData, RtcCore}; -use super::local_track::{LocalAudioTrack, LocalTrack, LocalVideoTrack}; use super::proto::models::TrackType; -use super::remote_track::{RemoteParticipant, RemoteTrack}; use super::subscriptions::{SubscriptionConfig, SubscriptionTarget}; +use super::tracks::{LocalAudioTrack, LocalTrack, LocalVideoTrack, RemoteParticipant, RemoteTrack}; /// Boxed future returned by an RTC [`TokenProvider`]. pub type TokenFuture = Pin> + Send + 'static>>; diff --git a/src/rtc/codecs/vpx.rs b/src/rtc/codecs/vpx.rs index c9b52c1..7fc50c6 100644 --- a/src/rtc/codecs/vpx.rs +++ b/src/rtc/codecs/vpx.rs @@ -1,7 +1,7 @@ //! Minimal libvpx VP8/VP9 encoder and decoder for the outbound and inbound //! video paths. //! -//! [`LocalVideoTrack`](crate::rtc::local_track::LocalVideoTrack) needs to turn raw +//! [`LocalVideoTrack`](crate::rtc::LocalVideoTrack) needs to turn raw //! I420 frames into encoded VP8/VP9 for the SFU. The `vpx-encode` crate is too //! restrictive for realtime streaming — it hardcodes the encoder config, so it //! cannot set `g_lag_in_frames = 0` (VP9 otherwise buffers frames and emits @@ -10,7 +10,7 @@ //! that late subscribers miss). This module binds `libvpx` directly (via //! `env-libvpx-sys`, exposed as `vpx_sys`) with the correct realtime config. //! -//! [`RemoteTrack::next_video_frame`](crate::rtc::remote_track::RemoteTrack::next_video_frame) +//! [`RemoteTrack::next_video_frame`](crate::rtc::RemoteTrack::next_video_frame) //! needs raw frames, not RTP: an agent that wants to *see* the call has to turn //! reassembled VP8/VP9 samples into pixels. The decoder symbols come from the //! same `env-libvpx-sys` bindings, so inbound video adds no new native diff --git a/src/rtc/join/mod.rs b/src/rtc/join/mod.rs index 71d1a86..319195d 100644 --- a/src/rtc/join/mod.rs +++ b/src/rtc/join/mod.rs @@ -53,7 +53,6 @@ use super::coordinator::ws::{ConnectUserDetails, CoordinatorEvent, WsAuthMessage use super::coordinator::{self, Credentials, JoinCallRequest, StatsOptions}; use super::error::{Result, RtcError, SfuJoinError, SfuTimeoutError}; use super::identity; -use super::local_track::LocalTrack; use super::peer::{self, PendingIce, negotiate_subscriber, publisher, register_ice_trickle}; use super::proto::event::{self, JoinRequest, JoinResponse, ReconnectDetails, SfuEvent, sfu_event}; use super::proto::models::{self, PeerType, TrackType}; @@ -63,12 +62,12 @@ use super::reconnect::{ self, FailureCaps, ReconnectStrategy, SlidingWindowRateLimiter, escalate_strategy, strategy_after_signal_close, }; -use super::remote_track::{RemoteParticipant, RemoteTrack}; use super::sfu::signal::SignalClient; use super::sfu::ws::{self, SfuReceiver, SfuSender}; use super::stats::{self, StatsReporter, StatsReporterParts}; use super::subscriptions::{SubscriptionConfig, SubscriptionTarget, TrackKey}; use super::tracer::Tracer; +use super::tracks::{LocalTrack, RemoteParticipant, RemoteTrack}; use serde_json::json; diff --git a/src/rtc/join/publication.rs b/src/rtc/join/publication.rs index c6a4914..c50e424 100644 --- a/src/rtc/join/publication.rs +++ b/src/rtc/join/publication.rs @@ -3,10 +3,10 @@ use std::collections::{HashMap, HashSet}; use super::super::error::Result; -use super::super::local_track::LocalTrack; use super::super::proto::event; use super::super::proto::models::{self, TrackType}; use crate::rtc::peer::publisher; +use crate::rtc::tracks::LocalTrack; #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(super) enum PublicationStatus { diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 290446c..2ca636e 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -26,18 +26,16 @@ pub mod coordinator; pub mod error; pub mod identity; pub mod join; -mod layers; -pub mod local_track; pub mod pcm; pub mod peer; pub mod proto; mod publish_options; pub mod reconnect; -pub mod remote_track; pub mod sfu; pub mod stats; pub mod subscriptions; pub mod tracer; +mod tracks; pub mod video_frame; pub use client::{RtcCall, RtcClient, TokenFuture, TokenProvider}; @@ -53,10 +51,6 @@ pub use error::{ }; pub use identity::{CLIENT_TYPE, SDK_TYPE, client_details, client_header}; pub use join::{CallEvent, CallStateSnapshot, CallingState, JoinCallData, RtcCore}; -pub use local_track::{ - LocalAudioTrack, LocalAudioTrackConfig, LocalTrack, LocalVideoTrack, LocalVideoTrackConfig, - RtpPacket, VideoLayering, audio_level_dbov, -}; pub use pcm::chunk::Pad; pub use pcm::convert::G711_SAMPLE_RATE; pub use pcm::{ @@ -66,12 +60,16 @@ pub use publish_options::{ClientPublishOptions, PreferredVideoCodec}; pub use reconnect::{ DEFAULT_MAX_JOIN_RETRIES, JoinAttemptOutcome, ReconnectStrategy, retry_interval, }; -pub use remote_track::{Codec, RemoteParticipant, RemoteTrack}; pub use sfu::signal::SignalClient; pub use sfu::ws::{SfuReceiver, SfuSender}; pub use stats::{DEFAULT_REPORTING_INTERVAL_MS, reporting_interval}; pub use subscriptions::{SubscriptionConfig, SubscriptionTarget}; pub use tracer::{TraceRecord, Tracer}; +pub use tracks::{ + Codec, LocalAudioTrack, LocalAudioTrackConfig, LocalTrack, LocalVideoTrack, + LocalVideoTrackConfig, RemoteParticipant, RemoteTrack, RtpPacket, VideoLayering, + audio_level_dbov, +}; pub use video_frame::VideoFrame; #[cfg(test)] diff --git a/src/rtc/peer/publisher.rs b/src/rtc/peer/publisher.rs index fc21705..d7df3d2 100644 --- a/src/rtc/peer/publisher.rs +++ b/src/rtc/peer/publisher.rs @@ -16,10 +16,10 @@ use webrtc::peer_connection::sdp::session_description::RTCSessionDescription; use webrtc::peer_connection::signaling_state::RTCSignalingState; use crate::rtc::error::{NegotiationError, Result, RtcError}; -use crate::rtc::local_track::LocalTrack; use crate::rtc::proto::models::{PublishOption, TrackInfo, TrackType}; use crate::rtc::proto::signal::SetPublisherRequest; use crate::rtc::sfu::signal::SignalClient; +use crate::rtc::tracks::LocalTrack; /// Renegotiate the publisher PeerConnection with the SFU for `tracks`. /// @@ -297,10 +297,10 @@ pub(crate) async fn build_track_infos( #[cfg(test)] mod tests { use super::*; - use crate::rtc::local_track::{LocalVideoTrack, LocalVideoTrackConfig}; use crate::rtc::peer; use crate::rtc::proto::event::VideoLayerSetting; use crate::rtc::proto::models::{Codec, VideoDimension}; + use crate::rtc::tracks::{LocalVideoTrack, LocalVideoTrackConfig}; fn video_option(name: &str) -> PublishOption { PublishOption { diff --git a/src/rtc/layers.rs b/src/rtc/tracks/layers.rs similarity index 98% rename from src/rtc/layers.rs rename to src/rtc/tracks/layers.rs index 07e10e2..93dee60 100644 --- a/src/rtc/layers.rs +++ b/src/rtc/tracks/layers.rs @@ -6,7 +6,7 @@ use std::num::NonZeroU8; -use super::proto::models::{PublishOption, VideoDimension, VideoLayer, VideoQuality}; +use crate::rtc::proto::models::{PublishOption, VideoDimension, VideoLayer, VideoQuality}; pub(crate) const SIMULCAST_RIDS: [&str; 3] = ["q", "h", "f"]; diff --git a/src/rtc/local_track.rs b/src/rtc/tracks/local.rs similarity index 99% rename from src/rtc/local_track.rs rename to src/rtc/tracks/local.rs index 8d75233..5ca186a 100644 --- a/src/rtc/local_track.rs +++ b/src/rtc/tracks/local.rs @@ -43,15 +43,15 @@ use webrtc::track::track_local::TrackLocal; use webrtc::track::track_local::TrackLocalWriter; use webrtc::track::track_local::track_local_static_rtp::TrackLocalStaticRTP; -use super::codecs::h264::{H264Encoder, validate_h264_encode_request}; -use super::codecs::rtp_h264::H264RtpPacketizer; -use super::codecs::rtp_vpx::VpxRtpPacketizer; -use super::codecs::vpx::{Vp9SvcMode, VpxCodec, VpxEncoder, VpxSvcEncoder}; -use super::error::{Result, RtcError}; use super::layers::{PlannedVideoLayer, simulcast_layers, single_layer}; -use super::pcm::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame, StreamResampler, rms_i16}; -use super::proto::event::VideoLayerSetting; -use super::proto::models::{PublishOption, TrackType, VideoLayer}; +use crate::rtc::codecs::h264::{H264Encoder, validate_h264_encode_request}; +use crate::rtc::codecs::rtp_h264::H264RtpPacketizer; +use crate::rtc::codecs::rtp_vpx::VpxRtpPacketizer; +use crate::rtc::codecs::vpx::{Vp9SvcMode, VpxCodec, VpxEncoder, VpxSvcEncoder}; +use crate::rtc::error::{Result, RtcError}; +use crate::rtc::pcm::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame, StreamResampler, rms_i16}; +use crate::rtc::proto::event::VideoLayerSetting; +use crate::rtc::proto::models::{PublishOption, TrackType, VideoLayer}; /// A raw RTP packet (`webrtc::rtp::packet::Packet`), used by the RTP-forward /// republish path. @@ -908,7 +908,7 @@ impl LocalVideoTrack { mime_type: MIME_TYPE_H264.to_owned(), clock_rate: 90_000, channels: 0, - sdp_fmtp_line: super::publish_options::H264_FMTP.to_owned(), + sdp_fmtp_line: crate::rtc::publish_options::H264_FMTP.to_owned(), rtcp_feedback: vec![], }, VideoCodec::H264, @@ -1992,6 +1992,7 @@ impl LocalTrack { #[cfg(test)] mod tests { use super::*; + use crate::rtc::proto::models::{Codec, VideoDimension}; /// One 20 ms frame of 440 Hz tone: FEC and DTX both key off whether the /// frame carries signal, so silence would not exercise either. @@ -2290,7 +2291,7 @@ mod tests { PublishOption { id: 41, track_type: track_type as i32, - codec: Some(super::super::proto::models::Codec { + codec: Some(Codec { name: codec.to_owned(), ..Default::default() }), @@ -2298,7 +2299,7 @@ mod tests { fps: 30, max_spatial_layers: 3, max_temporal_layers: 3, - video_dimension: Some(super::super::proto::models::VideoDimension { + video_dimension: Some(VideoDimension { width: 1280, height: 720, }), @@ -2494,7 +2495,7 @@ mod tests { fn vp9_svc_reconfiguration_preserves_rtp_counters_and_emits_new_ss() { let track = LocalVideoTrack::vp9_svc().expect("VP9 SVC"); let mut option = layered_option(TrackType::Video, "VP9"); - option.video_dimension = Some(super::super::proto::models::VideoDimension { + option.video_dimension = Some(VideoDimension { width: 320, height: 240, }); diff --git a/src/rtc/tracks/mod.rs b/src/rtc/tracks/mod.rs new file mode 100644 index 0000000..f1900bb --- /dev/null +++ b/src/rtc/tracks/mod.rs @@ -0,0 +1,12 @@ +//! Outbound ([`LocalAudioTrack`], [`LocalVideoTrack`]) and inbound +//! ([`RemoteTrack`]) media tracks. + +mod layers; +mod local; +mod remote; + +pub use local::{ + LocalAudioTrack, LocalAudioTrackConfig, LocalTrack, LocalVideoTrack, LocalVideoTrackConfig, + RtpPacket, VideoLayering, audio_level_dbov, +}; +pub use remote::{Codec, RemoteParticipant, RemoteTrack}; diff --git a/src/rtc/remote_track.rs b/src/rtc/tracks/remote.rs similarity index 98% rename from src/rtc/remote_track.rs rename to src/rtc/tracks/remote.rs index fb0c643..e17a28b 100644 --- a/src/rtc/remote_track.rs +++ b/src/rtc/tracks/remote.rs @@ -32,14 +32,14 @@ use webrtc::rtp::codecs::vp8::Vp8Packet; use webrtc::rtp::codecs::vp9::Vp9Packet; use webrtc::track::track_remote::TrackRemote; -use super::codecs::h264::{H264Decoder, access_unit_has_idr}; -use super::codecs::rtp_h264::H264Depacketizer; -use super::codecs::vpx::{VpxCodec, VpxDecoder}; -use super::error::{Result, RtcError}; -use super::local_track::RtpPacket; -use super::pcm::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame}; -use super::proto::models::{self, TrackType}; -use super::video_frame::VideoFrame; +use super::local::RtpPacket; +use crate::rtc::codecs::h264::{H264Decoder, access_unit_has_idr}; +use crate::rtc::codecs::rtp_h264::H264Depacketizer; +use crate::rtc::codecs::vpx::{VpxCodec, VpxDecoder}; +use crate::rtc::error::{Result, RtcError}; +use crate::rtc::pcm::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame}; +use crate::rtc::proto::models::{self, TrackType}; +use crate::rtc::video_frame::VideoFrame; /// How many packets the video [`SampleBuilder`] buffers while waiting for gaps /// to be filled (by NACK/RTX or a reordered arrival) before giving up on a diff --git a/src/rtc/video_frame.rs b/src/rtc/video_frame.rs index 4b6745a..08a5eaa 100644 --- a/src/rtc/video_frame.rs +++ b/src/rtc/video_frame.rs @@ -1,8 +1,8 @@ //! Decoded video frames ([`VideoFrame`]) produced by -//! [`RemoteTrack::next_video_frame`](super::remote_track::RemoteTrack::next_video_frame). +//! [`RemoteTrack::next_video_frame`](crate::rtc::RemoteTrack::next_video_frame). //! //! Frames are packed I420 (YUV 4:2:0 planar) — the same layout -//! [`LocalVideoTrack::write_i420`](super::local_track::LocalVideoTrack::write_i420) +//! [`LocalVideoTrack::write_i420`](crate::rtc::LocalVideoTrack::write_i420) //! consumes, so a decode → transform → re-encode bridge needs no conversion in //! between. //! From 8829b6d01c76f4d4100260a04ecd86601e9a5cc4 Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Wed, 23 Sep 2026 15:29:42 +0200 Subject: [PATCH 6/7] refactor: move the code out of pcm/mod.rs --- src/rtc/pcm/frame.rs | 142 ++++++++++++++++++++++++++++++++++++++++++ src/rtc/pcm/mod.rs | 144 +------------------------------------------ 2 files changed, 145 insertions(+), 141 deletions(-) create mode 100644 src/rtc/pcm/frame.rs diff --git a/src/rtc/pcm/frame.rs b/src/rtc/pcm/frame.rs new file mode 100644 index 0000000..774281f --- /dev/null +++ b/src/rtc/pcm/frame.rs @@ -0,0 +1,142 @@ +//! The [`PcmFrame`] interchange type and the SFU's native audio constants. + +use std::time::Duration; + +/// The SFU's native audio sample rate (Opus internal clock). +pub const OPUS_SAMPLE_RATE: u32 = 48_000; +/// Samples per channel in a 20 ms frame at 48 kHz (the Opus frame we pace on). +pub const FRAME_SAMPLES_20MS: usize = (OPUS_SAMPLE_RATE as usize) / 50; + +/// A block of interleaved 16-bit PCM samples. +/// +/// `samples` is interleaved when `channels > 1` (L, R, L, R, …). This is the +/// public type produced by [`RemoteTrack::next_pcm`](crate::rtc::RemoteTrack) +/// and consumed by [`LocalAudioTrack::write_pcm`](crate::rtc::LocalAudioTrack). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PcmFrame { + /// Interleaved 16-bit samples. + pub samples: Vec, + /// Samples per second (per channel). + pub sample_rate: u32, + /// Channel count (1 = mono, 2 = stereo). + pub channels: u16, +} + +impl PcmFrame { + /// Build a frame from interleaved samples. + pub fn new(samples: Vec, sample_rate: u32, channels: u16) -> Self { + Self { + samples, + sample_rate, + channels: channels.max(1), + } + } + + /// Build a mono frame. + pub fn mono(samples: Vec, sample_rate: u32) -> Self { + Self::new(samples, sample_rate, 1) + } + + /// A silent frame of `frames` samples per channel. + pub fn silence(frames: usize, sample_rate: u32, channels: u16) -> Self { + let channels = channels.max(1); + Self::new(vec![0; frames * channels as usize], sample_rate, channels) + } + + /// Number of samples per channel. + pub fn frames(&self) -> usize { + self.samples.len() / (self.channels.max(1) as usize) + } + + /// Whether the frame carries no samples. + pub fn is_empty(&self) -> bool { + self.samples.is_empty() + } + + /// Whether the frame carries two channels. + pub fn is_stereo(&self) -> bool { + self.channels == 2 + } + + /// The playback duration of this block. + pub fn duration(&self) -> Duration { + if self.sample_rate == 0 { + return Duration::ZERO; + } + Duration::from_secs_f64(self.frames() as f64 / f64::from(self.sample_rate)) + } + + /// The playback duration of this block in fractional milliseconds. + pub fn duration_ms(&self) -> f64 { + self.duration().as_secs_f64() * 1000.0 + } + + /// Root-mean-square amplitude across all samples, normalized to `[0, 1]`. + /// + /// Handy for asserting a republished stream is non-silent (energy above a + /// small threshold) without pulling in a DSP crate. + pub fn rms(&self) -> f64 { + rms_i16(&self.samples) + } + + /// Number of samples per channel a `duration` of audio occupies at this + /// frame's sample rate. + pub(crate) fn frames_in(&self, duration: Duration) -> usize { + (duration.as_secs_f64() * f64::from(self.sample_rate)) as usize + } +} + +/// Root-mean-square amplitude of `samples`, normalized to `[0, 1]`. +/// +/// The slice form lets the outbound pacer measure the 20 ms block it is about +/// to encode without building a [`PcmFrame`] around it. +pub(crate) fn rms_i16(samples: &[i16]) -> f64 { + if samples.is_empty() { + return 0.0; + } + let sum_sq: f64 = samples + .iter() + .map(|&s| { + let v = f64::from(s) / f64::from(i16::MAX); + v * v + }) + .sum(); + (sum_sq / samples.len() as f64).sqrt() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn rms_of_silence_is_zero_and_tone_is_positive() { + assert_eq!(PcmFrame::mono(vec![0; 480], 48_000).rms(), 0.0); + let tone: Vec = (0..480) + .map(|i| ((i as f64 * 0.1).sin() * 10_000.0) as i16) + .collect(); + assert!(PcmFrame::mono(tone, 48_000).rms() > 0.05); + } + + #[test] + fn duration_of_20ms_frame() { + let f = PcmFrame::mono(vec![0; FRAME_SAMPLES_20MS], OPUS_SAMPLE_RATE); + assert_eq!(f.duration(), Duration::from_millis(20)); + assert_eq!(f.duration_ms(), 20.0); + } + + #[test] + fn frames_counts_per_channel_not_total_samples() { + let stereo = PcmFrame::new(vec![1, 2, 3, 4, 5, 6], 48_000, 2); + assert_eq!(stereo.frames(), 3); + assert!(stereo.is_stereo()); + assert_eq!(stereo.duration_ms(), 3.0 / 48.0); + } + + #[test] + fn silence_is_sized_per_channel() { + let s = PcmFrame::silence(480, 48_000, 2); + assert_eq!(s.samples.len(), 960); + assert_eq!(s.frames(), 480); + assert_eq!(s.rms(), 0.0); + } +} diff --git a/src/rtc/pcm/mod.rs b/src/rtc/pcm/mod.rs index bb22509..d8116e4 100644 --- a/src/rtc/pcm/mod.rs +++ b/src/rtc/pcm/mod.rs @@ -69,148 +69,10 @@ pub mod chunk; pub mod convert; +mod frame; pub mod resample; -use std::time::Duration; - pub use convert::G711Mapping; +pub(crate) use frame::rms_i16; +pub use frame::{FRAME_SAMPLES_20MS, OPUS_SAMPLE_RATE, PcmFrame}; pub use resample::{Resampler, StreamResampler}; - -/// The SFU's native audio sample rate (Opus internal clock). -pub const OPUS_SAMPLE_RATE: u32 = 48_000; -/// Samples per channel in a 20 ms frame at 48 kHz (the Opus frame we pace on). -pub const FRAME_SAMPLES_20MS: usize = (OPUS_SAMPLE_RATE as usize) / 50; - -/// A block of interleaved 16-bit PCM samples. -/// -/// `samples` is interleaved when `channels > 1` (L, R, L, R, …). This is the -/// public type produced by [`RemoteTrack::next_pcm`](crate::rtc::RemoteTrack) -/// and consumed by [`LocalAudioTrack::write_pcm`](crate::rtc::LocalAudioTrack). -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PcmFrame { - /// Interleaved 16-bit samples. - pub samples: Vec, - /// Samples per second (per channel). - pub sample_rate: u32, - /// Channel count (1 = mono, 2 = stereo). - pub channels: u16, -} - -impl PcmFrame { - /// Build a frame from interleaved samples. - pub fn new(samples: Vec, sample_rate: u32, channels: u16) -> Self { - Self { - samples, - sample_rate, - channels: channels.max(1), - } - } - - /// Build a mono frame. - pub fn mono(samples: Vec, sample_rate: u32) -> Self { - Self::new(samples, sample_rate, 1) - } - - /// A silent frame of `frames` samples per channel. - pub fn silence(frames: usize, sample_rate: u32, channels: u16) -> Self { - let channels = channels.max(1); - Self::new(vec![0; frames * channels as usize], sample_rate, channels) - } - - /// Number of samples per channel. - pub fn frames(&self) -> usize { - self.samples.len() / (self.channels.max(1) as usize) - } - - /// Whether the frame carries no samples. - pub fn is_empty(&self) -> bool { - self.samples.is_empty() - } - - /// Whether the frame carries two channels. - pub fn is_stereo(&self) -> bool { - self.channels == 2 - } - - /// The playback duration of this block. - pub fn duration(&self) -> Duration { - if self.sample_rate == 0 { - return Duration::ZERO; - } - Duration::from_secs_f64(self.frames() as f64 / f64::from(self.sample_rate)) - } - - /// The playback duration of this block in fractional milliseconds. - pub fn duration_ms(&self) -> f64 { - self.duration().as_secs_f64() * 1000.0 - } - - /// Root-mean-square amplitude across all samples, normalized to `[0, 1]`. - /// - /// Handy for asserting a republished stream is non-silent (energy above a - /// small threshold) without pulling in a DSP crate. - pub fn rms(&self) -> f64 { - rms_i16(&self.samples) - } - - /// Number of samples per channel a `duration` of audio occupies at this - /// frame's sample rate. - pub(crate) fn frames_in(&self, duration: Duration) -> usize { - (duration.as_secs_f64() * f64::from(self.sample_rate)) as usize - } -} - -/// Root-mean-square amplitude of `samples`, normalized to `[0, 1]`. -/// -/// The slice form lets the outbound pacer measure the 20 ms block it is about -/// to encode without building a [`PcmFrame`] around it. -pub(crate) fn rms_i16(samples: &[i16]) -> f64 { - if samples.is_empty() { - return 0.0; - } - let sum_sq: f64 = samples - .iter() - .map(|&s| { - let v = f64::from(s) / f64::from(i16::MAX); - v * v - }) - .sum(); - (sum_sq / samples.len() as f64).sqrt() -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn rms_of_silence_is_zero_and_tone_is_positive() { - assert_eq!(PcmFrame::mono(vec![0; 480], 48_000).rms(), 0.0); - let tone: Vec = (0..480) - .map(|i| ((i as f64 * 0.1).sin() * 10_000.0) as i16) - .collect(); - assert!(PcmFrame::mono(tone, 48_000).rms() > 0.05); - } - - #[test] - fn duration_of_20ms_frame() { - let f = PcmFrame::mono(vec![0; FRAME_SAMPLES_20MS], OPUS_SAMPLE_RATE); - assert_eq!(f.duration(), Duration::from_millis(20)); - assert_eq!(f.duration_ms(), 20.0); - } - - #[test] - fn frames_counts_per_channel_not_total_samples() { - let stereo = PcmFrame::new(vec![1, 2, 3, 4, 5, 6], 48_000, 2); - assert_eq!(stereo.frames(), 3); - assert!(stereo.is_stereo()); - assert_eq!(stereo.duration_ms(), 3.0 / 48.0); - } - - #[test] - fn silence_is_sized_per_channel() { - let s = PcmFrame::silence(480, 48_000, 2); - assert_eq!(s.samples.len(), 960); - assert_eq!(s.frames(), 480); - assert_eq!(s.rms(), 0.0); - } -} From 80287960ac7895fcf29b27f4c4b97943add4cb27 Mon Sep 17 00:00:00 2001 From: Daniil Gusev Date: Wed, 23 Sep 2026 15:53:32 +0200 Subject: [PATCH 7/7] refactor: move the code out of rtc/peer/mod.rs --- src/rtc/peer/connection.rs | 376 ++++++++++++++++++++++++++++++++++++ src/rtc/peer/mod.rs | 377 +------------------------------------ 2 files changed, 378 insertions(+), 375 deletions(-) create mode 100644 src/rtc/peer/connection.rs diff --git a/src/rtc/peer/connection.rs b/src/rtc/peer/connection.rs new file mode 100644 index 0000000..f4a8936 --- /dev/null +++ b/src/rtc/peer/connection.rs @@ -0,0 +1,376 @@ +//! webrtc-rs PeerConnection construction and the throwaway generic SDPs. + +use std::sync::Arc; + +use serde_json::json; +use webrtc::api::interceptor_registry::register_default_interceptors; +use webrtc::api::media_engine::{ + MIME_TYPE_H264, MIME_TYPE_OPUS, MIME_TYPE_VP8, MIME_TYPE_VP9, MediaEngine, +}; +use webrtc::api::{API, APIBuilder}; +use webrtc::ice_transport::ice_server::RTCIceServer; +use webrtc::interceptor::registry::Registry; +use webrtc::peer_connection::RTCPeerConnection; +use webrtc::peer_connection::configuration::RTCConfiguration; +use webrtc::rtp_transceiver::RTCPFeedback; +use webrtc::rtp_transceiver::rtp_codec::{ + RTCRtpCodecCapability, RTCRtpCodecParameters, RTCRtpHeaderExtensionCapability, RTPCodecType, +}; +use webrtc::rtp_transceiver::rtp_transceiver_direction::RTCRtpTransceiverDirection; +use webrtc::sdp::extmap::{ + AUDIO_LEVEL_URI, SDES_MID_URI, SDES_REPAIR_RTP_STREAM_ID_URI, SDES_RTP_STREAM_ID_URI, +}; + +use crate::rtc::coordinator::IceServer; +use crate::rtc::error::Result; +use crate::rtc::publish_options::H264_FMTP; +use crate::rtc::tracer::Tracer; + +const OPUS_PAYLOAD_TYPE: u8 = 111; +const VP8_PAYLOAD_TYPE: u8 = 96; +const VP9_PAYLOAD_TYPE: u8 = 98; +const H264_PAYLOAD_TYPE: u8 = 125; + +/// Register exactly the codecs that the SDK can encode or decode. +/// +/// `MediaEngine::register_default_codecs` also advertises legacy audio, VP9 +/// profile 1, AV1, and HEVC. Negotiating any of those would deliver a track +/// that the decoded [`RemoteTrack`](crate::rtc::RemoteTrack) APIs cannot consume. +fn register_supported_codecs(media_engine: &mut MediaEngine) -> Result<()> { + media_engine.register_codec( + RTCRtpCodecParameters { + capability: RTCRtpCodecCapability { + mime_type: MIME_TYPE_OPUS.to_owned(), + clock_rate: 48_000, + channels: 2, + sdp_fmtp_line: "minptime=10;useinbandfec=1".to_owned(), + rtcp_feedback: vec![], + }, + payload_type: OPUS_PAYLOAD_TYPE, + ..Default::default() + }, + RTPCodecType::Audio, + )?; + + let video_feedback = vec![ + RTCPFeedback { + typ: "goog-remb".to_owned(), + parameter: String::new(), + }, + RTCPFeedback { + typ: "ccm".to_owned(), + parameter: "fir".to_owned(), + }, + RTCPFeedback { + typ: "nack".to_owned(), + parameter: String::new(), + }, + RTCPFeedback { + typ: "nack".to_owned(), + parameter: "pli".to_owned(), + }, + ]; + for (mime_type, payload_type, fmtp) in [ + (MIME_TYPE_VP8, VP8_PAYLOAD_TYPE, ""), + (MIME_TYPE_VP9, VP9_PAYLOAD_TYPE, "profile-id=0"), + (MIME_TYPE_H264, H264_PAYLOAD_TYPE, H264_FMTP), + ] { + media_engine.register_codec( + RTCRtpCodecParameters { + capability: RTCRtpCodecCapability { + mime_type: mime_type.to_owned(), + clock_rate: 90_000, + channels: 0, + sdp_fmtp_line: fmtp.to_owned(), + rtcp_feedback: video_feedback.clone(), + }, + payload_type, + ..Default::default() + }, + RTPCodecType::Video, + )?; + } + + Ok(()) +} + +/// Build a webrtc-rs [`API`] with the supported codecs and default interceptors +/// (NACK, RTCP reports, receiver-side TWCC). +/// +/// Known gap versus Pion/videosdk: webrtc-rs ships no publisher-side congestion +/// controller (TWCC *sender* estimator / GCC) and no RTX/NACK retransmission +/// sender. Opus audio is unaffected — low bitrate, loss-tolerant, single layer — +/// but high-bitrate video publishing runs without bandwidth estimation or +/// retransmission. The default interceptor set below is wired as-is rather than +/// worked around. +fn build_api() -> Result { + let mut media_engine = MediaEngine::default(); + register_supported_codecs(&mut media_engine)?; + let registry = register_default_interceptors(Registry::new(), &mut media_engine)?; + // RFC 6464 audio levels (videosdk `media_engine.go`). Registered *after* the + // default interceptors so TWCC keeps its usual id and this takes the next + // free one. Audio + send-only scopes it to our publisher's audio m-lines: + // an answerer may not introduce an extension the offerer never offered, so + // without this the SFU can never observe our level and no participant of + // ours is ever marked speaking. + media_engine.register_header_extension( + RTCRtpHeaderExtensionCapability { + uri: AUDIO_LEVEL_URI.to_owned(), + }, + RTPCodecType::Audio, + Some(RTCRtpTransceiverDirection::Sendonly), + )?; + // RID simulcast needs MID + RTP stream identifiers on each outbound packet. + // webrtc-rs uses these registrations while binding `new_with_rid` tracks. + for uri in [ + SDES_MID_URI, + SDES_RTP_STREAM_ID_URI, + SDES_REPAIR_RTP_STREAM_ID_URI, + ] { + media_engine.register_header_extension( + RTCRtpHeaderExtensionCapability { + uri: uri.to_owned(), + }, + RTPCodecType::Video, + Some(RTCRtpTransceiverDirection::Sendonly), + )?; + } + Ok(APIBuilder::new() + .with_media_engine(media_engine) + .with_interceptor_registry(registry) + .build()) +} + +/// Map coordinator ICE servers to webrtc-rs [`RTCIceServer`]s. +pub fn to_rtc_ice_servers(servers: &[IceServer]) -> Vec { + servers + .iter() + .filter(|s| !s.urls.is_empty()) + .map(|s| RTCIceServer { + urls: s.urls.clone(), + username: s.username.clone(), + credential: s.password.clone(), + }) + .collect() +} + +/// Create a fresh PeerConnection wired with the join credentials' ICE servers. +pub async fn new_peer_connection(ice: &[IceServer]) -> Result> { + let api = build_api()?; + let config = RTCConfiguration { + ice_servers: to_rtc_ice_servers(ice), + ..Default::default() + }; + let pc = api.new_peer_connection(config).await?; + Ok(Arc::new(pc)) +} + +/// Wire the PeerConnection lifecycle events into `tracer`, mirroring JS +/// `traceRTCPeerConnection` tag names (`signalingstatechange`, +/// `icegatheringstatechange`, `iceconnectionstatechange`, `negotiationneeded`, +/// `datachannel`). The `onicecandidate` / `ontrack` / `connectionstatechange` +/// tags are emitted from the join module's existing handlers for those events +/// (each webrtc-rs `on_*` handler can be registered only once), so this covers +/// only the events that would otherwise have no handler. +pub fn trace_peer_events(pc: &Arc, tracer: Arc) { + let t = tracer.clone(); + pc.on_signaling_state_change(Box::new(move |state| { + let t = t.clone(); + Box::pin(async move { + t.trace("signalingstatechange", json!(state.to_string())); + }) + })); + + let t = tracer.clone(); + pc.on_ice_gathering_state_change(Box::new(move |state| { + let t = t.clone(); + Box::pin(async move { + t.trace("icegatheringstatechange", json!(state.to_string())); + }) + })); + + let t = tracer.clone(); + pc.on_ice_connection_state_change(Box::new(move |state| { + let t = t.clone(); + Box::pin(async move { + t.trace("iceconnectionstatechange", json!(state.to_string())); + }) + })); + + let t = tracer.clone(); + pc.on_negotiation_needed(Box::new(move || { + let t = t.clone(); + Box::pin(async move { + t.trace("negotiationneeded", serde_json::Value::Null); + }) + })); + + let t = tracer; + pc.on_data_channel(Box::new(move |channel| { + let t = t.clone(); + Box::pin(async move { + t.trace("datachannel", json!([channel.id(), channel.label()])); + }) + })); +} + +/// Generate a throwaway SDP offer with `audio` + `video` m-lines in `direction`, +/// used purely so the SFU can extract our codec capabilities on the join +/// request. The temporary PeerConnection is closed before returning. +/// +/// `sendonly` produces the publisher SDP; `recvonly` produces the subscriber SDP. +pub async fn generic_sdp(direction: RTCRtpTransceiverDirection) -> Result { + let api = build_api()?; + let pc = api.new_peer_connection(RTCConfiguration::default()).await?; + + let result = build_generic_offer(&pc, direction).await; + // Always tear the temp PC down, even if offer creation failed. + let _ = pc.close().await; + result +} + +async fn build_generic_offer( + pc: &RTCPeerConnection, + direction: RTCRtpTransceiverDirection, +) -> Result { + use webrtc::rtp_transceiver::RTCRtpTransceiverInit; + + pc.add_transceiver_from_kind( + RTPCodecType::Video, + Some(RTCRtpTransceiverInit { + direction, + send_encodings: vec![], + }), + ) + .await?; + pc.add_transceiver_from_kind( + RTPCodecType::Audio, + Some(RTCRtpTransceiverInit { + direction, + send_encodings: vec![], + }), + ) + .await?; + + let offer = pc.create_offer(None).await?; + Ok(offer.sdp) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn ice_servers_filter_empty_and_map_credential() { + let servers = vec![ + IceServer { + urls: vec!["stun:stun.l.google.com:19302".into()], + username: String::new(), + password: String::new(), + }, + IceServer { + urls: vec!["turn:turn.example.com:3478".into()], + username: "user".into(), + password: "pass".into(), + }, + IceServer::default(), // no urls -> filtered out + ]; + let mapped = to_rtc_ice_servers(&servers); + assert_eq!(mapped.len(), 2); + assert_eq!(mapped[1].username, "user"); + assert_eq!(mapped[1].credential, "pass"); + } + + #[tokio::test] + async fn generic_sdp_has_audio_and_video_mlines() { + let sdp = generic_sdp(RTCRtpTransceiverDirection::Recvonly) + .await + .expect("generic sdp"); + assert!(sdp.contains("m=audio"), "expected audio m-line"); + assert!(sdp.contains("m=video"), "expected video m-line"); + } + + #[tokio::test] + async fn generic_sdp_advertises_only_decodable_codecs() { + let sdp = generic_sdp(RTCRtpTransceiverDirection::Recvonly) + .await + .expect("generic sdp"); + let audio = media_section(&sdp, "audio"); + let video = media_section(&sdp, "video"); + + assert!(audio.contains("opus/48000/2"), "missing Opus:\n{audio}"); + for unsupported in ["PCMU/8000", "PCMA/8000", "G722/8000"] { + assert!( + !audio.contains(unsupported), + "advertised unsupported audio codec {unsupported}:\n{audio}" + ); + } + + for supported in ["VP8/90000", "VP9/90000", "H264/90000"] { + assert!( + video.contains(supported), + "missing supported video codec {supported}:\n{video}" + ); + } + for unsupported in ["AV1/90000", "H265/90000"] { + assert!( + !video.contains(unsupported), + "advertised unsupported video codec {unsupported}:\n{video}" + ); + } + assert!( + video.contains("profile-id=0"), + "missing VP9 profile 0:\n{video}" + ); + assert!( + !video.contains("profile-id=1"), + "advertised unsupported VP9 profile 1:\n{video}" + ); + } + + /// Split an SDP into its per-m-line sections, keyed by media kind. + fn media_section<'a>(sdp: &'a str, kind: &str) -> &'a str { + let start = sdp + .find(&format!("m={kind}")) + .unwrap_or_else(|| panic!("no m={kind} section in:\n{sdp}")); + let rest = &sdp[start..]; + match rest[1..].find("\r\nm=") { + Some(end) => &rest[..end + 1], + None => rest, + } + } + + #[tokio::test] + async fn sendonly_sdp_offers_audio_level_on_audio_only() { + let sdp = generic_sdp(RTCRtpTransceiverDirection::Sendonly) + .await + .expect("generic sdp"); + let audio = media_section(&sdp, "audio"); + let video = media_section(&sdp, "video"); + let extmap = audio + .lines() + .map(str::trim_end) + .find(|l| l.starts_with("a=extmap:") && l.ends_with(AUDIO_LEVEL_URI)); + assert!( + extmap.is_some(), + "publisher audio m-line must offer a=extmap: {AUDIO_LEVEL_URI}:\n{audio}" + ); + assert!( + !video.contains(AUDIO_LEVEL_URI), + "audio-level extmap must not appear on the video m-line:\n{video}" + ); + } + + #[tokio::test] + async fn recvonly_sdp_omits_audio_level() { + // The subscriber PC is recvonly and never publishes, so the send-only + // registration must keep the extension off it. + let sdp = generic_sdp(RTCRtpTransceiverDirection::Recvonly) + .await + .expect("generic sdp"); + assert!( + !sdp.contains(AUDIO_LEVEL_URI), + "recvonly SDP must not offer {AUDIO_LEVEL_URI}:\n{sdp}" + ); + } +} diff --git a/src/rtc/peer/mod.rs b/src/rtc/peer/mod.rs index 8e8ad6e..0c8307a 100644 --- a/src/rtc/peer/mod.rs +++ b/src/rtc/peer/mod.rs @@ -7,384 +7,11 @@ //! SFU inspects to learn our codec capabilities on the `JoinRequest` //! (JS `getGenericSdp`, stream-py `create_join_request`). +mod connection; mod ice; pub mod publisher; mod subscriber; -use std::sync::Arc; - -use serde_json::json; -use webrtc::api::interceptor_registry::register_default_interceptors; -use webrtc::api::media_engine::{ - MIME_TYPE_H264, MIME_TYPE_OPUS, MIME_TYPE_VP8, MIME_TYPE_VP9, MediaEngine, -}; -use webrtc::api::{API, APIBuilder}; -use webrtc::ice_transport::ice_server::RTCIceServer; -use webrtc::interceptor::registry::Registry; -use webrtc::peer_connection::RTCPeerConnection; -use webrtc::peer_connection::configuration::RTCConfiguration; -use webrtc::rtp_transceiver::RTCPFeedback; -use webrtc::rtp_transceiver::rtp_codec::{ - RTCRtpCodecCapability, RTCRtpCodecParameters, RTCRtpHeaderExtensionCapability, RTPCodecType, -}; -use webrtc::rtp_transceiver::rtp_transceiver_direction::RTCRtpTransceiverDirection; -use webrtc::sdp::extmap::{ - AUDIO_LEVEL_URI, SDES_MID_URI, SDES_REPAIR_RTP_STREAM_ID_URI, SDES_RTP_STREAM_ID_URI, -}; - -use super::coordinator::IceServer; -use super::error::Result; -use super::publish_options::H264_FMTP; -use super::tracer::Tracer; - +pub use connection::{generic_sdp, new_peer_connection, to_rtc_ice_servers, trace_peer_events}; pub(super) use ice::{PendingIce, register_ice_trickle}; pub(super) use subscriber::negotiate_subscriber; - -const OPUS_PAYLOAD_TYPE: u8 = 111; -const VP8_PAYLOAD_TYPE: u8 = 96; -const VP9_PAYLOAD_TYPE: u8 = 98; -const H264_PAYLOAD_TYPE: u8 = 125; - -/// Register exactly the codecs that the SDK can encode or decode. -/// -/// `MediaEngine::register_default_codecs` also advertises legacy audio, VP9 -/// profile 1, AV1, and HEVC. Negotiating any of those would deliver a track -/// that the decoded [`RemoteTrack`](super::RemoteTrack) APIs cannot consume. -fn register_supported_codecs(media_engine: &mut MediaEngine) -> Result<()> { - media_engine.register_codec( - RTCRtpCodecParameters { - capability: RTCRtpCodecCapability { - mime_type: MIME_TYPE_OPUS.to_owned(), - clock_rate: 48_000, - channels: 2, - sdp_fmtp_line: "minptime=10;useinbandfec=1".to_owned(), - rtcp_feedback: vec![], - }, - payload_type: OPUS_PAYLOAD_TYPE, - ..Default::default() - }, - RTPCodecType::Audio, - )?; - - let video_feedback = vec![ - RTCPFeedback { - typ: "goog-remb".to_owned(), - parameter: String::new(), - }, - RTCPFeedback { - typ: "ccm".to_owned(), - parameter: "fir".to_owned(), - }, - RTCPFeedback { - typ: "nack".to_owned(), - parameter: String::new(), - }, - RTCPFeedback { - typ: "nack".to_owned(), - parameter: "pli".to_owned(), - }, - ]; - for (mime_type, payload_type, fmtp) in [ - (MIME_TYPE_VP8, VP8_PAYLOAD_TYPE, ""), - (MIME_TYPE_VP9, VP9_PAYLOAD_TYPE, "profile-id=0"), - (MIME_TYPE_H264, H264_PAYLOAD_TYPE, H264_FMTP), - ] { - media_engine.register_codec( - RTCRtpCodecParameters { - capability: RTCRtpCodecCapability { - mime_type: mime_type.to_owned(), - clock_rate: 90_000, - channels: 0, - sdp_fmtp_line: fmtp.to_owned(), - rtcp_feedback: video_feedback.clone(), - }, - payload_type, - ..Default::default() - }, - RTPCodecType::Video, - )?; - } - - Ok(()) -} - -/// Build a webrtc-rs [`API`] with the supported codecs and default interceptors -/// (NACK, RTCP reports, receiver-side TWCC). -/// -/// Known gap versus Pion/videosdk: webrtc-rs ships no publisher-side congestion -/// controller (TWCC *sender* estimator / GCC) and no RTX/NACK retransmission -/// sender. Opus audio is unaffected — low bitrate, loss-tolerant, single layer — -/// but high-bitrate video publishing runs without bandwidth estimation or -/// retransmission. The default interceptor set below is wired as-is rather than -/// worked around. -fn build_api() -> Result { - let mut media_engine = MediaEngine::default(); - register_supported_codecs(&mut media_engine)?; - let registry = register_default_interceptors(Registry::new(), &mut media_engine)?; - // RFC 6464 audio levels (videosdk `media_engine.go`). Registered *after* the - // default interceptors so TWCC keeps its usual id and this takes the next - // free one. Audio + send-only scopes it to our publisher's audio m-lines: - // an answerer may not introduce an extension the offerer never offered, so - // without this the SFU can never observe our level and no participant of - // ours is ever marked speaking. - media_engine.register_header_extension( - RTCRtpHeaderExtensionCapability { - uri: AUDIO_LEVEL_URI.to_owned(), - }, - RTPCodecType::Audio, - Some(RTCRtpTransceiverDirection::Sendonly), - )?; - // RID simulcast needs MID + RTP stream identifiers on each outbound packet. - // webrtc-rs uses these registrations while binding `new_with_rid` tracks. - for uri in [ - SDES_MID_URI, - SDES_RTP_STREAM_ID_URI, - SDES_REPAIR_RTP_STREAM_ID_URI, - ] { - media_engine.register_header_extension( - RTCRtpHeaderExtensionCapability { - uri: uri.to_owned(), - }, - RTPCodecType::Video, - Some(RTCRtpTransceiverDirection::Sendonly), - )?; - } - Ok(APIBuilder::new() - .with_media_engine(media_engine) - .with_interceptor_registry(registry) - .build()) -} - -/// Map coordinator ICE servers to webrtc-rs [`RTCIceServer`]s. -pub fn to_rtc_ice_servers(servers: &[IceServer]) -> Vec { - servers - .iter() - .filter(|s| !s.urls.is_empty()) - .map(|s| RTCIceServer { - urls: s.urls.clone(), - username: s.username.clone(), - credential: s.password.clone(), - }) - .collect() -} - -/// Create a fresh PeerConnection wired with the join credentials' ICE servers. -pub async fn new_peer_connection(ice: &[IceServer]) -> Result> { - let api = build_api()?; - let config = RTCConfiguration { - ice_servers: to_rtc_ice_servers(ice), - ..Default::default() - }; - let pc = api.new_peer_connection(config).await?; - Ok(Arc::new(pc)) -} - -/// Wire the PeerConnection lifecycle events into `tracer`, mirroring JS -/// `traceRTCPeerConnection` tag names (`signalingstatechange`, -/// `icegatheringstatechange`, `iceconnectionstatechange`, `negotiationneeded`, -/// `datachannel`). The `onicecandidate` / `ontrack` / `connectionstatechange` -/// tags are emitted from the join module's existing handlers for those events -/// (each webrtc-rs `on_*` handler can be registered only once), so this covers -/// only the events that would otherwise have no handler. -pub fn trace_peer_events(pc: &Arc, tracer: Arc) { - let t = tracer.clone(); - pc.on_signaling_state_change(Box::new(move |state| { - let t = t.clone(); - Box::pin(async move { - t.trace("signalingstatechange", json!(state.to_string())); - }) - })); - - let t = tracer.clone(); - pc.on_ice_gathering_state_change(Box::new(move |state| { - let t = t.clone(); - Box::pin(async move { - t.trace("icegatheringstatechange", json!(state.to_string())); - }) - })); - - let t = tracer.clone(); - pc.on_ice_connection_state_change(Box::new(move |state| { - let t = t.clone(); - Box::pin(async move { - t.trace("iceconnectionstatechange", json!(state.to_string())); - }) - })); - - let t = tracer.clone(); - pc.on_negotiation_needed(Box::new(move || { - let t = t.clone(); - Box::pin(async move { - t.trace("negotiationneeded", serde_json::Value::Null); - }) - })); - - let t = tracer; - pc.on_data_channel(Box::new(move |channel| { - let t = t.clone(); - Box::pin(async move { - t.trace("datachannel", json!([channel.id(), channel.label()])); - }) - })); -} - -/// Generate a throwaway SDP offer with `audio` + `video` m-lines in `direction`, -/// used purely so the SFU can extract our codec capabilities on the join -/// request. The temporary PeerConnection is closed before returning. -/// -/// `sendonly` produces the publisher SDP; `recvonly` produces the subscriber SDP. -pub async fn generic_sdp(direction: RTCRtpTransceiverDirection) -> Result { - let api = build_api()?; - let pc = api.new_peer_connection(RTCConfiguration::default()).await?; - - let result = build_generic_offer(&pc, direction).await; - // Always tear the temp PC down, even if offer creation failed. - let _ = pc.close().await; - result -} - -async fn build_generic_offer( - pc: &RTCPeerConnection, - direction: RTCRtpTransceiverDirection, -) -> Result { - use webrtc::rtp_transceiver::RTCRtpTransceiverInit; - - pc.add_transceiver_from_kind( - RTPCodecType::Video, - Some(RTCRtpTransceiverInit { - direction, - send_encodings: vec![], - }), - ) - .await?; - pc.add_transceiver_from_kind( - RTPCodecType::Audio, - Some(RTCRtpTransceiverInit { - direction, - send_encodings: vec![], - }), - ) - .await?; - - let offer = pc.create_offer(None).await?; - Ok(offer.sdp) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn ice_servers_filter_empty_and_map_credential() { - let servers = vec![ - IceServer { - urls: vec!["stun:stun.l.google.com:19302".into()], - username: String::new(), - password: String::new(), - }, - IceServer { - urls: vec!["turn:turn.example.com:3478".into()], - username: "user".into(), - password: "pass".into(), - }, - IceServer::default(), // no urls -> filtered out - ]; - let mapped = to_rtc_ice_servers(&servers); - assert_eq!(mapped.len(), 2); - assert_eq!(mapped[1].username, "user"); - assert_eq!(mapped[1].credential, "pass"); - } - - #[tokio::test] - async fn generic_sdp_has_audio_and_video_mlines() { - let sdp = generic_sdp(RTCRtpTransceiverDirection::Recvonly) - .await - .expect("generic sdp"); - assert!(sdp.contains("m=audio"), "expected audio m-line"); - assert!(sdp.contains("m=video"), "expected video m-line"); - } - - #[tokio::test] - async fn generic_sdp_advertises_only_decodable_codecs() { - let sdp = generic_sdp(RTCRtpTransceiverDirection::Recvonly) - .await - .expect("generic sdp"); - let audio = media_section(&sdp, "audio"); - let video = media_section(&sdp, "video"); - - assert!(audio.contains("opus/48000/2"), "missing Opus:\n{audio}"); - for unsupported in ["PCMU/8000", "PCMA/8000", "G722/8000"] { - assert!( - !audio.contains(unsupported), - "advertised unsupported audio codec {unsupported}:\n{audio}" - ); - } - - for supported in ["VP8/90000", "VP9/90000", "H264/90000"] { - assert!( - video.contains(supported), - "missing supported video codec {supported}:\n{video}" - ); - } - for unsupported in ["AV1/90000", "H265/90000"] { - assert!( - !video.contains(unsupported), - "advertised unsupported video codec {unsupported}:\n{video}" - ); - } - assert!( - video.contains("profile-id=0"), - "missing VP9 profile 0:\n{video}" - ); - assert!( - !video.contains("profile-id=1"), - "advertised unsupported VP9 profile 1:\n{video}" - ); - } - - /// Split an SDP into its per-m-line sections, keyed by media kind. - fn media_section<'a>(sdp: &'a str, kind: &str) -> &'a str { - let start = sdp - .find(&format!("m={kind}")) - .unwrap_or_else(|| panic!("no m={kind} section in:\n{sdp}")); - let rest = &sdp[start..]; - match rest[1..].find("\r\nm=") { - Some(end) => &rest[..end + 1], - None => rest, - } - } - - #[tokio::test] - async fn sendonly_sdp_offers_audio_level_on_audio_only() { - let sdp = generic_sdp(RTCRtpTransceiverDirection::Sendonly) - .await - .expect("generic sdp"); - let audio = media_section(&sdp, "audio"); - let video = media_section(&sdp, "video"); - let extmap = audio - .lines() - .map(str::trim_end) - .find(|l| l.starts_with("a=extmap:") && l.ends_with(AUDIO_LEVEL_URI)); - assert!( - extmap.is_some(), - "publisher audio m-line must offer a=extmap: {AUDIO_LEVEL_URI}:\n{audio}" - ); - assert!( - !video.contains(AUDIO_LEVEL_URI), - "audio-level extmap must not appear on the video m-line:\n{video}" - ); - } - - #[tokio::test] - async fn recvonly_sdp_omits_audio_level() { - // The subscriber PC is recvonly and never publishes, so the send-only - // registration must keep the extension off it. - let sdp = generic_sdp(RTCRtpTransceiverDirection::Recvonly) - .await - .expect("generic sdp"); - assert!( - !sdp.contains(AUDIO_LEVEL_URI), - "recvonly SDP must not offer {AUDIO_LEVEL_URI}:\n{sdp}" - ); - } -}