use broadcast_common::{Unpackage, crc32_mpeg2};
use mpeg_ts::ts::{TS_PACKET_SIZE, TsHeader};
use transmux::pipeline::TrackSpec;
use transmux::{AbandonReason, CodecConfig, DemuxEvent, StreamingTsDemux, TsDemux};
const PAT_PID: u16 = 0x0000;
const PMT_PID: u16 = 0x1000;
const PMT_PID_2: u16 = 0x1001;
const TRANSPORT_STREAM_ID: u16 = 1;
const PROGRAM_NUMBER: u16 = 1;
const PROGRAM_NUMBER_2: u16 = 2;
const PID_V: u16 = 0x0101; const PID_A: u16 = 0x0102; const PID_B: u16 = 0x0103; const PID_C: u16 = 0x0104;
const STREAM_TYPE_V: u8 = 0x90;
const STREAM_TYPE_A: u8 = 0x91;
const STREAM_TYPE_B: u8 = 0x92;
const STREAM_TYPE_C: u8 = 0x93;
const ISO_639_LANGUAGE_DESCRIPTOR_TAG: u8 = 0x0A;
fn lang_descriptor(lang: &[u8; 3], audio_type: u8) -> Vec<u8> {
let mut d = vec![ISO_639_LANGUAGE_DESCRIPTOR_TAG, 4];
d.extend_from_slice(lang);
d.push(audio_type);
d
}
fn finish_section(table_id: u8, body: &[u8]) -> Vec<u8> {
const CRC32_LEN: usize = 4;
let section_length = body.len() + CRC32_LEN;
let mut section = Vec::with_capacity(3 + section_length);
section.push(table_id);
section.push(0xB0 | ((section_length >> 8) as u8 & 0x0F));
section.push((section_length & 0xFF) as u8);
section.extend_from_slice(body);
let crc = crc32_mpeg2::compute(§ion);
section.extend_from_slice(&crc.to_be_bytes());
section
}
fn build_pat() -> Vec<u8> {
build_pat_full(0, true, &[(PROGRAM_NUMBER, PMT_PID)])
}
fn build_pat_full(version: u8, current_next_indicator: bool, programs: &[(u16, u16)]) -> Vec<u8> {
const TABLE_ID_PAT: u8 = 0x00;
let mut body = Vec::new();
body.extend_from_slice(&TRANSPORT_STREAM_ID.to_be_bytes());
body.push(0xC0 | (version << 1) | u8::from(current_next_indicator));
body.push(0); body.push(0); for &(program_number, pmt_pid) in programs {
body.extend_from_slice(&program_number.to_be_bytes());
body.push(0xE0 | ((pmt_pid >> 8) as u8 & 0x1F));
body.push((pmt_pid & 0xFF) as u8);
}
finish_section(TABLE_ID_PAT, &body)
}
fn build_pmt(version: u8, current_next_indicator: bool, entries: &[(u16, u8, &[u8])]) -> Vec<u8> {
build_pmt_for(PROGRAM_NUMBER, version, current_next_indicator, entries)
}
fn build_pmt_for(
program_number: u16,
version: u8,
current_next_indicator: bool,
entries: &[(u16, u8, &[u8])],
) -> Vec<u8> {
const TABLE_ID_PMT: u8 = 0x02;
const NO_PCR_PID: u16 = 0x1FFF;
let mut body = Vec::new();
body.extend_from_slice(&program_number.to_be_bytes()); body.push(0xC0 | (version << 1) | u8::from(current_next_indicator));
body.push(0); body.push(0); body.push(0xE0 | ((NO_PCR_PID >> 8) as u8 & 0x1F));
body.push((NO_PCR_PID & 0xFF) as u8);
body.push(0xF0); body.push(0);
for &(pid, stream_type, descriptors) in entries {
body.push(stream_type);
body.push(0xE0 | ((pid >> 8) as u8 & 0x1F));
body.push((pid & 0xFF) as u8);
let len = descriptors.len();
body.push(0xF0 | ((len >> 8) as u8 & 0x0F));
body.push((len & 0xFF) as u8);
body.extend_from_slice(descriptors);
}
finish_section(TABLE_ID_PMT, &body)
}
fn ts_packet(pid: u16, pusi: bool, cc: u8, payload: &[u8]) -> [u8; TS_PACKET_SIZE] {
let mut pkt = [0xFFu8; TS_PACKET_SIZE];
let header = TsHeader {
tei: false,
pusi,
pid,
scrambling: 0,
has_adaptation: false,
has_payload: true,
continuity_counter: cc,
};
header
.serialize_into(&mut pkt)
.expect("4-byte TsHeader always fits");
let n = payload.len().min(TS_PACKET_SIZE - 4);
assert!(
payload.len() <= TS_PACKET_SIZE - 4,
"fixture payload must fit in one TS packet"
);
pkt[4..4 + n].copy_from_slice(&payload[..n]);
pkt
}
fn psi_packet(pid: u16, cc: u8, section: &[u8]) -> [u8; TS_PACKET_SIZE] {
let mut payload = Vec::with_capacity(1 + section.len());
payload.push(0); payload.extend_from_slice(section);
ts_packet(pid, true, cc, &payload)
}
fn pes_bytes(payload: &[u8]) -> Vec<u8> {
let mut b = Vec::with_capacity(6 + payload.len());
b.extend_from_slice(&[0x00, 0x00, 0x01, 0xBE]);
b.extend_from_slice(&(payload.len() as u16).to_be_bytes());
b.extend_from_slice(payload);
b
}
fn pes_bytes_with_pts(payload: &[u8], pts: u64) -> Vec<u8> {
const OPTIONAL_HEADER_LEN: usize = 3 + 5;
let mut b = Vec::with_capacity(6 + OPTIONAL_HEADER_LEN + payload.len());
b.extend_from_slice(&[0x00, 0x00, 0x01, 0xBD]);
b.extend_from_slice(&((OPTIONAL_HEADER_LEN + payload.len()) as u16).to_be_bytes());
b.push(0x80); b.push(0x80); b.push(5); b.push(0x21 | ((((pts >> 30) & 0x07) as u8) << 1));
b.push(((pts >> 22) & 0xFF) as u8);
b.push(((((pts >> 15) & 0x7F) as u8) << 1) | 0x01);
b.push(((pts >> 7) & 0xFF) as u8);
b.push((((pts & 0x7F) as u8) << 1) | 0x01);
b.extend_from_slice(payload);
b
}
fn pes_ts_packets(pid: u16, cc_start: u8, pes: &[u8]) -> Vec<u8> {
let mut out = Vec::new();
let mut off = 0usize;
let mut cc = cc_start;
while off < pes.len() {
let n = (pes.len() - off).min(TS_PACKET_SIZE - 4);
out.extend_from_slice(&ts_packet(pid, off == 0, cc & 0x0F, &pes[off..off + n]));
off += n;
cc = cc.wrapping_add(1);
}
out
}
fn corrupt_crc(mut section: Vec<u8>) -> Vec<u8> {
let last = section.len() - 1;
section[last] ^= 0x01;
section
}
fn feed_locked(demux: &mut StreamingTsDemux, bytes: &[u8]) {
let mut buf = bytes.to_vec();
while buf.len() / TS_PACKET_SIZE < mpeg_ts::resync::LOCK_CONFIRMATIONS + 1 {
buf.extend_from_slice(&ts_packet(0x1FFF, false, 0, &[]));
}
demux.feed(&buf);
}
fn feed_bootstrap(demux: &mut StreamingTsDemux, version: u8, entries: &[(u16, u8, &[u8])]) {
let mut bytes = Vec::new();
bytes.extend_from_slice(&psi_packet(PAT_PID, 0, &build_pat()));
bytes.extend_from_slice(&psi_packet(PMT_PID, 0, &build_pmt(version, true, entries)));
for (i, &(pid, _, _)) in entries.iter().enumerate() {
let cc = (i as u8) * 2;
bytes.extend_from_slice(&ts_packet(pid, true, cc, &pes_bytes(b"au1")));
bytes.extend_from_slice(&ts_packet(pid, true, cc + 1, &pes_bytes(b"au2")));
}
while bytes.len() / TS_PACKET_SIZE < mpeg_ts::resync::LOCK_CONFIRMATIONS + 1 {
bytes.extend_from_slice(&ts_packet(0x1FFF, false, 0, &[]));
}
demux.feed(&bytes);
}
fn drain(demux: &mut StreamingTsDemux) -> Vec<DemuxEvent> {
let mut events = Vec::new();
while let Some(ev) = demux.poll_event() {
events.push(ev);
}
events
}
fn bootstrap_three_tracks() -> (StreamingTsDemux, Vec<DemuxEvent>) {
let mut demux = StreamingTsDemux::new();
feed_bootstrap(
&mut demux,
1,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
(PID_B, STREAM_TYPE_B, &[][..]),
],
);
let events = drain(&mut demux);
(demux, events)
}
fn track_added_specs(events: &[DemuxEvent]) -> Vec<TrackSpec> {
events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackAdded(spec) => Some(spec.clone()),
_ => None,
})
.collect()
}
fn track_id_for_pid(specs: &[TrackSpec], pid: u16) -> u32 {
specs
.iter()
.find(|s| s.source_pid == Some(pid))
.unwrap_or_else(|| panic!("no TrackAdded for pid {pid:#06x}"))
.track_id
}
#[test]
fn bootstrap_yields_exactly_three_tracks_no_samples_yet() {
let (_demux, events) = bootstrap_three_tracks();
let specs = track_added_specs(&events);
assert_eq!(
specs.len(),
3,
"expected exactly 3 TrackAdded (video/audioA/audioB), got {events:?}"
);
assert_eq!(
specs.iter().map(|s| s.source_pid).collect::<Vec<_>>(),
vec![Some(PID_V), Some(PID_A), Some(PID_B)],
"TrackAdded must fire in PMT declaration order"
);
assert!(
!events
.iter()
.any(|e| matches!(e, DemuxEvent::Sample { .. })),
"bootstrap's second access unit per PID is deliberately left \
mid-flight (one-behind, unemitted) — no Sample yet"
);
let resolved_count = events
.iter()
.filter(|e| matches!(e, DemuxEvent::TracksResolved { .. }))
.count();
assert_eq!(resolved_count, 1, "TracksResolved must fire exactly once");
}
#[test]
fn carousel_repeat_of_identical_version_emits_nothing() {
let (mut demux, _bootstrap_events) = bootstrap_three_tracks();
let repeat = psi_packet(
PMT_PID,
1,
&build_pmt(
1,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
(PID_B, STREAM_TYPE_B, &[][..]),
],
),
);
demux.feed(&repeat);
let events = drain(&mut demux);
assert!(
events.is_empty(),
"an identical-version PMT repeat must be a complete no-op, got {events:?}"
);
}
#[test]
fn next_table_cni_zero_is_not_applied() {
let (mut demux, _bootstrap_events) = bootstrap_three_tracks();
let next = psi_packet(
PMT_PID,
1,
&build_pmt(
5,
false,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
],
),
);
demux.feed(&next);
let events = drain(&mut demux);
assert!(
events.is_empty(),
"current_next_indicator=0 must never be applied/diffed, got {events:?}"
);
}
#[test]
fn removal_and_update_and_rearm_from_one_version_bump() {
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let specs = track_added_specs(&bootstrap_events);
let track_id_b = track_id_for_pid(&specs, PID_B);
let track_id_a = track_id_for_pid(&specs, PID_A);
let resolved_gen_before = bootstrap_events
.iter()
.find_map(|e| match e {
DemuxEvent::TracksResolved { generation, .. } => Some(*generation),
_ => None,
})
.expect("bootstrap must fire TracksResolved");
let v2 = psi_packet(
PMT_PID,
1,
&build_pmt(
2,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"qaa", 0)),
],
),
);
demux.feed(&v2);
let events = drain(&mut demux);
let removed: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackRemoved {
track_id,
provenance,
..
} => Some((*track_id, provenance.pid)),
_ => None,
})
.collect();
assert_eq!(
removed,
vec![(track_id_b, Some(PID_B))],
"exactly one TrackRemoved, for audioB's real track_id and pid"
);
let updated: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackUpdated(spec) => Some(spec.clone()),
_ => None,
})
.collect();
assert_eq!(
updated.len(),
1,
"exactly one TrackUpdated, got {updated:?}"
);
assert_eq!(updated[0].track_id, track_id_a);
assert_eq!(
updated[0].es_info_descriptors,
lang_descriptor(b"qaa", 0),
"TrackUpdated must carry audioA's NEW descriptor bytes"
);
let resolved_gen_after = events.iter().find_map(|e| match e {
DemuxEvent::TracksResolved { generation, .. } => Some(*generation),
_ => None,
});
assert_eq!(
resolved_gen_after,
Some(resolved_gen_before + 1),
"TracksResolved must re-arm with a bumped generation after the removal, got {events:?}"
);
let stray_pes = ts_packet(PID_B, true, 9, &pes_bytes(b"post-removal"));
demux.feed(&stray_pes);
demux.finish();
let trailing = drain(&mut demux);
assert!(
!trailing
.iter()
.any(|e| matches!(e, DemuxEvent::Sample { track_id, .. } if *track_id == track_id_b)),
"no Sample for the removed track_id may ever follow its TrackRemoved, got {trailing:?}"
);
}
#[test]
fn tracks_resolved_rearms_even_when_known_pid_count_returns_to_its_prior_value() {
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let gen_before = bootstrap_events
.iter()
.find_map(|e| match e {
DemuxEvent::TracksResolved { generation, .. } => Some(*generation),
_ => None,
})
.expect("bootstrap must fire TracksResolved");
let v2 = psi_packet(
PMT_PID,
1,
&build_pmt(
2,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
(PID_C, STREAM_TYPE_C, &[][..]),
],
),
);
demux.feed(&v2);
let mut events = drain(&mut demux);
assert!(
!events
.iter()
.any(|e| matches!(e, DemuxEvent::TracksResolved { .. })),
"must NOT resolve yet — the new PID C hasn't promoted to Live"
);
let mut more = Vec::new();
more.extend_from_slice(&ts_packet(PID_C, true, 20, &pes_bytes(b"au1")));
more.extend_from_slice(&ts_packet(PID_C, true, 21, &pes_bytes(b"au2")));
demux.feed(&more);
events.extend(drain(&mut demux));
let saw_c_added = events
.iter()
.any(|e| matches!(e, DemuxEvent::TrackAdded(spec) if spec.source_pid == Some(PID_C)));
assert!(saw_c_added, "TrackAdded must fire for the new PID C");
let gen_after = events.iter().find_map(|e| match e {
DemuxEvent::TracksResolved { generation, .. } => Some(*generation),
_ => None,
});
assert_eq!(
gen_after,
Some(gen_before + 1),
"TracksResolved must re-fire once the known-PID count (V, A, C = 3, same as \
the original V/A/B count) is fully resolved again — this is the bug a \
count-keyed de-dup key would silently swallow; got {events:?}"
);
}
#[test]
fn version_wrap_31_to_0_is_treated_as_a_change() {
let mut demux = StreamingTsDemux::new();
feed_bootstrap(
&mut demux,
31,
&[(PID_V, STREAM_TYPE_V, &lang_descriptor(b"eng", 0))],
);
let bootstrap_events = drain(&mut demux);
let specs = track_added_specs(&bootstrap_events);
assert_eq!(specs.len(), 1, "expected the single bootstrapped track");
let track_id_v = specs[0].track_id;
let wrapped = psi_packet(
PMT_PID,
1,
&build_pmt(
0,
true,
&[(PID_V, STREAM_TYPE_V, &lang_descriptor(b"qaa", 0))],
),
);
demux.feed(&wrapped);
let events = drain(&mut demux);
let updated: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackUpdated(spec) => Some(spec.clone()),
_ => None,
})
.collect();
assert_eq!(
updated.len(),
1,
"version 31 -> 0 must be treated as a real change (inequality, not `>`), \
got events {events:?}"
);
assert_eq!(updated[0].track_id, track_id_v);
assert_eq!(updated[0].es_info_descriptors, lang_descriptor(b"qaa", 0));
}
#[test]
fn track_abandoned_config_unrecoverable_is_actually_emitted() {
const STREAM_TYPE_AVC: u8 = 0x1B;
let mut demux = StreamingTsDemux::new();
let mut bytes = Vec::new();
bytes.extend_from_slice(&psi_packet(PAT_PID, 0, &build_pat()));
bytes.extend_from_slice(&psi_packet(
PMT_PID,
0,
&build_pmt(1, true, &[(PID_V, STREAM_TYPE_AVC, &[][..])]),
));
for (i, cc) in (0u8..3).enumerate() {
let au = [0x00, 0x00, 0x00, 0x01, 0x41, 0xAA, 0xBB, i as u8];
bytes.extend_from_slice(&pes_ts_packets(PID_V, cc, &pes_bytes(&au)));
}
feed_locked(&mut demux, &bytes);
let before = drain(&mut demux);
assert!(
!before
.iter()
.any(|e| matches!(e, DemuxEvent::TrackAbandoned { .. })),
"abandonment is an end-of-input conclusion here, not a mid-stream one: {before:?}"
);
demux.finish();
let events = drain(&mut demux);
let abandoned: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackAbandoned {
track_id,
reason,
provenance,
..
} => Some((*track_id, *reason, provenance.pid)),
_ => None,
})
.collect();
assert_eq!(
abandoned,
vec![(None, AbandonReason::ConfigUnrecoverable, Some(PID_V))],
"exactly one TrackAbandoned, with no track_id (TrackAdded never fired), \
the ConfigUnrecoverable reason, and the real PID; got {events:?}"
);
assert!(
!events
.iter()
.any(|e| matches!(e, DemuxEvent::TrackAdded(_))),
"a track whose config never resolved must never have been added"
);
}
const AC3_DESCRIPTOR: [u8; 3] = [0x6A, 0x01, 0x00];
const STREAM_TYPE_PES_PRIVATE: u8 = 0x06;
const STREAM_TYPE_SCTE35: u8 = 0x86;
const STREAM_TYPE_AVC: u8 = 0x1B;
fn real_ac3_syncframes() -> Vec<Vec<u8>> {
let path =
std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../fixtures/ts/dolby/ac3.ts");
let bytes =
std::fs::read(&path).unwrap_or_else(|e| panic!("read fixture {}: {e}", path.display()));
let media = TsDemux::new()
.unpackage(&bytes)
.expect("demux fixtures/ts/dolby/ac3.ts");
let track = media
.tracks
.iter()
.find(|t| matches!(t.spec.config, CodecConfig::Ac3 { .. }))
.expect("ac3.ts carries an AC-3 track");
track.samples.iter().map(|s| s.data.to_vec()).collect()
}
#[test]
fn codec_reclassification_rebuilds_the_probe_instead_of_panicking() {
let mut demux = StreamingTsDemux::new();
let mut boot = Vec::new();
boot.extend_from_slice(&psi_packet(PAT_PID, 0, &build_pat()));
boot.extend_from_slice(&psi_packet(
PMT_PID,
0,
&build_pmt(1, true, &[(PID_V, STREAM_TYPE_PES_PRIVATE, &[][..])]),
));
feed_locked(&mut demux, &boot);
let boot_events = drain(&mut demux);
assert!(
!boot_events
.iter()
.any(|e| matches!(e, DemuxEvent::TrackAdded(_))),
"no access unit yet, so the PID must still be probing: {boot_events:?}"
);
demux.feed(&psi_packet(
PMT_PID,
1,
&build_pmt(
2,
true,
&[(PID_V, STREAM_TYPE_PES_PRIVATE, &AC3_DESCRIPTOR[..])],
),
));
drain(&mut demux);
let frames = real_ac3_syncframes();
assert!(
frames.len() >= 3,
"fixture must carry several real syncframes, got {}",
frames.len()
);
for (i, frame) in frames.iter().take(3).enumerate() {
demux.feed(&pes_ts_packets(PID_V, (i as u8) * 8, &pes_bytes(frame)));
}
demux.finish();
let events = drain(&mut demux);
let added: Vec<TrackSpec> = track_added_specs(&events);
assert_eq!(
added.len(),
1,
"the reclassified PID must resolve to exactly one track, got {events:?}"
);
assert_eq!(added[0].source_pid, Some(PID_V));
assert!(
matches!(added[0].config, CodecConfig::Ac3 { .. }),
"the reclassified track must be typed AC-3 (config recovered from the \
real syncframes), got {:?}",
added[0].config
);
assert_eq!(
added[0].es_info_descriptors,
AC3_DESCRIPTOR.to_vec(),
"the new track must carry the v2 ES_info descriptor loop"
);
}
#[test]
fn carrier_is_rebuilt_when_a_section_carried_pid_becomes_pes_carried() {
let mut demux = StreamingTsDemux::new();
let mut boot = Vec::new();
boot.extend_from_slice(&psi_packet(PAT_PID, 0, &build_pat()));
boot.extend_from_slice(&psi_packet(
PMT_PID,
0,
&build_pmt(1, true, &[(PID_V, STREAM_TYPE_SCTE35, &[][..])]),
));
feed_locked(&mut demux, &boot);
drain(&mut demux);
demux.feed(&psi_packet(
PMT_PID,
1,
&build_pmt(2, true, &[(PID_V, STREAM_TYPE_AVC, &[][..])]),
));
drain(&mut demux);
let mut au = Vec::new();
for nal in [
&[0x67u8, 0x42, 0x00, 0x1E, 0xAB, 0x40][..],
&[0x68, 0xCE, 0x3C, 0x80][..],
&[0x65, 0x88, 0x84, 0x00, 0x11, 0x22][..],
] {
au.extend_from_slice(&[0x00, 0x00, 0x00, 0x01]);
au.extend_from_slice(nal);
}
for cc in 0u8..3 {
demux.feed(&pes_ts_packets(PID_V, cc * 4, &pes_bytes(&au)));
}
demux.finish();
let events = drain(&mut demux);
let added = track_added_specs(&events);
assert_eq!(added.len(), 1, "expected one H.264 track, got {events:?}");
assert!(
matches!(added[0].config, CodecConfig::Avc { .. }),
"the reclassified track must be typed AVC, got {:?}",
added[0].config
);
let samples: Vec<_> = events
.iter()
.filter(|e| matches!(e, DemuxEvent::Sample { .. }))
.collect();
assert!(
!samples.is_empty(),
"a PES-carried H.264 track must produce samples — a stale \
SectionReassembler yields exactly zero, silently; got {events:?}"
);
}
#[test]
fn corrupt_crc_pmt_is_dropped_and_disturbs_nothing() {
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let specs = track_added_specs(&bootstrap_events);
let track_id_b = track_id_for_pid(&specs, PID_B);
let bad = build_pmt(
2,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
],
);
demux.feed(&psi_packet(PMT_PID, 1, &corrupt_crc(bad)));
let events = drain(&mut demux);
assert!(
events.is_empty(),
"a CRC-failing PMT must be dropped silently, changing nothing: {events:?}"
);
let good = build_pmt(
2,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
],
);
demux.feed(&psi_packet(PMT_PID, 2, &good));
let events = drain(&mut demux);
let removed: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackRemoved { track_id, .. } => Some(*track_id),
_ => None,
})
.collect();
assert_eq!(
removed,
vec![track_id_b],
"the valid same-version section must still apply after the corrupt one \
was dropped; got {events:?}"
);
}
#[test]
fn corrupt_crc_pat_cannot_hijack_a_live_es_pid() {
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let specs = track_added_specs(&bootstrap_events);
let track_id_b = track_id_for_pid(&specs, PID_B);
let bad = build_pat_full(
1,
true,
&[(PROGRAM_NUMBER, PMT_PID), (PROGRAM_NUMBER_2, PID_B)],
);
demux.feed(&psi_packet(PAT_PID, 1, &corrupt_crc(bad)));
drain(&mut demux);
for cc in 0u8..3 {
demux.feed(&pes_ts_packets(PID_B, cc, &pes_bytes(b"still-audio")));
}
let events = drain(&mut demux);
assert!(
events
.iter()
.any(|e| matches!(e, DemuxEvent::Sample { track_id, .. } if *track_id == track_id_b)),
"a CRC-failing PAT must not divert a live ES PID into PMT reassembly; got {events:?}"
);
}
#[test]
fn next_pat_cni_zero_cannot_hijack_a_live_es_pid() {
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let specs = track_added_specs(&bootstrap_events);
let track_id_b = track_id_for_pid(&specs, PID_B);
let next = build_pat_full(
1,
false,
&[(PROGRAM_NUMBER, PMT_PID), (PROGRAM_NUMBER_2, PID_B)],
);
demux.feed(&psi_packet(PAT_PID, 1, &next));
drain(&mut demux);
for cc in 0u8..3 {
demux.feed(&pes_ts_packets(PID_B, cc, &pes_bytes(b"still-audio")));
}
let events = drain(&mut demux);
assert!(
events
.iter()
.any(|e| matches!(e, DemuxEvent::Sample { track_id, .. } if *track_id == track_id_b)),
"a current_next_indicator=0 PAT must never be applied; got {events:?}"
);
}
#[test]
fn pat_remap_of_a_pmt_pid_is_honoured() {
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let specs = track_added_specs(&bootstrap_events);
let track_id_b = track_id_for_pid(&specs, PID_B);
demux.feed(&psi_packet(
PAT_PID,
1,
&build_pat_full(1, true, &[(PROGRAM_NUMBER_2, PMT_PID)]),
));
drain(&mut demux);
demux.feed(&psi_packet(
PMT_PID,
1,
&build_pmt_for(
PROGRAM_NUMBER_2,
2,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
],
),
));
let events = drain(&mut demux);
let removed: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackRemoved { track_id, .. } => Some(*track_id),
_ => None,
})
.collect();
assert_eq!(
removed,
vec![track_id_b],
"a PMT under the remapped program_number must be applied, not rejected \
forever by a frozen PAT binding; got {events:?}"
);
}
#[test]
fn post_removal_payload_is_not_replayed_into_the_re_added_track() {
const ORPHAN_PTS: u64 = 90_000;
const REPLACEMENT_PTS: u64 = 90_000 * 3600;
const PTS_STEP: u64 = 3_000;
let (mut demux, bootstrap_events) = bootstrap_three_tracks();
let specs = track_added_specs(&bootstrap_events);
let old_track_id_b = track_id_for_pid(&specs, PID_B);
demux.feed(&psi_packet(
PMT_PID,
1,
&build_pmt(
2,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
],
),
));
let removal = drain(&mut demux);
assert!(
removal
.iter()
.any(|e| matches!(e, DemuxEvent::TrackRemoved { track_id, .. } if *track_id == old_track_id_b)),
"sanity: audioB must have been removed; got {removal:?}"
);
for i in 0u8..4 {
let pts = ORPHAN_PTS + u64::from(i) * PTS_STEP;
demux.feed(&pes_ts_packets(
PID_B,
i * 4,
&pes_bytes_with_pts(b"orphan", pts),
));
}
drain(&mut demux);
demux.feed(&psi_packet(
PMT_PID,
2,
&build_pmt(
3,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &lang_descriptor(b"eng", 0)),
(PID_B, STREAM_TYPE_B, &[][..]),
],
),
));
let mut events = drain(&mut demux);
for i in 0u8..3 {
let pts = REPLACEMENT_PTS + u64::from(i) * PTS_STEP;
demux.feed(&pes_ts_packets(
PID_B,
i * 4,
&pes_bytes_with_pts(b"fresh", pts),
));
}
demux.finish();
events.extend(drain(&mut demux));
let new_track_id = events
.iter()
.find_map(|e| match e {
DemuxEvent::TrackAdded(spec) if spec.source_pid == Some(PID_B) => Some(spec.track_id),
_ => None,
})
.expect("the re-added PID must fire a fresh TrackAdded");
assert_ne!(
new_track_id, old_track_id_b,
"a re-added PID is a new track, with a new track_id"
);
let first = events
.iter()
.find_map(|e| match e {
DemuxEvent::Sample {
track_id, sample, ..
} if *track_id == new_track_id => Some(sample),
_ => None,
})
.expect("the re-added track must produce samples");
assert_eq!(
first.data.as_ref(),
b"fresh",
"the re-added track's first sample must be post-re-add traffic, never a \
replayed pre-removal orphan payload"
);
assert_eq!(
first.dts,
Some(REPLACEMENT_PTS as i64),
"start_decode_time must not be anchored in the past by a replayed orphan"
);
}
#[test]
fn es_pid_declared_by_two_programs_survives_one_program_dropping_it() {
let mut demux = StreamingTsDemux::new();
let mut boot = Vec::new();
boot.extend_from_slice(&psi_packet(
PAT_PID,
0,
&build_pat_full(
0,
true,
&[(PROGRAM_NUMBER, PMT_PID), (PROGRAM_NUMBER_2, PMT_PID_2)],
),
));
boot.extend_from_slice(&psi_packet(
PMT_PID,
0,
&build_pmt_for(
PROGRAM_NUMBER,
1,
true,
&[
(PID_V, STREAM_TYPE_V, &[][..]),
(PID_A, STREAM_TYPE_A, &[][..]),
],
),
));
boot.extend_from_slice(&psi_packet(
PMT_PID_2,
0,
&build_pmt_for(
PROGRAM_NUMBER_2,
1,
true,
&[(PID_A, STREAM_TYPE_A, &[][..])],
),
));
for (i, pid) in [PID_V, PID_A].iter().enumerate() {
let cc = (i as u8) * 4;
boot.extend_from_slice(&pes_ts_packets(*pid, cc, &pes_bytes(b"au1")));
boot.extend_from_slice(&pes_ts_packets(*pid, cc + 1, &pes_bytes(b"au2")));
}
feed_locked(&mut demux, &boot);
let bootstrap_events = drain(&mut demux);
let specs = track_added_specs(&bootstrap_events);
let shared_track_id = track_id_for_pid(&specs, PID_A);
demux.feed(&psi_packet(
PMT_PID,
1,
&build_pmt_for(PROGRAM_NUMBER, 2, true, &[(PID_V, STREAM_TYPE_V, &[][..])]),
));
let events = drain(&mut demux);
assert!(
!events
.iter()
.any(|e| matches!(e, DemuxEvent::TrackRemoved { .. })),
"the shared PID is still declared by program 2 — no TrackRemoved yet; got {events:?}"
);
for cc in 0u8..3 {
demux.feed(&pes_ts_packets(PID_A, cc, &pes_bytes(b"still-shared")));
}
let events = drain(&mut demux);
assert!(
events.iter().any(
|e| matches!(e, DemuxEvent::Sample { track_id, .. } if *track_id == shared_track_id)
),
"the surviving shared track must keep producing samples; got {events:?}"
);
demux.feed(&psi_packet(
PMT_PID_2,
1,
&build_pmt_for(PROGRAM_NUMBER_2, 2, true, &[]),
));
let events = drain(&mut demux);
let removed: Vec<_> = events
.iter()
.filter_map(|e| match e {
DemuxEvent::TrackRemoved { track_id, .. } => Some(*track_id),
_ => None,
})
.collect();
assert_eq!(
removed,
vec![shared_track_id],
"the last declaring PMT dropping the PID must remove it; got {events:?}"
);
}