use crate::AdaptationField;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct Pcr {
pub ticks_27mhz: u64,
}
pub const PCR_MODULUS_27MHZ: u64 = (1u64 << 33) * 300;
pub const PCR_TOLERANCE_27MHZ: u32 = 14;
impl Pcr {
pub fn from_base_ext(base_90khz: u64, extension: u16) -> Self {
Self {
ticks_27mhz: (base_90khz & ((1u64 << 33) - 1)) * 300 + (extension as u64 % 300),
}
}
pub fn from_ticks_27mhz(ticks: u64) -> Self {
Self {
ticks_27mhz: ticks % PCR_MODULUS_27MHZ,
}
}
pub fn base_90khz(self) -> u64 {
self.ticks_27mhz / 300
}
pub fn extension(self) -> u16 {
(self.ticks_27mhz % 300) as u16
}
pub fn delta_ticks(self, other: Pcr) -> i64 {
let half = (PCR_MODULUS_27MHZ / 2) as i64;
let m = PCR_MODULUS_27MHZ as i64;
let raw = self.ticks_27mhz as i64 - other.ticks_27mhz as i64;
if raw > half {
raw - m
} else if raw < -half {
raw + m
} else {
raw
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DiscontinuityReason {
Signalled,
JumpExceedsTolerance {
observed_delta: i64,
predicted_delta: i64,
},
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum PcrEvent {
Sample {
pcr: Pcr,
jitter_27mhz: i64,
bitrate_bps: u64,
},
Discontinuity {
pcr: Pcr,
reason: DiscontinuityReason,
},
}
#[derive(Debug, Clone)]
pub struct PcrTracker {
pub max_jitter_27mhz: u32,
last: Option<(Pcr, u64)>,
prev: Option<(Pcr, u64)>,
pending_signalled_discontinuity: bool,
}
impl PcrTracker {
pub const DEFAULT_MAX_JITTER_27MHZ: u32 = 2_700_000;
pub fn new() -> Self {
Self::with_jitter_window(Self::DEFAULT_MAX_JITTER_27MHZ)
}
pub fn with_jitter_window(max_jitter_27mhz: u32) -> Self {
Self {
max_jitter_27mhz,
last: None,
prev: None,
pending_signalled_discontinuity: false,
}
}
pub fn observe(&mut self, af: &AdaptationField<'_>, byte_offset: u64) -> Option<PcrEvent> {
if af.discontinuity_indicator {
self.pending_signalled_discontinuity = true;
}
let (base, ext) = match (af.pcr_base, af.pcr_extension) {
(Some(b), Some(e)) => (b, e),
_ => return None,
};
let pcr = Pcr::from_base_ext(base, ext);
if self.pending_signalled_discontinuity {
self.pending_signalled_discontinuity = false;
self.last = Some((pcr, byte_offset));
self.prev = None;
return Some(PcrEvent::Discontinuity {
pcr,
reason: DiscontinuityReason::Signalled,
});
}
let event = match (self.prev, self.last) {
(None, None) => PcrEvent::Sample {
pcr,
jitter_27mhz: 0,
bitrate_bps: 0,
},
(None, Some((last_pcr, last_off))) => {
let dpcr = pcr.delta_ticks(last_pcr);
let dbytes = byte_offset.saturating_sub(last_off);
let bitrate = compute_bitrate_bps(dpcr, dbytes);
PcrEvent::Sample {
pcr,
jitter_27mhz: 0,
bitrate_bps: bitrate,
}
}
(Some((prev_pcr, prev_off)), Some((last_pcr, last_off))) => {
let prev_dpcr = last_pcr.delta_ticks(prev_pcr);
let prev_dbytes = last_off.saturating_sub(prev_off).max(1);
let cur_dbytes = byte_offset.saturating_sub(last_off);
let predicted_delta =
((prev_dpcr as i128) * (cur_dbytes as i128) / (prev_dbytes as i128)) as i64;
let observed_delta = pcr.delta_ticks(last_pcr);
let jitter = observed_delta - predicted_delta;
if jitter.unsigned_abs() > self.max_jitter_27mhz as u64 {
self.last = Some((pcr, byte_offset));
self.prev = None;
return Some(PcrEvent::Discontinuity {
pcr,
reason: DiscontinuityReason::JumpExceedsTolerance {
observed_delta,
predicted_delta,
},
});
}
let bitrate = compute_bitrate_bps(observed_delta, cur_dbytes);
PcrEvent::Sample {
pcr,
jitter_27mhz: jitter,
bitrate_bps: bitrate,
}
}
(Some(_), None) => PcrEvent::Sample {
pcr,
jitter_27mhz: 0,
bitrate_bps: 0,
},
};
self.prev = self.last;
self.last = Some((pcr, byte_offset));
Some(event)
}
pub fn last_pcr(&self) -> Option<Pcr> {
self.last.map(|(p, _)| p)
}
pub fn reset(&mut self) {
self.last = None;
self.prev = None;
self.pending_signalled_discontinuity = false;
}
}
impl Default for PcrTracker {
fn default() -> Self {
Self::new()
}
}
fn compute_bitrate_bps(delta_pcr_27mhz: i64, delta_bytes: u64) -> u64 {
if delta_pcr_27mhz <= 0 || delta_bytes == 0 {
return 0;
}
let num = (delta_bytes as u128) * 8 * 27_000_000;
let den = delta_pcr_27mhz as u128;
(num / den) as u64
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ContinuityEvent {
Continuous,
Duplicate,
NoPayload,
Dropped { gap: u8 },
Discontinuity,
}
#[derive(Debug, Default, Clone)]
pub struct ContinuityTracker {
last_payload_cc: Option<u8>,
duplicate_already_seen: bool,
}
impl ContinuityTracker {
pub fn new() -> Self {
Self::default()
}
pub fn observe(
&mut self,
cc: u8,
has_payload: bool,
discontinuity_indicator: bool,
) -> ContinuityEvent {
let cc = cc & 0x0F;
if discontinuity_indicator {
if has_payload {
self.last_payload_cc = Some(cc);
} else {
self.last_payload_cc = None;
}
self.duplicate_already_seen = false;
return ContinuityEvent::Discontinuity;
}
if !has_payload {
return ContinuityEvent::NoPayload;
}
let prev = match self.last_payload_cc {
None => {
self.last_payload_cc = Some(cc);
self.duplicate_already_seen = false;
return ContinuityEvent::Continuous;
}
Some(p) => p,
};
let expected = (prev + 1) & 0x0F;
if cc == expected {
self.last_payload_cc = Some(cc);
self.duplicate_already_seen = false;
ContinuityEvent::Continuous
} else if cc == prev && !self.duplicate_already_seen {
self.duplicate_already_seen = true;
ContinuityEvent::Duplicate
} else {
let gap = cc.wrapping_sub(expected) & 0x0F;
self.last_payload_cc = Some(cc);
self.duplicate_already_seen = false;
ContinuityEvent::Dropped { gap: gap.max(1) }
}
}
pub fn reset(&mut self) {
self.last_payload_cc = None;
self.duplicate_already_seen = false;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::AdaptationField;
fn af_with_pcr(base: u64, ext: u16, discontinuity: bool) -> AdaptationField<'static> {
AdaptationField {
length: 7,
discontinuity_indicator: discontinuity,
random_access_indicator: false,
elementary_stream_priority_indicator: false,
pcr_flag: true,
opcr_flag: false,
splicing_point_flag: false,
transport_private_data_flag: false,
adaptation_field_extension_flag: false,
pcr_base: Some(base),
pcr_extension: Some(ext),
opcr_base: None,
opcr_extension: None,
splice_countdown: None,
transport_private_data: None,
adaptation_field_extension: None,
raw: &[],
}
}
fn af_no_pcr(discontinuity: bool) -> AdaptationField<'static> {
AdaptationField {
length: 1,
discontinuity_indicator: discontinuity,
random_access_indicator: false,
elementary_stream_priority_indicator: false,
pcr_flag: false,
opcr_flag: false,
splicing_point_flag: false,
transport_private_data_flag: false,
adaptation_field_extension_flag: false,
pcr_base: None,
pcr_extension: None,
opcr_base: None,
opcr_extension: None,
splice_countdown: None,
transport_private_data: None,
adaptation_field_extension: None,
raw: &[],
}
}
#[test]
fn pcr_round_trips_base_ext() {
let p = Pcr::from_base_ext(0x1_2345_6789, 200);
assert_eq!(p.base_90khz(), 0x1_2345_6789);
assert_eq!(p.extension(), 200);
let q = Pcr::from_ticks_27mhz(p.ticks_27mhz);
assert_eq!(p, q);
}
#[test]
fn pcr_delta_handles_wraparound() {
let near_max = Pcr::from_ticks_27mhz(PCR_MODULUS_27MHZ - 1000);
let just_after = Pcr::from_ticks_27mhz(500);
assert_eq!(just_after.delta_ticks(near_max), 1500);
assert_eq!(near_max.delta_ticks(just_after), -1500);
}
#[test]
fn pcr_tracker_bootstraps_quietly() {
let mut t = PcrTracker::new();
let pcr1 = Pcr::from_ticks_27mhz(0);
let pcr2 = Pcr::from_ticks_27mhz(216 * 1880); let pcr3 = Pcr::from_ticks_27mhz(216 * 1880 * 2);
let af1 = af_with_pcr(pcr1.base_90khz(), pcr1.extension(), false);
let af2 = af_with_pcr(pcr2.base_90khz(), pcr2.extension(), false);
let af3 = af_with_pcr(pcr3.base_90khz(), pcr3.extension(), false);
let e1 = t.observe(&af1, 0).unwrap();
match e1 {
PcrEvent::Sample { jitter_27mhz, .. } => assert_eq!(jitter_27mhz, 0),
_ => panic!("expected sample"),
}
let e2 = t.observe(&af2, 1880).unwrap();
match e2 {
PcrEvent::Sample {
jitter_27mhz,
bitrate_bps,
..
} => {
assert_eq!(jitter_27mhz, 0);
assert_eq!(bitrate_bps, 1_000_000);
}
_ => panic!("expected sample"),
}
let e3 = t.observe(&af3, 1880 * 2).unwrap();
match e3 {
PcrEvent::Sample {
jitter_27mhz,
bitrate_bps,
..
} => {
assert_eq!(jitter_27mhz, 0);
assert_eq!(bitrate_bps, 1_000_000);
}
_ => panic!("expected sample"),
}
}
#[test]
fn pcr_tracker_detects_signalled_discontinuity() {
let mut t = PcrTracker::new();
let pcr1 = Pcr::from_ticks_27mhz(0);
let pcr2 = Pcr::from_ticks_27mhz(216 * 1880);
let pcr3 = Pcr::from_ticks_27mhz(216 * 1880 * 2);
t.observe(&af_with_pcr(pcr1.base_90khz(), pcr1.extension(), false), 0)
.unwrap();
t.observe(
&af_with_pcr(pcr2.base_90khz(), pcr2.extension(), false),
1880,
)
.unwrap();
let e3 = t
.observe(
&af_with_pcr(pcr3.base_90khz(), pcr3.extension(), true),
1880 * 2,
)
.unwrap();
match e3 {
PcrEvent::Discontinuity { reason, .. } => {
assert_eq!(reason, DiscontinuityReason::Signalled);
}
_ => panic!("expected discontinuity, got {e3:?}"),
}
assert_eq!(t.last_pcr(), Some(pcr3));
}
#[test]
fn pcr_tracker_indicator_on_packet_before_new_pcr() {
let mut t = PcrTracker::new();
let pcr1 = Pcr::from_ticks_27mhz(0);
let pcr2 = Pcr::from_ticks_27mhz(1000);
t.observe(&af_with_pcr(pcr1.base_90khz(), pcr1.extension(), false), 0)
.unwrap();
assert!(t.observe(&af_no_pcr(true), 188).is_none());
let e2 = t
.observe(
&af_with_pcr(pcr2.base_90khz(), pcr2.extension(), false),
376,
)
.unwrap();
match e2 {
PcrEvent::Discontinuity { reason, .. } => {
assert_eq!(reason, DiscontinuityReason::Signalled);
}
_ => panic!("expected discontinuity"),
}
}
#[test]
fn pcr_tracker_detects_unsignalled_jump() {
let mut t = PcrTracker::new();
let pcr1 = Pcr::from_ticks_27mhz(0);
let pcr2 = Pcr::from_ticks_27mhz(216 * 1880);
let bad = Pcr::from_ticks_27mhz(216 * 1880 * 2 + 27_000_000);
t.observe(&af_with_pcr(pcr1.base_90khz(), pcr1.extension(), false), 0)
.unwrap();
t.observe(
&af_with_pcr(pcr2.base_90khz(), pcr2.extension(), false),
1880,
)
.unwrap();
let e3 = t
.observe(
&af_with_pcr(bad.base_90khz(), bad.extension(), false),
1880 * 2,
)
.unwrap();
match e3 {
PcrEvent::Discontinuity {
reason: DiscontinuityReason::JumpExceedsTolerance { .. },
..
} => {}
_ => panic!("expected jump-exceeds-tolerance, got {e3:?}"),
}
}
#[test]
fn continuity_tracker_classifies_typical_flow() {
let mut t = ContinuityTracker::new();
assert_eq!(t.observe(5, true, false), ContinuityEvent::Continuous);
assert_eq!(t.observe(6, true, false), ContinuityEvent::Continuous);
assert_eq!(t.observe(6, true, false), ContinuityEvent::Duplicate);
assert_eq!(
t.observe(9, true, false),
ContinuityEvent::Dropped { gap: 2 }
);
}
#[test]
fn continuity_tracker_wraps_at_15() {
let mut t = ContinuityTracker::new();
t.observe(14, true, false);
assert_eq!(t.observe(15, true, false), ContinuityEvent::Continuous);
assert_eq!(t.observe(0, true, false), ContinuityEvent::Continuous);
assert_eq!(t.observe(1, true, false), ContinuityEvent::Continuous);
}
#[test]
fn continuity_tracker_no_payload_does_not_advance() {
let mut t = ContinuityTracker::new();
t.observe(3, true, false);
assert_eq!(t.observe(3, false, false), ContinuityEvent::NoPayload);
assert_eq!(t.observe(4, true, false), ContinuityEvent::Continuous);
}
#[test]
fn continuity_tracker_discontinuity_indicator_resets() {
let mut t = ContinuityTracker::new();
t.observe(3, true, false);
assert_eq!(t.observe(10, true, true), ContinuityEvent::Discontinuity);
assert_eq!(t.observe(11, true, false), ContinuityEvent::Continuous);
}
}