use crate::core::context::frame_source::{FrameSource, FrameSourceParams};
use crate::core::context::obj_pool::ObjPool;
use crate::core::context::{null_frame, FrameBox, FrameData};
use crate::core::scheduler::ffmpeg_scheduler::{
is_stopping, set_scheduler_error, wait_until_not_paused,
};
use crate::error::{AllocFrameError, Error};
use crate::util::ffmpeg_utils::av_err2str;
use crate::util::thread_synchronizer::{ThreadDoneGuard, ThreadSynchronizer};
use crossbeam_channel::{RecvTimeoutError, SendTimeoutError, Sender};
use ffmpeg_next::Frame;
use ffmpeg_sys_next::{av_frame_get_buffer, av_image_copy, av_image_fill_arrays, AVRational};
use log::{debug, error};
use std::ptr::null_mut;
use std::sync::atomic::AtomicUsize;
use std::sync::{Arc, Mutex};
use std::time::Duration;
pub(crate) fn frame_source_init(
index: usize,
frame_source: FrameSource,
frame_pool: ObjPool<Frame>,
scheduler_status: Arc<AtomicUsize>,
thread_sync: ThreadSynchronizer,
scheduler_result: Arc<Mutex<Option<crate::error::Result<()>>>>,
) -> crate::error::Result<()> {
thread_sync.thread_start();
let thread_done_guard = ThreadDoneGuard::adopt(
thread_sync.clone(),
scheduler_status.clone(),
scheduler_result.clone(),
);
let result = std::thread::Builder::new()
.name(format!("framesource{index}"))
.spawn(move || {
let _thread_done = thread_done_guard;
let ingress = frame_source.ingress;
let fg_sender = frame_source.fg_sender;
let params = frame_source.params;
let frame_pool = frame_pool;
let mut nb_frames: i64 = 0;
loop {
let result = ingress.recv_timeout(Duration::from_millis(100));
if is_stopping(wait_until_not_paused(&scheduler_status)) {
debug!("Frame source received end command, finishing.");
return;
}
let data = match result {
Ok(data) => data,
Err(RecvTimeoutError::Timeout) => continue,
Err(RecvTimeoutError::Disconnected) => break,
};
let frame = match build_video_frame(&frame_pool, ¶ms, &data, nb_frames) {
Ok(frame) => frame,
Err(e) => {
error!("Frame source failed to build a frame: {e}");
set_scheduler_error(&scheduler_status, &scheduler_result, e);
return;
}
};
nb_frames += 1;
let frame_box = FrameBox {
frame,
frame_data: frame_data_for(¶ms),
};
if !send_with_status_poll(&fg_sender, frame_box, &scheduler_status, &frame_pool) {
return;
}
}
let eof_marker = FrameBox {
frame: null_frame(),
frame_data: frame_data_for(¶ms),
};
send_with_status_poll(&fg_sender, eof_marker, &scheduler_status, &frame_pool);
debug!("Frame source finished after {nb_frames} frame(s).");
});
if let Err(e) = result {
error!("Frame source thread exited with error: {e}");
return Err(Error::FrameSourceThreadExited);
}
Ok(())
}
fn send_with_status_poll(
sender: &Sender<FrameBox>,
mut frame_box: FrameBox,
scheduler_status: &Arc<AtomicUsize>,
frame_pool: &ObjPool<Frame>,
) -> bool {
loop {
match sender.send_timeout(frame_box, Duration::from_millis(100)) {
Ok(()) => return true,
Err(SendTimeoutError::Timeout(returned)) => {
if is_stopping(wait_until_not_paused(scheduler_status)) {
debug!("Frame source received end command while sending.");
frame_pool.release(returned.frame);
return false;
}
frame_box = returned;
}
Err(SendTimeoutError::Disconnected(returned)) => {
debug!("Frame source: filtergraph receiver is gone.");
frame_pool.release(returned.frame);
return false;
}
}
}
}
fn frame_data_for(params: &FrameSourceParams) -> FrameData {
FrameData {
framerate: Some(AVRational {
num: params.fps_num,
den: params.fps_den,
}),
bits_per_raw_sample: 0,
input_stream_width: params.width,
input_stream_height: params.height,
subtitle_header: None,
fg_input_index: 0,
side_data: None,
}
}
fn build_video_frame(
frame_pool: &ObjPool<Frame>,
params: &FrameSourceParams,
data: &[u8],
pts: i64,
) -> crate::error::Result<Frame> {
let mut frame = frame_pool.get()?;
unsafe {
let f = frame.as_mut_ptr();
(*f).format = params.pix_fmt as i32;
(*f).width = params.width;
(*f).height = params.height;
let ret = av_frame_get_buffer(f, 0);
if ret < 0 {
error!("av_frame_get_buffer failed: {}", av_err2str(ret));
frame_pool.release(frame);
return Err(AllocFrameError::OutOfMemory.into());
}
let mut src_data: [*mut u8; 4] = [null_mut(); 4];
let mut src_linesize: [libc::c_int; 4] = [0; 4];
let ret = av_image_fill_arrays(
src_data.as_mut_ptr(),
src_linesize.as_mut_ptr(),
data.as_ptr(),
params.pix_fmt,
params.width,
params.height,
1,
);
if ret < 0 {
error!("av_image_fill_arrays failed: {}", av_err2str(ret));
frame_pool.release(frame);
return Err(Error::Bug);
}
av_image_copy(
(*f).data.as_ptr(),
(*f).linesize.as_ptr(),
src_data.as_ptr() as *const *const u8,
src_linesize.as_ptr(),
params.pix_fmt,
params.width,
params.height,
);
(*f).pts = pts;
(*f).duration = 1;
(*f).time_base = AVRational {
num: params.fps_den,
den: params.fps_num,
};
}
Ok(frame)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::scheduler::ffmpeg_scheduler::{frame_is_null, unref_frame};
use ffmpeg_sys_next::AVPixelFormat::{AV_PIX_FMT_GRAY8, AV_PIX_FMT_YUV420P};
use ffmpeg_sys_next::{av_frame_alloc, av_image_get_buffer_size};
fn test_new_frame() -> crate::error::Result<Frame> {
let f = unsafe { av_frame_alloc() };
assert!(!f.is_null(), "av_frame_alloc failed in test");
Ok(unsafe { Frame::wrap(f) })
}
fn test_pool() -> ObjPool<Frame> {
ObjPool::new(1, test_new_frame, unref_frame, frame_is_null).expect("frame pool")
}
fn params(pix_fmt: ffmpeg_sys_next::AVPixelFormat, w: i32, h: i32) -> FrameSourceParams {
FrameSourceParams {
width: w,
height: h,
pix_fmt,
fps_num: 30,
fps_den: 1,
}
}
fn tight_size(pix_fmt: ffmpeg_sys_next::AVPixelFormat, w: i32, h: i32) -> usize {
unsafe { av_image_get_buffer_size(pix_fmt, w, h, 1) as usize }
}
#[test]
fn fill_respects_linesize_for_odd_width() {
let pool = test_pool();
let p = params(AV_PIX_FMT_GRAY8, 65, 3);
let data: Vec<u8> = (0..tight_size(p.pix_fmt, 65, 3))
.map(|i| (i % 251) as u8)
.collect();
let frame = build_video_frame(&pool, &p, &data, 7).expect("build");
unsafe {
let f = frame.as_ptr();
assert!((*f).linesize[0] >= 65, "padded linesize expected");
for y in 0..3usize {
let row =
std::slice::from_raw_parts((*f).data[0].add(y * (*f).linesize[0] as usize), 65);
assert_eq!(row, &data[y * 65..y * 65 + 65], "row {y} content");
}
assert_eq!((*f).pts, 7);
assert_eq!((*f).duration, 1);
assert_eq!((*f).time_base.num, 1);
assert_eq!((*f).time_base.den, 30);
}
}
#[test]
fn fill_copies_all_planes_for_odd_yuv420p() {
let pool = test_pool();
let (w, h) = (65i32, 49i32);
let p = params(AV_PIX_FMT_YUV420P, w, h);
let (cw, ch) = (33usize, 25usize);
let y_size = (w * h) as usize;
let c_size = cw * ch;
let mut data = vec![0u8; tight_size(p.pix_fmt, w, h)];
assert_eq!(data.len(), y_size + 2 * c_size);
for (i, b) in data.iter_mut().enumerate() {
*b = (i * 7 % 253) as u8;
}
let frame = build_video_frame(&pool, &p, &data, 0).expect("build");
unsafe {
let f = frame.as_ptr();
let planes = [
(0usize, w as usize, h as usize, 0usize),
(1, cw, ch, y_size),
(2, cw, ch, y_size + c_size),
];
for (idx, pw, ph, base) in planes {
let ls = (*f).linesize[idx] as usize;
assert!(ls >= pw, "plane {idx} linesize");
for y in 0..ph {
let row = std::slice::from_raw_parts((*f).data[idx].add(y * ls), pw);
assert_eq!(
row,
&data[base + y * pw..base + y * pw + pw],
"plane {idx} row {y}"
);
}
}
}
}
#[test]
fn blocked_send_observes_terminal_status() {
use crate::core::scheduler::ffmpeg_scheduler::STATUS_END;
use std::sync::atomic::AtomicUsize;
use std::time::Instant;
let pool = test_pool();
let p = params(AV_PIX_FMT_GRAY8, 8, 2);
let boxed = |pool: &ObjPool<Frame>| FrameBox {
frame: pool.get().unwrap(),
frame_data: frame_data_for(&p),
};
let (tx, rx) = crossbeam_channel::bounded::<FrameBox>(1);
tx.send(boxed(&pool)).unwrap(); let status = Arc::new(AtomicUsize::new(STATUS_END));
let start = Instant::now();
let delivered = send_with_status_poll(&tx, boxed(&pool), &status, &pool);
assert!(!delivered, "terminal status must abort a blocked send");
assert!(
start.elapsed() < Duration::from_secs(10),
"the abort must land within a few poll intervals, took {:?}",
start.elapsed()
);
drop(rx);
}
#[test]
fn disconnected_filter_channel_fails_send() {
use crate::core::scheduler::ffmpeg_scheduler::STATUS_RUN;
use std::sync::atomic::AtomicUsize;
let pool = test_pool();
let p = params(AV_PIX_FMT_GRAY8, 8, 2);
let (tx, rx) = crossbeam_channel::bounded::<FrameBox>(1);
drop(rx);
let status = Arc::new(AtomicUsize::new(STATUS_RUN));
let frame_box = FrameBox {
frame: pool.get().unwrap(),
frame_data: frame_data_for(&p),
};
assert!(!send_with_status_poll(&tx, frame_box, &status, &pool));
}
#[test]
fn worker_parked_in_full_send_exits_on_terminal_status() {
use crate::core::scheduler::ffmpeg_scheduler::{STATUS_END, STATUS_RUN};
use crate::util::thread_synchronizer::ThreadSynchronizer;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Mutex;
use std::time::Instant;
let pool = ObjPool::new(4, test_new_frame, unref_frame, frame_is_null).expect("pool");
let p = params(AV_PIX_FMT_GRAY8, 8, 2);
let size = tight_size(p.pix_fmt, 8, 2);
let (ingress_tx, ingress_rx) = crossbeam_channel::bounded::<Vec<u8>>(4);
let (fg_tx, fg_rx) = crossbeam_channel::bounded::<FrameBox>(1);
let status = Arc::new(AtomicUsize::new(STATUS_RUN));
let thread_sync = ThreadSynchronizer::new();
let result = Arc::new(Mutex::new(None));
frame_source_init(
0,
FrameSource {
ingress: ingress_rx,
fg_sender: fg_tx,
params: p,
},
pool,
status.clone(),
thread_sync.clone(),
result.clone(),
)
.expect("spawn");
ingress_tx.send(vec![1u8; size]).unwrap();
ingress_tx.send(vec![2u8; size]).unwrap();
let deadline = Instant::now() + Duration::from_secs(20);
while !(fg_rx.is_full() && ingress_tx.is_empty()) {
assert!(
Instant::now() < deadline,
"worker never reached the parked send"
);
std::thread::sleep(Duration::from_millis(1));
}
std::thread::sleep(Duration::from_millis(300));
status.store(STATUS_END, Ordering::Release);
let (tx, rx) = std::sync::mpsc::channel();
let sync2 = thread_sync.clone();
std::thread::spawn(move || {
sync2.wait_for_all_threads();
let _ = tx.send(());
});
rx.recv_timeout(Duration::from_secs(30))
.expect("worker parked in a full-channel send must exit on terminal status");
assert!(
result.lock().unwrap().is_none(),
"a status-driven exit is not an error"
);
drop(fg_rx);
drop(ingress_tx);
}
#[test]
fn recycled_shell_refills_cleanly() {
let pool = test_pool();
let p = params(AV_PIX_FMT_GRAY8, 8, 2);
let data_a = vec![0xAA; tight_size(p.pix_fmt, 8, 2)];
let frame = build_video_frame(&pool, &p, &data_a, 0).expect("first build");
pool.release(frame); let data_b = vec![0x55; tight_size(p.pix_fmt, 8, 2)];
let frame = build_video_frame(&pool, &p, &data_b, 1).expect("recycled build");
unsafe {
let f = frame.as_ptr();
let row = std::slice::from_raw_parts((*f).data[0], 8);
assert_eq!(row, &data_b[..8]);
assert_eq!((*f).pts, 1);
}
}
}