use alloc::collections::btree_map::BTreeMap;
use alloc::string::String;
use alloc::vec::Vec;
use core::fmt::Write as _;
use core::time::Duration;
use broadcast_common::Parse;
use dvb_conformance::ConformanceMonitor;
use dvb_si::demux::SiDemux;
use dvb_si::tables::AnyTableSection;
use dvb_si::tables::pmt::StreamType;
use mpeg_pes::{PesAssembler, PesPacket};
use mpeg_ts::resync::TsResync;
use mpeg_ts::ts::{SectionReassembler, TS_PACKET_SIZE, TsPacket};
use transmux::{iter_annexb_nals, parse_adts_header};
const SCTE35_TABLE_ID: u8 = 0xFC;
const PTS_MODULUS: u64 = 1u64 << 33;
const PTS_HALF: u64 = 1u64 << 32;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum EsKind {
Video,
AudioAdts,
}
struct EsTrack {
kind: EsKind,
assembler: PesAssembler,
any_au: bool,
structured: bool,
prev_decode: Option<u64>,
decode_anomaly: bool,
}
impl EsTrack {
fn new(kind: EsKind) -> Self {
Self {
kind,
assembler: PesAssembler::new(),
any_au: false,
structured: false,
prev_decode: None,
decode_anomaly: false,
}
}
}
#[derive(Default)]
struct Scte35Track {
reassembler: SectionReassembler,
events: BTreeMap<u32, bool>,
}
struct ConformanceCount {
priority: &'static str,
clause: &'static str,
count: u64,
}
pub struct WatchState {
resync: TsResync,
conformance: ConformanceMonitor,
demux: SiDemux,
es_tracks: BTreeMap<u16, EsTrack>,
scte35_tracks: BTreeMap<u16, Scte35Track>,
conformance_counts: BTreeMap<&'static str, ConformanceCount>,
datagrams_total: u64,
scte35_events_total: u64,
pts_dts_anomalies_total: u64,
last_clock: Duration,
}
impl Default for WatchState {
fn default() -> Self {
Self::new()
}
}
impl WatchState {
#[must_use]
pub fn new() -> Self {
Self {
resync: TsResync::new(),
conformance: ConformanceMonitor::new(),
demux: SiDemux::builder().build(),
es_tracks: BTreeMap::new(),
scte35_tracks: BTreeMap::new(),
conformance_counts: BTreeMap::new(),
datagrams_total: 0,
scte35_events_total: 0,
pts_dts_anomalies_total: 0,
last_clock: Duration::ZERO,
}
}
pub fn feed_datagram(&mut self, payload: &[u8], clock: Duration) {
self.datagrams_total += 1;
let packets = self.resync.feed(payload);
for packet in &packets {
self.feed_ts_packet(packet, clock);
}
}
fn feed_ts_packet(&mut self, packet: &[u8; TS_PACKET_SIZE], clock: Duration) {
self.last_clock = clock;
for ev in self.conformance.feed(packet, clock) {
let entry = self
.conformance_counts
.entry(ev.indicator.name())
.or_insert_with(|| ConformanceCount {
priority: ev.priority.name(),
clause: ev.indicator.clause(),
count: 0,
});
entry.count += 1;
}
let Ok(ts_packet) = TsPacket::parse(packet) else {
return;
};
let pid = ts_packet.header.pid;
let events: Vec<_> = self.demux.feed(packet).collect();
for ev in events {
if let Ok(AnyTableSection::PmtSection(pmt)) = ev.table_section() {
for stream in &pmt.streams {
match stream.stream_type {
StreamType::Scte35 => {
self.scte35_tracks.entry(stream.elementary_pid).or_default();
}
StreamType::H264 | StreamType::Hevc => {
self.es_tracks
.entry(stream.elementary_pid)
.or_insert_with(|| EsTrack::new(EsKind::Video));
}
StreamType::AacAdts => {
self.es_tracks
.entry(stream.elementary_pid)
.or_insert_with(|| EsTrack::new(EsKind::AudioAdts));
}
_ => {}
}
}
}
}
if self.scte35_tracks.contains_key(&pid) {
self.feed_scte35(pid, &ts_packet);
}
if self.es_tracks.contains_key(&pid) {
self.feed_es(pid, &ts_packet);
}
}
fn feed_scte35(&mut self, pid: u16, ts_packet: &TsPacket<'_>) {
let Some(payload) = ts_packet.payload else {
return;
};
let Some(track) = self.scte35_tracks.get_mut(&pid) else {
return;
};
track.reassembler.feed(payload, ts_packet.header.pusi);
while let Some(section) = track.reassembler.pop_section() {
if section.is_empty() || section[0] != SCTE35_TABLE_ID {
continue;
}
let Ok(sis) = scte35_splice::SpliceInfoSection::parse(§ion[..]) else {
continue;
};
let Some(ref clear) = sis.clear else {
continue;
};
let scte35_splice::commands::AnyCommand::SpliceInsert(si) = &clear.command else {
continue;
};
if si.splice_event_cancel_indicator {
continue;
}
self.scte35_events_total += 1;
track
.events
.insert(si.splice_event_id, si.out_of_network_indicator);
}
}
fn feed_es(&mut self, pid: u16, ts_packet: &TsPacket<'_>) {
let mut discontinuity = false;
if ts_packet.header.has_adaptation
&& let Some(Ok(af)) = ts_packet.adaptation_field()
{
discontinuity = af.discontinuity_indicator;
}
let Some(track) = self.es_tracks.get_mut(&pid) else {
return;
};
if discontinuity {
track.assembler = PesAssembler::new();
track.prev_decode = None;
return;
}
let Some(payload) = ts_packet.payload else {
return;
};
if payload.is_empty() {
return;
}
let pes_bytes = track.assembler.feed(ts_packet.header.pusi, payload);
let Some(pes_bytes) = pes_bytes else {
return;
};
self.process_es_unit(pid, &pes_bytes);
}
fn process_es_unit(&mut self, pid: u16, pes_bytes: &[u8]) {
let Ok(pes) = PesPacket::parse(pes_bytes) else {
return;
};
let Some(track) = self.es_tracks.get_mut(&pid) else {
return;
};
track.any_au = true;
match track.kind {
EsKind::Video => {
if iter_annexb_nals(pes.payload).next().is_some() {
track.structured = true;
}
}
EsKind::AudioAdts => {
if has_adts_sync(pes.payload) {
track.structured = true;
}
}
}
let Some(header) = pes.header else {
return;
};
let (raw, present) = match (header.dts, header.pts) {
(Some(dts), _) => (dts.ticks(), true),
(None, Some(pts)) => (pts.ticks(), true),
(None, None) => (0, false),
};
if !present {
return;
}
if let Some(prev) = track.prev_decode {
let delta = raw.wrapping_sub(prev) & (PTS_MODULUS - 1);
if delta != 0 && delta > PTS_HALF {
track.decode_anomaly = true;
self.pts_dts_anomalies_total += 1;
}
}
track.prev_decode = Some(raw);
}
#[must_use]
pub fn render_prometheus(&self) -> String {
let mut out = String::new();
let conformance_stats = self.conformance.stats();
let resync_stats = self.resync.stats();
metric_header(
&mut out,
"media_doctor_packets_total",
"Total well-formed 188-byte TS packets processed (ISO/IEC 13818-1 section 2.4.3.2).",
"counter",
);
let _ = writeln!(
out,
"media_doctor_packets_total {}",
conformance_stats.packets
);
metric_header(
&mut out,
"media_doctor_datagrams_total",
"Total ingest datagrams fed (e.g. UDP payloads).",
"counter",
);
let _ = writeln!(out, "media_doctor_datagrams_total {}", self.datagrams_total);
metric_header(
&mut out,
"media_doctor_resync_events_total",
"Times TS byte-stream sync was lost and reacquired (mpeg_ts::resync::TsResync).",
"counter",
);
let _ = writeln!(
out,
"media_doctor_resync_events_total {}",
resync_stats.resyncs
);
metric_header(
&mut out,
"media_doctor_dropped_bytes_total",
"Bytes dropped before/while reacquiring TS packet sync.",
"counter",
);
let _ = writeln!(
out,
"media_doctor_dropped_bytes_total {}",
resync_stats.dropped_bytes
);
metric_header(
&mut out,
"media_doctor_conformance_in_sync",
"Whether the ETSI TR 101 290 monitor currently considers the stream in sync (1) or not (0).",
"gauge",
);
let _ = writeln!(
out,
"media_doctor_conformance_in_sync {}",
u8::from(conformance_stats.in_sync)
);
metric_header(
&mut out,
"media_doctor_conformance_events_total",
"ETSI TR 101 290 indicator events observed, by indicator and priority tier.",
"counter",
);
for (name, c) in &self.conformance_counts {
let _ = writeln!(
out,
"media_doctor_conformance_events_total{{indicator=\"{}\",priority=\"{}\"}} {}",
escape_label(name),
escape_label(c.priority),
c.count,
);
}
if !self.conformance_counts.is_empty() {
out.push_str("# clauses: ");
let mut first = true;
for (name, c) in &self.conformance_counts {
if !first {
out.push_str(", ");
}
first = false;
let _ = write!(out, "{name}={}", c.clause);
}
out.push('\n');
}
metric_header(
&mut out,
"media_doctor_scte35_events_total",
"Total SCTE-35 splice_insert events observed (ANSI/SCTE 35 section 9.7.3.1), excluding cancelled events.",
"counter",
);
let _ = writeln!(
out,
"media_doctor_scte35_events_total {}",
self.scte35_events_total
);
let scte35_open: u64 = self
.scte35_tracks
.values()
.flat_map(|t| t.events.values())
.filter(|&&open| open)
.count() as u64;
metric_header(
&mut out,
"media_doctor_scte35_open_events",
"Currently-unmatched (\"out\" with no \"in\" yet) SCTE-35 splice_insert events.",
"gauge",
);
let _ = writeln!(out, "media_doctor_scte35_open_events {scte35_open}");
metric_header(
&mut out,
"media_doctor_pts_dts_anomalies_total",
"Non-monotonic decode-timestamp (DTS, else PTS) events observed on tracked PES PIDs.",
"counter",
);
let _ = writeln!(
out,
"media_doctor_pts_dts_anomalies_total {}",
self.pts_dts_anomalies_total
);
metric_header(
&mut out,
"media_doctor_codec_signalling_mismatch",
"Whether a PMT-declared codec PID has ever shown bitstream framing disagreeing with \
the declared stream_type (1) or not (0); only emitted once at least one access unit \
has been observed on that PID (ISO/IEC 13818-1 Table 2-34).",
"gauge",
);
for (&pid, track) in &self.es_tracks {
if track.any_au {
let mismatch = u8::from(!track.structured);
let _ = writeln!(
out,
"media_doctor_codec_signalling_mismatch{{pid=\"0x{pid:04X}\"}} {mismatch}"
);
}
}
metric_header(
&mut out,
"media_doctor_pts_dts_anomaly",
"Whether a tracked PES PID has ever shown a non-monotonic decode timestamp (1) or not \
(0); only emitted once a decode timestamp has been observed on that PID.",
"gauge",
);
for (&pid, track) in &self.es_tracks {
if track.prev_decode.is_some() {
let _ = writeln!(
out,
"media_doctor_pts_dts_anomaly{{pid=\"0x{pid:04X}\"}} {}",
u8::from(track.decode_anomaly)
);
}
}
metric_header(
&mut out,
"media_doctor_last_packet_clock_seconds",
"Elapsed ingest wall-clock time (seconds) of the most recently processed TS packet.",
"gauge",
);
let _ = writeln!(
out,
"media_doctor_last_packet_clock_seconds {}",
self.last_clock.as_secs_f64()
);
out
}
}
fn metric_header(out: &mut String, name: &str, help: &str, ty: &str) {
let _ = writeln!(out, "# HELP {name} {help}");
let _ = writeln!(out, "# TYPE {name} {ty}");
}
fn escape_label(v: &str) -> String {
let mut out = String::with_capacity(v.len());
for c in v.chars() {
match c {
'\\' => out.push_str("\\\\"),
'"' => out.push_str("\\\""),
'\n' => out.push_str("\\n"),
_ => out.push(c),
}
}
out
}
fn has_adts_sync(payload: &[u8]) -> bool {
const ADTS_MIN: usize = 7;
if payload.len() < ADTS_MIN {
return false;
}
(0..=payload.len() - ADTS_MIN).any(|off| {
payload[off] == 0xFF
&& (payload[off + 1] & 0xF0) == 0xF0
&& parse_adts_header(&payload[off..]).is_ok()
})
}
#[cfg(test)]
mod tests {
use super::*;
fn fixture(rel: &str) -> Vec<u8> {
let path = format!(concat!(env!("CARGO_MANIFEST_DIR"), "/../fixtures/{}"), rel);
std::fs::read(&path).unwrap_or_else(|e| panic!("read fixture {path}: {e}"))
}
fn chunks(bytes: &[u8], size: usize) -> Vec<&[u8]> {
bytes.chunks(size).collect()
}
#[test]
fn watches_real_fixture_and_produces_sane_metrics() {
let bytes = fixture("ts/m6-single.ts");
assert_eq!(
bytes.len() % TS_PACKET_SIZE,
0,
"fixture must be whole TS packets"
);
let expected_packets = bytes.len() / TS_PACKET_SIZE;
let mut state = WatchState::new();
let mut clock = Duration::ZERO;
for datagram in chunks(&bytes, 7 * TS_PACKET_SIZE) {
state.feed_datagram(datagram, clock);
clock += Duration::from_millis(1);
}
let text = state.render_prometheus();
assert_eq!(
metric_value(&text, "media_doctor_packets_total"),
Some(expected_packets as f64),
"packets_total must match the fixture's real packet count:\n{text}"
);
assert_eq!(
metric_value(&text, "media_doctor_datagrams_total"),
Some(chunks(&bytes, 7 * TS_PACKET_SIZE).len() as f64)
);
assert_eq!(
metric_value(&text, "media_doctor_resync_events_total"),
Some(0.0)
);
assert_eq!(
metric_value(&text, "media_doctor_dropped_bytes_total"),
Some(0.0)
);
assert!(text.contains("media_doctor_conformance_in_sync"));
for name in [
"media_doctor_packets_total",
"media_doctor_datagrams_total",
"media_doctor_scte35_events_total",
"media_doctor_scte35_open_events",
"media_doctor_pts_dts_anomalies_total",
"media_doctor_last_packet_clock_seconds",
] {
assert!(
text.contains(name),
"missing metric family {name} in:\n{text}"
);
}
}
#[test]
fn misaligned_datagram_resyncs_and_still_counts_packets() {
let bytes = fixture("ts/m6-single.ts");
let mut misaligned = alloc::vec![0xAAu8; 37];
misaligned.extend_from_slice(&bytes[..20 * TS_PACKET_SIZE]);
let mut state = WatchState::new();
state.feed_datagram(&misaligned, Duration::ZERO);
let text = state.render_prometheus();
let packets = metric_value(&text, "media_doctor_packets_total").unwrap();
assert!(
packets >= 15.0,
"expected the resynchroniser to recover most of the 20 fed packets, got {packets}"
);
assert_eq!(
metric_value(&text, "media_doctor_dropped_bytes_total"),
Some(37.0),
"the 37 leading garbage bytes must be counted as dropped"
);
}
#[test]
fn garbage_datagram_never_panics() {
let mut state = WatchState::new();
state.feed_datagram(&[0u8; 4096], Duration::ZERO);
let text = state.render_prometheus();
assert_eq!(metric_value(&text, "media_doctor_packets_total"), Some(0.0));
}
#[test]
fn scte35_declared_in_pmt_is_reassembled_and_counted() {
use broadcast_common::Serialize;
use dvb_si::descriptors::DescriptorLoop;
use dvb_si::tables::pat::{PatEntry, PatSection};
use dvb_si::tables::pmt::{PmtSection, PmtStream};
use mpeg_ts::pid::well_known as wk;
use mpeg_ts::ts::TsHeader;
use scte35_splice::SpliceInfoSection;
use scte35_splice::commands::{AnyCommand, SpliceInsert};
const PMT_PID: u16 = 0x0100;
const SCTE35_PID: u16 = 0x0101;
fn ts_packet(pid: u16, section: &[u8]) -> [u8; TS_PACKET_SIZE] {
let mut pkt = [0xFFu8; TS_PACKET_SIZE];
let header = TsHeader {
tei: false,
pusi: true,
pid,
scrambling: 0,
has_adaptation: false,
has_payload: true,
continuity_counter: 0,
};
header.serialize_into(&mut pkt).unwrap();
pkt[4] = 0x00; pkt[5..5 + section.len()].copy_from_slice(section);
pkt
}
let pat = PatSection {
transport_stream_id: 1,
version_number: 0,
current_next_indicator: true,
section_number: 0,
last_section_number: 0,
entries: alloc::vec![PatEntry {
program_number: 1,
pid: PMT_PID,
}],
};
let mut pat_bytes = alloc::vec![0u8; pat.serialized_len()];
pat.serialize_into(&mut pat_bytes).unwrap();
let pmt = PmtSection::new(
1,
0,
true,
0,
0,
wk::NULL.value(),
DescriptorLoop::new(&[]),
alloc::vec![PmtStream {
stream_type: StreamType::Scte35,
elementary_pid: SCTE35_PID,
es_info: DescriptorLoop::new(&[]),
}],
);
let mut pmt_bytes = alloc::vec![0u8; pmt.serialized_len()];
pmt.serialize_into(&mut pmt_bytes).unwrap();
let splice_out = SpliceInfoSection::new_clear(
AnyCommand::SpliceInsert(SpliceInsert {
splice_event_id: 42,
out_of_network_indicator: true,
program_splice_flag: true,
splice_immediate_flag: true,
..SpliceInsert::default()
}),
&[],
);
let out_bytes = splice_out.to_bytes();
let mut state = WatchState::new();
let mut clock = Duration::ZERO;
for packet in [
ts_packet(wk::PAT.value(), &pat_bytes),
ts_packet(PMT_PID, &pmt_bytes),
ts_packet(SCTE35_PID, &out_bytes),
ts_packet(wk::NULL.value(), &[]),
ts_packet(wk::NULL.value(), &[]),
] {
state.feed_datagram(&packet, clock);
clock += Duration::from_millis(1);
}
let text = state.render_prometheus();
assert_eq!(
metric_value(&text, "media_doctor_scte35_events_total"),
Some(1.0)
);
assert_eq!(
metric_value(&text, "media_doctor_scte35_open_events"),
Some(1.0),
"the 'out' has no matching 'in' yet, so it must show as open:\n{text}"
);
}
fn metric_value(text: &str, name: &str) -> Option<f64> {
for line in text.lines() {
if let Some(rest) = line.strip_prefix(name) {
let rest = rest.trim_start();
if rest.starts_with('{') {
continue;
}
return rest.trim().parse().ok();
}
}
None
}
}