use std::io::{self, Write};
mod packet;
use packet::{Packet, PacketWriter};
const PID_PAT: u16 = 0x0000;
const PID_PMT: u16 = 0x1000;
const PID_VIDEO: u16 = 0x0100;
const PID_AUDIO: u16 = 0x0101;
#[allow(dead_code)]
const PID_NULL: u16 = 0x1FFF;
const PSI_INTERVAL_PACKETS: u64 = 250;
const PCR_INTERVAL_PACKETS: u64 = 40;
const PCR_LEAD_90KHZ: u64 = 90_000 / 5;
use crate::consts::coding_type;
const STREAM_TYPE_HEVC: u8 = coding_type::HEVC;
const STREAM_TYPE_AC3: u8 = coding_type::AC3;
const STREAM_TYPE_TRUEHD: u8 = coding_type::TRUEHD;
#[derive(Debug, Clone, Copy)]
pub enum AudioCodec {
Ac3,
TrueHd,
}
impl AudioCodec {
fn stream_type(self) -> u8 {
match self {
AudioCodec::Ac3 => STREAM_TYPE_AC3,
AudioCodec::TrueHd => STREAM_TYPE_TRUEHD,
}
}
}
pub struct M2tsMux<W: Write> {
out: PacketWriter<W>,
video_codec_private: Option<Vec<u8>>,
params_written: bool,
audio: Option<AudioCodec>,
base_pts_90k: Option<u64>,
cc_video: u8,
cc_audio: u8,
cc_pat: u8,
cc_pmt: u8,
packets_written: u64,
video_packets_since_pcr: u64,
first_video_written: bool,
}
impl<W: Write> M2tsMux<W> {
pub fn new(writer: W) -> Self {
Self {
out: PacketWriter::new(writer),
video_codec_private: None,
params_written: false,
audio: None,
base_pts_90k: None,
cc_video: 0,
cc_audio: 0,
cc_pat: 0,
cc_pmt: 0,
packets_written: 0,
video_packets_since_pcr: 0,
first_video_written: false,
}
}
pub fn set_video_codec_private(&mut self, hvcc: Vec<u8>) {
self.video_codec_private = Some(hvcc);
}
pub fn set_audio(&mut self, codec: AudioCodec) {
self.audio = Some(codec);
}
pub fn write_video(&mut self, pts_ns: i64, keyframe: bool, data: &[u8]) -> io::Result<()> {
let pts_90k = self.base_relative_pts(pts_ns, true);
let pcr = pts_90k.saturating_sub(PCR_LEAD_90KHZ);
let mut es = Vec::with_capacity(data.len() + 64);
if keyframe && !self.params_written {
if let Some(cp) = &self.video_codec_private {
if let Some(params) = super::hevc::hvcc_to_annex_b(cp) {
es.extend_from_slice(¶ms);
}
}
self.params_written = true;
}
super::hevc::append_length_prefixed_as_annex_b(&mut es, data);
let pes = build_video_pes(pts_90k, &es);
self.write_pes(PID_VIDEO, &pes, Some(pcr), keyframe)
}
pub fn write_audio(&mut self, pts_ns: i64, data: &[u8]) -> io::Result<()> {
if self.audio.is_none() {
return Ok(());
}
let pts_90k = self.base_relative_pts(pts_ns, false);
let pes = build_audio_pes(pts_90k, data);
self.write_pes(PID_AUDIO, &pes, None, false)
}
pub fn finish(&mut self) -> io::Result<()> {
self.out.flush()
}
fn base_relative_pts(&mut self, pts_ns: i64, may_seed_base: bool) -> u64 {
let raw_90k = if pts_ns > 0 {
(((pts_ns as u128) * 9 / 100_000) as u64) & 0x1_FFFF_FFFF
} else {
0
};
if may_seed_base {
self.base_pts_90k.get_or_insert(raw_90k);
}
let base = self.base_pts_90k.unwrap_or(raw_90k);
let delta = raw_90k.wrapping_sub(base) & 0x1_FFFF_FFFF;
if delta >= (1 << 32) { 0 } else { delta }
}
fn write_pes(
&mut self,
pid: u16,
pes: &[u8],
pcr: Option<u64>,
is_keyframe_video: bool,
) -> io::Result<()> {
let mut offset = 0;
let mut first = true;
while offset < pes.len() {
self.maybe_emit_psi()?;
let attach_pcr = (pid == PID_VIDEO)
&& (pcr.is_some())
&& (!self.first_video_written
|| self.video_packets_since_pcr >= PCR_INTERVAL_PACKETS);
let attach_rai = first && is_keyframe_video && pid == PID_VIDEO;
let mut af_body: Vec<u8> = if attach_pcr {
build_pcr_adaptation(pcr.unwrap_or(0))
} else {
Vec::new()
};
if attach_rai {
if af_body.is_empty() {
af_body.push(0x40); } else {
af_body[0] |= 0x40; }
}
let remaining = pes.len() - offset;
let (af_present, payload_len, stuffing): (bool, usize, usize) = if !af_body.is_empty() {
let max_payload = 184 - 1 - af_body.len();
let p = remaining.min(max_payload);
let s = max_payload - p;
(true, p, s)
} else if remaining >= 184 {
(false, 184, 0)
} else {
let max_payload = 182;
let p = remaining.min(max_payload);
let s = max_payload - p;
af_body.push(0x00); (true, p, s)
};
let cc = self.advance_cc(pid);
let mut packet = Packet::new();
packet.set_header(pid, first, true, af_present, cc);
if af_present {
packet.append_adaptation(&af_body, stuffing)?;
}
packet.append_payload(&pes[offset..offset + payload_len])?;
self.out.write_packet(&packet)?;
self.packets_written += 1;
if pid == PID_VIDEO {
self.first_video_written = true;
if attach_pcr {
self.video_packets_since_pcr = 0;
} else {
self.video_packets_since_pcr += 1;
}
}
offset += payload_len;
first = false;
}
Ok(())
}
fn advance_cc(&mut self, pid: u16) -> u8 {
let slot = match pid {
PID_VIDEO => &mut self.cc_video,
PID_AUDIO => &mut self.cc_audio,
PID_PAT => &mut self.cc_pat,
PID_PMT => &mut self.cc_pmt,
_ => return 0,
};
let cc = *slot;
*slot = (*slot + 1) & 0x0F;
cc
}
fn maybe_emit_psi(&mut self) -> io::Result<()> {
if self.packets_written == 0 || self.packets_written % PSI_INTERVAL_PACKETS == 0 {
self.emit_pat()?;
self.emit_pmt()?;
}
Ok(())
}
fn emit_pat(&mut self) -> io::Result<()> {
let payload = build_pat(PID_PMT);
let cc = self.advance_cc(PID_PAT);
let mut packet = Packet::new();
packet.set_header(PID_PAT, true, true, false, cc);
packet.append_payload(&payload)?;
packet.pad_to_188();
self.out.write_packet(&packet)?;
self.packets_written += 1;
Ok(())
}
fn emit_pmt(&mut self) -> io::Result<()> {
let payload = build_pmt(self.audio);
let cc = self.advance_cc(PID_PMT);
let mut packet = Packet::new();
packet.set_header(PID_PMT, true, true, false, cc);
packet.append_payload(&payload)?;
packet.pad_to_188();
self.out.write_packet(&packet)?;
self.packets_written += 1;
Ok(())
}
}
fn build_video_pes(pts_90k: u64, es: &[u8]) -> Vec<u8> {
build_pes_packet(
crate::consts::pes_stream_id::VIDEO,
pts_90k,
es,
false,
)
}
fn build_audio_pes(pts_90k: u64, es: &[u8]) -> Vec<u8> {
build_pes_packet(
crate::consts::pes_stream_id::PRIVATE_STREAM_1,
pts_90k,
es,
true,
)
}
fn build_pes_packet(stream_id: u8, pts_90k: u64, es: &[u8], length_in_header: bool) -> Vec<u8> {
let mut out = Vec::with_capacity(es.len() + 14);
out.extend_from_slice(&[0x00, 0x00, 0x01, stream_id]);
let pes_len = 8 + es.len();
if length_in_header && pes_len <= u16::MAX as usize {
out.extend_from_slice(&(pes_len as u16).to_be_bytes());
} else {
out.extend_from_slice(&[0x00, 0x00]);
}
out.push(0x80);
out.push(0x80);
out.push(5);
let pts = pts_90k & 0x1_FFFF_FFFF;
out.push(0x21 | (((pts >> 29) & 0x0E) as u8));
out.push(((pts >> 22) & 0xFF) as u8);
out.push(0x01 | (((pts >> 14) & 0xFE) as u8));
out.push(((pts >> 7) & 0xFF) as u8);
out.push(0x01 | (((pts << 1) & 0xFE) as u8));
out.extend_from_slice(es);
out
}
fn build_pat(pmt_pid: u16) -> Vec<u8> {
let mut section = Vec::new();
section.push(0x00); section.extend_from_slice(&[0xB0, 13]);
section.extend_from_slice(&[0x00, 0x01]); section.push(0xC1); section.push(0x00); section.push(0x00); section.extend_from_slice(&[0x00, 0x01]); let pid_bytes = (0xE000u16 | (pmt_pid & 0x1FFF)).to_be_bytes();
section.extend_from_slice(&pid_bytes);
let crc = mpegts_crc32(§ion);
section.extend_from_slice(&crc.to_be_bytes());
let mut payload = Vec::with_capacity(section.len() + 1);
payload.push(0x00);
payload.extend_from_slice(§ion);
payload
}
fn build_pmt(audio: Option<AudioCodec>) -> Vec<u8> {
let mut section = Vec::new();
section.push(0x02); let len_placeholder = section.len();
section.extend_from_slice(&[0xB0, 0x00]);
section.extend_from_slice(&1u16.to_be_bytes()); section.push(0xC1); section.push(0x00); section.push(0x00); let pcr_pid = (0xE000u16 | (PID_VIDEO & 0x1FFF)).to_be_bytes();
section.extend_from_slice(&pcr_pid);
section.extend_from_slice(&[0xF0, 0x00]);
section.push(STREAM_TYPE_HEVC);
let v_pid = (0xE000u16 | (PID_VIDEO & 0x1FFF)).to_be_bytes();
section.extend_from_slice(&v_pid);
section.extend_from_slice(&[0xF0, 0x00]);
if let Some(codec) = audio {
section.push(codec.stream_type());
let a_pid = (0xE000u16 | (PID_AUDIO & 0x1FFF)).to_be_bytes();
section.extend_from_slice(&a_pid);
section.extend_from_slice(&[0xF0, 0x00]);
}
let section_len = section.len() - 3 + 4;
section[len_placeholder] = 0xB0 | ((section_len >> 8) as u8 & 0x0F);
section[len_placeholder + 1] = section_len as u8;
let crc = mpegts_crc32(§ion);
section.extend_from_slice(&crc.to_be_bytes());
let mut payload = Vec::with_capacity(section.len() + 1);
payload.push(0x00);
payload.extend_from_slice(§ion);
payload
}
fn build_pcr_adaptation(pcr_90k: u64) -> Vec<u8> {
let mut af = vec![0x10]; let pcr_base = pcr_90k & 0x1_FFFF_FFFF; let pcr_ext: u16 = 0; af.push((pcr_base >> 25) as u8);
af.push((pcr_base >> 17) as u8);
af.push((pcr_base >> 9) as u8);
af.push((pcr_base >> 1) as u8);
af.push(((pcr_base << 7) as u8 & 0x80) | 0x7E | ((pcr_ext >> 8) as u8 & 0x01));
af.push(pcr_ext as u8);
af
}
fn mpegts_crc32(data: &[u8]) -> u32 {
let mut crc: u32 = 0xFFFF_FFFF;
for &b in data {
crc ^= (b as u32) << 24;
for _ in 0..8 {
if crc & 0x8000_0000 != 0 {
crc = (crc << 1) ^ 0x04C1_1DB7;
} else {
crc <<= 1;
}
}
}
crc
}
#[cfg(test)]
mod tests {
use super::*;
fn assert_ts_well_formed(buf: &[u8]) {
assert_eq!(
buf.len() % 188,
0,
"stream not packet-aligned: {} bytes",
buf.len()
);
for (i, chunk) in buf.chunks(188).enumerate() {
assert_eq!(chunk[0], 0x47, "packet {} missing sync byte", i);
}
}
fn extract_pids(buf: &[u8]) -> Vec<u16> {
buf.chunks(188)
.map(|p| u16::from_be_bytes([p[1] & 0x1F, p[2]]))
.collect()
}
#[test]
fn crc32_is_self_validating() {
let a = [
0u8, 0xB0, 0x0D, 0x00, 0x01, 0xC1, 0x00, 0x00, 0x00, 0x01, 0xE1, 0x00,
];
let mut b = a;
b[5] ^= 0x01; let crc_a = mpegts_crc32(&a);
let crc_b = mpegts_crc32(&b);
assert_ne!(crc_a, crc_b);
assert_eq!(crc_a, mpegts_crc32(&a)); let crc_zero = mpegts_crc32(&[0u8; 12]);
assert_ne!(crc_zero, crc_a);
}
#[test]
fn video_only_mux_emits_pat_pmt_then_video() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
drop(mux);
assert_ts_well_formed(&sink);
let pids = extract_pids(&sink);
assert_eq!(pids[0], PID_PAT);
assert_eq!(pids[1], PID_PMT);
assert!(pids.iter().any(|p| *p == PID_VIDEO));
}
#[test]
fn first_video_pes_carries_pcr() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
drop(mux);
let pkt = sink
.chunks(188)
.find(|p| u16::from_be_bytes([p[1] & 0x1F, p[2]]) == PID_VIDEO && (p[1] & 0x40) != 0)
.expect("video PUSI packet exists");
let af = af_body(pkt).expect("first video PES must carry an adaptation field");
assert!(!af.is_empty(), "AF flags byte present");
assert_eq!(af[0] & 0x10, 0x10, "PCR flag set on first video PES");
}
#[test]
fn audio_track_appears_in_pmt_and_stream() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
mux.set_audio(AudioCodec::Ac3);
let mut frame = Vec::new();
frame.extend_from_slice(&3u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C]);
mux.write_video(0, true, &frame).unwrap();
mux.write_audio(20_000_000, &[0x0B, 0x77, 0x12, 0x34])
.unwrap();
mux.finish().unwrap();
drop(mux);
assert_ts_well_formed(&sink);
let pids = extract_pids(&sink);
assert!(pids.iter().any(|p| *p == PID_VIDEO));
assert!(pids.iter().any(|p| *p == PID_AUDIO));
}
#[test]
fn psi_re_emits_at_interval() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let big: Vec<u8> = (0..(60 * 1024)).map(|i| (i & 0xff) as u8).collect();
let mut frame = Vec::new();
frame.extend_from_slice(&(big.len() as u32).to_be_bytes());
frame.extend_from_slice(&big);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
drop(mux);
assert_ts_well_formed(&sink);
let pids = extract_pids(&sink);
let pat_count = pids.iter().filter(|p| **p == PID_PAT).count();
let pmt_count = pids.iter().filter(|p| **p == PID_PMT).count();
assert!(pat_count >= 2, "expected ≥2 PAT, got {}", pat_count);
assert!(pmt_count >= 2, "expected ≥2 PMT, got {}", pmt_count);
}
#[test]
fn continuity_counter_increments_per_pid() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
for pts in [0i64, 40_000_000, 80_000_000] {
let mut frame = Vec::new();
frame.extend_from_slice(&3u32.to_be_bytes());
frame.extend_from_slice(&[0xAA, 0xBB, 0xCC]);
mux.write_video(pts, pts == 0, &frame).unwrap();
}
mux.finish().unwrap();
drop(mux);
let ccs: Vec<u8> = sink
.chunks(188)
.filter(|p| u16::from_be_bytes([p[1] & 0x1F, p[2]]) == PID_VIDEO)
.map(|p| p[3] & 0x0F)
.collect();
for w in ccs.windows(2) {
assert_eq!(w[1], (w[0] + 1) & 0x0F);
}
}
fn af_body(packet: &[u8]) -> Option<Vec<u8>> {
let afc = (packet[3] >> 4) & 0x03;
if afc & 0b10 == 0 {
return None;
}
let af_len = packet[4] as usize;
if af_len == 0 {
return Some(Vec::new());
}
Some(packet[5..5 + af_len].to_vec())
}
#[test]
fn stuffing_only_tail_packet_is_spec_valid() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
mux.set_audio(AudioCodec::Ac3);
let mut vframe = Vec::new();
vframe.extend_from_slice(&4u32.to_be_bytes());
vframe.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &vframe).unwrap();
let audio: Vec<u8> = (0..200u32).map(|i| (i & 0xFF) as u8).collect();
mux.write_audio(20_000_000, &audio).unwrap();
mux.finish().unwrap();
drop(mux);
assert_ts_well_formed(&sink);
let mut saw_stuffing_af = false;
for pkt in sink.chunks(188) {
let pid = u16::from_be_bytes([pkt[1] & 0x1F, pkt[2]]);
if pid != PID_AUDIO {
continue;
}
let afc = (pkt[3] >> 4) & 0x03;
if afc & 0b10 == 0 {
continue; }
let af_len = pkt[4] as usize;
assert!(
af_len >= 1,
"stuffing AF must include the mandatory flags byte"
);
assert_eq!(
pkt[5], 0x00,
"stuffing-only AF flags byte must be 0x00, not 0x{:02X}",
pkt[5]
);
assert!(af_len <= 183, "AF length overflows the 184-byte body");
saw_stuffing_af = true;
}
assert!(
saw_stuffing_af,
"expected at least one audio packet with a stuffing AF"
);
}
#[test]
fn rai_set_on_keyframe_pes_packet() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
drop(mux);
let pkt = sink
.chunks(188)
.find(|p| u16::from_be_bytes([p[1] & 0x1F, p[2]]) == PID_VIDEO && (p[1] & 0x40) != 0)
.expect("video PUSI packet exists");
let af = af_body(pkt).expect("AF present on first packet of keyframe video PES");
assert!(!af.is_empty(), "AF flags byte present");
assert_eq!(af[0] & 0x40, 0x40, "RAI bit set");
}
#[test]
fn pcr_packet_without_keyframe_has_rai_clear() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let mut frame0 = Vec::new();
frame0.extend_from_slice(&4u32.to_be_bytes());
frame0.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &frame0).unwrap();
let big: Vec<u8> = (0..(50 * 1024)).map(|i| (i & 0xff) as u8).collect();
for i in 1..3 {
let mut frame = Vec::new();
frame.extend_from_slice(&(big.len() as u32).to_be_bytes());
frame.extend_from_slice(&big);
mux.write_video((i as i64) * 40_000_000, false, &frame)
.unwrap();
}
mux.finish().unwrap();
drop(mux);
let video_pkts: Vec<&[u8]> = sink
.chunks(188)
.filter(|p| u16::from_be_bytes([p[1] & 0x1F, p[2]]) == PID_VIDEO)
.collect();
assert!(
video_pkts.len() >= 2,
"expected ≥2 video packets, got {}",
video_pkts.len()
);
let later_pcr = video_pkts
.iter()
.skip(1)
.find_map(|p| {
let af = af_body(p)?;
if !af.is_empty() && (af[0] & 0x10) != 0 {
Some(af)
} else {
None
}
})
.expect("a later PCR-bearing video packet exists");
assert_eq!(
later_pcr[0] & 0x40,
0,
"RAI must be clear on a non-keyframe-start PCR packet"
);
}
#[test]
fn pcr_restamped_mid_pes_within_interval() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let big: Vec<u8> = (0..(60 * 1024)).map(|i| (i & 0xff) as u8).collect();
let mut frame = Vec::new();
frame.extend_from_slice(&(big.len() as u32).to_be_bytes());
frame.extend_from_slice(&big);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
drop(mux);
assert_ts_well_formed(&sink);
let mut video_idx = 0usize;
let mut pcr_indices: Vec<usize> = Vec::new();
let mut total_video = 0usize;
for pkt in sink.chunks(188) {
let pid = u16::from_be_bytes([pkt[1] & 0x1F, pkt[2]]);
if pid != PID_VIDEO {
continue;
}
total_video += 1;
if let Some(af) = af_body(pkt) {
if !af.is_empty() && (af[0] & 0x10) != 0 {
pcr_indices.push(video_idx);
}
}
video_idx += 1;
}
assert!(
total_video > PCR_INTERVAL_PACKETS as usize,
"test needs a PES spanning more than one PCR interval, got {total_video} video packets"
);
assert!(
pcr_indices.len() >= 2,
"PCR must be re-stamped mid-PES, but only {} PCR-bearing packet(s) \
appeared across {} video packets of one PES",
pcr_indices.len(),
total_video
);
assert_eq!(pcr_indices[0], 0, "first video packet must carry PCR");
let max_gap = PCR_INTERVAL_PACKETS + 1;
for w in pcr_indices.windows(2) {
assert!(
(w[1] - w[0]) as u64 <= max_gap,
"PCR gap {} exceeds the {}-packet bound",
w[1] - w[0],
max_gap
);
}
let tail = total_video - 1 - *pcr_indices.last().unwrap();
assert!(
tail as u64 <= max_gap,
"trailing run after the last PCR ({tail}) exceeds the {max_gap}-packet bound"
);
}
#[test]
fn keyframe_video_with_pcr_combines_flags() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
let mut small = Vec::new();
small.extend_from_slice(&4u32.to_be_bytes());
small.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &small).unwrap();
mux.finish().unwrap();
drop(mux);
let video_pusi: Vec<&[u8]> = sink
.chunks(188)
.filter(|p| u16::from_be_bytes([p[1] & 0x1F, p[2]]) == PID_VIDEO && (p[1] & 0x40) != 0)
.collect();
assert!(!video_pusi.is_empty(), "a video PES start exists");
let af = af_body(video_pusi[0]).expect("AF present on first keyframe PES");
assert!(!af.is_empty(), "AF flags byte present");
assert_eq!(af[0], 0x50, "flags == RAI | PCR on the first keyframe PES");
}
fn find_pkt(buf: &[u8], pid: u16, pusi: bool) -> Option<&[u8]> {
buf.chunks(188).find(|p| {
u16::from_be_bytes([p[1] & 0x1F, p[2]]) == pid && (!pusi || (p[1] & 0x40) != 0)
})
}
fn psi_section(pkt: &[u8]) -> &[u8] {
let pointer = pkt[4] as usize;
&pkt[5 + pointer..]
}
#[test]
fn crc32_residue_over_section_plus_crc_is_zero() {
let data = [
0x00u8, 0xB0, 0x0D, 0x00, 0x01, 0xC1, 0x00, 0x00, 0x00, 0x01, 0xE1, 0x00,
];
let crc = mpegts_crc32(&data);
assert_eq!(crc, 0xE8F9_5E7D, "CRC-32/MPEG-2 known-answer vector");
let mut with_crc = data.to_vec();
with_crc.extend_from_slice(&crc.to_be_bytes());
assert_eq!(
mpegts_crc32(&with_crc),
0,
"CRC residue over message+CRC must be 0"
);
}
#[test]
fn emitted_pat_pmt_crc_is_valid() {
let mut sink: Vec<u8> = Vec::new();
{
let mut mux = M2tsMux::new(&mut sink);
mux.set_audio(AudioCodec::Ac3);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
}
for pid in [PID_PAT, PID_PMT] {
let pkt = find_pkt(&sink, pid, true).expect("PSI packet present");
let sec = psi_section(pkt);
let section_len = (((sec[1] & 0x0F) as usize) << 8) | sec[2] as usize;
let total = 3 + section_len;
assert!(sec.len() >= total, "section fits in payload");
assert_eq!(
mpegts_crc32(&sec[..total]),
0,
"PID {pid:#06x} section CRC must validate (residue 0)"
);
}
}
#[test]
fn pat_points_at_pmt_pid() {
let pat = build_pat(PID_PMT);
let sec = &pat[1..]; assert_eq!(sec[0], 0x00, "table_id = PAT");
let prog_num = u16::from_be_bytes([sec[8], sec[9]]);
let pmt_pid = u16::from_be_bytes([sec[10] & 0x1F, sec[11]]);
assert_eq!(prog_num, 1, "program_number 1");
assert_eq!(pmt_pid, PID_PMT, "PAT points at PMT PID");
}
#[test]
fn pmt_advertises_video_and_audio_stream_types() {
let pmt = build_pmt(Some(AudioCodec::Ac3));
let sec = &pmt[1..]; assert_eq!(sec[0], 0x02, "table_id = PMT");
let section_len = (((sec[1] & 0x0F) as usize) << 8) | sec[2] as usize;
let prog_info_len = (((sec[10] & 0x0F) as usize) << 8) | sec[11] as usize;
let mut pos = 12 + prog_info_len;
let end = 3 + section_len - 4; let mut types = Vec::new();
while pos + 5 <= end {
types.push(sec[pos]);
let es_info = (((sec[pos + 3] & 0x0F) as usize) << 8) | sec[pos + 4] as usize;
pos += 5 + es_info;
}
assert!(types.contains(&STREAM_TYPE_HEVC), "HEVC video in PMT");
assert!(types.contains(&STREAM_TYPE_AC3), "AC-3 audio in PMT");
}
#[test]
fn pmt_video_only_omits_audio_entry() {
let pmt = build_pmt(None);
let sec = &pmt[1..];
let section_len = (((sec[1] & 0x0F) as usize) << 8) | sec[2] as usize;
let prog_info_len = (((sec[10] & 0x0F) as usize) << 8) | sec[11] as usize;
let mut pos = 12 + prog_info_len;
let end = 3 + section_len - 4;
let mut count = 0;
while pos + 5 <= end {
count += 1;
let es_info = (((sec[pos + 3] & 0x0F) as usize) << 8) | sec[pos + 4] as usize;
pos += 5 + es_info;
}
assert_eq!(count, 1, "video-only PMT has exactly one ES entry");
}
#[test]
fn truehd_audio_uses_stream_type_0x83() {
let pmt = build_pmt(Some(AudioCodec::TrueHd));
let sec = &pmt[1..];
let section_len = (((sec[1] & 0x0F) as usize) << 8) | sec[2] as usize;
let end = 3 + section_len - 4;
let mut pos = 12; let mut found = false;
while pos + 5 <= end {
if sec[pos] == STREAM_TYPE_TRUEHD {
found = true;
}
let es_info = (((sec[pos + 3] & 0x0F) as usize) << 8) | sec[pos + 4] as usize;
pos += 5 + es_info;
}
assert!(found, "TrueHD stream_type 0x83 must appear in PMT");
}
#[test]
fn pmt_pcr_pid_is_video_pid() {
let pmt = build_pmt(None);
let sec = &pmt[1..];
let pcr_pid = u16::from_be_bytes([sec[8] & 0x1F, sec[9]]);
assert_eq!(pcr_pid, PID_VIDEO, "PCR_PID advertised as the video PID");
}
#[test]
fn pcr_base_round_trips_through_adaptation_field() {
let pcr: u64 = 0x1_2345_6789 & ((1 << 33) - 1);
let af = build_pcr_adaptation(pcr);
assert_eq!(af[0], 0x10, "PCR_flag set, others clear");
let base = ((af[1] as u64) << 25)
| ((af[2] as u64) << 17)
| ((af[3] as u64) << 9)
| ((af[4] as u64) << 1)
| ((af[5] as u64 >> 7) & 0x01);
assert_eq!(base, pcr, "PCR base must round-trip through the AF");
}
#[test]
fn first_video_pcr_leads_pts_by_lead_time() {
let mut sink: Vec<u8> = Vec::new();
{
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(1_000_000_000, true, &frame).unwrap();
mux.finish().unwrap();
}
let pkt = find_pkt(&sink, PID_VIDEO, true).unwrap();
let af = af_body(pkt).unwrap();
assert_eq!(af[0] & 0x10, 0x10);
let base = ((af[1] as u64) << 25)
| ((af[2] as u64) << 17)
| ((af[3] as u64) << 9)
| ((af[4] as u64) << 1)
| ((af[5] as u64 >> 7) & 0x01);
assert_eq!(base, 0, "first frame PCR clamps to 0 (no underflow)");
}
#[test]
fn extreme_pts_does_not_overflow_and_clamps_to_33bit() {
let mut sink: Vec<u8> = Vec::new();
{
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(i64::MAX, true, &frame).unwrap();
mux.finish().unwrap();
}
assert_ts_well_formed(&sink);
let pkt = find_pkt(&sink, PID_VIDEO, true).unwrap();
let af_len = pkt[4] as usize;
let pes = &pkt[4 + 1 + af_len..];
let pts = ((((pes[9] >> 1) & 0x07) as u64) << 30)
| ((pes[10] as u64) << 22)
| (((pes[11] >> 1) as u64) << 15)
| ((pes[12] as u64) << 7)
| ((pes[13] >> 1) as u64);
assert!(pts < (1u64 << 33), "PTS stays within the 33-bit field");
}
#[test]
fn base_relative_pts_wraps_across_33bit_clock_rollover() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
mux.base_pts_90k = Some((1u64 << 33) - 100);
let pts_ns = 1_000_000i64;
let rel = mux.base_relative_pts(pts_ns, false);
assert_eq!(
rel, 190,
"a 33-bit clock wrap must produce the true forward delta, not 0"
);
mux.base_pts_90k = Some(200);
let rel0 = mux.base_relative_pts(1_000_000i64, false); assert_eq!(rel0, 0, "a frame before the base must still floor to 0");
}
#[test]
fn negative_pts_ns_encodes_zero() {
let mut sink: Vec<u8> = Vec::new();
{
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&4u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C, 0x01]);
mux.write_video(-5, true, &frame).unwrap();
mux.finish().unwrap();
}
let pkt = find_pkt(&sink, PID_VIDEO, true).unwrap();
let af_len = pkt[4] as usize;
let pes = &pkt[4 + 1 + af_len..];
let pts = ((((pes[9] >> 1) & 0x07) as u64) << 30)
| ((pes[10] as u64) << 22)
| (((pes[11] >> 1) as u64) << 15)
| ((pes[12] as u64) << 7)
| ((pes[13] >> 1) as u64);
assert_eq!(pts, 0, "negative pts_ns encodes PTS 0");
}
#[test]
fn write_audio_without_track_is_silently_dropped() {
let mut sink: Vec<u8> = Vec::new();
{
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&3u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C]);
mux.write_video(0, true, &frame).unwrap();
mux.write_audio(0, &[0x0B, 0x77]).unwrap(); mux.finish().unwrap();
}
let pids = extract_pids(&sink);
assert!(
!pids.iter().any(|p| *p == PID_AUDIO),
"no audio track configured → no audio PID emitted"
);
}
#[test]
fn finish_without_frames_emits_nothing() {
let mut sink: Vec<u8> = Vec::new();
let mut mux = M2tsMux::new(&mut sink);
mux.finish().unwrap();
drop(mux);
assert!(sink.is_empty(), "no frames → no output");
}
#[test]
fn pat_always_on_pid_zero() {
let mut sink: Vec<u8> = Vec::new();
{
let mut mux = M2tsMux::new(&mut sink);
let mut frame = Vec::new();
frame.extend_from_slice(&3u32.to_be_bytes());
frame.extend_from_slice(&[0x40, 0x01, 0x0C]);
mux.write_video(0, true, &frame).unwrap();
mux.finish().unwrap();
}
assert_eq!(
extract_pids(&sink)[0],
0x0000,
"first packet is PAT on PID 0"
);
}
}