use std::fs::{self, File, OpenOptions};
use std::io::Write;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use media_plane::egress::SegmentEgress;
use media_plane::trunk::{ArchiveOverrun, SegmentCursor, SegmentCursorItem, Trunk};
use serde::{Deserialize, Serialize};
use tracing;
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,
}
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)]
struct IndexEntry {
seq: u32,
start_pts_ns: u64,
byte_offset: u64,
byte_len: u64,
}
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,
}
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());
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,
})
}
pub fn overrun_policy(&self) -> ArchiveOverrun {
self.config.overrun.into()
}
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 {
if 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();
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,
});
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(())
}
#[cfg(test)]
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 + 4 <= 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,
});
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);
}
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,
}
}
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 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,
};
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,
};
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,
};
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);
}
}