use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
use alloc::vec::Vec;
use core::marker::PhantomData;
use broadcast_common::{Demand, Serialize, Stage, Timestamp, Unpackage};
use mpeg_pes::{PesAssembler, PesPacket};
use mpeg_ts::resync::TsResync;
use mpeg_ts::ts::{SectionReassembler, TS_PACKET_SIZE, TsPacket};
use crate::aac_asc::{AdtsHeader, AudioSpecificConfig, parse_adts_header};
use crate::ac3::{
AC3_SAMPLES_PER_SYNCFRAME, Ac3SyncframeInfo, Ec3SyncframeInfo, split_ac3_syncframes,
split_eac3_syncframes,
};
use crate::annexb::{annexb_to_length_prefixed, iter_annexb_nals};
use crate::avc_config::{AVCConfigurationBox, AVCDecoderConfigurationRecord};
use crate::dts::{DtsCoreFrameInfo, split_dts_core_frames};
use crate::error::{Error, Result};
use crate::hevc_config::{HEVCConfigurationBox, HEVCDecoderConfigurationRecord};
use crate::media::{Media, PcrSample, Track};
use crate::mp4esds::{
DecoderConfigDescriptor, DecoderSpecificInfo, ESDescriptor, EsdsBox, ObjectTypeIndication,
SLConfigDescriptor, StreamType as EsdsStreamType,
};
use crate::mpeg_legacy::{Mpeg2SeqHeader, MpegAudioFrameHeader};
use crate::mpegh::{MHADecoderConfigurationRecord, find_mpegh3da_config};
use crate::nal::{NalCodec, access_unit_is_rap, is_keyframe_nal, nal_unit_type};
use crate::nalu_types::{AvcPps, AvcSps, HevcNalArray, HevcNalUnit};
use crate::pipeline::{CodecConfig, DataCarriage, Provenance, Sample, SampleFlags, TrackSpec};
const PAT_PID: u16 = 0x0000;
const TABLE_ID_PAT: u8 = 0x00;
const TABLE_ID_PMT: u8 = 0x02;
const SECTION_HEADER_LEN: usize = 8;
const VERSION_NUMBER_MASK: u8 = 0x1F;
const CURRENT_NEXT_INDICATOR_BIT: u8 = 0x01;
const CRC32_LEN: usize = 4;
const SECTION_SYNTAX_INDICATOR_BIT: u8 = 0x80;
const SECTION_LENGTH_HI_MASK: u8 = 0x0F;
const PID_HI_MASK: u8 = 0x1F;
const PAT_ENTRY_LEN: usize = 4;
const INFO_LENGTH_HI_MASK: u8 = 0x0F;
const NETWORK_PROGRAM_NUMBER: u16 = 0x0000;
const NULL_PACKET_PID: u16 = 0x1FFF;
const MAX_UNATTRIBUTED_BYTES: usize = 4 * 1024 * 1024;
const TS_MAX_PAYLOAD_BYTES: usize = TS_PACKET_SIZE - 4;
const MAX_PES_BUFFER_BYTES: usize = 4 * 1024 * 1024;
const MAX_PROBE_BACKLOG_BYTES: usize = 4 * 1024 * 1024;
const STREAM_TYPE_MPEG2_VIDEO: u8 = 0x02;
const STREAM_TYPE_MPEG1_AUDIO: u8 = 0x03;
const STREAM_TYPE_MPEG2_AUDIO: u8 = 0x04;
const STREAM_TYPE_AVC: u8 = 0x1B;
const STREAM_TYPE_HEVC: u8 = 0x24;
const STREAM_TYPE_AAC_ADTS: u8 = 0x0F;
const STREAM_TYPE_AC3: u8 = 0x81;
const STREAM_TYPE_EAC3: u8 = 0x87;
const STREAM_TYPE_DTS_82: u8 = 0x82;
const STREAM_TYPE_DTS_85: u8 = 0x85;
const STREAM_TYPE_DTS_8A: u8 = 0x8A;
const STREAM_TYPE_MPEGH: u8 = 0x2D;
const STREAM_TYPE_PES_PRIVATE: u8 = 0x06;
const STREAM_TYPE_METADATA_PES: u8 = 0x15;
const DESC_TAG_AC3: u8 = 0x6A;
const DESC_TAG_ENHANCED_AC3: u8 = 0x7A;
const DESC_TAG_DTS: u8 = 0x7B;
const STREAM_TYPE_PRIVATE_SECTIONS: u8 = 0x05;
const STREAM_TYPE_DSMCC_TYPE_A: u8 = 0x0A;
const STREAM_TYPE_DSMCC_TYPE_B: u8 = 0x0B;
const STREAM_TYPE_DSMCC_TYPE_C: u8 = 0x0C;
const STREAM_TYPE_DSMCC_TYPE_D: u8 = 0x0D;
const STREAM_TYPE_DSMCC_SYNC_DOWNLOAD: u8 = 0x14;
const STREAM_TYPE_SCTE35: u8 = 0x86;
const NAL_LENGTH_SIZE_MINUS_ONE: u8 = 3;
const H264_NAL_SPS: u8 = 7;
const H264_NAL_PPS: u8 = 8;
const H264_NAL_TYPE_MASK: u8 = 0x1F;
const H265_NAL_VPS: u8 = 32;
const H265_NAL_SPS: u8 = 33;
const H265_NAL_PPS: u8 = 34;
const HVCC_CONFIGURATION_VERSION: u8 = 1;
const HVCC_CONSTANT_FRAME_RATE_UNSPEC: u8 = 0;
const HVCC_NUM_TEMPORAL_LAYERS: u8 = 1;
const HVCC_PARALLELISM_TYPE_UNKNOWN: u8 = 0;
const HVCC_AVG_FRAME_RATE_UNSPEC: u16 = 0;
const HVCC_MIN_SPATIAL_SEGMENTATION_UNSPEC: u16 = 0;
const OTI_MPEG4_AUDIO: u8 = 0x40;
const OTI_MPEG2_VIDEO_MAIN: u8 = 0x61;
const OTI_MPEG1_AUDIO: u8 = 0x6B;
const OTI_MPEG2_AUDIO: u8 = 0x69;
const STREAM_TYPE_AUDIO: u8 = 0x05;
const STREAM_TYPE_VISUAL: u8 = 0x04;
const ESDS_ES_ID: u16 = 1;
const ESDS_VIDEO_ES_ID: u16 = 2;
const SL_CONFIG_PREDEFINED_MP4: u8 = 0x02;
const AUDIO_SAMPLE_SIZE_BITS: u16 = 16;
const MPEGH_REFERENCE_CHANNEL_LAYOUT_UNSPECIFIED: u8 = 0;
const MPEGH_CHANNEL_COUNT_UNSPECIFIED: u16 = 0;
const MPEGH_SAMPLE_RATE_UNSPECIFIED: u32 = 0;
const VIDEO_TIMESCALE: u32 = 90_000;
const AAC_SAMPLES_PER_FRAME: u32 = 1024;
const ADTS_HEADER_SIZE: usize = 7;
const MPEG2_PICTURE_START_CODE: u8 = 0x00;
const MPEG2_PICTURE_CODING_TYPE_I: u8 = 0x01;
const TS_WRAP: u64 = 1 << 33;
const TS_WRAP_HALF: u64 = TS_WRAP / 2;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Codec {
H264,
Hevc,
Mpeg2Video,
MpegAudio(bool),
Aac,
Ac3,
Eac3,
Dts,
MpegH,
Data(u8),
}
impl Codec {
fn from_stream_type(stream_type: u8) -> Self {
match stream_type {
STREAM_TYPE_MPEG2_VIDEO => Codec::Mpeg2Video,
STREAM_TYPE_MPEG1_AUDIO => Codec::MpegAudio(false),
STREAM_TYPE_MPEG2_AUDIO => Codec::MpegAudio(true),
STREAM_TYPE_AVC => Codec::H264,
STREAM_TYPE_HEVC => Codec::Hevc,
STREAM_TYPE_AAC_ADTS => Codec::Aac,
STREAM_TYPE_AC3 => Codec::Ac3,
STREAM_TYPE_EAC3 => Codec::Eac3,
STREAM_TYPE_DTS_82 | STREAM_TYPE_DTS_85 | STREAM_TYPE_DTS_8A => Codec::Dts,
STREAM_TYPE_MPEGH => Codec::MpegH,
_ => Codec::Data(stream_type),
}
}
fn refine_with_descriptors(self, stream_type: u8, descriptors: &[u8]) -> Self {
if !matches!(
stream_type,
STREAM_TYPE_PES_PRIVATE | STREAM_TYPE_METADATA_PES
) {
return self;
}
let mut off = 0usize;
while off + 2 <= descriptors.len() {
let tag = descriptors[off];
let len = descriptors[off + 1] as usize;
match tag {
DESC_TAG_AC3 => return Codec::Ac3,
DESC_TAG_ENHANCED_AC3 => return Codec::Eac3,
DESC_TAG_DTS => return Codec::Dts,
_ => {}
}
off += 2 + len;
}
self
}
}
fn data_carriage(stream_type: u8) -> DataCarriage {
match stream_type {
STREAM_TYPE_PRIVATE_SECTIONS
| STREAM_TYPE_DSMCC_TYPE_A
| STREAM_TYPE_DSMCC_TYPE_B
| STREAM_TYPE_DSMCC_TYPE_C
| STREAM_TYPE_DSMCC_TYPE_D
| STREAM_TYPE_DSMCC_SYNC_DOWNLOAD
| STREAM_TYPE_SCTE35 => DataCarriage::Sections,
_ => DataCarriage::Pes,
}
}
fn unwrap_ts(prev_unwrapped: i128, prev_raw: u64, raw: u64) -> i128 {
let mut delta = raw as i128 - prev_raw as i128;
if delta > TS_WRAP_HALF as i128 {
delta -= TS_WRAP as i128; } else if delta < -(TS_WRAP_HALF as i128) {
delta += TS_WRAP as i128; }
prev_unwrapped + delta
}
fn rescale_to_track(anchor_90k: i128, timescale: u32) -> i64 {
let ts = timescale.max(1) as i128;
let scaled = if ts == VIDEO_TIMESCALE as i128 {
anchor_90k
} else {
(anchor_90k * ts).div_euclid(VIDEO_TIMESCALE as i128)
};
to_ticks(scaled)
}
fn mpeg2_is_sync(au: &[u8]) -> bool {
let mut i = 0usize;
while i + 4 <= au.len() {
if au[i] == 0x00 && au[i + 1] == 0x00 && au[i + 2] == 0x01 {
let code = au[i + 3];
if code == crate::mpeg_legacy::SEQUENCE_HEADER_CODE[3] {
return true;
}
if code == MPEG2_PICTURE_START_CODE && i + 6 <= au.len() {
let pct = (au[i + 5] >> 3) & 0x07;
return pct == MPEG2_PICTURE_CODING_TYPE_I;
}
}
i += 1;
}
false
}
fn find_mpeg_audio_sync(data: &[u8]) -> Option<(usize, MpegAudioFrameHeader)> {
let mut off = 0usize;
while off + 4 <= data.len() {
if let Ok(hdr) = MpegAudioFrameHeader::parse(&data[off..]) {
return Some((off, hdr));
}
off += 1;
}
None
}
fn split_mpeg_audio_frames(payload: &[u8]) -> Vec<&[u8]> {
let mut frames = Vec::new();
let mut off = 0usize;
while off + 4 <= payload.len() {
let Some((sync_off, hdr)) = find_mpeg_audio_sync(&payload[off..]) else {
break;
};
off += sync_off;
let flen = hdr.frame_length;
if flen < 4 || off + flen > payload.len() {
break;
}
frames.push(&payload[off..off + flen]);
off += flen;
}
frames
}
fn find_adts_sync(data: &[u8]) -> Option<(usize, AdtsHeader)> {
let mut off = 0usize;
while off + ADTS_HEADER_SIZE <= data.len() {
if let Ok(hdr) = parse_adts_header(&data[off..]) {
return Some((off, hdr));
}
off += 1;
}
None
}
fn split_adts_frames(payload: &[u8]) -> Vec<&[u8]> {
let mut frames = Vec::new();
let mut off = 0usize;
while off + ADTS_HEADER_SIZE <= payload.len() {
let Some((sync_off, hdr)) = find_adts_sync(&payload[off..]) else {
break;
};
off += sync_off;
let frame_len = hdr.frame_length as usize;
if frame_len < ADTS_HEADER_SIZE || off + frame_len > payload.len() {
break;
}
frames.push(&payload[off..off + frame_len]);
off += frame_len;
}
frames
}
fn sfi_to_hz(sfi: u8) -> Option<u32> {
Some(match sfi {
0 => 96000,
1 => 88200,
2 => 64000,
3 => 48000,
4 => 44100,
5 => 32000,
6 => 24000,
7 => 22050,
8 => 16000,
9 => 12000,
10 => 11025,
11 => 8000,
12 => 7350,
_ => return None,
})
}
fn parse_pat(section: &[u8]) -> Result<Vec<(u16, u16)>> {
if section.first().copied() != Some(TABLE_ID_PAT) {
return Ok(Vec::new());
}
let body = section_body(section, "PAT")?;
let mut programs = Vec::new();
let mut off = 0usize;
while off + PAT_ENTRY_LEN <= body.len() {
let program_number = u16::from_be_bytes([body[off], body[off + 1]]);
let pid = (((body[off + 2] & PID_HI_MASK) as u16) << 8) | body[off + 3] as u16;
if program_number != NETWORK_PROGRAM_NUMBER {
programs.push((program_number, pid));
}
off += PAT_ENTRY_LEN;
}
Ok(programs)
}
struct PmtSectionHeader {
program_number: u16,
version: u8,
current_next: bool,
section_number: u8,
last_section_number: u8,
}
fn parse_pmt_section_header(section: &[u8]) -> Result<PmtSectionHeader> {
if section.first().copied() != Some(TABLE_ID_PMT) {
return Err(Error::InvalidValue {
field: "table_id",
value: section.first().copied().unwrap_or(0) as u64,
reason: "not a PMT section",
});
}
if section.len() < SECTION_HEADER_LEN {
return Err(Error::BufferTooShort {
need: SECTION_HEADER_LEN,
have: section.len(),
what: "PMT section header",
});
}
Ok(PmtSectionHeader {
program_number: u16::from_be_bytes([section[3], section[4]]),
version: (section[5] >> 1) & VERSION_NUMBER_MASK,
current_next: section[5] & CURRENT_NEXT_INDICATOR_BIT != 0,
section_number: section[6],
last_section_number: section[7],
})
}
fn parse_pmt(section: &[u8]) -> Result<Vec<(u16, Codec, Vec<u8>)>> {
if section.first().copied() != Some(TABLE_ID_PMT) {
return Ok(Vec::new());
}
let body = section_body(section, "PMT")?;
if body.len() < 4 {
return Err(Error::BufferTooShort {
need: 4,
have: body.len(),
what: "PMT program-info prefix",
});
}
let program_info_length = (((body[2] & INFO_LENGTH_HI_MASK) as usize) << 8) | body[3] as usize;
let mut off = 4 + program_info_length;
let mut out = Vec::new();
while off + 5 <= body.len() {
let stream_type = body[off];
let es_pid = (((body[off + 1] & PID_HI_MASK) as u16) << 8) | body[off + 2] as u16;
let es_info_length =
(((body[off + 3] & INFO_LENGTH_HI_MASK) as usize) << 8) | body[off + 4] as usize;
let desc_start = off + 5;
let desc_end = (desc_start + es_info_length).min(body.len());
let descriptors = body[desc_start..desc_end].to_vec();
let codec =
Codec::from_stream_type(stream_type).refine_with_descriptors(stream_type, &descriptors);
out.push((es_pid, codec, descriptors));
off += 5 + es_info_length;
}
Ok(out)
}
fn section_body<'a>(section: &'a [u8], what: &'static str) -> Result<&'a [u8]> {
if section.len() < SECTION_HEADER_LEN + CRC32_LEN {
return Err(Error::BufferTooShort {
need: SECTION_HEADER_LEN + CRC32_LEN,
have: section.len(),
what,
});
}
let section_length =
(((section[1] & SECTION_LENGTH_HI_MASK) as usize) << 8) | section[2] as usize;
let total = 3 + section_length;
let end = total.min(section.len());
if end < SECTION_HEADER_LEN + CRC32_LEN {
return Err(Error::BufferTooShort {
need: SECTION_HEADER_LEN + CRC32_LEN,
have: end,
what,
});
}
Ok(§ion[SECTION_HEADER_LEN..end - CRC32_LEN])
}
fn psi_section_crc_ok(section: &[u8]) -> bool {
if section.len() < SECTION_HEADER_LEN + CRC32_LEN {
return false;
}
if section[1] & SECTION_SYNTAX_INDICATOR_BIT == 0 {
return false;
}
let section_length =
(((section[1] & SECTION_LENGTH_HI_MASK) as usize) << 8) | section[2] as usize;
let total = 3 + section_length;
if total > section.len() || total < SECTION_HEADER_LEN + CRC32_LEN {
return false;
}
let (covered, trailer) = section[..total].split_at(total - CRC32_LEN);
let declared = u32::from_be_bytes([trailer[0], trailer[1], trailer[2], trailer[3]]);
broadcast_common::crc32_mpeg2::compute(covered) == declared
}
fn section_current_next(section: &[u8]) -> bool {
section
.get(5)
.is_some_and(|b| b & CURRENT_NEXT_INDICATOR_BIT != 0)
}
struct BufferedAu {
data: Vec<u8>,
pts_uw: i128,
dts_uw: i128,
}
enum ConfigProbe {
H264 {
sps: Option<Vec<u8>>,
pps: Option<Vec<u8>>,
},
Hevc {
vps: Option<Vec<u8>>,
sps: Option<Vec<u8>>,
pps: Option<Vec<u8>>,
},
Mpeg2Video,
MpegAudio {
is_mpeg2: bool,
},
Aac,
Ac3,
Eac3,
Dts,
MpegH,
Data,
}
fn initial_probe(codec: Codec) -> ConfigProbe {
match codec {
Codec::H264 => ConfigProbe::H264 {
sps: None,
pps: None,
},
Codec::Hevc => ConfigProbe::Hevc {
vps: None,
sps: None,
pps: None,
},
Codec::Mpeg2Video => ConfigProbe::Mpeg2Video,
Codec::MpegAudio(is_mpeg2) => ConfigProbe::MpegAudio { is_mpeg2 },
Codec::Aac => ConfigProbe::Aac,
Codec::Ac3 => ConfigProbe::Ac3,
Codec::Eac3 => ConfigProbe::Eac3,
Codec::Dts => ConfigProbe::Dts,
Codec::MpegH => ConfigProbe::MpegH,
Codec::Data(_) => ConfigProbe::Data,
}
}
#[derive(Clone, Copy)]
enum VideoCodec {
H264,
Hevc,
Mpeg2,
}
enum AudioKind {
Aac,
Ac3,
Eac3,
Dts,
MpegAudio { samples_per_frame: u32 },
}
#[derive(Default)]
struct AudioAnchor {
seed: Option<AudioAnchorSeed>,
}
#[derive(Clone, Copy)]
struct AudioAnchorSeed {
next_dts: i64,
next_pts: i64,
anchor_dts_uw: i128,
anchor_pts_uw: i128,
ticks_since_anchor: i64,
}
const AUDIO_REANCHOR_THRESHOLD_MS: i128 = 20;
pub(crate) fn audio_discontinuity_threshold_90k(sample_rate: u32) -> i128 {
let one_sample_period = (VIDEO_TIMESCALE as u128).div_ceil(sample_rate.max(1) as u128) as i128;
let ms_bound = (VIDEO_TIMESCALE as i128 * AUDIO_REANCHOR_THRESHOLD_MS) / 1000;
ms_bound.max(one_sample_period * 2)
}
struct PendingOneBehind {
data: Vec<u8>,
is_sync: bool,
pts_uw: i128,
dts_uw: i128,
}
fn to_ticks(uw: i128) -> i64 {
uw.clamp(i64::MIN as i128, i64::MAX as i128) as i64
}
fn ts_provenance(dts: i64, pts: i64) -> Provenance {
Provenance {
wire_dts: Some((dts as u64) % TS_WRAP),
wire_pts: Some((pts as u64) % TS_WRAP),
}
}
enum LiveKind {
Video {
pending: Option<PendingOneBehind>,
last_duration: u32,
codec: VideoCodec,
},
Audio {
sample_rate: u32,
kind: AudioKind,
anchor: AudioAnchor,
},
Data {
pending: Option<PendingOneBehind>,
last_duration: u32,
},
MpegH {
pending: Option<PendingOneBehind>,
last_duration: u32,
},
Section,
}
struct LiveTrack {
track_id: u32,
kind: LiveKind,
config: CodecConfig,
timescale: u32,
}
enum TrackState {
Probing {
probe: ConfigProbe,
backlog: Vec<BufferedAu>,
},
Parked {
config: CodecConfig,
timescale: u32,
kind: LiveKind,
backlog: Vec<BufferedAu>,
},
Live(LiveTrack),
Abandoned,
}
#[derive(Default)]
struct WrapState {
initialized: bool,
dts_seen_real: bool,
pts_seen_real: bool,
prev_dts_raw: u64,
prev_dts_uw: i128,
prev_pts_raw: u64,
prev_pts_uw: i128,
}
impl WrapState {
fn push(&mut self, raw_pts: u64, raw_dts: u64) -> (i128, i128) {
if !self.initialized {
self.initialized = true;
self.dts_seen_real = raw_dts != 0;
self.pts_seen_real = raw_pts != 0;
self.prev_dts_raw = raw_dts;
self.prev_dts_uw = raw_dts as i128;
self.prev_pts_raw = raw_pts;
self.prev_pts_uw = raw_pts as i128;
return (self.prev_pts_uw, self.prev_dts_uw);
}
let dts_uw = if self.dts_seen_real {
unwrap_ts(self.prev_dts_uw, self.prev_dts_raw, raw_dts)
} else {
self.dts_seen_real = raw_dts != 0;
raw_dts as i128
};
let pts_uw = if self.pts_seen_real {
unwrap_ts(self.prev_pts_uw, self.prev_pts_raw, raw_pts)
} else {
self.pts_seen_real = raw_pts != 0;
raw_pts as i128
};
self.prev_dts_raw = raw_dts;
self.prev_dts_uw = dts_uw;
self.prev_pts_raw = raw_pts;
self.prev_pts_uw = pts_uw;
(pts_uw, dts_uw)
}
}
enum Carrier {
Pes(PesAssembler),
Section(SectionReassembler),
}
fn initial_carrier(codec: Codec) -> Carrier {
match codec {
Codec::Data(stream_type) if data_carriage(stream_type) == DataCarriage::Sections => {
Carrier::Section(SectionReassembler::default())
}
_ => Carrier::Pes(PesAssembler::new()),
}
}
struct StreamState {
codec: Codec,
descriptors: Vec<u8>,
carrier: Carrier,
pes_bytes: usize,
fallback: (u64, u64),
has_any: bool,
wrap: WrapState,
track: Option<TrackState>,
backlog_bytes: usize,
}
#[allow(clippy::too_many_arguments)]
fn advance_one_behind(
pending: &mut Option<PendingOneBehind>,
last_duration: &mut u32,
data: Vec<u8>,
is_sync: bool,
pts_uw: i128,
dts_uw: i128,
duration_from_pts: bool,
track_id: u32,
events: &mut VecDeque<DemuxEvent>,
) {
if let Some(prev) = pending.take() {
let duration = if duration_from_pts {
(pts_uw - prev.pts_uw).max(0) as u32
} else {
(dts_uw - prev.dts_uw).max(0) as u32
};
*last_duration = duration;
let dts = to_ticks(prev.dts_uw);
let pts = to_ticks(prev.pts_uw);
events.push_back(DemuxEvent::Sample {
track_id,
sample: Sample {
data: prev.data.into(),
dts: Some(dts),
pts: Some(pts),
duration: Some(duration),
flags: SampleFlags::new(prev.is_sync),
provenance: Some(ts_provenance(dts, pts)),
},
});
}
*pending = Some(PendingOneBehind {
data,
is_sync,
pts_uw,
dts_uw,
});
}
fn flush_one_behind(
pending: &mut Option<PendingOneBehind>,
last_duration: u32,
track_id: u32,
events: &mut VecDeque<DemuxEvent>,
) {
if let Some(p) = pending.take() {
let dts = to_ticks(p.dts_uw);
let pts = to_ticks(p.pts_uw);
events.push_back(DemuxEvent::Sample {
track_id,
sample: Sample {
data: p.data.into(),
dts: Some(dts),
pts: Some(pts),
duration: Some(last_duration),
flags: SampleFlags::new(p.is_sync),
provenance: Some(ts_provenance(dts, pts)),
},
});
}
}
fn video_sample_bytes(codec: VideoCodec, au_data: &[u8]) -> (Vec<u8>, bool) {
match codec {
VideoCodec::H264 => {
let is_rap = access_unit_is_rap(NalCodec::Avc, au_data, false);
(annexb_to_length_prefixed(au_data), is_rap)
}
VideoCodec::Hevc => {
let mut irap = false;
for nal in iter_annexb_nals(au_data) {
if is_keyframe_nal(NalCodec::Hevc, nal) {
irap = true;
}
}
(annexb_to_length_prefixed(au_data), irap)
}
VideoCodec::Mpeg2 => (au_data.to_vec(), mpeg2_is_sync(au_data)),
}
}
#[allow(clippy::too_many_arguments)]
fn emit_audio_au(
kind: &AudioKind,
sample_rate: u32,
anchor: &mut AudioAnchor,
au_data: &[u8],
pts_uw: i128,
dts_uw: i128,
track_id: u32,
events: &mut VecDeque<DemuxEvent>,
) {
let seed = anchor.seed;
let fresh_anchor = match &seed {
None => true,
Some(seed) => {
let expected_dts_uw = seed.anchor_dts_uw
+ (seed.ticks_since_anchor as i128 * VIDEO_TIMESCALE as i128)
/ sample_rate.max(1) as i128;
(dts_uw - expected_dts_uw).abs() > audio_discontinuity_threshold_90k(sample_rate)
}
};
if fresh_anchor && seed.is_some() {
events.push_back(DemuxEvent::Discontinuity {
track: Some(track_id),
kind: DiscontinuityKind::TimelineReanchored,
provenance: EventProvenance::default(),
});
}
let (dts0, pts0) = match (fresh_anchor, &seed) {
(false, Some(seed)) => (seed.next_dts, seed.next_pts),
_ => (
rescale_to_track(dts_uw, sample_rate),
rescale_to_track(pts_uw, sample_rate),
),
};
let mut elapsed = 0u64;
let au_provenance = ts_provenance(to_ticks(dts_uw), to_ticks(pts_uw));
let audio_sample = |data: Vec<u8>, duration: u32, elapsed: u64| -> Sample {
let dts = dts0 + elapsed as i64;
let pts = pts0 + elapsed as i64;
Sample::from_raw(data, Some(dts), Some(pts), Some(duration)).with_provenance(au_provenance)
};
match kind {
AudioKind::Aac => {
for frame in split_adts_frames(au_data) {
if frame.len() > ADTS_HEADER_SIZE {
events.push_back(DemuxEvent::Sample {
track_id,
sample: audio_sample(
frame[ADTS_HEADER_SIZE..].to_vec(),
AAC_SAMPLES_PER_FRAME,
elapsed,
),
});
}
elapsed += AAC_SAMPLES_PER_FRAME as u64;
}
}
AudioKind::Ac3 => {
for frame in split_ac3_syncframes(au_data) {
events.push_back(DemuxEvent::Sample {
track_id,
sample: audio_sample(frame.to_vec(), AC3_SAMPLES_PER_SYNCFRAME, elapsed),
});
elapsed += AC3_SAMPLES_PER_SYNCFRAME as u64;
}
}
AudioKind::Eac3 => {
for split in split_eac3_syncframes(au_data) {
let duration = split.info.samples_per_frame();
events.push_back(DemuxEvent::Sample {
track_id,
sample: audio_sample(split.data, duration, elapsed),
});
elapsed += duration as u64;
}
}
AudioKind::Dts => {
for frame in split_dts_core_frames(au_data) {
events.push_back(DemuxEvent::Sample {
track_id,
sample: audio_sample(frame.data.to_vec(), frame.samples, elapsed),
});
elapsed += frame.samples as u64;
}
}
AudioKind::MpegAudio { samples_per_frame } => {
for frame in split_mpeg_audio_frames(au_data) {
events.push_back(DemuxEvent::Sample {
track_id,
sample: audio_sample(frame.to_vec(), *samples_per_frame, elapsed),
});
elapsed += *samples_per_frame as u64;
}
}
}
let (anchor_dts_uw, anchor_pts_uw, ticks_since_anchor) = match (fresh_anchor, &seed) {
(false, Some(seed)) => (
seed.anchor_dts_uw,
seed.anchor_pts_uw,
seed.ticks_since_anchor,
),
_ => (dts_uw, pts_uw, 0i64),
};
anchor.seed = Some(AudioAnchorSeed {
next_dts: dts0 + elapsed as i64,
next_pts: pts0 + elapsed as i64,
anchor_dts_uw,
anchor_pts_uw,
ticks_since_anchor: ticks_since_anchor + elapsed as i64,
});
}
fn push_live_au(
live: &mut LiveTrack,
data: &[u8],
pts_uw: i128,
dts_uw: i128,
events: &mut VecDeque<DemuxEvent>,
) {
let track_id = live.track_id;
match &mut live.kind {
LiveKind::Video {
pending,
last_duration,
codec,
} => {
let (bytes, is_sync) = video_sample_bytes(*codec, data);
advance_one_behind(
pending,
last_duration,
bytes,
is_sync,
pts_uw,
dts_uw,
false,
track_id,
events,
);
}
LiveKind::Data {
pending,
last_duration,
} => {
advance_one_behind(
pending,
last_duration,
data.to_vec(),
true,
pts_uw,
dts_uw,
true,
track_id,
events,
);
}
LiveKind::Audio {
sample_rate,
kind,
anchor,
} => {
emit_audio_au(
kind,
*sample_rate,
anchor,
data,
pts_uw,
dts_uw,
track_id,
events,
);
}
LiveKind::MpegH {
pending,
last_duration,
} => {
let is_sync = find_mpegh3da_config(data).is_some();
advance_one_behind(
pending,
last_duration,
data.to_vec(),
is_sync,
pts_uw,
dts_uw,
true,
track_id,
events,
);
}
LiveKind::Section => {
events.push_back(DemuxEvent::Sample {
track_id,
sample: Sample::from_raw(data.to_vec(), None, None, None),
});
}
}
}
fn finalize_probe(
codec: Codec,
descriptors: &[u8],
probe: &mut ConfigProbe,
backlog: &[BufferedAu],
) -> Option<(CodecConfig, u32, LiveKind)> {
let latest = backlog.last()?;
match probe {
ConfigProbe::Data => {
let Codec::Data(stream_type) = codec else {
return None;
};
let carriage = data_carriage(stream_type);
let kind = match carriage {
DataCarriage::Pes => LiveKind::Data {
pending: None,
last_duration: 0,
},
DataCarriage::Sections => LiveKind::Section,
};
Some((
CodecConfig::Data {
stream_type,
descriptors: descriptors.to_vec(),
carriage,
},
VIDEO_TIMESCALE,
kind,
))
}
ConfigProbe::H264 { sps, pps } => {
for nal in iter_annexb_nals(&latest.data) {
match nal[0] & H264_NAL_TYPE_MASK {
H264_NAL_SPS if sps.is_none() => *sps = Some(nal.to_vec()),
H264_NAL_PPS if pps.is_none() => *pps = Some(nal.to_vec()),
_ => {}
}
}
let (sps_bytes, pps_bytes) = (sps.as_ref()?, pps.as_ref()?);
if sps_bytes.len() < 4 {
return None;
}
let info = crate::sps::decode_avc_sps(sps_bytes).ok();
let (width, height) = info
.as_ref()
.map(|i| (i.width as u16, i.height as u16))
.unwrap_or((0, 0));
let ext = info
.as_ref()
.filter(|i| crate::sps::is_high_profile(i.profile_idc));
let record = AVCDecoderConfigurationRecord {
configuration_version: 1,
profile_indication: sps_bytes[1],
profile_compatibility: sps_bytes[2],
level_indication: sps_bytes[3],
length_size_minus_one: NAL_LENGTH_SIZE_MINUS_ONE,
sps: alloc::vec![AvcSps(sps_bytes.clone())],
pps: alloc::vec![AvcPps(pps_bytes.clone())],
chroma_format: ext.map(|i| i.chroma_format_idc),
bit_depth_luma_minus8: ext.map(|i| i.bit_depth_luma.saturating_sub(8)),
bit_depth_chroma_minus8: ext.map(|i| i.bit_depth_chroma.saturating_sub(8)),
sps_ext: alloc::vec![],
};
Some((
CodecConfig::Avc {
config: AVCConfigurationBox::new(record),
width,
height,
},
VIDEO_TIMESCALE,
LiveKind::Video {
pending: None,
last_duration: 0,
codec: VideoCodec::H264,
},
))
}
ConfigProbe::Hevc { vps, sps, pps } => {
for nal in iter_annexb_nals(&latest.data) {
match nal_unit_type(NalCodec::Hevc, nal) {
Some(H265_NAL_VPS) if vps.is_none() => *vps = Some(nal.to_vec()),
Some(H265_NAL_SPS) if sps.is_none() => *sps = Some(nal.to_vec()),
Some(H265_NAL_PPS) if pps.is_none() => *pps = Some(nal.to_vec()),
_ => {}
}
}
let sps_bytes = sps.as_ref()?;
let info = crate::sps::decode_hevc_sps(sps_bytes).ok()?;
let width = info.width.min(u16::MAX as u32) as u16;
let height = info.height.min(u16::MAX as u32) as u16;
let mut arrays: Vec<HevcNalArray> = Vec::new();
if let Some(vps_nal) = vps.clone() {
arrays.push(HevcNalArray::new(
true,
H265_NAL_VPS,
alloc::vec![HevcNalUnit::new(vps_nal)],
));
}
arrays.push(HevcNalArray::new(
true,
H265_NAL_SPS,
alloc::vec![HevcNalUnit::new(sps_bytes.clone())],
));
if let Some(pps_nal) = pps.clone() {
arrays.push(HevcNalArray::new(
true,
H265_NAL_PPS,
alloc::vec![HevcNalUnit::new(pps_nal)],
));
}
let record = HEVCDecoderConfigurationRecord {
configuration_version: HVCC_CONFIGURATION_VERSION,
general_profile_space: info.general_profile_space,
general_tier_flag: info.general_tier_flag,
general_profile_idc: info.general_profile_idc,
general_profile_compatibility_flags: info.general_profile_compatibility_flags,
general_constraint_indicator_flags: info.general_constraint_indicator_flags,
general_level_idc: info.general_level_idc,
min_spatial_segmentation_idc: HVCC_MIN_SPATIAL_SEGMENTATION_UNSPEC,
parallelism_type: HVCC_PARALLELISM_TYPE_UNKNOWN,
chroma_format_idc: info.chroma_format_idc,
bit_depth_luma_minus8: info.bit_depth_luma.saturating_sub(8),
bit_depth_chroma_minus8: info.bit_depth_chroma.saturating_sub(8),
avg_frame_rate: HVCC_AVG_FRAME_RATE_UNSPEC,
constant_frame_rate: HVCC_CONSTANT_FRAME_RATE_UNSPEC,
num_temporal_layers: HVCC_NUM_TEMPORAL_LAYERS,
temporal_id_nested: false,
length_size_minus_one: NAL_LENGTH_SIZE_MINUS_ONE,
arrays,
};
Some((
CodecConfig::Hevc {
config: HEVCConfigurationBox::new(record),
width,
height,
},
VIDEO_TIMESCALE,
LiveKind::Video {
pending: None,
last_duration: 0,
codec: VideoCodec::Hevc,
},
))
}
ConfigProbe::Mpeg2Video => {
let seq = backlog
.iter()
.find_map(|au| Mpeg2SeqHeader::find(&au.data).ok())?;
let esds = EsdsBox::new(ESDescriptor {
es_id: ESDS_VIDEO_ES_ID,
stream_dependence_flag: false,
url_flag: false,
ocr_stream_flag: false,
stream_priority: 0,
depends_on_es_id: None,
url: None,
ocr_es_id: None,
decoder_config: Some(DecoderConfigDescriptor {
object_type_indication: ObjectTypeIndication(OTI_MPEG2_VIDEO_MAIN),
stream_type: EsdsStreamType(STREAM_TYPE_VISUAL),
up_stream: false,
buffer_size_db: 0,
max_bitrate: 0,
avg_bitrate: 0,
decoder_specific_info: None,
}),
sl_config: Some(SLConfigDescriptor {
body: alloc::vec![SL_CONFIG_PREDEFINED_MP4],
}),
});
Some((
CodecConfig::Mpeg2Video {
esds,
width: seq.width,
height: seq.height,
},
VIDEO_TIMESCALE,
LiveKind::Video {
pending: None,
last_duration: 0,
codec: VideoCodec::Mpeg2,
},
))
}
ConfigProbe::MpegAudio { is_mpeg2 } => {
let first = backlog
.iter()
.find_map(|au| find_mpeg_audio_sync(&au.data).map(|(_, hdr)| hdr))?;
let sample_rate = first.sample_rate;
let channel_count = first.channels;
let samples_per_frame = first.samples_per_frame;
let oti = if *is_mpeg2 {
OTI_MPEG2_AUDIO
} else {
OTI_MPEG1_AUDIO
};
let esds = EsdsBox::new(ESDescriptor {
es_id: ESDS_ES_ID,
stream_dependence_flag: false,
url_flag: false,
ocr_stream_flag: false,
stream_priority: 0,
depends_on_es_id: None,
url: None,
ocr_es_id: None,
decoder_config: Some(DecoderConfigDescriptor {
object_type_indication: ObjectTypeIndication(oti),
stream_type: EsdsStreamType(STREAM_TYPE_AUDIO),
up_stream: false,
buffer_size_db: 0,
max_bitrate: 0,
avg_bitrate: 0,
decoder_specific_info: None,
}),
sl_config: Some(SLConfigDescriptor {
body: alloc::vec![SL_CONFIG_PREDEFINED_MP4],
}),
});
Some((
CodecConfig::MpegAudio {
esds,
layer: first.layer,
channel_count,
sample_rate,
sample_size: AUDIO_SAMPLE_SIZE_BITS,
},
sample_rate,
LiveKind::Audio {
sample_rate,
kind: AudioKind::MpegAudio { samples_per_frame },
anchor: AudioAnchor::default(),
},
))
}
ConfigProbe::Aac => {
let first_hdr = backlog
.iter()
.find_map(|au| find_adts_sync(&au.data).map(|(_, hdr)| hdr))?;
let asc = AudioSpecificConfig::from_adts_header(&first_hdr);
let sample_rate = sfi_to_hz(first_hdr.sampling_frequency_index)?;
let channel_count = first_hdr.channel_configuration as u16;
let esds = EsdsBox::new(ESDescriptor {
es_id: ESDS_ES_ID,
stream_dependence_flag: false,
url_flag: false,
ocr_stream_flag: false,
stream_priority: 0,
depends_on_es_id: None,
url: None,
ocr_es_id: None,
decoder_config: Some(DecoderConfigDescriptor {
object_type_indication: ObjectTypeIndication(OTI_MPEG4_AUDIO),
stream_type: EsdsStreamType(STREAM_TYPE_AUDIO),
up_stream: false,
buffer_size_db: 0,
max_bitrate: 0,
avg_bitrate: 0,
decoder_specific_info: Some(DecoderSpecificInfo {
data: asc.to_bytes(),
}),
}),
sl_config: Some(SLConfigDescriptor {
body: alloc::vec![SL_CONFIG_PREDEFINED_MP4],
}),
});
Some((
CodecConfig::Aac {
esds,
channel_count,
sample_rate,
sample_size: AUDIO_SAMPLE_SIZE_BITS,
},
sample_rate,
LiveKind::Audio {
sample_rate,
kind: AudioKind::Aac,
anchor: AudioAnchor::default(),
},
))
}
ConfigProbe::Ac3 => {
let info = backlog
.iter()
.find_map(|au| Ac3SyncframeInfo::from_es(&au.data).ok())?;
let sample_rate = info.sample_rate;
let channel_count = info.channel_count() as u16;
let config = info.into_dac3();
Some((
CodecConfig::Ac3 {
config,
channel_count,
sample_rate,
sample_size: AUDIO_SAMPLE_SIZE_BITS,
},
sample_rate,
LiveKind::Audio {
sample_rate,
kind: AudioKind::Ac3,
anchor: AudioAnchor::default(),
},
))
}
ConfigProbe::Eac3 => {
let info = backlog
.iter()
.find_map(|au| Ec3SyncframeInfo::from_es(&au.data).ok())?;
let sample_rate = info.sample_rate;
let channel_count = info.channel_count() as u16;
let config = info.into_dec3();
Some((
CodecConfig::Eac3 {
config,
channel_count,
sample_rate,
sample_size: AUDIO_SAMPLE_SIZE_BITS,
},
sample_rate,
LiveKind::Audio {
sample_rate,
kind: AudioKind::Eac3,
anchor: AudioAnchor::default(),
},
))
}
ConfigProbe::Dts => {
let info = backlog
.iter()
.find_map(|au| DtsCoreFrameInfo::from_es(&au.data).ok())?;
let sample_rate = info.sample_rate;
let channel_count = info.channels as u16;
let config = info.into_ddts();
Some((
CodecConfig::Dts {
config,
codec_fourcc: crate::dts::DTSC_FOURCC,
channel_count,
sample_rate,
sample_size: AUDIO_SAMPLE_SIZE_BITS,
},
sample_rate,
LiveKind::Audio {
sample_rate,
kind: AudioKind::Dts,
anchor: AudioAnchor::default(),
},
))
}
ConfigProbe::MpegH => {
let config_bytes = backlog
.iter()
.find_map(|au| find_mpegh3da_config(&au.data))?;
let profile_level_indication = *config_bytes.first()?;
let config = MHADecoderConfigurationRecord::new(
profile_level_indication,
MPEGH_REFERENCE_CHANNEL_LAYOUT_UNSPECIFIED,
config_bytes.to_vec(),
);
Some((
CodecConfig::MpegH {
config,
channel_count: MPEGH_CHANNEL_COUNT_UNSPECIFIED,
sample_rate: MPEGH_SAMPLE_RATE_UNSPECIFIED,
sample_size: AUDIO_SAMPLE_SIZE_BITS,
},
VIDEO_TIMESCALE,
LiveKind::MpegH {
pending: None,
last_duration: 0,
},
))
}
}
}
fn abandon_backlog(
stream: &mut StreamState,
pid: u16,
events: &mut VecDeque<DemuxEvent>,
) -> TrackState {
stream.backlog_bytes = 0;
events.push_back(DemuxEvent::TrackAbandoned {
track_id: None,
reason: AbandonReason::BudgetExceeded,
provenance: EventProvenance {
pid: Some(pid),
packet_index: None,
},
});
TrackState::Abandoned
}
fn advance_track(
stream: &mut StreamState,
pid: u16,
data: Vec<u8>,
pts_uw: i128,
dts_uw: i128,
events: &mut VecDeque<DemuxEvent>,
) {
let Some(track) = stream.track.take() else {
return;
};
let new_track = match track {
TrackState::Live(mut live) => {
push_live_au(&mut live, &data, pts_uw, dts_uw, events);
TrackState::Live(live)
}
TrackState::Abandoned => TrackState::Abandoned,
TrackState::Parked {
config,
timescale,
kind,
mut backlog,
} => {
stream.backlog_bytes = stream.backlog_bytes.saturating_add(data.len());
backlog.push(BufferedAu {
data,
pts_uw,
dts_uw,
});
if stream.backlog_bytes > MAX_PROBE_BACKLOG_BYTES {
abandon_backlog(stream, pid, events)
} else {
TrackState::Parked {
config,
timescale,
kind,
backlog,
}
}
}
TrackState::Probing {
mut probe,
mut backlog,
} => {
stream.backlog_bytes = stream.backlog_bytes.saturating_add(data.len());
backlog.push(BufferedAu {
data,
pts_uw,
dts_uw,
});
match finalize_probe(stream.codec, &stream.descriptors, &mut probe, &backlog) {
Some((config, timescale, kind)) => TrackState::Parked {
config,
timescale,
kind,
backlog,
},
None if stream.backlog_bytes > MAX_PROBE_BACKLOG_BYTES => {
abandon_backlog(stream, pid, events)
}
None => TrackState::Probing { probe, backlog },
}
}
};
stream.track = Some(new_track);
}
fn feed_pes_bounded(
stream: &mut StreamState,
pid: u16,
pusi: bool,
payload: &[u8],
events: &mut VecDeque<DemuxEvent>,
) -> Option<Vec<u8>> {
let Carrier::Pes(assembler) = &mut stream.carrier else {
return None;
};
if pusi {
stream.pes_bytes = payload.len();
} else if stream.pes_bytes > 0 {
stream.pes_bytes = stream.pes_bytes.saturating_add(payload.len());
}
let completed = assembler.feed(pusi, payload);
if stream.pes_bytes > MAX_PES_BUFFER_BYTES {
let dropped_bytes = stream.pes_bytes as u64;
let _ = assembler.flush();
stream.pes_bytes = 0;
let track = match stream.track.as_ref() {
Some(TrackState::Live(live)) => Some(live.track_id),
_ => None,
};
events.push_back(DemuxEvent::Discontinuity {
track,
kind: DiscontinuityKind::BudgetExceeded {
bytes: dropped_bytes,
},
provenance: EventProvenance {
pid: Some(pid),
packet_index: None,
},
});
}
completed
}
fn on_completed_pes(
stream: &mut StreamState,
pid: u16,
pes_bytes: &[u8],
events: &mut VecDeque<DemuxEvent>,
) {
let Ok(pes) = PesPacket::parse(pes_bytes) else {
return;
};
if pes.payload.is_empty() {
return;
}
let fallback = if stream.has_any {
stream.fallback
} else {
(0, 0)
};
let (pts, dts) = match pes.header.as_ref() {
Some(h) => {
let hp = h.pts.map(|p| p.0);
let hd = h.dts.map(|d| d.0);
let pts = hp.or(hd).unwrap_or(fallback.0);
let dts = hd.unwrap_or(pts);
(pts, dts)
}
None => fallback,
};
stream.fallback = (pts, dts);
stream.has_any = true;
let (pts_uw, dts_uw) = stream.wrap.push(pts, dts);
advance_track(stream, pid, pes.payload.to_vec(), pts_uw, dts_uw, events);
}
fn on_completed_section(
stream: &mut StreamState,
pid: u16,
section: &[u8],
events: &mut VecDeque<DemuxEvent>,
) {
if section.is_empty() {
return;
}
advance_track(stream, pid, section.to_vec(), 0, 0, events);
}
pub use crate::ir::{AbandonReason, DemuxEvent, DiscontinuityKind, EventProvenance};
const PCR_CLOCK_HZ: u32 = 27_000_000;
pub struct StreamingTsDemux {
resync: TsResync,
packet_index: u64,
pat_reasm: SectionReassembler,
pmt_reasm: BTreeMap<u16, PmtState>,
es_seen: BTreeSet<u16>,
streams: BTreeMap<u16, StreamState>,
unattributed: BTreeMap<u16, VecDeque<(bool, Vec<u8>)>>,
removed_pids: BTreeSet<u16>,
es_declarers: BTreeMap<u16, BTreeSet<u16>>,
unattributed_order: VecDeque<u16>,
unattributed_bytes: usize,
codec_order: Vec<u16>,
data_order: Vec<u16>,
resolved: BTreeSet<u16>,
next_track_id: u32,
events: VecDeque<DemuxEvent>,
generation: u32,
tracks_resolved_signalled_at: Option<u32>,
}
struct PmtState {
reasm: SectionReassembler,
program_number: u16,
last_applied_version: Option<u8>,
applied_es: BTreeSet<u16>,
}
impl Default for StreamingTsDemux {
fn default() -> Self {
Self::new()
}
}
impl StreamingTsDemux {
pub fn new() -> Self {
Self {
resync: TsResync::new(),
packet_index: 0,
pat_reasm: SectionReassembler::default(),
pmt_reasm: BTreeMap::new(),
es_seen: BTreeSet::new(),
streams: BTreeMap::new(),
unattributed: BTreeMap::new(),
removed_pids: BTreeSet::new(),
es_declarers: BTreeMap::new(),
unattributed_order: VecDeque::new(),
unattributed_bytes: 0,
codec_order: Vec::new(),
data_order: Vec::new(),
resolved: BTreeSet::new(),
next_track_id: 1,
events: VecDeque::new(),
generation: 0,
tracks_resolved_signalled_at: None,
}
}
pub fn feed(&mut self, data: &[u8]) {
let packets = self.resync.feed(data);
for raw in &packets {
self.process_packet(raw);
}
}
fn live_track_id(&self, pid: u16) -> Option<u32> {
match self.streams.get(&pid)?.track.as_ref()? {
TrackState::Live(live) => Some(live.track_id),
_ => None,
}
}
fn process_packet(&mut self, raw: &[u8; TS_PACKET_SIZE]) {
let idx = self.packet_index;
self.packet_index += 1;
let Ok(pkt) = TsPacket::parse(raw) else {
return;
};
if let Some(Ok(af)) = pkt.adaptation_field() {
let provenance = EventProvenance {
pid: Some(pkt.header.pid),
packet_index: Some(idx),
};
if af.discontinuity_indicator {
self.events.push_back(DemuxEvent::Discontinuity {
track: self.live_track_id(pkt.header.pid),
kind: DiscontinuityKind::Signalled,
provenance,
});
}
if let Some(pcr) = af.pcr {
self.events.push_back(DemuxEvent::ClockReference {
ticks: pcr.as_27mhz(),
clock_hz: PCR_CLOCK_HZ,
discontinuous: af.discontinuity_indicator,
provenance,
});
}
}
let pid = pkt.header.pid;
let pusi = pkt.header.pusi;
let Some(payload) = pkt.payload else {
return;
};
if pid == PAT_PID {
self.pat_reasm.feed(payload, pusi);
while let Some(section) = self.pat_reasm.pop_section() {
if !psi_section_crc_ok(§ion) {
continue;
}
if !section_current_next(§ion) {
continue;
}
if let Ok(programs) = parse_pat(§ion) {
for (program_number, pmt_pid) in programs {
self.learn_pmt_pid(pmt_pid, program_number);
}
}
}
return;
}
if let Some(pmt_state) = self.pmt_reasm.get_mut(&pid) {
pmt_state.reasm.feed(payload, pusi);
let mut sections: Vec<Vec<u8>> = Vec::new();
while let Some(section) = pmt_state.reasm.pop_section() {
sections.push(section.to_vec());
}
let program_number = pmt_state.program_number;
let mut to_apply: Option<Vec<(u16, Codec, Vec<u8>)>> = None;
for section in §ions {
if !psi_section_crc_ok(section) {
continue;
}
let Ok(header) = parse_pmt_section_header(section) else {
continue;
};
if header.program_number != program_number {
continue;
}
if header.section_number != 0 || header.last_section_number != 0 {
continue;
}
if !header.current_next {
continue;
}
if pmt_state.last_applied_version == Some(header.version) {
continue;
}
pmt_state.last_applied_version = Some(header.version);
if let Ok(es_list) = parse_pmt(section) {
to_apply = Some(es_list);
}
}
if let Some(es_list) = to_apply {
let old_applied_es = pmt_state.applied_es.clone();
let new_applied_es: BTreeSet<u16> = es_list.iter().map(|(p, _, _)| *p).collect();
self.apply_pmt_diff(pid, &old_applied_es, es_list);
if let Some(pmt_state) = self.pmt_reasm.get_mut(&pid) {
pmt_state.applied_es = new_applied_es;
}
}
self.try_promote_ready();
return;
}
if let Some(stream) = self.streams.get_mut(&pid) {
let mut sections: Vec<Vec<u8>> = Vec::new();
let completed_pes = if matches!(stream.carrier, Carrier::Pes(_)) {
feed_pes_bounded(stream, pid, pusi, payload, &mut self.events)
} else if let Carrier::Section(reasm) = &mut stream.carrier {
reasm.feed(payload, pusi);
while let Some(s) = reasm.pop_section() {
sections.push(s.to_vec());
}
None
} else {
None
};
if let Some(completed) = completed_pes {
on_completed_pes(stream, pid, &completed, &mut self.events);
}
for s in sections {
on_completed_section(stream, pid, &s, &mut self.events);
}
} else if pid != NULL_PACKET_PID && !self.removed_pids.contains(&pid) {
self.unattributed
.entry(pid)
.or_default()
.push_back((pusi, payload.to_vec()));
self.unattributed_order.push_back(pid);
self.unattributed_bytes += payload.len();
self.evict_unattributed();
}
self.try_promote_ready();
}
fn learn_pmt_pid(&mut self, pmt_pid: u16, program_number: u16) {
match self.pmt_reasm.get_mut(&pmt_pid) {
Some(state) if state.program_number != program_number => {
state.program_number = program_number;
state.last_applied_version = None;
}
Some(_) => {}
None => {
self.pmt_reasm.insert(
pmt_pid,
PmtState {
reasm: SectionReassembler::default(),
program_number,
last_applied_version: None,
applied_es: BTreeSet::new(),
},
);
}
}
}
fn evict_unattributed(&mut self) {
while self.unattributed_bytes > MAX_UNATTRIBUTED_BYTES {
let Some(pid) = self.unattributed_order.pop_front() else {
break;
};
if let Some(buf) = self.unattributed.get_mut(&pid) {
if let Some((_, payload)) = buf.pop_front() {
self.unattributed_bytes = self.unattributed_bytes.saturating_sub(payload.len());
}
if buf.is_empty() {
self.unattributed.remove(&pid);
}
}
}
}
fn register_new_es(&mut self, es_pid: u16, codec: Codec, descriptors: Vec<u8>) {
self.register_new_es_at(es_pid, codec, descriptors, None);
}
fn register_new_es_at(
&mut self,
es_pid: u16,
codec: Codec,
descriptors: Vec<u8>,
reinsert_at: Option<usize>,
) {
self.removed_pids.remove(&es_pid);
let order = if matches!(codec, Codec::Data(_)) {
&mut self.data_order
} else {
&mut self.codec_order
};
match reinsert_at {
Some(idx) if idx <= order.len() => order.insert(idx, es_pid),
_ => order.push(es_pid),
}
let mut stream = StreamState {
codec,
descriptors,
carrier: initial_carrier(codec),
pes_bytes: 0,
fallback: (0, 0),
has_any: false,
wrap: WrapState::default(),
track: Some(TrackState::Probing {
probe: initial_probe(codec),
backlog: Vec::new(),
}),
backlog_bytes: 0,
};
if let Some(buffered) = self.unattributed.remove(&es_pid) {
for (buf_pusi, buf_payload) in buffered {
self.unattributed_bytes = self.unattributed_bytes.saturating_sub(buf_payload.len());
let mut sections: Vec<Vec<u8>> = Vec::new();
let completed_pes = if matches!(stream.carrier, Carrier::Pes(_)) {
feed_pes_bounded(
&mut stream,
es_pid,
buf_pusi,
&buf_payload,
&mut self.events,
)
} else if let Carrier::Section(reasm) = &mut stream.carrier {
reasm.feed(&buf_payload, buf_pusi);
while let Some(s) = reasm.pop_section() {
sections.push(s.to_vec());
}
None
} else {
None
};
if let Some(completed) = completed_pes {
on_completed_pes(&mut stream, es_pid, &completed, &mut self.events);
}
for s in sections {
on_completed_section(&mut stream, es_pid, &s, &mut self.events);
}
}
}
self.streams.insert(es_pid, stream);
}
fn remove_track(&mut self, pid: u16) {
self.es_seen.remove(&pid);
self.codec_order.retain(|&p| p != pid);
self.data_order.retain(|&p| p != pid);
self.resolved.remove(&pid);
if let Some(buffered) = self.unattributed.remove(&pid) {
for (_, payload) in &buffered {
self.unattributed_bytes = self.unattributed_bytes.saturating_sub(payload.len());
}
}
self.removed_pids.insert(pid);
if let Some(stream) = self.streams.remove(&pid) {
if let Some(TrackState::Live(live)) = stream.track {
self.events.push_back(DemuxEvent::TrackRemoved {
track_id: live.track_id,
provenance: EventProvenance {
pid: Some(pid),
packet_index: None,
},
});
}
}
}
fn apply_pmt_diff(
&mut self,
pmt_pid: u16,
old_applied_es: &BTreeSet<u16>,
es_list: Vec<(u16, Codec, Vec<u8>)>,
) {
let new_pids: BTreeSet<u16> = es_list.iter().map(|(p, _, _)| *p).collect();
let removed: Vec<u16> = old_applied_es
.iter()
.copied()
.filter(|p| !new_pids.contains(p))
.collect();
for pid in removed {
let declarers = self.es_declarers.get_mut(&pid);
let still_declared = match declarers {
Some(declarers) => {
declarers.remove(&pmt_pid);
let empty = declarers.is_empty();
if empty {
self.es_declarers.remove(&pid);
}
!empty
}
None => false,
};
if !still_declared {
self.remove_track(pid);
}
}
for (es_pid, codec, descriptors) in es_list {
self.es_declarers.entry(es_pid).or_default().insert(pmt_pid);
if self.es_seen.insert(es_pid) {
self.register_new_es(es_pid, codec, descriptors);
continue;
}
let Some(existing) = self.streams.get(&es_pid) else {
continue;
};
let old_codec = existing.codec;
let codec_changed = old_codec != codec;
let descriptors_changed = existing.descriptors != descriptors;
if !codec_changed && !descriptors_changed {
continue;
}
if codec_changed {
let other_declarers_remain = self
.es_declarers
.get(&es_pid)
.is_some_and(|declarers| declarers.iter().any(|&p| p != pmt_pid));
if other_declarers_remain {
continue;
}
let old_is_data = matches!(old_codec, Codec::Data(_));
let new_is_data = matches!(codec, Codec::Data(_));
let order_slot = if old_is_data == new_is_data {
let order = if old_is_data {
&self.data_order
} else {
&self.codec_order
};
order.iter().position(|&p| p == es_pid)
} else {
None
};
self.remove_track(es_pid);
self.es_seen.insert(es_pid);
self.register_new_es_at(es_pid, codec, descriptors, order_slot);
continue;
}
let Some(stream) = self.streams.get_mut(&es_pid) else {
continue;
};
stream.descriptors = descriptors.clone();
if let Some(TrackState::Live(live)) = stream.track.as_ref() {
let spec = TrackSpec::new(live.track_id, live.timescale, live.config.clone())
.with_source(es_pid, descriptors);
self.events.push_back(DemuxEvent::TrackUpdated(spec));
}
}
self.generation = self.generation.wrapping_add(1);
}
fn try_promote_ready(&mut self) {
while let Some(&next_pid) = self
.codec_order
.iter()
.chain(self.data_order.iter())
.find(|p| !self.resolved.contains(p))
{
let Some(stream) = self.streams.get_mut(&next_pid) else {
break;
};
let Some(track) = stream.track.take() else {
break;
};
match track {
TrackState::Parked {
config,
timescale,
kind,
backlog,
} => {
let track_id = self.next_track_id;
self.next_track_id += 1;
let spec = TrackSpec::new(track_id, timescale, config.clone())
.with_source(next_pid, stream.descriptors.clone());
self.events.push_back(DemuxEvent::TrackAdded(spec));
let mut live = LiveTrack {
track_id,
kind,
config,
timescale,
};
for au in backlog {
push_live_au(&mut live, &au.data, au.pts_uw, au.dts_uw, &mut self.events);
}
stream.track = Some(TrackState::Live(live));
stream.backlog_bytes = 0;
self.resolved.insert(next_pid);
}
other @ TrackState::Probing { .. } => {
stream.track = Some(other);
break; }
other @ TrackState::Live(_) => {
stream.track = Some(other);
self.resolved.insert(next_pid);
}
TrackState::Abandoned => {
stream.track = Some(TrackState::Abandoned);
self.resolved.insert(next_pid);
}
}
}
self.maybe_signal_tracks_resolved();
}
fn maybe_signal_tracks_resolved(&mut self) {
let known = self.codec_order.len() + self.data_order.len();
if known == 0 {
return;
}
if self.resolved.len() == known
&& self.tracks_resolved_signalled_at != Some(self.generation)
{
self.tracks_resolved_signalled_at = Some(self.generation);
self.events.push_back(DemuxEvent::TracksResolved {
generation: self.generation,
});
}
}
pub fn poll_event(&mut self) -> Option<DemuxEvent> {
self.events.pop_front()
}
pub fn finish(&mut self) {
for (&pid, stream) in self.streams.iter_mut() {
let completed = match &mut stream.carrier {
Carrier::Pes(assembler) => assembler.flush(),
Carrier::Section(_) => None,
};
if let Some(completed) = completed {
on_completed_pes(stream, pid, &completed, &mut self.events);
}
}
self.try_promote_ready();
while let Some(&next_pid) = self
.codec_order
.iter()
.chain(self.data_order.iter())
.find(|p| !self.resolved.contains(p))
{
match self.streams.get(&next_pid).and_then(|s| s.track.as_ref()) {
Some(TrackState::Probing { .. }) => {
self.resolved.insert(next_pid);
self.events.push_back(DemuxEvent::TrackAbandoned {
track_id: None,
reason: AbandonReason::ConfigUnrecoverable,
provenance: EventProvenance {
pid: Some(next_pid),
packet_index: None,
},
});
self.try_promote_ready();
}
_ => break,
}
}
for stream in self.streams.values_mut() {
if let Some(TrackState::Live(live)) = &mut stream.track {
match &mut live.kind {
LiveKind::Video {
pending,
last_duration,
..
} => {
flush_one_behind(pending, *last_duration, live.track_id, &mut self.events);
}
LiveKind::Data {
pending,
last_duration,
}
| LiveKind::MpegH {
pending,
last_duration,
} => {
flush_one_behind(pending, *last_duration, live.track_id, &mut self.events);
}
LiveKind::Audio { .. } => {}
LiveKind::Section => {}
}
}
}
}
}
impl Stage for StreamingTsDemux {
type In<'a> = &'a [u8];
type Out = DemuxEvent;
type Error = core::convert::Infallible;
fn feed(&mut self, input: &[u8], _now: Timestamp) -> core::result::Result<(), Self::Error> {
self.feed(input);
Ok(())
}
fn poll(&mut self) -> Option<Self::Out> {
self.poll_event()
}
fn finish(&mut self) -> core::result::Result<(), Self::Error> {
self.finish();
Ok(())
}
fn next_deadline(&self) -> Option<Timestamp> {
None
}
fn on_deadline(&mut self, _now: Timestamp) {}
fn demand(&self) -> Demand {
if self.unattributed_bytes.saturating_add(TS_MAX_PAYLOAD_BYTES) > MAX_UNATTRIBUTED_BYTES {
Demand::saturated()
} else {
Demand::new(TS_PACKET_SIZE)
}
}
}
#[derive(Debug, Default, Clone)]
pub struct TsDemux<'a> {
_marker: PhantomData<&'a [u8]>,
}
impl<'a> TsDemux<'a> {
pub fn new() -> Self {
Self {
_marker: PhantomData,
}
}
pub fn demux(&mut self, input: &'a [u8]) -> Result<Media> {
let mut demux = StreamingTsDemux::new();
demux.feed(input);
demux.finish();
let mut tracks: Vec<Track> = Vec::new();
let mut index_by_id: BTreeMap<u32, usize> = BTreeMap::new();
let mut pcr: Vec<PcrSample> = Vec::new();
while let Some(event) = demux.poll_event() {
match event {
DemuxEvent::TrackAdded(spec) => {
index_by_id.insert(spec.track_id, tracks.len());
tracks.push(Track::new(spec, Vec::new()));
}
DemuxEvent::TrackUpdated(spec) => {
if let Some(&i) = index_by_id.get(&spec.track_id) {
tracks[i].spec = spec;
}
}
DemuxEvent::TrackRemoved { .. } => {}
DemuxEvent::TrackAbandoned { .. } => {}
DemuxEvent::Sample { track_id, sample } => {
if let Some(&i) = index_by_id.get(&track_id) {
let track = &mut tracks[i];
if track.samples.is_empty() {
if let Some(dts) = sample.dts {
track.start_decode_time = dts as u64;
}
}
track.samples.push(sample);
}
}
DemuxEvent::ClockReference {
ticks,
discontinuous,
provenance,
..
} => pcr.push(PcrSample {
pcr_27mhz: ticks,
pid: provenance.pid.unwrap_or(0),
packet_index: provenance.packet_index.unwrap_or(0),
discontinuity: discontinuous,
}),
DemuxEvent::Discontinuity { .. } => {}
DemuxEvent::TracksResolved { .. } => {}
}
}
Ok(Media::new(tracks, VIDEO_TIMESCALE).with_pcr(pcr))
}
}
impl<'a> Unpackage for TsDemux<'a> {
type Input = &'a [u8];
type Media = Media;
type Error = Error;
fn unpackage(&mut self, input: &'a [u8]) -> Result<Media> {
self.demux(input)
}
}
#[cfg(test)]
mod tests {
use super::*;
use broadcast_common::Parse;
#[test]
fn data_carriage_classifies_every_known_stream_type() {
assert_eq!(data_carriage(0x06), DataCarriage::Pes, "PES private data");
assert_eq!(data_carriage(0x15), DataCarriage::Pes, "metadata in PES");
assert_eq!(data_carriage(0x7F), DataCarriage::Pes, "unrecognised → PES");
for &st in &[
STREAM_TYPE_PRIVATE_SECTIONS,
STREAM_TYPE_DSMCC_TYPE_A,
STREAM_TYPE_DSMCC_TYPE_B,
STREAM_TYPE_DSMCC_TYPE_C,
STREAM_TYPE_DSMCC_TYPE_D,
STREAM_TYPE_DSMCC_SYNC_DOWNLOAD,
STREAM_TYPE_SCTE35,
] {
assert_eq!(
data_carriage(st),
DataCarriage::Sections,
"stream_type {st:#04X} must be section-carried"
);
}
}
#[test]
fn from_stream_type_unknown_becomes_opaque_data() {
assert_eq!(Codec::from_stream_type(STREAM_TYPE_AVC), Codec::H264);
assert_eq!(Codec::from_stream_type(0x7F), Codec::Data(0x7F));
assert_eq!(
Codec::from_stream_type(STREAM_TYPE_SCTE35),
Codec::Data(STREAM_TYPE_SCTE35)
);
}
const PACKET_PAYLOAD_LEN: usize = TS_PACKET_SIZE - 4;
fn payload_only_packet(pid: u16, cc: u8) -> [u8; TS_PACKET_SIZE] {
let mut p = [0xFFu8; TS_PACKET_SIZE];
p[0] = 0x47; p[1] = ((pid >> 8) as u8) & PID_HI_MASK; p[2] = (pid & 0xFF) as u8; p[3] = 0x10 | (cc & 0x0F); p
}
#[test]
fn unattributed_buffer_is_bounded_for_never_claimed_pid() {
let target_bytes = MAX_UNATTRIBUTED_BYTES * 3;
let packet_count = target_bytes / PACKET_PAYLOAD_LEN + 1;
let unclaimed_pid: u16 = 0x0123;
let mut demux = StreamingTsDemux::new();
for i in 0..packet_count {
demux.feed(&payload_only_packet(unclaimed_pid, i as u8));
}
assert!(
demux.unattributed_bytes <= MAX_UNATTRIBUTED_BYTES,
"unattributed_bytes {} exceeded cap {}",
demux.unattributed_bytes,
MAX_UNATTRIBUTED_BYTES
);
assert!(
demux.unattributed_bytes > 0,
"expected the never-claimed PID's payload to be buffered"
);
let actual: usize = demux
.unattributed
.values()
.flat_map(|q| q.iter())
.map(|(_, payload)| payload.len())
.sum();
assert_eq!(
actual, demux.unattributed_bytes,
"unattributed_bytes drifted from the real retained size"
);
}
#[test]
fn runaway_pes_without_payload_unit_start_is_bounded_not_unbounded() {
use crate::TsMux;
use crate::media::{Media, Track};
use crate::pipeline::{CodecConfig, Sample, TrackSpec};
use crate::rtp_sdp::avc_config_from_sprop;
use broadcast_common::Package;
const ES_PID: u16 = 0x0100;
let avc = avc_config_from_sprop("Z0IAKeKQFAe2AtwEBAaQeJEV,aM48gA==").unwrap();
let spec = TrackSpec::new(
1,
VIDEO_TIMESCALE,
CodecConfig::Avc {
config: avc,
width: 0,
height: 0,
},
);
let frame_dur = VIDEO_TIMESCALE / 30;
let samples: Vec<Sample> = (0..3u32)
.map(|i| {
let nal = [0x65u8, 0xAA, i as u8];
let mut data = (nal.len() as u32).to_be_bytes().to_vec();
data.extend_from_slice(&nal);
let dts = i64::from(i) * i64::from(frame_dur);
Sample::new(data, Some(dts), Some(dts), Some(frame_dur), i == 0)
})
.collect();
let track = Track::new(spec, samples);
let media = Media::new(vec![track], VIDEO_TIMESCALE);
let ts_bytes = TsMux::default().package(&media).expect("mux to TS");
let mut demux = StreamingTsDemux::new();
demux.feed(&ts_bytes);
while demux.poll_event().is_some() {}
assert!(
demux.streams.contains_key(&ES_PID),
"expected the muxed video ES at PID {ES_PID:#06X} — TsMux's ES_PID_BASE convention"
);
let mut cc: u8 = 0;
let mut hit_cap = false;
for _ in 0..30_000u32 {
let mut pkt = [0xFFu8; TS_PACKET_SIZE];
pkt[0] = 0x47;
pkt[1] = ((ES_PID >> 8) as u8) & PID_HI_MASK; pkt[2] = (ES_PID & 0xFF) as u8;
pkt[3] = 0x10 | (cc & 0x0F);
cc = cc.wrapping_add(1);
demux.feed(&pkt);
if matches!(
demux.poll_event(),
Some(DemuxEvent::Discontinuity { provenance, .. }) if provenance.pid == Some(ES_PID)
) {
hit_cap = true;
break;
}
}
assert!(
hit_cap,
"expected MAX_PES_BUFFER_BYTES to trip well within 30000 continuation \
packets (never grow unbounded)"
);
assert_eq!(
demux.streams.get(&ES_PID).unwrap().pes_bytes,
0,
"pes_bytes must reset to 0 once the cap trips"
);
const PUSI_BIT: u8 = 0x40; let mut pkt = [0xFFu8; TS_PACKET_SIZE];
pkt[0] = 0x47;
pkt[1] = PUSI_BIT | (((ES_PID >> 8) as u8) & PID_HI_MASK);
pkt[2] = (ES_PID & 0xFF) as u8;
pkt[3] = 0x10 | (cc & 0x0F);
pkt[4] = 0; demux.feed(&pkt);
assert_eq!(
demux.streams.get(&ES_PID).unwrap().pes_bytes,
PACKET_PAYLOAD_LEN,
"a fresh payload_unit_start must be accepted and start a new count"
);
}
#[test]
fn probe_backlog_is_bounded_for_both_the_never_resolving_pid_and_its_collateral_pid() {
use crate::TsMux;
use crate::media::{Media, Track};
use crate::pipeline::{CodecConfig, DataCarriage, Sample, TrackSpec};
use crate::rtp_sdp::avc_config_from_sprop;
use broadcast_common::Package;
const SAMPLE_BYTES: usize = 4096;
const SAMPLE_COUNT: u32 = 1200;
let frame_dur = VIDEO_TIMESCALE / 30;
let avc = avc_config_from_sprop("Z0IAKeKQFAe2AtwEBAaQeJEV,aM48gA==").unwrap();
let video_spec = TrackSpec::new(
1,
VIDEO_TIMESCALE,
CodecConfig::Avc {
config: avc,
width: 0,
height: 0,
},
);
let video_samples: Vec<Sample> = (0..SAMPLE_COUNT)
.map(|i| {
let mut nal = alloc::vec![0x41u8]; nal.resize(SAMPLE_BYTES, 0xAA);
let mut data = (nal.len() as u32).to_be_bytes().to_vec();
data.extend_from_slice(&nal);
let dts = i64::from(i) * i64::from(frame_dur);
Sample::new(data, Some(dts), Some(dts), Some(frame_dur), false)
})
.collect();
let video_track = Track::new(video_spec, video_samples);
let data_spec = TrackSpec::new(
2,
VIDEO_TIMESCALE,
CodecConfig::Data {
stream_type: 0x7F,
descriptors: Vec::new(),
carriage: DataCarriage::Pes,
},
);
let data_samples: Vec<Sample> = (0..SAMPLE_COUNT)
.map(|i| {
let payload = alloc::vec![0xBBu8; SAMPLE_BYTES];
let dts = i64::from(i) * i64::from(frame_dur);
Sample::new(payload, Some(dts), Some(dts), Some(frame_dur), true)
})
.collect();
let data_track = Track::new(data_spec, data_samples);
let media = Media::new(vec![video_track, data_track], VIDEO_TIMESCALE);
let ts_bytes = TsMux::default().package(&media).expect("mux to TS");
let mut demux = StreamingTsDemux::new();
demux.feed(&ts_bytes);
let mut abandoned_pids: Vec<u16> = Vec::new();
while let Some(ev) = demux.poll_event() {
if let DemuxEvent::TrackAbandoned {
reason: AbandonReason::BudgetExceeded,
provenance,
..
} = ev
{
if let Some(pid) = provenance.pid {
abandoned_pids.push(pid);
}
}
}
assert!(
!abandoned_pids.is_empty(),
"expected at least one TrackAbandoned{{BudgetExceeded}} from an abandoned probe backlog \
(issue #774: this replaced the mis-typed Discontinuity this path used to emit)"
);
let mut demux2 = StreamingTsDemux::new();
demux2.feed(&ts_bytes);
let saw_discontinuity_for_abandoned_pid = std::iter::from_fn(|| demux2.poll_event())
.any(|ev| matches!(ev, DemuxEvent::Discontinuity { provenance, .. } if provenance.pid == Some(PID_A)));
assert!(
!saw_discontinuity_for_abandoned_pid,
"a probe-backlog-budget abandonment (issue #774) must be a TrackAbandoned, \
never a Discontinuity"
);
for (&pid, stream) in demux.streams.iter() {
assert!(
stream.backlog_bytes <= MAX_PROBE_BACKLOG_BYTES,
"PID {pid:#06X} backlog_bytes {} exceeded cap {MAX_PROBE_BACKLOG_BYTES}",
stream.backlog_bytes
);
}
const PID_A: u16 = 0x0100; let abandoned = matches!(
demux.streams.get(&PID_A).and_then(|s| s.track.as_ref()),
Some(TrackState::Abandoned)
);
assert!(abandoned, "PID A must be Abandoned, never Live");
const PID_B: u16 = 0x0101;
let pid_b_resolved = matches!(
demux.streams.get(&PID_B).and_then(|s| s.track.as_ref()),
Some(TrackState::Live(_)) | Some(TrackState::Abandoned)
);
assert!(
pid_b_resolved,
"PID B must reach a final disposition (Live or Abandoned), not stay wedged"
);
}
#[test]
fn track_abandoned_config_unrecoverable_fires_at_finish() {
use crate::TsMux;
use crate::media::{Media, Track};
use crate::pipeline::{CodecConfig, Sample, TrackSpec};
use crate::rtp_sdp::avc_config_from_sprop;
use broadcast_common::Package;
const PID_A: u16 = 0x0100; let frame_dur = VIDEO_TIMESCALE / 30;
let avc = avc_config_from_sprop("Z0IAKeKQFAe2AtwEBAaQeJEV,aM48gA==").unwrap();
let video_spec = TrackSpec::new(
1,
VIDEO_TIMESCALE,
CodecConfig::Avc {
config: avc,
width: 0,
height: 0,
},
);
let video_samples: Vec<Sample> = (0..5u32)
.map(|i| {
let nal = alloc::vec![0x41u8, 0xAA, 0xBB]; let mut data = (nal.len() as u32).to_be_bytes().to_vec();
data.extend_from_slice(&nal);
let dts = i64::from(i) * i64::from(frame_dur);
Sample::new(data, Some(dts), Some(dts), Some(frame_dur), false)
})
.collect();
let video_track = Track::new(video_spec, video_samples);
let media = Media::new(vec![video_track], VIDEO_TIMESCALE);
let ts_bytes = TsMux::default().package(&media).expect("mux to TS");
let mut demux = StreamingTsDemux::new();
demux.feed(&ts_bytes);
assert!(
!matches!(
demux.streams.get(&PID_A).and_then(|s| s.track.as_ref()),
Some(TrackState::Abandoned)
),
"sanity: PID A must still be Probing before finish() — the byte cap must not \
have tripped (this test is about the end-of-input path, not the budget one)"
);
while demux.poll_event().is_some() {}
demux.finish();
let mut saw_config_unrecoverable = false;
while let Some(ev) = demux.poll_event() {
if let DemuxEvent::TrackAbandoned {
track_id,
reason: AbandonReason::ConfigUnrecoverable,
provenance,
} = ev
{
assert_eq!(
track_id, None,
"a track abandoned before ever resolving has no track_id to report"
);
assert_eq!(provenance.pid, Some(PID_A));
saw_config_unrecoverable = true;
}
}
assert!(
saw_config_unrecoverable,
"expected TrackAbandoned{{ConfigUnrecoverable}} once finish() concludes PID A's \
config will never resolve"
);
}
#[test]
fn rescale_to_track_preserves_a_negative_anchor_for_audio_as_it_does_for_video() {
const SAMPLE_RATE: u32 = 48_000;
const NEGATIVE_90K: i128 = -90_000;
assert_eq!(
rescale_to_track(NEGATIVE_90K, VIDEO_TIMESCALE),
-90_000,
"the 90 kHz identity path already preserved this"
);
assert_eq!(
rescale_to_track(NEGATIVE_90K, SAMPLE_RATE),
-48_000,
"the audio rescale must preserve it too — one second before zero is \
-48000 ticks at 48 kHz, not 0"
);
assert_eq!(rescale_to_track(-1, SAMPLE_RATE), -1);
assert_eq!(rescale_to_track(1, SAMPLE_RATE), 0);
}
fn aac_44100_access_unit() -> Vec<u8> {
const ADTS_PROFILE_AAC_LC: u8 = 1;
const SFI_44100: u8 = 4;
const CHANNELS_STEREO: u8 = 2;
const PAYLOAD_BYTES: usize = 128;
let frame_len = (ADTS_HEADER_SIZE + PAYLOAD_BYTES) as u16;
let header = crate::aac_asc::build_adts_header(
ADTS_PROFILE_AAC_LC,
SFI_44100,
CHANNELS_STEREO,
frame_len,
);
let mut au = header.to_vec();
au.resize(ADTS_HEADER_SIZE + PAYLOAD_BYTES, 0x21);
au
}
#[test]
fn constant_increment_44100_aac_emits_no_timeline_reanchor() {
const PES_INCREMENT_TICKS: i128 = 2090;
const SAMPLE_RATE: u32 = 44_100;
const FRAMES: i128 = 500;
let au = aac_44100_access_unit();
let mut anchor = AudioAnchor::default();
let mut events: VecDeque<DemuxEvent> = VecDeque::new();
for n in 0..FRAMES {
let ts = n * PES_INCREMENT_TICKS;
emit_audio_au(
&AudioKind::Aac,
SAMPLE_RATE,
&mut anchor,
&au,
ts,
ts,
1,
&mut events,
);
}
assert!(
events
.iter()
.any(|e| matches!(e, DemuxEvent::Sample { .. })),
"sanity: the synthesised ADTS frames must actually split into samples"
);
let reanchors = events
.iter()
.filter(|e| {
matches!(
e,
DemuxEvent::Discontinuity {
kind: DiscontinuityKind::TimelineReanchored,
..
}
)
})
.count();
assert_eq!(
reanchors, 0,
"a constant-increment 44.1 kHz stream is continuous — the \
rounding drift a real muxer accrues by construction must never be \
reported as a discontinuity"
);
}
#[test]
fn a_real_timeline_gap_emits_exactly_one_reanchor() {
const PES_INCREMENT_TICKS: i128 = 2090;
const SAMPLE_RATE: u32 = 44_100;
const FRAMES_BEFORE: i128 = 50;
const FRAMES_AFTER: i128 = 50;
const GAP_TICKS: i128 = 180_000;
let au = aac_44100_access_unit();
let mut anchor = AudioAnchor::default();
let mut events: VecDeque<DemuxEvent> = VecDeque::new();
let emit = |ts: i128, anchor: &mut AudioAnchor, events: &mut VecDeque<DemuxEvent>| {
emit_audio_au(&AudioKind::Aac, SAMPLE_RATE, anchor, &au, ts, ts, 1, events);
};
for n in 0..FRAMES_BEFORE {
emit(n * PES_INCREMENT_TICKS, &mut anchor, &mut events);
}
let resume = FRAMES_BEFORE * PES_INCREMENT_TICKS + GAP_TICKS;
for n in 0..FRAMES_AFTER {
emit(resume + n * PES_INCREMENT_TICKS, &mut anchor, &mut events);
}
let reanchors = events
.iter()
.filter(|e| {
matches!(
e,
DemuxEvent::Discontinuity {
kind: DiscontinuityKind::TimelineReanchored,
..
}
)
})
.count();
assert_eq!(
reanchors, 1,
"one genuine gap must produce exactly one TimelineReanchored — not \
zero (the anchor silently absorbing a real splice) and not one \
per following access unit"
);
}
#[test]
fn adts_resyncs_across_pes_boundaries() {
let mut path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"));
path.push("..");
path.push("fixtures");
path.push("ts");
path.push("h264_aac.ts");
let ts_bytes =
std::fs::read(&path).unwrap_or_else(|e| panic!("read fixture {}: {e}", path.display()));
let mut demux = TsDemux::new();
let media = demux.demux(&ts_bytes).expect("demux h264_aac.ts");
let audio = media
.tracks
.iter()
.find(|t| matches!(t.config(), CodecConfig::Aac { .. }))
.expect("h264_aac.ts has an AAC track");
let esds = match audio.config() {
CodecConfig::Aac { esds, .. } => esds,
_ => unreachable!(),
};
let dsi = esds
.es_descriptor
.decoder_config
.as_ref()
.and_then(|dc| dc.decoder_specific_info.as_ref())
.expect("AAC esds carries a DecoderSpecificInfo");
let asc = AudioSpecificConfig::parse(&dsi.data).expect("parse real AudioSpecificConfig");
let mut es: Vec<u8> = Vec::new();
for s in &audio.samples {
let frame_len = (ADTS_HEADER_SIZE + s.data.len()) as u16;
let hdr = asc
.to_adts_header(frame_len)
.expect("build real ADTS header");
es.extend_from_slice(&hdr);
es.extend_from_slice(&s.data);
}
assert!(
audio.samples.len() >= 10,
"fixture too small to be a meaningful resync test"
);
const CHUNK: usize = 2000;
let mut recovered = 0usize;
let mut off = 0usize;
while off < es.len() {
let end = (off + CHUNK).min(es.len());
recovered += split_adts_frames(&es[off..end]).len();
off = end;
}
assert!(
recovered * 2 >= audio.samples.len(),
"resync must recover most real AAC frames across misaligned \
chunks (got {recovered} of {} real frames)",
audio.samples.len()
);
}
#[test]
fn refine_with_descriptors_recognises_every_dolby_dts_tag() {
let ac3_desc = [DESC_TAG_AC3, 1, 0x40];
assert_eq!(
Codec::Data(STREAM_TYPE_PES_PRIVATE)
.refine_with_descriptors(STREAM_TYPE_PES_PRIVATE, &ac3_desc),
Codec::Ac3
);
let dts_desc = [DESC_TAG_DTS, 1, 0x00];
assert_eq!(
Codec::Data(STREAM_TYPE_METADATA_PES)
.refine_with_descriptors(STREAM_TYPE_METADATA_PES, &dts_desc),
Codec::Dts,
"0x15 (metadata in PES) is also descriptor-disambiguated"
);
let eac3_desc = [DESC_TAG_ENHANCED_AC3, 1, 0x00];
assert_eq!(
Codec::Data(STREAM_TYPE_PES_PRIVATE)
.refine_with_descriptors(STREAM_TYPE_PES_PRIVATE, &eac3_desc),
Codec::Eac3
);
}
#[test]
fn refine_with_descriptors_leaves_non_audio_0x06_as_data() {
const DESC_TAG_SUBTITLING: u8 = 0x59;
let subtitle_desc = [DESC_TAG_SUBTITLING, 3, 0x65, 0x6E, 0x67];
assert_eq!(
Codec::Data(STREAM_TYPE_PES_PRIVATE)
.refine_with_descriptors(STREAM_TYPE_PES_PRIVATE, &subtitle_desc),
Codec::Data(STREAM_TYPE_PES_PRIVATE)
);
}
#[test]
fn refine_with_descriptors_ignores_other_stream_types() {
let eac3_desc = [DESC_TAG_ENHANCED_AC3, 1, 0x00];
assert_eq!(
Codec::H264.refine_with_descriptors(STREAM_TYPE_AVC, &eac3_desc),
Codec::H264
);
}
#[test]
fn codec_change_on_shared_pid_does_not_tear_down_other_program_track() {
const PMT_A: u16 = 0x1000;
const PMT_B: u16 = 0x1001;
const SHARED_PID: u16 = 0x0050;
const STREAM_TYPE: u8 = 0x7F;
let mut demux = StreamingTsDemux::new();
demux.apply_pmt_diff(
PMT_A,
&BTreeSet::new(),
alloc::vec![(SHARED_PID, Codec::Data(STREAM_TYPE), Vec::new())],
);
demux.apply_pmt_diff(
PMT_B,
&BTreeSet::new(),
alloc::vec![(SHARED_PID, Codec::Data(STREAM_TYPE), Vec::new())],
);
assert_eq!(
demux.es_declarers.get(&SHARED_PID).map(|d| d.len()),
Some(2),
"both PMTs must be recorded as declarers of the shared PID"
);
let carriage = data_carriage(STREAM_TYPE);
assert_eq!(carriage, DataCarriage::Pes);
demux.streams.get_mut(&SHARED_PID).unwrap().track = Some(TrackState::Parked {
config: CodecConfig::Data {
stream_type: STREAM_TYPE,
descriptors: Vec::new(),
carriage,
},
timescale: VIDEO_TIMESCALE,
kind: LiveKind::Data {
pending: None,
last_duration: 0,
},
backlog: Vec::new(),
});
demux.try_promote_ready();
let track_id_before = match demux.streams.get(&SHARED_PID).unwrap().track.as_ref() {
Some(TrackState::Live(live)) => live.track_id,
_ => panic!("expected the shared PID to be Live after promotion"),
};
while demux.poll_event().is_some() {}
let mut pmt_a_applied = BTreeSet::new();
pmt_a_applied.insert(SHARED_PID);
demux.apply_pmt_diff(
PMT_A,
&pmt_a_applied,
alloc::vec![(SHARED_PID, Codec::Ac3, Vec::new())],
);
match demux
.streams
.get(&SHARED_PID)
.and_then(|s| s.track.as_ref())
{
Some(TrackState::Live(live)) => assert_eq!(
live.track_id, track_id_before,
"shared track must keep its original track_id"
),
_ => {
panic!("expected the shared track to survive PMT A's reclassification, still Live")
}
}
assert!(
!demux
.events
.iter()
.any(|ev| matches!(ev, DemuxEvent::TrackRemoved { .. })),
"PMT A's reclassification must not remove a track PMT B still declares"
);
assert_eq!(
demux.streams.get(&SHARED_PID).unwrap().codec,
Codec::Data(STREAM_TYPE),
"the existing classification wins while another declarer disagrees"
);
let mut pmt_b_applied = BTreeSet::new();
pmt_b_applied.insert(SHARED_PID);
demux.apply_pmt_diff(PMT_B, &pmt_b_applied, Vec::new());
assert_eq!(
demux.es_declarers.get(&SHARED_PID).map(|d| d.len()),
Some(1),
"PMT A must be the sole remaining declarer"
);
demux.apply_pmt_diff(
PMT_A,
&pmt_a_applied,
alloc::vec![(SHARED_PID, Codec::Aac, Vec::new())],
);
assert!(
demux
.events
.iter()
.any(|ev| matches!(ev, DemuxEvent::TrackRemoved { .. })),
"once PMT A is the last declarer, its reclassification must actually tear down \
the old track"
);
assert_eq!(demux.streams.get(&SHARED_PID).unwrap().codec, Codec::Aac);
}
#[test]
fn codec_change_reregistration_preserves_declaration_order_slot() {
const PMT: u16 = 0x1000;
const PID_X: u16 = 0x0050;
const PID_Y: u16 = 0x0051;
let mut demux = StreamingTsDemux::new();
demux.apply_pmt_diff(
PMT,
&BTreeSet::new(),
alloc::vec![
(PID_X, Codec::Data(0x06), Vec::new()),
(PID_Y, Codec::Data(0x07), Vec::new()),
],
);
assert_eq!(demux.data_order, alloc::vec![PID_X, PID_Y]);
let mut old_applied = BTreeSet::new();
old_applied.insert(PID_X);
old_applied.insert(PID_Y);
demux.apply_pmt_diff(
PMT,
&old_applied,
alloc::vec![
(PID_X, Codec::Data(0x08), Vec::new()),
(PID_Y, Codec::Data(0x07), Vec::new()),
],
);
assert_eq!(
demux.data_order,
alloc::vec![PID_X, PID_Y],
"PID X must keep its original (first) declaration-order slot, not move to the back"
);
assert_eq!(demux.streams.get(&PID_X).unwrap().codec, Codec::Data(0x08));
}
}