use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::path::PathBuf;
use broadcast_common::{Package, Unpackage};
use mpeg_ts::ts::{SectionReassembler, TsPacket};
use transmux::pipeline::{CodecConfig, DataCarriage};
use transmux::{CmafMux, Fmp4Demux, Media, Track, TsDemux, TsHlsPackager, TsMux};
const TS: usize = 188;
fn fixtures_ts_dir() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../fixtures/ts")
}
fn read_fixture(name: &str) -> Vec<u8> {
let path = fixtures_ts_dir().join(name);
std::fs::read(&path).unwrap_or_else(|e| panic!("{}: {e}", path.display()))
}
fn demux(data: &[u8]) -> Media {
TsDemux::new()
.unpackage(data)
.expect("TS demux must succeed")
}
fn pusi_seen_pids(data: &[u8]) -> HashMap<u16, usize> {
let mut counts = HashMap::new();
for chunk in data.chunks_exact(TS) {
let Ok(pkt) = TsPacket::parse(chunk) else {
continue;
};
if pkt.header.pusi {
*counts.entry(pkt.header.pid).or_insert(0usize) += 1;
}
}
counts
}
fn find_pmt_pid(data: &[u8]) -> u16 {
const PAT_PID: u16 = 0x0000;
const TABLE_ID_PAT: u8 = 0x00;
let mut reasm = SectionReassembler::default();
for chunk in data.chunks_exact(TS) {
let Ok(pkt) = TsPacket::parse(chunk) else {
continue;
};
if pkt.header.pid != PAT_PID {
continue;
}
let Some(payload) = pkt.payload else {
continue;
};
reasm.feed(payload, pkt.header.pusi);
while let Some(section) = reasm.pop_section() {
if section.first().copied() != Some(TABLE_ID_PAT) || section.len() < 12 {
continue;
}
let section_length = (((section[1] & 0x0F) as usize) << 8) | section[2] as usize;
let end = (3 + section_length).min(section.len());
if end < 12 {
continue;
}
let body = §ion[8..end - 4];
let mut off = 0usize;
while off + 4 <= body.len() {
let program_number = u16::from_be_bytes([body[off], body[off + 1]]);
let pid = (((body[off + 2] & 0x1F) as u16) << 8) | body[off + 3] as u16;
if program_number != 0 {
return pid;
}
off += 4;
}
}
}
panic!("no program_map_PID found via an independent PAT walk");
}
fn collect_pmt_es(data: &[u8], pmt_pid: u16) -> Vec<(u8, u16, Vec<u8>)> {
const TABLE_ID_PMT: u8 = 0x02;
let mut reasm = SectionReassembler::default();
let mut out = Vec::new();
for chunk in data.chunks_exact(TS) {
let Ok(pkt) = TsPacket::parse(chunk) else {
continue;
};
if pkt.header.pid != pmt_pid {
continue;
}
let Some(payload) = pkt.payload else {
continue;
};
reasm.feed(payload, pkt.header.pusi);
while let Some(section) = reasm.pop_section() {
if section.first().copied() != Some(TABLE_ID_PMT) || section.len() < 12 {
continue;
}
let section_length = (((section[1] & 0x0F) as usize) << 8) | section[2] as usize;
let end = (3 + section_length).min(section.len());
if end < 12 {
continue;
}
let body = §ion[8..end - 4];
if body.len() < 4 {
continue;
}
let program_info_length = (((body[2] & 0x0F) as usize) << 8) | body[3] as usize;
let mut off = 4 + program_info_length;
while off + 5 <= body.len() {
let stream_type = body[off];
let pid = (((body[off + 1] & 0x1F) as u16) << 8) | body[off + 2] as u16;
let es_info_length =
(((body[off + 3] & 0x0F) as usize) << 8) | body[off + 4] as usize;
let ds = off + 5;
let de = (ds + es_info_length).min(body.len());
out.push((stream_type, pid, body[ds..de].to_vec()));
off += 5 + es_info_length;
}
}
}
out
}
struct LiveEs {
pid: u16,
stream_type: u8,
descriptor_variants: BTreeSet<Vec<u8>>,
}
fn live_pmt_es(data: &[u8]) -> Vec<LiveEs> {
let pmt_pid = find_pmt_pid(data);
let pusi_counts = pusi_seen_pids(data);
let mut by_pid: BTreeMap<u16, (u8, BTreeSet<Vec<u8>>)> = BTreeMap::new();
for (stream_type, pid, descriptors) in collect_pmt_es(data, pmt_pid) {
if pusi_counts.get(&pid).copied().unwrap_or(0) == 0 {
continue; }
let entry = by_pid.entry(pid).or_insert((stream_type, BTreeSet::new()));
assert_eq!(
entry.0, stream_type,
"PID {pid:#06x} must not change stream_type across PMT repeats"
);
entry.1.insert(descriptors);
}
by_pid
.into_iter()
.map(|(pid, (stream_type, descriptor_variants))| LiveEs {
pid,
stream_type,
descriptor_variants,
})
.collect()
}
fn expected_carriage(stream_type: u8) -> DataCarriage {
match stream_type {
0x05 | 0x0A | 0x0B | 0x0C | 0x0D | 0x14 | 0x86 => DataCarriage::Sections,
_ => DataCarriage::Pes,
}
}
fn expected_dolby_dts_tag(stream_type: u8, descriptors: &[u8]) -> Option<u8> {
if !matches!(stream_type, 0x06 | 0x15) {
return None;
}
let mut off = 0usize;
while off + 2 <= descriptors.len() {
let tag = descriptors[off];
let len = descriptors[off + 1] as usize;
if matches!(tag, 0x6A | 0x7A | 0x7B) {
return Some(tag);
}
off += 2 + len;
}
None
}
fn find_track_by_pid(media: &Media, pid: u16) -> &Track {
media
.tracks
.iter()
.find(|t| t.spec.source_pid == Some(pid))
.unwrap_or_else(|| panic!("no track with source_pid {pid:#06x}"))
}
fn is_valid_long_form_section(bytes: &[u8]) -> bool {
if bytes.len() < 3 {
return false;
}
let section_length = (((bytes[1] & 0x0F) as usize) << 8) | bytes[2] as usize;
3 + section_length == bytes.len()
}
fn find_data_track<'a>(media: &'a Media, es: &LiveEs) -> &'a Track {
media
.tracks
.iter()
.find(|t| match &t.spec.config {
CodecConfig::Data {
stream_type,
descriptors,
..
} => *stream_type == es.stream_type && es.descriptor_variants.contains(descriptors),
_ => false,
})
.unwrap_or_else(|| {
panic!(
"PID {:#06x}: no Data track for stream_type {:#04X} matching any \
of the {} observed ES_info variants",
es.pid,
es.stream_type,
es.descriptor_variants.len()
)
})
}
fn find_data_track_exact<'a>(media: &'a Media, stream_type: u8, descriptors: &[u8]) -> &'a Track {
media
.tracks
.iter()
.find(|t| match &t.spec.config {
CodecConfig::Data {
stream_type: st,
descriptors: d,
..
} => *st == stream_type && d.as_slice() == descriptors,
_ => false,
})
.unwrap_or_else(|| {
panic!("no Data track for stream_type {stream_type:#04X} with the exact descriptors")
})
}
#[test]
fn demux_completeness_every_live_pmt_stream_becomes_a_track() {
let data = read_fixture("m6-single.ts");
let media = demux(&data);
let live = live_pmt_es(&data);
assert_eq!(
live.len(),
6,
"m6-single.ts must have exactly 6 live PMT-listed PIDs"
);
assert_eq!(
media.tracks.len(),
6,
"every live PMT stream must become exactly one track, got {:?}",
media
.tracks
.iter()
.map(|t| &t.spec.config)
.collect::<Vec<_>>()
);
let mut n_pes = 0usize;
let mut n_sections = 0usize;
let mut n_dolby = 0usize;
for es in &live {
let track = find_track_by_pid(&media, es.pid);
let expected_tag = es
.descriptor_variants
.iter()
.find_map(|d| expected_dolby_dts_tag(es.stream_type, d));
match (expected_tag, &track.spec.config) {
(
None,
CodecConfig::Data {
carriage,
descriptors,
..
},
) => {
assert!(
!descriptors.is_empty() || es.descriptor_variants.contains(&Vec::new()),
"PID {:#06x}: descriptors must equal the (non-empty, in this fixture) \
PMT ES_info bytes",
es.pid
);
assert_eq!(
*carriage,
expected_carriage(es.stream_type),
"PID {:#06x} (stream_type {:#04X}) carriage classification",
es.pid,
es.stream_type
);
match carriage {
DataCarriage::Pes => n_pes += 1,
DataCarriage::Sections => n_sections += 1,
_ => {}
}
}
(Some(0x6A), CodecConfig::Ac3 { .. })
| (Some(0x7A), CodecConfig::Eac3 { .. })
| (Some(0x7B), CodecConfig::Dts { .. }) => n_dolby += 1,
(tag, other) => panic!(
"PID {:#06x}: expected_dolby_tag={tag:?} but track config is {other:?} \
-- issue #641",
es.pid
),
}
}
assert_eq!(
n_pes, 1,
"expected 1 PES-carried Data track (0x8C subtitle) -- 0x82/0x83/0x84 now \
classify as E-AC-3, not opaque data (issue #641)"
);
assert_eq!(
n_sections, 2,
"expected 2 section-carried Data tracks (0xAA/0xAB — 0xAC never starts a \
section in this excerpt, see module docs)"
);
assert_eq!(
n_dolby, 3,
"expected 3 E-AC-3 tracks (0x82/0x83/0x84: main, audio-description, and a \
third variant, all carrying a real enhanced_AC3_descriptor -- issue #641)"
);
}
#[test]
fn section_tracks_carry_valid_reassembled_sections() {
let data = read_fixture("m6-single.ts");
let media = demux(&data);
let live = live_pmt_es(&data);
let mut checked = 0usize;
for es in &live {
if expected_carriage(es.stream_type) != DataCarriage::Sections {
continue;
}
let track = find_data_track(&media, es);
assert!(
!track.samples.is_empty(),
"section-carried stream_type {:#04X} must have >= 1 sample",
es.stream_type
);
for (i, sample) in track.samples.iter().enumerate() {
assert!(
is_valid_long_form_section(&sample.data),
"stream_type {:#04X} sample {i} is not a structurally valid \
long-form section (len {}), proving it was NOT reassembled",
es.stream_type,
sample.data.len()
);
assert!(
sample.dts.is_none(),
"a section sample must carry no DTS (dts: None), never a fabricated one"
);
assert!(
sample.pts.is_none(),
"a section sample must carry no PTS (pts: None), never a fabricated one"
);
assert!(
sample.duration.is_none(),
"a section sample must carry no duration either"
);
}
checked += 1;
}
assert_eq!(checked, 2, "expected to check both section-carried tracks");
}
#[test]
fn ts_ir_ts_round_trip_is_payload_lossless_for_data_and_sections() {
let data = read_fixture("m6-single.ts");
let media = demux(&data);
let ts2 = TsMux::new()
.package(&media)
.expect("TsMux must carry every Data and E-AC-3 track, not error");
let media2 = demux(&ts2);
assert_eq!(
media2.tracks.len(),
media.tracks.len(),
"re-demux must recover the same number of tracks"
);
let out_pmt_pid = find_pmt_pid(&ts2);
let out_es = collect_pmt_es(&ts2, out_pmt_pid);
for track in &media.tracks {
let orig_payloads: Vec<&[u8]> = track.samples.iter().map(|s| s.data.as_ref()).collect();
match &track.spec.config {
CodecConfig::Data {
stream_type,
descriptors,
..
} => {
assert!(
out_es
.iter()
.any(|(st, _pid, d)| st == stream_type && d == descriptors),
"re-emitted PMT must list stream_type {stream_type:#04X} with its \
preserved ES_info descriptors"
);
let round = find_data_track_exact(&media2, *stream_type, descriptors);
let round_payloads: Vec<&[u8]> =
round.samples.iter().map(|s| s.data.as_ref()).collect();
assert_eq!(
orig_payloads, round_payloads,
"stream_type {stream_type:#04X}: sample payloads must round-trip byte-for-byte"
);
}
CodecConfig::Eac3 { .. } => {
let first = orig_payloads
.first()
.expect("E-AC-3 track must have at least one sample");
let round = media2
.tracks
.iter()
.find(|t| {
matches!(t.spec.config, CodecConfig::Eac3 { .. })
&& t.samples.first().map(|s| s.data.as_ref()) == Some(*first)
})
.unwrap_or_else(|| panic!("no re-demuxed E-AC-3 track matching first sample"));
let round_payloads: Vec<&[u8]> =
round.samples.iter().map(|s| s.data.as_ref()).collect();
assert_eq!(
orig_payloads, round_payloads,
"E-AC-3 track: sample payloads must round-trip byte-for-byte"
);
}
other => panic!(
"m6-single.ts must produce only Data or E-AC-3 tracks in this excerpt, \
got {other:?}"
),
}
}
}
#[test]
fn ts_to_fmp4_errors_naming_data_track_then_succeeds_once_filtered() {
let av_media = demux(&read_fixture("h264_aac.ts"));
assert!(
av_media
.tracks
.iter()
.any(|t| matches!(t.spec.config, CodecConfig::Avc { .. }))
);
assert!(
av_media
.tracks
.iter()
.any(|t| matches!(t.spec.config, CodecConfig::Aac { .. }))
);
let data_source = read_fixture("m6-single.ts");
let data_media = demux(&data_source);
let live = live_pmt_es(&data_source);
let section_es = live
.iter()
.find(|es| expected_carriage(es.stream_type) == DataCarriage::Sections)
.expect("m6-single.ts must have a section-carried Data track");
let mut data_track = find_data_track(&data_media, section_es).clone();
let mut tracks = av_media.tracks.clone();
let data_track_id = tracks.iter().map(|t| t.spec.track_id).max().unwrap_or(0) + 1;
data_track.spec.track_id = data_track_id;
let data_stream_type = match &data_track.spec.config {
CodecConfig::Data { stream_type, .. } => *stream_type,
other => panic!("expected CodecConfig::Data, got {other:?}"),
};
tracks.push(data_track);
let mixed = Media::new(tracks, av_media.movie_timescale);
assert_eq!(
mixed.tracks.len(),
3,
"video + audio + one opaque Data track"
);
let err = CmafMux::default()
.package(&mixed)
.expect_err("CmafMux must reject a Media containing a CodecConfig::Data track");
match err {
transmux::Error::UnmuxableDataTrack {
track_id,
stream_type,
} => {
assert_eq!(track_id, data_track_id, "error must name the Data track");
assert_eq!(
stream_type, data_stream_type,
"error must carry the Data track's stream_type"
);
}
other => panic!("expected UnmuxableDataTrack, got {other:?}"),
}
let av_only = mixed
.select_tracks_by(|t| !matches!(t.spec.config, CodecConfig::Data { .. }))
.expect("select_tracks_by must keep the 2 carriable tracks");
assert_eq!(av_only.tracks.len(), 2, "video + audio only, once filtered");
let out = CmafMux::default()
.package(&av_only)
.expect("CmafMux must succeed once the Data track is explicitly filtered out");
let reparsed: Media = Fmp4Demux::new()
.unpackage(&out)
.expect("re-parse the fMP4 output");
assert_eq!(
reparsed.tracks.len(),
2,
"only video+audio survive, got {:?}",
reparsed
.tracks
.iter()
.map(|t| &t.spec.config)
.collect::<Vec<_>>()
);
assert!(
reparsed
.tracks
.iter()
.any(|t| matches!(t.spec.config, CodecConfig::Avc { .. })),
"the video track must survive"
);
assert!(
reparsed
.tracks
.iter()
.any(|t| matches!(t.spec.config, CodecConfig::Aac { .. })),
"the audio track must survive"
);
assert!(
!reparsed
.tracks
.iter()
.any(|t| matches!(t.spec.config, CodecConfig::Data { .. })),
"no Data track may survive into the fMP4 output"
);
}
#[test]
fn ts_hls_carries_every_data_and_section_track_in_every_segment_pmt() {
let data = read_fixture("m6-single.ts");
let media = demux(&data);
let live = live_pmt_es(&data);
assert_eq!(live.len(), 6, "sanity: 6 live PMT-listed PIDs (see test 1)");
let out = TsHlsPackager::new(1)
.package(&media)
.expect("TS-HLS packaging must carry every Data/section/E-AC-3 track, not error");
assert!(
!out.segments.is_empty(),
"must produce at least one segment"
);
assert!(out.playlist.starts_with("#EXTM3U"));
const STREAM_TYPE_EAC3: u8 = 0x87;
enum Expect {
Data(u8, Vec<u8>),
Eac3,
}
let track_ids: Vec<Expect> = media
.tracks
.iter()
.map(|t| match &t.spec.config {
CodecConfig::Data {
stream_type,
descriptors,
..
} => Expect::Data(*stream_type, descriptors.clone()),
CodecConfig::Eac3 { .. } => Expect::Eac3,
other => panic!("m6-single.ts must produce only Data or E-AC-3 tracks, got {other:?}"),
})
.collect();
for (i, seg) in out.segments.iter().enumerate() {
assert_eq!(seg.len() % TS, 0, "segment {i} must be whole TS packets");
let seg_pmt_pid = find_pmt_pid(seg);
let seg_es = collect_pmt_es(seg, seg_pmt_pid);
for expect in &track_ids {
match expect {
Expect::Data(stream_type, descriptors) => assert!(
seg_es
.iter()
.any(|(st, _pid, d)| st == stream_type && d == descriptors),
"segment {i}'s PMT must list stream_type {stream_type:#04X} \
with its preserved ES_info descriptors, got {seg_es:?}"
),
Expect::Eac3 => assert!(
seg_es.iter().any(|(st, ..)| *st == STREAM_TYPE_EAC3),
"segment {i}'s PMT must list an E-AC-3 (stream_type 0x87) \
elementary stream, got {seg_es:?}"
),
}
}
}
let mut concat = Vec::new();
for seg in &out.segments {
concat.extend_from_slice(seg);
}
let media2 = demux(&concat);
for track in &media.tracks {
let orig_payloads: Vec<&[u8]> = track.samples.iter().map(|s| s.data.as_ref()).collect();
match &track.spec.config {
CodecConfig::Data {
stream_type,
descriptors,
..
} => {
let round = find_data_track_exact(&media2, *stream_type, descriptors);
let round_payloads: Vec<&[u8]> =
round.samples.iter().map(|s| s.data.as_ref()).collect();
assert_eq!(
orig_payloads, round_payloads,
"stream_type {stream_type:#04X}: payload-lossless through TS-HLS segmentation"
);
}
CodecConfig::Eac3 { .. } => {
let first = orig_payloads
.first()
.expect("E-AC-3 track must have at least one sample");
let round = media2
.tracks
.iter()
.find(|t| {
matches!(t.spec.config, CodecConfig::Eac3 { .. })
&& t.samples.first().map(|s| s.data.as_ref()) == Some(*first)
})
.unwrap_or_else(|| panic!("no re-demuxed E-AC-3 track matching first sample"));
let round_payloads: Vec<&[u8]> =
round.samples.iter().map(|s| s.data.as_ref()).collect();
assert_eq!(
orig_payloads, round_payloads,
"E-AC-3 track: payload-lossless through TS-HLS segmentation"
);
}
_ => unreachable!("checked above"),
}
}
}