#![cfg_attr(not(feature = "std"), no_std)]
#![cfg_attr(docsrs, feature(doc_cfg))]
#![doc = "\n# Examples\n"]
#![doc = "Two runnable examples ship with this crate (`cargo run -p dvb-conformance --example <name>`).\n"]
#![doc = "\n## `monitor_stream`\n\n```rust,ignore"]
#![doc = include_str!("../examples/monitor_stream.rs")]
#![doc = "```\n\n## `priority_breakdown`\n\n```rust,ignore"]
#![doc = include_str!("../examples/priority_breakdown.rs")]
#![doc = "```"]
extern crate alloc;
use alloc::collections::BTreeMap;
use alloc::format;
use alloc::string::String;
use alloc::vec::Vec;
use core::time::Duration;
use dvb_common::Parse;
use dvb_si::tables::pat::{PatSection, TABLE_ID as PAT_TABLE_ID};
use dvb_si::tables::pmt::PmtSection;
use mpeg_ts::section::Section;
use mpeg_ts::ts::{SectionReassembler, TsPacket};
const PID_PAT: u16 = 0x0000;
const PID_CAT: u16 = 0x0001;
const PID_NIT: u16 = 0x0010;
const PID_SDT_BAT: u16 = 0x0011;
const PID_EIT: u16 = 0x0012;
const PID_TDT_TOT: u16 = 0x0014;
const PID_NULL: u16 = 0x1FFF;
const SYNC_BYTE: u8 = 0x47;
const SI_PIDS: [u16; 6] = [PID_PAT, PID_CAT, PID_NIT, PID_SDT_BAT, PID_EIT, PID_TDT_TOT];
const DEFAULT_PAT_MAX_INTERVAL_MS: u64 = 500;
const DEFAULT_PMT_MAX_INTERVAL_MS: u64 = 500;
const DEFAULT_PID_ERROR_PERIOD_SECS: u64 = 5;
const DEFAULT_SYNC_ACQUIRE_PACKETS: u8 = 5;
const DEFAULT_SYNC_LOSS_PACKETS: u8 = 2;
const DEFAULT_PCR_REPETITION_LIMIT_MS: u64 = 100;
const DEFAULT_PCR_DISCONTINUITY_LIMIT_MS: u64 = 100;
const DEFAULT_PTS_REPETITION_LIMIT_MS: u64 = 700;
const DEFAULT_SI_NIT_INTERVAL_SECS: u64 = 10;
const DEFAULT_SI_SDT_INTERVAL_SECS: u64 = 2;
const DEFAULT_SI_EIT_PF_INTERVAL_SECS: u64 = 2;
const DEFAULT_SI_TDT_INTERVAL_SECS: u64 = 30;
const PCR_MODULUS_27MHZ: u64 = (1u64 << 33) * 300;
const CLOCK_27MHZ: u64 = 27_000_000;
const PES_PREFIX_0: u8 = 0x00;
const PES_PREFIX_1: u8 = 0x00;
const PES_PREFIX_2: u8 = 0x01;
const PES_FLAGS_OFFSET: usize = 6;
const PES_PTS_DTS_FLAGS_MASK: u8 = 0b1100_0000;
const PES_PTS_PRESENT: u8 = 0b1000_0000;
const CAT_TABLE_ID: u8 = dvb_si::table_id::TableId::Cat as u8;
const NIT_ACTUAL_TABLE_ID: u8 = dvb_si::table_id::TableId::NetworkInformationActual as u8;
const SDT_ACTUAL_TABLE_ID: u8 = dvb_si::table_id::TableId::ServiceDescriptionActual as u8;
const EIT_PF_ACTUAL_TABLE_ID: u8 = dvb_si::table_id::TableId::EventInformationPfActual as u8;
const TDT_TABLE_ID: u8 = dvb_si::table_id::TableId::TimeAndDate as u8;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize))]
#[non_exhaustive]
pub enum Priority {
First,
Second,
Third,
}
impl Priority {
#[must_use]
pub fn name(&self) -> &'static str {
match self {
Self::First => "first priority",
Self::Second => "second priority",
Self::Third => "third priority",
}
}
}
dvb_common::impl_spec_display!(Priority);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[cfg_attr(feature = "serde", derive(serde::Serialize))]
#[non_exhaustive]
pub enum Indicator {
TsSyncLoss,
SyncByteError,
PatError2,
ContinuityCountError,
PmtError2,
PidError,
TransportError,
CrcError,
PcrRepetitionError,
PcrDiscontinuityError,
PtsError,
CatError,
SiRepetitionError,
}
impl Indicator {
#[must_use]
pub fn priority(self) -> Priority {
match self {
Self::TsSyncLoss
| Self::SyncByteError
| Self::PatError2
| Self::ContinuityCountError
| Self::PmtError2
| Self::PidError => Priority::First,
Self::TransportError
| Self::CrcError
| Self::PcrRepetitionError
| Self::PcrDiscontinuityError
| Self::PtsError
| Self::CatError => Priority::Second,
Self::SiRepetitionError => Priority::Third,
}
}
#[must_use]
pub fn name(self) -> &'static str {
match self {
Self::TsSyncLoss => "TS_sync_loss",
Self::SyncByteError => "Sync_byte_error",
Self::PatError2 => "PAT_error_2",
Self::ContinuityCountError => "Continuity_count_error",
Self::PmtError2 => "PMT_error_2",
Self::PidError => "PID_error",
Self::TransportError => "Transport_error",
Self::CrcError => "CRC_error",
Self::PcrRepetitionError => "PCR_repetition_error",
Self::PcrDiscontinuityError => "PCR_discontinuity_indicator_error",
Self::PtsError => "PTS_error",
Self::CatError => "CAT_error",
Self::SiRepetitionError => "SI_repetition_error",
}
}
#[must_use]
pub fn clause(self) -> &'static str {
match self {
Self::TsSyncLoss => "TR 101 290 v1.4.1 Table 5.0a indicator 1.1",
Self::SyncByteError => "TR 101 290 v1.4.1 Table 5.0a indicator 1.2",
Self::PatError2 => "TR 101 290 v1.4.1 Table 5.0a indicator 1.3.a",
Self::ContinuityCountError => "TR 101 290 v1.4.1 Table 5.0a indicator 1.4",
Self::PmtError2 => "TR 101 290 v1.4.1 Table 5.0a indicator 1.5.a",
Self::PidError => "TR 101 290 v1.4.1 Table 5.0a indicator 1.6",
Self::TransportError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.1",
Self::CrcError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.2",
Self::PcrRepetitionError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.3a",
Self::PcrDiscontinuityError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.3b",
Self::PtsError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.5",
Self::CatError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.6",
Self::SiRepetitionError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.2",
}
}
}
dvb_common::impl_spec_display!(Indicator);
#[derive(Debug, Clone, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize))]
#[non_exhaustive]
pub struct ConformanceEvent {
pub indicator: Indicator,
pub priority: Priority,
pub pid: Option<u16>,
pub at: Duration,
pub detail: String,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize))]
#[non_exhaustive]
pub struct Stats {
pub packets: u64,
pub events: u64,
pub in_sync: bool,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct Config {
pub pat_max_interval: Duration,
pub pmt_max_interval: Duration,
pub pid_error_period: Duration,
pub sync_acquire_packets: u8,
pub sync_loss_packets: u8,
pub pcr_repetition_limit: Duration,
pub pcr_discontinuity_limit: Duration,
pub pts_repetition_limit: Duration,
pub si_nit_interval: Duration,
pub si_sdt_interval: Duration,
pub si_eit_pf_interval: Duration,
pub si_tdt_interval: Duration,
}
impl Default for Config {
fn default() -> Self {
Self {
pat_max_interval: Duration::from_millis(DEFAULT_PAT_MAX_INTERVAL_MS),
pmt_max_interval: Duration::from_millis(DEFAULT_PMT_MAX_INTERVAL_MS),
pid_error_period: Duration::from_secs(DEFAULT_PID_ERROR_PERIOD_SECS),
sync_acquire_packets: DEFAULT_SYNC_ACQUIRE_PACKETS,
sync_loss_packets: DEFAULT_SYNC_LOSS_PACKETS,
pcr_repetition_limit: Duration::from_millis(DEFAULT_PCR_REPETITION_LIMIT_MS),
pcr_discontinuity_limit: Duration::from_millis(DEFAULT_PCR_DISCONTINUITY_LIMIT_MS),
pts_repetition_limit: Duration::from_millis(DEFAULT_PTS_REPETITION_LIMIT_MS),
si_nit_interval: Duration::from_secs(DEFAULT_SI_NIT_INTERVAL_SECS),
si_sdt_interval: Duration::from_secs(DEFAULT_SI_SDT_INTERVAL_SECS),
si_eit_pf_interval: Duration::from_secs(DEFAULT_SI_EIT_PF_INTERVAL_SECS),
si_tdt_interval: Duration::from_secs(DEFAULT_SI_TDT_INTERVAL_SECS),
}
}
}
struct CcState {
last_cc: u8,
had_payload: bool,
dup_used: bool,
initialised: bool,
}
struct PresenceTimer {
last_seen: Duration,
reported: bool,
}
struct PmtTracking {
timer: PresenceTimer,
reassembler: SectionReassembler,
}
struct EsTracking {
timer: PresenceTimer,
}
struct PcrState {
last_pcr_27mhz: u64,
last_pcr_time: Duration,
initialised: bool,
}
struct PtsState {
last_pts_time: Duration,
armed: bool,
}
struct SiReassembly {
reassembler: SectionReassembler,
}
struct SiRepetitionTimer {
last_seen: Duration,
reported: bool,
armed: bool,
}
pub struct ConformanceMonitor {
config: Config,
events: Vec<ConformanceEvent>,
stats: Stats,
in_sync: bool,
good_run: u8,
bad_run: u8,
cc_states: BTreeMap<u16, CcState>,
pat_reassembler: SectionReassembler,
pat_timer: PresenceTimer,
pmt_trackings: BTreeMap<u16, PmtTracking>,
es_trackings: BTreeMap<u16, EsTracking>,
si_reassemblies: BTreeMap<u16, SiReassembly>,
pcr_states: BTreeMap<u16, PcrState>,
pts_states: BTreeMap<u16, PtsState>,
cat_seen: bool,
scrambled_without_cat_reported: bool,
si_timers: BTreeMap<u8, SiRepetitionTimer>,
}
impl ConformanceMonitor {
pub fn new() -> Self {
Self::with_config(Config::default())
}
pub fn with_config(config: Config) -> Self {
let mut si_reassemblies = BTreeMap::new();
for &pid in &SI_PIDS {
si_reassemblies.insert(
pid,
SiReassembly {
reassembler: SectionReassembler::default(),
},
);
}
Self {
config,
events: Vec::new(),
stats: Stats {
packets: 0,
events: 0,
in_sync: false,
},
in_sync: false,
good_run: 0,
bad_run: 0,
cc_states: BTreeMap::new(),
pat_reassembler: SectionReassembler::default(),
pat_timer: PresenceTimer {
last_seen: Duration::ZERO,
reported: false,
},
pmt_trackings: BTreeMap::new(),
es_trackings: BTreeMap::new(),
si_reassemblies,
pcr_states: BTreeMap::new(),
pts_states: BTreeMap::new(),
cat_seen: false,
scrambled_without_cat_reported: false,
si_timers: BTreeMap::new(),
}
}
pub fn feed(&mut self, ts_packet: &[u8], t: Duration) -> &[ConformanceEvent] {
self.events.clear();
self.stats.packets += 1;
let sync_ok = !ts_packet.is_empty() && ts_packet[0] == SYNC_BYTE;
if !sync_ok {
self.emit(Indicator::SyncByteError, None, t, "sync_byte != 0x47");
}
if sync_ok {
self.good_run = self.good_run.saturating_add(1);
self.bad_run = 0;
if !self.in_sync && self.good_run >= self.config.sync_acquire_packets {
self.in_sync = true;
}
} else {
self.bad_run = self.bad_run.saturating_add(1);
self.good_run = 0;
if self.in_sync && self.bad_run >= self.config.sync_loss_packets {
self.in_sync = false;
self.emit(
Indicator::TsSyncLoss,
None,
t,
"sync lost after hysteresis threshold",
);
}
}
if !self.in_sync {
return &self.events;
}
let packet = match TsPacket::parse(ts_packet) {
Ok(p) => p,
Err(_) => return &self.events,
};
let header = &packet.header;
let pid = header.pid;
if header.tei {
self.emit(
Indicator::TransportError,
Some(pid),
t,
format!("transport_error_indicator set on PID 0x{:04X}", pid),
);
}
if pid != PID_NULL {
self.check_cc(
pid,
header.continuity_counter,
header.has_payload,
t,
ts_packet,
);
}
if pid == PID_PAT && header.scrambling != 0 {
self.emit(
Indicator::PatError2,
Some(PID_PAT),
t,
format!(
"scrambling_control_field != 00 on PID 0x0000 (got {})",
header.scrambling
),
);
}
if self.pmt_trackings.contains_key(&pid) && header.scrambling != 0 {
self.emit(
Indicator::PmtError2,
Some(pid),
t,
format!(
"scrambling_control_field != 00 on program_map_PID 0x{:04X}",
pid
),
);
}
if header.scrambling != 0 && !self.cat_seen && !self.scrambled_without_cat_reported {
self.scrambled_without_cat_reported = true;
self.emit(
Indicator::CatError,
Some(pid),
t,
format!(
"scrambled packet on PID 0x{:04X} but no CAT seen on PID 0x0001",
pid
),
);
}
if pid == PID_PAT && header.has_payload {
if let Some(payload) = packet.payload {
self.pat_reassembler.feed(payload, header.pusi);
}
self.pat_timer.last_seen = t;
self.pat_timer.reported = false;
while let Some(section_bytes) = self.pat_reassembler.pop_section() {
self.check_crc_and_process_pat(§ion_bytes, pid, t);
}
}
if self.pmt_trackings.contains_key(&pid) && header.has_payload {
if let Some(payload) = packet.payload {
if let Some(tracking) = self.pmt_trackings.get_mut(&pid) {
tracking.reassembler.feed(payload, header.pusi);
}
}
let sections: Vec<_> = if let Some(tracking) = self.pmt_trackings.get_mut(&pid) {
tracking.timer.last_seen = t;
tracking.timer.reported = false;
core::iter::from_fn(|| tracking.reassembler.pop_section()).collect()
} else {
Vec::new()
};
for section_bytes in §ions {
self.check_crc_and_process_pmt(section_bytes, pid, t);
}
}
if pid != PID_PAT
&& !self.pmt_trackings.contains_key(&pid)
&& self.si_reassemblies.contains_key(&pid)
&& header.has_payload
{
if let Some(payload) = packet.payload {
if let Some(si_ra) = self.si_reassemblies.get_mut(&pid) {
si_ra.reassembler.feed(payload, header.pusi);
}
}
let sections: Vec<_> = if let Some(si_ra) = self.si_reassemblies.get_mut(&pid) {
core::iter::from_fn(|| si_ra.reassembler.pop_section()).collect()
} else {
Vec::new()
};
for section_bytes in §ions {
self.check_crc_for_si(section_bytes, pid, t);
self.check_cat_table_id(section_bytes, pid, t);
self.update_si_repetition(section_bytes, pid, t);
}
}
if let Some(tracking) = self.es_trackings.get_mut(&pid) {
tracking.timer.last_seen = t;
tracking.timer.reported = false;
}
if let Some(Ok(af)) = packet.adaptation_field() {
if let Some(pcr) = af.pcr {
self.check_pcr(pid, pcr.as_27mhz(), af.discontinuity_indicator, t);
}
}
if header.pusi
&& header.scrambling == 0
&& self.es_trackings.contains_key(&pid)
&& header.has_payload
{
if let Some(payload) = packet.payload {
self.check_pts(pid, payload, t);
}
}
self.check_presence_timeouts(t);
&self.events
}
pub fn stats(&self) -> Stats {
Stats {
in_sync: self.in_sync,
..self.stats
}
}
fn emit(
&mut self,
indicator: Indicator,
pid: Option<u16>,
at: Duration,
detail: impl Into<String>,
) {
let event = ConformanceEvent {
indicator,
priority: indicator.priority(),
pid,
at,
detail: detail.into(),
};
self.stats.events += 1;
self.events.push(event);
}
fn check_cc(&mut self, pid: u16, cc: u8, has_payload: bool, t: Duration, raw: &[u8]) {
let discontinuity = if raw.len() >= 5 {
let b3 = raw[3];
let has_adaptation = (b3 & 0x20) != 0;
if has_adaptation {
let af_len = raw[4] as usize;
if af_len > 0 && raw.len() > 5 {
(raw[5] & 0x80) != 0
} else {
false
}
} else {
false
}
} else {
false
};
let (expected, is_duplicate, should_emit_dup, should_emit_cc) = {
let state = self.cc_states.entry(pid).or_insert_with(|| CcState {
last_cc: cc,
had_payload: has_payload,
dup_used: false,
initialised: false,
});
if !state.initialised {
state.last_cc = cc;
state.had_payload = has_payload;
state.dup_used = false;
state.initialised = true;
return;
}
if discontinuity {
(0u8, false, false, false)
} else {
let is_duplicate = cc == state.last_cc && has_payload;
let mut should_emit_dup = false;
let mut should_emit_cc = false;
if is_duplicate {
if state.dup_used {
should_emit_dup = true;
}
} else {
state.dup_used = false;
let expected = if has_payload {
(state.last_cc.wrapping_add(1)) & 0x0F
} else {
state.last_cc
};
if cc != expected {
should_emit_cc = true;
}
}
(
if has_payload {
(state.last_cc.wrapping_add(1)) & 0x0F
} else {
state.last_cc
},
is_duplicate,
should_emit_dup,
should_emit_cc,
)
}
};
if should_emit_dup {
self.emit(
Indicator::ContinuityCountError,
Some(pid),
t,
format!(
"second consecutive duplicate on PID 0x{:04X} (cc={})",
pid, cc
),
);
}
if should_emit_cc {
self.emit(
Indicator::ContinuityCountError,
Some(pid),
t,
format!("expected cc={}, got {} on PID 0x{:04X}", expected, cc, pid),
);
}
let state = self.cc_states.get_mut(&pid).unwrap();
if discontinuity {
state.last_cc = cc;
state.had_payload = has_payload;
state.dup_used = false;
} else if is_duplicate {
state.dup_used = true;
} else {
state.dup_used = false;
state.last_cc = cc;
state.had_payload = has_payload;
}
}
fn check_crc_and_process_pat(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
self.check_crc_for_section(section_bytes, pid, t);
self.process_pat_section(section_bytes, t);
}
fn check_crc_and_process_pmt(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
self.check_crc_for_section(section_bytes, pid, t);
self.process_pmt_section(section_bytes, pid, t);
}
fn check_crc_for_section(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
if let Err(mpeg_ts::error::Error::CrcMismatch { .. }) = section.validate_crc(section_bytes)
{
self.emit(
Indicator::CrcError,
Some(pid),
t,
format!(
"CRC-32 mismatch on PID 0x{:04X} (table_id 0x{:02X})",
pid, section.table_id
),
);
}
}
fn check_crc_for_si(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
self.check_crc_for_section(section_bytes, pid, t);
}
fn check_cat_table_id(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
if pid != PID_CAT {
return;
}
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
if section.table_id == CAT_TABLE_ID {
self.cat_seen = true;
self.scrambled_without_cat_reported = false;
} else {
self.emit(
Indicator::CatError,
Some(PID_CAT),
t,
format!(
"section with table_id 0x{:02X} on PID 0x0001 (expected 0x01 for CAT)",
section.table_id
),
);
}
}
fn process_pat_section(&mut self, section_bytes: &[u8], t: Duration) {
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
if section.table_id != PAT_TABLE_ID {
self.emit(
Indicator::PatError2,
Some(PID_PAT),
t,
format!(
"section with table_id 0x{:02X} on PID 0x0000 (expected 0x00)",
section.table_id
),
);
return;
}
let pat = match PatSection::parse(section_bytes) {
Ok(p) => p,
Err(_) => return,
};
for entry in pat.programmes() {
let pmt_pid = entry.pid;
self.pmt_trackings
.entry(pmt_pid)
.or_insert_with(|| PmtTracking {
timer: PresenceTimer {
last_seen: t,
reported: false,
},
reassembler: SectionReassembler::default(),
});
}
}
fn process_pmt_section(&mut self, section_bytes: &[u8], _pid: u16, t: Duration) {
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
let pmt_table_id: u8 = dvb_si::tables::pmt::TABLE_ID;
if section.table_id != pmt_table_id {
return;
}
let pmt = match PmtSection::parse(section_bytes) {
Ok(p) => p,
Err(_) => return,
};
let mut new_es_pids: Vec<u16> = Vec::new();
if pmt.pcr_pid != PID_NULL && !self.es_trackings.contains_key(&pmt.pcr_pid) {
new_es_pids.push(pmt.pcr_pid);
}
for stream in &pmt.streams {
let es_pid = stream.elementary_pid;
if !self.es_trackings.contains_key(&es_pid) {
new_es_pids.push(es_pid);
}
}
for es_pid in new_es_pids {
self.es_trackings.insert(
es_pid,
EsTracking {
timer: PresenceTimer {
last_seen: t,
reported: false,
},
},
);
}
}
fn check_pcr(&mut self, pid: u16, pcr_27mhz: u64, discontinuity: bool, t: Duration) {
let state = self.pcr_states.entry(pid).or_insert_with(|| PcrState {
last_pcr_27mhz: 0,
last_pcr_time: Duration::ZERO,
initialised: false,
});
if !state.initialised {
state.last_pcr_27mhz = pcr_27mhz;
state.last_pcr_time = t;
state.initialised = true;
return;
}
let last_pcr_time = state.last_pcr_time;
let last_pcr_27mhz = state.last_pcr_27mhz;
let rep_interval = t.saturating_sub(last_pcr_time);
let should_emit_rep = rep_interval > self.config.pcr_repetition_limit;
let delta =
(pcr_27mhz.wrapping_add(PCR_MODULUS_27MHZ) - last_pcr_27mhz) % PCR_MODULUS_27MHZ;
let delta_ms = delta * 1000 / CLOCK_27MHZ;
let limit_ms = self.config.pcr_discontinuity_limit.as_millis() as u64;
let should_emit_disc = delta_ms > limit_ms && !discontinuity;
if should_emit_rep {
self.emit(
Indicator::PcrRepetitionError,
Some(pid),
t,
format!(
"PCR interval {} ms exceeds limit {} ms on PID 0x{:04X}",
rep_interval.as_millis(),
self.config.pcr_repetition_limit.as_millis(),
pid
),
);
}
if should_emit_disc {
self.emit(
Indicator::PcrDiscontinuityError,
Some(pid),
t,
format!(
"PCR delta {} ms exceeds limit {} ms on PID 0x{:04X} without discontinuity_indicator",
delta_ms, limit_ms, pid
),
);
}
let state = self.pcr_states.get_mut(&pid).unwrap();
state.last_pcr_27mhz = pcr_27mhz;
state.last_pcr_time = t;
}
fn check_pts(&mut self, pid: u16, payload: &[u8], t: Duration) {
if payload.len() < PES_FLAGS_OFFSET + 2 {
return;
}
if payload[0] != PES_PREFIX_0 || payload[1] != PES_PREFIX_1 || payload[2] != PES_PREFIX_2 {
return;
}
let flags_byte = payload[PES_FLAGS_OFFSET];
if (flags_byte >> 6) != 0b10 {
return;
}
let pts_dts_flags = payload[PES_FLAGS_OFFSET + 1] & PES_PTS_DTS_FLAGS_MASK;
let pts_present = (pts_dts_flags & PES_PTS_PRESENT) != 0;
if !pts_present {
return;
}
let state = self.pts_states.entry(pid).or_insert_with(|| PtsState {
last_pts_time: Duration::ZERO,
armed: false,
});
if !state.armed {
state.last_pts_time = t;
state.armed = true;
return;
}
let last_pts_time = state.last_pts_time;
let pts_interval = t.saturating_sub(last_pts_time);
let should_emit = pts_interval > self.config.pts_repetition_limit;
if should_emit {
self.emit(
Indicator::PtsError,
Some(pid),
t,
format!(
"PTS interval {} ms exceeds limit {} ms on PID 0x{:04X}",
pts_interval.as_millis(),
self.config.pts_repetition_limit.as_millis(),
pid
),
);
}
let state = self.pts_states.get_mut(&pid).unwrap();
state.last_pts_time = t;
}
fn update_si_repetition(&mut self, section_bytes: &[u8], _pid: u16, t: Duration) {
let table_id = match Section::parse(section_bytes) {
Ok(s) => s.table_id,
Err(_) => return,
};
let is_tracked = table_id == NIT_ACTUAL_TABLE_ID
|| table_id == SDT_ACTUAL_TABLE_ID
|| table_id == EIT_PF_ACTUAL_TABLE_ID
|| table_id == TDT_TABLE_ID;
if !is_tracked {
return;
}
let timer = self
.si_timers
.entry(table_id)
.or_insert_with(|| SiRepetitionTimer {
last_seen: Duration::ZERO,
reported: false,
armed: false,
});
timer.last_seen = t;
timer.reported = false;
timer.armed = true;
}
fn check_presence_timeouts(&mut self, t: Duration) {
if t.saturating_sub(self.pat_timer.last_seen) > self.config.pat_max_interval
&& !self.pat_timer.reported
{
self.pat_timer.reported = true;
self.emit(
Indicator::PatError2,
Some(PID_PAT),
t,
format!(
"no PAT section within {} ms",
self.config.pat_max_interval.as_millis()
),
);
}
let pmt_timeouts: Vec<(u16, u64)> = self
.pmt_trackings
.iter()
.filter_map(|(&pid, tracking)| {
if t.saturating_sub(tracking.timer.last_seen) > self.config.pmt_max_interval
&& !tracking.timer.reported
{
Some((pid, self.config.pmt_max_interval.as_millis() as u64))
} else {
None
}
})
.collect();
for (pid, interval_ms) in pmt_timeouts {
if let Some(tracking) = self.pmt_trackings.get_mut(&pid) {
tracking.timer.reported = true;
}
self.emit(
Indicator::PmtError2,
Some(pid),
t,
format!(
"no PMT section on program_map_PID 0x{:04X} within {} ms",
pid, interval_ms
),
);
}
let pid_timeouts: Vec<(u16, u64)> = self
.es_trackings
.iter()
.filter_map(|(&pid, tracking)| {
if t.saturating_sub(tracking.timer.last_seen) > self.config.pid_error_period
&& !tracking.timer.reported
{
Some((pid, self.config.pid_error_period.as_secs()))
} else {
None
}
})
.collect();
for (pid, period_secs) in pid_timeouts {
if let Some(tracking) = self.es_trackings.get_mut(&pid) {
tracking.timer.reported = true;
}
self.emit(
Indicator::PidError,
Some(pid),
t,
format!(
"referenced PID 0x{:04X} absent for > {} s",
pid, period_secs
),
);
}
let si_timeouts: Vec<(u8, u64, u16, u64)> = self
.si_timers
.iter()
.filter_map(|(&table_id, timer)| {
if !timer.armed || timer.reported {
return None;
}
let (limit, pid) = match table_id {
NIT_ACTUAL_TABLE_ID => (self.config.si_nit_interval, PID_NIT),
SDT_ACTUAL_TABLE_ID => (self.config.si_sdt_interval, PID_SDT_BAT),
EIT_PF_ACTUAL_TABLE_ID => (self.config.si_eit_pf_interval, PID_EIT),
TDT_TABLE_ID => (self.config.si_tdt_interval, PID_TDT_TOT),
_ => return None,
};
let interval = t.saturating_sub(timer.last_seen);
if interval > limit {
Some((
table_id,
interval.as_millis() as u64,
pid,
limit.as_millis() as u64,
))
} else {
None
}
})
.collect();
for (table_id, interval_ms, pid, limit_ms) in si_timeouts {
if let Some(timer) = self.si_timers.get_mut(&table_id) {
timer.reported = true;
}
let table_name = match table_id {
NIT_ACTUAL_TABLE_ID => "NIT_actual",
SDT_ACTUAL_TABLE_ID => "SDT_actual",
EIT_PF_ACTUAL_TABLE_ID => "EIT_P/F_actual",
TDT_TABLE_ID => "TDT",
_ => "unknown",
};
self.emit(
Indicator::SiRepetitionError,
Some(pid),
t,
format!(
"{} repetition interval {} ms exceeds {} ms",
table_name, interval_ms, limit_ms
),
);
}
}
}
impl Default for ConformanceMonitor {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests;