use std::collections::{HashMap, VecDeque};
use std::future::Future;
use std::num::NonZeroUsize;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::task::{Context, Poll};
use std::time::Duration;
use broadcast_common::stage::Timestamp;
use bytes::Bytes;
use event_listener::{Event, EventListener, Listener};
use timed_metadata::{MediaTime, PTS_HZ, TimeAnchor, TimedEvent};
use transmux::{Sample, SegmentMeta, TrackSpec};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum RetentionClass {
Timed,
Sparse,
}
#[derive(Debug, Clone, Copy)]
#[non_exhaustive]
pub struct TrunkConfig {
pub timed_capacity: NonZeroUsize,
pub sparse_capacity: NonZeroUsize,
pub segment_capacity: NonZeroUsize,
pub event_capacity: NonZeroUsize,
pub part_capacity: NonZeroUsize,
}
impl TrunkConfig {
pub fn new(
timed_capacity: NonZeroUsize,
sparse_capacity: NonZeroUsize,
segment_capacity: NonZeroUsize,
event_capacity: NonZeroUsize,
part_capacity: NonZeroUsize,
) -> Self {
TrunkConfig {
timed_capacity,
sparse_capacity,
segment_capacity,
event_capacity,
part_capacity,
}
}
}
struct ClassLog {
entries: VecDeque<(u32, Sample)>,
base: u64,
published: u64,
capacity: usize,
}
impl ClassLog {
fn new(capacity: usize) -> Self {
ClassLog {
entries: VecDeque::with_capacity(capacity),
base: 0,
published: 0,
capacity,
}
}
fn push(&mut self, track_id: u32, sample: Sample) {
if self.entries.len() == self.capacity {
self.entries.pop_front();
self.base += 1;
}
self.entries.push_back((track_id, sample));
self.published += 1;
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct SegmentEntry {
pub bytes: Bytes,
pub sequence_number: u32,
pub duration: Duration,
pub timeline_position: Timestamp,
pub meta: SegmentMeta,
}
impl SegmentEntry {
pub fn new(
bytes: impl Into<Bytes>,
sequence_number: u32,
duration: Duration,
timeline_position: Timestamp,
meta: SegmentMeta,
) -> Self {
SegmentEntry {
bytes: bytes.into(),
sequence_number,
duration,
timeline_position,
meta,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum ArchiveOverrun {
Gap,
StallIngest,
Terminate,
}
impl Default for ArchiveOverrun {
fn default() -> Self {
ArchiveOverrun::Gap
}
}
struct PinState {
consumed: u64,
policy: ArchiveOverrun,
terminated: bool,
}
struct SegmentLog {
entries: VecDeque<SegmentEntry>,
base: u64,
published: u64,
capacity: usize,
pins: HashMap<u64, PinState>,
next_pin_id: u64,
}
impl SegmentLog {
fn new(capacity: usize) -> Self {
SegmentLog {
entries: VecDeque::with_capacity(capacity),
base: 0,
published: 0,
capacity,
pins: HashMap::new(),
next_pin_id: 0,
}
}
fn push(&mut self, entry: SegmentEntry) {
if self.entries.len() == self.capacity {
self.entries.pop_front();
self.base += 1;
}
self.entries.push_back(entry);
self.published += 1;
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum EventAnchor {
Media(MediaTime),
Segment {
segment_number: u32,
delta: u64,
},
Utc {
utc_epoch_ms: i64,
},
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct EventEntry {
pub event: TimedEvent,
pub anchor: EventAnchor,
}
struct EventLog {
entries: VecDeque<EventEntry>,
base: u64,
published: u64,
capacity: usize,
segment_starts: VecDeque<(u32, MediaTime)>,
time_anchor: Option<TimeAnchor>,
}
impl EventLog {
fn new(capacity: usize) -> Self {
EventLog {
entries: VecDeque::with_capacity(capacity),
base: 0,
published: 0,
capacity,
segment_starts: VecDeque::with_capacity(capacity),
time_anchor: None,
}
}
fn try_resolve(&self, anchor: EventAnchor) -> EventAnchor {
match anchor {
EventAnchor::Segment {
segment_number,
delta,
} => self
.segment_starts
.iter()
.find(|(n, _)| *n == segment_number)
.map(|(_, start)| EventAnchor::Media(MediaTime(start.0.saturating_add(delta))))
.unwrap_or(anchor),
EventAnchor::Utc { utc_epoch_ms } => self
.time_anchor
.as_ref()
.map(|a| EventAnchor::Media(epoch_ms_to_media(a, utc_epoch_ms)))
.unwrap_or(anchor),
EventAnchor::Media(_) => anchor,
}
}
fn push(&mut self, event: TimedEvent, anchor: EventAnchor) {
let anchor = self.try_resolve(anchor);
if self.entries.len() == self.capacity {
self.entries.pop_front();
self.base += 1;
}
self.entries.push_back(EventEntry { event, anchor });
self.published += 1;
}
fn note_segment_start(&mut self, segment_number: u32, start: MediaTime) {
if self.segment_starts.len() == self.capacity {
self.segment_starts.pop_front();
}
self.segment_starts.push_back((segment_number, start));
for entry in &mut self.entries {
if let EventAnchor::Segment {
segment_number: n,
delta,
} = entry.anchor
&& n == segment_number
{
entry.anchor = EventAnchor::Media(MediaTime(start.0.saturating_add(delta)));
}
}
}
fn set_time_anchor(&mut self, anchor: TimeAnchor) {
self.time_anchor = Some(anchor);
for entry in &mut self.entries {
if let EventAnchor::Utc { utc_epoch_ms } = entry.anchor {
entry.anchor = EventAnchor::Media(epoch_ms_to_media(&anchor, utc_epoch_ms));
}
}
}
}
fn epoch_ms_to_media(anchor: &TimeAnchor, utc_epoch_ms: i64) -> MediaTime {
let delta_ms = i128::from(utc_epoch_ms) - i128::from(anchor.utc_epoch_ms);
let delta_ticks = delta_ms * i128::from(PTS_HZ) / 1000;
let media = i128::from(anchor.pts_90k) + delta_ticks;
MediaTime(media.clamp(0, i128::from(u64::MAX)) as u64)
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct PartEntry {
pub bytes: Bytes,
pub segment_number: u32,
pub part_index: u32,
pub duration: Duration,
pub independent: bool,
}
impl PartEntry {
pub fn new(
bytes: impl Into<Bytes>,
segment_number: u32,
part_index: u32,
duration: Duration,
independent: bool,
) -> Self {
PartEntry {
bytes: bytes.into(),
segment_number,
part_index,
duration,
independent,
}
}
}
struct PartLog {
entries: VecDeque<PartEntry>,
base: u64,
published: u64,
capacity: usize,
}
impl PartLog {
fn new(capacity: usize) -> Self {
PartLog {
entries: VecDeque::with_capacity(capacity),
base: 0,
published: 0,
capacity,
}
}
fn push(&mut self, entry: PartEntry) {
if self.entries.len() == self.capacity {
self.entries.pop_front();
self.base += 1;
}
self.entries.push_back(entry);
self.published += 1;
}
}
struct TrunkState {
timed: ClassLog,
sparse: ClassLog,
segments: SegmentLog,
events: EventLog,
parts: PartLog,
tracks: Arc<[TrackSpec]>,
track_generation: u64,
}
pub struct Trunk {
state: Mutex<TrunkState>,
segment_pin_released: Condvar,
writer_taken: AtomicBool,
segment_writer_taken: AtomicBool,
progress: Event,
waiter_count: AtomicUsize,
part_waiter_cap: usize,
}
impl Trunk {
pub fn new(config: TrunkConfig) -> Arc<Trunk> {
Arc::new(Trunk {
state: Mutex::new(TrunkState {
timed: ClassLog::new(config.timed_capacity.get()),
sparse: ClassLog::new(config.sparse_capacity.get()),
segments: SegmentLog::new(config.segment_capacity.get()),
events: EventLog::new(config.event_capacity.get()),
parts: PartLog::new(config.part_capacity.get()),
tracks: Arc::from(Vec::new()),
track_generation: 0,
}),
segment_pin_released: Condvar::new(),
writer_taken: AtomicBool::new(false),
segment_writer_taken: AtomicBool::new(false),
progress: Event::new(),
waiter_count: AtomicUsize::new(0),
part_waiter_cap: config.part_capacity.get(),
})
}
pub fn writer(self: &Arc<Self>) -> Option<TrunkWriter> {
self.writer_taken
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.ok()
.map(|_| TrunkWriter {
trunk: Arc::clone(self),
})
}
pub fn segment_writer(self: &Arc<Self>) -> Option<SegmentWriter> {
self.segment_writer_taken
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.ok()
.map(|_| SegmentWriter {
trunk: Arc::clone(self),
})
}
pub fn subscribe(self: &Arc<Self>) -> SampleCursor {
let state = self.state.lock().expect("Trunk state lock poisoned");
SampleCursor {
trunk: Arc::clone(self),
timed_consumed: state.timed.published,
sparse_consumed: state.sparse.published,
}
}
pub fn subscribe_from_backlog(self: &Arc<Self>) -> SampleCursor {
let state = self.state.lock().expect("Trunk state lock poisoned");
SampleCursor {
trunk: Arc::clone(self),
timed_consumed: state.timed.base,
sparse_consumed: state.sparse.base,
}
}
pub fn timed_len(&self) -> usize {
self.state
.lock()
.expect("Trunk state lock poisoned")
.timed
.entries
.len()
}
pub fn sparse_len(&self) -> usize {
self.state
.lock()
.expect("Trunk state lock poisoned")
.sparse
.entries
.len()
}
pub fn subscribe_segments(self: &Arc<Self>) -> SegmentCursor {
let state = self.state.lock().expect("Trunk state lock poisoned");
SegmentCursor {
trunk: Arc::clone(self),
consumed: state.segments.published,
pin_id: None,
done: false,
}
}
pub fn pin_segments(self: &Arc<Self>, on_overrun: ArchiveOverrun) -> SegmentCursor {
let mut state = self.state.lock().expect("Trunk state lock poisoned");
let pin_id = state.segments.next_pin_id;
state.segments.next_pin_id += 1;
let consumed = state.segments.published;
state.segments.pins.insert(
pin_id,
PinState {
consumed,
policy: on_overrun,
terminated: false,
},
);
SegmentCursor {
trunk: Arc::clone(self),
consumed: 0,
pin_id: Some(pin_id),
done: false,
}
}
pub fn segment_len(&self) -> usize {
self.state
.lock()
.expect("Trunk state lock poisoned")
.segments
.entries
.len()
}
pub fn subscribe_events(self: &Arc<Self>) -> EventCursor {
let state = self.state.lock().expect("Trunk state lock poisoned");
EventCursor {
trunk: Arc::clone(self),
consumed: state.events.published,
}
}
pub fn events_between(&self, from: MediaTime, to: MediaTime) -> Vec<EventEntry> {
let state = self.state.lock().expect("Trunk state lock poisoned");
state
.events
.entries
.iter()
.filter(|e| matches!(e.anchor, EventAnchor::Media(t) if t.0 >= from.0 && t.0 < to.0))
.cloned()
.collect()
}
pub fn events_in_segment(&self, segment_number: u32) -> Vec<EventEntry> {
let state = self.state.lock().expect("Trunk state lock poisoned");
let log = &state.events;
let Some(&(_, start)) = log
.segment_starts
.iter()
.find(|(n, _)| *n == segment_number)
else {
return Vec::new();
};
let end = log
.segment_starts
.iter()
.find(|(n, _)| *n == segment_number + 1)
.map(|&(_, s)| s.0);
log.entries
.iter()
.filter(|e| match e.anchor {
EventAnchor::Media(t) => t.0 >= start.0 && end.map(|e2| t.0 < e2).unwrap_or(true),
_ => false,
})
.cloned()
.collect()
}
pub fn event_len(&self) -> usize {
self.state
.lock()
.expect("Trunk state lock poisoned")
.events
.entries
.len()
}
pub fn part_bytes(&self, segment_number: u32, part_index: u32) -> Option<Bytes> {
let state = self.state.lock().expect("Trunk state lock poisoned");
state
.parts
.entries
.iter()
.find(|p| p.segment_number == segment_number && p.part_index == part_index)
.map(|p| p.bytes.clone())
}
pub fn parts_in_segment(&self, segment_number: u32) -> Vec<PartEntry> {
let state = self.state.lock().expect("Trunk state lock poisoned");
state
.parts
.entries
.iter()
.filter(|p| p.segment_number == segment_number)
.cloned()
.collect()
}
pub fn part_len(&self) -> usize {
self.state
.lock()
.expect("Trunk state lock poisoned")
.parts
.entries
.len()
}
pub fn waiter_count(&self) -> usize {
self.waiter_count.load(Ordering::Acquire)
}
pub fn last_closed_segment(&self) -> Option<u32> {
self.state
.lock()
.expect("Trunk state lock poisoned")
.segments
.entries
.back()
.map(|e| e.sequence_number)
}
pub fn tracks(&self) -> Arc<[TrackSpec]> {
Arc::clone(&self.state.lock().expect("Trunk state lock poisoned").tracks)
}
pub fn track_generation(&self) -> u64 {
self.state
.lock()
.expect("Trunk state lock poisoned")
.track_generation
}
pub fn listen(self: &Arc<Self>) -> Option<ProgressListener> {
loop {
let current = self.waiter_count.load(Ordering::Acquire);
if current >= self.part_waiter_cap {
return None;
}
if self
.waiter_count
.compare_exchange(current, current + 1, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
break;
}
}
Some(ProgressListener {
_slot: WaiterSlot(Arc::clone(self)),
listener: self.progress.listen(),
})
}
}
struct WaiterSlot(Arc<Trunk>);
impl Drop for WaiterSlot {
fn drop(&mut self) {
self.0.waiter_count.fetch_sub(1, Ordering::AcqRel);
}
}
pub struct ProgressListener {
_slot: WaiterSlot,
listener: EventListener,
}
impl ProgressListener {
pub fn wait_deadline(self, deadline: std::time::Instant) -> bool {
let ProgressListener { _slot, listener } = self;
listener.wait_deadline(deadline).is_some()
}
}
impl Future for ProgressListener {
type Output = ();
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let this = self.get_mut();
Pin::new(&mut this.listener).poll(cx)
}
}
pub struct TrunkWriter {
trunk: Arc<Trunk>,
}
impl TrunkWriter {
pub fn publish(&self, track_id: u32, retention: RetentionClass, sample: Sample) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
match retention {
RetentionClass::Timed => state.timed.push(track_id, sample),
RetentionClass::Sparse => state.sparse.push(track_id, sample),
}
}
pub fn publish_event(&self, event: TimedEvent, anchor: EventAnchor) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
state.events.push(event, anchor);
}
pub fn set_tracks(&self, tracks: Vec<TrackSpec>) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
state.tracks = Arc::from(tracks);
state.track_generation += 1;
drop(state);
self.trunk.progress.notify(usize::MAX);
}
}
pub struct SegmentWriter {
trunk: Arc<Trunk>,
}
impl SegmentWriter {
pub fn publish_segment(&self, entry: SegmentEntry) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
loop {
if state.segments.entries.len() < state.segments.capacity {
break;
}
let oldest = state.segments.base;
let mut must_wait = false;
for pin in state.segments.pins.values_mut() {
if pin.terminated || pin.consumed > oldest {
continue;
}
match pin.policy {
ArchiveOverrun::Gap => {}
ArchiveOverrun::Terminate => pin.terminated = true,
ArchiveOverrun::StallIngest => must_wait = true,
}
}
if !must_wait {
break;
}
state = self
.trunk
.segment_pin_released
.wait(state)
.expect("Trunk segment_pin_released condvar poisoned");
}
state.segments.push(entry);
drop(state);
self.trunk.progress.notify(usize::MAX);
}
pub fn publish_part(&self, entry: PartEntry) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
state.parts.push(entry);
drop(state);
self.trunk.progress.notify(usize::MAX);
}
pub fn note_segment_start(&self, segment_number: u32, start: MediaTime) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
state.events.note_segment_start(segment_number, start);
}
pub fn set_time_anchor(&self, anchor: TimeAnchor) {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
state.events.set_time_anchor(anchor);
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum SampleCursorItem {
Timed {
track_id: u32,
sample: Sample,
},
Sparse {
track_id: u32,
sample: Sample,
},
Lagged {
skipped: u64,
},
Degraded {
skipped: u64,
},
}
pub struct SampleCursor {
trunk: Arc<Trunk>,
timed_consumed: u64,
sparse_consumed: u64,
}
impl SampleCursor {
pub fn poll(&mut self) -> Option<SampleCursorItem> {
let state = self.trunk.state.lock().expect("Trunk state lock poisoned");
if self.timed_consumed < state.timed.base {
let skipped = state.timed.base - self.timed_consumed;
self.timed_consumed = state.timed.base;
return Some(SampleCursorItem::Lagged { skipped });
}
if self.sparse_consumed < state.sparse.base {
let skipped = state.sparse.base - self.sparse_consumed;
self.sparse_consumed = state.sparse.base;
return Some(SampleCursorItem::Degraded { skipped });
}
let sparse_idx = (self.sparse_consumed - state.sparse.base) as usize;
if let Some((track_id, sample)) = state.sparse.entries.get(sparse_idx) {
self.sparse_consumed += 1;
return Some(SampleCursorItem::Sparse {
track_id: *track_id,
sample: sample.clone(),
});
}
let timed_idx = (self.timed_consumed - state.timed.base) as usize;
if let Some((track_id, sample)) = state.timed.entries.get(timed_idx) {
self.timed_consumed += 1;
return Some(SampleCursorItem::Timed {
track_id: *track_id,
sample: sample.clone(),
});
}
None
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum SegmentCursorItem {
Segment(SegmentEntry),
Lagged {
skipped: u64,
},
Gap {
skipped: u64,
},
Terminated,
}
pub struct SegmentCursor {
trunk: Arc<Trunk>,
consumed: u64,
pin_id: Option<u64>,
done: bool,
}
impl SegmentCursor {
pub fn poll(&mut self) -> Option<SegmentCursorItem> {
if self.done {
return None;
}
let Some(pin_id) = self.pin_id else {
let state = self.trunk.state.lock().expect("Trunk state lock poisoned");
if self.consumed < state.segments.base {
let skipped = state.segments.base - self.consumed;
self.consumed = state.segments.base;
return Some(SegmentCursorItem::Lagged { skipped });
}
let idx = (self.consumed - state.segments.base) as usize;
return if let Some(entry) = state.segments.entries.get(idx) {
self.consumed += 1;
Some(SegmentCursorItem::Segment(entry.clone()))
} else {
None
};
};
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
let Some(pin) = state.segments.pins.get(&pin_id) else {
self.done = true;
return None;
};
if pin.terminated {
state.segments.pins.remove(&pin_id);
self.pin_id = None;
self.done = true;
return Some(SegmentCursorItem::Terminated);
}
let consumed = pin.consumed;
if consumed < state.segments.base {
let skipped = state.segments.base - consumed;
state
.segments
.pins
.get_mut(&pin_id)
.expect("pin_id was resolved from this same locked state, so its entry exists")
.consumed = state.segments.base;
drop(state);
self.trunk.segment_pin_released.notify_all();
return Some(SegmentCursorItem::Gap { skipped });
}
let idx = (consumed - state.segments.base) as usize;
if let Some(entry) = state.segments.entries.get(idx) {
let item = entry.clone();
state
.segments
.pins
.get_mut(&pin_id)
.expect("pin_id was resolved from this same locked state, so its entry exists")
.consumed += 1;
drop(state);
self.trunk.segment_pin_released.notify_all();
return Some(SegmentCursorItem::Segment(item));
}
None
}
}
impl Drop for SegmentCursor {
fn drop(&mut self) {
if let Some(pin_id) = self.pin_id.take() {
let mut state = self.trunk.state.lock().expect("Trunk state lock poisoned");
state.segments.pins.remove(&pin_id);
drop(state);
self.trunk.segment_pin_released.notify_all();
}
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum EventCursorItem {
Event(EventEntry),
Lagged {
skipped: u64,
},
}
pub struct EventCursor {
trunk: Arc<Trunk>,
consumed: u64,
}
impl EventCursor {
pub fn poll(&mut self) -> Option<EventCursorItem> {
let state = self.trunk.state.lock().expect("Trunk state lock poisoned");
let log = &state.events;
if self.consumed < log.base {
let skipped = log.base - self.consumed;
self.consumed = log.base;
return Some(EventCursorItem::Lagged { skipped });
}
let idx = (self.consumed - log.base) as usize;
if let Some(entry) = log.entries.get(idx) {
self.consumed += 1;
return Some(EventCursorItem::Event(entry.clone()));
}
None
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::mpsc;
use std::thread;
use transmux::pipeline::{CodecConfig, DataCarriage};
fn opaque_track(track_id: u32) -> TrackSpec {
TrackSpec::new(
track_id,
90_000,
CodecConfig::Data {
stream_type: 0x06,
descriptors: Vec::new(),
carriage: DataCarriage::Pes,
},
)
}
fn nz(n: usize) -> NonZeroUsize {
NonZeroUsize::new(n).expect("test capacity must be non-zero")
}
fn sample(byte: u8, len: usize) -> Sample {
Sample::new(Bytes::from(vec![byte; len]), Some(0), Some(0), None, true)
}
fn timed_data(item: &SampleCursorItem) -> Option<(u32, &Sample)> {
match item {
SampleCursorItem::Timed { track_id, sample } => Some((*track_id, sample)),
_ => None,
}
}
fn segment_entry(byte: u8, seq: u32) -> SegmentEntry {
SegmentEntry::new(
Bytes::from(vec![byte; 16]),
seq,
Duration::from_secs(2),
Timestamp::from_nanos(u64::from(seq) * 2_000_000_000),
SegmentMeta {
discontinuous: false,
},
)
}
fn segment_data(item: &SegmentCursorItem) -> Option<&SegmentEntry> {
match item {
SegmentCursorItem::Segment(entry) => Some(entry),
_ => None,
}
}
fn drain_segments(cursor: &mut SegmentCursor, n: usize) -> Vec<SegmentCursorItem> {
let mut out = Vec::new();
for _ in 0..n {
match cursor.poll() {
Some(item) => out.push(item),
None => break,
}
}
out
}
fn drain(cursor: &mut SampleCursor, n: usize) -> Vec<SampleCursorItem> {
let mut out = Vec::new();
for _ in 0..n {
match cursor.poll() {
Some(item) => out.push(item),
None => break,
}
}
out
}
#[test]
fn multiple_cursors_see_every_sample_in_order_with_no_dup_or_loss() {
let trunk = Trunk::new(TrunkConfig::new(nz(100), nz(10), nz(4), nz(8), nz(8)));
let mut c1 = trunk.subscribe();
let mut c2 = trunk.subscribe();
let mut c3 = trunk.subscribe();
let writer = trunk.writer().unwrap();
for i in 0u8..5 {
writer.publish(7, RetentionClass::Timed, sample(i, 16));
}
for cursor in [&mut c1, &mut c2, &mut c3] {
let items = drain(cursor, 5);
assert_eq!(items.len(), 5, "each cursor must see exactly 5 samples");
let bytes: Vec<u8> = items
.iter()
.map(|item| timed_data(item).unwrap().1.data[0])
.collect();
assert_eq!(bytes, vec![0, 1, 2, 3, 4], "must be in publish order");
assert!(cursor.poll().is_none(), "no extra/duplicated items");
}
}
#[test]
fn slow_reader_lags_but_writer_completes_regardless() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(10), nz(4), nz(8), nz(8)));
let mut slow = trunk.subscribe();
let writer = trunk.writer().unwrap();
for i in 0u8..=255u8 {
for _ in 0..4 {
writer.publish(1, RetentionClass::Timed, sample(i, 8));
}
}
assert_eq!(
trunk.timed_len(),
4,
"writer unblocked: ring stayed bounded"
);
let first = slow.poll().unwrap();
assert!(
matches!(first, SampleCursorItem::Lagged { skipped: 1020 }),
"expected Lagged{{skipped: 1020}}, got {first:?}"
);
}
#[test]
fn lag_is_reported_with_an_accurate_skipped_count() {
let trunk = Trunk::new(TrunkConfig::new(nz(3), nz(10), nz(4), nz(8), nz(8)));
let mut cursor = trunk.subscribe();
let writer = trunk.writer().unwrap();
for i in 0u8..9 {
writer.publish(2, RetentionClass::Timed, sample(i, 4));
}
let first = cursor.poll().unwrap();
assert!(
matches!(first, SampleCursorItem::Lagged { skipped: 6 }),
"expected Lagged{{skipped: 6}}, got {first:?}"
);
let items = drain(&mut cursor, 3);
let bytes: Vec<u8> = items
.iter()
.map(|item| timed_data(item).unwrap().1.data[0])
.collect();
assert_eq!(bytes, vec![6, 7, 8]);
assert!(cursor.poll().is_none());
}
#[test]
fn sparse_reader_loses_data_reports_degraded_distinguishable_from_timed_lagged() {
let trunk = Trunk::new(TrunkConfig::new(nz(2), nz(2), nz(4), nz(8), nz(8)));
let mut cursor = trunk.subscribe();
let writer = trunk.writer().unwrap();
for i in 0u8..5 {
writer.publish(3, RetentionClass::Timed, sample(i, 4));
}
for i in 0u8..4 {
writer.publish(9, RetentionClass::Sparse, sample(100 + i, 4));
}
let timed_loss = cursor.poll().unwrap();
assert!(
matches!(timed_loss, SampleCursorItem::Lagged { skipped: 3 }),
"expected ordinary Lagged{{skipped: 3}} for the Timed ring, got {timed_loss:?}"
);
let sparse_loss = cursor.poll().unwrap();
assert!(
matches!(sparse_loss, SampleCursorItem::Degraded { skipped: 2 }),
"expected escalated Degraded{{skipped: 2}} for the Sparse ring, got {sparse_loss:?}"
);
assert_ne!(
core::mem::discriminant(&timed_loss),
core::mem::discriminant(&sparse_loss),
"Lagged and Degraded must be distinct variants, not merely different field values"
);
}
#[test]
fn ring_is_bounded_under_flood_on_both_classes() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(3), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
for i in 0u32..50_000 {
writer.publish(5, RetentionClass::Timed, sample((i % 256) as u8, 2));
assert!(
trunk.timed_len() <= 4,
"Timed ring exceeded its cap mid-flood"
);
if i % 7 == 0 {
writer.publish(6, RetentionClass::Sparse, sample((i % 256) as u8, 2));
assert!(
trunk.sparse_len() <= 3,
"Sparse ring exceeded its cap mid-flood"
);
}
}
assert_eq!(trunk.timed_len(), 4);
assert_eq!(trunk.sparse_len(), 3);
}
#[test]
fn subscribe_from_backlog_replays_exact_resident_entries_both_classes() {
let trunk = Trunk::new(TrunkConfig::new(nz(100), nz(100), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
for i in 10u8..13 {
writer.publish(1, RetentionClass::Timed, sample(i, 4));
}
for i in 20u8..22 {
writer.publish(2, RetentionClass::Sparse, sample(i, 4));
}
let mut cursor = trunk.subscribe_from_backlog();
let sparse_items = drain(&mut cursor, 2);
let sparse_bytes: Vec<u8> = sparse_items
.iter()
.map(|item| match item {
SampleCursorItem::Sparse { sample, .. } => sample.data[0],
other => panic!("expected Sparse, got {other:?}"),
})
.collect();
assert_eq!(sparse_bytes, vec![20, 21], "exact resident Sparse backlog");
let timed_items = drain(&mut cursor, 3);
let timed_bytes: Vec<u8> = timed_items
.iter()
.map(|item| timed_data(item).unwrap().1.data[0])
.collect();
assert_eq!(
timed_bytes,
vec![10, 11, 12],
"exact resident Timed backlog"
);
assert!(
cursor.poll().is_none(),
"no extra items beyond the resident backlog"
);
writer.publish(1, RetentionClass::Timed, sample(99, 4));
let live = cursor.poll().unwrap();
assert_eq!(timed_data(&live).unwrap().1.data[0], 99);
}
#[test]
fn subscribe_from_backlog_reports_lagged_when_backlog_already_overwritten() {
let trunk = Trunk::new(TrunkConfig::new(nz(3), nz(10), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
for i in 0u8..9 {
writer.publish(1, RetentionClass::Timed, sample(i, 4));
}
let mut cursor = trunk.subscribe_from_backlog();
for i in 9u8..15 {
writer.publish(1, RetentionClass::Timed, sample(i, 4));
}
let first = cursor.poll().unwrap();
assert!(
matches!(first, SampleCursorItem::Lagged { skipped: 6 }),
"expected Lagged{{skipped: 6}}, got {first:?}"
);
let items = drain(&mut cursor, 3);
let bytes: Vec<u8> = items
.iter()
.map(|item| timed_data(item).unwrap().1.data[0])
.collect();
assert_eq!(bytes, vec![12, 13, 14]);
assert!(cursor.poll().is_none());
}
#[test]
fn payload_is_shared_not_copied_across_cursors() {
let trunk = Trunk::new(TrunkConfig::new(nz(8), nz(8), nz(4), nz(8), nz(8)));
let mut c1 = trunk.subscribe();
let mut c2 = trunk.subscribe();
let mut c3 = trunk.subscribe();
let writer = trunk.writer().unwrap();
writer.publish(4, RetentionClass::Timed, sample(0xAB, 65536));
let i1 = c1.poll().unwrap();
let i2 = c2.poll().unwrap();
let i3 = c3.poll().unwrap();
let p1 = timed_data(&i1).unwrap().1.data.as_ptr();
let p2 = timed_data(&i2).unwrap().1.data.as_ptr();
let p3 = timed_data(&i3).unwrap().1.data.as_ptr();
assert_eq!(
p1, p2,
"cursor 2's payload must be the SAME allocation as cursor 1's"
);
assert_eq!(
p2, p3,
"cursor 3's payload must be the SAME allocation as cursor 1's"
);
assert_eq!(timed_data(&i1).unwrap().1.data.len(), 65536);
}
#[test]
fn second_writer_is_refused() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let _first = trunk.writer().unwrap();
assert!(
trunk.writer().is_none(),
"a Trunk has exactly one sample/event writer"
);
}
#[test]
fn second_segment_writer_is_refused() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let _first = trunk.segment_writer().unwrap();
assert!(
trunk.segment_writer().is_none(),
"a Trunk has exactly one segment/part writer"
);
}
#[test]
fn sample_and_segment_writers_can_be_held_simultaneously() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let samples = trunk.writer().unwrap();
let segments = trunk.segment_writer().expect(
"the segment/part writer must still be takeable while the sample/event writer is held",
);
samples.publish(1, RetentionClass::Timed, sample(1, 4));
segments.publish_segment(segment_entry(1, 1));
assert_eq!(trunk.timed_len(), 1);
assert_eq!(trunk.segment_len(), 1);
}
#[test]
fn segmenter_holds_sample_cursor_and_segment_writer_at_once() {
let trunk = Trunk::new(TrunkConfig::new(nz(8), nz(4), nz(4), nz(8), nz(8)));
let ingest = trunk.writer().unwrap();
let mut samples = trunk.subscribe();
let segmenter = trunk.segment_writer().expect(
"a segmenter must be able to take the segment/part writer \
while the ingest driver still holds the sample/event writer",
);
let mut seg_cursor = trunk.subscribe_segments();
for i in 0u8..3 {
ingest.publish(1, RetentionClass::Timed, sample(i, 4));
}
let mut muxed = Vec::new();
for _ in 0..3 {
match samples.poll() {
Some(SampleCursorItem::Timed { sample, .. }) => {
muxed.push(sample.data[0]);
}
other => panic!("expected a Timed sample, got {other:?}"),
}
}
assert_eq!(muxed, vec![0, 1, 2]);
segmenter.publish_part(part_entry(0xAB, 1, 0));
segmenter.publish_segment(segment_entry(0xAA, 1));
let got = seg_cursor
.poll()
.expect("the segmenter's published segment must be visible");
assert_eq!(segment_data(&got).unwrap().sequence_number, 1);
assert_eq!(
trunk.part_bytes(1, 0),
Some(Bytes::from(vec![0xAB; 8])),
"the segmenter's published part must be individually addressable too"
);
}
#[test]
fn subscribe_starts_from_now_not_from_history() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
writer.publish(1, RetentionClass::Timed, sample(1, 4));
writer.publish(1, RetentionClass::Timed, sample(2, 4));
let mut cursor = trunk.subscribe();
assert!(cursor.poll().is_none());
writer.publish(1, RetentionClass::Timed, sample(3, 4));
let item = cursor.poll().unwrap();
assert_eq!(timed_data(&item).unwrap().1.data[0], 3);
}
#[test]
fn multiple_segment_cursors_see_every_segment_in_order_with_no_dup_or_loss() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(100), nz(8), nz(8)));
let mut c1 = trunk.subscribe_segments();
let mut c2 = trunk.subscribe_segments();
let mut c3 = trunk.subscribe_segments();
let writer = trunk.segment_writer().unwrap();
for i in 0u32..5 {
writer.publish_segment(segment_entry(i as u8, i + 1));
}
for cursor in [&mut c1, &mut c2, &mut c3] {
let items = drain_segments(cursor, 5);
assert_eq!(items.len(), 5, "each cursor must see exactly 5 segments");
let seqs: Vec<u32> = items
.iter()
.map(|item| segment_data(item).unwrap().sequence_number)
.collect();
assert_eq!(seqs, vec![1, 2, 3, 4, 5], "must be in playlist order");
assert!(cursor.poll().is_none(), "no extra/duplicated items");
}
}
#[test]
fn non_pinning_slow_segment_reader_lags_but_writer_completes_regardless() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let mut slow = trunk.subscribe_segments();
let writer = trunk.segment_writer().unwrap();
for i in 0u32..1024 {
writer.publish_segment(segment_entry((i % 256) as u8, i + 1));
}
assert_eq!(
trunk.segment_len(),
4,
"writer unblocked: segment log stayed bounded"
);
let first = slow.poll().unwrap();
assert!(
matches!(first, SegmentCursorItem::Lagged { skipped: 1020 }),
"expected Lagged{{skipped: 1020}}, got {first:?}"
);
}
#[test]
fn pinning_reader_receives_every_segment_while_non_pinning_reader_lags() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(2), nz(8), nz(8)));
let mut slow = trunk.subscribe_segments(); let mut archive = trunk.pin_segments(ArchiveOverrun::StallIngest); let writer = Arc::new(trunk.segment_writer().unwrap());
writer.publish_segment(segment_entry(1, 1));
writer.publish_segment(segment_entry(2, 2));
let (done_tx, done_rx) = mpsc::channel();
let blocked_writer = Arc::clone(&writer);
let handle = thread::spawn(move || {
blocked_writer.publish_segment(segment_entry(3, 3));
done_tx.send(()).unwrap();
});
assert!(
done_rx.recv_timeout(Duration::from_millis(200)).is_err(),
"publish_segment must still be blocked: archive has not consumed seq 1 yet"
);
let first = archive.poll().unwrap();
assert_eq!(segment_data(&first).unwrap().sequence_number, 1);
done_rx
.recv_timeout(Duration::from_secs(60))
.expect("publish_segment must unblock once the pin advances");
handle.join().unwrap();
let second = archive.poll().unwrap();
assert_eq!(segment_data(&second).unwrap().sequence_number, 2);
let third = archive.poll().unwrap();
assert_eq!(segment_data(&third).unwrap().sequence_number, 3);
assert!(archive.poll().is_none());
let lag = slow.poll().unwrap();
assert!(
matches!(lag, SegmentCursorItem::Lagged { skipped: 1 }),
"expected Lagged{{skipped: 1}}, got {lag:?}"
);
let remaining: Vec<u32> = drain_segments(&mut slow, 2)
.iter()
.map(|item| segment_data(item).unwrap().sequence_number)
.collect();
assert_eq!(remaining, vec![2, 3]);
}
#[test]
fn pinning_is_bounded_an_unacking_consumer_cannot_grow_memory_without_limit() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let _archive = trunk.pin_segments(ArchiveOverrun::default());
let writer = trunk.segment_writer().unwrap();
for i in 0u32..50_000 {
writer.publish_segment(segment_entry((i % 256) as u8, i + 1));
assert!(
trunk.segment_len() <= 4,
"segment log exceeded its cap mid-flood despite an un-acking pinning cursor"
);
}
assert_eq!(trunk.segment_len(), 4);
}
#[test]
fn archive_overrun_gap_evicts_and_reports_gap() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(2), nz(8), nz(8)));
let mut archive = trunk.pin_segments(ArchiveOverrun::Gap);
let writer = trunk.segment_writer().unwrap();
for i in 0u32..5 {
writer.publish_segment(segment_entry(i as u8, i + 1));
}
assert_eq!(trunk.segment_len(), 2);
let gap = archive.poll().unwrap();
assert!(
matches!(gap, SegmentCursorItem::Gap { skipped: 3 }),
"expected Gap{{skipped: 3}}, got {gap:?}"
);
let remaining: Vec<u32> = drain_segments(&mut archive, 2)
.iter()
.map(|item| segment_data(item).unwrap().sequence_number)
.collect();
assert_eq!(remaining, vec![4, 5]);
assert!(archive.poll().is_none());
}
#[test]
fn archive_overrun_stall_ingest_actually_blocks_the_writer() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(1), nz(8), nz(8)));
let mut archive = trunk.pin_segments(ArchiveOverrun::StallIngest);
let writer = Arc::new(trunk.segment_writer().unwrap());
writer.publish_segment(segment_entry(1, 1));
let (done_tx, done_rx) = mpsc::channel();
let blocked_writer = Arc::clone(&writer);
let handle = thread::spawn(move || {
blocked_writer.publish_segment(segment_entry(2, 2));
done_tx.send(()).unwrap();
});
assert!(
done_rx.recv_timeout(Duration::from_millis(200)).is_err(),
"publish_segment must block: the pin has not consumed seq 1 yet"
);
let first = archive.poll().unwrap();
assert_eq!(segment_data(&first).unwrap().sequence_number, 1);
done_rx
.recv_timeout(Duration::from_secs(60))
.expect("publish_segment must unblock once the pin advances");
handle.join().unwrap();
}
#[test]
fn archive_overrun_terminate_drops_the_cursor() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(2), nz(8), nz(8)));
let mut archive = trunk.pin_segments(ArchiveOverrun::Terminate);
let writer = trunk.segment_writer().unwrap();
for i in 0u32..5 {
writer.publish_segment(segment_entry(i as u8, i + 1));
}
assert_eq!(trunk.segment_len(), 2, "writer unblocked despite Terminate");
let item = archive.poll().unwrap();
assert!(
matches!(item, SegmentCursorItem::Terminated),
"expected Terminated, got {item:?}"
);
assert!(archive.poll().is_none());
assert!(archive.poll().is_none());
writer.publish_segment(segment_entry(9, 6));
assert_eq!(trunk.segment_len(), 2);
}
#[test]
fn segment_bytes_are_shared_not_copied_across_cursors() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(8), nz(8), nz(8)));
let mut c1 = trunk.subscribe_segments();
let mut c2 = trunk.subscribe_segments();
let mut c3 = trunk.subscribe_segments();
let writer = trunk.segment_writer().unwrap();
writer.publish_segment(SegmentEntry::new(
Bytes::from(vec![0xCDu8; 65536]),
1,
Duration::from_secs(2),
Timestamp::from_nanos(0),
SegmentMeta {
discontinuous: false,
},
));
let i1 = c1.poll().unwrap();
let i2 = c2.poll().unwrap();
let i3 = c3.poll().unwrap();
let p1 = segment_data(&i1).unwrap().bytes.as_ptr();
let p2 = segment_data(&i2).unwrap().bytes.as_ptr();
let p3 = segment_data(&i3).unwrap().bytes.as_ptr();
assert_eq!(
p1, p2,
"cursor 2's segment payload must be the SAME allocation as cursor 1's"
);
assert_eq!(
p2, p3,
"cursor 3's segment payload must be the SAME allocation as cursor 1's"
);
assert_eq!(segment_data(&i1).unwrap().bytes.len(), 65536);
}
use timed_metadata::{EventKind, SourcePayload};
fn basic_event(id: u32) -> TimedEvent {
TimedEvent {
id: Some(id),
kind: EventKind::BreakStart,
at: None,
duration: None,
source: SourcePayload::Scte35 { raw: Vec::new() },
}
}
fn event_id(item: &EventCursorItem) -> Option<u32> {
match item {
EventCursorItem::Event(e) => e.event.id,
_ => None,
}
}
fn splice_insert_bytes(event_id: u32, pts_time: u64) -> Vec<u8> {
use broadcast_common::Serialize;
use scte35_splice::SpliceInfoSection;
use scte35_splice::commands::AnyCommand;
use scte35_splice::commands::splice_insert::SpliceInsert;
use scte35_splice::time::SpliceTime;
let si = SpliceInsert {
splice_event_id: event_id,
out_of_network_indicator: true,
splice_time: Some(SpliceTime::with_pts(pts_time)),
..SpliceInsert::default()
};
let section = SpliceInfoSection::new_clear(AnyCommand::SpliceInsert(si), &[]);
section.to_bytes()
}
#[test]
fn events_between_returns_exactly_the_half_open_range() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
for (id, ticks) in [(1u32, 1_000u64), (2, 2_000), (3, 3_000), (4, 4_000)] {
writer.publish_event(basic_event(id), EventAnchor::Media(MediaTime(ticks)));
}
let got = trunk.events_between(MediaTime(2_000), MediaTime(4_000));
let ids: Vec<u32> = got.iter().map(|e| e.event.id.unwrap()).collect();
assert_eq!(
ids,
vec![2, 3],
"start (2_000) inclusive, end (4_000) exclusive"
);
}
#[test]
fn segment_relative_event_resolves_at_publish_time_when_boundary_already_known() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
let segment_writer = trunk.segment_writer().unwrap();
segment_writer.note_segment_start(3, MediaTime(300_000));
writer.publish_event(
basic_event(9),
EventAnchor::Segment {
segment_number: 3,
delta: 1_500,
},
);
let got = trunk.events_in_segment(3);
assert_eq!(got.len(), 1);
assert_eq!(got[0].event.id, Some(9));
assert!(matches!(
got[0].anchor,
EventAnchor::Media(MediaTime(t)) if t == 301_500
));
}
#[test]
fn segment_relative_event_resolves_to_the_named_segment_not_whichever_is_open() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
let segment_writer = trunk.segment_writer().unwrap();
let mut cursor = trunk.subscribe_events();
writer.publish_event(
basic_event(42),
EventAnchor::Segment {
segment_number: 2,
delta: 1_000,
},
);
segment_writer.note_segment_start(1, MediaTime(0));
let item = cursor.poll().unwrap();
let entry = match item {
EventCursorItem::Event(e) => e,
other => panic!("expected Event, got {other:?}"),
};
assert!(
matches!(
entry.anchor,
EventAnchor::Segment {
segment_number: 2,
delta: 1_000
}
),
"must stay pending on segment 2 — segment 1 being open must not \
resolve it against the wrong boundary: {:?}",
entry.anchor
);
assert!(
trunk.events_in_segment(2).is_empty(),
"not resolved yet: must not appear under segment 2 either"
);
assert!(trunk.events_in_segment(1).is_empty());
segment_writer.note_segment_start(2, MediaTime(90_000));
let in_seg2 = trunk.events_in_segment(2);
assert_eq!(in_seg2.len(), 1);
assert_eq!(in_seg2[0].event.id, Some(42));
assert!(matches!(
in_seg2[0].anchor,
EventAnchor::Media(MediaTime(t)) if t == 91_000
));
assert!(
trunk.events_in_segment(1).is_empty(),
"must not ALSO appear under segment 1"
);
}
#[test]
fn utc_only_event_stays_unanchored_until_a_time_anchor_arrives() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
let segment_writer = trunk.segment_writer().unwrap();
let mut cursor = trunk.subscribe_events();
writer.publish_event(
basic_event(7),
EventAnchor::Utc {
utc_epoch_ms: 5_000,
},
);
let item = cursor.poll().unwrap();
let entry = match item {
EventCursorItem::Event(e) => e,
other => panic!("expected Event, got {other:?}"),
};
assert!(
matches!(
entry.anchor,
EventAnchor::Utc {
utc_epoch_ms: 5_000
}
),
"must stay honestly unanchored — NO fabricated media time: {:?}",
entry.anchor
);
assert!(
trunk
.events_between(MediaTime(0), MediaTime(u64::MAX))
.is_empty(),
"an unanchored event must not appear in a media-time query"
);
segment_writer.set_time_anchor(TimeAnchor {
pts_90k: 0,
utc_epoch_ms: 1_000,
});
let resolved = trunk.events_between(MediaTime(0), MediaTime(u64::MAX));
assert_eq!(resolved.len(), 1);
assert_eq!(resolved[0].event.id, Some(7));
assert!(
matches!(resolved[0].anchor, EventAnchor::Media(MediaTime(t)) if t == 360_000),
"expected MediaTime(360_000), got {:?}",
resolved[0].anchor
);
}
#[test]
fn a_33_bit_pts_wrap_does_not_corrupt_event_log_ordering() {
const PTS_WRAP: u64 = 1u64 << 33;
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
let mut timeline = timed_metadata::Timeline::new();
let before_wrap = splice_insert_bytes(1, PTS_WRAP - 10);
let ev1 = timeline.push_scte35(&before_wrap).unwrap();
let at1 = ev1.at.unwrap();
writer.publish_event(ev1, EventAnchor::Media(at1));
let after_wrap = splice_insert_bytes(2, 5);
let ev2 = timeline.push_scte35(&after_wrap).unwrap();
let at2 = ev2.at.unwrap();
writer.publish_event(ev2, EventAnchor::Media(at2));
assert!(
at2.0 > at1.0,
"Timeline itself must unroll monotonically: at1={}, at2={}",
at1.0,
at2.0
);
let got = trunk.events_between(MediaTime(0), MediaTime(u64::MAX));
assert_eq!(got.len(), 2);
assert_eq!(got[0].event.id, Some(1));
assert_eq!(got[1].event.id, Some(2));
assert!(
matches!(got[0].anchor, EventAnchor::Media(t) if t.0 == at1.0),
"event 1's stored anchor must equal Timeline's unrolled value \
exactly, got {:?}",
got[0].anchor
);
assert!(
matches!(got[1].anchor, EventAnchor::Media(t) if t.0 == at2.0),
"event 2's stored (post-wrap) anchor must equal Timeline's \
unrolled value exactly — not re-masked back into 33 bits, got {:?}",
got[1].anchor
);
}
#[test]
fn event_log_is_bounded_under_flood() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(3), nz(8)));
let writer = trunk.writer().unwrap();
for i in 0u32..50_000 {
writer.publish_event(basic_event(i), EventAnchor::Media(MediaTime(u64::from(i))));
assert!(
trunk.event_len() <= 3,
"event log exceeded its cap mid-flood"
);
}
assert_eq!(trunk.event_len(), 3);
}
#[test]
fn event_cursor_lag_is_reported_with_an_accurate_skipped_count_writer_never_blocks() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(3), nz(8)));
let mut cursor = trunk.subscribe_events();
let writer = trunk.writer().unwrap();
for i in 0u32..9 {
writer.publish_event(basic_event(i), EventAnchor::Media(MediaTime(u64::from(i))));
}
assert_eq!(
trunk.event_len(),
3,
"writer unblocked: event log stayed bounded"
);
let first = cursor.poll().unwrap();
assert!(
matches!(first, EventCursorItem::Lagged { skipped: 6 }),
"expected Lagged{{skipped: 6}}, got {first:?}"
);
let mut ids = Vec::new();
for _ in 0..3 {
ids.push(event_id(&cursor.poll().unwrap()).unwrap());
}
assert_eq!(ids, vec![6, 7, 8]);
assert!(cursor.poll().is_none());
}
fn round_trip_anchor() -> TimeAnchor {
TimeAnchor {
pts_90k: 900_000, utc_epoch_ms: 1_700_000_000_000, }
}
const TICKS_PER_EPOCH_MS: u64 = PTS_HZ / 1000;
const MEDIA_ROUND_TRIP_MAX_TICKS: u64 = TICKS_PER_EPOCH_MS - 1;
#[test]
fn epoch_ms_round_trip_through_media_time_is_exact() {
let anchor = round_trip_anchor();
for offset_ms in [
0i64,
1,
-1,
1_000,
-1_000,
86_400_000,
-9_000,
200_000_000_000_000,
] {
let epoch_ms = anchor.utc_epoch_ms + offset_ms;
let media = epoch_ms_to_media(&anchor, epoch_ms);
let back = anchor.media_to_epoch_ms(media);
assert_eq!(
back, epoch_ms,
"epoch_ms -> media -> epoch_ms must be EXACT at offset {offset_ms} ms \
(media = {media:?})"
);
}
}
#[test]
fn media_round_trip_is_lossy_by_at_most_one_millisecond() {
let anchor = round_trip_anchor();
let mut worst = 0u64;
for offset_ticks in [
0i64,
1,
-1,
89,
-89,
90,
-90,
91,
-91,
18_000_000_000_000_037,
] {
let media = MediaTime((anchor.pts_90k as i64 + offset_ticks) as u64);
let epoch_ms = anchor.media_to_epoch_ms(media);
let back = epoch_ms_to_media(&anchor, epoch_ms);
let diff = media.0.abs_diff(back.0);
assert!(
diff <= MEDIA_ROUND_TRIP_MAX_TICKS,
"media -> epoch_ms -> media lost {diff} ticks at offset \
{offset_ticks} (bound is {MEDIA_ROUND_TRIP_MAX_TICKS}, i.e. \
< 1 ms): {media:?} -> {epoch_ms} -> {back:?}"
);
worst = worst.max(diff);
}
assert_eq!(
worst, MEDIA_ROUND_TRIP_MAX_TICKS,
"the derived bound must be attained, not merely respected — \
otherwise it is a loose tolerance, not a property"
);
}
#[test]
fn epoch_before_the_timeline_origin_clamps_to_zero_which_is_safe_not_correct() {
let anchor = round_trip_anchor();
let epoch_ms = anchor.utc_epoch_ms - 20_000;
let media = epoch_ms_to_media(&anchor, epoch_ms);
assert_eq!(
media,
MediaTime(0),
"a pre-origin epoch must clamp to the start of the timeline"
);
assert!(
media.0 < u64::from(u32::MAX),
"must not have wrapped into the far future: {media:?}"
);
let back = anchor.media_to_epoch_ms(media);
assert_ne!(
back, epoch_ms,
"clamping is lossy by construction — this asserts the honest \
claim (safe) rather than the false one (correct)"
);
assert_eq!(
back,
anchor.utc_epoch_ms - 10_000,
"clamped media time 0 maps back to the timeline origin (10 s \
before the anchor), not to the requested instant"
);
}
fn part_entry(byte: u8, segment_number: u32, part_index: u32) -> PartEntry {
PartEntry::new(
Bytes::from(vec![byte; 8]),
segment_number,
part_index,
Duration::from_millis(200),
part_index == 0,
)
}
#[test]
fn part_is_addressable_and_readable_before_its_parent_segment_closes() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.segment_writer().unwrap();
assert!(
trunk.last_closed_segment().is_none(),
"sanity check: nothing has closed yet"
);
writer.publish_part(part_entry(0xAB, 9, 0));
let bytes = trunk
.part_bytes(9, 0)
.expect("a published part of an open segment must be addressable now");
assert_eq!(bytes, Bytes::from(vec![0xAB; 8]));
assert!(
trunk.last_closed_segment().is_none(),
"the part landed without any segment ever closing"
);
}
#[test]
fn waiter_is_woken_when_the_awaited_part_lands_and_resolves_to_it() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = Arc::new(trunk.segment_writer().unwrap());
let listener = trunk.listen().expect("first registration must succeed");
assert!(trunk.part_bytes(9, 0).is_none(), "not published yet");
let bg_writer = Arc::clone(&writer);
let handle = thread::spawn(move || {
thread::sleep(Duration::from_millis(50));
bg_writer.publish_part(part_entry(0xCC, 9, 1));
bg_writer.publish_part(part_entry(0xAB, 9, 0));
});
let woken = listener.wait_deadline(std::time::Instant::now() + Duration::from_secs(60));
assert!(
woken,
"listener must wake on publish_part, not park forever"
);
handle.join().unwrap();
let bytes = trunk
.part_bytes(9, 0)
.expect("the awaited part must now be readable");
assert_eq!(
bytes,
Bytes::from(vec![0xAB; 8]),
"must resolve to the awaited part (index 0), not the decoy (index 1)"
);
}
#[test]
fn waiter_whose_target_never_arrives_is_bounded_not_parked_forever() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let listener = trunk.listen().unwrap();
let start = std::time::Instant::now();
let woken = listener.wait_deadline(start + Duration::from_millis(150));
let elapsed = start.elapsed();
assert!(!woken, "must report timeout, not a fabricated wake-up");
assert!(
elapsed < Duration::from_secs(60),
"must actually return at the deadline, not hang: took {elapsed:?}"
);
}
#[test]
fn writer_never_blocks_with_a_registered_never_serviced_waiter() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = Arc::new(trunk.segment_writer().unwrap());
let _never_serviced = trunk.listen().unwrap();
let (done_tx, done_rx) = mpsc::channel();
let bg_writer = Arc::clone(&writer);
thread::spawn(move || {
bg_writer.publish_part(part_entry(1, 1, 0));
bg_writer.publish_segment(segment_entry(2, 1));
done_tx.send(()).unwrap();
});
done_rx.recv_timeout(Duration::from_secs(60)).expect(
"publish_part/publish_segment must complete promptly even with \
a live, never-serviced waiter registered",
);
}
#[test]
fn waiter_set_is_bounded_a_flood_of_listen_calls_cannot_grow_without_limit() {
let cap = nz(4).get();
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(cap)));
let mut held: Vec<ProgressListener> = Vec::new();
for _ in 0..cap {
held.push(trunk.listen().expect("must succeed up to the cap"));
}
assert_eq!(trunk.waiter_count(), cap);
assert!(
trunk.listen().is_none(),
"must refuse a registration beyond part_capacity"
);
for _ in 0..50_000 {
let l = trunk.listen();
assert!(
trunk.waiter_count() <= cap,
"waiter count exceeded part_capacity mid-flood"
);
drop(l);
}
assert_eq!(
trunk.waiter_count(),
cap,
"the held registrations are still exactly at the cap"
);
held.pop();
assert_eq!(trunk.waiter_count(), cap - 1);
assert!(
trunk.listen().is_some(),
"a released slot must be re-usable"
);
}
#[test]
fn parts_remain_addressable_after_segment_close_until_ordinary_eviction_reclaims_them() {
let part_cap = nz(4).get();
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(part_cap)));
let writer = trunk.segment_writer().unwrap();
writer.publish_part(part_entry(0xAB, 1, 0));
writer.publish_segment(segment_entry(1, 1));
assert_eq!(
trunk.part_bytes(1, 0),
Some(Bytes::from(vec![0xAB; 8])),
"a just-closed segment's part must still be individually fetchable"
);
assert_eq!(trunk.last_closed_segment(), Some(1));
for i in 0..part_cap as u32 {
writer.publish_part(part_entry(0xFF, 99, i));
}
assert!(
trunk.part_bytes(1, 0).is_none(),
"the part is gone once ordinary part_capacity eviction reclaims \
it — NOT because its segment closed, but because the ring's own \
bound was exceeded, exactly like every other ring in this module"
);
}
#[test]
fn fresh_trunk_has_no_tracks_and_generation_zero() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
assert_eq!(trunk.tracks().len(), 0, "nothing has ever set a track set");
assert_eq!(trunk.track_generation(), 0);
}
#[test]
fn set_tracks_replaces_the_whole_set_and_bumps_generation_by_one_per_call() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
writer.set_tracks(vec![opaque_track(1)]);
assert_eq!(
trunk
.tracks()
.iter()
.map(|t| t.track_id)
.collect::<Vec<_>>(),
vec![1]
);
assert_eq!(trunk.track_generation(), 1);
writer.set_tracks(vec![opaque_track(7), opaque_track(9)]);
assert_eq!(
trunk
.tracks()
.iter()
.map(|t| t.track_id)
.collect::<Vec<_>>(),
vec![7, 9],
"set_tracks must replace the set wholesale, not append to it"
);
assert_eq!(
trunk.track_generation(),
2,
"generation must advance by exactly one per set_tracks call"
);
}
#[test]
fn generation_is_stable_across_unrelated_activity() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = trunk.writer().unwrap();
writer.set_tracks(vec![opaque_track(1)]);
assert_eq!(trunk.track_generation(), 1);
for i in 0u8..10 {
writer.publish(1, RetentionClass::Timed, sample(i, 4));
}
writer.publish_event(basic_event(1), EventAnchor::Media(MediaTime(0)));
assert_eq!(
trunk.track_generation(),
1,
"publishing samples/events must never bump track_generation"
);
}
#[test]
fn set_tracks_wakes_a_registered_listener() {
let trunk = Trunk::new(TrunkConfig::new(nz(4), nz(4), nz(4), nz(8), nz(8)));
let writer = Arc::new(trunk.writer().unwrap());
let listener = trunk.listen().expect("first registration must succeed");
let bg_writer = Arc::clone(&writer);
let handle = thread::spawn(move || {
thread::sleep(Duration::from_millis(50));
bg_writer.set_tracks(vec![opaque_track(1)]);
});
let woken = listener.wait_deadline(std::time::Instant::now() + Duration::from_secs(60));
assert!(woken, "listener must wake on set_tracks, not park forever");
handle.join().unwrap();
assert_eq!(trunk.track_generation(), 1);
}
}