use std::{
path::Path,
sync::{Arc, Mutex},
};
use crate::pp_log::{PpLog, pp_error};
use ffmpeg_next as ffmpeg;
use thiserror::Error as ThisError;
use crate::{
buffer::MediaBuffer,
contract::{InputContract, MediaKind, PortContract},
control::{ControlMsg, SeekRejectReason},
element::{Element, ElementType, Sink, element_pp_log},
error::Result,
};
#[derive(Debug, ThisError)]
pub enum FileMuxerError {
#[error("FileMuxer stream sinks only accept Packet or Eos buffers, got {0}")]
UnsupportedBuffer(&'static str),
#[error("ffmpeg error: {0}")]
Ffmpeg(#[from] ffmpeg::Error),
}
struct PendingStream {
name: Arc<str>,
input_time_base: ffmpeg::Rational,
kind: Option<MediaKind>,
}
pub struct FileMuxer {
output: ffmpeg::format::context::Output,
streams: Vec<PendingStream>,
}
impl FileMuxer {
pub fn create(path: impl AsRef<Path>) -> Result<Self> {
let output = ffmpeg::format::output(&path).map_err(FileMuxerError::from)?;
Ok(Self {
output,
streams: Vec::new(),
})
}
pub fn add_stream(
&mut self,
name: impl Into<String>,
parameters: ffmpeg::codec::Parameters,
time_base: ffmpeg::Rational,
) -> Result<()> {
let mut stream = self
.output
.add_stream(parameters.id())
.map_err(FileMuxerError::from)?;
let kind = MediaKind::packet_for(parameters.medium());
stream.set_time_base(time_base);
stream.set_parameters(parameters);
self.streams.push(PendingStream {
name: name.into().into(),
input_time_base: time_base,
kind,
});
Ok(())
}
pub fn open(mut self) -> Result<Vec<Box<dyn Sink>>> {
self.output.write_header().map_err(FileMuxerError::from)?;
let total = self.streams.len();
let shared = Arc::new(FileMuxerShared {
state: Mutex::new(MuxerState {
output: self.output,
done: 0,
finished: false,
}),
total,
});
Ok(self
.streams
.into_iter()
.enumerate()
.map(|(index, stream)| -> Box<dyn Sink> {
Box::new(FileMuxerStreamSink {
pp_log: element_pp_log(ElementType::FileMuxer, &stream.name, None),
name: stream.name,
shared: shared.clone(),
stream_index: index,
input_time_base: stream.input_time_base,
kind: stream.kind,
done: false,
})
})
.collect())
}
}
struct MuxerState {
output: ffmpeg::format::context::Output,
done: usize,
finished: bool,
}
struct FileMuxerShared {
state: Mutex<MuxerState>,
total: usize,
}
impl FileMuxerShared {
fn write_packet(
&self,
stream_index: usize,
input_time_base: ffmpeg::Rational,
packet: &ffmpeg::Packet,
) -> Result<()> {
let mut state = self.state.lock().unwrap();
if state.finished {
return Ok(());
}
let mut packet = packet.clone();
let output_time_base = state
.output
.stream(stream_index)
.expect("stream was added in FileMuxer::add_stream")
.time_base();
packet.rescale_ts(input_time_base, output_time_base);
packet.set_stream(stream_index);
packet.set_position(-1);
packet
.write_interleaved(&mut state.output)
.map_err(FileMuxerError::from)?;
Ok(())
}
fn finish_track(&self) -> Result<()> {
let mut state = self.state.lock().unwrap();
state.done += 1;
if state.finished || state.done < self.total {
return Ok(());
}
state.finished = true;
state.output.write_trailer().map_err(FileMuxerError::from)?;
Ok(())
}
}
pub struct FileMuxerStreamSink {
pp_log: PpLog,
name: Arc<str>,
shared: Arc<FileMuxerShared>,
stream_index: usize,
input_time_base: ffmpeg::Rational,
kind: Option<MediaKind>,
done: bool,
}
impl FileMuxerStreamSink {
fn finish(&mut self) -> Result<()> {
if self.done {
return Ok(());
}
self.done = true;
self.shared
.finish_track()
.inspect_err(|error| pp_error!(self, "write_trailer failed: {error}"))
}
}
impl Element for FileMuxerStreamSink {
fn name(&self) -> Arc<str> {
self.name.clone()
}
fn element_type(&self) -> ElementType {
ElementType::FileMuxer
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Sink for FileMuxerStreamSink {
fn input_contract(&self) -> InputContract {
match self.kind {
Some(kind) => InputContract::Fixed(PortContract::packet(kind)),
None => InputContract::Unknown,
}
}
fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
match buf {
MediaBuffer::Packet(packet) => self
.shared
.write_packet(self.stream_index, self.input_time_base, &packet)
.inspect_err(|error| pp_error!(self, "write_interleaved failed: {error}")),
MediaBuffer::Eos => self.finish(),
other => Err(FileMuxerError::UnsupportedBuffer(other.kind()).into()),
}
}
fn control(&mut self, msg: ControlMsg) -> Result<()> {
if let ControlMsg::CheckSeek(context) = &msg {
context.reject(
self.element_type(),
self.name(),
SeekRejectReason::ElementNotSeekable,
);
}
if msg == ControlMsg::Stop {
self.finish()?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::control::{SeekCheckContext, SeekRejectReason};
use crate::element::Source;
use crate::elements::{AudioCodec, SwAudioEncoder, SwAudioEncoderOptions};
fn open_aac_encoder(sample_rate: u32, channels: u16) -> SwAudioEncoder {
SwAudioEncoder::new(
"encoder",
SwAudioEncoderOptions {
codec: AudioCodec::Aac,
sample_rate,
channels,
time_base: ffmpeg::Rational::new(1, sample_rate as i32),
bit_rate: 64_000,
},
)
.expect("aac encoder must be available")
}
fn open_opus_encoder(sample_rate: u32, channels: u16) -> Option<SwAudioEncoder> {
SwAudioEncoder::new(
"encoder",
SwAudioEncoderOptions {
codec: AudioCodec::Opus,
sample_rate,
channels,
time_base: ffmpeg::Rational::new(1, sample_rate as i32),
bit_rate: 64_000,
},
)
.ok()
}
fn silent_frame(
sample_rate: u32,
channels: u16,
samples: usize,
pts: i64,
) -> ffmpeg::frame::Audio {
let mut frame = ffmpeg::frame::Audio::new(
ffmpeg::format::Sample::F32(ffmpeg::format::sample::Type::Packed),
samples,
ffmpeg::ChannelLayout::default(channels as i32),
);
frame.set_rate(sample_rate);
frame.set_pts(Some(pts));
frame.data_mut(0).fill(0);
frame
}
#[test]
fn seek_check_rejects_a_muxer_track() {
let encoder = open_aac_encoder(48000, 1);
let path = std::env::temp_dir().join(format!(
"file_muxer_seek_check_test_{}.mp4",
std::process::id()
));
let mut muxer = FileMuxer::create(&path).expect("create muxer");
muxer
.add_stream(
"audio",
encoder.parameters(),
ffmpeg::Rational::new(1, 48000),
)
.expect("add stream");
let mut sink = muxer.open().expect("open muxer").pop().unwrap();
let context = Arc::new(SeekCheckContext::new());
sink.control(ControlMsg::CheckSeek(Arc::clone(&context)))
.expect("check control");
let error = context.result().expect_err("muxer must reject seek");
assert_eq!(error.rejections().len(), 1);
assert_eq!(
error.rejections()[0].reason,
SeekRejectReason::ElementNotSeekable
);
sink.control(ControlMsg::Stop).expect("finalize muxer");
std::fs::remove_file(path).ok();
}
#[test]
fn single_track_still_produces_a_playable_file() {
let mut encoder = open_aac_encoder(48000, 1);
let dir = std::env::temp_dir();
let path = dir.join(format!("file_muxer_single_test_{}.mp4", std::process::id()));
let mut muxer = FileMuxer::create(&path).expect("the muxer must open");
muxer
.add_stream(
"audio",
encoder.parameters(),
ffmpeg::Rational::new(1, 48000),
)
.expect("add_stream must succeed");
let mut sinks = muxer.open().expect("open must write the header");
assert_eq!(sinks.len(), 1);
encoder.src_pads()[0].link(sinks.pop().unwrap());
for tick in 0..20i64 {
encoder
.consume(MediaBuffer::Audio(Arc::new(silent_frame(
48000,
1,
960,
tick * 960,
))))
.expect("consume must succeed");
}
encoder
.consume(MediaBuffer::Eos)
.expect("eos must flush cleanly");
drop(encoder);
let input = ffmpeg::format::input(&path).expect("muxed file must be readable back");
assert_eq!(input.streams().count(), 1);
std::fs::remove_file(&path).ok();
}
#[test]
fn dropping_every_sink_without_eos_or_stop_still_releases_the_file() {
let encoder = open_aac_encoder(48000, 1);
let dir = std::env::temp_dir();
let path = dir.join(format!("file_muxer_drop_test_{}.mp4", std::process::id()));
let mut muxer = FileMuxer::create(&path).expect("the muxer must open");
muxer
.add_stream(
"audio",
encoder.parameters(),
ffmpeg::Rational::new(1, 48000),
)
.expect("add_stream must succeed");
let sinks = muxer.open().expect("open must write the header");
drop(sinks);
drop(encoder);
std::fs::remove_file(&path)
.expect("file handle must be released once every sink is dropped");
}
#[test]
fn muxes_two_independent_tracks_without_finalizing_early() {
let mut encoder_a = open_aac_encoder(48000, 2);
let mut encoder_b = open_aac_encoder(44100, 1);
let dir = std::env::temp_dir();
let path = dir.join(format!("file_muxer_multi_test_{}.mp4", std::process::id()));
let mut muxer = FileMuxer::create(&path).expect("the muxer must open");
muxer
.add_stream("a", encoder_a.parameters(), ffmpeg::Rational::new(1, 48000))
.expect("add_stream a");
muxer
.add_stream("b", encoder_b.parameters(), ffmpeg::Rational::new(1, 44100))
.expect("add_stream b");
let mut sinks = muxer.open().expect("open must write the header");
assert_eq!(sinks.len(), 2);
let sink_b = sinks.pop().unwrap();
let sink_a = sinks.pop().unwrap();
encoder_a.src_pads()[0].link(sink_a);
encoder_b.src_pads()[0].link(sink_b);
for tick in 0..10i64 {
encoder_a
.consume(MediaBuffer::Audio(Arc::new(silent_frame(
48000,
2,
960,
tick * 960,
))))
.expect("consume must succeed");
}
encoder_a
.consume(MediaBuffer::Eos)
.expect("eos must flush cleanly");
for tick in 0..10i64 {
encoder_b
.consume(MediaBuffer::Audio(Arc::new(silent_frame(
44100,
1,
882,
tick * 882,
))))
.expect("consume must succeed");
}
encoder_b
.consume(MediaBuffer::Eos)
.expect("eos must flush cleanly");
drop(encoder_a);
drop(encoder_b);
let mut input = ffmpeg::format::input(&path).expect("muxed file must be readable back");
assert_eq!(input.streams().count(), 2, "expected two tracks");
let mut counts = [0usize; 2];
let mut packet = ffmpeg::Packet::empty();
while packet.read(&mut input).is_ok() {
counts[packet.stream()] += 1;
packet = ffmpeg::Packet::empty();
}
assert!(counts[0] > 0, "track a has no packets: {counts:?}");
assert!(
counts[1] > 0,
"track b has no packets: {counts:?} — trailer was written before track b finished"
);
std::fs::remove_file(&path).ok();
}
#[test]
fn the_container_follows_the_path_extension() {
let mut encoder = open_aac_encoder(48000, 1);
let dir = std::env::temp_dir();
let path = dir.join(format!("muxer_container_test_{}.mkv", std::process::id()));
let _ = std::fs::remove_file(&path);
let mut muxer = FileMuxer::create(&path).expect("the muxer must open a .mkv path");
muxer
.add_stream(
"audio",
encoder.parameters(),
ffmpeg::Rational::new(1, 48000),
)
.expect("add_stream must succeed");
let mut sinks = muxer.open().expect("open must write the header");
encoder.src_pads()[0].link(sinks.pop().expect("exactly one stream was added"));
for tick in 0..20i64 {
encoder
.consume(MediaBuffer::Audio(Arc::new(silent_frame(
48000,
1,
960,
tick * 960,
))))
.expect("consume must succeed");
}
encoder
.consume(MediaBuffer::Eos)
.expect("eos must flush cleanly");
drop(encoder);
let input = ffmpeg::format::input(&path).expect("the file must be readable");
let format = input.format();
assert!(
format.name().contains("matroska"),
"a .mkv path should have produced Matroska, got {:?}",
format.name()
);
assert_eq!(input.streams().count(), 1);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_video_track_carries_its_extradata_into_a_matroska_header() {
use crate::elements::{SwEncoder, SwEncoderOptions, VideoCodec};
let time_base = ffmpeg::Rational::new(1, 30);
let options = |codec| SwEncoderOptions {
codec,
width: 320,
height: 180,
time_base,
frame_rate: ffmpeg::Rational::new(30, 1),
bit_rate: 400_000,
gop_size: 30,
max_b_frames: None,
};
let Some(encoder) = [VideoCodec::OpenH264, VideoCodec::H264]
.into_iter()
.find_map(|codec| SwEncoder::new("video", options(codec)).ok())
else {
eprintln!("skipping: this FFmpeg build has no libopenh264 or libx264");
return;
};
let dir = std::env::temp_dir();
let path = dir.join(format!("file_muxer_mkv_video_{}.mkv", std::process::id()));
let _ = std::fs::remove_file(&path);
let mut muxer = FileMuxer::create(&path).expect("the muxer must open");
muxer
.add_stream("video", encoder.parameters(), time_base)
.expect("add_stream must succeed");
let sinks = muxer
.open()
.expect("Matroska must accept a video track's header");
drop(sinks);
std::fs::remove_file(&path).ok();
}
#[test]
fn opus_carries_its_own_header_into_matroska() {
let Some(mut encoder) = open_opus_encoder(48_000, 2) else {
eprintln!("skipping: this FFmpeg build has no libopus");
return;
};
let dir = std::env::temp_dir();
let path = dir.join(format!("file_muxer_opus_mkv_{}.mkv", std::process::id()));
let _ = std::fs::remove_file(&path);
let mut muxer = FileMuxer::create(&path).expect("the muxer must open");
muxer
.add_stream(
"audio",
encoder.parameters(),
ffmpeg::Rational::new(1, 48_000),
)
.expect("add_stream must succeed");
let mut sinks = muxer
.open()
.expect("Matroska must accept an Opus track's header");
encoder.src_pads()[0].link(sinks.pop().expect("exactly one stream was added"));
for tick in 0..20i64 {
encoder
.consume(MediaBuffer::Audio(Arc::new(silent_frame(
48_000,
2,
960,
tick * 960,
))))
.expect("consume must succeed");
}
encoder
.consume(MediaBuffer::Eos)
.expect("eos must flush cleanly");
drop(encoder);
let input = ffmpeg::format::input(&path).expect("the file must be readable");
assert!(
input.format().name().contains("matroska"),
"got {:?}",
input.format().name()
);
let stream = input
.streams()
.best(ffmpeg::media::Type::Audio)
.expect("the file must hold the audio track it was given");
assert_eq!(
stream.parameters().id(),
ffmpeg::codec::Id::OPUS,
"the track must read back as Opus, not as whatever the header defaulted to"
);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_reordered_stream_keeps_its_decode_order_through_the_muxer() {
use crate::elements::FileDemuxer;
use crate::pipeline::Pipeline;
let fixture = crate::test_support::synthesize_reordered("reordered", 3.0);
let source = fixture.path.to_string_lossy().into_owned();
let arriving = timestamps(&source);
assert!(
arriving.iter().any(|(dts, pts)| dts != pts),
"the fixture carries no reordering to preserve: {:?}",
&arriving[..arriving.len().min(8)]
);
let path = std::env::temp_dir().join("media-pp-reordered-remux.mp4");
let _ = std::fs::remove_file(&path);
let (demuxer, streams) = FileDemuxer::open("demuxer", &source).expect("open the fixture");
let video = streams
.iter()
.find(|stream| stream.kind == ffmpeg::media::Type::Video)
.expect("the fixture has video")
.index;
let parameters = demuxer.stream_parameters(video).expect("video parameters");
let time_base = demuxer.stream_time_base(video).expect("video time base");
let mut muxer = FileMuxer::create(&path).expect("create the remux");
muxer
.add_stream("video", parameters, time_base)
.expect("add the video stream");
let sink = muxer
.open()
.expect("write the header")
.pop()
.expect("one stream was added");
let pipeline = Pipeline::new("remux", demuxer, move |source, context| {
let branch = context.branch().to(sink)?;
context.attach(source, video, branch)?;
Ok(())
})
.expect("wire the remux");
pipeline.run().expect("run the remux");
for event in pipeline.bus().iter() {
if matches!(event, crate::bus::BusEvent::Eos { .. }) {
break;
}
}
pipeline.stop();
let written = timestamps(&path.to_string_lossy());
assert_eq!(
written.len(),
arriving.len(),
"every packet that came in has to come out"
);
assert!(
written.iter().any(|(dts, pts)| dts != pts),
"the reordering did not survive being written"
);
assert!(
written.windows(2).all(|pair| pair[0].0 <= pair[1].0),
"decode order is what `dts` is for and it has to keep rising: {:?}",
&written[..written.len().min(8)]
);
assert_eq!(
written, arriving,
"a remux copies packets; it does not restamp them"
);
std::fs::remove_file(&path).ok();
}
fn timestamps(path: &str) -> Vec<(i64, i64)> {
let mut input = ffmpeg::format::input(path).expect("the file opens");
let video = input
.streams()
.find(|stream| stream.parameters().medium() == ffmpeg::media::Type::Video)
.expect("it has video")
.index();
input
.packets()
.filter(|(stream, _)| stream.index() == video)
.filter_map(|(_, packet)| Some((packet.dts()?, packet.pts()?)))
.collect()
}
#[derive(Debug, Clone, Copy)]
struct Shape {
video_frames: usize,
video_seconds: f64,
audio_seconds: f64,
video_start: f64,
audio_start: f64,
}
impl Shape {
fn video_fps(self) -> f64 {
self.video_frames as f64 / self.video_seconds
}
}
fn shape_of(path: &str) -> Shape {
let input = ffmpeg::format::input(path).expect("the file opens");
let stream = |medium| {
input
.streams()
.find(|stream| stream.parameters().medium() == medium)
.unwrap_or_else(|| panic!("{path} carries no {medium:?}"))
};
let seconds = |stream: &ffmpeg::format::stream::Stream<'_>| {
stream.duration() as f64 * f64::from(stream.time_base())
};
let start = |stream: &ffmpeg::format::stream::Stream<'_>| {
stream.start_time() as f64 * f64::from(stream.time_base())
};
let video = stream(ffmpeg::media::Type::Video);
let audio = stream(ffmpeg::media::Type::Audio);
Shape {
video_frames: video.frames() as usize,
video_seconds: seconds(&video),
audio_seconds: seconds(&audio),
video_start: start(&video),
audio_start: start(&audio),
}
}
#[test]
fn a_transcoded_file_keeps_the_shape_of_what_went_in() {
use crate::elements::{FileDemuxer, SwDecoder, SwEncoder, SwEncoderOptions, VideoCodec};
use crate::pipeline::Pipeline;
let fixture = crate::test_support::synthesize("transcode-shape", 4.0, 44_100);
let source_path = fixture.path.to_string_lossy().into_owned();
let arriving = shape_of(&source_path);
let path = std::env::temp_dir().join("media-pp-transcode-shape.mp4");
let _ = std::fs::remove_file(&path);
let (demuxer, streams) = FileDemuxer::open("demuxer", &source_path).expect("open");
let index = |medium| {
streams
.iter()
.find(|stream| stream.kind == medium)
.unwrap_or_else(|| panic!("the fixture carries no {medium:?}"))
.index
};
let video = index(ffmpeg::media::Type::Video);
let audio = index(ffmpeg::media::Type::Audio);
let video_decoder = SwDecoder::new(
"video-decoder",
demuxer.stream_parameters(video).expect("video parameters"),
)
.expect("open the video decoder");
let audio_decoder = SwDecoder::new(
"audio-decoder",
demuxer.stream_parameters(audio).expect("audio parameters"),
)
.expect("open the audio decoder");
let width = 320;
let height = 240;
let video_time_base = demuxer.stream_time_base(video).expect("video time base");
let video_encoder = SwEncoder::new(
"video-encoder",
SwEncoderOptions {
codec: VideoCodec::OpenH264,
width,
height,
time_base: video_time_base,
frame_rate: ffmpeg::Rational::new(30, 1),
bit_rate: 800_000,
gop_size: 30,
max_b_frames: None,
},
)
.expect("open the video encoder");
let audio_encoder = open_aac_encoder(48_000, 2);
let mut muxer = FileMuxer::create(&path).expect("create the output");
muxer
.add_stream("video", video_encoder.parameters(), video_time_base)
.expect("add the video stream");
muxer
.add_stream(
"audio",
audio_encoder.parameters(),
audio_encoder.time_base(),
)
.expect("add the audio stream");
let mut sinks = muxer.open().expect("write the header");
let audio_sink = sinks.pop().expect("audio was added second");
let video_sink = sinks.pop().expect("video was added first");
let scaler = crate::elements::SwScaler::new(
"to-yuv",
ffmpeg::format::Pixel::YUV420P,
width,
height,
ffmpeg::software::scaling::Flags::BILINEAR,
);
let pipeline = Pipeline::new("transcode", demuxer, move |source, context| {
let picture = context
.branch()
.pipe(video_decoder)
.pipe(scaler)
.pipe(video_encoder)
.to(video_sink)?;
context.attach(source, video, picture)?;
let sound = context
.branch()
.pipe(audio_decoder)
.pipe(audio_encoder)
.to(audio_sink)?;
context.attach(source, audio, sound)?;
Ok(())
})
.expect("wire the transcode");
pipeline.run().expect("run the transcode");
for event in pipeline.bus().iter() {
if matches!(event, crate::bus::BusEvent::Eos { .. }) {
break;
}
}
pipeline.stop();
let written = shape_of(&path.to_string_lossy());
eprintln!("SHAPE in : {arriving:?} fps={:.3}", arriving.video_fps());
eprintln!("SHAPE out: {written:?} fps={:.3}", written.video_fps());
assert!(
(written.video_fps() - arriving.video_fps()).abs() < 0.5,
"the picture came out at a different rate: {:.3} in, {:.3} out",
arriving.video_fps(),
written.video_fps()
);
assert!(
(written.video_seconds - arriving.video_seconds).abs() < 0.25,
"the picture came out a different length: {:.3}s in, {:.3}s out",
arriving.video_seconds,
written.video_seconds
);
assert!(
(written.audio_seconds - written.video_seconds).abs() < 0.25,
"the sound and the picture came out different lengths: \
{:.3}s of sound against {:.3}s of picture",
written.audio_seconds,
written.video_seconds
);
assert!(
(written.audio_start - written.video_start).abs()
<= (arriving.audio_start - arriving.video_start).abs() + 0.05,
"the two tracks no longer start where they did: \
in {:.3}s/{:.3}s, out {:.3}s/{:.3}s",
arriving.video_start,
arriving.audio_start,
written.video_start,
written.audio_start
);
std::fs::remove_file(&path).ok();
}
}