use std::fs::{self, File, OpenOptions};
use std::io::Write;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use dvb_si::demux::{SectionEvent, SiDemux};
use dvb_si::tables::eit::{EitKind, PID as EIT_PID};
use dvb_si::tables::{AnyTableSection, RunningStatus};
use media_plane::egress::SegmentEgress;
use media_plane::trunk::{ArchiveOverrun, SegmentCursor, SegmentCursorItem, Trunk};
use mpeg_ts::pid::Pid;
use serde::{Deserialize, Serialize};
use tracing;
const TS_PACKET_LEN: usize = 188;
const DEFAULT_PERIOD_DURATION_SECS: u64 = 10800;
#[derive(Debug, Clone, PartialEq, Default, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DvrConfig {
#[serde(default)]
pub enabled: bool,
#[serde(default)]
pub archive_root: String,
#[serde(default = "default_period_duration_secs")]
pub period_duration_secs: u64,
#[serde(default)]
pub retention_periods: usize,
#[serde(default)]
pub retention_bytes: u64,
#[serde(default)]
pub overrun: ArchiveOverrunSerde,
#[serde(default)]
pub dvb_service_id: Option<u16>,
}
fn default_period_duration_secs() -> u64 {
DEFAULT_PERIOD_DURATION_SECS
}
#[derive(Debug, Clone, Copy, PartialEq, Default, Deserialize)]
#[non_exhaustive]
#[serde(rename_all = "lowercase")]
pub enum ArchiveOverrunSerde {
#[default]
Gap,
Stall,
Terminate,
}
impl ArchiveOverrunSerde {
pub fn name(&self) -> &'static str {
match self {
ArchiveOverrunSerde::Gap => "gap",
ArchiveOverrunSerde::Stall => "stall",
ArchiveOverrunSerde::Terminate => "terminate",
}
}
}
broadcast_common::impl_spec_display!(ArchiveOverrunSerde);
impl From<ArchiveOverrunSerde> for ArchiveOverrun {
fn from(v: ArchiveOverrunSerde) -> Self {
match v {
ArchiveOverrunSerde::Gap => ArchiveOverrun::Gap,
ArchiveOverrunSerde::Stall => ArchiveOverrun::StallIngest,
ArchiveOverrunSerde::Terminate => ArchiveOverrun::Terminate,
}
}
}
impl DvrConfig {
pub fn validate(&self) -> Result<(), String> {
if !self.enabled {
return Ok(());
}
if self.archive_root.is_empty() {
return Err("archive_root must be set when DVR is enabled".to_string());
}
if self.retention_periods == 0 && self.retention_bytes == 0 {
return Err(
"at least one of retention_periods or retention_bytes must be > 0 \
when DVR is enabled (unbounded growth would fill the disk)"
.to_string(),
);
}
Ok(())
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[non_exhaustive]
pub struct EitProgramme {
pub event_id: u16,
pub service_id: u16,
pub title: Option<String>,
pub start: Option<String>,
pub duration_secs: Option<u64>,
}
impl EitProgramme {
fn from_event(event: &dvb_si::tables::eit::EitEvent<'_>, service_id: u16) -> Self {
let title = event.descriptors.iter().find_map(|d| match d {
Ok(dvb_si::descriptors::AnyDescriptor::ShortEvent(se)) => {
Some(se.event_name.decode().into_owned())
}
_ => None,
});
let start = event.start_time().map(|dt| {
format!(
"{:04}-{:02}-{:02}T{:02}:{:02}:{:02}Z",
dt.year, dt.month, dt.day, dt.hour, dt.minute, dt.second
)
});
let duration_secs = event.duration().map(|d| d.as_secs());
EitProgramme {
event_id: event.event_id,
service_id,
title,
start,
duration_secs,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[non_exhaustive]
pub struct IndexEntry {
pub seq: u32,
pub start_pts_ns: u64,
pub byte_offset: u64,
pub byte_len: u64,
#[serde(default)]
pub duration_ns: u64,
#[serde(default)]
pub discontinuous: bool,
}
struct PeriodRecord {
num: u32,
file_bytes: u64,
}
pub struct DvrRecorder {
route_name: String,
archive_dir: PathBuf,
config: DvrConfig,
ext: String,
cursor: SegmentCursor,
current_file: Option<File>,
period: u32,
period_opened_at: Option<SystemTime>,
write_offset: u64,
init_len: u64,
last_init: Option<Vec<u8>>,
index: Vec<IndexEntry>,
periods: Vec<PeriodRecord>,
total_bytes: u64,
gaps: u64,
terminated: bool,
si_demux: Option<SiDemux>,
target_service_id: Option<u16>,
last_present_event_id: Option<u16>,
current_programme: Option<EitProgramme>,
si_carry: Vec<u8>,
}
impl DvrRecorder {
pub fn new(
route_name: String,
config: DvrConfig,
ext: &str,
trunk: &Arc<Trunk>,
) -> Result<Self, String> {
config.validate()?;
let archive_dir = PathBuf::from(&config.archive_root).join(&route_name);
let cursor = trunk.pin_segments(config.overrun.into());
let si_demux = config.dvb_service_id.map(|_| {
SiDemux::builder()
.dvb_si_pids(false)
.pid(Pid::new(EIT_PID))
.build()
});
let target_service_id = config.dvb_service_id;
Ok(DvrRecorder {
route_name,
archive_dir,
config,
ext: ext.to_string(),
cursor,
current_file: None,
period: 0,
period_opened_at: None,
write_offset: 0,
init_len: 0,
last_init: None,
index: Vec::new(),
periods: Vec::new(),
total_bytes: 0,
gaps: 0,
terminated: false,
si_demux,
target_service_id,
last_present_event_id: None,
current_programme: None,
si_carry: Vec::new(),
})
}
pub fn overrun_policy(&self) -> ArchiveOverrun {
self.config.overrun.into()
}
pub fn current_programme(&self) -> Option<&EitProgramme> {
self.current_programme.as_ref()
}
pub fn feed_si(&mut self, ts_bytes: &[u8]) -> Result<(), String> {
if self.si_demux.is_none() {
return Ok(());
}
let mut buf = std::mem::take(&mut self.si_carry);
buf.extend_from_slice(ts_bytes);
let mut offset = 0;
while offset + TS_PACKET_LEN <= buf.len() {
let events: Vec<SectionEvent> = self
.si_demux
.as_mut()
.expect("checked Some above")
.feed(&buf[offset..offset + TS_PACKET_LEN])
.collect();
for event in events {
self.handle_si_event(event)?;
}
offset += TS_PACKET_LEN;
}
self.si_carry = buf[offset..].to_vec();
Ok(())
}
fn handle_si_event(&mut self, event: SectionEvent) -> Result<(), String> {
let Some(target) = self.target_service_id else {
return Ok(());
};
let section = match event.table_section() {
Ok(AnyTableSection::EitSection(s)) => s,
_ => return Ok(()),
};
if section.kind != EitKind::PresentFollowingActual || section.service_id != target {
return Ok(());
}
let Some(present) = section
.events
.iter()
.find(|e| e.running_status == RunningStatus::Running)
else {
return Ok(());
};
let event_id = present.event_id;
let programme = EitProgramme::from_event(present, section.service_id);
match self.last_present_event_id {
None => {
self.last_present_event_id = Some(event_id);
self.current_programme = Some(programme);
if self.current_file.is_some() {
self.write_programme_sidecar()?;
}
}
Some(last) if last != event_id => {
tracing::info!(
route = %self.route_name,
old_event_id = last,
new_event_id = event_id,
"DVR EIT present/following transition — rolling period"
);
self.last_present_event_id = Some(event_id);
self.current_programme = Some(programme);
if self.current_file.is_some() {
let last_init = self.last_init.clone();
self.start_period(last_init.as_deref())?;
}
}
_ => {}
}
Ok(())
}
fn write_programme_sidecar(&self) -> Result<(), String> {
let Some(programme) = &self.current_programme else {
return Ok(());
};
let json = serde_json::to_vec(programme)
.map_err(|e| format!("serializing programme metadata: {e}"))?;
let tmp = self.archive_dir.join(".event.tmp");
let dst = self
.archive_dir
.join(format!("p{}.event.json", self.period));
fs::write(&tmp, &json).map_err(|e| format!("writing programme metadata: {e}"))?;
fs::rename(&tmp, &dst).map_err(|e| format!("renaming programme metadata: {e}"))?;
Ok(())
}
pub fn poll_and_persist(&mut self, init_bytes: Option<&[u8]>) -> Result<(), String> {
if self.ext == ".m4s" {
match (init_bytes, &self.last_init) {
(Some(new), None) => {
self.start_period(Some(new))?;
}
(Some(new), Some(old)) if new != old.as_slice() => {
tracing::info!(
route = %self.route_name,
"DVR init segment changed — rolling period"
);
self.start_period(Some(new))?;
}
_ => {}
}
}
if let Some(opened) = self.period_opened_at {
let elapsed = opened.elapsed().unwrap_or(Duration::ZERO).as_secs();
let limit = if self.config.period_duration_secs > 0 {
self.config.period_duration_secs
} else {
u64::MAX
};
if elapsed >= limit && self.current_file.is_some() {
tracing::info!(
route = %self.route_name,
period = self.period,
elapsed_secs = elapsed,
"DVR period duration reached — rolling period"
);
let last_init = self.last_init.clone();
self.start_period(last_init.as_deref())?;
}
}
while let Some(item) = self.cursor.poll() {
self.on_segment(&item)?;
}
Ok(())
}
fn start_period(&mut self, init_bytes: Option<&[u8]>) -> Result<(), String> {
if let Some(file) = self.current_file.take() {
drop(file);
}
if !self.index.is_empty() || self.write_offset > 0 {
let file_bytes = self.write_offset;
self.periods.push(PeriodRecord {
num: self.period,
file_bytes,
});
self.total_bytes += file_bytes;
self.period += 1;
}
let path = self.period_path();
fs::create_dir_all(&self.archive_dir).map_err(|e| format!("creating archive dir: {e}"))?;
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(&path)
.map_err(|e| format!("opening period file {}: {e}", path.display()))?;
self.write_offset = file.metadata().map(|m| m.len()).unwrap_or(0);
self.init_len = 0;
if let Some(init) = init_bytes
&& self.ext == ".m4s"
&& !init.is_empty()
&& self.write_offset == 0
{
file.write_all(init)
.map_err(|e| format!("writing init: {e}"))?;
file.flush().map_err(|e| format!("flushing init: {e}"))?;
self.init_len = init.len() as u64;
self.write_offset = self.init_len;
self.last_init = Some(init.to_vec());
tracing::debug!(
route = %self.route_name,
period = self.period,
len = init.len(),
"wrote fMP4 init at period file head"
);
}
self.current_file = Some(file);
self.period_opened_at = Some(SystemTime::now());
self.index.clear();
self.write_programme_sidecar()?;
tracing::info!(
route = %self.route_name,
period = self.period,
"DVR period file opened"
);
Ok(())
}
fn append_segment(&mut self, entry: &media_plane::trunk::SegmentEntry) -> Result<(), String> {
if self.ext == ".m4s" && self.last_init.is_none() {
tracing::debug!(
route = %self.route_name,
seq = entry.sequence_number,
"skipping segment — fMP4 init not yet available"
);
return Ok(());
}
if self.current_file.is_none() {
self.start_period(None)?;
}
let file = self.current_file.as_mut().expect("current_file set above");
file.write_all(&entry.bytes).map_err(|e| {
format!(
"writing segment {} to period file: {e}",
entry.sequence_number,
)
})?;
file.flush().map_err(|e| format!("flushing segment: {e}"))?;
let byte_offset = self.write_offset;
let byte_len = entry.bytes.len() as u64;
self.write_offset += byte_len;
self.index.push(IndexEntry {
seq: entry.sequence_number,
start_pts_ns: entry.timeline_position.as_nanos(),
byte_offset,
byte_len,
duration_ns: entry.duration.as_nanos() as u64,
discontinuous: entry.meta.discontinuous,
});
self.flush_index()?;
self.enforce_retention()?;
tracing::debug!(
route = %self.route_name,
period = self.period,
seq = entry.sequence_number,
bytes = byte_len,
offset = byte_offset,
"appended segment"
);
Ok(())
}
fn flush_index(&self) -> Result<(), String> {
let json =
serde_json::to_vec(&self.index).map_err(|e| format!("serializing index: {e}"))?;
let tmp = self.period_dir_path().join(".idx.tmp");
let dst = self.index_path();
fs::write(&tmp, &json).map_err(|e| format!("writing index: {e}"))?;
fs::rename(&tmp, &dst).map_err(|e| format!("renaming index: {e}"))?;
Ok(())
}
pub fn rebuild_index(&self) -> Result<Vec<IndexEntry>, String> {
let data = fs::read(self.period_path())
.map_err(|e| format!("reading period file for index rebuild: {e}"))?;
let data = &data[self.init_len as usize..];
if self.ext == ".ts" {
return Err("TS index rebuild not yet implemented".to_string());
}
let mut entries = Vec::new();
let mut offset: usize = 0;
let init_len = self.init_len as usize;
while offset + 8 <= data.len() {
let size = u32::from_be_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]) as usize;
let box_type = &data[offset + 4..offset + 8];
if size < 8 || offset + size > data.len() {
break;
}
if box_type == b"moof" {
let start = (init_len + offset) as u64;
let len = size as u64;
let mdat_offset = offset + size;
let mdat_size = if mdat_offset + 8 <= data.len() {
u32::from_be_bytes([
data[mdat_offset],
data[mdat_offset + 1],
data[mdat_offset + 2],
data[mdat_offset + 3],
]) as usize
} else {
0
};
let total_len = if mdat_offset + 8 <= data.len()
&& &data[mdat_offset + 4..mdat_offset + 8] == b"mdat"
&& mdat_size >= 8
{
(size + mdat_size) as u64
} else {
len
};
entries.push(IndexEntry {
seq: 0,
start_pts_ns: 0,
byte_offset: start,
byte_len: total_len,
duration_ns: 0,
discontinuous: false,
});
offset += total_len as usize;
} else {
offset += size;
}
}
Ok(entries)
}
fn period_path(&self) -> PathBuf {
self.archive_dir
.join(format!("p{}.{}", self.period, &self.ext[1..]))
}
fn index_path(&self) -> PathBuf {
self.archive_dir.join(format!("p{}.idx", self.period))
}
fn period_dir_path(&self) -> PathBuf {
self.archive_dir.clone()
}
fn enforce_retention(&mut self) -> Result<(), String> {
if self.config.retention_bytes > 0 {
while self.total_bytes > self.config.retention_bytes && !self.periods.is_empty() {
self.evict_oldest_period();
}
}
if self.config.retention_periods > 0 {
while self.periods.len() > self.config.retention_periods {
self.evict_oldest_period();
}
}
Ok(())
}
fn evict_oldest_period(&mut self) {
if self.periods.is_empty() {
return;
}
let num = self.periods[0].num;
let file_bytes = self.periods[0].file_bytes;
let ext_no_dot = &self.ext[1..];
let file_path = self.archive_dir.join(format!("p{}.{}", num, ext_no_dot));
let idx_path = self.archive_dir.join(format!("p{}.idx", num));
if let Err(e) = fs::remove_file(&file_path) {
tracing::warn!(
route = %self.route_name,
period = num,
path = %file_path.display(),
error = %e,
"failed to remove evicted period file"
);
}
if let Err(e) = fs::remove_file(&idx_path) {
tracing::warn!(
route = %self.route_name,
period = num,
path = %idx_path.display(),
error = %e,
"failed to remove evicted period index"
);
}
self.total_bytes = self.total_bytes.saturating_sub(file_bytes);
self.periods.remove(0);
tracing::debug!(
route = %self.route_name,
period = num,
bytes = file_bytes,
"evicted period (retention)"
);
}
}
impl SegmentEgress for DvrRecorder {
type Error = String;
fn on_segment(&mut self, item: &SegmentCursorItem) -> Result<(), Self::Error> {
if self.terminated {
return Ok(());
}
match item {
SegmentCursorItem::Segment(entry) if let Err(e) = self.append_segment(entry) => {
tracing::error!(
route = %self.route_name,
seq = entry.sequence_number,
error = %e,
"DVR append failed"
);
return Err(e);
}
SegmentCursorItem::Segment(entry) => {
tracing::debug!(
route = %self.route_name,
seq = entry.sequence_number,
bytes = entry.bytes.len(),
"appended segment"
);
}
SegmentCursorItem::Gap { skipped } => {
self.gaps += *skipped;
tracing::warn!(
route = %self.route_name,
skipped,
gaps_total = self.gaps,
"DVR gap: live ring evicted segment(s) before recorder consumed them"
);
}
SegmentCursorItem::Lagged { skipped } => {
self.gaps += *skipped;
tracing::warn!(
route = %self.route_name,
skipped,
"DVR unexpected Lagged (pinning cursor should not produce this)"
);
}
SegmentCursorItem::Terminated => {
self.terminated = true;
self.current_file.take();
tracing::info!(
route = %self.route_name,
"DVR recording terminated (ArchiveOverrun::Terminate)"
);
}
_ => {}
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use media_plane::trunk::{SegmentEntry, Trunk, TrunkConfig};
use std::num::NonZeroUsize;
use std::time::Duration;
fn nz(n: usize) -> NonZeroUsize {
NonZeroUsize::new(n).expect("test capacity must be non-zero")
}
fn trunk_config() -> TrunkConfig {
TrunkConfig::new(nz(4), nz(4), nz(8), nz(4), nz(4))
}
fn dummy_segment(seq: u32, byte: u8) -> SegmentEntry {
SegmentEntry::new(
bytes::Bytes::from(vec![byte; 16]),
seq,
Duration::from_secs(2),
broadcast_common::Timestamp::from_nanos(u64::from(seq) * 2_000_000_000),
transmux::SegmentMeta {
discontinuous: false,
},
)
}
fn dvr_config(tmp: &std::path::Path, retention: usize) -> DvrConfig {
DvrConfig {
enabled: true,
archive_root: tmp.to_string_lossy().to_string(),
retention_periods: retention,
retention_bytes: 0,
period_duration_secs: 3600, overrun: ArchiveOverrunSerde::Gap,
dvb_service_id: None,
}
}
fn temp_dir() -> std::path::PathBuf {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let dir = std::env::temp_dir().join(format!("multimux-dvr-{}-{}", std::process::id(), n));
let _ = std::fs::create_dir_all(&dir);
dir
}
fn cleanup_temp(dir: &std::path::Path) {
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn rebuild_index_does_not_panic_on_a_period_file_truncated_after_moof() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = dvr_config(&tmp, 5);
let mut recorder =
DvrRecorder::new("test".to_string(), cfg, ".m4s", &trunk).expect("recorder");
let init_bytes = b"INIT";
recorder.poll_and_persist(Some(init_bytes)).expect("init");
let mut moof = Vec::new();
moof.extend_from_slice(&16u32.to_be_bytes());
moof.extend_from_slice(b"moof");
moof.extend_from_slice(&[0u8; 8]);
moof.extend_from_slice(&[0u8; 5]);
writer.publish_segment(SegmentEntry::new(
bytes::Bytes::from(moof),
0,
Duration::from_secs(1),
broadcast_common::Timestamp::from_nanos(0),
transmux::SegmentMeta {
discontinuous: false,
},
));
recorder.poll_and_persist(None).expect("persist truncated");
let entries = recorder.rebuild_index().expect("rebuild must not panic");
assert!(
entries.len() <= 1,
"a truncated period yields at most the one recoverable moof, got {}",
entries.len()
);
cleanup_temp(&tmp);
}
#[test]
fn period_file_is_independently_playable_fmp4() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = dvr_config(&tmp, 5);
let mut recorder =
DvrRecorder::new("test".to_string(), cfg, ".m4s", &trunk).expect("recorder");
let init_bytes = b"FAKE_FTYP_MOOV_HEADER";
recorder
.poll_and_persist(Some(init_bytes))
.expect("poll with init");
let entry = dummy_segment(1, 0xAB);
let expected_bytes = entry.bytes.clone();
writer.publish_segment(entry);
recorder
.poll_and_persist(Some(init_bytes))
.expect("persist segment");
let period_path = tmp.join("test").join("p0.m4s");
assert!(period_path.exists(), "period file must exist");
let on_disk = std::fs::read(&period_path).expect("read period file");
assert!(
on_disk.starts_with(init_bytes),
"init must be at head of period file"
);
let seg_start = init_bytes.len();
assert_eq!(
&on_disk[seg_start..seg_start + expected_bytes.len()],
expected_bytes.as_ref(),
"segment bytes must match what was published"
);
let idx_path = tmp.join("test").join("p0.idx");
assert!(idx_path.exists(), "index file must exist");
let idx_data = std::fs::read(&idx_path).expect("read index");
let idx_entries: Vec<IndexEntry> = serde_json::from_slice(&idx_data).expect("parse index");
assert_eq!(idx_entries.len(), 1);
assert_eq!(idx_entries[0].seq, 1);
assert_eq!(idx_entries[0].byte_offset, seg_start as u64);
assert_eq!(idx_entries[0].byte_len, expected_bytes.len() as u64);
cleanup_temp(&tmp);
}
#[test]
fn index_offsets_are_exact() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = dvr_config(&tmp, 10);
let mut recorder =
DvrRecorder::new("test".to_string(), cfg.clone(), ".m4s", &trunk).expect("recorder");
let init = b"INIT_BYTES";
recorder.poll_and_persist(Some(init)).expect("poll init");
let segments: Vec<SegmentEntry> =
(1..=3).map(|seq| dummy_segment(seq, seq as u8)).collect();
let expected: Vec<bytes::Bytes> = segments.iter().map(|s| s.bytes.clone()).collect();
for seg in &segments {
writer.publish_segment(seg.clone());
}
recorder
.poll_and_persist(Some(init))
.expect("persist segments");
let idx_path = tmp.join("test").join("p0.idx");
let idx_data = std::fs::read(&idx_path).expect("read index");
let idx_entries: Vec<IndexEntry> = serde_json::from_slice(&idx_data).expect("parse index");
assert_eq!(idx_entries.len(), 3);
let period_path = tmp.join("test").join("p0.m4s");
let on_disk = std::fs::read(&period_path).expect("read period file");
for (i, entry) in idx_entries.iter().enumerate() {
let start = entry.byte_offset as usize;
let end = start + entry.byte_len as usize;
let slice = &on_disk[start..end];
assert_eq!(
slice,
expected[i].as_ref(),
"index offset {start}..{end} does not match segment {} bytes",
entry.seq
);
assert_eq!(entry.seq, (i + 1) as u32);
}
cleanup_temp(&tmp);
}
fn moof_mdat_segment(content: &[u8]) -> Vec<u8> {
let mdat_size = 8 + content.len();
let moof_size: u32 = 8; let total = moof_size as usize + mdat_size;
let mut buf = Vec::with_capacity(total);
buf.extend_from_slice(&moof_size.to_be_bytes());
buf.extend_from_slice(b"moof");
buf.extend_from_slice(&(mdat_size as u32).to_be_bytes());
buf.extend_from_slice(b"mdat");
buf.extend_from_slice(content);
buf
}
fn moof_segment_entry(seq: u32, content: &[u8]) -> SegmentEntry {
SegmentEntry::new(
bytes::Bytes::from(moof_mdat_segment(content)),
seq,
Duration::from_secs(2),
broadcast_common::Timestamp::from_nanos(u64::from(seq) * 2_000_000_000),
transmux::SegmentMeta {
discontinuous: false,
},
)
}
#[test]
fn index_rebuild_works() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = dvr_config(&tmp, 10);
let mut recorder =
DvrRecorder::new("test".to_string(), cfg, ".m4s", &trunk).expect("recorder");
let init = b"INIT";
recorder.poll_and_persist(Some(init)).expect("poll init");
for seq in 1..=3 {
let payload = vec![seq as u8; 32];
writer.publish_segment(moof_segment_entry(seq, &payload));
}
recorder.poll_and_persist(Some(init)).expect("persist");
let idx_path = tmp.join("test").join("p0.idx");
let original_data = std::fs::read(&idx_path).expect("read original index");
let original: Vec<IndexEntry> =
serde_json::from_slice(&original_data).expect("parse original");
assert_eq!(original.len(), 3);
let rebuilt = recorder.rebuild_index().expect("rebuild index");
assert_eq!(
rebuilt.len(),
original.len(),
"rebuilt index must have same entry count as original"
);
for (a, b) in original.iter().zip(rebuilt.iter()) {
assert_eq!(a.byte_offset, b.byte_offset, "byte_offset must match");
assert_eq!(a.byte_len, b.byte_len, "byte_len must match");
}
cleanup_temp(&tmp);
}
#[test]
fn init_change_rolls_the_file() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = dvr_config(&tmp, 10);
let mut recorder =
DvrRecorder::new("test".to_string(), cfg, ".m4s", &trunk).expect("recorder");
let init_a = b"INIT_A";
recorder
.poll_and_persist(Some(init_a))
.expect("poll init A");
writer.publish_segment(dummy_segment(1, 0xAA));
recorder
.poll_and_persist(Some(init_a))
.expect("persist seg 1");
let init_b = b"INIT_B";
recorder
.poll_and_persist(Some(init_b))
.expect("poll init B");
writer.publish_segment(dummy_segment(2, 0xBB));
recorder
.poll_and_persist(Some(init_b))
.expect("persist seg 2");
let p0_path = tmp.join("test").join("p0.m4s");
let p1_path = tmp.join("test").join("p1.m4s");
assert!(p0_path.exists(), "period 0 file must exist");
assert!(p1_path.exists(), "period 1 file must exist");
let p0_data = std::fs::read(&p0_path).expect("read p0");
assert!(p0_data.starts_with(init_a), "p0 must start with init A");
let p1_data = std::fs::read(&p1_path).expect("read p1");
assert!(p1_data.starts_with(init_b), "p1 must start with init B");
assert!(
p0_data.len() > init_a.len(),
"p0 must contain segments beyond init"
);
assert!(
p1_data.len() > init_b.len(),
"p1 must contain segments beyond init"
);
assert!(
tmp.join("test").join("p0.idx").exists(),
"p0.idx must exist"
);
assert!(
tmp.join("test").join("p1.idx").exists(),
"p1.idx must exist"
);
cleanup_temp(&tmp);
}
#[test]
fn period_rollover_and_retention() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = DvrConfig {
enabled: true,
archive_root: tmp.to_string_lossy().to_string(),
retention_periods: 2,
retention_bytes: 0,
period_duration_secs: 86400, overrun: ArchiveOverrunSerde::Gap,
dvb_service_id: None,
};
let mut recorder =
DvrRecorder::new("test".to_string(), cfg, ".m4s", &trunk).expect("recorder");
let init = b"INIT";
recorder.poll_and_persist(Some(init)).expect("poll init");
let inits: [&[u8]; 3] = [b"INIT_0", b"INIT_1", b"INIT_2"];
for (i, period_init) in inits.iter().enumerate() {
recorder
.poll_and_persist(Some(period_init))
.expect("poll init");
for seq_offset in 0..2 {
let seq = (i * 2 + seq_offset) as u32 + 1;
writer.publish_segment(dummy_segment(seq, seq as u8));
}
recorder
.poll_and_persist(Some(period_init))
.expect("persist");
}
let p0_path = tmp.join("test").join("p0.m4s");
let p1_path = tmp.join("test").join("p1.m4s");
let p2_path = tmp.join("test").join("p2.m4s");
assert!(!p0_path.exists(), "period 0 must be evicted by retention");
assert!(p1_path.exists(), "period 1 must survive");
assert!(p2_path.exists(), "period 2 must survive");
assert!(!tmp.join("test").join("p0.idx").exists(), "p0.idx evicted");
assert!(tmp.join("test").join("p1.idx").exists(), "p1.idx survives");
assert!(tmp.join("test").join("p2.idx").exists(), "p2.idx survives");
cleanup_temp(&tmp);
}
#[test]
fn recording_does_not_perturb_live_serving() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = dvr_config(&tmp, 10);
let mut recorder =
DvrRecorder::new("test".to_string(), cfg, ".m4s", &trunk).expect("recorder");
let mut live_cursor = trunk.subscribe_segments();
let init = b"INIT";
recorder.poll_and_persist(Some(init)).expect("poll init");
writer.publish_segment(dummy_segment(1, 0xAA));
writer.publish_segment(dummy_segment(2, 0xBB));
writer.publish_segment(dummy_segment(3, 0xCC));
recorder.poll_and_persist(Some(init)).expect("persist");
let mut live_seqs = Vec::new();
while let Some(item) = live_cursor.poll() {
if let SegmentCursorItem::Segment(entry) = item {
live_seqs.push(entry.sequence_number);
}
}
assert_eq!(
live_seqs,
vec![1, 2, 3],
"live-serving cursor must see all segments, unaffected by DVR"
);
let period_path = tmp.join("test").join("p0.m4s");
assert!(period_path.exists(), "period file must exist");
let on_disk = std::fs::read(&period_path).expect("read");
assert_eq!(
on_disk.len(),
init.len() + 3 * 16,
"period file must contain init + all three 16-byte segments"
);
cleanup_temp(&tmp);
}
#[test]
fn ingest_pipeline_records_and_demux_from_disk() {
use crate::source::ts_program::{TsIngestSession, test_support::build_ts_bytes};
use crate::source::{DriverProgress, advance_route};
use media_plane::DEFAULT_MAX_PROGRAMS;
use media_plane::ingress::{HandshakePolicy, IngestDriver};
let tmp = temp_dir();
let dvr_cfg = DvrConfig {
enabled: true,
archive_root: tmp.to_string_lossy().to_string(),
period_duration_secs: 0,
retention_periods: 8,
retention_bytes: 0,
overrun: ArchiveOverrunSerde::Gap,
dvb_service_id: None,
};
let route = crate::route::RouteHandle::new(1.0, 250, 64)
.with_name("ingest-test")
.with_dvr(dvr_cfg);
fn nz2(n: usize) -> NonZeroUsize {
NonZeroUsize::new(n).expect("n > 0")
}
let mut driver = IngestDriver::new(
TsIngestSession::new(),
TrunkConfig::new(nz2(8), nz2(8), nz2(64), nz2(8), nz2(8)),
HandshakePolicy::establish_by(broadcast_common::Timestamp::from_nanos(u64::MAX)),
DEFAULT_MAX_PROGRAMS,
);
let mut progress = DriverProgress::new();
let ts1 = build_ts_bytes(1, 0xAB, 90);
let ts2 = build_ts_bytes(1, 0xCD, 90);
driver.feed(&ts1, broadcast_common::Timestamp::ZERO);
advance_route(&driver, &route, &mut progress);
driver.feed(&ts2, broadcast_common::Timestamp::from_nanos(1));
advance_route(&driver, &route, &mut progress);
driver.finish();
advance_route(&driver, &route, &mut progress);
let archive_dir = tmp.join("ingest-test");
assert!(
archive_dir.exists(),
"archive directory must exist after recording; tmp contents: {:?}",
std::fs::read_dir(&tmp)
.map(|d| d
.filter_map(|e| e.ok())
.map(|e| e.path().display().to_string())
.collect::<Vec<_>>())
.unwrap_or_default()
);
let period_files: Vec<PathBuf> = {
let mut files: Vec<_> = std::fs::read_dir(&archive_dir)
.expect("read archive dir")
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| {
p.extension().is_some_and(|ext| ext == "m4s")
&& p.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n.starts_with('p'))
})
.collect();
files.sort();
files
};
assert!(
!period_files.is_empty(),
"at least one period file must exist"
);
let total_periods = period_files.len();
let mut total_tracks = 0usize;
let mut total_samples = 0usize;
for period_path in &period_files {
let data = std::fs::read(period_path).expect("read period file");
use broadcast_common::Unpackage;
let mut demux = transmux::Fmp4Demux::new();
let media = demux
.unpackage(&data)
.expect("Fmp4Demux must succeed on period file — init at head + fragments");
assert!(
!media.tracks.is_empty(),
"demux must recover at least one track from {}",
period_path.display()
);
for track in &media.tracks {
assert!(
!track.samples.is_empty(),
"track {} must have at least one decodable sample",
track.spec.track_id
);
eprintln!(
" track {} — {} samples",
track.spec.track_id,
track.samples.len()
);
total_tracks += 1;
total_samples += track.samples.len();
}
}
assert!(total_tracks > 0, "must recover at least one track");
assert!(
total_samples > 0,
"must recover at least one decodable sample"
);
eprintln!(
"DVR pipeline test: {} periods, {} tracks, {} samples",
total_periods, total_tracks, total_samples
);
cleanup_temp(&tmp);
}
#[test]
fn real_dvbt_capture_records_and_replays_from_disk_only() {
let fixture_path = concat!(
env!("CARGO_MANIFEST_DIR"),
"/../private/fixtures/ts/france-tnt-dvbt-20s.ts"
);
if !std::path::Path::new(fixture_path).exists() {
eprintln!(
"SKIP real_dvbt_capture_records_and_replays_from_disk_only: \
private fixture not found at {fixture_path} \
(run `git submodule update --init private`)"
);
return;
}
use crate::source::ts_program::TsIngestSession;
use crate::source::{DriverProgress, advance_route};
use media_plane::DEFAULT_MAX_PROGRAMS;
use media_plane::ingress::{HandshakePolicy, IngestDriver};
let tmp = temp_dir();
let dvr_cfg = DvrConfig {
enabled: true,
archive_root: tmp.to_string_lossy().to_string(),
period_duration_secs: 5, retention_periods: 16,
retention_bytes: 0,
overrun: ArchiveOverrunSerde::Gap,
dvb_service_id: None,
};
let route = crate::route::RouteHandle::new(2.0, 250, 64)
.with_name("france-tnt")
.with_dvr(dvr_cfg);
fn nz2(n: usize) -> NonZeroUsize {
NonZeroUsize::new(n).expect("n > 0")
}
let mut driver = IngestDriver::new(
TsIngestSession::new(),
TrunkConfig::new(nz2(8), nz2(8), nz2(64), nz2(8), nz2(8)),
HandshakePolicy::establish_by(broadcast_common::Timestamp::from_nanos(u64::MAX)),
DEFAULT_MAX_PROGRAMS,
);
let ts_bytes = std::fs::read(fixture_path).expect("read fixture");
let n_packets = ts_bytes.len() / 188;
eprintln!(
" feeding {} TS packets in chunks of ~1316 bytes (7 packets)",
n_packets
);
let chunk_bytes = 1316;
let mut progress = DriverProgress::new();
let mut offset = 0usize;
while offset < ts_bytes.len() {
let end = (offset + chunk_bytes).min(ts_bytes.len());
let chunk = &ts_bytes[offset..end];
let t = broadcast_common::Timestamp::from_nanos((offset / 188) as u64 * 40_000);
driver.feed(chunk, t);
advance_route(&driver, &route, &mut progress);
offset = end;
}
driver.finish();
advance_route(&driver, &route, &mut progress);
let program_ids: Vec<_> = driver.programs().collect();
assert!(
!program_ids.is_empty(),
"at least one programme must be discovered from the real DVB-T capture"
);
eprintln!(
" programmes discovered: {} ({:?})",
program_ids.len(),
program_ids
);
let archive_dir = tmp.join("france-tnt");
assert!(
program_ids.len() >= 2,
"expected at least 2 programmes from the 5-service DVB-T capture, got {}",
program_ids.len()
);
assert!(
archive_dir.exists(),
"DVR archive directory should exist after MPTS ingest"
);
eprintln!(
" #906 FIXED: {} programme(s) discovered from a 5-service MPTS capture",
program_ids.len()
);
cleanup_temp(&tmp);
}
#[test]
fn hard_cap_rolls_on_clock_even_with_eit_configured_and_never_fed() {
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = DvrConfig {
enabled: true,
archive_root: tmp.to_string_lossy().to_string(),
retention_periods: 8,
retention_bytes: 0,
period_duration_secs: 2, overrun: ArchiveOverrunSerde::Gap,
dvb_service_id: Some(0x0601), };
let mut recorder =
DvrRecorder::new("hardcap".to_string(), cfg, ".ts", &trunk).expect("recorder");
writer.publish_segment(dummy_segment(1, 0xAA));
recorder.poll_and_persist(None).expect("persist seg 1");
let p0_path = tmp.join("hardcap").join("p0.ts");
assert!(p0_path.exists(), "period 0 must exist");
recorder.period_opened_at = Some(SystemTime::now() - Duration::from_secs(3));
writer.publish_segment(dummy_segment(2, 0xBB));
recorder
.poll_and_persist(None)
.expect("persist seg 2 — hard cap must roll first");
let p1_path = tmp.join("hardcap").join("p1.ts");
assert!(
p1_path.exists(),
"hard cap must roll to a new period even though dvb_service_id is \
set and no EIT was ever fed"
);
cleanup_temp(&tmp);
}
#[test]
fn eit_transition_rolls_period_using_real_fixture_following_event() {
use broadcast_common::{Parse as WireParse, Serialize as WireSerialize};
const TF1_SERVICE_ID: u16 = 0x0601;
let fixture_path = concat!(
env!("CARGO_MANIFEST_DIR"),
"/../fixtures/dvb-si/tnt-5w-12732v-isi6-10s.ts"
);
let ts_bytes = std::fs::read(fixture_path).expect("read real DVB-T fixture");
let mut discover = SiDemux::builder()
.dvb_si_pids(false)
.pid(Pid::new(EIT_PID))
.build();
let mut genuine_section_bytes: Vec<bytes::Bytes> = Vec::new();
for chunk in ts_bytes.chunks_exact(TS_PACKET_LEN) {
for event in discover.feed(chunk) {
if let Ok(AnyTableSection::EitSection(section)) = event.table_section()
&& section.kind == EitKind::PresentFollowingActual
&& section.service_id == TF1_SERVICE_ID
{
genuine_section_bytes.push(event.bytes().clone());
}
}
}
assert!(
!genuine_section_bytes.is_empty(),
"fixture must carry a genuine TF1 EIT p/f actual section — see PROVENANCE"
);
let genuine_sections: Vec<_> = genuine_section_bytes
.iter()
.map(|b| {
dvb_si::tables::eit::EitSection::parse(b)
.expect("parse genuine TF1 EIT p/f section")
})
.collect();
let genuine = &genuine_sections[0];
let present = genuine_sections
.iter()
.flat_map(|s| s.events.iter())
.find(|e| e.running_status == RunningStatus::Running)
.expect("genuine sections must have a present (running) event");
let mut following = genuine_sections
.iter()
.flat_map(|s| s.events.iter())
.find(|e| e.running_status != RunningStatus::Running)
.cloned()
.expect("genuine sections must have a following (not-running) event");
assert_eq!(present.event_id, 0x7857);
assert_eq!(following.event_id, 0x7858);
let expected_new_title = following.descriptors.iter().find_map(|d| match d {
Ok(dvb_si::descriptors::AnyDescriptor::ShortEvent(se)) => {
Some(se.event_name.decode().into_owned())
}
_ => None,
});
assert_eq!(
expected_new_title.as_deref(),
Some("Plus belle la vie, encore...")
);
let expected_duration_secs = following.duration().map(|d| d.as_secs());
let tmp = temp_dir();
let trunk = Trunk::new(trunk_config());
let writer = trunk.segment_writer().expect("segment writer");
let cfg = DvrConfig {
enabled: true,
archive_root: tmp.to_string_lossy().to_string(),
retention_periods: 8,
retention_bytes: 0,
period_duration_secs: 3600,
overrun: ArchiveOverrunSerde::Gap,
dvb_service_id: Some(TF1_SERVICE_ID),
};
let mut recorder =
DvrRecorder::new("tf1".to_string(), cfg, ".ts", &trunk).expect("recorder");
recorder.feed_si(&ts_bytes).expect("feed genuine fixture");
assert_eq!(
recorder.current_programme().map(|p| p.event_id),
Some(0x7857),
"baseline present event must be the genuine TF1 present event"
);
writer.publish_segment(dummy_segment(1, 0xAA));
recorder.poll_and_persist(None).expect("persist seg 1");
assert!(
tmp.join("tf1").join("p0.ts").exists(),
"period 0 must exist"
);
let p0_programme: EitProgramme = serde_json::from_slice(
&std::fs::read(tmp.join("tf1").join("p0.event.json")).expect("read p0.event.json"),
)
.expect("parse p0.event.json");
assert_eq!(p0_programme.event_id, 0x7857);
following.running_status = RunningStatus::Running;
let transitioned = dvb_si::tables::eit::EitSection {
kind: genuine.kind,
table_id: genuine.table_id,
service_id: genuine.service_id,
version_number: (genuine.version_number + 1) % 32,
current_next_indicator: true,
section_number: 0,
last_section_number: 0,
transport_stream_id: genuine.transport_stream_id,
original_network_id: genuine.original_network_id,
segment_last_section_number: 0,
last_table_id: genuine.table_id,
events: vec![following],
};
let mut section_buf = vec![0u8; WireSerialize::serialized_len(&transitioned)];
WireSerialize::serialize_into(&transitioned, &mut section_buf)
.expect("serialize transitioned section");
let mut packetiser = mpeg_ts::mux::SectionPacketiser::new(EIT_PID);
let packets = packetiser.packetise(&[§ion_buf]);
let mut packet_bytes = Vec::new();
for p in &packets {
packet_bytes.extend_from_slice(p);
}
recorder
.feed_si(&packet_bytes)
.expect("feed transitioned section");
assert_eq!(
recorder.current_programme().map(|p| p.event_id),
Some(0x7858),
"present event must now be the (formerly following) event"
);
assert!(
tmp.join("tf1").join("p1.ts").exists(),
"EIT p/f transition must roll to a new period file"
);
let p1_programme: EitProgramme = serde_json::from_slice(
&std::fs::read(tmp.join("tf1").join("p1.event.json")).expect("read p1.event.json"),
)
.expect("parse p1.event.json");
assert_eq!(p1_programme.event_id, 0x7858);
assert_eq!(p1_programme.service_id, TF1_SERVICE_ID);
assert_eq!(p1_programme.title, expected_new_title);
assert_eq!(p1_programme.duration_secs, expected_duration_secs);
cleanup_temp(&tmp);
}
}