Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions doc/bin/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 9 additions & 0 deletions doc/bin/srt.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 3 additions & 3 deletions quest/m1/ts-psi-reassembly.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
23 changes: 21 additions & 2 deletions rs/moq-cli/src/args.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -782,7 +782,7 @@ pub struct TsImport {
pub program: Option<TsProgram>,
}

/// 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.
Expand Down Expand Up @@ -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,
Expand Down
6 changes: 4 additions & 2 deletions rs/moq-cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) => {
Expand Down
179 changes: 3 additions & 176 deletions rs/moq-cli/src/publish.rs
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ enum PublishDecoder {
// verbatim, so it uses the `mpegts` catalog extension rather than the media-only `()`.
Ts(Box<ts::Import<ts::Ext>>),
/// `import ts --program all`: one importer, broadcast, and catalog per program.
TsPrograms(Box<TsPrograms>),
TsPrograms(Box<ts::Programs>),
Flv(Box<flv::Import>),
}

Expand Down Expand Up @@ -218,7 +218,7 @@ fn suggest_program(err: anyhow::Error) -> anyhow::Error {
enum PublishCatalog {
Media(moq_mux::catalog::Producer),
Ts(moq_mux::catalog::Producer<ts::Ext>),
/// `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,
}

Expand Down Expand Up @@ -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<ts::Ext>,
catalog: moq_mux::catalog::Producer<ts::Ext>,
}

/// `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<u8>,
/// Empty until the PAT arrives.
programs: Vec<ProgramImport>,
}

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)]
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<u8> {
Expand Down Expand Up @@ -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<u8> {
let track = consumer.track(name).unwrap().subscribe(None).await.unwrap();
Expand Down
Loading
Loading