use alloc::collections::BTreeMap;
use mpeg_ts::owned::OwnedTsPacket;
use mpeg_ts::ts::{Pcr, TS_PACKET_SIZE, TsPacket};
use crate::ops::{Op, StreamModel};
const PCR_27MHZ_MODULUS: u64 = (1u64 << 33) * 300;
#[non_exhaustive]
#[derive(Debug, Clone)]
pub enum PcrRestamp {
Interpolate,
FromBitrate {
bps: u64,
},
}
impl PcrRestamp {
pub fn interpolate() -> Self {
Self::Interpolate
}
pub fn from_bitrate(bps: u64) -> Self {
Self::FromBitrate { bps }
}
}
#[derive(Clone, Copy)]
struct Anchor {
anchor_27mhz: u64,
anchor_pkt: u64,
last_obs_pkt: u64,
last_obs_27mhz: u64,
}
pub(crate) struct PcrRestampOp {
anchors: BTreeMap<u16, Anchor>,
mode: PcrRestamp,
}
impl PcrRestampOp {
pub(crate) fn new(mode: PcrRestamp) -> Self {
Self {
anchors: BTreeMap::new(),
mode,
}
}
fn ticks_per_packet(bps: u64) -> u64 {
let num = 188u64 * 8 * 27_000_000u64;
if bps == 0 || bps >= num {
1
} else {
(num / bps).max(1)
}
}
fn read_pcr(packet: &[u8]) -> Option<(u16, Pcr, bool)> {
let pkt = TsPacket::parse(packet).ok()?;
let af = pkt.adaptation_field().and_then(|r| r.ok())?;
let pcr = af.pcr?;
Some((pkt.header.pid, pcr, af.discontinuity_indicator))
}
}
impl Op for PcrRestampOp {
fn process(&mut self, packet: &[u8], model: &mut StreamModel, out: &mut dyn FnMut(&[u8])) {
if packet.len() != TS_PACKET_SIZE {
out(packet);
return;
}
let Some((pid, current, discontinuity)) = Self::read_pcr(packet) else {
out(packet);
return;
};
let now = model.packet_count;
if discontinuity {
let a = Anchor {
anchor_27mhz: current.as_27mhz(),
anchor_pkt: now,
last_obs_pkt: now,
last_obs_27mhz: current.as_27mhz(),
};
self.anchors.insert(pid, a);
model.timing.has_anchor = true;
model.timing.clock_27mhz = current.as_27mhz();
out(packet);
return;
}
let Some(anchor) = self.anchors.get_mut(&pid) else {
let a = Anchor {
anchor_27mhz: current.as_27mhz(),
anchor_pkt: now,
last_obs_pkt: now,
last_obs_27mhz: current.as_27mhz(),
};
self.anchors.insert(pid, a);
model.timing.has_anchor = true;
model.timing.clock_27mhz = current.as_27mhz();
out(packet);
return;
};
let new_27mhz = match &self.mode {
PcrRestamp::FromBitrate { bps } => {
let delta = now.saturating_sub(anchor.anchor_pkt);
anchor
.anchor_27mhz
.wrapping_add(Self::ticks_per_packet(*bps) * delta)
% PCR_27MHZ_MODULUS
}
PcrRestamp::Interpolate => {
let obs = current.as_27mhz();
let pkt_delta = now.saturating_sub(anchor.last_obs_pkt);
let fwd = obs.wrapping_sub(anchor.last_obs_27mhz) % PCR_27MHZ_MODULUS;
if fwd > 0 && fwd < PCR_27MHZ_MODULUS / 2 && pkt_delta > 0 {
anchor.last_obs_pkt = now;
anchor.last_obs_27mhz = obs;
obs
} else {
let span_pkt = anchor.last_obs_pkt.saturating_sub(anchor.anchor_pkt).max(1);
let span_ticks = anchor.last_obs_27mhz.saturating_sub(anchor.anchor_27mhz);
let rate = (span_ticks / span_pkt).max(1);
let delta = now.saturating_sub(anchor.anchor_pkt);
anchor.anchor_27mhz.wrapping_add(rate * delta) % PCR_27MHZ_MODULUS
}
}
};
let mut buf = [0u8; TS_PACKET_SIZE];
buf.copy_from_slice(packet);
if OwnedTsPacket::set_pcr(&mut buf, Pcr::from_27mhz(new_27mhz)).is_ok() {
out(&buf);
} else {
out(packet);
}
}
fn flush(&mut self, _model: &mut StreamModel, _out: &mut dyn FnMut(&[u8])) {
}
}