From a3494a004d8fffb0a64fb4d512d0ddfb26b2bf91 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sun, 20 Sep 2026 08:23:11 -0700 Subject: [PATCH 1/2] chore: claim gateway API quest Co-Authored-By: GPT-5 From 9e4d6c79318990d8acfbf1298a8befe6b78a423e Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sun, 20 Sep 2026 09:22:19 -0700 Subject: [PATCH 2/2] feat(gateway)!: align embedding APIs Co-Authored-By: GPT-5 --- Cargo.lock | 4 +- quest/m1/README.md | 1 - quest/m1/api-gateways.md | 55 ---------- quest/m1/api-review-gate.md | 3 +- quest/m2/gateway-embed.md | 1 - quest/m2/qos/stats/schema.md | 4 +- rs/moq-cli/src/hls.rs | 17 ++++ rs/moq-cli/src/rtc.rs | 11 +- rs/moq-cli/src/srt.rs | 21 ++-- rs/moq-hls/Cargo.toml | 3 - rs/moq-hls/src/error.rs | 18 ---- rs/moq-hls/src/export/mod.rs | 12 ++- rs/moq-hls/src/export/playlist.rs | 14 +-- rs/moq-hls/src/export/rendition.rs | 25 ++--- rs/moq-hls/src/export/segments.rs | 18 ++-- rs/moq-hls/src/import.rs | 42 ++------ rs/moq-hls/src/lib.rs | 10 -- rs/moq-hls/src/server/routes.rs | 2 +- rs/moq-relay/src/stats.rs | 6 +- rs/moq-room/src/claims.rs | 10 +- rs/moq-rtc/Cargo.toml | 1 - rs/moq-rtc/src/client/whep.rs | 14 +-- rs/moq-rtc/src/client/whip.rs | 18 +--- rs/moq-rtc/src/codec/av1.rs | 3 +- rs/moq-rtc/src/codec/h264.rs | 3 +- rs/moq-rtc/src/codec/h265.rs | 3 +- rs/moq-rtc/src/codec/mod.rs | 36 ++++--- rs/moq-rtc/src/codec/opus.rs | 3 +- rs/moq-rtc/src/codec/vp8.rs | 3 +- rs/moq-rtc/src/codec/vp9.rs | 3 +- rs/moq-rtc/src/egress.rs | 6 +- rs/moq-rtc/src/error.rs | 36 ++++++- rs/moq-rtc/src/lib.rs | 14 +-- rs/moq-rtc/src/server/mod.rs | 54 ++-------- rs/moq-rtc/src/server/whep.rs | 46 +++++---- rs/moq-rtc/src/server/whip.rs | 41 +++++--- rs/moq-rtmp/README.md | 2 +- rs/moq-rtmp/src/error.rs | 9 +- rs/moq-rtmp/src/listen.rs | 41 +++++--- rs/moq-rtmp/src/server.rs | 4 +- rs/moq-srt/Cargo.toml | 2 +- rs/moq-srt/README.md | 6 +- rs/moq-srt/src/dial.rs | 156 ++++++++++++++++++----------- rs/moq-srt/src/error.rs | 17 ++-- rs/moq-srt/src/lib.rs | 9 +- rs/moq-srt/src/listen.rs | 24 ++--- rs/moq-srt/src/server.rs | 56 +++++++---- rs/moq-srt/src/ts.rs | 4 +- rs/moq-stats/src/aggregate.rs | 6 +- rs/moq-stats/src/consume.rs | 36 +++---- rs/moq-stats/src/lib.rs | 8 +- rs/moq-stats/src/produce.rs | 18 ++-- 52 files changed, 456 insertions(+), 503 deletions(-) delete mode 100644 quest/m1/api-gateways.md diff --git a/Cargo.lock b/Cargo.lock index 5aca0d7c55..f94aec16fb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4967,7 +4967,6 @@ dependencies = [ name = "moq-hls" version = "0.4.16" dependencies = [ - "anyhow", "axum", "bytes", "hang", @@ -5172,7 +5171,6 @@ dependencies = [ name = "moq-rtc" version = "0.2.11" dependencies = [ - "anyhow", "aws-lc-rs", "axum", "bytes", @@ -5230,13 +5228,13 @@ dependencies = [ name = "moq-srt" version = "0.2.11" dependencies = [ - "anyhow", "bytes", "futures", "hang", "moq-mux", "moq-net", "moq-tokio", + "srt-protocol", "srt-tokio", "thiserror 2.0.20", "tokio", diff --git a/quest/m1/README.md b/quest/m1/README.md index 0b183926c8..9b232fb087 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -31,7 +31,6 @@ the transport line in m2 assumes a single stack. - [Bindings announce match](/quest/m1/api-origin-scopes.md) - every binding takes a pattern scope and reports the announce match with its captures - [PathPrefixes](/quest/m1/api-path-prefixes.md) - the unused moq_net::PathPrefixes type is deleted before the release - [Rendition ownership](/quest/m1/api-mux-rendition.md) - one handle publishes a media track and reports its estimate, instead of five -- [Gateway types](/quest/m1/api-gateways.md) - no `anyhow` in a gateway `Error`, `PathOwned` prefixes, `Duration` segments, `moq_rtc::Server::new(config)`, an SRT reject with a reason - [Cluster -01](/quest/m1/cluster-01/README.md) - rs/moq-net and js/net speak the revised cluster extension (HOP_ID, REQUEST_UPDATE repricing) and -01 is published - [API review gate](/quest/m1/api-review-gate.md) - each `api-*` quest above is landed or deferred by the maintainer before the merge PR opens - [Merge dev](/quest/m1/merge-dev.md) - dev lands on main with a closing keyword for every issue it fixed diff --git a/quest/m1/api-gateways.md b/quest/m1/api-gateways.md deleted file mode 100644 index f9d470d301..0000000000 --- a/quest/m1/api-gateways.md +++ /dev/null @@ -1,55 +0,0 @@ -# [M] The gateway crates share moq-net's types and errors - -## Goal - -`moq-rtmp`, `moq-srt`, `moq-rtc`, `moq-hls`, `moq-stats`, and `moq-room` -release with typed errors an embedder can map, paths and durations in -moq-net's types, and constructors that take only what they use. Every one -of these crates is already breaking on dev, so the ride-along costs nothing -now and a full bump later. - -## Plan - -- `moq_rtc::Server::new(config)`; the origin pair moves to - `publish_router(origin)`/`subscribe_router(origin)`, the only place it is - read. moq.pro spawns an unused origin driver to satisfy the constructor. -- Delete `Error::Other(anyhow)` from all four gateway crates. moq-rtc wraps - `Unauthorized` and "not announced" as `Other`, so moq.pro answers WHIP - refusals with 500 and WHEP internal failures with 404; route them through - `Error::Moq` and add the few typed variants missing (`CatalogTimeout`, - `NoRenditions`, an RTMP `Session(String)`). moq-hls keeps `reqwest` and - `url` errors opaque and drops the two `pub use`s, as #3243 did elsewhere. -- `listen::Config.prefix` is a `PathOwned` on moq-rtmp and moq-srt (a - `String` that needs a trailing slash today); moq-stats already does this. -- `moq_hls::Segment.duration: Duration` (an `f64` of seconds built from a - `Duration`); `Kind` parses through a real error or `Kind::parse -> - Option`; `import::Config.playlist` is a `Url` (or a `Playlist` enum) - instead of a string parsed at run time. -- moq-stats: `produce::Config`, `consume::Config`, `consume::{Traffic, - Sessions}` replace the root `ProducerConfig`/`ConsumerConfig`/ - `TrafficConsumer`/`SessionsConsumer`, matching `aggregate::Config`. -- `moq_room::claims(room: impl AsPath, identity: impl AsPath)`; the room is - documented as a path prefix and taken as a `String`. -- `moq_srt::{Publish, Subscribe}::reject(self, Reject)` with - `Reject::{Unauthorized, Forbidden, Unavailable, BadRequest}` mapped to the - SRT extended codes (1401, 1403, 1503, 1400); today `reject()` sends a - fixed `Forbidden`. A `moq_srt` enum rather than a re-export of - `srt_tokio::ServerRejectReason`, so a backend swap is not an API - migration. Prove it over the wire for a rejected publish and a rejected - subscribe with a non-default reason, beside the existing clean-close - coverage, and update `rs/moq-cli/src/srt.rs` and the `rs/moq-srt/README.md` - embedder example. -- `moq_rtmp::Play::with_max_age(Duration)` and - `Publish::with_max_age(Option)` take the same shape. -- Decide the moq-srt dial shape: free functions over a `Config` today where - moq-rtmp has a `Client` builder with the bring-your-own-transport seam - moq.pro uses. Recommended: the `Client`. - -Public API: breaking on every crate named, so on dev. Wire: an SRT -client sees the extended reject code that matches the verdict; run -`just test smoke-full`. Consumers: moq-cli, moq.pro's edge (all four -gateways in-process). - -## Related - -- [Gateway embedding](/quest/m2/gateway-embed.md) - the additive entry points an in-process embedder still lacks diff --git a/quest/m1/api-review-gate.md b/quest/m1/api-review-gate.md index 2b0e904e63..00a995decd 100644 --- a/quest/m1/api-review-gate.md +++ b/quest/m1/api-review-gate.md @@ -17,8 +17,7 @@ file is deleted on completion) or deletes the quest with a note in quest is deleted too. No code. The list: [Announce event](/quest/m1/api-net-announce.md), -[Rendition ownership](/quest/m1/api-mux-rendition.md), -[Gateway types](/quest/m1/api-gateways.md). +[Rendition ownership](/quest/m1/api-mux-rendition.md). ## Related diff --git a/quest/m2/gateway-embed.md b/quest/m2/gateway-embed.md index feb8715c86..116a4ed2e2 100644 --- a/quest/m2/gateway-embed.md +++ b/quest/m2/gateway-embed.md @@ -33,7 +33,6 @@ Public API: additive. Wire: none. ## Required -- [Gateway types](/quest/m1/api-gateways.md) - the constructors this extends - [Merge dev](/quest/m1/merge-dev.md) - starts on main ## Related diff --git a/quest/m2/qos/stats/schema.md b/quest/m2/qos/stats/schema.md index 40c774c668..9581a2c937 100644 --- a/quest/m2/qos/stats/schema.md +++ b/quest/m2/qos/stats/schema.md @@ -20,10 +20,10 @@ what a subscriber received and played, per audio and video. The relay is `[/]sessions.json[.z]`, a `BTreeMap` keyed by auth root, read through `Consumer::sessions` whatever `E` is; a client publishes it only when it holds sessions worth counting. -- An exact-path mode. `ProducerConfig` treats its path as a prefix and +- An exact-path mode. `produce::Config` treats its path as a prefix and advertises `/node[/]`, so a client asking for `room/alice.stats` would publish `room/alice.stats/node`, which no longer - ends in `.stats`. `ProducerConfig::at(path)` publishes the broadcast at + ends in `.stats`. `produce::Config::at(path)` publishes the broadcast at exactly that path with no category segment, refusing a path that does not end in `.stats`; the relay keeps the prefix layout. Test the advertised path for both modes. diff --git a/rs/moq-cli/src/hls.rs b/rs/moq-cli/src/hls.rs index d8f5b37145..f09a193cde 100644 --- a/rs/moq-cli/src/hls.rs +++ b/rs/moq-cli/src/hls.rs @@ -2,10 +2,12 @@ //! DASH over HTTP from MoQ broadcasts (export), fetching media groups on demand. use std::net::SocketAddr; +use std::path::PathBuf; use anyhow::Context; use axum::http::Method; use hang::moq_net; +use url::Url; use crate::moq::{ImportTarget, notify_ready}; @@ -60,6 +62,7 @@ pub async fn import(target: ImportTarget, playlist: String) -> anyhow::Result<() .announce(Default::default()) .context("failed to announce broadcast")?; + let playlist = playlist_url(&playlist)?; let mut importer = moq_hls::import::Import::new(producer, catalog, moq_hls::import::Config::new(playlist))?; tracing::info!(%name, "importing HLS"); @@ -69,6 +72,20 @@ pub async fn import(target: ImportTarget, playlist: String) -> anyhow::Result<() Ok(importer.run().await?) } +fn playlist_url(playlist: &str) -> anyhow::Result { + if playlist.starts_with("http://") || playlist.starts_with("https://") { + return Url::parse(playlist).context("invalid HLS playlist URL"); + } + + let path = PathBuf::from(playlist); + let absolute = if path.is_absolute() { + path + } else { + std::env::current_dir()?.join(path) + }; + Url::from_file_path(&absolute).map_err(|_| anyhow::anyhow!("invalid HLS playlist path: {}", absolute.display())) +} + /// Serve HLS and DASH over HTTP for the single broadcast `name` (reached at /// `//master.m3u8` and `//manifest.mpd`); other broadcasts in the /// Origin are not served. diff --git a/rs/moq-cli/src/rtc.rs b/rs/moq-cli/src/rtc.rs index 99a31db27a..8449a01890 100644 --- a/rs/moq-cli/src/rtc.rs +++ b/rs/moq-cli/src/rtc.rs @@ -61,8 +61,8 @@ pub async fn listen_import(target: ImportTarget, listen: Listen) -> anyhow::Resu let mut config = server_config(&listen); config.max_age = target.max_age; config.bandwidth = target.bandwidth; - let server = moq_rtc::Server::new(config, publisher, target.origin.consume()); - serve(server.publish_router(), "WHIP", listen).await + let server = moq_rtc::Server::new(config); + serve(server.publish_router(publisher), "WHIP", listen).await } /// WHEP server: serve WebRTC plays of `name` from the Origin (export). @@ -73,11 +73,8 @@ pub async fn listen_export(origin: moq_net::origin::Consumer, name: String, list let subscriber = origin .scope("", &scope) .with_context(|| format!("failed to scope origin to broadcast `{name}`"))?; - // A WHEP server only reads; it still needs a publisher handle for the shared - // glue, so hand it an unused, empty Origin producer. - let publisher = moq_tokio::origin::spawn(); - let server = moq_rtc::Server::new(server_config(&listen), publisher, subscriber); - serve(server.subscribe_router(), "WHEP", listen).await + let server = moq_rtc::Server::new(server_config(&listen)); + serve(server.subscribe_router(subscriber), "WHEP", listen).await } /// Restrict a producer to the single broadcast `name` so a WHIP peer can only publish it. diff --git a/rs/moq-cli/src/srt.rs b/rs/moq-cli/src/srt.rs index 01756fa37b..7afbe3ff7f 100644 --- a/rs/moq-cli/src/srt.rs +++ b/rs/moq-cli/src/srt.rs @@ -6,7 +6,7 @@ use std::time::Duration; use anyhow::Context; use hang::moq_net; -use moq_srt::{Request, Server}; +use moq_srt::{Reject, Request, Server}; use moq_tokio::RedactedUrl; use url::Url; @@ -62,7 +62,7 @@ pub async fn listen_import(target: ImportTarget, addr: SocketAddr, latency: Dura } Request::Subscribe(subscribe) => { tokio::spawn(async move { - let _ = subscribe.reject().await; + let _ = subscribe.reject(Reject::Forbidden).await; }); } _ => {} @@ -96,7 +96,7 @@ pub async fn listen_export( } Request::Publish(publish) => { tokio::spawn(async move { - let _ = publish.reject().await; + let _ = publish.reject(Reject::Forbidden).await; }); } _ => {} @@ -113,11 +113,11 @@ pub async fn connect_import(target: ImportTarget, url: Url, latency: Duration) - tracing::info!(url = %RedactedUrl::new(&url), %name, "SRT client pulling"); notify_ready(); - let mut config = moq_srt::dial::Config::new(addr, resource); - config.latency = latency; - config.max_age = target.max_age; - config.bandwidth = target.bandwidth; - Ok(moq_srt::dial::pull(&config, &target.origin, name).await?) + let client = moq_srt::Client::new(addr, resource) + .with_latency(latency) + .with_max_age(target.max_age) + .with_bandwidth(target.bandwidth); + Ok(client.pull(&target.origin, name).await?) } /// Push a broadcast from the Origin to a remote SRT server (export). @@ -131,9 +131,8 @@ pub async fn connect_export( tracing::info!(url = %RedactedUrl::new(&url), %name, "SRT client pushing"); notify_ready(); - let mut config = moq_srt::dial::Config::new(addr, resource); - config.latency = latency; - Ok(moq_srt::dial::publish(&config, &origin, &name).await?) + let client = moq_srt::Client::new(addr, resource).with_latency(latency); + Ok(client.publish(&origin, &name).await?) } /// Parse `srt://host:port?streamid=` into a resolved address and resource. diff --git a/rs/moq-hls/Cargo.toml b/rs/moq-hls/Cargo.toml index b33d53d9a7..ba14521b31 100644 --- a/rs/moq-hls/Cargo.toml +++ b/rs/moq-hls/Cargo.toml @@ -17,9 +17,6 @@ default = ["server"] server = ["dep:axum"] [dependencies] -# Always required (import + export library). -anyhow = { workspace = true, features = ["backtrace"] } - # Only needed by the HTTP export server (gated by `server`). axum = { workspace = true, optional = true } bytes = { workspace = true } diff --git a/rs/moq-hls/src/error.rs b/rs/moq-hls/src/error.rs index 50055295e3..d854fe4011 100644 --- a/rs/moq-hls/src/error.rs +++ b/rs/moq-hls/src/error.rs @@ -39,14 +39,6 @@ pub enum Error { #[error("mux: {0}")] Mux(#[from] moq_mux::Error), - /// The playlist argument looked like an HTTP(S) URL but failed to parse. - #[error("invalid playlist URL")] - InvalidPlaylistUrl, - - /// The playlist argument was a local path that could not be made into a `file://` URL. - #[error("invalid file path")] - InvalidFilePath, - /// A `file://` URL could not be turned back into a filesystem path. #[error("invalid file URL")] InvalidFileUrl, @@ -127,10 +119,6 @@ pub enum Error { /// I/O error while reading a local playlist or segment. #[error("io: {0}")] Io(std::sync::Arc), - - /// Catch-all for gateway logic that reports via `anyhow`. - #[error("{0}")] - Other(std::sync::Arc), } impl Error { @@ -159,12 +147,6 @@ impl From for Error { } } -impl From for Error { - fn from(err: anyhow::Error) -> Self { - Error::Other(std::sync::Arc::new(err)) - } -} - /// Convenience alias for results from the HLS gateway. pub type Result = std::result::Result; diff --git a/rs/moq-hls/src/export/mod.rs b/rs/moq-hls/src/export/mod.rs index 2996bc3d4b..565108a2ce 100644 --- a/rs/moq-hls/src/export/mod.rs +++ b/rs/moq-hls/src/export/mod.rs @@ -740,7 +740,7 @@ mod tests { let playlist = rendition.playlist(); assert_eq!(playlist.segments.len(), 2, "the live-edge group is not listed"); assert_eq!(playlist.segments[0].segment, 0); - assert_eq!(playlist.segments[0].duration, 2.0); + assert_eq!(playlist.segments[0].duration, Duration::from_secs(2)); assert_eq!(playlist.segments[1].segment, 1); assert_eq!( playlist.target_duration, 2, @@ -1107,7 +1107,7 @@ mod tests { let _ = tokio::time::timeout(Duration::from_secs(5), rendition.playable()).await; let playlist = rendition.playlist(); - assert_eq!(playlist.segments[0].duration, 3.0); + assert_eq!(playlist.segments[0].duration, Duration::from_secs(3)); assert_eq!( playlist.target_duration, 3, "no bound was declared, so the target duration must still cover the 3s segment" @@ -1179,7 +1179,11 @@ mod tests { let audio_segments: Vec = audio_playlist.segments.iter().map(|s| s.segment).collect(); assert_eq!(video_segments, vec![0, 1]); assert_eq!(audio_segments, vec![0, 1], "audio lists the same aligned segments"); - assert_eq!(audio_playlist.segments[0].duration, 2.0, "cut at the video boundary"); + assert_eq!( + audio_playlist.segments[0].duration, + Duration::from_secs(2), + "cut at the video boundary" + ); // The same URI names the same span of content time on either rendition. let rendered = audio_rendition.media_playlist(None).expect("playable"); @@ -1926,7 +1930,7 @@ mod tests { .expect("a segment, not end"); assert_eq!(first.segment, 0); assert_eq!(&first.media[4..8], b"moof", "the segment carries its transmuxed media"); - assert_eq!(first.duration, 2.0); + assert_eq!(first.duration, Duration::from_secs(2)); assert!(!first.discontinuity, "a clean start is not a discontinuity"); let second = segments.next().await.unwrap().expect("second segment"); diff --git a/rs/moq-hls/src/export/playlist.rs b/rs/moq-hls/src/export/playlist.rs index 362cfafc43..ec88dab909 100644 --- a/rs/moq-hls/src/export/playlist.rs +++ b/rs/moq-hls/src/export/playlist.rs @@ -7,7 +7,7 @@ //! directory. use std::fmt::Write; -use std::time::SystemTime; +use std::time::{Duration, SystemTime}; /// fMP4 segments via `EXT-X-MAP` require protocol version 6. const VERSION: u32 = 6; @@ -35,8 +35,8 @@ pub(crate) struct Snapshot { pub(crate) struct Segment { /// The aligned segment number; the URI is `seg/{segment}.m4s`. pub segment: u64, - /// `EXTINF` duration in seconds. - pub duration: f64, + /// `EXTINF` duration. + pub duration: Duration, /// The rendition has no content for this span (`EXT-X-GAP`): the segment keeps its slot in /// the aligned numbering, but a player should not request it. pub gap: bool, @@ -82,7 +82,7 @@ pub(crate) fn render_media(snapshot: &Snapshot, query: Option<&str>) -> String { if segment.gap { let _ = writeln!(out, "#EXT-X-GAP"); } - let _ = writeln!(out, "#EXTINF:{:.5},", segment.duration); + let _ = writeln!(out, "#EXTINF:{:.5},", segment.duration.as_secs_f64()); let _ = writeln!(out, "seg/{}.m4s{suffix}", segment.segment); } @@ -107,13 +107,13 @@ mod tests { segments: vec![ Segment { segment: 10, - duration: 2.0, + duration: Duration::from_secs(2), gap: false, discontinuity: false, }, Segment { segment: 11, - duration: 1.96, + duration: Duration::from_millis(1960), gap: false, discontinuity: false, }, @@ -146,7 +146,7 @@ mod tests { media_sequence: 0, segments: vec![Segment { segment: 0, - duration: 4.0, + duration: Duration::from_secs(4), gap: false, discontinuity: false, }], diff --git a/rs/moq-hls/src/export/rendition.rs b/rs/moq-hls/src/export/rendition.rs index 506376cb2c..b99830eac1 100644 --- a/rs/moq-hls/src/export/rendition.rs +++ b/rs/moq-hls/src/export/rendition.rs @@ -44,6 +44,15 @@ pub enum Kind { } impl Kind { + /// Parse a rendition URL path component. + pub fn parse(value: &str) -> Option { + match value { + "video" => Some(Self::Video), + "audio" => Some(Self::Audio), + _ => None, + } + } + /// The URL path component for this kind (`"video"` / `"audio"`). pub fn as_str(self) -> &'static str { match self { @@ -53,18 +62,6 @@ impl Kind { } } -impl std::str::FromStr for Kind { - type Err = (); - - fn from_str(s: &str) -> std::result::Result { - match s { - "video" => Ok(Kind::Video), - "audio" => Ok(Kind::Audio), - _ => Err(()), - } - } -} - /// The rendition's catalog config, kept whole so a [`Muxer`] can be built per request. enum Config { Video(VideoConfig), @@ -318,7 +315,7 @@ impl Rendition { index, segment: entry.segment, ranges: entry.tracks.get(&self.name).cloned().unwrap_or_default(), - duration: entry.duration.as_secs_f64(), + duration: entry.duration, pts: entry.pts, end: Duration::from(entry.pts) + entry.duration, }; @@ -403,7 +400,7 @@ impl Rendition { let observed = window .segments .iter() - .map(|s| s.duration.ceil().max(0.0) as u64) + .map(|s| s.duration.as_secs_f64().ceil() as u64) .max() .unwrap_or(0); let target_duration = declared.max(observed).max(1); diff --git a/rs/moq-hls/src/export/segments.rs b/rs/moq-hls/src/export/segments.rs index d9617d892c..2a5a14488a 100644 --- a/rs/moq-hls/src/export/segments.rs +++ b/rs/moq-hls/src/export/segments.rs @@ -56,8 +56,8 @@ pub(crate) struct Row { /// This rendition's group ranges within the segment. Empty means the rendition has no /// content for the span (`EXT-X-GAP`). pub ranges: Vec, - /// Presentation duration in seconds. - pub duration: f64, + /// Presentation duration. + pub duration: Duration, /// The segment's starting presentation timestamp. pub pts: moq_net::Timestamp, /// The segment's ending presentation timestamp (`pts + duration`), for window eviction @@ -289,8 +289,8 @@ pub struct Segment { pub segment: u64, /// The transmuxed CMAF fragment (`moof`+`mdat`), fetched on demand by [`Consumer::next`]. pub media: Bytes, - /// Presentation duration in seconds. - pub duration: f64, + /// Presentation duration. + pub duration: Duration, /// Wall-clock start time, when the timeline advertises an anchor. pub program_date_time: Option, /// The media timeline is broken before this segment: one or more segments were skipped @@ -397,7 +397,7 @@ mod tests { index: segment, segment, ranges: vec![Range::new(group, group)], - duration: duration_ms as f64 / 1000.0, + duration: Duration::from_millis(duration_ms), pts, end: Duration::from(pts) + Duration::from_millis(duration_ms), } @@ -414,7 +414,7 @@ mod tests { assert!(!window.ended); assert_eq!(window.segments.len(), 2, "rows are complete segments; all are listed"); assert_eq!(window.segments[0].segment, 0); - assert_eq!(window.segments[0].duration, 2.0); + assert_eq!(window.segments[0].duration, Duration::from_secs(2)); live.end(); assert!(live.window().ended); @@ -432,8 +432,8 @@ mod tests { // Segments still cover >= 4s after eviction, and the sequence is the first listed // segment's aligned number. assert!(snapshot.sequence > 0); - let span: f64 = snapshot.segments.iter().map(|s| s.duration).sum(); - assert!(span >= 4.0); + let span: Duration = snapshot.segments.iter().map(|s| s.duration).sum(); + assert!(span >= Duration::from_secs(4)); assert_eq!(snapshot.segments.first().unwrap().segment, snapshot.sequence); } @@ -491,7 +491,7 @@ mod tests { index: 1, segment: 1, ranges: Vec::new(), - duration: 1.0, + duration: Duration::from_secs(1), pts: moq_net::Timestamp::from_millis(1_000).unwrap(), end: Duration::from_millis(2_000), }, diff --git a/rs/moq-hls/src/import.rs b/rs/moq-hls/src/import.rs index c97379518a..165872e63d 100644 --- a/rs/moq-hls/src/import.rs +++ b/rs/moq-hls/src/import.rs @@ -8,7 +8,6 @@ use std::collections::HashMap; use std::collections::hash_map::Entry; use std::io::SeekFrom; -use std::path::PathBuf; use std::time::Duration; use bytes::Bytes; @@ -57,40 +56,22 @@ const MAX_RENDITION_FAILURES: usize = 3; #[derive(Clone)] #[non_exhaustive] pub struct Config { - /// The master or media playlist URL or file path to import. - pub playlist: String, + /// The master or media playlist URL to import. Local files use a `file://` URL. + pub playlist: Url, /// HTTP client used to fetch the playlist and segments, for example one carrying /// credentials for an authenticated origin. Defaults to a plain client with a 30 /// second per-request timeout. /// - /// This is a [`reqwest::Client`], re-exported as [`crate::reqwest`]; a major version - /// bump of that dependency is a breaking change for this field. + /// This is a [`reqwest::Client`]. pub client: Option, } impl Config { /// Create an import configuration for `playlist` using the default HTTP client. - pub fn new(playlist: String) -> Self { + pub fn new(playlist: Url) -> Self { Self { playlist, client: None } } - - /// Parse the playlist string into a URL. - /// If it starts with http:// or https://, parse as URL. - /// Otherwise, treat as a file path and convert to file:// URL. - fn parse_playlist(&self) -> Result { - if self.playlist.starts_with("http://") || self.playlist.starts_with("https://") { - Url::parse(&self.playlist).map_err(|_| Error::InvalidPlaylistUrl) - } else { - let path = PathBuf::from(&self.playlist); - let absolute = if path.is_absolute() { - path - } else { - std::env::current_dir()?.join(path) - }; - Url::from_file_path(&absolute).map_err(|_| Error::InvalidFilePath) - } - } } /// Result of a single import step. @@ -611,7 +592,7 @@ pub struct Import { impl Import { /// Create a new HLS import that will write into the given broadcast. pub fn new(broadcast: moq_net::broadcast::Producer, catalog: CatalogProducer, cfg: Config) -> Result { - let base_url = cfg.parse_playlist()?; + let base_url = cfg.playlist; Ok(Self { sink: Sink { broadcast, catalog }, fetcher: Fetcher::new(cfg.client)?, @@ -1033,7 +1014,7 @@ fn moq_sequence(discontinuity_sequence: u64, media_sequence: u64) -> Result #[cfg(test)] mod tests { use super::*; - use std::path::Path; + use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicUsize, Ordering}; use tokio::io::AsyncWriteExt as _; use tokio::net::TcpListener; @@ -1054,7 +1035,7 @@ mod tests { let mut broadcast = moq_net::broadcast::Info::new().produce(); let catalog = CatalogProducer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap(); - let cfg = Config::new(playlist_path.to_string_lossy().into_owned()); + let cfg = Config::new(Url::from_file_path(&playlist_path).unwrap()); let import = Import::new(broadcast, catalog.clone(), cfg).unwrap(); (import, catalog) } @@ -1152,7 +1133,7 @@ mod tests { #[test] fn hls_config_new_sets_fields() { - let url = "https://example.com/stream.m3u8".to_string(); + let url = Url::parse("https://example.com/stream.m3u8").unwrap(); let cfg = Config::new(url.clone()); assert_eq!(cfg.playlist, url); } @@ -1194,7 +1175,7 @@ mod tests { fn hls_import_starts_without_tracks() { let mut broadcast = moq_net::broadcast::Info::new().produce(); let catalog = CatalogProducer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap(); - let url = "https://example.com/master.m3u8".to_string(); + let url = Url::parse("https://example.com/master.m3u8").unwrap(); let cfg = Config::new(url); let hls = Import::new(broadcast, catalog, cfg).unwrap(); @@ -1256,7 +1237,7 @@ mod tests { let mut broadcast = moq_net::broadcast::Info::new().produce(); let catalog = CatalogProducer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap(); - let cfg = Config::new(path.to_string_lossy().into_owned()); + let cfg = Config::new(Url::from_file_path(&path).unwrap()); let mut import = Import::new(broadcast, catalog, cfg).unwrap(); assert!(matches!(import.ensure_tracks().await, Err(Error::NoVariants))); @@ -1543,8 +1524,7 @@ mod tests { let mut broadcast = moq_net::broadcast::Info::new().produce(); let catalog = CatalogProducer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap(); - // `Config` takes a filesystem path for non-http inputs. - let cfg = Config::new(path.to_str().unwrap().to_string()); + let cfg = Config::new(Url::from_file_path(&path).unwrap()); let mut hls = Import::new(broadcast, catalog, cfg).unwrap(); hls.ensure_tracks().await.unwrap(); hls diff --git a/rs/moq-hls/src/lib.rs b/rs/moq-hls/src/lib.rs index 29e544a04a..ae42fa0894 100644 --- a/rs/moq-hls/src/lib.rs +++ b/rs/moq-hls/src/lib.rs @@ -37,13 +37,3 @@ pub use server::Server; /// breaking change for this crate. #[cfg(feature = "server")] pub use axum; - -/// Re-export of the HTTP client used by [`import`], so consumers can name the -/// [`reqwest::Error`] carried by [`Error::Reqwest`] without adding their own reqwest -/// dependency. A major reqwest bump is therefore a breaking change for this crate. -pub use reqwest; - -/// Re-export of the URL parser, so consumers can name the [`url::Url`] and -/// [`url::ParseError`] carried by [`Error`] without adding their own url dependency. -/// A major url bump is therefore a breaking change for this crate. -pub use url; diff --git a/rs/moq-hls/src/server/routes.rs b/rs/moq-hls/src/server/routes.rs index 275a63646c..72bd73ea7c 100644 --- a/rs/moq-hls/src/server/routes.rs +++ b/rs/moq-hls/src/server/routes.rs @@ -214,7 +214,7 @@ async fn segment(server: &Server, broadcast: &str, kind: &str, rendition: &str, /// Resolve a rendition, waiting for the catalog to populate. async fn rendition_for(server: &Server, broadcast: &str, kind: &str, rendition: &str) -> Option> { - let kind = kind.parse::().ok()?; + let kind = Kind::parse(kind)?; let broadcaster = server.broadcaster(broadcast).await?; let _ = tokio::time::timeout(READY_TIMEOUT, broadcaster.ready()).await; broadcaster.rendition(kind, rendition) diff --git a/rs/moq-relay/src/stats.rs b/rs/moq-relay/src/stats.rs index 9a1d53e7ba..6fb8bfbb9a 100644 --- a/rs/moq-relay/src/stats.rs +++ b/rs/moq-relay/src/stats.rs @@ -69,7 +69,7 @@ pub struct Config { /// single `/node/` broadcast for the whole node. Set to 1 to /// publish a per-first-segment broadcast (e.g. per tenant), so a consumer can /// announce-scope to just that group rather than slurping every node's full - /// stats. See [`moq_stats::ProducerConfig::depth`]. + /// stats. See [`moq_stats::produce::Config::depth`]. #[usage( long = "stats-depth", env = "MOQ_STATS_DEPTH", @@ -101,14 +101,14 @@ impl Config { /// last clone of the producer drops). pub fn build(&self, origin: origin::Producer) -> moq_stats::Producer { if !self.enabled { - return moq_stats::Producer::new(moq_stats::ProducerConfig::new()); + return moq_stats::Producer::new(moq_stats::produce::Config::new()); } let prefix = self.prefix.clone(); let interval = Duration::from_secs(self.interval.max(1)); let node = self.node.clone().map(PathOwned::from); let depth = self.depth; tracing::info!(prefix, interval_secs = interval.as_secs(), node = ?node, depth, "stats publishing enabled"); - let config = moq_stats::ProducerConfig::new() + let config = moq_stats::produce::Config::new() .with_origin(origin) .with_prefix(prefix) .with_interval(interval) diff --git a/rs/moq-room/src/claims.rs b/rs/moq-room/src/claims.rs index e8fac7e8d7..038273ab37 100644 --- a/rs/moq-room/src/claims.rs +++ b/rs/moq-room/src/claims.rs @@ -9,13 +9,15 @@ /// and `publish` is `/**` so a participant cannot publish at anyone /// else's paths. An identity that is empty after normalization or contains `*` is /// refused, since a wildcard would let the participant publish as someone else. -pub fn claims(room: impl Into, identity: &str) -> Result { - if moq_net::Path::new(identity).is_empty() { +pub fn claims(room: impl moq_net::AsPath, identity: impl moq_net::AsPath) -> Result { + let room = room.as_path(); + let identity = identity.as_path(); + if identity.is_empty() { return Err(crate::Error::EmptyIdentity); } - let publish = moq_auth::Pattern::subtree(identity)?; + let publish = moq_auth::Pattern::subtree(identity.as_str())?; Ok(moq_auth::Claims::default() - .with_root(room) + .with_root(room.as_str()) .with_subscribe([moq_auth::Pattern::all()]) .with_publish([publish])) } diff --git a/rs/moq-rtc/Cargo.toml b/rs/moq-rtc/Cargo.toml index 424767a87d..4be0eba1a3 100644 --- a/rs/moq-rtc/Cargo.toml +++ b/rs/moq-rtc/Cargo.toml @@ -16,7 +16,6 @@ categories = ["multimedia", "network-programming", "web-programming"] # The CLI lives in `moq-cli` (the `rtc` subcommand); embedders mount the routers # / dial with `Client` against their own origin. [dependencies] -anyhow = { workspace = true, features = ["backtrace"] } axum = { workspace = true } bytes = { workspace = true } hang = { workspace = true } diff --git a/rs/moq-rtc/src/client/whep.rs b/rs/moq-rtc/src/client/whep.rs index 17cf177162..c1a3e764bc 100644 --- a/rs/moq-rtc/src/client/whep.rs +++ b/rs/moq-rtc/src/client/whep.rs @@ -37,9 +37,7 @@ pub(crate) async fn dial(client: &Client, url: Url, broadcast: moq_net::broadcas api.add_media(MediaKind::Audio, Direction::RecvOnly, None, None, None), api.add_media(MediaKind::Video, Direction::RecvOnly, None, None, None), ]; - let (offer, pending) = api - .apply() - .ok_or_else(|| Error::Other(anyhow::anyhow!("no SDP changes to apply")))?; + let (offer, pending) = api.apply().ok_or(Error::NoSdpChanges)?; let res = client .http() @@ -48,17 +46,13 @@ pub(crate) async fn dial(client: &Client, url: Url, broadcast: moq_net::broadcas .header(reqwest::header::ACCEPT, "application/sdp") .body(offer.to_sdp_string()) .send() - .await - .map_err(|err| Error::Other(anyhow::anyhow!("WHEP POST failed: {err}")))?; + .await?; if !res.status().is_success() { - return Err(Error::Other(anyhow::anyhow!("WHEP server returned {}", res.status()))); + return Err(Error::HttpStatus(res.status().as_u16())); } - let body = res - .text() - .await - .map_err(|err| Error::Other(anyhow::anyhow!("reading WHEP answer body: {err}")))?; + let body = res.text().await?; let answer = SdpAnswer::from_sdp_string(&body).map_err(|err| Error::InvalidSdp(err.to_string()))?; rtc.sdp_api().accept_answer(pending, answer).map_err(Error::rtc)?; diff --git a/rs/moq-rtc/src/client/whip.rs b/rs/moq-rtc/src/client/whip.rs index 0d6cec9ff9..2e02db0a8d 100644 --- a/rs/moq-rtc/src/client/whip.rs +++ b/rs/moq-rtc/src/client/whip.rs @@ -26,9 +26,7 @@ pub(crate) async fn dial( let source = EgressSource::new(moq_mux::Source::new(origin, path)).await?; let codecs = source.catalog_codecs(); if codecs.is_empty() { - return Err(Error::Other(anyhow::anyhow!( - "catalog has no codecs we can egress (Opus / H.264 / H.265 / VP8 / VP9 / AV1)" - ))); + return Err(Error::NoRenditions); } let (socket, candidates) = session::bind_udp(&client.config().ice_candidates).await?; @@ -54,9 +52,7 @@ pub(crate) async fn dial( { mids.push(api.add_media(MediaKind::Video, Direction::SendOnly, None, None, None)); } - let (offer, pending) = api - .apply() - .ok_or_else(|| Error::Other(anyhow::anyhow!("no SDP changes to apply")))?; + let (offer, pending) = api.apply().ok_or(Error::NoSdpChanges)?; let res = client .http() @@ -65,17 +61,13 @@ pub(crate) async fn dial( .header(reqwest::header::ACCEPT, "application/sdp") .body(offer.to_sdp_string()) .send() - .await - .map_err(|err| Error::Other(anyhow::anyhow!("WHIP POST failed: {err}")))?; + .await?; if !res.status().is_success() { - return Err(Error::Other(anyhow::anyhow!("WHIP server returned {}", res.status()))); + return Err(Error::HttpStatus(res.status().as_u16())); } - let body = res - .text() - .await - .map_err(|err| Error::Other(anyhow::anyhow!("reading WHIP answer body: {err}")))?; + let body = res.text().await?; let answer = SdpAnswer::from_sdp_string(&body).map_err(|err| Error::InvalidSdp(err.to_string()))?; rtc.sdp_api().accept_answer(pending, answer).map_err(Error::rtc)?; diff --git a/rs/moq-rtc/src/codec/av1.rs b/rs/moq-rtc/src/codec/av1.rs index 4effeb6dc8..b16898b421 100644 --- a/rs/moq-rtc/src/codec/av1.rs +++ b/rs/moq-rtc/src/codec/av1.rs @@ -24,8 +24,7 @@ impl Bridge { impl codec::Bridge for Bridge { fn push(&mut self, frame: codec::Frame) -> Result<()> { - let pts = moq_net::Timestamp::from_micros(frame.timestamp_us) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("invalid timestamp: {err}")))?; + let pts = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_mux::Error::from)?; // str0m hands over one whole temporal unit per frame, so flush to emit it. let mut frames = self.split.decode(&frame.payload, Some(pts))?; frames.extend(self.split.flush(Some(pts))?); diff --git a/rs/moq-rtc/src/codec/h264.rs b/rs/moq-rtc/src/codec/h264.rs index d11170cae0..961a692255 100644 --- a/rs/moq-rtc/src/codec/h264.rs +++ b/rs/moq-rtc/src/codec/h264.rs @@ -23,8 +23,7 @@ impl Bridge { impl codec::Bridge for Bridge { fn push(&mut self, frame: codec::Frame) -> Result<()> { - let pts = moq_net::Timestamp::from_micros(frame.timestamp_us) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("invalid timestamp: {err}")))?; + let pts = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_mux::Error::from)?; // str0m hands over one whole access unit per frame, so flush to emit it. let mut frames = self.split.decode(&frame.payload, Some(pts))?; frames.extend(self.split.flush(Some(pts))?); diff --git a/rs/moq-rtc/src/codec/h265.rs b/rs/moq-rtc/src/codec/h265.rs index 14ce21316a..606cad6af3 100644 --- a/rs/moq-rtc/src/codec/h265.rs +++ b/rs/moq-rtc/src/codec/h265.rs @@ -25,8 +25,7 @@ impl Bridge { impl codec::Bridge for Bridge { fn push(&mut self, frame: codec::Frame) -> Result<()> { - let pts = moq_net::Timestamp::from_micros(frame.timestamp_us) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("invalid timestamp: {err}")))?; + let pts = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_mux::Error::from)?; // str0m hands over one whole access unit per frame, so flush to emit it. let mut frames = self.split.decode(&frame.payload, Some(pts))?; frames.extend(self.split.flush(Some(pts))?); diff --git a/rs/moq-rtc/src/codec/mod.rs b/rs/moq-rtc/src/codec/mod.rs index 5b3f4857e9..d4b7b23f9b 100644 --- a/rs/moq-rtc/src/codec/mod.rs +++ b/rs/moq-rtc/src/codec/mod.rs @@ -145,9 +145,7 @@ impl DeferredVideo { } let DeferredState::Pending(pending) = std::mem::replace(&mut self.state, DeferredState::Poisoned) else { - return Err(crate::Error::Other(anyhow::anyhow!( - "video bridge initialization already failed" - ))); + return Err(crate::Error::BridgeFailed); }; let reserved = pending.catalog.reserve(); let abort = pending.track.clone(); @@ -269,7 +267,7 @@ impl Track { } => { let prefix = frame.keyframe.then(|| keyframe_prefix.as_ref()); moq_mux::codec::annexb::from_length_prefixed(&frame.payload, *length_size, prefix) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("annexb: {err}")))? + .map_err(moq_mux::Error::from)? } }; if payload.is_empty() { @@ -292,17 +290,17 @@ fn h264_convert(config: &VideoConfig) -> Result { let Some(avcc) = config.description.as_ref().filter(|d| !d.is_empty()) else { return Ok(TrackConvert::Passthrough); }; - let params = moq_mux::codec::h264::Avcc::parse(avcc) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("avcc parse: {err}")))?; + let params = moq_mux::codec::h264::Avcc::parse(avcc).map_err(moq_mux::Error::from)?; // Without SPS+PPS the keyframe prefix would be empty and every keyframe // would reach the peer without inline parameter sets, i.e. undecodable. // Fail loudly instead, matching moq-mux's `h264::Export`. if params.sps.is_empty() || params.pps.is_empty() { - return Err(crate::Error::Other(anyhow::anyhow!( - "avc1 avcC is missing parameter sets (sps={}, pps={})", - params.sps.len(), - params.pps.len() - ))); + return Err(moq_mux::Error::H264(moq_mux::codec::h264::Error::MissingParamSets { + name: "WebRTC rendition".to_string(), + sps: params.sps.len(), + pps: params.pps.len(), + }) + .into()); } let keyframe_prefix = moq_mux::codec::annexb::build_prefix(params.sps.iter().chain(params.pps.iter())); Ok(TrackConvert::LengthPrefixed { @@ -320,17 +318,17 @@ fn h265_convert(config: &VideoConfig) -> Result { let Some(hvcc) = config.description.as_ref().filter(|d| !d.is_empty()) else { return Ok(TrackConvert::Passthrough); }; - let params = moq_mux::codec::h265::Hvcc::parse(hvcc) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("hvcc parse: {err}")))?; + let params = moq_mux::codec::h265::Hvcc::parse(hvcc).map_err(moq_mux::Error::from)?; // Same reasoning as `h264_convert`: a keyframe with no inline VPS/SPS/PPS // is undecodable, so reject an hvcC that omits any of them. if params.vps.is_empty() || params.sps.is_empty() || params.pps.is_empty() { - return Err(crate::Error::Other(anyhow::anyhow!( - "hvc1 hvcC is missing parameter sets (vps={}, sps={}, pps={})", - params.vps.len(), - params.sps.len(), - params.pps.len() - ))); + return Err(moq_mux::Error::H265(moq_mux::codec::h265::Error::MissingParamSets { + name: "WebRTC rendition".to_string(), + vps: params.vps.len(), + sps: params.sps.len(), + pps: params.pps.len(), + }) + .into()); } let keyframe_prefix = moq_mux::codec::annexb::build_prefix(params.vps.iter().chain(params.sps.iter()).chain(params.pps.iter())); diff --git a/rs/moq-rtc/src/codec/opus.rs b/rs/moq-rtc/src/codec/opus.rs index 38b3196790..dae4454f4e 100644 --- a/rs/moq-rtc/src/codec/opus.rs +++ b/rs/moq-rtc/src/codec/opus.rs @@ -25,8 +25,7 @@ impl Bridge { impl codec::Bridge for Bridge { fn push(&mut self, frame: codec::Frame) -> Result<()> { - let pts = moq_net::Timestamp::from_micros(frame.timestamp_us) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("invalid timestamp: {err}")))?; + let pts = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_mux::Error::from)?; self.import.decode(&frame.payload, Some(pts))?; // The importer accumulates; cut each packet into its own group (one QUIC stream) so the // relay forwards it without waiting for the next. diff --git a/rs/moq-rtc/src/codec/vp8.rs b/rs/moq-rtc/src/codec/vp8.rs index 7013d8d331..4cede5be7e 100644 --- a/rs/moq-rtc/src/codec/vp8.rs +++ b/rs/moq-rtc/src/codec/vp8.rs @@ -21,8 +21,7 @@ impl Bridge { impl codec::Bridge for Bridge { fn push(&mut self, frame: codec::Frame) -> Result<()> { - let pts = moq_net::Timestamp::from_micros(frame.timestamp_us) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("invalid timestamp: {err}")))?; + let pts = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_mux::Error::from)?; self.import.decode(frame.payload, pts) } diff --git a/rs/moq-rtc/src/codec/vp9.rs b/rs/moq-rtc/src/codec/vp9.rs index f2d60b3754..0bd213df9c 100644 --- a/rs/moq-rtc/src/codec/vp9.rs +++ b/rs/moq-rtc/src/codec/vp9.rs @@ -21,8 +21,7 @@ impl Bridge { impl codec::Bridge for Bridge { fn push(&mut self, frame: codec::Frame) -> Result<()> { - let pts = moq_net::Timestamp::from_micros(frame.timestamp_us) - .map_err(|err| crate::Error::Other(anyhow::anyhow!("invalid timestamp: {err}")))?; + let pts = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_mux::Error::from)?; self.import.decode(frame.payload, pts) } diff --git a/rs/moq-rtc/src/egress.rs b/rs/moq-rtc/src/egress.rs index dd3313614b..7c07f88ec5 100644 --- a/rs/moq-rtc/src/egress.rs +++ b/rs/moq-rtc/src/egress.rs @@ -107,11 +107,7 @@ impl EgressSource { /// which takes the receiver via [`Self::take_writes`]. pub async fn new(source: moq_mux::Source) -> Result { let mut consumer = source.catalog::<()>(moq_mux::catalog::CatalogFormat::Hang).await?; - let catalog = consumer - .next() - .await - .map_err(|err| Error::Other(anyhow::anyhow!("catalog subscribe: {err}")))? - .ok_or_else(|| Error::Other(anyhow::anyhow!("catalog closed before first snapshot")))?; + let catalog = consumer.next().await?.ok_or(Error::CatalogClosed)?; let (tx, rx) = mpsc::channel(64); Ok(Self { diff --git a/rs/moq-rtc/src/error.rs b/rs/moq-rtc/src/error.rs index af832ec185..76faaf5832 100644 --- a/rs/moq-rtc/src/error.rs +++ b/rs/moq-rtc/src/error.rs @@ -22,6 +22,22 @@ pub enum Error { #[error("ICE did not connect before the establishment deadline")] IceTimeout, + /// A broadcast did not produce a catalog before the negotiation deadline. + #[error("catalog did not arrive before the negotiation deadline")] + CatalogTimeout, + + /// The catalog track closed before publishing its first snapshot. + #[error("catalog closed before its first snapshot")] + CatalogClosed, + + /// The catalog has no rendition this gateway can send over WebRTC. + #[error("catalog has no WebRTC-compatible renditions")] + NoRenditions, + + /// The WebRTC engine produced no SDP changes for an offer. + #[error("no SDP changes to apply")] + NoSdpChanges, + /// I/O error on the media socket (bind, send, or receive). #[error("io error: {0}")] Io(#[from] std::io::Error), @@ -34,6 +50,14 @@ pub enum Error { #[error("mux error: {0}")] Mux(#[from] moq_mux::Error), + /// HTTP transport failed while dialing a WHIP or WHEP endpoint. + #[error("http error: {0}")] + Http(std::sync::Arc), + + /// A WHIP or WHEP endpoint rejected the HTTP request. + #[error("HTTP endpoint returned status {0}")] + HttpStatus(u16), + /// Error from the WebRTC engine (SDP negotiation, DTLS, media state). #[error("rtc error: {0}")] Rtc(String), @@ -42,9 +66,15 @@ pub enum Error { #[error("rtc input error: {0}")] RtcInput(String), - /// Catch-all for gateway logic that reports via `anyhow`. - #[error(transparent)] - Other(#[from] anyhow::Error), + /// An internal video bridge could not be initialized. + #[error("video bridge initialization failed")] + BridgeFailed, +} + +impl From for Error { + fn from(err: reqwest::Error) -> Self { + Self::Http(std::sync::Arc::new(err)) + } } impl Error { diff --git a/rs/moq-rtc/src/lib.rs b/rs/moq-rtc/src/lib.rs index f8789c494a..c394fa19bb 100644 --- a/rs/moq-rtc/src/lib.rs +++ b/rs/moq-rtc/src/lib.rs @@ -15,9 +15,7 @@ //! //! ## Embedding //! -//! Build a [`Server`] over your own -//! [`OriginProducer`](moq_net::origin::Producer) / -//! [`OriginConsumer`](moq_net::origin::Consumer) and merge +//! Build a [`Server`] and pass your own origin handles when merging //! [`Server::publish_router`] / [`Server::subscribe_router`] into your own axum //! app, or dial out with [`Client`]. A command-line interface is provided by the //! `moq-cli` binary, on top of this library. @@ -114,14 +112,10 @@ mod tests { drop(announcements); let server_origin = moq_tokio::origin::spawn(); - let server = Server::new( - server::Config::default(), - server_origin.clone(), - server_origin.consume(), - ); + let server = Server::new(server::Config::default()); let app = Router::new() - .nest("/whip", server.publish_router()) - .nest("/whep", server.subscribe_router()); + .nest("/whip", server.publish_router(server_origin.clone())) + .nest("/whep", server.subscribe_router(server_origin.consume())); let listener = tokio::net::TcpListener::bind("127.0.0.1:0") .await .expect("bind HTTP listener"); diff --git a/rs/moq-rtc/src/server/mod.rs b/rs/moq-rtc/src/server/mod.rs index 055a15448b..3585f15b8a 100644 --- a/rs/moq-rtc/src/server/mod.rs +++ b/rs/moq-rtc/src/server/mod.rs @@ -17,7 +17,6 @@ use std::sync::{Arc, Mutex}; use std::time::Duration; use axum::Router; -use axum::extract::{Path, State}; use axum::http::{HeaderValue, StatusCode, Uri}; use tokio::sync::{OnceCell, oneshot}; @@ -175,11 +174,7 @@ impl Default for Config { } } -/// Glue that owns the moq-net origin pair and hands axum routers to the caller. -/// -/// `publisher` is where `server publish` (WHIP) writes ingested broadcasts; -/// `subscriber` is what `server subscribe` (WHEP) reads from. They're -/// typically the two halves of the same upstream [`moq_net::Session`]. +/// Shared WebRTC media state that hands axum routers to the caller. #[derive(Clone)] pub struct Server { inner: Arc, @@ -187,9 +182,6 @@ pub struct Server { struct Inner { config: Config, - publisher: moq_net::origin::Producer, - /// Source for `server subscribe` (WHEP) egress. - subscriber: moq_net::origin::Consumer, /// The shared media socket + demux, bound lazily on the first accept so /// `Server::new` can stay synchronous (and an idle server binds no port). mux: OnceCell, @@ -200,14 +192,11 @@ struct Inner { } impl Server { - /// Build a server. `publisher` receives WHIP broadcasts; `subscriber` - /// is the source for WHEP egress. - pub fn new(config: Config, publisher: moq_net::origin::Producer, subscriber: moq_net::origin::Consumer) -> Self { + /// Build a server with shared ICE and media settings. + pub fn new(config: Config) -> Self { Self { inner: Arc::new(Inner { config, - publisher, - subscriber, mux: OnceCell::new(), sessions: Mutex::new(HashMap::new()), }), @@ -229,8 +218,8 @@ impl Server { /// no authentication. To own the route and authorize requests yourself /// (resolving the broadcast name from a verified token), skip the router and /// call [`whip::accept`] directly from your own handler. - pub fn publish_router(&self) -> Router { - whip::router(self.clone()) + pub fn publish_router(&self, publisher: moq_net::origin::Producer) -> Router { + whip::router(self.clone(), publisher) } /// Router for `server subscribe` (WHEP). Mount under whichever HTTP path @@ -240,22 +229,14 @@ impl Server { /// no authentication. To own the route and authorize requests yourself /// (resolving the broadcast name from a verified token), skip the router and /// call [`whep::accept`] directly from your own handler. - pub fn subscribe_router(&self) -> Router { - whep::router(self.clone()) + pub fn subscribe_router(&self, subscriber: moq_net::origin::Consumer) -> Router { + whep::router(self.clone(), subscriber) } pub(crate) fn config(&self) -> &Config { &self.inner.config } - pub(crate) fn publisher(&self) -> &moq_net::origin::Producer { - &self.inner.publisher - } - - pub(crate) fn subscriber(&self) -> &moq_net::origin::Consumer { - &self.inner.subscriber - } - /// Register a session under its resource id, returning the cancel receiver. /// Called by [`whip::accept`] / [`whep::accept`] before returning the /// negotiated session runner. @@ -288,8 +269,8 @@ impl Server { /// Shared `DELETE` handler for both bundled routers: parse the resource id from /// the trailing path segment and terminate the matching session. -pub(crate) async fn delete(State(server): State, Path(path): Path) -> StatusCode { - match crate::sdp::parse_resource_id(&path) { +pub(crate) fn delete(server: &Server, path: &str) -> StatusCode { + match crate::sdp::parse_resource_id(path) { Ok(id) if server.terminate(&id.to_string()) => StatusCode::OK, Ok(_) => StatusCode::NOT_FOUND, Err(_) => StatusCode::BAD_REQUEST, @@ -298,25 +279,10 @@ pub(crate) async fn delete(State(server): State, Path(path): Path moq_net::origin::Producer { - let (producer, driver) = moq_net::origin::Producer::new(moq_net::origin::Config::default()); - if tokio::runtime::Handle::try_current().is_ok() { - tokio::spawn(driver.run(moq_tokio::runtime::Runtime::<()>::new())); - } else { - // A sync test: nothing polls the driver, and dropping it would tear - // the origin down, so leak it and rely on the synchronous half. - std::mem::forget(driver); - } - producer - } - use super::*; fn server() -> Server { - let publisher = produce_origin(); - let subscriber = produce_origin().consume(); - Server::new(Config::default(), publisher, subscriber) + Server::new(Config::default()) } #[test] diff --git a/rs/moq-rtc/src/server/whep.rs b/rs/moq-rtc/src/server/whep.rs index 400b122a4a..4a3b3c704c 100644 --- a/rs/moq-rtc/src/server/whep.rs +++ b/rs/moq-rtc/src/server/whep.rs @@ -19,6 +19,12 @@ use crate::{Error, Result, egress::EgressSource, sdp, server::Server, session}; pub use crate::server::Response; +#[derive(Clone)] +struct RouterState { + server: Server, + subscriber: moq_net::origin::Consumer, +} + /// How long WHEP negotiation waits for the broadcast's first catalog snapshot /// before failing the request. A broadcast can be announced (or served by a /// dynamic origin fallback) yet never publish a catalog; without a bound the @@ -26,21 +32,21 @@ pub use crate::server::Response; const CATALOG_TIMEOUT: Duration = Duration::from_secs(5); /// Build the WHEP axum router. -pub fn router(server: Server) -> Router { +pub fn router(server: Server, subscriber: moq_net::origin::Consumer) -> Router { Router::new() - .route("/{*path}", post(handle).delete(crate::server::delete)) - .with_state(server) + .route("/{*path}", post(handle).delete(delete)) + .with_state(RouterState { server, subscriber }) } async fn handle( - server: State, + state: State, path: Path, OriginalUri(uri): OriginalUri, headers: HeaderMap, body: Bytes, ) -> HttpResponse { - let (server, path) = (server.0, path.0); - match accept_offer(&server, &path, &headers, body).await { + let (state, path) = (state.0, path.0); + match accept_offer(&state.server, &state.subscriber, &path, &headers, body).await { Ok(response) => { let Response { resource_id, @@ -66,12 +72,22 @@ async fn handle( /// Router glue: enforce the WHEP `Content-Type` then hand the raw offer to /// [`accept`], using the request path as the (unauthenticated) broadcast name. -async fn accept_offer(server: &Server, path: &str, headers: &HeaderMap, body: Bytes) -> Result { +async fn accept_offer( + server: &Server, + subscriber: &moq_net::origin::Consumer, + path: &str, + headers: &HeaderMap, + body: Bytes, +) -> Result { if !is_sdp(headers) { return Err(Error::InvalidSdp("expected Content-Type: application/sdp".into())); } let offer = std::str::from_utf8(&body).map_err(|err| Error::InvalidSdp(err.to_string()))?; - accept(server, server.subscriber(), path, offer).await + accept(server, subscriber, path, offer).await +} + +async fn delete(State(state): State, Path(path): Path) -> StatusCode { + crate::server::delete(&state.server, &path) } /// Accept a WHEP SDP offer and egress the MoQ broadcast `broadcast` (a path @@ -93,7 +109,7 @@ async fn accept_offer(server: &Server, path: &str, headers: &HeaderMap, body: By /// `offer` is the raw SDP body; the caller is responsible for checking the /// `Content-Type: application/sdp` request header. Fails with [`Error::InvalidSdp`] /// on a malformed offer, and surfaces a not-announced broadcast (or one outside -/// `subscriber`'s scope) as [`Error::Other`]. +/// `subscriber`'s scope) as [`Error::Moq`]. pub async fn accept( server: &Server, subscriber: &moq_net::origin::Consumer, @@ -113,16 +129,10 @@ pub async fn accept( // announced-but-catalog-less one, would otherwise park this handler forever. let source = tokio::time::timeout(CATALOG_TIMEOUT, EgressSource::new(source)) .await - .map_err(|_| { - Error::Other(anyhow::anyhow!( - "broadcast {broadcast} did not resolve with a catalog within {CATALOG_TIMEOUT:?}" - )) - })??; + .map_err(|_| Error::CatalogTimeout)??; let codecs = source.catalog_codecs(); if codecs.is_empty() { - return Err(Error::Other(anyhow::anyhow!( - "catalog has no codecs we can egress (Opus / H.264 / H.265 / VP8 / VP9 / AV1)" - ))); + return Err(Error::NoRenditions); } // Register a session on the shared media mux (see whip::accept). Restrict our @@ -175,6 +185,8 @@ fn status_for(err: &Error) -> StatusCode { Error::InvalidSdp(_) => StatusCode::BAD_REQUEST, Error::UnsupportedCodec(_) => StatusCode::UNSUPPORTED_MEDIA_TYPE, Error::SessionNotFound => StatusCode::NOT_FOUND, + Error::Moq(moq_net::Error::Unauthorized) => StatusCode::UNAUTHORIZED, + Error::Moq(moq_net::Error::NotFound | moq_net::Error::Unroutable) => StatusCode::NOT_FOUND, _ => StatusCode::INTERNAL_SERVER_ERROR, } } diff --git a/rs/moq-rtc/src/server/whip.rs b/rs/moq-rtc/src/server/whip.rs index e0d7c877c8..6d618aad81 100644 --- a/rs/moq-rtc/src/server/whip.rs +++ b/rs/moq-rtc/src/server/whip.rs @@ -18,21 +18,27 @@ use crate::{Error, Result, ingest::IngestSink, sdp, server::Server, session}; pub use crate::server::Response; +#[derive(Clone)] +struct RouterState { + server: Server, + publisher: moq_net::origin::Producer, +} + /// Build the WHIP axum router. -pub fn router(server: Server) -> Router { +pub fn router(server: Server, publisher: moq_net::origin::Producer) -> Router { Router::new() - .route("/{*path}", post(handle).delete(crate::server::delete)) - .with_state(server) + .route("/{*path}", post(handle).delete(delete)) + .with_state(RouterState { server, publisher }) } async fn handle( - State(server): State, + State(state): State, Path(path): Path, OriginalUri(uri): OriginalUri, headers: HeaderMap, body: Bytes, ) -> HttpResponse { - match accept_offer(&server, &path, &headers, body).await { + match accept_offer(&state.server, &state.publisher, &path, &headers, body).await { Ok(response) => { let Response { resource_id, @@ -58,12 +64,22 @@ async fn handle( /// Router glue: enforce the WHIP `Content-Type` then hand the raw offer to /// [`accept`], using the request path as the (unauthenticated) broadcast name. -async fn accept_offer(server: &Server, path: &str, headers: &HeaderMap, body: Bytes) -> Result { +async fn accept_offer( + server: &Server, + publisher: &moq_net::origin::Producer, + path: &str, + headers: &HeaderMap, + body: Bytes, +) -> Result { if !is_sdp(headers) { return Err(Error::InvalidSdp("expected Content-Type: application/sdp".into())); } let offer = std::str::from_utf8(&body).map_err(|err| Error::InvalidSdp(err.to_string()))?; - accept(server, server.publisher(), path, offer).await + accept(server, publisher, path, offer).await +} + +async fn delete(State(state): State, Path(path): Path) -> StatusCode { + crate::server::delete(&state.server, &path) } /// Accept a WHIP SDP offer and publish the negotiated WebRTC media into @@ -85,7 +101,7 @@ async fn accept_offer(server: &Server, path: &str, headers: &HeaderMap, body: By /// `offer` is the raw SDP body; the caller is responsible for checking the /// `Content-Type: application/sdp` request header. Fails with /// [`Error::InvalidSdp`] on a malformed offer and surfaces -/// [`moq_net::Error::Unauthorized`] (as [`Error::Other`]) if `broadcast` is +/// [`moq_net::Error::Unauthorized`] if `broadcast` is /// outside `publisher`'s scope. pub async fn accept( server: &Server, @@ -99,12 +115,8 @@ pub async fn accept( // Create the broadcast on the publish origin before negotiating, so a // fast subscriber doesn't see a 404 in the gap between the SDP answer // and the first RTP packet. - let producer = publisher - .create_broadcast(&broadcast) - .map_err(|err| Error::Other(anyhow::anyhow!("failed to create broadcast: {err}")))?; - producer - .announce(moq_net::origin::Route::default()) - .map_err(|err| Error::Other(anyhow::anyhow!("failed to announce broadcast: {err}")))?; + let producer = publisher.create_broadcast(&broadcast)?; + producer.announce(moq_net::origin::Route::default())?; let handle = producer.clone(); let config = moq_mux::catalog::Config::default() @@ -161,6 +173,7 @@ fn status_for(err: &Error) -> StatusCode { Error::InvalidSdp(_) => StatusCode::BAD_REQUEST, Error::UnsupportedCodec(_) => StatusCode::UNSUPPORTED_MEDIA_TYPE, Error::SessionNotFound => StatusCode::NOT_FOUND, + Error::Moq(moq_net::Error::Unauthorized) => StatusCode::UNAUTHORIZED, _ => StatusCode::INTERNAL_SERVER_ERROR, } } diff --git a/rs/moq-rtmp/README.md b/rs/moq-rtmp/README.md index e00b762d07..4d3197262b 100644 --- a/rs/moq-rtmp/README.md +++ b/rs/moq-rtmp/README.md @@ -34,7 +34,7 @@ local with no extra hop: ```rust let mut rtmp = moq_rtmp::Config::default(); rtmp.listen = Some("0.0.0.0:1935".parse()?); -rtmp.prefix = "live/".to_string(); +rtmp.prefix = "live".into(); // `origin` is your relay's local origin (e.g. `cluster.origin.clone()`). tokio::select! { diff --git a/rs/moq-rtmp/src/error.rs b/rs/moq-rtmp/src/error.rs index 5dfe66343d..9f0f691b07 100644 --- a/rs/moq-rtmp/src/error.rs +++ b/rs/moq-rtmp/src/error.rs @@ -14,10 +14,9 @@ pub enum Error { #[error("io: {0}")] Io(Arc), - /// Catch-all for ingest logic that reports via `anyhow` (the RTMP session and - /// the moq-mux demuxer surface their errors this way). - #[error("{0}")] - Other(Arc), + /// The RTMP handshake or session state machine failed. + #[error("rtmp session: {0}")] + Session(String), } impl From for Error { @@ -28,7 +27,7 @@ impl From for Error { impl From for Error { fn from(err: anyhow::Error) -> Self { - Error::Other(Arc::new(err)) + Error::Session(err.to_string()) } } diff --git a/rs/moq-rtmp/src/listen.rs b/rs/moq-rtmp/src/listen.rs index 7c34b22d14..44cb337601 100644 --- a/rs/moq-rtmp/src/listen.rs +++ b/rs/moq-rtmp/src/listen.rs @@ -23,7 +23,7 @@ use std::net::SocketAddr; use std::sync::{Arc, Mutex}; use std::time::Duration; -use moq_net::origin; +use moq_net::{Path, PathOwned, origin}; use crate::Result; use crate::server::{Request, Server}; @@ -41,9 +41,9 @@ pub struct Config { /// RTMP ingest is disabled. pub listen: Option, - /// Prefix prepended to every ingested broadcast path. Lets one listener - /// namespace all of its streams (e.g. `live/`). - pub prefix: String, + /// Path prefix prepended to every ingested broadcast path. Lets one listener + /// namespace all of its streams (e.g. `live`). + pub prefix: PathOwned, /// How long a play's FLV muxer waits for a stalled group before skipping to a /// newer one (the moq-level frame-drop latency). Defaults to @@ -79,7 +79,7 @@ impl Default for Config { fn default() -> Self { Self { listen: None, - prefix: String::new(), + prefix: Path::empty().to_owned(), export_max_age: crate::DEFAULT_MAX_AGE, import_max_age: None, #[cfg(feature = "tls")] @@ -130,7 +130,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { // of clobbering the live one. This lives on Config so cloned RTMP/RTMPS // listeners share the same claim table. let active = config.active.clone(); - let prefix = Arc::new(config.prefix); + let prefix = config.prefix; let export_max_age = config.export_max_age; let import_max_age = config.import_max_age; // Players are served out of the same origin the publishers write into. @@ -154,7 +154,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { }; // Claim the path before accepting; the guard releases it when the // connection task ends (success, error, or panic). - let Some(_guard) = active.claim(&path) else { + let Some(_guard) = active.claim(path.as_str()) else { tracing::warn!(%peer, %path, "rejecting RTMP publish: path already being published"); let _ = publish.reject("path already being published").await; return; @@ -183,7 +183,9 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { } } - Err(anyhow::anyhow!("RTMP listener stopped accepting connections").into()) + Err(crate::Error::Session( + "listener stopped accepting connections".to_string(), + )) } /// Derive a broadcast path from an RTMP app and stream key, applying `prefix`. @@ -191,7 +193,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { /// `rtmp://host//` maps to `/`, falling back to just /// the app (or just the key) when the other half is empty. Returns `None` when /// there's nothing usable to route on. -pub(crate) fn resolve_path(prefix: &str, app: &str, key: &str) -> Option { +pub(crate) fn resolve_path(prefix: &Path, app: &str, key: &str) -> Option { let app = app.trim_matches('/').trim(); let key = key.trim_matches('/').trim(); let name = match (app.is_empty(), key.is_empty()) { @@ -200,7 +202,7 @@ pub(crate) fn resolve_path(prefix: &str, app: &str, key: &str) -> Option (true, false) => key.to_string(), (false, false) => format!("{app}/{key}"), }; - Some(format!("{prefix}{name}")) + Some(prefix.join(name)) } /// The set of broadcast paths with a live ingest, used to reject duplicate @@ -241,27 +243,36 @@ mod tests { #[test] fn app_and_key() { - assert_eq!(resolve_path("", "live", "cam0").as_deref(), Some("live/cam0")); + assert_eq!( + resolve_path(&Path::empty(), "live", "cam0").unwrap().as_str(), + "live/cam0" + ); } #[test] fn app_only() { - assert_eq!(resolve_path("", "cam0", "").as_deref(), Some("cam0")); + assert_eq!(resolve_path(&Path::empty(), "cam0", "").unwrap().as_str(), "cam0"); } #[test] fn prefix_is_prepended() { - assert_eq!(resolve_path("live/", "cam0", "").as_deref(), Some("live/cam0")); + assert_eq!( + resolve_path(&Path::new("live"), "cam0", "").unwrap().as_str(), + "live/cam0" + ); } #[test] fn slashes_are_trimmed() { - assert_eq!(resolve_path("", "/live/", "/cam0/").as_deref(), Some("live/cam0")); + assert_eq!( + resolve_path(&Path::empty(), "/live/", "/cam0/").unwrap().as_str(), + "live/cam0" + ); } #[test] fn empty_is_rejected() { - assert_eq!(resolve_path("", "", ""), None); + assert_eq!(resolve_path(&Path::empty(), "", ""), None); } #[test] diff --git a/rs/moq-rtmp/src/server.rs b/rs/moq-rtmp/src/server.rs index 398c0b36f1..250eac7c5e 100644 --- a/rs/moq-rtmp/src/server.rs +++ b/rs/moq-rtmp/src/server.rs @@ -534,8 +534,8 @@ impl Publish { /// segmented egress (HLS/DASH) reading the broadcast downstream, which may only /// advertise segments that are still fetchable. Lower it when nothing reads history /// and the memory matters. - pub fn with_max_age(mut self, max_age: impl Into>) -> Self { - self.max_age = max_age.into(); + pub fn with_max_age(mut self, max_age: Option) -> Self { + self.max_age = max_age; self } diff --git a/rs/moq-srt/Cargo.toml b/rs/moq-srt/Cargo.toml index 0893c71652..98e89fee91 100644 --- a/rs/moq-srt/Cargo.toml +++ b/rs/moq-srt/Cargo.toml @@ -16,7 +16,6 @@ categories = ["multimedia", "network-programming", "web-programming"] # subcommand); a relay can also embed `moq_srt::run` / `Server` against its own # origin. [dependencies] -anyhow = { workspace = true, features = ["backtrace"] } bytes = { workspace = true } futures = { workspace = true } moq-mux = { workspace = true } @@ -31,3 +30,4 @@ tracing = { workspace = true } # Builds the synthetic broadcasts the egress tests mux back to TS. hang = { workspace = true } moq-tokio = { workspace = true } +srt-protocol = "0.4" diff --git a/rs/moq-srt/README.md b/rs/moq-srt/README.md index 4939ddf85f..b7eb3cf0f1 100644 --- a/rs/moq-srt/README.md +++ b/rs/moq-srt/README.md @@ -25,7 +25,7 @@ directly (see [Auth](#auth) below). ```rust let mut srt = moq_srt::Config::default(); srt.listen = Some("0.0.0.0:9000".parse()?); -srt.prefix = "live/".to_string(); +srt.prefix = "live".into(); // `origin` is your relay's local origin (e.g. `cluster.origin.clone()`). tokio::select! { @@ -86,8 +86,8 @@ while let Some(request) = server.accept().await { tokio::spawn(subscribe.accept(&consumer, "live/cam0")); } } - // ...or call `.reject()` on the `Publish` / `Subscribe` instead of `.accept()` - // to deny it. + // ...or call `.reject(moq_srt::Reject::Forbidden)` on the `Publish` / + // `Subscribe` instead of `.accept()` to deny it. } ``` diff --git a/rs/moq-srt/src/dial.rs b/rs/moq-srt/src/dial.rs index 9759f21c9e..2195f0bbb4 100644 --- a/rs/moq-srt/src/dial.rs +++ b/rs/moq-srt/src/dial.rs @@ -4,10 +4,10 @@ //! callers, this *dials* a remote `srt://host:port` as an SRT caller and bridges //! MPEG-TS in one of two directions, selected by the stream-id `m=` mode it sends: //! -//! - **[`publish`] (push / restream)**: call with `m=publish`, read a MoQ +//! - **[`Client::publish`] (push / restream)**: call with `m=publish`, read a MoQ //! broadcast from an origin, re-mux it to MPEG-TS with [`moq_mux`], and send it //! to the remote listener. This restreams MoQ out to a remote SRT ingest. -//! - **[`pull`] (ingest)**: call with `m=request`, receive the remote's +//! - **[`Client::pull`] (ingest)**: call with `m=request`, receive the remote's //! MPEG-TS, demux it with [`moq_mux`], and publish the result into an origin as //! an ordinary MoQ broadcast. This ingests a remote SRT source. //! @@ -26,13 +26,12 @@ use srt_tokio::SrtSocket; use crate::Result; use crate::server::{DEFAULT_LATENCY, configure_buffers, serve_publish, serve_subscribe}; -/// Where to dial and how, shared by [`publish`] and [`pull`]. +/// An SRT caller that can publish a MoQ broadcast or pull a remote stream. /// -/// Construct via [`Config::new`] and set the fields you need, so new options stay -/// additive. +/// Construct via [`Client::new`] and chain the `with_*` setters. #[derive(Debug, Clone)] #[non_exhaustive] -pub struct Config { +pub struct Client { /// The remote SRT listener to call. pub addr: SocketAddr, @@ -42,7 +41,7 @@ pub struct Config { /// SRT receive latency, negotiated at handshake time: the buffer that trades delay /// for loss recovery. It doubles as [`publish`]'s egress skip threshold. - pub latency: Duration, + latency: Duration, /// How long relays keep a non-latest group of an ingested media track fetchable, or /// `None` for hang's own default. @@ -53,14 +52,14 @@ pub struct Config { /// advertise segments that are still fetchable. Lower it when nothing reads history /// and the memory matters. [`pull`] only; [`publish`] reads a broadcast someone else /// declared. - pub max_age: Option, + max_age: Option, /// Connection allocator each ingested track claims its peak-hold bitrate on. /// [`pull`] only; [`publish`] reads a broadcast someone else declared. - pub bandwidth: moq_net::bandwidth::Allocator, + bandwidth: moq_net::bandwidth::Allocator, } -impl Config { +impl Client { /// Dial `addr` for `resource`, with the default SRT latency (500ms) and the /// publisher's own media retention. pub fn new(addr: SocketAddr, resource: impl Into) -> Self { @@ -72,56 +71,55 @@ impl Config { bandwidth: moq_net::bandwidth::Allocator::unlimited(), } } -} -/// Push a MoQ broadcast out to the remote: connect as an SRT caller requesting the -/// remote receive on [`Config::resource`] (`m=publish`), re-mux `path` from `origin` -/// to MPEG-TS, and send it until the broadcast ends. -/// -/// This future resolves when the broadcast ends, so callers usually run it on its own -/// task. -pub async fn publish(config: &Config, origin: &origin::Consumer, path: impl moq_net::AsPath) -> Result<()> { - let path = path.as_path(); - let socket = call(config, Mode::Publish).await?; - serve_subscribe(origin, path.as_str(), socket, config.latency).await -} + /// Override the SRT receive latency negotiated at handshake time. + pub fn with_latency(mut self, latency: Duration) -> Self { + self.latency = latency; + self + } -/// Pull a remote stream into `origin`: connect as an SRT caller requesting the remote -/// send on [`Config::resource`] (`m=request`), demux its MPEG-TS, and publish the -/// result at `path` until the remote ends. -/// -/// This future resolves when the remote stream ends, so callers usually run it on its -/// own task. -pub async fn pull(config: &Config, origin: &origin::Producer, path: impl moq_net::AsPath) -> Result<()> { - let path = path.as_path(); - let socket = call(config, Mode::Request).await?; - let catalog = moq_mux::catalog::Config::default() - .with_max_age(config.max_age) - .with_bandwidth(config.bandwidth.clone()); - serve_publish(origin, path.as_str(), socket, catalog).await -} + /// Set how long non-latest groups created by [`pull`](Self::pull) remain fetchable. + pub fn with_max_age(mut self, max_age: Option) -> Self { + self.max_age = max_age; + self + } -/// Dial as an SRT caller, sending the standard `#!::r=,m=` stream id -/// and returning the connected socket. -/// -/// `mode` is the *remote's* role, the inverse of the local direction (the remote -/// receives on `m=publish`, sends on `m=request`). -async fn call(config: &Config, mode: Mode) -> Result { - let Config { addr, resource, .. } = config; - // `,` and `=` delimit the `#!::r=,m=` stream id, so a resource - // carrying either would corrupt it and misroute at the listener. Reject rather - // than silently produce a broken id (MoQ paths never contain these). - if resource.contains([',', '=']) { - return Err(anyhow::anyhow!("srt resource must not contain ',' or '=': {resource:?}").into()); + /// Claim each track ingested by [`pull`](Self::pull) on `bandwidth`. + pub fn with_bandwidth(mut self, bandwidth: moq_net::bandwidth::Allocator) -> Self { + self.bandwidth = bandwidth; + self + } + + /// Push a MoQ broadcast out to the remote as MPEG-TS until the broadcast ends. + pub async fn publish(&self, origin: &origin::Consumer, path: impl moq_net::AsPath) -> Result<()> { + let path = path.as_path(); + let socket = self.call(Mode::Publish).await?; + serve_subscribe(origin, path.as_str(), socket, self.latency).await + } + + /// Pull a remote MPEG-TS stream into `origin` at `path` until the remote ends. + pub async fn pull(&self, origin: &origin::Producer, path: impl moq_net::AsPath) -> Result<()> { + let path = path.as_path(); + let socket = self.call(Mode::Request).await?; + let catalog = moq_mux::catalog::Config::default() + .with_max_age(self.max_age) + .with_bandwidth(self.bandwidth.clone()); + serve_publish(origin, path.as_str(), socket, catalog).await + } + + async fn call(&self, mode: Mode) -> Result { + if self.resource.contains([',', '=']) { + return Err(crate::Error::InvalidResource(self.resource.clone())); + } + let stream_id = format!("#!::r={},m={}", self.resource, mode.as_str()); + let socket = SrtSocket::builder() + .latency(self.latency) + .set(configure_buffers) + .call(self.addr, Some(&stream_id)) + .await?; + tracing::info!(addr = %self.addr, resource = %self.resource, mode = mode.as_str(), "SRT caller connected"); + Ok(socket) } - let stream_id = format!("#!::r={resource},m={}", mode.as_str()); - let socket = SrtSocket::builder() - .latency(config.latency) - .set(configure_buffers) - .call(*addr, Some(&stream_id)) - .await?; - tracing::info!(%addr, %resource, mode = mode.as_str(), "SRT caller connected"); - Ok(socket) } /// The SRT stream-id `m=` mode sent to the remote, i.e. the remote's role. @@ -157,11 +155,15 @@ mod tests { producer } + use std::io; use std::net::SocketAddr; use std::time::Duration; + use srt_protocol::protocol::pending_connection::ConnectionReject; + use srt_tokio::access::RejectReason; + use super::*; - use crate::server::{Request, Server}; + use crate::server::{Reject, Request, Server}; /// Grab a free UDP port by binding `:0` and releasing it. Racy in principle, but /// the window before the SRT server rebinds it is tiny; good enough for a test. @@ -195,7 +197,7 @@ mod tests { }); // Caller: dial with m=publish, then drop (we only assert connect + routing). - let caller = tokio::spawn(async move { call(&Config::new(addr, "cam0"), Mode::Publish).await }); + let caller = tokio::spawn(async move { Client::new(addr, "cam0").call(Mode::Publish).await }); let socket = tokio::time::timeout(Duration::from_secs(10), caller) .await @@ -234,7 +236,7 @@ mod tests { (resource, is_subscribe) }); - let caller = tokio::spawn(async move { call(&Config::new(addr, "cam0"), Mode::Request).await }); + let caller = tokio::spawn(async move { Client::new(addr, "cam0").call(Mode::Request).await }); let socket = tokio::time::timeout(Duration::from_secs(10), caller) .await @@ -250,4 +252,40 @@ mod tests { assert_eq!(resource, "cam0"); assert!(is_subscribe, "m=request should route to a server Subscribe request"); } + + async fn rejected(mode: Mode, reason: Reject, code: i32) { + let addr = free_udp_addr().await; + let mut server = Server::bind(addr, None).await.unwrap(); + let server_task = tokio::spawn(async move { + match (mode, server.accept().await.expect("a request")) { + (Mode::Publish, Request::Publish(request)) => request.reject(reason).await.unwrap(), + (Mode::Request, Request::Subscribe(request)) => request.reject(reason).await.unwrap(), + _ => panic!("request routed in the wrong direction"), + } + }); + + let err = match Client::new(addr, "cam0").call(mode).await { + Ok(_) => panic!("rejected SRT caller connected"), + Err(err) => err, + }; + server_task.await.unwrap(); + let crate::Error::Io(err) = err else { + panic!("SRT rejection was not an I/O error: {err}"); + }; + assert_eq!(err.kind(), io::ErrorKind::ConnectionRefused); + assert_eq!( + err.get_ref().and_then(|err| err.downcast_ref::()), + Some(&ConnectionReject::Rejected(RejectReason::CoreUnrecognized(code))) + ); + } + + #[tokio::test] + async fn publish_rejection_carries_unauthorized_code() { + rejected(Mode::Publish, Reject::Unauthorized, 1401).await; + } + + #[tokio::test] + async fn subscribe_rejection_carries_unavailable_code() { + rejected(Mode::Request, Reject::Unavailable, 1503).await; + } } diff --git a/rs/moq-srt/src/error.rs b/rs/moq-srt/src/error.rs index 33ccd283e2..98c7f7ddb2 100644 --- a/rs/moq-srt/src/error.rs +++ b/rs/moq-srt/src/error.rs @@ -18,10 +18,13 @@ pub enum Error { #[error("io: {0}")] Io(Arc), - /// Catch-all for ingest logic that reports via `anyhow` (the moq-mux - /// demuxer surfaces its errors this way). - #[error("{0}")] - Other(Arc), + /// The remote resource contains delimiters reserved by the SRT stream-id syntax. + #[error("invalid SRT resource: {0}")] + InvalidResource(String), + + /// The listener stopped accepting connections. + #[error("SRT listener stopped accepting connections")] + ListenerClosed, } impl From for Error { @@ -30,11 +33,5 @@ impl From for Error { } } -impl From for Error { - fn from(err: anyhow::Error) -> Self { - Error::Other(Arc::new(err)) - } -} - /// Result alias for the SRT ingest gateway. pub type Result = std::result::Result; diff --git a/rs/moq-srt/src/lib.rs b/rs/moq-srt/src/lib.rs index a8ba498115..ffa1351086 100644 --- a/rs/moq-srt/src/lib.rs +++ b/rs/moq-srt/src/lib.rs @@ -24,9 +24,9 @@ //! JWT and scoping the origin per token) plugs its policy in. It mirrors //! `moq-tokio`'s `Server` / `Request`. //! -//! Beyond the listener, the [`dial`] module is the *dial-out* (client) role: build a -//! [`dial::Config`] naming a remote SRT listener and either [`dial::publish`] a MoQ -//! broadcast to it (restream MoQ out to a remote SRT ingest) or [`dial::pull`] a +//! Beyond the listener, [`Client`] is the *dial-out* role: name a remote SRT +//! listener and either [`Client::publish`] a MoQ broadcast to it (restream MoQ +//! out to a remote SRT ingest) or [`Client::pull`] a //! remote stream into an origin (ingest a remote SRT source). It reuses the same //! MPEG-TS <-> moq bridge; only the SRT caller transport is new. //! @@ -44,6 +44,7 @@ mod listen; mod server; mod ts; +pub use dial::Client; pub use error::{Error, Result}; pub use listen::{Config, run}; -pub use server::{Publish, Request, Server, Subscribe}; +pub use server::{Publish, Reject, Request, Server, Subscribe}; diff --git a/rs/moq-srt/src/listen.rs b/rs/moq-srt/src/listen.rs index d0651b93db..6b5f3d84df 100644 --- a/rs/moq-srt/src/listen.rs +++ b/rs/moq-srt/src/listen.rs @@ -29,7 +29,7 @@ use std::net::SocketAddr; use std::sync::{Arc, Mutex}; use std::time::Duration; -use moq_net::origin; +use moq_net::{Path, PathOwned, origin}; use crate::Result; use crate::server::{Request, Server}; @@ -47,9 +47,9 @@ pub struct Config { /// gateway is disabled. pub listen: Option, - /// Prefix prepended to every broadcast path, for both publish and request. - /// Lets one listener namespace all of its streams (e.g. `live/`). - pub prefix: String, + /// Path prefix prepended to every broadcast path, for both publish and request. + /// Lets one listener namespace all of its streams (e.g. `live`). + pub prefix: PathOwned, /// SRT receive latency: the negotiated buffer that trades delay for loss /// recovery. @@ -70,7 +70,7 @@ impl Default for Config { fn default() -> Self { Self { listen: None, - prefix: String::new(), + prefix: Path::empty().to_owned(), latency: crate::server::DEFAULT_LATENCY, max_age: None, } @@ -108,7 +108,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { // RTMP stream key) instead of being silently parked as a backup that could // take over the path when the first publisher drops. let active = ActivePaths::default(); - let prefix = Arc::new(config.prefix); + let prefix = config.prefix; let max_age = config.max_age; while let Some(request) = server.accept().await { @@ -122,12 +122,12 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { // publishers. tokio::spawn(async move { let peer = publish.peer(); - let path = format!("{prefix}{}", publish.resource()); + let path = prefix.join(publish.resource()); // Claim the path before accepting; the guard releases it when the // connection task ends (success, error, or panic). - let Some(_guard) = active.claim(&path) else { + let Some(_guard) = active.claim(path.as_str()) else { tracing::warn!(%peer, %path, "rejecting SRT publish: path already being ingested"); - let _ = publish.reject().await; + let _ = publish.reject(crate::Reject::Unavailable).await; return; }; if let Err(err) = publish.with_max_age(max_age).accept(&origin, &path).await { @@ -143,7 +143,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { // don't claim an `ActivePaths` slot. tokio::spawn(async move { let peer = subscribe.peer(); - let path = format!("{prefix}{}", subscribe.resource()); + let path = prefix.join(subscribe.resource()); if let Err(err) = subscribe.accept(&consumer, &path).await { tracing::warn!(%peer, %path, %err, "SRT request ended with error"); } else { @@ -154,9 +154,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { } } - Err(crate::Error::from(anyhow::anyhow!( - "SRT listener stopped accepting connections" - ))) + Err(crate::Error::ListenerClosed) } /// The set of broadcast paths with a live ingest, used to reject duplicate diff --git a/rs/moq-srt/src/server.rs b/rs/moq-srt/src/server.rs index 83e34b8f59..4d5bb556e6 100644 --- a/rs/moq-srt/src/server.rs +++ b/rs/moq-srt/src/server.rs @@ -25,14 +25,38 @@ use std::time::{Duration, Instant}; use futures::{SinkExt, StreamExt}; use moq_mux::container::Frame; use moq_net::origin; -use srt_tokio::access::{ - AccessControlList, ConnectionMode, RejectReason, ServerRejectReason, StandardAccessControlEntry, -}; +use srt_tokio::access::{AccessControlList, ConnectionMode, RejectReason, StandardAccessControlEntry}; use srt_tokio::options::{PacketCount, SocketOptions, StreamId}; use srt_tokio::{ConnectionRequest, SrtIncoming, SrtListener, SrtSocket}; use crate::Result; +/// Why an SRT publish or subscribe was refused. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[non_exhaustive] +pub enum Reject { + /// Authentication failed. + Unauthorized, + /// The authenticated caller is not allowed to access the resource. + Forbidden, + /// The service or resource is temporarily unavailable. + Unavailable, + /// The request or stream-id is malformed. + BadRequest, +} + +impl Reject { + fn reason(self) -> RejectReason { + let code = match self { + Self::Unauthorized => 1401, + Self::Forbidden => 1403, + Self::Unavailable => 1503, + Self::BadRequest => 1400, + }; + RejectReason::CoreUnrecognized(code) + } +} + /// Default SRT receive latency: the negotiated buffer that trades delay for loss /// recovery. Override per-server with [`Server::bind`]'s `latency` argument. pub(crate) const DEFAULT_LATENCY: Duration = Duration::from_millis(500); @@ -194,7 +218,7 @@ impl Server { let peer = request.remote(); let Some((resource, mode)) = parse_stream_id(request.stream_id()) else { tracing::warn!(%peer, stream_id = ?request.stream_id(), "rejecting SRT: no usable stream id"); - reject_log(request, ServerRejectReason::BadRequest, peer).await; + reject_log(request, Reject::BadRequest, peer).await; continue; }; @@ -354,13 +378,9 @@ impl Publish { serve_publish(origin, path.as_str(), socket, config).await } - /// Reject the publish, sending the client a `Forbidden` rejection. - pub async fn reject(self) -> Result<()> { - Ok(self - .0 - .request - .reject(RejectReason::Server(ServerRejectReason::Forbidden)) - .await?) + /// Reject the publish with a verdict the client can distinguish on the wire. + pub async fn reject(self, reason: Reject) -> Result<()> { + Ok(self.0.request.reject(reason.reason()).await?) } } @@ -405,20 +425,16 @@ impl Subscribe { serve_subscribe(origin, path.as_str(), socket, self.0.latency).await } - /// Reject the subscribe, sending the client a `Forbidden` rejection. - pub async fn reject(self) -> Result<()> { - Ok(self - .0 - .request - .reject(RejectReason::Server(ServerRejectReason::Forbidden)) - .await?) + /// Reject the subscribe with a verdict the client can distinguish on the wire. + pub async fn reject(self, reason: Reject) -> Result<()> { + Ok(self.0.request.reject(reason.reason()).await?) } } /// Reject a connection request, logging (but not propagating) a send failure. /// Used for connections the server drops itself, before they reach the caller. -async fn reject_log(request: ConnectionRequest, reason: ServerRejectReason, peer: SocketAddr) { - if let Err(err) = request.reject(RejectReason::Server(reason)).await { +async fn reject_log(request: ConnectionRequest, reason: Reject, peer: SocketAddr) { + if let Err(err) = request.reject(reason.reason()).await { tracing::debug!(%peer, %err, "failed to send SRT rejection"); } } diff --git a/rs/moq-srt/src/ts.rs b/rs/moq-srt/src/ts.rs index 9646e423e8..ce8a380fa3 100644 --- a/rs/moq-srt/src/ts.rs +++ b/rs/moq-srt/src/ts.rs @@ -58,13 +58,13 @@ impl Publisher { /// `decode` drains `data` fully, buffering any partial trailing packet in /// its own internal scratch, so there's nothing to retain here. pub fn feed(&mut self, data: Bytes) -> Result<()> { - Ok(self.importer.decode(&data)?) + Ok(self.importer.decode(&data).map_err(moq_mux::Error::from)?) } /// Flush any buffered media, close out the broadcast's open groups, and end /// the broadcast so the origin unannounces it immediately. pub fn finish(&mut self) -> Result<()> { - self.importer.finish()?; + self.importer.finish().map_err(moq_mux::Error::from)?; self.broadcast.finish(); Ok(()) } diff --git a/rs/moq-stats/src/aggregate.rs b/rs/moq-stats/src/aggregate.rs index bcd6009f34..d178f0cd35 100644 --- a/rs/moq-stats/src/aggregate.rs +++ b/rs/moq-stats/src/aggregate.rs @@ -21,7 +21,7 @@ use crate::{Result, SessionsFrame, TrafficFrame, parse_node_path, sessions_track /// the `with_*` setters. /// /// The `prefix` and `depth` must match the producing side's -/// [`ProducerConfig`](crate::ProducerConfig): they are how announced paths are +/// [`produce::Config`](crate::produce::Config): they are how announced paths are /// recognized as node broadcasts and filtered from sibling categories under the /// same prefix. #[derive(Debug, Clone)] @@ -493,7 +493,7 @@ mod tests { use moq_net::{PathOwned, Timestamp, announce, broadcast, origin, track}; - use crate::{Producer, ProducerConfig}; + use crate::{Producer, produce}; use super::*; @@ -502,7 +502,7 @@ mod tests { /// `.stats//node/`). fn node_producer(origin: &origin::Producer, node: &str) -> Producer { Producer::new( - ProducerConfig::new() + produce::Config::new() .with_origin(origin.clone()) .with_node(PathOwned::from(node.to_string())) .with_depth(1), diff --git a/rs/moq-stats/src/consume.rs b/rs/moq-stats/src/consume.rs index 6703ad582c..f621373224 100644 --- a/rs/moq-stats/src/consume.rs +++ b/rs/moq-stats/src/consume.rs @@ -5,18 +5,18 @@ use moq_net::stats::{Role, Tier}; use crate::{Result, SessionsFrame, TrafficFrame, sessions_track, traffic_track}; -/// Configuration for a [`Consumer`]. Construct with [`ConsumerConfig::new`] +/// Configuration for a [`Consumer`]. Construct with [`Config::new`] /// and chain the `with_*` setters. #[derive(Debug, Clone, Default)] #[non_exhaustive] -pub struct ConsumerConfig { +pub struct Config { /// Read the compressed `.json.z` tracks instead of the plain `.json` ones. /// Same data for a fraction of the bytes, but requires a producer that /// publishes them. Defaults to `false`. pub compression: bool, } -impl ConsumerConfig { +impl Config { /// A config with default settings: the plain `.json` tracks. pub fn new() -> Self { Self::default() @@ -38,30 +38,30 @@ impl ConsumerConfig { /// immediately, so callers typically subscribe the tiers they know exist. pub struct Consumer { broadcast: broadcast::Consumer, - config: ConsumerConfig, + config: Config, } impl Consumer { /// Wrap a stats broadcast. The broadcast is whatever the announce at a /// stats path resolved to; parse the path with [`crate::parse_node_path`]. - pub fn new(broadcast: broadcast::Consumer, config: ConsumerConfig) -> Self { + pub fn new(broadcast: broadcast::Consumer, config: Config) -> Self { Self { broadcast, config } } /// Subscribe to the traffic track for `(tier, role)`, awaiting the /// subscription handshake. - pub async fn traffic(&self, tier: &Tier, role: Role) -> Result { + pub async fn traffic(&self, tier: &Tier, role: Role) -> Result { let name = traffic_track(tier, role, self.config.compression); - Ok(TrafficConsumer { + Ok(Traffic { inner: self.subscribe(&name).await?, }) } /// Subscribe to the sessions track for `tier`, awaiting the subscription /// handshake. - pub async fn sessions(&self, tier: &Tier) -> Result { + pub async fn sessions(&self, tier: &Tier) -> Result { let name = sessions_track(tier, self.config.compression); - Ok(SessionsConsumer { + Ok(Sessions { inner: self.subscribe(&name).await?, }) } @@ -79,23 +79,23 @@ impl Consumer { /// A typed reader over one traffic track. Yields the latest [`TrafficFrame`]; /// intermediate frames a slow reader missed are collapsed, which is safe /// because the counters are cumulative. -pub struct TrafficConsumer { +pub struct Traffic { inner: moq_json::snapshot::Consumer, } -impl TrafficConsumer { +impl Traffic { /// The next frame, or `None` once the track ends (the producer went away). pub async fn next(&mut self) -> Result> { Ok(self.inner.next().await?) } } -/// A typed reader over one sessions track; see [`TrafficConsumer`]. -pub struct SessionsConsumer { +/// A typed reader over one sessions track; see [`Traffic`]. +pub struct Sessions { inner: moq_json::snapshot::Consumer, } -impl SessionsConsumer { +impl Sessions { /// The next frame, or `None` once the track ends (the producer went away). pub async fn next(&mut self) -> Result> { Ok(self.inner.next().await?) @@ -121,14 +121,14 @@ mod tests { use moq_net::{Consume, PathOwned, Timestamp, announce, broadcast, origin, track}; - use crate::{Producer, ProducerConfig, Tier}; + use crate::{Producer, Tier, produce}; use super::*; fn test_producer() -> (Producer, origin::Producer) { let origin = produce_origin(); let producer = Producer::new( - ProducerConfig::new() + produce::Config::new() .with_origin(origin.clone()) .with_node(PathOwned::from("sjc")), ); @@ -213,8 +213,8 @@ mod tests { drive_tick().await; let broadcast = announced(&origin).await; - let plain = Consumer::new(broadcast.consume(), ConsumerConfig::new()); - let compressed = Consumer::new(broadcast.consume(), ConsumerConfig::new().with_compression(true)); + let plain = Consumer::new(broadcast.consume(), Config::new()); + let compressed = Consumer::new(broadcast.consume(), Config::new().with_compression(true)); let mut plain_traffic = plain.traffic(&tier, Role::Publisher).await.expect("subscribe plain"); let mut z_traffic = compressed diff --git a/rs/moq-stats/src/lib.rs b/rs/moq-stats/src/lib.rs index f744b41592..d1613d7217 100644 --- a/rs/moq-stats/src/lib.rs +++ b/rs/moq-stats/src/lib.rs @@ -52,11 +52,11 @@ //! should treat a decrease as a fresh segment. pub mod aggregate; -mod consume; -mod produce; +pub mod consume; +pub mod produce; -pub use consume::{Consumer, ConsumerConfig, SessionsConsumer, TrafficConsumer}; -pub use produce::{Producer, ProducerConfig}; +pub use consume::Consumer; +pub use produce::Producer; use std::collections::BTreeMap; diff --git a/rs/moq-stats/src/produce.rs b/rs/moq-stats/src/produce.rs index 02c1a46ac7..bfce26ccfe 100644 --- a/rs/moq-stats/src/produce.rs +++ b/rs/moq-stats/src/produce.rs @@ -13,17 +13,17 @@ use web_async::spawn; use crate::{COMPRESSED_SUFFIX, SessionsFrame, TrafficFrame, sessions_track, traffic_track}; -/// Settings for a [`Producer`]. Construct with [`ProducerConfig::new`] and chain +/// Settings for a [`Producer`]. Construct with [`Config::new`] and chain /// the `with_*` setters (e.g. -/// `ProducerConfig::new().with_origin(origin).with_prefix(".foo")`), then hand it +/// `Config::new().with_origin(origin).with_prefix(".foo")`), then hand it /// to [`Producer::new`]. /// /// With no origin set the resulting producer is a no-op: its registry is /// disabled (bumps are dropped) and no task spawns. Call -/// [`ProducerConfig::with_origin`] to publish. +/// [`Config::with_origin`] to publish. #[derive(Clone)] #[non_exhaustive] -pub struct ProducerConfig { +pub struct Config { /// Origin the stats broadcasts are created on. /// When `None`, [`Producer::new`] spawns no task and publishes nothing. pub origin: Option, @@ -50,7 +50,7 @@ pub struct ProducerConfig { pub depth: usize, } -impl ProducerConfig { +impl Config { /// A config with default settings: no origin (no-op), `.stats` prefix, 1s /// interval, and no node suffix. Call [`Self::with_origin`] to actually /// publish. @@ -96,7 +96,7 @@ impl ProducerConfig { } } -impl Default for ProducerConfig { +impl Default for Config { fn default() -> Self { Self::new() } @@ -148,8 +148,8 @@ impl Producer { /// [`Producer`] clone is dropped. With no origin the producer is a no-op /// (its registry is disabled, nothing is published) and no task spawns, so /// it's safe to build outside an async runtime. - pub fn new(config: ProducerConfig) -> Self { - let ProducerConfig { + pub fn new(config: Config) -> Self { + let Config { origin, prefix, node, @@ -902,7 +902,7 @@ mod tests { fn test_producer(node: Option<&str>) -> (Producer, origin::Producer) { let origin = produce_origin(); let producer = Producer::new( - ProducerConfig::new() + Config::new() .with_origin(origin.clone()) .with_node(node.map(|s| PathOwned::from(s.to_string()))), );