use super::{CodecParser, Frame, PesPacket, pts_to_ns};
const DTS_CORE_SYNC: [u8; 4] = [0x7F, 0xFE, 0x80, 0x01];
#[cfg(test)]
const DTS_HD_EXT_SYNC: [u8; 4] = [0x64, 0x58, 0x20, 0x25];
pub struct DtsParser {
buf: Vec<u8>,
pending_pts: i64,
pts_marks: Vec<(usize, i64)>,
}
impl Default for DtsParser {
fn default() -> Self {
Self::new()
}
}
impl DtsParser {
pub fn new() -> Self {
Self {
buf: Vec::with_capacity(32768),
pending_pts: 0,
pts_marks: Vec::new(),
}
}
fn drain_front(&mut self, n: usize) {
if n == 0 {
return;
}
self.buf.drain(..n);
for m in &mut self.pts_marks {
m.0 = m.0.saturating_sub(n);
}
let last_zero = self
.pts_marks
.iter()
.rposition(|&(off, _)| off == 0)
.filter(|&i| i > 0);
if let Some(i) = last_zero {
self.pts_marks.drain(..i);
}
}
fn front_pts(&self) -> i64 {
self.pts_marks
.iter()
.rev()
.find(|&&(off, _)| off == 0)
.map(|&(_, pts)| pts)
.unwrap_or(self.pending_pts)
}
}
const MAX_AU_BYTES: usize = 65536;
const CORE_HEADER_MIN_BYTES: usize = 10;
const MIN_CORE_FRAME_BYTES: usize = 96;
const PTS_UNSET: i64 = -1;
impl CodecParser for DtsParser {
fn parse(&mut self, pes: &PesPacket) -> Vec<Frame> {
if pes.data.is_empty() {
return Vec::new();
}
let pts_ns = pes.pts.map(pts_to_ns).unwrap_or(0);
if self.buf.is_empty() || self.pending_pts == PTS_UNSET {
self.pending_pts = pts_ns;
}
self.pts_marks.push((self.buf.len(), pts_ns));
self.buf.extend_from_slice(&pes.data);
let mut frames = Vec::new();
loop {
let Some(start) = find_sync(&self.buf, &DTS_CORE_SYNC) else {
if self.buf.len() > 3 {
let tail = self.buf.len() - 3;
self.drain_front(tail);
}
break;
};
if start > 0 {
self.drain_front(start);
debug_assert_eq!(
find_sync(&self.buf, &DTS_CORE_SYNC),
Some(0),
"drain_front(start) must leave the core sync at offset 0"
);
}
if self.buf.len() < CORE_HEADER_MIN_BYTES {
break;
}
let core_size = dts_core_frame_size(&self.buf);
if !(MIN_CORE_FRAME_BYTES..=MAX_AU_BYTES).contains(&core_size) {
self.drain_front(4);
continue;
}
if self.buf.len() < core_size {
break; }
let mut forced = false;
let au_end = match next_core_boundary(&self.buf, core_size) {
NextCore::Found(end) => end,
NextCore::NeedMore => break, NextCore::None => {
if self.buf.len() <= MAX_AU_BYTES {
break;
}
forced = true;
self.buf.len()
}
};
let au: Vec<u8> = self.buf[..au_end].to_vec();
let au_pts = self.front_pts();
frames.push(Frame {
pts_ns: au_pts,
keyframe: true,
data: au,
duration_ns: None,
});
self.drain_front(au_end);
self.pending_pts = self.front_pts();
if forced {
self.pending_pts = PTS_UNSET;
self.pts_marks.clear();
}
}
if self.buf.is_empty() {
self.pts_marks.clear();
}
frames
}
fn flush(&mut self) -> Vec<Frame> {
if find_sync(&self.buf, &DTS_CORE_SYNC) != Some(0) || self.buf.len() < CORE_HEADER_MIN_BYTES
{
self.buf.clear();
return Vec::new();
}
let core_size = dts_core_frame_size(&self.buf);
if core_size < MIN_CORE_FRAME_BYTES || self.buf.len() < core_size {
self.buf.clear();
return Vec::new();
}
let front = self.front_pts();
let pts_ns = if front == PTS_UNSET { 0 } else { front };
let au = std::mem::take(&mut self.buf);
self.pts_marks.clear();
vec![Frame {
pts_ns,
keyframe: true,
data: au,
duration_ns: None,
}]
}
fn codec_private(&self) -> Option<Vec<u8>> {
None
}
}
fn find_sync(data: &[u8], pattern: &[u8; 4]) -> Option<usize> {
if data.len() < 4 {
return None;
}
(0..=data.len() - 4).find(|&i| data[i..i + 4] == *pattern)
}
enum NextCore {
Found(usize),
NeedMore,
None,
}
fn next_core_boundary(buf: &[u8], core_size: usize) -> NextCore {
let mut from = core_size;
while let Some(rel) = find_sync(&buf[from..], &DTS_CORE_SYNC) {
let pos = from + rel;
if buf.len() - pos < CORE_HEADER_MIN_BYTES {
return NextCore::NeedMore;
}
let sz = dts_core_frame_size(&buf[pos..]);
if (MIN_CORE_FRAME_BYTES..=MAX_AU_BYTES).contains(&sz) {
return NextCore::Found(pos);
}
from = pos + 4;
}
NextCore::None
}
fn dts_core_frame_size(data: &[u8]) -> usize {
if data.len() < CORE_HEADER_MIN_BYTES {
return 0;
}
let fsize =
((data[5] as usize & 0x03) << 12) | ((data[6] as usize) << 4) | ((data[7] as usize) >> 4);
fsize + 1
}
#[cfg(test)]
mod tests {
use super::*;
use crate::mux::ts::PesPacket;
fn make_pes(data: Vec<u8>, pts: Option<i64>) -> PesPacket {
PesPacket {
pid: 0x1100,
pts,
dts: None,
data,
}
}
fn make_dts_core(size: usize) -> Vec<u8> {
let fsize = size - 1;
let mut data = vec![0u8; size];
data[0..4].copy_from_slice(&DTS_CORE_SYNC);
data[5] = (data[5] & 0xFC) | ((fsize >> 12) & 0x03) as u8;
data[6] = ((fsize >> 4) & 0xFF) as u8;
data[7] = (data[7] & 0x0F) | (((fsize & 0x0F) << 4) as u8);
data
}
#[test]
fn parse_empty_pes() {
let mut parser = DtsParser::new();
let pes = make_pes(Vec::new(), Some(0));
assert!(parser.parse(&pes).is_empty());
}
#[test]
fn parse_single_frame() {
let mut parser = DtsParser::new();
let frame = make_dts_core(512);
let pes = make_pes(frame, Some(90000));
assert!(parser.parse(&pes).is_empty());
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].data.len(), 512);
}
#[test]
fn parse_frame_spanning_two_pes() {
let mut parser = DtsParser::new();
let frame = make_dts_core(512);
let mid = 256;
let pes1 = make_pes(frame[..mid].to_vec(), Some(90000));
assert!(parser.parse(&pes1).is_empty());
let pes2 = make_pes(frame[mid..].to_vec(), Some(93000));
assert!(parser.parse(&pes2).is_empty());
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].data.len(), 512);
}
#[test]
fn two_cores_back_to_back_emit_first_on_boundary() {
let mut parser = DtsParser::new();
let mut stream = make_dts_core(512);
stream.extend_from_slice(&make_dts_core(640));
let f = parser.parse(&make_pes(stream, Some(90000)));
assert_eq!(f.len(), 1);
assert_eq!(f[0].data.len(), 512);
assert_eq!(f[0].pts_ns, pts_to_ns(90000), "AU1 keeps its PES PTS");
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].data.len(), 640);
assert_eq!(
tail[0].pts_ns,
pts_to_ns(90000),
"AU2 came in the same PES → same PTS"
);
}
#[test]
fn two_aus_flushed_in_one_call_keep_their_own_pts() {
let mut parser = DtsParser::new();
let f0 = parser.parse(&make_pes(make_dts_core(512), Some(100)));
assert!(f0.is_empty(), "core1 held awaiting next core");
let mut pes_b = make_dts_core(600);
pes_b.extend_from_slice(&make_dts_core(640));
let f = parser.parse(&make_pes(pes_b, Some(200)));
assert_eq!(f.len(), 2, "AU1 and AU2 both close in this call");
assert_eq!(f[0].data.len(), 512, "AU1 = core1");
assert_eq!(
f[0].pts_ns,
pts_to_ns(100),
"AU1 keeps PES A's PTS, not the later PES B PTS"
);
assert_eq!(f[1].data.len(), 600, "AU2 = core2");
assert_eq!(
f[1].pts_ns,
pts_to_ns(200),
"AU2's core arrived in PES B → PES B PTS"
);
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].pts_ns, pts_to_ns(200));
}
#[test]
fn second_au_with_core_in_earlier_pes_keeps_that_pts() {
let mut parser = DtsParser::new();
let mut pes_a = make_dts_core(512);
pes_a.extend_from_slice(&make_dts_core(600));
let f = parser.parse(&make_pes(pes_a, Some(100)));
assert_eq!(f.len(), 1);
assert_eq!(f[0].pts_ns, pts_to_ns(100), "AU1 PES A PTS");
let f2 = parser.parse(&make_pes(make_dts_core(640), Some(200)));
assert_eq!(f2.len(), 1);
assert_eq!(f2[0].data.len(), 600, "AU2 = core2");
assert_eq!(
f2[0].pts_ns,
pts_to_ns(100),
"AU2's core arrived in PES A → must keep PES A's PTS, not PES B's"
);
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].pts_ns, pts_to_ns(200), "AU3 = core3, PES B PTS");
}
fn make_dts_ext(size: usize) -> Vec<u8> {
let mut e = vec![0u8; size];
e[0..4].copy_from_slice(&DTS_HD_EXT_SYNC);
e
}
#[test]
fn keeps_dts_hd_extension_in_separate_pes_packets() {
let mut parser = DtsParser::new();
assert!(
parser
.parse(&make_pes(make_dts_core(512), Some(90000)))
.is_empty(),
"core alone: must wait for any following extension"
);
assert!(
parser
.parse(&make_pes(make_dts_ext(256), Some(91000)))
.is_empty(),
"first extension PES: still waiting for the unit to close"
);
assert!(
parser
.parse(&make_pes(make_dts_ext(200), Some(91500)))
.is_empty(),
"second extension PES: unit still not closed (no next core yet)"
);
let f = parser.parse(&make_pes(make_dts_core(512), Some(93000)));
assert_eq!(f.len(), 1);
assert_eq!(
f[0].data.len(),
512 + 256 + 200,
"frame must include core + every extension substream"
);
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].data.len(), 512);
}
#[test]
fn extension_split_across_pes_is_preserved() {
let mut parser = DtsParser::new();
let ext = make_dts_ext(300);
assert!(
parser
.parse(&make_pes(make_dts_core(512), Some(90000)))
.is_empty()
);
assert!(
parser
.parse(&make_pes(ext[..150].to_vec(), Some(91000)))
.is_empty()
);
assert!(
parser
.parse(&make_pes(ext[150..].to_vec(), Some(91000)))
.is_empty()
);
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].data.len(), 512 + 300);
}
fn bogus_tiny_core_sync() -> Vec<u8> {
let mut v = vec![0u8; 10];
v[0..4].copy_from_slice(&DTS_CORE_SYNC);
assert_eq!(dts_core_frame_size(&v), 1);
v
}
#[test]
fn bogus_tiny_core_sync_does_not_split_or_drop_real_au() {
let mut parser = DtsParser::new();
let mut ext = make_dts_ext(256);
let bogus = bogus_tiny_core_sync();
ext[64..64 + bogus.len()].copy_from_slice(&bogus);
let mut frame1 = make_dts_core(512);
frame1.extend_from_slice(&ext);
assert!(
parser.parse(&make_pes(frame1, Some(90000))).is_empty(),
"bogus tiny core sync must not close the AU; wait for a real core"
);
let f = parser.parse(&make_pes(make_dts_core(640), Some(93000)));
assert_eq!(f.len(), 1, "exactly one real access unit emitted");
assert_eq!(
f[0].data.len(),
512 + 256,
"AU must be the full core + extension, not split at the bogus sync"
);
assert_eq!(f[0].pts_ns, pts_to_ns(90000), "AU keeps the core's PTS");
let tail = parser.flush();
assert_eq!(tail.len(), 1);
assert_eq!(tail[0].data.len(), 640);
}
#[test]
fn sub_spec_core_size_is_rejected_as_false_sync() {
let false_size = 64usize;
assert!(
(CORE_HEADER_MIN_BYTES..MIN_CORE_FRAME_BYTES).contains(&false_size),
"test fixture must sit in the widened reject window"
);
let mut parser = DtsParser::new();
let mut ext = make_dts_ext(256);
let bogus = make_dts_core(false_size); ext[64..64 + bogus.len()].copy_from_slice(&bogus);
let mut frame1 = make_dts_core(512);
frame1.extend_from_slice(&ext);
assert!(
parser.parse(&make_pes(frame1, Some(90000))).is_empty(),
"sub-spec core size must not close the AU"
);
let f = parser.parse(&make_pes(make_dts_core(640), Some(93000)));
assert_eq!(f.len(), 1);
assert_eq!(
f[0].data.len(),
512 + 256,
"AU must not be split at the sub-spec false sync"
);
}
#[test]
fn forced_emit_does_not_corrupt_next_au_pts() {
let mut parser = DtsParser::new();
let core_pts = 90000i64;
assert!(
parser
.parse(&make_pes(make_dts_core(512), Some(core_pts)))
.is_empty()
);
let ext_pts = 120000i64; let big_ext = make_dts_ext(MAX_AU_BYTES + 1024);
let f = parser.parse(&make_pes(big_ext, Some(ext_pts)));
assert_eq!(f.len(), 1, "oversized buffer force-emits one AU");
assert_eq!(
f[0].pts_ns,
pts_to_ns(core_pts),
"forced AU keeps the core PTS"
);
let next_core_pts = 150000i64;
assert!(
parser
.parse(&make_pes(make_dts_core(512), Some(next_core_pts)))
.is_empty()
);
let next_next_pts = 180000i64;
let f2 = parser.parse(&make_pes(make_dts_core(512), Some(next_next_pts)));
assert_eq!(f2.len(), 1);
assert_eq!(
f2[0].pts_ns,
pts_to_ns(next_core_pts),
"AU after a forced emit must use the next core's PTS, not the \
stale extension PTS"
);
}
#[test]
fn codec_private_none() {
let parser = DtsParser::new();
assert!(parser.codec_private().is_none());
}
}