const BD_TS_PACKET_SIZE: usize = 192;
const TS_PACKET_SIZE: usize = 188;
const SYNC_BYTE: u8 = 0x47;
#[derive(Debug)]
pub struct PesPacket {
pub pid: u16,
pub pts: Option<i64>,
pub dts: Option<i64>,
pub data: Vec<u8>,
}
struct PesAssembler {
pid: u16,
buffer: Vec<u8>,
pts: Option<i64>,
dts: Option<i64>,
active: bool,
header_remaining: usize,
}
const PES_BUFFER_INIT_CAP: usize = 16 * 1024;
impl PesAssembler {
fn new(pid: u16) -> Self {
Self {
pid,
buffer: Vec::with_capacity(PES_BUFFER_INIT_CAP),
pts: None,
dts: None,
active: false,
header_remaining: 0,
}
}
fn start(&mut self, pts: Option<i64>, dts: Option<i64>) -> Option<PesPacket> {
let completed = if self.active && !self.buffer.is_empty() {
Some(PesPacket {
pid: self.pid,
pts: self.pts,
dts: self.dts,
data: std::mem::replace(&mut self.buffer, Vec::with_capacity(PES_BUFFER_INIT_CAP)),
})
} else {
self.buffer.clear();
None
};
self.pts = pts;
self.dts = dts;
self.active = true;
completed
}
fn push(&mut self, data: &[u8]) {
if self.active {
self.buffer.extend_from_slice(data);
}
}
fn flush(&mut self) -> Option<PesPacket> {
if self.active && !self.buffer.is_empty() {
self.active = false;
Some(PesPacket {
pid: self.pid,
pts: self.pts,
dts: self.dts,
data: std::mem::take(&mut self.buffer),
})
} else {
None
}
}
}
pub struct TsDemuxer {
assemblers: Vec<PesAssembler>,
pid_index: Vec<i16>, remainder: Vec<u8>, }
impl TsDemuxer {
pub fn new(pids: &[u16]) -> Self {
debug_assert!(
pids.len() <= i16::MAX as usize,
"TsDemuxer: too many PIDs for an i16 index table"
);
let max_pid = pids.iter().copied().max().unwrap_or(0) as usize;
let table_size = (max_pid + 1).max(8192);
let mut pid_index = vec![-1i16; table_size];
let mut assemblers = Vec::with_capacity(pids.len());
for (i, &pid) in pids.iter().enumerate() {
pid_index[pid as usize] = i as i16;
assemblers.push(PesAssembler::new(pid));
}
Self {
assemblers,
pid_index,
remainder: Vec::new(),
}
}
pub fn feed(&mut self, data: &[u8]) -> Vec<PesPacket> {
let mut completed = Vec::with_capacity(4);
let mut offset = 0;
if !self.remainder.is_empty() {
let need = BD_TS_PACKET_SIZE - self.remainder.len();
if data.len() < need {
self.remainder.extend_from_slice(data);
return completed;
}
let mut boundary = [0u8; BD_TS_PACKET_SIZE];
boundary[..self.remainder.len()].copy_from_slice(&self.remainder);
boundary[self.remainder.len()..].copy_from_slice(&data[..need]);
self.remainder.clear();
self.process_packet(&boundary, &mut completed);
offset = need;
}
while offset + BD_TS_PACKET_SIZE <= data.len() {
let packet = &data[offset..offset + BD_TS_PACKET_SIZE];
offset += BD_TS_PACKET_SIZE;
self.process_packet(packet, &mut completed);
}
if offset < data.len() {
let leftover = &data[offset..];
if leftover.len() < BD_TS_PACKET_SIZE {
self.remainder.extend_from_slice(leftover);
} else {
self.remainder.clear();
}
}
completed
}
fn process_packet(&mut self, packet: &[u8], completed: &mut Vec<PesPacket>) {
if packet[4] != SYNC_BYTE {
return;
}
let ts = &packet[4..];
let pid = (((ts[1] & 0x1F) as u16) << 8) | ts[2] as u16;
let pusi = ts[1] & 0x40 != 0; let adaptation = (ts[3] >> 4) & 0x03;
let idx = if (pid as usize) < self.pid_index.len() {
self.pid_index[pid as usize]
} else {
-1
};
if idx < 0 {
return;
}
if adaptation == 0x00 {
return;
}
let asm = &mut self.assemblers[idx as usize];
let payload_start = if adaptation == 0x03 || adaptation == 0x02 {
let af_len = ts[4] as usize;
if af_len > 183 {
return; }
5 + af_len
} else {
4
};
if payload_start >= TS_PACKET_SIZE {
return;
}
if adaptation == 0x02 {
return;
}
let payload = &ts[payload_start..];
if pusi {
let (pts, dts, header_len) = parse_pes_header(payload);
if let Some(prev) = asm.start(pts, dts) {
completed.push(prev);
}
if header_len == 0 {
asm.header_remaining = 0;
} else if header_len <= payload.len() {
asm.header_remaining = 0;
if header_len < payload.len() {
asm.push(&payload[header_len..]);
}
} else {
asm.header_remaining = header_len - payload.len();
}
} else if asm.header_remaining > 0 {
let skip = asm.header_remaining.min(payload.len());
asm.header_remaining -= skip;
if skip < payload.len() {
asm.push(&payload[skip..]);
}
} else {
asm.push(payload);
}
}
pub fn flush(&mut self) -> Vec<PesPacket> {
let mut completed = Vec::new();
for asm in &mut self.assemblers {
if let Some(pkt) = asm.flush() {
completed.push(pkt);
}
}
completed
}
}
fn parse_pes_header(data: &[u8]) -> (Option<i64>, Option<i64>, usize) {
if data.len() < 9 || data[0] != 0x00 || data[1] != 0x00 || data[2] != 0x01 {
return (None, None, 0);
}
let stream_id = data[3];
if stream_id == 0xBC
|| stream_id == 0xBE
|| stream_id == 0xBF
|| stream_id == 0xF0
|| stream_id == 0xF1
|| stream_id == 0xF2
|| stream_id == 0xF8
|| stream_id == 0xFF
{
return (None, None, 6);
}
let pts_dts_flags = (data[7] >> 6) & 0x03;
let header_data_len = data[8] as usize;
let header_len = 9 + header_data_len;
let mut pts = None;
let mut dts = None;
if pts_dts_flags >= 2 && header_data_len >= 5 && data.len() >= 14 {
pts = parse_timestamp(&data[9..14]);
}
if pts_dts_flags == 3 && header_data_len >= 10 && data.len() >= 19 {
dts = parse_timestamp(&data[14..19]);
}
(pts, dts, header_len)
}
fn parse_timestamp(data: &[u8]) -> Option<i64> {
if data.len() < 5 {
return None;
}
if (data[0] & 0x01) == 0 || (data[2] & 0x01) == 0 || (data[4] & 0x01) == 0 {
return None;
}
let b0 = data[0] as i64;
let b1 = data[1] as i64;
let b2 = data[2] as i64;
let b3 = data[3] as i64;
let b4 = data[4] as i64;
Some(((b0 >> 1) & 0x07) << 30 | b1 << 22 | (b2 >> 1) << 15 | b3 << 7 | b4 >> 1)
}
fn is_resync_point(data: &[u8], offset: usize) -> bool {
if data.get(offset + 4) != Some(&SYNC_BYTE) {
return false;
}
match data.get(offset + BD_TS_PACKET_SIZE + 4) {
Some(&b) => b == SYNC_BYTE,
None => true, }
}
fn psi_payload_base(pkt: &[u8]) -> Option<usize> {
let afc = (pkt[7] >> 4) & 0x03;
match afc {
0x01 => Some(8), 0x03 => {
let af_len = pkt[8] as usize;
let base = 9 + af_len; if base < BD_TS_PACKET_SIZE {
Some(base)
} else {
None }
}
_ => None,
}
}
fn collect_psi_section(data: &[u8], target_pid: u16, table_id: u8) -> Option<Vec<u8>> {
let mut offset = 0;
while offset + BD_TS_PACKET_SIZE <= data.len() {
if !is_resync_point(data, offset) {
offset += 1;
continue;
}
let pid = (((data[offset + 5] & 0x1F) as u16) << 8) | data[offset + 6] as u16;
let pusi = data[offset + 5] & 0x40 != 0;
if pid == target_pid && pusi {
let Some(payload_off) = psi_payload_base(&data[offset..offset + BD_TS_PACKET_SIZE])
else {
offset += BD_TS_PACKET_SIZE;
continue;
};
let payload = &data[offset + payload_off..offset + BD_TS_PACKET_SIZE];
let pointer = payload[0] as usize;
let sec_start = 1 + pointer;
if sec_start + 3 > payload.len() || payload[sec_start] != table_id {
offset += BD_TS_PACKET_SIZE;
continue;
}
let section_len =
(((payload[sec_start + 1] & 0x0F) as usize) << 8) | payload[sec_start + 2] as usize;
let total = 3 + section_len; let mut section = Vec::with_capacity(total);
section.extend_from_slice(&payload[sec_start..]);
if section.len() >= total {
section.truncate(total);
return Some(section);
}
let mut scan = offset + BD_TS_PACKET_SIZE;
while scan + BD_TS_PACKET_SIZE <= data.len() && section.len() < total {
if data[scan + 4] != SYNC_BYTE {
scan += 1;
continue;
}
let cpid = (((data[scan + 5] & 0x1F) as u16) << 8) | data[scan + 6] as u16;
let cpusi = data[scan + 5] & 0x40 != 0;
if cpid == target_pid && !cpusi {
if let Some(cbase) = psi_payload_base(&data[scan..scan + BD_TS_PACKET_SIZE]) {
section.extend_from_slice(&data[scan + cbase..scan + BD_TS_PACKET_SIZE]);
}
}
scan += BD_TS_PACKET_SIZE;
}
if section.len() >= total {
section.truncate(total);
return Some(section);
}
return None;
}
offset += BD_TS_PACKET_SIZE;
}
None
}
pub fn scan_streams(data: &[u8]) -> Option<Vec<crate::disc::Stream>> {
use crate::disc::*;
let pat = collect_psi_section(data, 0, 0x00)?;
let pat_section_len = (((pat[1] & 0x0F) as usize) << 8) | pat[2] as usize;
if pat_section_len < 4 {
return None;
}
let mut pat_pmt_pid: Option<u16> = None;
{
let entries_start = 8;
let entries_end = (3 + pat_section_len - 4).min(pat.len());
let mut e = entries_start;
while e + 4 <= entries_end {
let prog_num = ((pat[e] as u16) << 8) | pat[e + 1] as u16;
let p = (((pat[e + 2] & 0x1F) as u16) << 8) | pat[e + 3] as u16;
if prog_num != 0 {
pat_pmt_pid = Some(p);
break;
}
e += 4;
}
}
let pmt_pid = pat_pmt_pid?;
let mut streams = Vec::new();
let pmt = collect_psi_section(data, pmt_pid, 0x02)?;
if pmt.len() >= 12 {
let section_len = (((pmt[1] & 0x0F) as usize) << 8) | pmt[2] as usize;
if section_len < 4 {
return None;
}
let prog_info_len = (((pmt[10] & 0x0F) as usize) << 8) | pmt[11] as usize;
let mut pos = 12 + prog_info_len;
let end = (3 + section_len - 4).min(pmt.len());
while pos + 5 <= end {
let stream_type = pmt[pos];
let es_pid = (((pmt[pos + 1] & 0x1F) as u16) << 8) | pmt[pos + 2] as u16;
let es_info_len = (((pmt[pos + 3] & 0x0F) as usize) << 8) | pmt[pos + 4] as usize;
let codec = Codec::from_coding_type(stream_type);
let stream = match codec.kind() {
CodecKind::Video => {
let resolution = match codec {
Codec::Hevc => Resolution::R2160p,
Codec::Mpeg2 => Resolution::R1080i,
_ => Resolution::R1080p,
};
Some(Stream::Video(VideoStream {
pid: es_pid,
codec,
resolution,
frame_rate: FrameRate::Unknown,
hdr: HdrFormat::Sdr,
color_space: ColorSpace::Bt709,
secondary: false,
label: String::new(),
}))
}
CodecKind::Audio => Some(Stream::Audio(AudioStream {
pid: es_pid,
codec,
channels: AudioChannels::Surround51,
language: "und".into(),
sample_rate: SampleRate::S48,
secondary: false,
purpose: crate::disc::LabelPurpose::Normal,
label: String::new(),
})),
CodecKind::Subtitle => Some(Stream::Subtitle(SubtitleStream {
pid: es_pid,
codec,
language: "und".into(),
forced: false,
qualifier: crate::disc::LabelQualifier::None,
codec_data: None,
})),
CodecKind::Unknown => {
tracing::warn!(
target: "mux",
"dropping PMT stream entry with unknown stream_type {:#04x} (PID {:#06x})",
stream_type,
es_pid,
);
None
}
};
if let Some(s) = stream {
streams.push(s);
}
pos += 5 + es_info_len;
}
}
if streams.is_empty() {
None
} else {
Some(streams)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_parse_timestamp() {
let data = [0x21, 0x00, 0x01, 0x00, 0x01];
assert_eq!(parse_timestamp(&data), Some(0));
let data2 = [0x21, 0x00, 0x07, 0xE9, 0x01]; let pts = parse_timestamp(&data2);
assert!(pts.is_some() && pts.unwrap() >= 0);
let bad = [0x00, 0x00, 0x00, 0x00, 0x00]; assert_eq!(parse_timestamp(&bad), None);
}
#[test]
fn test_demuxer_empty() {
let mut demux = TsDemuxer::new(&[0x1011]);
let result = demux.feed(&[]);
assert!(result.is_empty());
}
fn bdts_packet(body: [u8; 184], pid: u16, pusi: bool) -> Vec<u8> {
let mut pkt = vec![0u8; BD_TS_PACKET_SIZE];
pkt[4] = SYNC_BYTE;
pkt[5] = ((pid >> 8) as u8) & 0x1F;
if pusi {
pkt[5] |= 0x40;
}
pkt[6] = (pid & 0xFF) as u8;
pkt[7] = 0x10; pkt[8..8 + 184].copy_from_slice(&body);
pkt
}
fn pat_packet(pmt_pid: u16) -> Vec<u8> {
let mut body = [0xFFu8; 184];
let mut i = 0;
body[i] = 0x00; i += 1;
body[i] = 0x00; body[i + 1] = 0xB0; body[i + 2] = 0x0D; body[i + 3] = 0x00; body[i + 4] = 0x01; body[i + 5] = 0xC1; body[i + 6] = 0x00; body[i + 7] = 0x00; body[i + 8] = 0x00;
body[i + 9] = 0x01;
body[i + 10] = 0xE0 | (((pmt_pid >> 8) as u8) & 0x1F);
body[i + 11] = (pmt_pid & 0xFF) as u8;
let _ = &mut i;
bdts_packet(body, 0, true)
}
fn pmt_packet(pmt_pid: u16, entries: &[(u8, u16)]) -> Vec<u8> {
let mut body = [0xFFu8; 184];
body[0] = 0x00; let s = 1; body[s] = 0x02; let entries_len = entries.len() * 5;
let section_length = 9 + entries_len + 4;
body[s + 1] = 0xB0 | (((section_length >> 8) as u8) & 0x0F);
body[s + 2] = (section_length & 0xFF) as u8;
body[s + 3] = 0x00; body[s + 4] = 0x01; body[s + 5] = 0xC1; body[s + 6] = 0x00; body[s + 7] = 0x00; body[s + 8] = 0xE0; body[s + 9] = 0x00; body[s + 10] = 0xF0; body[s + 11] = 0x00; let mut p = s + 12;
for &(stype, es_pid) in entries {
body[p] = stype;
body[p + 1] = 0xE0 | (((es_pid >> 8) as u8) & 0x1F);
body[p + 2] = (es_pid & 0xFF) as u8;
body[p + 3] = 0xF0; body[p + 4] = 0x00; p += 5;
}
bdts_packet(body, pmt_pid, true)
}
fn data_packet(pid: u16, pusi: bool, payload: &[u8]) -> Vec<u8> {
let mut pkt = vec![0u8; BD_TS_PACKET_SIZE];
pkt[4] = SYNC_BYTE;
pkt[5] = ((pid >> 8) as u8) & 0x1F;
if pusi {
pkt[5] |= 0x40;
}
pkt[6] = (pid & 0xFF) as u8;
pkt[7] = 0x10; let room = TS_PACKET_SIZE - 4; let n = payload.len().min(room);
pkt[8..8 + n].copy_from_slice(&payload[..n]);
pkt
}
fn pmt_packet_with_af(pmt_pid: u16, entries: &[(u8, u16)]) -> Vec<u8> {
let af_len: u8 = 2; let mut pkt = vec![0u8; BD_TS_PACKET_SIZE];
pkt[4] = SYNC_BYTE;
pkt[5] = (((pmt_pid >> 8) as u8) & 0x1F) | 0x40; pkt[6] = (pmt_pid & 0xFF) as u8;
pkt[7] = 0x30; pkt[8] = af_len; pkt[9] = 0x00; pkt[10] = 0xFF; let payload_off = 4 + 4 + 1 + af_len as usize;
let mut body = vec![0xFFu8; BD_TS_PACKET_SIZE - payload_off];
body[0] = 0x00; let s = 1;
body[s] = 0x02; let entries_len = entries.len() * 5;
let section_length = 9 + entries_len + 4;
body[s + 1] = 0xB0 | (((section_length >> 8) as u8) & 0x0F);
body[s + 2] = (section_length & 0xFF) as u8;
body[s + 3] = 0x00;
body[s + 4] = 0x01;
body[s + 5] = 0xC1;
body[s + 6] = 0x00;
body[s + 7] = 0x00;
body[s + 8] = 0xE0;
body[s + 9] = 0x00;
body[s + 10] = 0xF0;
body[s + 11] = 0x00;
let mut p = s + 12;
for &(stype, es_pid) in entries {
body[p] = stype;
body[p + 1] = 0xE0 | (((es_pid >> 8) as u8) & 0x1F);
body[p + 2] = (es_pid & 0xFF) as u8;
body[p + 3] = 0xF0;
body[p + 4] = 0x00;
p += 5;
}
pkt[payload_off..].copy_from_slice(&body);
pkt
}
#[test]
fn short_pes_payload_injects_no_header_bytes() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut garbage = vec![0xAAu8; 32];
garbage[8] = 0x00;
garbage[9] = 0x00;
garbage[10] = 0x01;
garbage[11] = 0x03;
let mut stream = demux.feed(&data_packet(pid, true, &garbage));
assert!(
stream.is_empty(),
"garbage PUSI packet must not complete a PES on its own"
);
let es = [0xDEu8, 0xAD, 0xBE, 0xEF];
stream.extend(demux.feed(&data_packet(pid, false, &es)));
stream.extend(demux.flush());
assert_eq!(stream.len(), 1, "one PES assembled from the continuation");
let pes = &stream[0];
assert!(
pes.data.windows(es.len()).any(|w| w == es),
"continuation ES bytes present, got {:02X?}",
pes.data
);
assert!(
!pes.data.iter().any(|&b| b == 0xAA),
"garbage PES-header bytes must not appear in the elementary stream"
);
assert!(
!pes.data.windows(3).any(|w| w == [0x00, 0x00, 0x01]),
"no injected start code leaked from the malformed PES header"
);
}
#[test]
fn scan_streams_handles_adaptation_field_in_pmt() {
use crate::disc::{Codec, Stream};
let pmt_pid = 0x0100;
let mut data = pat_packet(pmt_pid);
data.extend(pmt_packet_with_af(pmt_pid, &[(0x1B, 0x1011)]));
data.extend(pat_packet(pmt_pid));
let streams = scan_streams(&data).expect("PMT with AF should parse");
assert!(
streams
.iter()
.any(|s| matches!(s, Stream::Video(v) if v.codec == Codec::H264)),
"H.264 video must be found past the adaptation field"
);
}
#[test]
fn scan_streams_maps_lpcm_via_from_coding_type() {
use crate::disc::{Codec, Stream};
let pmt_pid = 0x0100;
let mut data = pat_packet(pmt_pid);
data.extend(pmt_packet(pmt_pid, &[(0x1B, 0x1011), (0x80, 0x1100)]));
let streams = scan_streams(&data).expect("PMT should parse");
assert_eq!(streams.len(), 2, "video + LPCM audio");
let lpcm = streams
.iter()
.find(|s| matches!(s, Stream::Audio(a) if a.pid == 0x1100))
.expect("LPCM audio stream present");
if let Stream::Audio(a) = lpcm {
assert_eq!(a.codec, Codec::Lpcm, "0x80 must map to LPCM");
}
assert!(
streams
.iter()
.any(|s| matches!(s, Stream::Video(v) if v.codec == Codec::H264)),
"H.264 video present"
);
}
fn pmt_two_packets(pmt_pid: u16, entries: &[(u8, u16)]) -> Vec<u8> {
let entries_len = entries.len() * 5;
let section_length = 9 + entries_len + 4; let mut section = Vec::new();
section.push(0x02); section.push(0xB0 | (((section_length >> 8) as u8) & 0x0F));
section.push((section_length & 0xFF) as u8);
section.extend_from_slice(&[0x00, 0x01]); section.push(0xC1); section.push(0x00); section.push(0x00); section.extend_from_slice(&[0xE0, 0x00]); section.extend_from_slice(&[0xF0, 0x00]); for &(stype, es_pid) in entries {
section.push(stype);
section.push(0xE0 | (((es_pid >> 8) as u8) & 0x1F));
section.push((es_pid & 0xFF) as u8);
section.extend_from_slice(&[0xF0, 0x00]); }
section.extend_from_slice(&[0xFF, 0xFF, 0xFF, 0xFF]);
let first_cap = 184 - 1; let head_len = first_cap.min(section.len());
let mut p0 = [0xFFu8; 184];
p0[0] = 0x00; p0[1..1 + head_len].copy_from_slice(§ion[..head_len]);
let pkt0 = bdts_packet(p0, pmt_pid, true);
let mut p1 = [0xFFu8; 184];
let tail = §ion[head_len..];
assert!(!tail.is_empty(), "test must actually span two packets");
p1[..tail.len()].copy_from_slice(tail);
let pkt1 = bdts_packet(p1, pmt_pid, false);
let mut out = pkt0;
out.extend(pkt1);
out
}
#[test]
fn scan_streams_reassembles_pmt_across_packets() {
use crate::disc::{Codec, Stream};
let pmt_pid = 0x0100;
let mut entries: Vec<(u8, u16)> = Vec::new();
entries.push((0x1B, 0x1011)); for i in 0..40u16 {
entries.push((0x80, 0x1100 + i)); }
let mut data = pat_packet(pmt_pid);
data.extend(pmt_two_packets(pmt_pid, &entries));
let streams = scan_streams(&data).expect("multi-packet PMT should parse");
assert_eq!(streams.len(), entries.len(), "every PMT entry reassembled");
assert!(
streams
.iter()
.any(|s| matches!(s, Stream::Video(v) if v.codec == Codec::H264)),
"video survives the split"
);
assert!(
streams.iter().any(
|s| matches!(s, Stream::Audio(a) if a.pid == 0x1100 + 39 && a.codec == Codec::Lpcm)
),
"trailing audio entry from the continuation packet survives"
);
}
fn es_packet_exact(pid: u16, pusi: bool, payload: &[u8]) -> Vec<u8> {
const TS_PAYLOAD: usize = 184;
assert!(payload.len() <= TS_PAYLOAD);
let mut pkt = vec![0u8; BD_TS_PACKET_SIZE];
pkt[4] = SYNC_BYTE;
pkt[5] = ((pid >> 8) as u8) & 0x1F;
if pusi {
pkt[5] |= 0x40;
}
pkt[6] = (pid & 0xFF) as u8;
let pad = TS_PAYLOAD - payload.len();
if pad == 0 {
pkt[7] = 0x10; pkt[8..8 + payload.len()].copy_from_slice(payload);
} else {
pkt[7] = 0x30; let af_field_len = pad - 1; pkt[8] = af_field_len as u8;
if af_field_len >= 1 {
pkt[9] = 0x00; for b in pkt.iter_mut().skip(10).take(af_field_len - 1) {
*b = 0xFF; }
}
let payload_off = 8 + pad;
pkt[payload_off..payload_off + payload.len()].copy_from_slice(payload);
}
pkt
}
fn encode_pts_i64(pts: i64, prefix: u8) -> [u8; 5] {
let p = pts as u64;
[
prefix | (((p >> 30) as u8) & 0x07) << 1 | 1,
((p >> 22) & 0xFF) as u8,
(((p >> 15) & 0x7F) as u8) << 1 | 1,
((p >> 7) & 0xFF) as u8,
(((p) & 0x7F) as u8) << 1 | 1,
]
}
#[test]
fn parse_timestamp_decodes_known_value_90000() {
let enc = encode_pts_i64(90_000, 0x20);
assert_eq!(parse_timestamp(&enc), Some(90_000));
}
#[test]
fn parse_timestamp_max_33bit_value() {
let max = (1i64 << 33) - 1;
let enc = encode_pts_i64(max, 0x20);
assert_eq!(parse_timestamp(&enc), Some(max));
}
#[test]
fn parse_timestamp_rejects_each_missing_marker_bit() {
let good = encode_pts_i64(12_345, 0x20);
for &byte_idx in &[0usize, 2, 4] {
let mut bad = good;
bad[byte_idx] &= 0xFE; assert_eq!(
parse_timestamp(&bad),
None,
"marker bit cleared in byte {byte_idx} must reject"
);
}
for &byte_idx in &[1usize, 3] {
let mut still_ok = good;
still_ok[byte_idx] &= 0xFE;
assert!(
parse_timestamp(&still_ok).is_some(),
"byte {byte_idx} has no marker bit; clearing LSB must still parse"
);
}
}
#[test]
fn parse_timestamp_too_short_returns_none() {
assert_eq!(parse_timestamp(&[0x21, 0x00, 0x01, 0x00]), None);
assert_eq!(parse_timestamp(&[]), None);
}
#[test]
fn parse_pes_header_rejects_bad_start_code() {
let mut buf = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x80, 0x05];
buf.extend_from_slice(&encode_pts_i64(0, 0x20));
let (pts, dts, hl) = parse_pes_header(&buf);
assert!(pts.is_some() && dts.is_none() && hl == 14);
buf[2] = 0x02;
assert_eq!(parse_pes_header(&buf), (None, None, 0));
}
#[test]
fn parse_pes_header_too_short_is_malformed() {
let short = [0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x80];
assert_eq!(parse_pes_header(&short), (None, None, 0));
}
#[test]
fn parse_pes_header_extension_less_stream_ids_report_len_6() {
for sid in [0xBCu8, 0xBE, 0xBF, 0xF0, 0xF1, 0xF2, 0xF8, 0xFF] {
let buf = [0x00, 0x00, 0x01, sid, 0x00, 0x00, 0x80, 0xC0, 0x0A];
let (pts, dts, hl) = parse_pes_header(&buf);
assert_eq!(
(pts, dts, hl),
(None, None, 6),
"stream_id {sid:#04x} must be extension-less (len 6, no timestamps)"
);
}
}
#[test]
fn parse_pes_header_pts_only_vs_pts_dts() {
let mut pts_only = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x80, 0x05];
pts_only.extend_from_slice(&encode_pts_i64(90_000, 0x20));
let (p, d, hl) = parse_pes_header(&pts_only);
assert_eq!((p, d, hl), (Some(90_000), None, 14));
let mut both = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0xC0, 0x0A];
both.extend_from_slice(&encode_pts_i64(180_000, 0x30));
both.extend_from_slice(&encode_pts_i64(90_000, 0x10));
let (p, d, hl) = parse_pes_header(&both);
assert_eq!((p, d, hl), (Some(180_000), Some(90_000), 19));
}
#[test]
fn parse_pes_header_dts_flag_without_room_skips_dts() {
let mut buf = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0xC0, 0x05];
buf.extend_from_slice(&encode_pts_i64(90_000, 0x30));
buf.extend_from_slice(&[0xAA; 10]);
let (p, d, hl) = parse_pes_header(&buf);
assert_eq!(p, Some(90_000), "PTS present");
assert_eq!(d, None, "DTS dropped: header_data_length too short for it");
assert_eq!(hl, 14, "header_len = 9 + header_data_length(5)");
}
#[test]
fn parse_pes_header_len_is_uncapped() {
let buf = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x80, 200];
let (_, _, hl) = parse_pes_header(&buf);
assert_eq!(
hl,
9 + 200,
"header_len uncapped at 209 even though slice is 9"
);
}
#[test]
fn untracked_pid_produces_nothing() {
let mut demux = TsDemuxer::new(&[0x1011]);
let mut pes = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
pes.extend_from_slice(&[0xDE, 0xAD]);
let out = demux.feed(&data_packet(0x1012, true, &pes)); assert!(out.is_empty());
assert!(demux.flush().is_empty());
}
#[test]
fn bad_sync_byte_skips_packet() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut pkt = data_packet(pid, true, &{
let mut v = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
v.extend_from_slice(&[0x11, 0x22, 0x33]);
v
});
pkt[4] = 0x46; let out = demux.feed(&pkt);
assert!(
out.is_empty(),
"bad sync byte must drop the packet entirely"
);
assert!(demux.flush().is_empty());
}
#[test]
fn afc_reserved_zero_drops_payload() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut pkt = data_packet(pid, true, &{
let mut v = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
v.extend_from_slice(&[0xCA, 0xFE]);
v
});
pkt[7] = 0x00; let out = demux.feed(&pkt);
assert!(out.is_empty());
assert!(demux.flush().is_empty(), "reserved AFC contributes no ES");
}
#[test]
fn afc_adaptation_only_carries_no_payload() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut start = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
start.extend_from_slice(&[0x01, 0x02, 0x03, 0x04]);
demux.feed(&es_packet_exact(pid, true, &start));
let mut afonly = vec![0u8; BD_TS_PACKET_SIZE];
afonly[4] = SYNC_BYTE;
afonly[5] = ((pid >> 8) as u8) & 0x1F; afonly[6] = (pid & 0xFF) as u8;
afonly[7] = 0x20; afonly[8] = 5; for b in afonly.iter_mut().skip(9).take(183) {
*b = 0xEE; }
demux.feed(&afonly);
let out = demux.flush();
assert_eq!(out.len(), 1);
assert!(
!out[0].data.iter().any(|&b| b == 0xEE),
"AF-only packet bytes must never be appended as ES"
);
assert_eq!(out[0].data, vec![0x01, 0x02, 0x03, 0x04]);
}
#[test]
fn adaptation_field_len_skipped_before_payload() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut pkt = vec![0u8; BD_TS_PACKET_SIZE];
pkt[4] = SYNC_BYTE;
pkt[5] = (((pid >> 8) as u8) & 0x1F) | 0x40; pkt[6] = (pid & 0xFF) as u8;
pkt[7] = 0x30; let pes = [
0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00, 0x77, 0x88,
];
let payload_area = 184usize;
let af_total = payload_area - pes.len(); let af_field_len = af_total - 1; pkt[8] = af_field_len as u8;
pkt[9] = 0x00; for b in pkt.iter_mut().skip(10).take(af_field_len - 1) {
*b = 0xBB; }
let payload_off = 4 + 4 + af_total;
pkt[payload_off..payload_off + pes.len()].copy_from_slice(&pes);
demux.feed(&pkt);
let out = demux.flush();
assert_eq!(out.len(), 1);
assert_eq!(out[0].data, vec![0x77, 0x88]);
assert!(
!out[0].data.iter().any(|&b| b == 0xBB),
"adaptation-field stuffing must not appear in the ES"
);
}
#[test]
fn malformed_af_length_over_183_drops_packet() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut pkt = vec![0u8; BD_TS_PACKET_SIZE];
pkt[4] = SYNC_BYTE;
pkt[5] = (((pid >> 8) as u8) & 0x1F) | 0x40;
pkt[6] = (pid & 0xFF) as u8;
pkt[7] = 0x30; pkt[8] = 184; let out = demux.feed(&pkt);
assert!(out.is_empty());
assert!(demux.flush().is_empty());
}
#[test]
fn pes_reassembled_from_continuation_packets() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut start = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
start.extend_from_slice(&[0xA1, 0xA2]);
let mut out = demux.feed(&es_packet_exact(pid, true, &start));
assert!(out.is_empty(), "first PES not yet completed");
out.extend(demux.feed(&es_packet_exact(pid, false, &[0xB1, 0xB2])));
out.extend(demux.feed(&es_packet_exact(pid, false, &[0xC1, 0xC2])));
let mut start2 = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
start2.extend_from_slice(&[0xD1]);
out.extend(demux.feed(&es_packet_exact(pid, true, &start2)));
assert_eq!(out.len(), 1, "previous PES completed by new PUSI");
assert_eq!(out[0].data, vec![0xA1, 0xA2, 0xB1, 0xB2, 0xC1, 0xC2]);
out.extend(demux.flush());
assert_eq!(out.last().unwrap().data, vec![0xD1]);
}
#[test]
fn pes_header_spanning_two_packets_is_fully_skipped() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut start = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 184];
start.extend(std::iter::repeat_n(0xAAu8, 175)); demux.feed(&es_packet_exact(pid, true, &start));
let mut cont = vec![0xAAu8; 9]; cont.extend_from_slice(&[0xEF, 0xBE]); demux.feed(&es_packet_exact(pid, false, &cont));
let out = demux.flush();
assert_eq!(out.len(), 1);
assert_eq!(
out[0].data,
vec![0xEF, 0xBE],
"only post-header ES survives; spillover header bytes skipped"
);
}
#[test]
fn unaligned_feed_reassembles_across_call_boundary() {
let pid = 0x1011;
let mut full = es_packet_exact(pid, true, &{
let mut v = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
v.extend_from_slice(&[0x10, 0x20, 0x30, 0x40]);
v
});
full.extend(es_packet_exact(pid, false, &[0x50, 0x60]));
let mut demux = TsDemuxer::new(&[pid]);
let cut = 100;
let mut out = demux.feed(&full[..cut]);
out.extend(demux.feed(&full[cut..]));
out.extend(demux.flush());
assert_eq!(out.len(), 1);
assert_eq!(out[0].data, vec![0x10, 0x20, 0x30, 0x40, 0x50, 0x60]);
}
#[test]
fn feed_holds_sub_packet_remainder_without_emitting() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let pkt = es_packet_exact(pid, true, &{
let mut v = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
v.extend_from_slice(&[0xAB, 0xCD]);
v
});
let out1 = demux.feed(&pkt[..50]); assert!(out1.is_empty());
let out2 = demux.feed(&pkt[50..100]); assert!(out2.is_empty(), "sub-packet remainder must not emit");
let mut out = demux.feed(&pkt[100..]);
out.extend(demux.flush());
assert_eq!(out.len(), 1);
assert_eq!(out[0].data, vec![0xAB, 0xCD]);
}
#[test]
fn two_pids_route_independently_no_collision() {
let (v, a) = (0x1011u16, 0x1100u16);
let mut demux = TsDemuxer::new(&[v, a]);
let mut vstart = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
vstart.extend_from_slice(&[0x11, 0x11]);
let mut astart = vec![0x00, 0x00, 0x01, 0xBD, 0x00, 0x00, 0x80, 0x00, 0x00];
astart.extend_from_slice(&[0x22, 0x22]);
let mut out = Vec::new();
out.extend(demux.feed(&es_packet_exact(v, true, &vstart)));
out.extend(demux.feed(&es_packet_exact(a, true, &astart)));
out.extend(demux.feed(&es_packet_exact(v, false, &[0x33])));
out.extend(demux.feed(&es_packet_exact(a, false, &[0x44])));
out.extend(demux.flush());
let vpes = out.iter().find(|p| p.pid == v).unwrap();
let apes = out.iter().find(|p| p.pid == a).unwrap();
assert_eq!(
vpes.data,
vec![0x11, 0x11, 0x33],
"video ES not contaminated"
);
assert_eq!(
apes.data,
vec![0x22, 0x22, 0x44],
"audio ES not contaminated"
);
}
#[test]
fn pusi_with_pts_is_extracted() {
let pid = 0x1011;
let mut demux = TsDemuxer::new(&[pid]);
let mut pes = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x80, 0x05];
pes.extend_from_slice(&encode_pts_i64(90_000, 0x20));
pes.extend_from_slice(&[0xFE, 0xED]);
demux.feed(&es_packet_exact(pid, true, &pes));
let out = demux.flush();
assert_eq!(out.len(), 1);
assert_eq!(out[0].pts, Some(90_000));
assert_eq!(out[0].data, vec![0xFE, 0xED]);
}
#[test]
fn flush_on_empty_assembler_yields_nothing() {
let mut demux = TsDemuxer::new(&[0x1011]);
assert!(demux.flush().is_empty());
}
#[test]
fn new_with_empty_pids_tracks_nothing() {
let mut demux = TsDemuxer::new(&[]);
let mut pes = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
pes.extend_from_slice(&[0xAA]);
assert!(demux.feed(&data_packet(0x1011, true, &pes)).is_empty());
assert!(demux.flush().is_empty());
}
#[test]
fn high_pid_above_table_floor_is_tracked() {
let pid = 0x1FFFu16; let mut demux = TsDemuxer::new(&[pid]);
let mut pes = vec![0x00, 0x00, 0x01, 0xE0, 0x00, 0x00, 0x80, 0x00, 0x00];
pes.extend_from_slice(&[0x5A, 0xA5]);
demux.feed(&es_packet_exact(pid, true, &pes));
let out = demux.flush();
assert_eq!(out.len(), 1);
assert_eq!(out[0].pid, pid);
assert_eq!(out[0].data, vec![0x5A, 0xA5]);
}
#[test]
fn scan_streams_no_pat_returns_none() {
let data = vec![0u8; BD_TS_PACKET_SIZE * 2]; assert!(scan_streams(&data).is_none());
}
#[test]
fn scan_streams_pat_but_no_pmt_returns_none() {
let pmt_pid = 0x0100;
let mut data = pat_packet(pmt_pid);
data.extend(pat_packet(pmt_pid)); assert!(scan_streams(&data).is_none());
}
#[test]
fn scan_streams_drops_unknown_stream_type() {
use crate::disc::Stream;
let pmt_pid = 0x0100;
let mut data = pat_packet(pmt_pid);
data.extend(pmt_packet(pmt_pid, &[(0x1B, 0x1011), (0x7F, 0x1500)]));
data.extend(pat_packet(pmt_pid)); let streams = scan_streams(&data).expect("known stream survives");
assert_eq!(streams.len(), 1, "unknown stream_type entry dropped");
assert!(matches!(streams[0], Stream::Video(_)));
}
#[test]
fn scan_streams_hevc_defaults_to_uhd_resolution() {
use crate::disc::{Resolution, Stream};
let pmt_pid = 0x0100;
let mut data = pat_packet(pmt_pid);
data.extend(pmt_packet(pmt_pid, &[(0x24, 0x1011)])); data.extend(pat_packet(pmt_pid));
let streams = scan_streams(&data).expect("HEVC video parses");
let v = streams
.iter()
.find_map(|s| match s {
Stream::Video(v) => Some(v),
_ => None,
})
.expect("video present");
assert_eq!(v.resolution, Resolution::R2160p, "HEVC defaults to UHD");
}
#[test]
fn scan_streams_mpeg2_defaults_to_1080i() {
use crate::disc::{Resolution, Stream};
let pmt_pid = 0x0100;
let mut data = pat_packet(pmt_pid);
data.extend(pmt_packet(pmt_pid, &[(0x02, 0x1011)])); data.extend(pat_packet(pmt_pid));
let streams = scan_streams(&data).expect("MPEG-2 video parses");
let v = streams
.iter()
.find_map(|s| match s {
Stream::Video(v) => Some(v),
_ => None,
})
.expect("video present");
assert_eq!(v.resolution, Resolution::R1080i, "MPEG-2 defaults to 1080i");
}
}