use alloc::collections::BTreeMap;
use alloc::vec;
use alloc::vec::Vec;
use broadcast_common::traits::{Parse, Serialize};
use dvb_si::tables::pat::{PatEntry, PatSection};
use dvb_si::tables::pmt::{self, PmtSection};
use mpeg_ts::mux::SectionPacketizer;
use mpeg_ts::ts::{SectionReassembler, TS_PACKET_SIZE, TsHeader, extract_ts_payload};
use crate::ops::{Op, StreamModel};
const PAT_PID: u16 = 0x0000;
const NULL_PID: u16 = 0x1FFF;
pub(crate) struct PsiRegenOp {
pmt_programs: BTreeMap<u16, u16>,
pmt_reasm: BTreeMap<u16, SectionReassembler>,
transport_stream_id: Option<u16>,
seen_pat_slot: bool,
emitted_pat: bool,
}
impl PsiRegenOp {
pub(crate) fn new() -> Self {
Self {
pmt_programs: BTreeMap::new(),
pmt_reasm: BTreeMap::new(),
transport_stream_id: None,
seen_pat_slot: false,
emitted_pat: false,
}
}
fn ts_payload_and_pusi(packet: &[u8]) -> Option<(&[u8], bool)> {
let header = TsHeader::parse(&packet[..4]).ok()?;
let payload = extract_ts_payload(packet)?;
Some((payload, header.pusi))
}
fn pid_from_packet(packet: &[u8]) -> u16 {
(((packet[1] & 0x1F) as u16) << 8) | packet[2] as u16
}
fn pusi_table_id(payload: &[u8]) -> Option<u8> {
let pointer = *payload.first()? as usize;
payload.get(1 + pointer).copied()
}
fn observe_pat_tsid(&mut self, payload: &[u8], pusi: bool) {
if self.transport_stream_id.is_some() {
return;
}
let section = if pusi {
let pointer = match payload.first() {
Some(&p) => p as usize,
None => return,
};
match payload.get(1 + pointer..) {
Some(s) => s,
None => return,
}
} else {
payload
};
if let Ok(pat) = PatSection::parse(section) {
self.transport_stream_id = Some(pat.transport_stream_id);
}
}
fn scan_pmt(&mut self, pid: u16, payload: &[u8], pusi: bool) -> bool {
let reasm = match self.pmt_reasm.entry(pid) {
alloc::collections::btree_map::Entry::Occupied(e) => e.into_mut(),
alloc::collections::btree_map::Entry::Vacant(slot) => {
if !pusi {
return false;
}
match Self::pusi_table_id(payload) {
Some(tid) if tid == pmt::TABLE_ID => {}
_ => return false,
}
slot.insert(SectionReassembler::default())
}
};
reasm.feed(payload, pusi);
let mut cycle_wrapped = false;
while let Some(section) = reasm.pop_section() {
if let Ok(p) = PmtSection::parse(§ion) {
let prev = self.pmt_programs.insert(p.program_number, pid);
if prev.is_some() {
cycle_wrapped = true;
}
}
}
cycle_wrapped
}
fn rebuild_pat(&self) -> Option<Vec<u8>> {
if self.pmt_programs.is_empty() {
return None;
}
let entries: Vec<PatEntry> = self
.pmt_programs
.iter()
.map(|(&program_number, &pmt_pid)| PatEntry {
program_number,
pid: pmt_pid,
})
.collect();
let pat = PatSection {
transport_stream_id: self.transport_stream_id.unwrap_or(0),
version_number: 0,
current_next_indicator: true,
section_number: 0,
last_section_number: 0,
entries,
};
let mut buf = vec![0u8; pat.serialized_len()];
pat.serialize_into(&mut buf).ok()?;
Some(buf)
}
fn emit_regen_pat(&mut self, out: &mut dyn FnMut(&[u8])) {
if let Some(section) = self.rebuild_pat() {
let mut packetizer = SectionPacketizer::new(PAT_PID);
for pkt in packetizer.packetize(&[§ion]) {
out(&pkt);
}
self.emitted_pat = true;
}
}
}
impl Op for PsiRegenOp {
fn process(&mut self, packet: &[u8], _model: &mut StreamModel, out: &mut dyn FnMut(&[u8])) {
if packet.len() != TS_PACKET_SIZE {
out(packet);
return;
}
let pid = Self::pid_from_packet(packet);
if pid == PAT_PID {
self.seen_pat_slot = true;
if let Some((payload, pusi)) = Self::ts_payload_and_pusi(packet) {
self.observe_pat_tsid(payload, pusi);
}
self.emit_regen_pat(out);
return;
}
if pid == NULL_PID {
out(packet);
return;
}
if let Some((payload, pusi)) = Self::ts_payload_and_pusi(packet) {
let cycle_wrapped = self.scan_pmt(pid, payload, pusi);
if !self.seen_pat_slot && !self.emitted_pat && cycle_wrapped {
self.emit_regen_pat(out);
}
}
out(packet);
}
fn flush(&mut self, _model: &mut StreamModel, out: &mut dyn FnMut(&[u8])) {
if !self.emitted_pat {
self.emit_regen_pat(out);
}
}
}