use std::{
path::PathBuf,
sync::{Arc, Mutex},
time::{Duration, Instant},
};
use crate::pp_log::{PpLog, pp_info};
use ffmpeg_next as ffmpeg;
use super::mp4_muxer::{Mp4Muxer, Mp4MuxerError};
use crate::{
buffer::MediaBuffer,
control::ControlMsg,
element::{Element, ElementType, Sink, element_pp_log},
error::Result,
};
#[derive(Debug, Clone, Copy)]
pub enum SegmentPolicy {
Duration(Duration),
}
struct StreamDef {
name: Arc<str>,
parameters: ffmpeg::codec::Parameters,
time_base: ffmpeg::Rational,
is_video: bool,
}
pub struct SegmentedMp4Muxer {
policy: SegmentPolicy,
naming: Box<dyn FnMut(u64) -> PathBuf + Send>,
streams: Vec<StreamDef>,
}
impl SegmentedMp4Muxer {
pub fn create(
policy: SegmentPolicy,
naming: impl FnMut(u64) -> PathBuf + Send + 'static,
) -> Self {
Self {
policy,
naming: Box::new(naming),
streams: Vec::new(),
}
}
pub fn add_stream(
&mut self,
name: impl Into<String>,
parameters: ffmpeg::codec::Parameters,
time_base: ffmpeg::Rational,
) {
let is_video = parameters.medium() == ffmpeg::media::Type::Video;
self.streams.push(StreamDef {
name: name.into().into(),
parameters,
time_base,
is_video,
});
}
pub fn open(mut self) -> Result<Vec<Box<dyn Sink>>> {
let names: Vec<Arc<str>> = self.streams.iter().map(|s| s.name.clone()).collect();
let path = (self.naming)(0);
let current_sinks = open_segment(&self.streams, path)?;
let group = Arc::new(SegmentGroup {
policy: self.policy,
naming: Mutex::new(self.naming),
state: Mutex::new(GroupState {
streams: self.streams,
current_sinks,
segment_index: 0,
segment_started: Instant::now(),
}),
});
Ok(names
.into_iter()
.enumerate()
.map(|(index, name)| -> Box<dyn Sink> {
Box::new(SegmentedTrackSink {
pp_log: element_pp_log(ElementType::SegmentedMp4Muxer, &name, None),
name,
track_index: index,
group: group.clone(),
})
})
.collect())
}
}
fn open_segment(streams: &[StreamDef], path: PathBuf) -> Result<Vec<Box<dyn Sink>>> {
let mut muxer = Mp4Muxer::create(&path)?;
for stream in streams {
muxer.add_stream(
stream.name.to_string(),
stream.parameters.clone(),
stream.time_base,
)?;
}
muxer.open()
}
struct GroupState {
streams: Vec<StreamDef>,
current_sinks: Vec<Box<dyn Sink>>,
segment_index: u64,
segment_started: Instant,
}
struct SegmentGroup {
policy: SegmentPolicy,
naming: Mutex<Box<dyn FnMut(u64) -> PathBuf + Send>>,
state: Mutex<GroupState>,
}
impl SegmentGroup {
fn consume_packet(
&self,
track_index: usize,
packet: Arc<ffmpeg::Packet>,
pp_log: &PpLog,
) -> Result<()> {
let mut state = self.state.lock().unwrap();
let SegmentPolicy::Duration(due_after) = self.policy;
let mut rotated_to = None;
if state.segment_started.elapsed() >= due_after {
let has_video = state.streams.iter().any(|s| s.is_video);
let this_is_video = state.streams[track_index].is_video;
let should_cut = if has_video {
this_is_video && packet.is_key()
} else {
true
};
if should_cut {
for sink in state.current_sinks.iter_mut() {
sink.control(ControlMsg::Stop)?;
}
let index = state.segment_index + 1;
let path = (self.naming.lock().unwrap())(index);
state.current_sinks = open_segment(&state.streams, path)?;
state.segment_index = index;
state.segment_started = Instant::now();
rotated_to = Some(index);
}
}
let result = state.current_sinks[track_index].consume(MediaBuffer::Packet(packet));
drop(state);
if let Some(index) = rotated_to {
pp_info!(pp_log: pp_log, "rotated segment_index={index}");
}
result
}
fn finish_eos(&self, track_index: usize) -> Result<()> {
let mut state = self.state.lock().unwrap();
state.current_sinks[track_index].consume(MediaBuffer::Eos)
}
fn finish_stop(&self, track_index: usize) -> Result<()> {
let mut state = self.state.lock().unwrap();
state.current_sinks[track_index].control(ControlMsg::Stop)
}
}
struct SegmentedTrackSink {
pp_log: PpLog,
name: Arc<str>,
track_index: usize,
group: Arc<SegmentGroup>,
}
impl Element for SegmentedTrackSink {
fn name(&self) -> Arc<str> {
self.name.clone()
}
fn element_type(&self) -> ElementType {
ElementType::SegmentedMp4Muxer
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Sink for SegmentedTrackSink {
fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
match buf {
MediaBuffer::Packet(packet) => {
self.group
.consume_packet(self.track_index, packet, &self.pp_log)
}
MediaBuffer::Eos => self.group.finish_eos(self.track_index),
other => Err(Mp4MuxerError::UnsupportedBuffer(other.kind()).into()),
}
}
fn control(&mut self, msg: ControlMsg) -> Result<()> {
if msg == ControlMsg::Stop {
self.group.finish_stop(self.track_index)?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
elements::{SwEncoder, SwEncoderOptions, TestVideoOptions, TestVideoSource, VideoCodec},
pipeline::Pipeline,
};
#[test]
fn rotates_into_multiple_valid_keyframe_aligned_segments() {
let video_options = TestVideoOptions {
width: 160,
height: 120,
framerate: ffmpeg::Rational::new(15, 1),
};
let video_source = TestVideoSource::new("video", video_options);
let time_base = video_source.time_base();
let encoder = SwEncoder::new(
"encoder",
SwEncoderOptions {
codec: VideoCodec::OpenH264,
width: video_options.width,
height: video_options.height,
time_base,
frame_rate: video_options.framerate,
bit_rate: 200_000,
gop_size: 8,
},
)
.expect("openh264 encoder must be available");
let dir = std::env::temp_dir();
let prefix = format!("segmented_mp4_test_{}", std::process::id());
let paths: Arc<Mutex<Vec<PathBuf>>> = Arc::new(Mutex::new(Vec::new()));
let recorded_paths = paths.clone();
let mut muxer = SegmentedMp4Muxer::create(
SegmentPolicy::Duration(Duration::from_millis(300)),
move |index| {
let path = dir.join(format!("{prefix}_{index:03}.mp4"));
recorded_paths.lock().unwrap().push(path.clone());
path
},
);
muxer.add_stream("video", encoder.parameters(), time_base);
let mut sinks = muxer.open().expect("open must succeed");
let sink = sinks.pop().expect("exactly one stream was added");
let pipeline = Pipeline::new("segmented-test", video_source, |source, ctx| {
let branch = ctx.branch().pipe(encoder).to(sink)?;
ctx.attach(source, 0, branch)?;
Ok(())
})
.expect("test pipeline wiring must succeed");
pipeline.run();
std::thread::sleep(Duration::from_secs(3));
pipeline.stop();
pipeline.bus().log_events();
let paths = paths.lock().unwrap().clone();
assert!(
paths.len() >= 2,
"expected at least 2 segments, got {}: {paths:?}",
paths.len()
);
for path in &paths {
let mut input = ffmpeg::format::input(path)
.unwrap_or_else(|error| panic!("segment {path:?} must be readable: {error}"));
assert_eq!(
input.streams().count(),
1,
"segment {path:?} should have exactly one stream"
);
let mut packet = ffmpeg::Packet::empty();
if packet.read(&mut input).is_err() {
panic!("segment {path:?} has no packets at all");
}
assert!(
packet.is_key(),
"segment {path:?}'s first packet must be a keyframe"
);
}
for path in &paths {
std::fs::remove_file(path).ok();
}
}
#[test]
fn rejects_a_buffer_type_the_wrapped_mp4_muxer_does_not_accept() {
let video_options = TestVideoOptions {
width: 160,
height: 120,
framerate: ffmpeg::Rational::new(15, 1),
};
let video_source = TestVideoSource::new("video", video_options);
let time_base = video_source.time_base();
let encoder = SwEncoder::new(
"encoder",
SwEncoderOptions {
codec: VideoCodec::OpenH264,
width: video_options.width,
height: video_options.height,
time_base,
frame_rate: video_options.framerate,
bit_rate: 200_000,
gop_size: 8,
},
)
.expect("openh264 encoder must be available");
let path = std::env::temp_dir().join(format!(
"segmented_mp4_reject_test_{}.mp4",
std::process::id()
));
let recorded_path = path.clone();
let mut muxer = SegmentedMp4Muxer::create(
SegmentPolicy::Duration(Duration::from_secs(3600)),
move |_index| recorded_path.clone(),
);
muxer.add_stream("video", encoder.parameters(), time_base);
let mut sinks = muxer.open().expect("open must succeed");
let mut sink = sinks.pop().expect("exactly one stream was added");
let error = sink
.consume(MediaBuffer::Audio(Arc::new(ffmpeg::frame::Audio::empty())))
.expect_err("an Audio buffer must be rejected, not silently dropped");
assert!(
matches!(
error,
crate::error::Error::Mp4MuxerError(Mp4MuxerError::UnsupportedBuffer("Audio"))
),
"unexpected error: {error:?}"
);
drop(sink);
std::fs::remove_file(&path).ok();
}
#[test]
fn old_segment_is_released_immediately_not_deferred_until_the_whole_recording_stops() {
let video_options = TestVideoOptions {
width: 160,
height: 120,
framerate: ffmpeg::Rational::new(15, 1),
};
let video_source = TestVideoSource::new("video", video_options);
let time_base = video_source.time_base();
let encoder = SwEncoder::new(
"encoder",
SwEncoderOptions {
codec: VideoCodec::OpenH264,
width: video_options.width,
height: video_options.height,
time_base,
frame_rate: video_options.framerate,
bit_rate: 200_000,
gop_size: 8, },
)
.expect("openh264 encoder must be available");
let dir = std::env::temp_dir();
let prefix = format!("segmented_mp4_release_test_{}", std::process::id());
let paths: Arc<Mutex<Vec<PathBuf>>> = Arc::new(Mutex::new(Vec::new()));
let recorded_paths = paths.clone();
let mut muxer = SegmentedMp4Muxer::create(
SegmentPolicy::Duration(Duration::from_millis(300)),
move |index| {
let path = dir.join(format!("{prefix}_{index:03}.mp4"));
recorded_paths.lock().unwrap().push(path.clone());
path
},
);
muxer.add_stream("video", encoder.parameters(), time_base);
let mut sinks = muxer.open().expect("open must succeed");
let sink = sinks.pop().expect("exactly one stream was added");
let pipeline = Pipeline::new("segmented-release-test", video_source, |source, ctx| {
let branch = ctx.branch().pipe(encoder).to(sink)?;
ctx.attach(source, 0, branch)?;
Ok(())
})
.expect("test pipeline wiring must succeed");
pipeline.run();
let waited = Instant::now();
loop {
if paths.lock().unwrap().len() >= 2 {
break;
}
assert!(
waited.elapsed() < Duration::from_secs(5),
"no rotation happened within 5s"
);
std::thread::sleep(Duration::from_millis(50));
}
let first_segment = paths.lock().unwrap()[0].clone();
let mut input = ffmpeg::format::input(&first_segment).unwrap_or_else(|error| {
panic!("segment 0 must already be readable while still recording segment 1: {error}")
});
let mut packet = ffmpeg::Packet::empty();
assert!(
packet.read(&mut input).is_ok(),
"segment 0 must have packets"
);
drop(input);
pipeline.stop();
pipeline.bus().log_events();
for path in paths.lock().unwrap().iter() {
std::fs::remove_file(path).ok();
}
}
}