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/README.md b/README.md index 5e67747..88eae92 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/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/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/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..7fc50c6 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::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::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/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/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/lifecycle.rs b/src/rtc/join/lifecycle.rs index f672bfa..8f5f568 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; @@ -694,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 ee9a285..319195d 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; @@ -51,27 +49,25 @@ 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; -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, }; -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; +use super::tracks::{LocalTrack, RemoteParticipant, RemoteTrack}; use serde_json::json; @@ -84,8 +80,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 +268,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 +276,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..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 super::super::publisher; +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/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/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 4a576d3..2ca636e 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,56 +14,43 @@ //! //! # Stability //! -//! The wire-layer modules — [`proto`], [`peer`], [`sfu_ws`], [`signal`], -//! [`publisher`], [`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. 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; -pub mod local_track; pub mod pcm; pub mod peer; pub mod proto; 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 sfu; pub mod stats; pub mod subscriptions; pub mod tracer; +mod tracks; pub mod video_frame; -mod vpx; -mod vpx_decode; 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, }; 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::{ @@ -73,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_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}; +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/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); - } -} diff --git a/src/rtc/peer.rs b/src/rtc/peer/connection.rs similarity index 95% rename from src/rtc/peer.rs rename to src/rtc/peer/connection.rs index 3899e0a..f4a8936 100644 --- a/src/rtc/peer.rs +++ b/src/rtc/peer/connection.rs @@ -1,11 +1,4 @@ //! webrtc-rs PeerConnection construction and the throwaway generic SDPs. -//! -//! The participant path uses two PeerConnections (JS / videosdk): the publisher -//! is the offerer (SetPublisher over Twirp) and the subscriber is the answerer -//! (answers the SFU's `subscriber_offer` over the WS). This module builds them -//! with the SDK-supported codec + interceptor set and produces the "generic" SDPs the -//! SFU inspects to learn our codec capabilities on the `JoinRequest` -//! (JS `getGenericSdp`, stream-py `create_join_request`). use std::sync::Arc; @@ -28,10 +21,10 @@ 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; +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; @@ -42,7 +35,7 @@ const H264_PAYLOAD_TYPE: u8 = 125; /// /// `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. +/// that the decoded [`RemoteTrack`](crate::rtc::RemoteTrack) APIs cannot consume. fn register_supported_codecs(media_engine: &mut MediaEngine) -> Result<()> { media_engine.register_codec( RTCRtpCodecParameters { diff --git a/src/rtc/peer/ice.rs b/src/rtc/peer/ice.rs new file mode 100644 index 0000000..1029f39 --- /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::sfu::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/mod.rs b/src/rtc/peer/mod.rs new file mode 100644 index 0000000..0c8307a --- /dev/null +++ b/src/rtc/peer/mod.rs @@ -0,0 +1,17 @@ +//! webrtc-rs PeerConnection construction and the throwaway generic SDPs. +//! +//! The participant path uses two PeerConnections (JS / videosdk): the publisher +//! is the offerer (SetPublisher over Twirp) and the subscriber is the answerer +//! (answers the SFU's `subscriber_offer` over the WS). This module builds them +//! with the SDK-supported codec + interceptor set and produces the "generic" SDPs the +//! 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; + +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; 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..d7df3d2 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::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`. /// @@ -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::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 { @@ -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..a58b5b3 --- /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::sfu::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(()) +} 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 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 a01ac7e..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::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}; +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 4a1f783..e17a28b 100644 --- a/src/rtc/remote_track.rs +++ b/src/rtc/tracks/remote.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::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; +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. //! 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()); - } -}