use std::sync::Arc;
use transmux::ll_hls::LlHlsSegmenter;
use transmux::pipeline::{Sample, TrackSpec};
use crate::Result;
use crate::store::StreamStore;
const MOVIE_TIMESCALE: u32 = 90_000;
#[allow(async_fn_in_trait)]
pub trait SampleSource {
fn track_specs(&self) -> Vec<TrackSpec>;
async fn next_samples(&mut self) -> Result<Option<Vec<(u32, Sample)>>>;
}
impl SampleSource for crate::source::rtsp::RtspSession {
fn track_specs(&self) -> Vec<TrackSpec> {
crate::source::rtsp::RtspSession::track_specs(self)
}
async fn next_samples(&mut self) -> Result<Option<Vec<(u32, Sample)>>> {
crate::source::rtsp::RtspSession::next_samples(self).await
}
}
pub async fn run_pipeline<S: SampleSource>(
store: Arc<StreamStore>,
target_duration_secs: f64,
part_target_ms: u32,
mut source: S,
) -> Result<()> {
let specs = source.track_specs();
let mut seg = LlHlsSegmenter::with_part_target(
specs,
MOVIE_TIMESCALE,
target_duration_secs,
part_target_ms,
)?;
store.set_init(seg.init_segment()?);
while let Some(batch) = source.next_samples().await? {
for (track_id, sample) in batch {
seg.push(track_id, sample)?;
}
for part in seg.take_ready_parts() {
store.add_part(part);
}
for segment in seg.take_ready_segments() {
store.add_segment(segment);
}
}
seg.flush()?;
for part in seg.take_ready_parts() {
store.add_part(part);
}
for segment in seg.take_ready_segments() {
store.add_segment(segment);
}
Ok(())
}
#[doc(hidden)]
#[cfg(any(test, feature = "testsupport"))]
pub struct MockSource {
specs: Vec<TrackSpec>,
batches: std::vec::IntoIter<Vec<(u32, Sample)>>,
}
#[doc(hidden)]
#[cfg(any(test, feature = "testsupport"))]
impl MockSource {
pub fn new(specs: Vec<TrackSpec>, batches: Vec<Vec<(u32, Sample)>>) -> Self {
MockSource {
specs,
batches: batches.into_iter(),
}
}
}
#[cfg(any(test, feature = "testsupport"))]
impl SampleSource for MockSource {
fn track_specs(&self) -> Vec<TrackSpec> {
self.specs.clone()
}
async fn next_samples(&mut self) -> Result<Option<Vec<(u32, Sample)>>> {
Ok(self.batches.next())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::StreamStore;
use transmux::avc_config_from_sprop;
use transmux::pipeline::CodecConfig;
const SPROP: &str = "Z0IAKeKQFAe2AtwEBAaQeJEV,aM48gA==";
const VIDEO_TIMESCALE: u32 = 90_000;
const FRAME_DUR: u32 = VIDEO_TIMESCALE / 30;
fn video_track_spec() -> TrackSpec {
let config = avc_config_from_sprop(SPROP).expect("valid sprop");
TrackSpec::new(
1,
VIDEO_TIMESCALE,
CodecConfig::Avc {
config,
width: 0,
height: 0,
},
)
}
#[tokio::test]
async fn drives_source_through_segmenter_into_store() {
let store = Arc::new(StreamStore::new(1.0, 500, 8));
let specs = vec![video_track_spec()];
let mut batches = Vec::new();
for i in 0..90u32 {
let is_sync = i == 0 || i == 45;
let data = vec![0xAAu8.wrapping_add(i as u8); 32];
let sample = Sample::new(data, FRAME_DUR, is_sync, 0);
batches.push(vec![(1u32, sample)]);
}
let source = MockSource::new(specs, batches);
run_pipeline(store.clone(), 1.0, 500, source)
.await
.expect("pipeline runs to completion");
assert!(store.init_bytes().is_some(), "init segment stored");
let playlist = store.media_playlist_m3u8(1);
assert!(
playlist.contains("seg-") || playlist.contains("#EXT-X-PART"),
"playlist has landed media: {playlist}"
);
}
#[tokio::test]
async fn empty_batches_are_a_no_op() {
let store = Arc::new(StreamStore::new(1.0, 500, 8));
let specs = vec![video_track_spec()];
let source = MockSource::new(specs, vec![Vec::new(), Vec::new()]);
run_pipeline(store.clone(), 1.0, 500, source)
.await
.expect("pipeline tolerates empty batches");
assert!(store.init_bytes().is_some());
}
#[tokio::test]
async fn eos_flush_emits_buffered_tail_segment() {
let store = Arc::new(StreamStore::new(1.0, 500, 8));
let specs = vec![video_track_spec()];
let mut batches = Vec::new();
for i in 0..60u32 {
let is_sync = i == 0 || i == 45;
let data = vec![0xCCu8.wrapping_add(i as u8); 32];
let sample = Sample::new(data, FRAME_DUR, is_sync, 0);
batches.push(vec![(1u32, sample)]);
}
let source = MockSource::new(specs, batches);
run_pipeline(store.clone(), 1.0, 500, source)
.await
.expect("pipeline runs to completion");
assert!(store.init_bytes().is_some(), "init segment stored");
let playlist = store.media_playlist_m3u8(1);
assert!(
playlist.contains("seg-1-2"),
"seg.flush() must emit buffered tail as seg-1-2, got playlist: {}",
playlist
);
}
}