use std::collections::VecDeque;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::SystemTime;
use event_listener::{Event, EventListener};
use transmux::ll_hls::{PartInfo, SegmentInfo};
use transmux::pipeline::TrackSpec;
const MIN_MAX_LIVE_PARTS: usize = 8;
const MAX_LIVE_PARTS_SAFETY_MARGIN: usize = 4;
pub(crate) fn compute_max_live_parts(target_duration_secs: f64, part_target_ms: u32) -> usize {
let part_target_secs = f64::from(part_target_ms) / 1000.0;
let nominal_parts = if part_target_secs > 0.0 {
(target_duration_secs / part_target_secs).ceil() as usize
} else {
0
};
(nominal_parts + MAX_LIVE_PARTS_SAFETY_MARGIN).max(MIN_MAX_LIVE_PARTS)
}
struct Inner {
init: Option<Vec<u8>>,
segments: VecDeque<SegmentInfo>,
live_parts: Vec<PartInfo>,
recent_parts: VecDeque<PartInfo>,
window_segments: usize,
health: HealthState,
max_segment_duration: f64,
track_specs: Vec<TrackSpec>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum HealthState {
Connecting,
Live,
Reconnecting,
Failed,
}
impl HealthState {
pub fn name(&self) -> &'static str {
match self {
HealthState::Connecting => "connecting",
HealthState::Live => "live",
HealthState::Reconnecting => "reconnecting",
HealthState::Failed => "failed",
}
}
}
impl std::fmt::Display for HealthState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.name())
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct SegmentWindowEntry {
pub segment_seq: u32,
pub duration_secs: f64,
}
pub struct MediaStore {
inner: Mutex<Inner>,
target_duration_secs: f64,
part_target_ms: u32,
max_live_parts: usize,
progress_version: AtomicU64,
progress_event: Event,
created_at: SystemTime,
}
impl MediaStore {
pub fn new(target_duration_secs: f64, part_target_ms: u32, window_segments: usize) -> Self {
MediaStore {
inner: Mutex::new(Inner {
init: None,
segments: VecDeque::new(),
live_parts: Vec::new(),
recent_parts: VecDeque::new(),
window_segments,
health: HealthState::Connecting,
max_segment_duration: 0.0,
track_specs: Vec::new(),
}),
target_duration_secs,
part_target_ms,
max_live_parts: compute_max_live_parts(target_duration_secs, part_target_ms),
progress_version: AtomicU64::new(0),
progress_event: Event::new(),
created_at: SystemTime::now(),
}
}
fn bump(&self) {
self.progress_version.fetch_add(1, Ordering::SeqCst);
self.progress_event.notify(usize::MAX);
}
pub fn set_init(&self, bytes: Vec<u8>) {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.init = Some(bytes);
self.bump();
}
pub fn add_part(&self, part: PartInfo) {
let mut g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
g.live_parts.push(part);
while g.live_parts.len() > self.max_live_parts {
g.live_parts.remove(0);
}
drop(g);
self.bump();
}
#[cfg(test)]
pub(crate) fn live_part_count(&self) -> usize {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.live_parts
.len()
}
pub fn add_segment(&self, seg: SegmentInfo) {
let mut g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let seq = seg.segment_seq;
g.max_segment_duration = g.max_segment_duration.max(seg.duration);
let (closed, still_live): (Vec<PartInfo>, Vec<PartInfo>) =
core::mem::take(&mut g.live_parts)
.into_iter()
.partition(|p| p.segment_seq <= seq);
g.live_parts = still_live;
for p in closed {
g.recent_parts.push_back(p);
}
while g.recent_parts.len() > self.max_live_parts {
g.recent_parts.pop_front();
}
g.segments.push_back(seg);
while g.segments.len() > g.window_segments {
g.segments.pop_front();
}
drop(g);
self.bump();
}
pub fn init_bytes(&self) -> Option<Vec<u8>> {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.init
.clone()
}
pub(crate) fn segment_bytes(&self, seq: u32) -> Option<Vec<u8>> {
let g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
g.segments
.iter()
.find(|s| s.segment_seq == seq)
.map(|s| s.bytes.clone())
}
pub(crate) fn part_bytes(&self, seq: u32, part_index: u32) -> Option<Vec<u8>> {
let g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let matches = |p: &&PartInfo| p.segment_seq == seq && p.part_index == part_index;
g.live_parts
.iter()
.find(matches)
.or_else(|| g.recent_parts.iter().find(matches))
.map(|p| p.bytes.clone())
}
pub fn latest_progress(&self) -> (u32, u32) {
let g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let last_closed_seg = g.segments.back().map(|s| s.segment_seq).unwrap_or(0);
let in_progress_seg = g
.live_parts
.last()
.map(|p| p.segment_seq)
.unwrap_or(last_closed_seg);
let part_count = g
.live_parts
.iter()
.filter(|p| p.segment_seq == in_progress_seg)
.count() as u32;
(in_progress_seg, part_count)
}
pub fn progress_version(&self) -> u64 {
self.progress_version.load(Ordering::SeqCst)
}
pub fn listen(&self) -> EventListener {
self.progress_event.listen()
}
pub fn set_health(&self, state: HealthState) {
let mut g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if g.health != state {
g.health = state;
drop(g);
self.bump();
}
}
pub fn health(&self) -> HealthState {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.health
}
pub fn target_duration_secs(&self) -> f64 {
self.target_duration_secs
}
pub fn part_target_ms(&self) -> u32 {
self.part_target_ms
}
pub fn created_at(&self) -> SystemTime {
self.created_at
}
pub fn set_track_specs(&self, specs: Vec<TrackSpec>) {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.track_specs = specs;
}
pub fn track_specs(&self) -> Vec<TrackSpec> {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.track_specs
.clone()
}
pub fn window_segments(&self) -> Vec<SegmentWindowEntry> {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.segments
.iter()
.map(|s| SegmentWindowEntry {
segment_seq: s.segment_seq,
duration_secs: s.duration,
})
.collect()
}
pub(crate) fn max_segment_duration(&self) -> f64 {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.max_segment_duration
}
pub(crate) fn last_closed_segment_seq(&self) -> u32 {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.segments
.back()
.map(|s| s.segment_seq)
.unwrap_or(0)
}
pub(crate) fn with_segments_and_parts<R>(
&self,
f: impl FnOnce(&VecDeque<SegmentInfo>, &[PartInfo]) -> R,
) -> R {
let g = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
f(&g.segments, &g.live_parts)
}
}
#[cfg(test)]
mod tests {
use super::*;
use event_listener::Listener;
fn seg(seq: u32, parts: u32) -> SegmentInfo {
SegmentInfo {
bytes: vec![seq as u8; 8],
duration: 4.0,
segment_seq: seq,
part_count: parts,
}
}
fn part(seq: u32, idx: u32) -> PartInfo {
PartInfo {
bytes: vec![idx as u8; 4],
duration: 0.5,
independent: idx == 0,
segment_seq: seq,
part_index: idx,
}
}
#[test]
fn window_evicts_oldest_and_serves_bytes() {
let s = MediaStore::new(4.0, 500, 2);
s.set_init(vec![0xAA; 10]);
s.add_segment(seg(1, 8));
s.add_segment(seg(2, 8));
s.add_segment(seg(3, 8)); assert!(s.segment_bytes(1).is_none(), "seq 1 evicted");
assert!(s.segment_bytes(2).is_some());
assert!(s.segment_bytes(3).is_some());
assert_eq!(s.init_bytes().unwrap(), vec![0xAA; 10]);
}
#[test]
fn recent_parts_bounded_across_many_closes() {
let s = MediaStore::new(4.0, 500, 4);
s.set_init(vec![0; 4]);
let cap = compute_max_live_parts(4.0, 500);
for seq in 1..=20u32 {
for idx in 0..4u32 {
s.add_part(part(seq, idx));
}
s.add_segment(seg(seq, 4));
}
assert!(s.part_bytes(1, 0).is_none(), "old closed part evicted");
assert!(
s.part_bytes(20, 3).is_some(),
"most-recent closed part retained (within the {cap}-part bound)"
);
}
#[test]
fn progress_version_bumps_on_new_data() {
let s = MediaStore::new(4.0, 500, 4);
let before = s.progress_version();
s.add_part(part(1, 0));
assert_ne!(s.progress_version(), before, "progress version changed");
}
#[test]
fn listen_wakes_on_new_data() {
let s = MediaStore::new(4.0, 500, 4);
let listener = s.listen();
s.add_part(part(1, 0));
assert!(
listener
.wait_deadline(std::time::Instant::now() + std::time::Duration::from_secs(2))
.is_some(),
"listener must wake within 2s of add_part"
);
}
#[test]
fn health_defaults_to_connecting() {
let s = MediaStore::new(4.0, 500, 4);
assert_eq!(s.health(), HealthState::Connecting);
}
#[test]
fn set_health_updates_and_bumps_progress_only_on_change() {
let s = MediaStore::new(4.0, 500, 4);
let before = s.progress_version();
s.set_health(HealthState::Connecting);
assert_eq!(
s.progress_version(),
before,
"unchanged state does not bump progress"
);
s.set_health(HealthState::Live);
assert_eq!(s.health(), HealthState::Live);
assert_ne!(
s.progress_version(),
before,
"state change bumps progress so blocked readers wake"
);
let mid = s.progress_version();
s.set_health(HealthState::Reconnecting);
assert_eq!(s.health(), HealthState::Reconnecting);
assert_ne!(s.progress_version(), mid);
}
#[test]
fn max_segment_duration_tracks_lifetime_max_not_just_current_window() {
let s = MediaStore::new(4.0, 500, 2);
s.set_init(vec![0; 4]);
assert_eq!(s.max_segment_duration(), 0.0, "nothing closed yet");
let mut over = seg(1, 8);
over.duration = 4.0;
s.add_segment(over);
assert_eq!(s.max_segment_duration(), 4.0);
let mut over = seg(2, 8);
over.duration = 7.5;
s.add_segment(over);
assert_eq!(s.max_segment_duration(), 7.5);
let mut small = seg(3, 8);
small.duration = 3.0;
s.add_segment(small); let mut small2 = seg(4, 8);
small2.duration = 3.0;
s.add_segment(small2); assert!(
s.segment_bytes(2).is_none(),
"seq 2 (the 7.5s segment) evicted from the window"
);
assert_eq!(
s.max_segment_duration(),
7.5,
"lifetime max must survive window eviction"
);
}
#[test]
fn last_closed_segment_seq_tracks_the_newest_close() {
let s = MediaStore::new(4.0, 500, 4);
assert_eq!(s.last_closed_segment_seq(), 0, "nothing closed yet");
s.add_segment(seg(1, 8));
assert_eq!(s.last_closed_segment_seq(), 1);
s.add_segment(seg(2, 8));
assert_eq!(s.last_closed_segment_seq(), 2);
}
#[test]
fn health_state_name_and_display_agree() {
for (state, label) in [
(HealthState::Connecting, "connecting"),
(HealthState::Live, "live"),
(HealthState::Reconnecting, "reconnecting"),
(HealthState::Failed, "failed"),
] {
assert_eq!(state.name(), label);
assert_eq!(state.to_string(), label);
}
}
#[test]
fn track_specs_round_trip_and_default_empty() {
use transmux::pipeline::CodecConfig;
let s = MediaStore::new(4.0, 500, 4);
assert!(
s.track_specs().is_empty(),
"no specs set yet -> empty, not a panic/placeholder"
);
let spec = TrackSpec::new(
1,
90_000,
CodecConfig::Vp8 {
width: 0,
height: 0,
},
);
s.set_track_specs(vec![spec.clone()]);
let got = s.track_specs();
assert_eq!(got.len(), 1);
assert_eq!(got[0].track_id, spec.track_id);
assert_eq!(got[0].timescale, spec.timescale);
}
#[test]
fn window_segments_reflects_closed_segments_oldest_first_and_evicts() {
let s = MediaStore::new(4.0, 500, 2);
assert!(s.window_segments().is_empty(), "nothing closed yet");
s.add_segment(seg(1, 4));
s.add_segment(seg(2, 4));
let window = s.window_segments();
assert_eq!(
window.iter().map(|e| e.segment_seq).collect::<Vec<_>>(),
vec![1, 2],
"oldest first"
);
assert_eq!(window[0].duration_secs, 4.0);
s.add_segment(seg(3, 4)); assert_eq!(
s.window_segments()
.iter()
.map(|e| e.segment_seq)
.collect::<Vec<_>>(),
vec![2, 3],
"eviction reflected in the snapshot"
);
}
#[test]
fn created_at_is_set_at_construction() {
let before = SystemTime::now();
let s = MediaStore::new(4.0, 500, 4);
let after = SystemTime::now();
assert!(s.created_at() >= before && s.created_at() <= after);
}
}