#![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;
mod tstd;
use alloc::collections::BTreeMap;
use alloc::format;
use alloc::string::String;
use alloc::vec::Vec;
use core::time::Duration;
use tstd::TstdModel;
use broadcast_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_RST: u16 = 0x0013;
const PID_TDT_TOT: u16 = 0x0014;
const PID_NULL: u16 = 0x1FFF;
const SYNC_BYTE: u8 = 0x47;
const SI_PIDS: [u16; 7] = [
PID_PAT,
PID_CAT,
PID_NIT,
PID_SDT_BAT,
PID_EIT,
PID_RST,
PID_TDT_TOT,
];
const RESERVED_PID_MIN: u16 = 0x0002;
const RESERVED_PID_MAX: u16 = 0x000F;
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 DEFAULT_UNREFERENCED_PID_PERIOD_MS: u64 = 500;
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;
const NIT_OTHER_TABLE_ID: u8 = dvb_si::table_id::TableId::NetworkInformationOther as u8;
const SDT_OTHER_TABLE_ID: u8 = dvb_si::table_id::TableId::ServiceDescriptionOther as u8;
const BAT_TABLE_ID: u8 = dvb_si::table_id::TableId::BouquetAssociation as u8;
const EIT_PF_OTHER_TABLE_ID: u8 = dvb_si::table_id::TableId::EventInformationPfOther as u8;
const RST_TABLE_ID: u8 = dvb_si::table_id::TableId::RunningStatus as u8;
const STUFFING_TABLE_ID: u8 = dvb_si::table_id::TableId::Stuffing as u8;
const TOT_TABLE_ID: u8 = dvb_si::table_id::TableId::TimeOffset 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",
}
}
}
broadcast_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,
NitError,
SiRepetitionError,
UnreferencedPid,
SdtError,
EitError,
RstError,
TdtError,
PcrAccuracyError,
BufferError,
EmptyBufferError,
DataDelayError,
}
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
| Self::PcrAccuracyError => Priority::Second,
Self::NitError
| Self::SiRepetitionError
| Self::UnreferencedPid
| Self::SdtError
| Self::EitError
| Self::RstError
| Self::TdtError
| Self::BufferError
| Self::EmptyBufferError
| Self::DataDelayError => 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::NitError => "NIT_error",
Self::SiRepetitionError => "SI_repetition_error",
Self::UnreferencedPid => "Unreferenced_PID",
Self::SdtError => "SDT_error",
Self::EitError => "EIT_error",
Self::RstError => "RST_error",
Self::TdtError => "TDT_error",
Self::PcrAccuracyError => "PCR_accuracy_error",
Self::BufferError => "Buffer_error",
Self::EmptyBufferError => "Empty_buffer_error",
Self::DataDelayError => "Data_delay_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::NitError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.1",
Self::SiRepetitionError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.2",
Self::UnreferencedPid => "TR 101 290 v1.4.1 Table 5.0c indicator 3.4",
Self::SdtError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.5",
Self::EitError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.6",
Self::RstError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.7",
Self::TdtError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.8",
Self::PcrAccuracyError => "TR 101 290 v1.4.1 Table 5.0b indicator 2.4",
Self::BufferError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.3",
Self::EmptyBufferError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.9",
Self::DataDelayError => "TR 101 290 v1.4.1 Table 5.0c indicator 3.10",
}
}
}
broadcast_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,
pub unreferenced_pid_period: 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),
unreferenced_pid_period: Duration::from_millis(DEFAULT_UNREFERENCED_PID_PERIOD_MS),
}
}
}
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,
}
struct UnreferencedPidTracking {
first_seen: Duration,
reported: 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>,
unreferenced_pid_timers: BTreeMap<u16, UnreferencedPidTracking>,
tstd: TstdModel,
}
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(),
unreferenced_pid_timers: BTreeMap::new(),
tstd: TstdModel::new(Duration::ZERO),
}
}
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{pid:04X}"),
);
}
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{pid:04X}"),
);
}
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{pid:04X} but no CAT seen on PID 0x0001"),
);
}
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.check_nit_table_id(section_bytes, pid, t);
self.check_sdt_table_id(section_bytes, pid, t);
self.check_eit_table_id(section_bytes, pid, t);
self.check_rst_table_id(section_bytes, pid, t);
self.check_tdt_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);
}
}
if pid != PID_NULL {
self.track_unreferenced_pid(pid, t);
}
self.process_tstd(header, ts_packet, 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{pid:04X} (cc={cc})"),
);
}
if should_emit_cc {
self.emit(
Indicator::ContinuityCountError,
Some(pid),
t,
format!("expected cc={expected}, got {cc} on PID 0x{pid:04X}"),
);
}
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,
};
self.feed_tbsys_section_bytes(section_bytes, pid, t);
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 feed_tbsys_section_bytes(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
let section_len = section_bytes.len() as u64;
self.tstd.drain_tb_sys(t);
if self.tstd.tb_sys.would_overflow(section_len, t) {
self.emit(
Indicator::BufferError,
Some(pid),
t,
format!(
"TBsys overflow on PID 0x{pid:04X}: section {section_len} bytes exceeds {} byte capacity at 1 Mbit/s drain",
tstd::TB_SYS_SIZE,
),
);
}
let _overflow = self.tstd.tb_sys.feed(section_len, t);
if self.tstd.tb_sys.check_empty_interval(t) && !self.tstd.tb_sys_empty_reported {
self.tstd.tb_sys_empty_reported = true;
self.emit(
Indicator::EmptyBufferError,
Some(pid),
t,
format!(
"TBsys not empty in the last {} s",
tstd::TB_SYS_EMPTY_INTERVAL_SECS
),
);
}
if let Some(delay) = self.tstd.tb_sys.delay_secs(t) {
if delay > tstd::DATA_DELAY_LIMIT_SECS as f64 {
self.emit(
Indicator::DataDelayError,
Some(pid),
t,
format!(
"TBsys data delay {delay:.2} s exceeds {} s on PID 0x{pid:04X}",
tstd::DATA_DELAY_LIMIT_SECS,
),
);
}
}
}
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 check_nit_table_id(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
if pid != PID_NIT {
return;
}
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
let allowed = section.table_id == NIT_ACTUAL_TABLE_ID
|| section.table_id == NIT_OTHER_TABLE_ID
|| section.table_id == STUFFING_TABLE_ID;
if !allowed {
self.emit(
Indicator::NitError,
Some(PID_NIT),
t,
format!(
"section with table_id 0x{:02X} on PID 0x0010 (expected NIT_actual/NIT_other/ST)",
section.table_id
),
);
}
}
fn check_sdt_table_id(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
if pid != PID_SDT_BAT {
return;
}
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
let allowed = section.table_id == SDT_ACTUAL_TABLE_ID
|| section.table_id == SDT_OTHER_TABLE_ID
|| section.table_id == BAT_TABLE_ID
|| section.table_id == STUFFING_TABLE_ID;
if !allowed {
self.emit(
Indicator::SdtError,
Some(PID_SDT_BAT),
t,
format!(
"section with table_id 0x{:02X} on PID 0x0011 (expected SDT_actual/SDT_other/BAT/ST)",
section.table_id
),
);
}
}
fn check_eit_table_id(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
if pid != PID_EIT {
return;
}
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
let table_id = section.table_id;
let allowed = table_id == EIT_PF_ACTUAL_TABLE_ID
|| table_id == EIT_PF_OTHER_TABLE_ID
|| (dvb_si::tables::eit::TABLE_ID_SCHEDULE_ACTUAL_FIRST
..=dvb_si::tables::eit::TABLE_ID_SCHEDULE_ACTUAL_LAST)
.contains(&table_id)
|| (dvb_si::tables::eit::TABLE_ID_SCHEDULE_OTHER_FIRST
..=dvb_si::tables::eit::TABLE_ID_SCHEDULE_OTHER_LAST)
.contains(&table_id)
|| table_id == STUFFING_TABLE_ID;
if !allowed {
self.emit(
Indicator::EitError,
Some(PID_EIT),
t,
format!(
"section with table_id 0x{table_id:02X} on PID 0x0012 (expected EIT P/F or schedule range or ST)"
),
);
}
}
fn check_rst_table_id(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
if pid != PID_RST {
return;
}
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
let allowed = section.table_id == RST_TABLE_ID || section.table_id == STUFFING_TABLE_ID;
if !allowed {
self.emit(
Indicator::RstError,
Some(PID_RST),
t,
format!(
"section with table_id 0x{:02X} on PID 0x0013 (expected RST/ST)",
section.table_id
),
);
}
}
fn check_tdt_table_id(&mut self, section_bytes: &[u8], pid: u16, t: Duration) {
if pid != PID_TDT_TOT {
return;
}
let section = match Section::parse(section_bytes) {
Ok(s) => s,
Err(_) => return,
};
let allowed = section.table_id == TDT_TABLE_ID
|| section.table_id == TOT_TABLE_ID
|| section.table_id == STUFFING_TABLE_ID;
if !allowed {
self.emit(
Indicator::TdtError,
Some(PID_TDT_TOT),
t,
format!(
"section with table_id 0x{:02X} on PID 0x0014 (expected TDT/TOT/ST)",
section.table_id
),
);
}
}
fn is_referenced_or_reserved_pid(&self, pid: u16) -> bool {
pid == PID_PAT
|| pid == PID_CAT
|| pid == PID_NIT
|| pid == PID_SDT_BAT
|| pid == PID_EIT
|| pid == PID_RST
|| pid == PID_TDT_TOT
|| pid == PID_NULL
|| (RESERVED_PID_MIN..=RESERVED_PID_MAX).contains(&pid)
|| self.pmt_trackings.contains_key(&pid)
|| self.es_trackings.contains_key(&pid)
}
fn track_unreferenced_pid(&mut self, pid: u16, t: Duration) {
if self.is_referenced_or_reserved_pid(pid) {
self.unreferenced_pid_timers.remove(&pid);
return;
}
self.unreferenced_pid_timers
.entry(pid)
.or_insert_with(|| UnreferencedPidTracking {
first_seen: t,
reported: false,
});
}
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(),
});
self.unreferenced_pid_timers.remove(&pmt_pid);
}
}
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,
},
},
);
self.unreferenced_pid_timers.remove(&es_pid);
}
}
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 {delta_ms} ms exceeds limit {limit_ms} ms on PID 0x{pid:04X} without discontinuity_indicator"
),
);
}
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 process_tstd(&mut self, header: &mpeg_ts::ts::TsHeader, ts_packet: &[u8], t: Duration) {
let pid = header.pid;
self.tstd.drain_tb_sys(t);
let is_si_pid = pid == PID_PAT
|| pid == PID_CAT
|| pid == PID_NIT
|| pid == PID_SDT_BAT
|| pid == PID_EIT
|| pid == PID_RST
|| pid == PID_TDT_TOT;
let packet_bytes = ts_packet.len() as u64;
if is_si_pid && header.has_payload {
}
if pid == PID_NULL || is_si_pid {
return;
}
struct TstdCheck {
buffer_overflow: bool,
empty_interval: bool,
delay_exceeded: bool,
}
let check = {
let entry = self.tstd.pid_buffers.entry(pid).or_insert_with(|| {
tstd::PidStdState::new(1_000_000u64, t)
});
entry.total_bytes += packet_bytes;
let leak = {
let elapsed_us = t.saturating_sub(entry.first_seen).as_micros() as u64;
if elapsed_us >= 1_000 {
(entry.total_bytes * 1_000_000 / elapsed_us).max(tstd::TB_LEAK_RATE_FLOOR)
} else {
5_000_000u64
}
};
entry.tb.set_leak_rate(leak);
entry.tb.drain_to(t);
let buffer_overflow = {
entry.tb.drain_to(t);
let _overflow = entry.tb.feed(packet_bytes, t);
false
};
let empty_interval = entry.tb.check_empty_interval(t) && !entry.empty_reported;
let delay_exceeded = if let Some(delay) = entry.tb.delay_secs(t) {
delay > tstd::DATA_DELAY_LIMIT_SECS as f64 && !entry.delay_reported
} else {
false
};
if empty_interval {
entry.empty_reported = true;
}
if delay_exceeded {
entry.delay_reported = true;
}
entry.last_packet_time = t;
TstdCheck {
buffer_overflow,
empty_interval,
delay_exceeded,
}
};
if check.buffer_overflow {
self.emit(
Indicator::BufferError,
Some(pid),
t,
format!(
"TBn overflow on PID 0x{pid:04X}: {} byte capacity exceeded",
tstd::TB_SIZE,
),
);
}
if check.empty_interval {
self.emit(
Indicator::EmptyBufferError,
Some(pid),
t,
format!(
"TBn not empty in the last {} s on PID 0x{pid:04X}",
tstd::TB_EMPTY_INTERVAL_SECS,
),
);
}
if check.delay_exceeded {
self.emit(
Indicator::DataDelayError,
Some(pid),
t,
format!(
"TBn data delay exceeds {} s on PID 0x{pid:04X}",
tstd::DATA_DELAY_LIMIT_SECS,
),
);
}
}
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{pid:04X} within {interval_ms} 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{pid:04X} absent for > {period_secs} s"),
);
}
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, group_indicator) = match table_id {
NIT_ACTUAL_TABLE_ID => ("NIT_actual", Some(Indicator::NitError)),
SDT_ACTUAL_TABLE_ID => ("SDT_actual", Some(Indicator::SdtError)),
EIT_PF_ACTUAL_TABLE_ID => ("EIT_P/F_actual", Some(Indicator::EitError)),
TDT_TABLE_ID => ("TDT", Some(Indicator::TdtError)),
_ => ("unknown", None),
};
self.emit(
Indicator::SiRepetitionError,
Some(pid),
t,
format!("{table_name} repetition interval {interval_ms} ms exceeds {limit_ms} ms"),
);
if let Some(indicator) = group_indicator {
self.emit(
indicator,
Some(pid),
t,
format!("no {table_name} section on PID 0x{pid:04X} within {limit_ms} ms"),
);
}
}
let unref_timeouts: Vec<(u16, u64)> = self
.unreferenced_pid_timers
.iter()
.filter_map(|(&pid, timer)| {
if timer.reported {
return None;
}
let elapsed = t.saturating_sub(timer.first_seen);
if elapsed > self.config.unreferenced_pid_period {
Some((pid, self.config.unreferenced_pid_period.as_millis() as u64))
} else {
None
}
})
.collect();
for (pid, period_ms) in unref_timeouts {
if let Some(timer) = self.unreferenced_pid_timers.get_mut(&pid) {
timer.reported = true;
}
self.emit(
Indicator::UnreferencedPid,
Some(pid),
t,
format!(
"PID 0x{pid:04X} present for > {period_ms} ms without being referenced by PAT/CAT/a PMT or a well-known SI PID"
),
);
}
}
}
impl Default for ConformanceMonitor {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests;