use std::collections::{HashMap, VecDeque};
use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, RwLock};
use std::time::{Duration, SystemTime};
use broadcast_common::Timestamp;
use bytes::Bytes;
use ll_hls_runtime::server::LlHlsOrigin;
use media_plane::trunk::{
PartEntry, SegmentCursor, SegmentCursorItem, SegmentEntry, SegmentWriter, TrunkConfig,
};
use media_plane::{ProgramId, Trunk};
use transmux::TrackSpec;
const DEFAULT_TIMED_CAPACITY: usize = 64;
const DEFAULT_SPARSE_CAPACITY: usize = 16;
const DEFAULT_EVENT_CAPACITY: usize = 64;
const DEFAULT_PART_CAPACITY: usize = 64;
const DEFAULT_ROUTE_NAME: &str = "unknown";
fn nz(n: usize) -> NonZeroUsize {
NonZeroUsize::new(n).expect("route.rs capacity constants are all non-zero")
}
pub(crate) const SPTS_PROGRAM_ID: ProgramId = ProgramId(0);
#[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",
}
}
}
broadcast_common::impl_spec_display!(HealthState);
#[non_exhaustive]
pub(crate) enum ProgramResolution {
Found(Arc<ProgramServing>),
NotYetAnnounced,
NotFound,
}
#[derive(Debug, Clone, Copy)]
pub struct DashWindowSegment {
pub segment_seq: u32,
pub duration_secs: f64,
}
pub(crate) struct DashState {
cursor: Mutex<SegmentCursor>,
window: Mutex<VecDeque<DashWindowSegment>>,
capacity: usize,
track_specs: Mutex<Vec<TrackSpec>>,
}
impl DashState {
fn new(trunk: &Arc<Trunk>, window_segments: NonZeroUsize) -> Self {
DashState {
cursor: Mutex::new(trunk.subscribe_segments()),
window: Mutex::new(VecDeque::new()),
capacity: window_segments.get(),
track_specs: Mutex::new(Vec::new()),
}
}
fn set_track_specs(&self, specs: Vec<TrackSpec>) {
*self
.track_specs
.lock()
.expect("DashState track_specs lock poisoned") = specs;
}
fn track_specs(&self) -> Vec<TrackSpec> {
self.track_specs
.lock()
.expect("DashState track_specs lock poisoned")
.clone()
}
fn drain(&self) {
let mut cursor = self
.cursor
.lock()
.expect("DashState segment cursor lock poisoned");
let mut window = self.window.lock().expect("DashState window lock poisoned");
while let Some(item) = cursor.poll() {
if let SegmentCursorItem::Segment(entry) = item {
if window.len() == self.capacity {
window.pop_front();
}
window.push_back(DashWindowSegment {
segment_seq: entry.sequence_number,
duration_secs: entry.duration.as_secs_f64(),
});
}
}
}
fn window_segments(&self) -> Vec<DashWindowSegment> {
self.drain();
self.window
.lock()
.expect("DashState window lock poisoned")
.iter()
.copied()
.collect()
}
}
pub(crate) struct ProgramServing {
trunk: Arc<Trunk>,
ll_hls: Arc<LlHlsOrigin>,
dash: Arc<DashState>,
segment_writer: Mutex<Option<SegmentWriter>>,
next_timeline_ns: AtomicU64,
}
impl ProgramServing {
fn new(
trunk: Arc<Trunk>,
target_duration_secs: f64,
part_target_ms: u32,
window_segments: NonZeroUsize,
) -> Arc<Self> {
let ll_hls = Arc::new(LlHlsOrigin::new(
Arc::clone(&trunk),
target_duration_secs,
part_target_ms,
window_segments,
));
let dash = Arc::new(DashState::new(&trunk, window_segments));
Arc::new(ProgramServing {
trunk,
ll_hls,
dash,
segment_writer: Mutex::new(None),
next_timeline_ns: AtomicU64::new(0),
})
}
pub(crate) fn trunk(&self) -> Arc<Trunk> {
Arc::clone(&self.trunk)
}
pub(crate) fn ll_hls(&self) -> Arc<LlHlsOrigin> {
Arc::clone(&self.ll_hls)
}
fn set_init(&self, bytes: impl Into<Bytes>) {
self.ll_hls.set_init(bytes);
}
fn init_bytes(&self) -> Option<Bytes> {
self.ll_hls.init_bytes()
}
fn set_track_specs(&self, specs: Vec<TrackSpec>) {
self.dash.set_track_specs(specs);
}
fn track_specs(&self) -> Vec<TrackSpec> {
self.dash.track_specs()
}
fn window_segments(&self) -> Vec<DashWindowSegment> {
self.dash.window_segments()
}
fn with_segment_writer<R>(&self, f: impl FnOnce(&SegmentWriter) -> R) -> Option<R> {
let mut guard = self
.segment_writer
.lock()
.expect("ProgramServing segment_writer lock poisoned");
if guard.is_none() {
*guard = self.trunk.segment_writer();
}
guard.as_ref().map(f)
}
fn add_part(&self, info: transmux::ll_hls::PartInfo) {
let published = self.with_segment_writer(|writer| {
writer.publish_part(PartEntry::new(
info.bytes,
info.segment_seq,
info.part_index,
Duration::from_secs_f64(info.duration),
info.independent,
));
});
if published.is_none() {
tracing::warn!(
"RouteHandle::add_part: this program's Trunk segment writer is unavailable \
(already taken by a real ProgramSegmenter?)"
);
}
}
fn add_segment(&self, info: transmux::ll_hls::SegmentInfo) {
let duration = Duration::from_secs_f64(info.duration);
let start_ns = self.next_timeline_ns.fetch_add(
u64::try_from(duration.as_nanos()).unwrap_or(u64::MAX),
Ordering::SeqCst,
);
let published = self.with_segment_writer(|writer| {
writer.publish_segment(SegmentEntry::new(
info.bytes,
info.segment_seq,
duration,
Timestamp::from_nanos(start_ns),
transmux::SegmentMeta {
discontinuous: false,
},
));
});
if published.is_none() {
tracing::warn!(
"RouteHandle::add_segment: this program's Trunk segment writer is unavailable \
(already taken by a real ProgramSegmenter?)"
);
}
}
pub(crate) fn latest_progress(&self) -> (u32, usize) {
let last_closed = self.trunk.last_closed_segment().unwrap_or(0);
let candidate = last_closed + 1;
let parts = self.trunk.parts_in_segment(candidate);
if parts.is_empty() {
(last_closed, 0)
} else {
(candidate, parts.len())
}
}
}
pub struct RouteHandle {
health: Mutex<HealthState>,
target_duration_secs: f64,
part_target_ms: u32,
window_segments_cap: NonZeroUsize,
created_at: SystemTime,
programs: RwLock<HashMap<ProgramId, Arc<ProgramServing>>>,
name: String,
}
impl RouteHandle {
pub fn new(target_duration_secs: f64, part_target_ms: u32, window_segments: usize) -> Self {
let window_segments = NonZeroUsize::new(window_segments).unwrap_or(NonZeroUsize::MIN);
RouteHandle {
health: Mutex::new(HealthState::Connecting),
target_duration_secs,
part_target_ms,
window_segments_cap: window_segments,
created_at: SystemTime::now(),
programs: RwLock::new(HashMap::new()),
name: DEFAULT_ROUTE_NAME.to_string(),
}
}
pub fn with_name(mut self, name: impl Into<String>) -> Self {
self.name = name.into();
self
}
pub fn name(&self) -> &str {
&self.name
}
pub fn publish_program(&self, program: ProgramId, trunk: Arc<Trunk>) {
let mut programs = self
.programs
.write()
.expect("RouteHandle::programs lock poisoned");
let already_bound = programs
.get(&program)
.is_some_and(|existing| Arc::ptr_eq(&existing.trunk, &trunk));
if already_bound {
return;
}
let serving = ProgramServing::new(
trunk,
self.target_duration_secs,
self.part_target_ms,
self.window_segments_cap,
);
programs.insert(program, serving);
}
pub fn publish_new_program(&self, program: ProgramId) -> Arc<Trunk> {
let trunk_config = TrunkConfig::new(
nz(DEFAULT_TIMED_CAPACITY),
nz(DEFAULT_SPARSE_CAPACITY),
self.window_segments_cap,
nz(DEFAULT_EVENT_CAPACITY),
nz(DEFAULT_PART_CAPACITY),
);
let trunk = Trunk::new(trunk_config);
self.publish_program(program, Arc::clone(&trunk));
trunk
}
pub(crate) fn serving(&self, program: ProgramId) -> Option<Arc<ProgramServing>> {
self.programs
.read()
.expect("RouteHandle::programs lock poisoned")
.get(&program)
.cloned()
}
pub(crate) fn resolve_program(&self, program: ProgramId) -> ProgramResolution {
let programs = self
.programs
.read()
.expect("RouteHandle::programs lock poisoned");
match programs.get(&program) {
Some(serving) => ProgramResolution::Found(Arc::clone(serving)),
None if programs.is_empty() => ProgramResolution::NotYetAnnounced,
None => ProgramResolution::NotFound,
}
}
#[cfg(test)]
pub(crate) fn ll_hls(&self, program: ProgramId) -> Option<Arc<LlHlsOrigin>> {
self.serving(program).map(|s| s.ll_hls())
}
pub fn set_init(&self, program: ProgramId, bytes: impl Into<Bytes>) {
match self.serving(program) {
Some(serving) => serving.set_init(bytes),
None => tracing::warn!(?program, "RouteHandle::set_init: program not published yet"),
}
}
pub fn init_bytes(&self, program: ProgramId) -> Option<Bytes> {
self.serving(program)?.init_bytes()
}
pub fn set_track_specs(&self, program: ProgramId, specs: Vec<TrackSpec>) {
match self.serving(program) {
Some(serving) => serving.set_track_specs(specs),
None => {
tracing::warn!(
?program,
"RouteHandle::set_track_specs: program not published yet"
)
}
}
}
pub fn track_specs(&self, program: ProgramId) -> Vec<TrackSpec> {
self.serving(program)
.map(|s| s.track_specs())
.unwrap_or_default()
}
pub fn add_part(&self, program: ProgramId, info: transmux::ll_hls::PartInfo) {
match self.serving(program) {
Some(serving) => serving.add_part(info),
None => tracing::warn!(?program, "RouteHandle::add_part: program not published yet"),
}
}
pub fn add_segment(&self, program: ProgramId, info: transmux::ll_hls::SegmentInfo) {
match self.serving(program) {
Some(serving) => serving.add_segment(info),
None => tracing::warn!(
?program,
"RouteHandle::add_segment: program not published yet"
),
}
}
pub fn window_segments(&self, program: ProgramId) -> Vec<DashWindowSegment> {
self.serving(program)
.map(|s| s.window_segments())
.unwrap_or_default()
}
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 health(&self) -> HealthState {
*self.health.lock().unwrap()
}
pub fn set_health(&self, state: HealthState) {
*self.health.lock().unwrap() = state;
}
}
#[cfg(test)]
mod program_registry_tests {
use super::*;
use transmux::CodecConfig;
fn test_trunk() -> Arc<Trunk> {
Trunk::new(TrunkConfig::new(nz(1), nz(1), nz(1), nz(1), nz(1)))
}
#[test]
fn publish_then_resolve_returns_same_trunk() {
let route = RouteHandle::new(4.0, 500, 4);
let trunk = test_trunk();
route.publish_program(ProgramId(1), Arc::clone(&trunk));
match route.resolve_program(ProgramId(1)) {
ProgramResolution::Found(resolved) => {
assert_eq!(Arc::as_ptr(&resolved.trunk()), Arc::as_ptr(&trunk));
}
_ => panic!(
"expected ProgramResolution::Found for a just-published program, got a variant \
that is not Found"
),
}
}
#[test]
fn resolve_before_any_program_is_not_yet_announced() {
let route = RouteHandle::new(4.0, 500, 4);
let result = route.resolve_program(ProgramId(1));
assert!(
matches!(result, ProgramResolution::NotYetAnnounced),
"expected NotYetAnnounced with an empty registry"
);
}
#[test]
fn resolve_unknown_program_among_others_is_not_found() {
let route = RouteHandle::new(4.0, 500, 4);
route.publish_program(ProgramId(1), test_trunk());
let result = route.resolve_program(ProgramId(2));
assert!(
matches!(result, ProgramResolution::NotFound),
"expected NotFound: program 2 was never published, but program 1 was"
);
}
#[test]
fn two_programs_resolve_independently_to_different_trunks() {
let route = RouteHandle::new(4.0, 500, 4);
let trunk_a = test_trunk();
let trunk_b = test_trunk();
route.publish_program(ProgramId(1), Arc::clone(&trunk_a));
route.publish_program(ProgramId(2), Arc::clone(&trunk_b));
let first = match route.resolve_program(ProgramId(1)) {
ProgramResolution::Found(s) => s,
_ => panic!("expected ProgramResolution::Found for program 1"),
};
let second = match route.resolve_program(ProgramId(2)) {
ProgramResolution::Found(s) => s,
_ => panic!("expected ProgramResolution::Found for program 2"),
};
assert_eq!(Arc::as_ptr(&first.trunk()), Arc::as_ptr(&trunk_a));
assert_eq!(Arc::as_ptr(&second.trunk()), Arc::as_ptr(&trunk_b));
assert_ne!(Arc::as_ptr(&first.trunk()), Arc::as_ptr(&second.trunk()));
}
fn video_spec(track_id: u32) -> TrackSpec {
TrackSpec::new(
track_id,
90_000,
CodecConfig::Vp8 {
width: 1280,
height: 720,
},
)
}
fn seg(seq: u32, duration: f64) -> transmux::ll_hls::SegmentInfo {
transmux::ll_hls::SegmentInfo {
bytes: vec![seq as u8; 8],
duration,
segment_seq: seq,
part_count: 1,
}
}
#[test]
fn two_programs_serve_distinct_media() {
let route = RouteHandle::new(1.0, 500, 8);
let program_a = ProgramId(1);
let program_b = ProgramId(2);
route.publish_new_program(program_a);
route.publish_new_program(program_b);
route.set_init(program_a, &b"init-1"[..]);
route.set_init(program_b, &b"init-2"[..]);
route.set_track_specs(program_a, vec![video_spec(1)]);
route.set_track_specs(program_b, vec![video_spec(2)]);
route.add_segment(program_a, seg(10, 2.0));
route.add_segment(program_b, seg(20, 2.0));
route.add_segment(program_b, seg(21, 2.0));
assert_eq!(
route.init_bytes(program_a),
Some(Bytes::from_static(b"init-1")),
"program 1 must serve its own init bytes"
);
assert_eq!(
route.init_bytes(program_b),
Some(Bytes::from_static(b"init-2")),
"program 2 must serve its own, DISTINCT init bytes"
);
let window_a = route.window_segments(program_a);
let window_b = route.window_segments(program_b);
assert_eq!(
window_a.iter().map(|s| s.segment_seq).collect::<Vec<_>>(),
vec![10],
"program 1 must show only its own segment"
);
assert_eq!(
window_b.iter().map(|s| s.segment_seq).collect::<Vec<_>>(),
vec![20, 21],
"program 2 must show only its own (different-count) segments"
);
assert_eq!(
route
.track_specs(program_a)
.iter()
.map(|s| s.track_id)
.collect::<Vec<_>>(),
vec![1],
"program 1's own track spec"
);
assert_eq!(
route
.track_specs(program_b)
.iter()
.map(|s| s.track_id)
.collect::<Vec<_>>(),
vec![2],
"program 2's own, DIFFERENT track spec"
);
}
#[test]
fn unannounced_program_is_not_yet_announced_not_not_found() {
let route = RouteHandle::new(4.0, 500, 4);
assert!(matches!(
route.resolve_program(ProgramId(7)),
ProgramResolution::NotYetAnnounced
));
}
#[test]
fn name_defaults_then_reflects_with_name() {
let route = RouteHandle::new(4.0, 500, 4);
assert_eq!(route.name(), DEFAULT_ROUTE_NAME);
let route = route.with_name("cam1");
assert_eq!(route.name(), "cam1");
}
}