diff --git a/doc/bin/cli.md b/doc/bin/cli.md index 9c5d484e05..65c66fa6a8 100644 --- a/doc/bin/cli.md +++ b/doc/bin/cli.md @@ -86,6 +86,7 @@ a PAT that adds a program mid-stream ends the import the same way. the first PAT lists as its own broadcast, with its own clock and catalog, keeping the catalog suffix last: `--broadcast event.hang` publishes `event/1.hang`, `event/2.hang`, and so on. `export ts` writes one program per broadcast. +`import srt` takes the same `--program`. ```bash moq --connect https://relay.example.com/anon --broadcast event.hang import ts --program all < mux.ts diff --git a/doc/bin/srt.md b/doc/bin/srt.md index 141ea5cbb4..a2abd327ab 100644 --- a/doc/bin/srt.md +++ b/doc/bin/srt.md @@ -25,6 +25,15 @@ ffplay srt://localhost:9000 moq --connect https://relay.example.com/anon --broadcast event.hang import srt --connect 'srt://encoder.example.com:9000?streamid=live/cam' ``` +A multi-program feed is refused, as with `import ts`, unless `--program` +picks one: `--program 2` imports program 2 alone, and `--program all` +publishes each program as its own broadcast (`event.hang` becomes +`event/1.hang`, `event/2.hang`, and so on). + +```bash +moq --connect https://relay.example.com/anon --broadcast event.hang import srt --listen '[::]:9000' --program all +``` + `--latency` sets the SRT receive buffer and doubles as the skip threshold on export. Export paces each SRT payload on the media clock, and re-anchors that pacing on a declared marker, so a restarted timeline plays out from the diff --git a/quest/m1/ts-psi-reassembly.md b/quest/m1/ts-psi-reassembly.md index 665a76a8da..d5fa4556ed 100644 --- a/quest/m1/ts-psi-reassembly.md +++ b/quest/m1/ts-psi-reassembly.md @@ -7,7 +7,7 @@ starts after a nonzero `pointer_field`, instead of aborting the import. Today the `mpeg2ts` 0.6.1 reader parses PSI from a single packet and rejects a nonzero `pointer_field`, and its error ends the whole import, so a valid long PMT (many audio languages, long descriptors) or a PAT with more than about 40 programs -cannot be imported. `ts::programs()` finds such a PAT too. +cannot be imported. `ts::Programs` finds such a PAT too. ## Plan @@ -53,13 +53,13 @@ Settled decisions: is `#[non_exhaustive]`, so the new field is additive; its docs (today per elementary stream) and `is_empty` widen to cover it, and `ts::stats::Log` reports it. -- `ts::programs()` reads through the same PAT path, so a PAT spanning packets +- `ts::Programs` reads through the same PAT path, so a PAT spanning packets is found before any program publishes. - One quest, because the demux refactor alone changes nothing observable. Keep every existing TS import test and fixture passing unchanged. Add tests for a PMT spanning two packets, a PAT behind a nonzero `pointer_field`, a -multi-packet PAT with enough programs to need it (read by `ts::programs()` and +multi-packet PAT with enough programs to need it (read by `ts::Programs` and by `with_program`), and a PAT and a PMT with a corrupt CRC, each between good repetitions, that is dropped and counted once while the import keeps its previous layout and a later good PMT revision still applies. A positive diff --git a/rs/moq-cli/src/args.rs b/rs/moq-cli/src/args.rs index b2c659e1eb..733a4eaff1 100644 --- a/rs/moq-cli/src/args.rs +++ b/rs/moq-cli/src/args.rs @@ -747,7 +747,7 @@ pub enum ImportSource { /// RTMP: pull a remote play (`--connect`) or accept incoming publishes (`--listen`). Rtmp(crate::rtmp::Args), /// SRT: pull a remote stream (`--connect`) or accept incoming publishes (`--listen`). - Srt(crate::srt::Args), + Srt(crate::srt::ImportArgs), /// WebRTC: WHEP client pulling a remote (`--connect`) or WHIP server accepting publishes (`--listen`). Rtc(crate::rtc::Args), /// Capture a local source (camera, display, window, app, microphone) and @@ -782,7 +782,7 @@ pub struct TsImport { pub program: Option, } -/// An `import ts --program` value. +/// An `import ts --program` or `import srt --program` value. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum TsProgram { /// The program with this PAT program number. @@ -1058,6 +1058,25 @@ mod tests { assert_eq!(program("two"), None); } + /// `import srt` takes the same `--program` as `import ts`; `export srt` has no program to pick. + #[test] + fn import_srt_takes_a_program() { + let cli = + Invocation::try_parse_from(["moq", "import", "srt", "--listen", "[::]:9000", "--program", "all"]).unwrap(); + let Command::Import(import) = &cli.stages[0] else { + panic!("an import stage"); + }; + let ImportSource::Srt(args) = &import.source else { + panic!("an import srt stage"); + }; + assert_eq!(args.program(), Some(moq_srt::Program::All)); + assert!(args.endpoint.listen.is_some()); + + assert!( + Invocation::try_parse_from(["moq", "export", "srt", "--listen", "[::]:9000", "--program", "2"]).is_err() + ); + } + /// A released spelling is refused, and the error names what to write instead. /// /// The alternative is what this replaced: the flag parsed onto a hidden field, diff --git a/rs/moq-cli/src/main.rs b/rs/moq-cli/src/main.rs index d339abe70c..fd1bb72bf2 100644 --- a/rs/moq-cli/src/main.rs +++ b/rs/moq-cli/src/main.rs @@ -661,11 +661,13 @@ fn spawn_import( } } ImportSource::Srt(srt) => { + let program = srt.program(); + let srt = srt.endpoint; if let Some(addr) = srt.listen { let name = require_broadcast(name, "import srt --listen")?; - tasks.spawn(srt::listen_import(target(name), addr, srt.latency.into_std())); + tasks.spawn(srt::listen_import(target(name), addr, srt.latency.into_std(), program)); } else if let Some(url) = srt.connect { - tasks.spawn(srt::connect_import(target(name), url, srt.latency.into_std())); + tasks.spawn(srt::connect_import(target(name), url, srt.latency.into_std(), program)); } } ImportSource::Rtc(rtc) => { diff --git a/rs/moq-cli/src/publish.rs b/rs/moq-cli/src/publish.rs index 5bf7c08916..3dc1adf963 100644 --- a/rs/moq-cli/src/publish.rs +++ b/rs/moq-cli/src/publish.rs @@ -142,7 +142,7 @@ enum PublishDecoder { // verbatim, so it uses the `mpegts` catalog extension rather than the media-only `()`. Ts(Box>), /// `import ts --program all`: one importer, broadcast, and catalog per program. - TsPrograms(Box), + TsPrograms(Box), Flv(Box), } @@ -218,7 +218,7 @@ fn suggest_program(err: anyhow::Error) -> anyhow::Error { enum PublishCatalog { Media(moq_mux::catalog::Producer), Ts(moq_mux::catalog::Producer), - /// `import ts --program all`: each program's catalog, finished by [`TsPrograms::finish`]. + /// `import ts --program all`: each program's catalog, finished by [`ts::Programs::finish`]. TsPrograms, } @@ -253,117 +253,6 @@ fn ts_import( Ok((import, catalog)) } -/// One program's broadcast under `import ts --program all`. -struct ProgramImport { - /// Keeps the origin-created broadcast alive; the importer's clone doesn't. - _broadcast: moq_net::broadcast::Producer, - import: ts::Import, - catalog: moq_mux::catalog::Producer, -} - -/// `import ts --program all`: every program the first PAT lists, each published as its own -/// broadcast on its own clock and catalog. -/// -/// Input ahead of that PAT is dropped, as a single importer drops media ahead of its PSI. A -/// program a later PAT adds is not published; one it removes ends the import. -struct TsPrograms { - origin: moq_net::origin::Producer, - name: String, - config: moq_mux::catalog::Config, - /// Input held until it contains a whole PAT. - pending: Vec, - /// Empty until the PAT arrives. - programs: Vec, -} - -impl TsPrograms { - fn new(origin: moq_net::origin::Producer, name: String, config: moq_mux::catalog::Config) -> Self { - Self { - origin, - name, - config, - pending: Vec::new(), - programs: Vec::new(), - } - } - - fn decode(&mut self, chunk: &[u8]) -> anyhow::Result<()> { - if !self.programs.is_empty() { - return self.feed(chunk); - } - self.pending.extend_from_slice(chunk); - let Some(programs) = ts::programs(&self.pending).filter(|programs| !programs.is_empty()) else { - // Only a packet's worth of tail can still hold the start of the PAT. - let keep = self.pending.len().saturating_sub(TS_PACKET_SIZE - 1); - self.pending.drain(..keep); - return Ok(()); - }; - for program in programs { - let name = program_broadcast(&self.name, program); - let mut broadcast = self - .origin - .create_broadcast(&name) - .with_context(|| format!("failed to create broadcast {name}"))?; - let (import, catalog) = ts_import(&mut broadcast, self.config.clone(), Some(program))?; - broadcast - .announce(Default::default()) - .with_context(|| format!("failed to announce broadcast {name}"))?; - self.programs.push(ProgramImport { - _broadcast: broadcast, - import, - catalog, - }); - } - let pending = std::mem::take(&mut self.pending); - self.feed(&pending) - } - - fn feed(&mut self, chunk: &[u8]) -> anyhow::Result<()> { - for program in &mut self.programs { - program.import.decode(chunk)?; - } - Ok(()) - } - - /// Every program's counters in one map: PIDs are unique across a multiplex. - fn stats(&self) -> ts::Stats { - let mut stats = ts::Stats::default(); - for program in &self.programs { - stats.streams.extend(program.import.stats().streams); - } - stats - } - - /// Finish each program's tracks, then its catalog while it still lists them. - fn finish(&mut self) -> anyhow::Result<()> { - for program in &mut self.programs { - program.import.finish()?; - program.catalog.finish()?; - } - Ok(()) - } - - fn abort(self, err: moq_net::Error) { - for program in self.programs { - program.import.abort(err.clone()); - } - } -} - -const TS_PACKET_SIZE: usize = 188; - -/// The broadcast one program of `name` publishes on, keeping the catalog suffix last so format -/// detection still sees it: `event.hang` becomes `event/2.hang`. -fn program_broadcast(name: &str, program: u16) -> String { - match moq_mux::catalog::CatalogFormat::detect(name) { - Some(format) => { - let stem = &name[..name.len() - format.extension().len()]; - format!("{stem}/{program}{}", format.extension()) - } - None => format!("{name}/{program}"), - } -} - // Exactly one Source exists per process, so the size gap between the small // Stream variant and the larger Capture config is irrelevant. #[allow(clippy::large_enum_variant)] @@ -460,7 +349,7 @@ impl Publish { pub fn ts_programs(origin: moq_net::origin::Producer, name: String, config: moq_mux::catalog::Config) -> Self { Self { source: Source::Stream { - decoder: PublishDecoder::TsPrograms(Box::new(TsPrograms::new(origin, name, config))), + decoder: PublishDecoder::TsPrograms(Box::new(ts::Programs::new(origin, name, config).live())), catalog: PublishCatalog::TsPrograms, }, broadcast: None, @@ -1031,13 +920,6 @@ mod tests { assert_eq!(last.audio.renditions.len(), 1, "the audio rendition is still listed"); } - #[test] - fn program_broadcasts_keep_the_catalog_suffix_last() { - assert_eq!(program_broadcast("event.hang", 2), "event/2.hang"); - assert_eq!(program_broadcast("demo/event.msf", 7), "demo/event/7.msf"); - assert_eq!(program_broadcast("event", 1), "event/1"); - } - /// A PAT listing two programs, then one MP2 PES of each: program 1 on PID `0x61` at 1 s /// with fill bytes `0xAA`/`0xBB`, program 2 on PID `0x71` an hour later with `0xCC`/`0xDD`. fn two_programs() -> Vec { @@ -1145,61 +1027,6 @@ mod tests { assert!(err.contains("--program") && err.contains("programs (1, 2)"), "{err}"); } - /// `import ts --program all` holds the input until the PAT is whole, then publishes each - /// program as its own broadcast carrying only that program, live on its own first frame - /// rather than an hour apart on one shared clock. - #[tokio::test] - async fn every_program_publishes_its_own_broadcast() { - let origin = moq_tokio::origin::spawn(); - let ago = Duration::from_secs(60); - let clock = moq_mux::Clock::at(std::time::Instant::now() - ago, std::time::SystemTime::now() - ago).unwrap(); - let config = moq_mux::catalog::Config::default().with_clock(clock); - let publish = Publish::ts_programs(origin.clone(), "event.hang".to_string(), config); - let Source::Stream { - decoder: PublishDecoder::TsPrograms(mut programs), - .. - } = publish.source - else { - panic!("expected the per-program decoder"); - }; - - let input = two_programs(); - let before = clock.now(); - programs.decode(&input[..100]).unwrap(); - assert!( - programs.programs.is_empty(), - "no program is known before the PAT is whole" - ); - programs.decode(&input[100..]).unwrap(); - programs.finish().unwrap(); - let after = clock.now(); - - for (path, fills) in [("event/1.hang", [0xAA, 0xBB]), ("event/2.hang", [0xCC, 0xDD])] { - let consumer = moq_mux::Source::new(origin.consume(), path).broadcast().await.unwrap(); - let catalog = hang::catalog::Catalog::<()>::subscribe(&consumer) - .await - .unwrap() - .next() - .await - .unwrap() - .expect("a catalog"); - let renditions: Vec<_> = catalog.audio.renditions.iter().collect(); - assert_eq!(renditions.len(), 1, "{path} carries its own program's one stream"); - let (name, config) = renditions[0]; - let track = consumer.track(name).unwrap().subscribe(None).await.unwrap(); - let container = Container::try_from(config).unwrap(); - let frame = Consumer::new(track, container).read().await.unwrap().expect("a frame"); - assert!(fills.contains(&frame.payload[4]), "{path} carries only its own program"); - let skew = Duration::from_secs(2).as_micros(); - assert!( - before.as_micros() - skew <= frame.timestamp.as_micros() - && frame.timestamp.as_micros() <= after.as_micros() + skew, - "{path} is live on arrival: {:?} not in {before:?}..={after:?}", - frame.timestamp - ); - } - } - /// Read the first frame of a verbatim track back as raw bytes. async fn read_frame(consumer: &moq_net::broadcast::Consumer, name: &str) -> Vec { let track = consumer.track(name).unwrap().subscribe(None).await.unwrap(); diff --git a/rs/moq-cli/src/srt.rs b/rs/moq-cli/src/srt.rs index 7afbe3ff7f..35993e691d 100644 --- a/rs/moq-cli/src/srt.rs +++ b/rs/moq-cli/src/srt.rs @@ -10,6 +10,7 @@ use moq_srt::{Reject, Request, Server}; use moq_tokio::RedactedUrl; use url::Url; +use crate::args::TsProgram; use crate::moq::{ImportTarget, notify_ready}; /// SRT endpoint args: exactly one of `--connect` (dial) / `--listen` (bind). @@ -31,8 +32,49 @@ pub struct Args { pub latency: crate::duration::Duration, } +/// SRT import args: the endpoint, plus which programs of a multiplex to publish. +#[derive(usage::Args, Clone)] +#[usage(unknown_flags = "error", args_override_self = false)] +pub struct ImportArgs { + #[usage(flatten)] + pub endpoint: Args, + + /// Import one program of a multi-program feed, by its PAT program number, or `all` to + /// publish each program as its own broadcast (`event.hang` becomes `event/1.hang`, + /// `event/2.hang`, ...). Without it, a feed carrying more than one program is refused. + #[usage(long)] + pub program: Option, +} + +impl ImportArgs { + /// The library's selection for `--program`. + pub fn program(&self) -> Option { + self.program.map(|program| match program { + TsProgram::One(program) => moq_srt::Program::One(program), + TsProgram::All => moq_srt::Program::All, + }) + } +} + +/// Point a multi-program refusal at the flag that resolves it. +fn suggest_program(err: moq_srt::Error) -> anyhow::Error { + let multiplex = matches!(&err, moq_srt::Error::Mux(moq_mux::Error::Other(inner)) + if inner.is::()); + let err = anyhow::Error::from(err); + if multiplex { + err.context("choose one with `--program `, or publish each with `--program all`") + } else { + err + } +} + /// Accept incoming SRT publishes into the Origin as `target.name`; reject requests (import). -pub async fn listen_import(target: ImportTarget, addr: SocketAddr, latency: Duration) -> anyhow::Result<()> { +pub async fn listen_import( + target: ImportTarget, + addr: SocketAddr, + latency: Duration, + program: Option, +) -> anyhow::Result<()> { let ImportTarget { origin, name, @@ -53,10 +95,12 @@ pub async fn listen_import(target: ImportTarget, addr: SocketAddr, latency: Dura if let Err(err) = publish .with_max_age(max_age) .with_bandwidth(bandwidth) + .with_program(program) .accept(&origin, &name) .await { - tracing::warn!(%name, %err, "SRT ingest ended with error"); + let err = suggest_program(err); + tracing::warn!(%name, err = format!("{err:#}"), "SRT ingest ended with error"); } }); } @@ -107,7 +151,12 @@ pub async fn listen_export( } /// Dial a remote SRT server and pull its stream into the Origin under `target.name` (import). -pub async fn connect_import(target: ImportTarget, url: Url, latency: Duration) -> anyhow::Result<()> { +pub async fn connect_import( + target: ImportTarget, + url: Url, + latency: Duration, + program: Option, +) -> anyhow::Result<()> { let (addr, resource) = parse_url(&url).await?; let name = &target.name; tracing::info!(url = %RedactedUrl::new(&url), %name, "SRT client pulling"); @@ -116,8 +165,9 @@ pub async fn connect_import(target: ImportTarget, url: Url, latency: Duration) - 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?) + .with_bandwidth(target.bandwidth) + .with_program(program); + client.pull(&target.origin, name).await.map_err(suggest_program) } /// Push a broadcast from the Origin to a remote SRT server (export). @@ -184,6 +234,18 @@ mod tests { assert!(parse("udp://127.0.0.1:9000").await.is_err()); } + /// A multiplex refusal names the flag that resolves it; other errors pass through unchanged. + #[test] + fn a_multiplex_suggests_the_program_flag() { + let refused = moq_mux::container::ts::MultipleProgramsError { programs: vec![1, 2] }; + let err = suggest_program(moq_srt::Error::Mux(anyhow::Error::from(refused).into())); + let err = format!("{err:#}"); + assert!(err.contains("--program") && err.contains("programs (1, 2)"), "{err}"); + + let err = format!("{:#}", suggest_program(moq_srt::Error::ListenerClosed)); + assert!(!err.contains("--program"), "{err}"); + } + #[tokio::test] async fn requires_port_and_resource() { assert!(parse("srt://127.0.0.1").await.is_err()); diff --git a/rs/moq-mux/src/container/ts/import.rs b/rs/moq-mux/src/container/ts/import.rs index d23e1160ef..5f4c10ef23 100644 --- a/rs/moq-mux/src/container/ts/import.rs +++ b/rs/moq-mux/src/container/ts/import.rs @@ -1000,7 +1000,8 @@ impl Import { /// A PAT listed more than one program, and the [`Import`] was not told which one to take. /// /// Importing them all onto one broadcast would put unrelated clocks on one timeline, so the -/// import stops instead. Pick one with [`Import::with_program`]. +/// import stops instead. Pick one with [`Import::with_program`], or import each as its own +/// broadcast with [`Programs`](super::Programs). #[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)] #[error("transport stream carries {} programs ({})", .programs.len(), list_programs(.programs))] pub struct MultipleProgramsError { @@ -1012,36 +1013,6 @@ fn list_programs(programs: &[u16]) -> String { programs.iter().map(u16::to_string).collect::>().join(", ") } -/// The program numbers listed by the first whole PAT in `data`, in PAT order, or `None` if -/// there is none yet. -/// -/// For choosing importers before any exists: each [`Import`] reads the PAT itself. -pub fn programs(data: &[u8]) -> Option> { - let mut off = 0; - while let Some(rel) = memchr::memchr(0x47, &data[off..]) { - off += rel; - let packet = data.get(off..off + TsPacket::SIZE)?; - // PID 0 opening a section. The section CRC rejects a sync byte found in payload. - if packet[1] & 0x5f == 0x40 - && packet[2] == 0 - && let Ok(Some(TsPacket { - payload: Some(TsPayload::Pat(pat)), - .. - })) = TsPacketReader::new(packet).read_ts_packet() - { - return Some( - pat.table - .iter() - .map(|entry| entry.program_num) - .filter(|&program| program != 0) - .collect(), - ); - } - off += 1; - } - None -} - /// What each demuxed elementary stream delivered, and the audio frame sync it lost or could /// not verify, keyed by elementary stream PID. /// @@ -3097,7 +3068,7 @@ impl Read for Feed { } #[cfg(test)] -mod test { +pub(super) mod test { use std::collections::BTreeMap; use std::time::Duration; @@ -5424,7 +5395,7 @@ mod test { /// A two-program multiplex whose clocks sit an hour apart: program 1's MP2 on /// `0x0061` starts at 1 s, program 2's on `0x0071` at 3601 s. - fn two_programs() -> Vec { + pub(in crate::container::ts) fn two_programs() -> Vec { let mut data = synth_programs( &[ (1, 0x0100, &[(StreamType::Mpeg1Audio, 0x0061)]), @@ -5504,17 +5475,6 @@ mod test { assert!(err.contains("no program 3") && err.contains("1, 2"), "{err}"); } - #[test] - fn programs_reads_the_first_whole_pat() { - let data = two_programs(); - assert_eq!(super::programs(&data), Some(vec![1, 2])); - // A sync byte in leading junk is skipped; a truncated PAT is not a PAT yet. - let mut shifted = vec![0x47, 0x40, 0x00, 0x10]; - shifted.extend_from_slice(&data); - assert_eq!(super::programs(&shifted), Some(vec![1, 2])); - assert_eq!(super::programs(&data[..100]), None); - } - /// Two MP2 renditions, the first of which the PMT designates as the PCR PID. const PCR_PID: u16 = 0x0061; const PEER_PID: u16 = 0x0062; diff --git a/rs/moq-mux/src/container/ts/mod.rs b/rs/moq-mux/src/container/ts/mod.rs index da0158c5b1..05858df90f 100644 --- a/rs/moq-mux/src/container/ts/mod.rs +++ b/rs/moq-mux/src/container/ts/mod.rs @@ -22,6 +22,7 @@ mod adts; mod export; mod import; mod mux_rate; +mod programs; mod si; // The `mpegts` catalog section (per-track PID + descriptors plus verbatim carriage @@ -32,6 +33,7 @@ mod catalog; pub use catalog::{Catalog, Descriptor, Ext, Framing, Mpegts, Program, SiEntry, Track, Verbatim}; pub use export::*; pub use import::*; +pub use programs::Programs; pub mod stats; diff --git a/rs/moq-mux/src/container/ts/programs.rs b/rs/moq-mux/src/container/ts/programs.rs new file mode 100644 index 0000000000..8f2a139b74 --- /dev/null +++ b/rs/moq-mux/src/container/ts/programs.rs @@ -0,0 +1,242 @@ +//! Every program of a multiplex, each imported as its own broadcast. + +use anyhow::Context; +use mpeg2ts::ts::{ReadTsPacket, TsPacket, TsPacketReader, TsPayload}; + +use super::{Ext, Import, Stats}; +use crate::catalog; + +/// Imports every program the first PAT lists, each as its own broadcast on `origin`, with its +/// own clock and catalog. +/// +/// Program `n` of `event.hang` publishes on `event/n.hang`, keeping the catalog suffix last so +/// format detection still sees it; a name without one gains a plain `/n`. Input ahead of that +/// PAT is dropped, as a single [`Import`] drops media ahead of its PSI. A program a later PAT +/// adds is not published; one it removes ends the import. +pub struct Programs { + origin: moq_net::origin::Producer, + name: moq_net::PathOwned, + config: catalog::Config, + live: bool, + /// Input held until it contains a whole PAT. + pending: Vec, + /// Empty until the PAT arrives. + programs: Vec, +} + +/// One program's broadcast. +struct Program { + /// Keeps the origin-created broadcast alive; the importer's clone doesn't. + _broadcast: moq_net::broadcast::Producer, + import: Import, + catalog: catalog::Producer, +} + +impl Programs { + /// Split the multiplex into broadcasts named after `name`, each catalog built from `config`. + pub fn new(origin: moq_net::origin::Producer, name: impl moq_net::AsPath, config: catalog::Config) -> Self { + Self { + origin, + name: name.as_path().to_owned(), + config, + live: false, + pending: Vec::new(), + programs: Vec::new(), + } + } + + /// Publish each program on its broadcast's clock, as [`Import::live`] does. + pub fn live(mut self) -> Self { + self.live = true; + self + } + + /// Demux a chunk of the multiplex. A trailing partial packet is retained for the next call. + pub fn decode(&mut self, data: &[u8]) -> anyhow::Result<()> { + if !self.programs.is_empty() { + return self.feed(data); + } + self.pending.extend_from_slice(data); + let Some(programs) = pat_programs(&self.pending).filter(|programs| !programs.is_empty()) else { + // Only a packet's worth of tail can still hold the start of the PAT. + let keep = self.pending.len().saturating_sub(TsPacket::SIZE - 1); + self.pending.drain(..keep); + return Ok(()); + }; + for program in programs { + let name = program_broadcast(self.name.as_str(), program); + let mut broadcast = self + .origin + .create_broadcast(&name) + .with_context(|| format!("failed to create broadcast {name}"))?; + // TS carries undecoded elementary streams verbatim, so each catalog carries the + // `mpegts` extension, and owns its broadcast's catalog tracks. + let config = self + .config + .clone() + .with_catalog(catalog::hang::Catalog::::default()); + let catalog = catalog::Producer::new(&mut broadcast, config)?; + let mut import = Import::new(broadcast.clone(), catalog.reserve()).with_program(program); + if self.live { + import = import.live(); + } + broadcast + .announce(Default::default()) + .with_context(|| format!("failed to announce broadcast {name}"))?; + self.programs.push(Program { + _broadcast: broadcast, + import, + catalog, + }); + } + let pending = std::mem::take(&mut self.pending); + self.feed(&pending) + } + + fn feed(&mut self, data: &[u8]) -> anyhow::Result<()> { + for program in &mut self.programs { + program.import.decode(data)?; + } + Ok(()) + } + + /// Every program's counters in one map: PIDs are unique across a multiplex. + pub fn stats(&self) -> Stats { + let mut stats = Stats::default(); + for program in &self.programs { + stats.streams.extend(program.import.stats().streams); + } + stats + } + + /// Finish each program's tracks, then its catalog while it still lists them. + pub fn finish(&mut self) -> anyhow::Result<()> { + for program in &mut self.programs { + program.import.finish()?; + program.catalog.finish()?; + } + Ok(()) + } + + /// Abort every program's tracks with `err`, as [`Import::abort`] does. Consumes the importer. + pub fn abort(self, err: moq_net::Error) { + for program in self.programs { + program.import.abort(err.clone()); + } + } +} + +/// The broadcast one program of `name` publishes on: `event.hang` becomes `event/2.hang`. +fn program_broadcast(name: &str, program: u16) -> String { + match catalog::CatalogFormat::detect(name) { + Some(format) => { + let stem = &name[..name.len() - format.extension().len()]; + format!("{stem}/{program}{}", format.extension()) + } + None => format!("{name}/{program}"), + } +} + +/// The program numbers listed by the first whole PAT in `data`, in PAT order, or `None` if +/// there is none yet. +fn pat_programs(data: &[u8]) -> Option> { + let mut off = 0; + while let Some(rel) = memchr::memchr(0x47, &data[off..]) { + off += rel; + let packet = data.get(off..off + TsPacket::SIZE)?; + // PID 0 opening a section. The section CRC rejects a sync byte found in payload. + if packet[1] & 0x5f == 0x40 + && packet[2] == 0 + && let Ok(Some(TsPacket { + payload: Some(TsPayload::Pat(pat)), + .. + })) = TsPacketReader::new(packet).read_ts_packet() + { + return Some( + pat.table + .iter() + .map(|entry| entry.program_num) + .filter(|&program| program != 0) + .collect(), + ); + } + off += 1; + } + None +} + +#[cfg(test)] +mod test { + use std::time::Duration; + + use super::*; + use crate::catalog::hang::Container; + use crate::container::Consumer; + use crate::container::ts::import::test::two_programs; + + #[test] + fn program_broadcasts_keep_the_catalog_suffix_last() { + assert_eq!(program_broadcast("event.hang", 2), "event/2.hang"); + assert_eq!(program_broadcast("demo/event.msf", 7), "demo/event/7.msf"); + assert_eq!(program_broadcast("event", 1), "event/1"); + } + + #[test] + fn pat_programs_reads_the_first_whole_pat() { + let data = two_programs(); + assert_eq!(pat_programs(&data), Some(vec![1, 2])); + // A sync byte in leading junk is skipped; a truncated PAT is not a PAT yet. + let mut shifted = vec![0x47, 0x40, 0x00, 0x10]; + shifted.extend_from_slice(&data); + assert_eq!(pat_programs(&shifted), Some(vec![1, 2])); + assert_eq!(pat_programs(&data[..100]), None); + } + + /// The input is held until the PAT is whole, then each program publishes as its own + /// broadcast carrying only that program, live on its own first frame rather than an hour + /// apart on one shared clock. + #[tokio::test] + async fn every_program_publishes_its_own_broadcast() { + let origin = crate::source::produce_origin(); + let ago = Duration::from_secs(60); + let clock = crate::Clock::at(std::time::Instant::now() - ago, std::time::SystemTime::now() - ago).unwrap(); + let config = catalog::Config::default().with_clock(clock); + let mut programs = Programs::new(origin.clone(), "event.hang", config).live(); + + let input = two_programs(); + let before = clock.now(); + programs.decode(&input[..100]).unwrap(); + assert!( + programs.programs.is_empty(), + "no program is known before the PAT is whole" + ); + programs.decode(&input[100..]).unwrap(); + programs.finish().unwrap(); + let after = clock.now(); + + for (path, fills) in [("event/1.hang", [0xAA, 0xBB]), ("event/2.hang", [0xCC, 0xDD])] { + let consumer = crate::Source::new(origin.consume(), path).broadcast().await.unwrap(); + let catalog = hang::catalog::Catalog::<()>::subscribe(&consumer) + .await + .unwrap() + .next() + .await + .unwrap() + .expect("a catalog"); + let renditions: Vec<_> = catalog.audio.renditions.iter().collect(); + assert_eq!(renditions.len(), 1, "{path} carries its own program's one stream"); + let (name, config) = renditions[0]; + let track = consumer.track(name).unwrap().subscribe(None).await.unwrap(); + let container = Container::try_from(config).unwrap(); + let frame = Consumer::new(track, container).read().await.unwrap().expect("a frame"); + assert!(fills.contains(&frame.payload[4]), "{path} carries only its own program"); + let skew = Duration::from_secs(2).as_micros(); + assert!( + before.as_micros() - skew <= frame.timestamp.as_micros() + && frame.timestamp.as_micros() <= after.as_micros() + skew, + "{path} is live on arrival: {:?} not in {before:?}..={after:?}", + frame.timestamp + ); + } + } +} diff --git a/rs/moq-srt/README.md b/rs/moq-srt/README.md index 36a25f5a68..c589fa8ceb 100644 --- a/rs/moq-srt/README.md +++ b/rs/moq-srt/README.md @@ -39,9 +39,10 @@ tokio::select! { ## CLI A command-line interface is provided by the [`moq-cli`](../moq-cli) binary, on -top of this library. +top of this library. Its listener bridges the single `--broadcast` and ignores +the stream id; see [`doc/bin/srt.md`](../../doc/bin/srt.md). -Feed any SRT source: +Feed any SRT source into a `run` listener with prefix `live`: ```bash # Publish: lands at broadcast `live/cam0`. @@ -58,16 +59,22 @@ the publisher does. ## Routing -Each connection's broadcast path and direction come from its SRT stream id: +Under `run`, each connection's broadcast path and direction come from its SRT +stream id: - Standard form `#!::r=,m=` -> ``, with `m=request` selecting egress and anything else (including absent) selecting ingest. - Otherwise the raw stream id (e.g. OBS-style `app/key`), always ingest. -`--srt-prefix` is prepended to namespace a listener's streams. First publisher on +`Config::prefix` is prepended to namespace a listener's streams. First publisher on a path wins; a second publish of the same path is rejected. Requests don't claim a path, so any number of players can pull the same broadcast. +A multi-program feed is refused unless `Config::program` picks one program +(`Program::One(n)`) or every program (`Program::All`), which publishes each on +its own broadcast under the path: `live/cam0/1`, `live/cam0/2`, and so on. To +choose per connection, drive `Server` and call `Publish::with_program`. + ## Auth `run` is unauthenticated: anyone who can reach the UDP port can publish or diff --git a/rs/moq-srt/src/dial.rs b/rs/moq-srt/src/dial.rs index 6c51666405..2367a85784 100644 --- a/rs/moq-srt/src/dial.rs +++ b/rs/moq-srt/src/dial.rs @@ -23,8 +23,8 @@ use std::time::Duration; use moq_net::origin; use srt_tokio::SrtSocket; -use crate::Result; use crate::server::{DEFAULT_LATENCY, configure_buffers, serve_publish, serve_subscribe}; +use crate::{Program, Result}; /// An SRT caller that can publish a MoQ broadcast or pull a remote stream. /// @@ -57,6 +57,9 @@ pub struct Client { /// Connection allocator each ingested track claims its peak-hold bitrate on. /// [`pull`] only; [`publish`] reads a broadcast someone else declared. bandwidth: moq_net::bandwidth::Allocator, + + /// The programs of a multi-program remote [`pull`] publishes, or `None` to refuse one. + program: Option, } impl Client { @@ -69,6 +72,7 @@ impl Client { latency: DEFAULT_LATENCY, max_age: None, bandwidth: moq_net::bandwidth::Allocator::unlimited(), + program: None, } } @@ -90,6 +94,13 @@ impl Client { self } + /// Publish one program of a multi-program stream [`pull`](Self::pull) receives, or each as + /// its own broadcast under the pulled path. `None` (the default) refuses a multiplex. + pub fn with_program(mut self, program: impl Into>) -> Self { + self.program = program.into(); + 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(); @@ -104,7 +115,7 @@ impl Client { 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 + serve_publish(origin, path.as_str(), socket, catalog, self.program).await } async fn call(&self, mode: Mode) -> Result { diff --git a/rs/moq-srt/src/lib.rs b/rs/moq-srt/src/lib.rs index ffa1351086..da8d3543bc 100644 --- a/rs/moq-srt/src/lib.rs +++ b/rs/moq-srt/src/lib.rs @@ -48,3 +48,4 @@ pub use dial::Client; pub use error::{Error, Result}; pub use listen::{Config, run}; pub use server::{Publish, Reject, Request, Server, Subscribe}; +pub use ts::Program; diff --git a/rs/moq-srt/src/listen.rs b/rs/moq-srt/src/listen.rs index 6b5f3d84df..5cc7f5476e 100644 --- a/rs/moq-srt/src/listen.rs +++ b/rs/moq-srt/src/listen.rs @@ -31,8 +31,8 @@ use std::time::Duration; use moq_net::{Path, PathOwned, origin}; -use crate::Result; use crate::server::{Request, Server}; +use crate::{Program, Result}; /// SRT gateway configuration. /// @@ -64,6 +64,11 @@ pub struct Config { /// advertise segments that are still fetchable. Lower it when nothing reads history /// and the memory matters. Only affects ingest (`m=publish`); egress ignores it. pub max_age: Option, + + /// The programs of a multi-program feed every ingest publishes, or `None` to refuse a + /// multiplex. To choose per path, drive [`Server`] and call + /// [`Publish::with_program`](crate::Publish::with_program) on each request. + pub program: Option, } impl Default for Config { @@ -73,6 +78,7 @@ impl Default for Config { prefix: Path::empty().to_owned(), latency: crate::server::DEFAULT_LATENCY, max_age: None, + program: None, } } } @@ -110,6 +116,7 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { let active = ActivePaths::default(); let prefix = config.prefix; let max_age = config.max_age; + let program = config.program; while let Some(request) = server.accept().await { let prefix = prefix.clone(); @@ -130,7 +137,12 @@ pub async fn run(origin: origin::Producer, config: Config) -> Result<()> { let _ = publish.reject(crate::Reject::Unavailable).await; return; }; - if let Err(err) = publish.with_max_age(max_age).accept(&origin, &path).await { + if let Err(err) = publish + .with_max_age(max_age) + .with_program(program) + .accept(&origin, &path) + .await + { tracing::warn!(%peer, %path, %err, "SRT ingest ended with error"); } else { tracing::info!(%peer, %path, "SRT ingest ended"); diff --git a/rs/moq-srt/src/server.rs b/rs/moq-srt/src/server.rs index ddbab0d555..e51667dc19 100644 --- a/rs/moq-srt/src/server.rs +++ b/rs/moq-srt/src/server.rs @@ -29,7 +29,7 @@ use srt_tokio::access::{AccessControlList, ConnectionMode, RejectReason, Standar use srt_tokio::options::{PacketCount, SocketOptions, StreamId}; use srt_tokio::{ConnectionRequest, SrtIncoming, SrtListener, SrtSocket}; -use crate::Result; +use crate::{Program, Result}; /// Why an SRT publish or subscribe was refused. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -248,6 +248,7 @@ impl Server { latency: self.latency, max_age: None, bandwidth: moq_net::bandwidth::Allocator::unlimited(), + program: None, }; // `m=request` reads a broadcast out; everything else publishes one in. @@ -280,6 +281,8 @@ struct Pending { /// Connection allocator each ingested track claims its peak-hold bitrate on. /// Override with [`Publish::with_bandwidth`]. bandwidth: moq_net::bandwidth::Allocator, + /// The programs of a multiplex an ingest publishes. Override with [`Publish::with_program`]. + program: Option, } /// What an accepted SRT connection wants: to contribute media ([`Publish`]) or to @@ -378,6 +381,13 @@ impl Publish { self } + /// Publish one program of a multi-program feed, or each as its own broadcast under the + /// accepted path. `None` (the default) fails an ingest whose PAT lists more than one. + pub fn with_program(mut self, program: impl Into>) -> Self { + self.0.program = program.into(); + self + } + /// Accept the publish: announce a broadcast at `path` in `origin` and pump the /// connection's MPEG-TS into it until the client disconnects. /// @@ -392,7 +402,7 @@ impl Publish { let config = moq_mux::catalog::Config::default() .with_max_age(self.0.max_age) .with_bandwidth(self.0.bandwidth); - serve_publish(origin, path.as_str(), socket, config).await + serve_publish(origin, path.as_str(), socket, config, self.0.program).await } /// Reject the publish with a verdict the client can distinguish on the wire. @@ -469,10 +479,11 @@ pub(crate) async fn serve_publish( path: &str, mut socket: SrtSocket, config: moq_mux::catalog::Config, + program: Option, ) -> Result<()> { use futures::TryStreamExt; - let mut publisher = crate::ts::Publisher::new(origin, path, config)?; + let mut publisher = crate::ts::Publisher::new(origin, path, config, program)?; // Run the read/feed loop so an error surfaces here instead of unwinding past // the publisher, which would drop it (and its tracks) with a bare Error::Dropped. diff --git a/rs/moq-srt/src/ts.rs b/rs/moq-srt/src/ts.rs index f81681c5e1..b470181c23 100644 --- a/rs/moq-srt/src/ts.rs +++ b/rs/moq-srt/src/ts.rs @@ -15,21 +15,29 @@ use moq_net::origin; use crate::Result; -/// Publishes an MPEG-TS source into the origin as a single broadcast. +/// Which programs of a multi-program MPEG-TS an ingest publishes. +/// +/// Without one, an ingest whose PAT lists more than one program fails with +/// [`ts::MultipleProgramsError`] rather than merging unrelated clocks onto one broadcast. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Program { + /// Only the program with this PAT program number, on the ingest's path. + One(u16), + /// Every program the first PAT lists, each as its own broadcast under the ingest's path: + /// `live/cam0` publishes `live/cam0/1`, `live/cam0/2`, ... (`event.hang` publishes + /// `event/1.hang`, keeping the catalog suffix last). + All, +} + +/// Publishes an MPEG-TS source into the origin: one broadcast, or one per program. /// /// Each chunk is handed straight to the TS importer, which consumes whole /// transport packets and retains any partial trailing packet internally for the /// next call (the same pattern `moq-cli import ... stdin ts` uses against stdin). -/// Either [`Self::finish`] or dropping the publisher ends the broadcast and -/// unannounces the path. +/// Either [`Self::finish`] or dropping the publisher ends the broadcasts and +/// unannounces their paths. pub struct Publisher { - // TS carries undecoded elementary streams (SCTE-35, teletext, DVB AC-3, ...) - // verbatim, so the importer uses the `mpegts` catalog extension rather than the - // media-only `()`, which would route those PIDs to `Stream::Ignored` and drop them. - importer: ts::Import, - // A clone of the importer's producer, so an end can close the broadcast - // (prompt unannounce) even though the importer owns it. - broadcast: moq_net::broadcast::Producer, + importer: Importer, // The importer's per-stream counters, logged as `moq import ts` logs them, under a // span naming the path, since one server carries many ingests. log: ts::stats::Log, @@ -37,24 +45,53 @@ pub struct Publisher { span: tracing::Span, } +enum Importer { + One { + // TS carries undecoded elementary streams (SCTE-35, teletext, DVB AC-3, ...) + // verbatim, so the importer uses the `mpegts` catalog extension rather than the + // media-only `()`, which would route those PIDs to `Stream::Ignored` and drop them. + import: Box>, + // A clone of the importer's producer, so an end can close the broadcast + // (prompt unannounce) even though the importer owns it. + broadcast: moq_net::broadcast::Producer, + }, + /// Each program's broadcast is created and announced once the first PAT names it. + All(Box), +} + impl Publisher { - /// Create the broadcast on `origin` at `path` and wire up the TS importer + - /// catalog. + /// Wire up the TS importer and catalog for `path` on `origin`, announcing the broadcast + /// now, or each program's once the PAT lists it for [`Program::All`]. /// /// `config` is the catalog the importer publishes into: retention /// (`with_max_age`) and the connection allocator passthrough tracks claim on /// (`with_bandwidth`). - pub fn new(origin: &origin::Producer, path: &str, config: moq_mux::catalog::Config) -> Result { - let mut broadcast = origin.publish(path, moq_net::origin::Route::default())?; - let config = config.with_catalog(moq_mux::catalog::hang::Catalog::::default()); - let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config)?; - let handle = broadcast.clone(); - let importer = ts::Import::new(broadcast, catalog.reserve()); - tracing::info!(%path, "publishing ingest broadcast"); + pub fn new( + origin: &origin::Producer, + path: &str, + config: moq_mux::catalog::Config, + program: Option, + ) -> Result { + let importer = match program { + Some(Program::All) => Importer::All(Box::new(ts::Programs::new(origin.clone(), path, config))), + Some(Program::One(_)) | None => { + let mut broadcast = origin.publish(path, moq_net::origin::Route::default())?; + let config = config.with_catalog(moq_mux::catalog::hang::Catalog::::default()); + let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config)?; + let mut import = ts::Import::new(broadcast.clone(), catalog.reserve()); + if let Some(Program::One(program)) = program { + import = import.with_program(program); + } + Importer::One { + import: Box::new(import), + broadcast, + } + } + }; + tracing::info!(%path, ?program, "publishing ingest broadcast"); Ok(Self { importer, - broadcast: handle, log: ts::stats::Log::default(), sampled: tokio::time::Instant::now(), span: tracing::info_span!("srt", %path), @@ -68,32 +105,53 @@ impl Publisher { /// [`ts::stats::Log::INTERVAL`] it also logs what moved in the importer's /// per-stream counters. pub fn feed(&mut self, data: Bytes) -> Result<()> { - self.importer.decode(&data).map_err(moq_mux::Error::from)?; + match &mut self.importer { + Importer::One { import, .. } => import.decode(&data), + Importer::All(programs) => programs.decode(&data), + } + .map_err(moq_mux::Error::from)?; if self.sampled.elapsed() >= ts::stats::Log::INTERVAL { self.sampled = tokio::time::Instant::now(); let _span = self.span.enter(); - self.log.sample(self.importer.stats()); + self.log.sample(self.stats()); } Ok(()) } + fn stats(&self) -> ts::Stats { + match &self.importer { + Importer::One { import, .. } => import.stats(), + Importer::All(programs) => programs.stats(), + } + } + /// 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().map_err(moq_mux::Error::from)?; + match &mut self.importer { + Importer::One { import, broadcast } => { + import.finish().map_err(moq_mux::Error::from)?; + broadcast.close(); + } + Importer::All(programs) => programs.finish().map_err(moq_mux::Error::from)?, + } // The drain at end of input can publish a frame nothing vouched for. - self.span.in_scope(|| self.log.finish(&self.importer.stats())); - self.broadcast.close(); + self.span.in_scope(|| self.log.finish(&self.stats())); Ok(()) } /// Abort the published tracks with `err` so subscribers see the real cause /// (the SRT caller dropped, a demux error) rather than a generic `Error::Dropped`. /// - /// Consumes the publisher and closes the broadcast. + /// Consumes the publisher and closes the broadcasts. pub fn abort(self, err: moq_net::Error) { - self.importer.abort(err); - self.broadcast.close(); + match self.importer { + Importer::One { import, broadcast } => { + import.abort(err); + broadcast.close(); + } + Importer::All(programs) => programs.abort(err), + } } } @@ -256,7 +314,7 @@ mod tests { /// Publish `ts` on `path` one SRT payload (7 packets) at a time, each delivered when the /// program clock says it is due. async fn ingest(origin: &moq_net::origin::Producer, path: &str, ts: &[u8]) { - let mut publisher = Publisher::new(origin, path, Default::default()).unwrap(); + let mut publisher = Publisher::new(origin, path, Default::default(), None).unwrap(); let mut clock = Duration::ZERO; for payload in timed(ts).collect::>().chunks(7) { let (due, _) = payload[0]; @@ -295,23 +353,101 @@ mod tests { section } - /// A PAT with one program (number 1) whose PMT lives on `pmt_pid`. - fn pat(pmt_pid: u16) -> Vec { - let mut s = vec![0x00, 0xb0, 0x0d, 0x00, 0x01, 0xc1, 0x00, 0x00]; - s.extend_from_slice(&[0x00, 0x01]); - s.extend_from_slice(&[0xe0 | (pmt_pid >> 8) as u8, pmt_pid as u8]); + /// A PAT listing each `(program_number, pmt_pid)`. + fn pat(programs: &[(u16, u16)]) -> Vec { + let len = 9 + 4 * programs.len(); + let mut s = vec![0x00, 0xb0, len as u8, 0x00, 0x01, 0xc1, 0x00, 0x00]; + for &(program, pmt_pid) in programs { + s.extend_from_slice(&program.to_be_bytes()); + s.extend_from_slice(&[0xe0 | (pmt_pid >> 8) as u8, pmt_pid as u8]); + } seal(s) } - /// A PMT for program 1 declaring a single H.264 elementary stream on `es_pid`. - fn pmt(es_pid: u16) -> Vec { - let mut s = vec![0x02, 0xb0, 0x12, 0x00, 0x01, 0xc1, 0x00, 0x00]; + /// A PMT for `program` declaring a single elementary stream of `stream_type` on `es_pid`. + fn pmt(program: u16, stream_type: u8, es_pid: u16) -> Vec { + let mut s = vec![0x02, 0xb0, 0x12]; + s.extend_from_slice(&program.to_be_bytes()); + s.extend_from_slice(&[0xc1, 0x00, 0x00]); s.extend_from_slice(&[0xe0 | (es_pid >> 8) as u8, es_pid as u8]); s.extend_from_slice(&[0xf0, 0x00]); - s.extend_from_slice(&[0x1b, 0xe0 | (es_pid >> 8) as u8, es_pid as u8, 0xf0, 0x00]); + s.extend_from_slice(&[stream_type, 0xe0 | (es_pid >> 8) as u8, es_pid as u8, 0xf0, 0x00]); seal(s) } + /// The PSI of a two-program multiplex, each program carrying one private PES stream + /// (carried verbatim, so the catalog lists it without waiting on media): program 1's on + /// PID 0x101, program 2's on 0x201. + fn two_programs() -> Bytes { + let mut ts = psi_packet(0x0000, &pat(&[(1, 0x0100), (2, 0x0200)])); + ts.extend_from_slice(&psi_packet(0x0100, &pmt(1, 0x06, 0x0101))); + ts.extend_from_slice(&psi_packet(0x0200, &pmt(2, 0x06, 0x0201))); + ts.into() + } + + /// The PAT program number and the stream PIDs the catalog at `path` records, once it is + /// announced. + async fn program(origin: &moq_net::origin::Producer, path: &str) -> (u16, Vec) { + let consumer = origin.consume(); + timeout(Duration::from_secs(5), consumer.routed(path)) + .await + .expect("announce timed out") + .expect("the broadcast is announced"); + let broadcast = consumer.request_broadcast(path).await.unwrap(); + let mut catalog = moq_mux::catalog::Consumer::::new(&broadcast, CatalogFormat::Hang) + .await + .unwrap(); + let snapshot = timeout(Duration::from_secs(5), catalog.next()) + .await + .expect("catalog timed out") + .unwrap() + .expect("a catalog"); + let program = snapshot.ext.mpegts.program.as_ref().expect("the PAT identity"); + let pids = snapshot.ext.mpegts.tracks.values().map(|track| track.pid).collect(); + (program.program_number, pids) + } + + /// Without a selection, a multiplex is refused rather than merged onto one broadcast. + #[tokio::test(start_paused = true)] + async fn publisher_refuses_a_multiplex() { + let origin = produce_origin(); + let mut publisher = Publisher::new(&origin, "ingest", Default::default(), None).unwrap(); + let err = publisher.feed(two_programs()).unwrap_err(); + let crate::Error::Mux(moq_mux::Error::Other(inner)) = &err else { + panic!("a demux error: {err}"); + }; + assert_eq!( + inner.downcast_ref::(), + Some(&ts::MultipleProgramsError { programs: vec![1, 2] }) + ); + } + + /// `Program::One` publishes the chosen program alone on the ingest's path. + #[tokio::test(start_paused = true)] + async fn publisher_imports_one_program() { + let origin = produce_origin(); + let mut publisher = Publisher::new(&origin, "ingest", Default::default(), Some(Program::One(2))).unwrap(); + publisher.feed(two_programs()).unwrap(); + assert_eq!(program(&origin, "ingest").await, (2, vec![0x0201])); + } + + /// `Program::All` publishes each program as its own broadcast under the ingest's path. + #[tokio::test(start_paused = true)] + async fn publisher_imports_every_program() { + let origin = produce_origin(); + let mut publisher = Publisher::new(&origin, "ingest", Default::default(), Some(Program::All)).unwrap(); + publisher.feed(two_programs()).unwrap(); + assert_eq!(program(&origin, "ingest/1").await, (1, vec![0x0101])); + assert_eq!(program(&origin, "ingest/2").await, (2, vec![0x0201])); + assert!( + timeout(Duration::from_secs(1), origin.consume().routed("ingest")) + .await + .is_err(), + "nothing is published on the bare path" + ); + publisher.finish().unwrap(); + } + /// The retention the caller configured has to reach the media tracks the TS importer /// mints off the PMT, not stop at the catalog producer it was set on. #[tokio::test] @@ -321,11 +457,12 @@ mod tests { &origin, "live/cam0", moq_mux::catalog::Config::default().with_max_age(Duration::from_secs(3)), + None, ) .unwrap(); - let mut ts = psi_packet(0x0000, &pat(0x0100)); - ts.extend_from_slice(&psi_packet(0x0100, &pmt(0x0101))); + let mut ts = psi_packet(0x0000, &pat(&[(1, 0x0100)])); + ts.extend_from_slice(&psi_packet(0x0100, &pmt(1, 0x1b, 0x0101))); publisher.feed(Bytes::from(ts)).unwrap(); let consumer = origin.consume(); @@ -373,7 +510,7 @@ mod tests { #[tokio::test(start_paused = true)] async fn publisher_preserves_scte35_cues() { let origin = produce_origin(); - let mut publisher = Publisher::new(&origin, "ingest", Default::default()).unwrap(); + let mut publisher = Publisher::new(&origin, "ingest", Default::default(), None).unwrap(); let consumer = origin.consume(); timeout(Duration::from_secs(5), consumer.routed("ingest"))