use std::sync::Arc;
use crate::pp_log::{PpLog, pp_error, pp_info};
use ffmpeg_next as ffmpeg;
use thiserror::Error as ThisError;
use crate::{
buffer::MediaBuffer,
control::ControlMsg,
element::{Element, ElementType, Sink, Source, element_pp_log},
error::Result,
pad::SrcPad,
pool::UnboundObjectPool,
};
const POOL_SIZE: usize = 4;
#[derive(Debug, ThisError)]
pub enum ScalerError {
#[error("ffmpeg error: {0}")]
Ffmpeg(#[from] ffmpeg::Error),
#[error(
"Scaler only converts/resizes decoded Video frames, got a {0}; \
link it straight after a decoder, not a demuxer"
)]
UnsupportedBuffer(&'static str),
}
pub struct Scaler {
pp_log: PpLog,
name: Arc<str>,
dst_format: ffmpeg::format::Pixel,
dst_width: u32,
dst_height: u32,
flags: ffmpeg::software::scaling::Flags,
context: Option<ffmpeg::software::scaling::Context>,
pool: UnboundObjectPool<ffmpeg::frame::Video>,
pad: SrcPad,
}
unsafe impl Send for Scaler {}
impl Scaler {
pub fn new(
name: impl Into<String>,
dst_format: ffmpeg::format::Pixel,
dst_width: u32,
dst_height: u32,
flags: ffmpeg::software::scaling::Flags,
) -> Self {
let name: Arc<str> = name.into().into();
let pp_log = element_pp_log(ElementType::Scaler, &name, None);
pp_info!(
pp_log: &pp_log,
"created: dst_format={dst_format:?}, dst={dst_width}x{dst_height}"
);
let pad = SrcPad::new(format!("{name}_src"));
let pool = UnboundObjectPool::new(
POOL_SIZE,
move || ffmpeg::frame::Video::new(dst_format, dst_width, dst_height),
|_| {},
);
Self {
name,
pp_log,
dst_format,
dst_width,
dst_height,
flags,
context: None,
pool,
pad,
}
}
fn context_matches(&self, frame: &ffmpeg::frame::Video) -> bool {
match &self.context {
Some(context) => {
let input = context.input();
input.format == frame.format()
&& input.width == frame.width()
&& input.height == frame.height()
}
None => false,
}
}
}
impl Element for Scaler {
fn name(&self) -> Arc<str> {
self.name.clone()
}
fn element_type(&self) -> ElementType {
ElementType::Scaler
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Source for Scaler {
fn src_pads(&mut self) -> &mut [SrcPad] {
std::slice::from_mut(&mut self.pad)
}
}
impl Sink for Scaler {
fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
match buf {
MediaBuffer::Video(frame) => {
if !self.context_matches(&frame) {
match &mut self.context {
Some(context) => context.cached(
frame.format(),
frame.width(),
frame.height(),
self.dst_format,
self.dst_width,
self.dst_height,
self.flags,
),
None => {
self.context = Some(
ffmpeg::software::scaling::Context::get(
frame.format(),
frame.width(),
frame.height(),
self.dst_format,
self.dst_width,
self.dst_height,
self.flags,
)
.inspect_err(|error| {
pp_error!(self, "failed to build scaling context: {error}")
})
.map_err(ScalerError::from)?,
);
}
}
}
let mut output = self.pool.get();
self.context
.as_mut()
.expect("built or confirmed matching above")
.run(&frame, &mut output)
.inspect_err(|error| pp_error!(self, "scale failed: {error}"))
.map_err(ScalerError::from)?;
output.set_pts(frame.pts());
self.pad.push(MediaBuffer::Video(Arc::new(output)))
}
MediaBuffer::Eos => self.pad.push(MediaBuffer::Eos),
MediaBuffer::Packet(_) => {
pp_error!(self, "unsupported buffer: Packet");
Err(ScalerError::UnsupportedBuffer("Packet").into())
}
MediaBuffer::Audio(_) => {
pp_error!(self, "unsupported buffer: Audio");
Err(ScalerError::UnsupportedBuffer("Audio").into())
}
}
}
fn control(&mut self, msg: ControlMsg) -> Result<()> {
self.pad.control(msg)
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use super::*;
struct CapturingSink {
pp_log: PpLog,
received: Arc<Mutex<Vec<MediaBuffer>>>,
}
impl Element for CapturingSink {
fn name(&self) -> Arc<str> {
"capture".into()
}
fn element_type(&self) -> ElementType {
ElementType::Other
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Sink for CapturingSink {
fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
self.received.lock().unwrap().push(buf);
Ok(())
}
fn control(&mut self, _msg: ControlMsg) -> Result<()> {
Ok(())
}
}
fn video_frame(
format: ffmpeg::format::Pixel,
width: u32,
height: u32,
pts: i64,
) -> MediaBuffer {
let pool = UnboundObjectPool::new(
0,
move || ffmpeg::frame::Video::new(format, width, height),
|_| {},
);
let mut frame = pool.get();
frame.set_pts(Some(pts));
MediaBuffer::Video(Arc::new(frame))
}
fn new_scaler(
dst_format: ffmpeg::format::Pixel,
dst_width: u32,
dst_height: u32,
) -> (Scaler, Arc<Mutex<Vec<MediaBuffer>>>) {
let mut scaler = Scaler::new(
"scaler",
dst_format,
dst_width,
dst_height,
ffmpeg::software::scaling::Flags::BILINEAR,
);
let received = Arc::new(Mutex::new(Vec::new()));
scaler.src_pads()[0].link(Box::new(CapturingSink {
received: received.clone(),
pp_log: element_pp_log(ElementType::Other, "capture", None),
}));
(scaler, received)
}
#[test]
fn converts_pixel_format_and_size_while_preserving_pts() {
let (mut scaler, received) = new_scaler(ffmpeg::format::Pixel::RGB24, 80, 60);
scaler
.consume(video_frame(ffmpeg::format::Pixel::YUV420P, 160, 120, 4242))
.expect("scale must succeed");
let received = received.lock().unwrap();
assert_eq!(received.len(), 1);
let MediaBuffer::Video(frame) = &received[0] else {
panic!("expected a Video buffer");
};
assert_eq!(frame.format(), ffmpeg::format::Pixel::RGB24);
assert_eq!(frame.width(), 80);
assert_eq!(frame.height(), 60);
assert_eq!(frame.pts(), Some(4242));
}
#[test]
fn rebuilds_its_scaling_context_when_input_dimensions_change_mid_stream() {
let (mut scaler, received) = new_scaler(ffmpeg::format::Pixel::RGB24, 80, 60);
scaler
.consume(video_frame(ffmpeg::format::Pixel::YUV420P, 160, 120, 0))
.expect("first frame must scale");
scaler
.consume(video_frame(ffmpeg::format::Pixel::YUV420P, 320, 240, 1))
.expect("a differently-sized second frame must still scale, not reuse a stale context");
let received = received.lock().unwrap();
assert_eq!(received.len(), 2);
for buf in received.iter() {
let MediaBuffer::Video(frame) = buf else {
panic!("expected a Video buffer");
};
assert_eq!(frame.width(), 80);
assert_eq!(frame.height(), 60);
}
}
#[test]
fn eos_forwards_downstream() {
let (mut scaler, received) = new_scaler(ffmpeg::format::Pixel::RGB24, 80, 60);
scaler
.consume(MediaBuffer::Eos)
.expect("eos must forward cleanly");
assert!(matches!(
received.lock().unwrap().as_slice(),
[MediaBuffer::Eos]
));
}
#[test]
fn rejects_packet_and_audio_buffers_with_a_clean_error_instead_of_scaling_garbage() {
let (mut scaler, _received) = new_scaler(ffmpeg::format::Pixel::RGB24, 80, 60);
let packet = MediaBuffer::Packet(Arc::new(ffmpeg::Packet::empty()));
assert!(
scaler.consume(packet).is_err(),
"Packet must be rejected, not silently accepted"
);
let audio = MediaBuffer::Audio(Arc::new(ffmpeg::frame::Audio::empty()));
assert!(
scaler.consume(audio).is_err(),
"Audio must be rejected, not silently accepted"
);
}
}