refactor(avhw): split encoder module
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -0,0 +1,302 @@
|
||||
use std::mem;
|
||||
use std::ptr;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
use anyhow::{bail, Result};
|
||||
use ffmpeg_next as ff;
|
||||
use ffmpeg_next::ffi;
|
||||
use ffmpeg_next::packet::Mut as _;
|
||||
|
||||
use super::encode_output::{self, FrameOutput, PacketOutput};
|
||||
use super::hash::hash_sampled_y_plane;
|
||||
use super::{
|
||||
ff_err, BitrateCommand, CpuNv12Frame, EncodeOutcome, ResolutionChange, SwEncodeTiming,
|
||||
};
|
||||
|
||||
pub struct SwEncEncode {
|
||||
pub(super) sws_ctx: *mut ffi::SwsContext,
|
||||
pub(super) enc_video: ff::codec::encoder::video::Video,
|
||||
pub(super) output: Option<FrameOutput>,
|
||||
pub(super) yuv_frame: *mut ffi::AVFrame,
|
||||
pub(super) last_frame_hash: u64,
|
||||
pub(super) frame_count: u64,
|
||||
pub(super) starting_timestamp: Option<i64>,
|
||||
pub(super) frames_written: bool,
|
||||
pub(super) webrtc_disconnected: bool,
|
||||
pub(super) webrtc_paused: Option<Arc<AtomicBool>>,
|
||||
pub(super) bitrate_rx: crossbeam_channel::Receiver<BitrateCommand>,
|
||||
pub(super) resolution_rx: crossbeam_channel::Receiver<ResolutionChange>,
|
||||
pub(super) enc_width: u32,
|
||||
pub(super) enc_height: u32,
|
||||
pub(super) fps: u32,
|
||||
pub(super) bitrate: u64,
|
||||
pub(super) gop_size: u32,
|
||||
/// Set true when WebRTC requests a keyframe. Forces the next frame to
|
||||
/// `AV_PICTURE_TYPE_I` and bypasses the dedup hash check. Cleared only
|
||||
/// after `avcodec_send_frame` accepts the forced frame.
|
||||
pub(super) force_keyframe_pending: bool,
|
||||
/// Last per-frame timing snapshot. Reset to `Default` at the start of
|
||||
/// every `encode_cpu_frame` call (even on early returns) so stale values
|
||||
/// from a previous frame can never leak out.
|
||||
pub(super) last_timing: SwEncodeTiming,
|
||||
/// Capture time of the frame currently being encoded. Saved from the
|
||||
/// input `CpuNv12Frame` so `drain_encoder` can propagate it into the
|
||||
/// emitted `EncodedH264Frame` for the frame_age stat (issue #20).
|
||||
pub(super) last_capture_time: Option<Instant>,
|
||||
}
|
||||
|
||||
/// WebRTC media clock frequency in Hz. Matches RTP clock for video (RFC 3551).
|
||||
/// Used as encoder time_base denominator for WebRTC mode (1/90000) so that
|
||||
/// PTS values directly become RTP timestamps with microsecond precision.
|
||||
/// MP4 mode keeps 1/fps time_base for file output simplicity.
|
||||
pub const WEBRTC_RTP_CLOCK_HZ: i128 = 90_000;
|
||||
|
||||
// SAFETY: SwEncEncode owns sws_ctx/yuv_frame/enc_video exclusively after construction.
|
||||
// It is moved to a single encode thread and only accessed through &mut self there.
|
||||
unsafe impl Send for SwEncEncode {}
|
||||
|
||||
impl SwEncEncode {
|
||||
pub fn flush(&mut self) -> Result<()> {
|
||||
// SAFETY: Sending a null frame flushes the opened software encoder;
|
||||
// no frame data is dereferenced. enc_video is exclusively borrowed via &mut self.
|
||||
unsafe {
|
||||
let ret = ffi::avcodec_send_frame(self.enc_video.as_mut_ptr(), ptr::null());
|
||||
if ret < 0 && ret != ffi::AVERROR_EOF {
|
||||
bail!("software encoder flush send failed: {}", ff_err(ret));
|
||||
}
|
||||
}
|
||||
let start_ts = self.starting_timestamp.unwrap_or(0);
|
||||
let _ = self.drain_encoder(start_ts)?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn take_timing(&mut self) -> SwEncodeTiming {
|
||||
mem::take(&mut self.last_timing)
|
||||
}
|
||||
|
||||
pub fn encode_cpu_frame(&mut self, frame: &CpuNv12Frame) -> Result<EncodeOutcome> {
|
||||
self.last_timing = SwEncodeTiming::default();
|
||||
// Save capture_time so drain_encoder can propagate it into the
|
||||
// EncodedH264Frame emitted via the WebRTC channel (issue #20).
|
||||
self.last_capture_time = Some(frame.capture_time);
|
||||
|
||||
if self.webrtc_disconnected {
|
||||
return Ok(EncodeOutcome::SkippedDisconnected);
|
||||
}
|
||||
|
||||
// Must drain before the stride check: the import thread emits
|
||||
// ResolutionChange before the new (smaller-stride) frame arrives.
|
||||
while let Ok(cmd) = self.bitrate_rx.try_recv() {
|
||||
match cmd {
|
||||
BitrateCommand::UpdateBitrate { target_bps } => {
|
||||
// #23 defensive guardrail: clamp to reasonable max even if policy layer
|
||||
// is bypassed. 50 Mbps is a hard ceiling; primary cap is enforced in
|
||||
// state_portal.rs webrtc_thread_loop via --max-bitrate flag.
|
||||
const ENCODER_BITRATE_HARD_CAP: u64 = 50_000_000;
|
||||
let target_bps = target_bps.min(ENCODER_BITRATE_HARD_CAP);
|
||||
tracing::info!(target_bps, "updating encoder bitrate from BWE feedback");
|
||||
self.bitrate = target_bps;
|
||||
// SAFETY: enc_video is an opened AVCodecContext exclusively owned by &mut self.
|
||||
unsafe {
|
||||
let ctx = self.enc_video.as_mut_ptr();
|
||||
(*ctx).bit_rate = target_bps as i64;
|
||||
}
|
||||
}
|
||||
BitrateCommand::UpdateResolution { .. } => {}
|
||||
BitrateCommand::ForceKeyframe => {
|
||||
self.force_keyframe_pending = true;
|
||||
tracing::debug!("encode thread: ForceKeyframe requested");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let force_this_frame = self.force_keyframe_pending;
|
||||
|
||||
while let Ok(change) = self.resolution_rx.try_recv() {
|
||||
self.recreate_encoder(change.width, change.height)?;
|
||||
}
|
||||
|
||||
if frame.y_stride < self.enc_width as usize || frame.uv_stride < self.enc_width as usize {
|
||||
bail!("CPU NV12 frame stride is smaller than encoder width");
|
||||
}
|
||||
if let Some(ref paused) = self.webrtc_paused {
|
||||
if paused.load(Ordering::Relaxed) {
|
||||
return Ok(EncodeOutcome::SkippedPaused);
|
||||
}
|
||||
}
|
||||
|
||||
let width = self.enc_width as usize;
|
||||
let height = self.enc_height as usize;
|
||||
let required_y_len = frame.y_stride * height.saturating_sub(1) + width;
|
||||
if frame.y_data.len() < required_y_len {
|
||||
bail!("CPU NV12 frame Y plane is smaller than encoder dimensions");
|
||||
}
|
||||
|
||||
let frame_index = self.frame_count;
|
||||
self.frame_count = self.frame_count.saturating_add(1);
|
||||
let current_hash = hash_sampled_y_plane(&frame.y_data, width, height, frame.y_stride);
|
||||
let force_gop_frame =
|
||||
self.gop_size > 0 && frame_index.is_multiple_of(u64::from(self.gop_size));
|
||||
if frame_index > 0
|
||||
&& !force_gop_frame
|
||||
&& !force_this_frame
|
||||
&& current_hash == self.last_frame_hash
|
||||
{
|
||||
tracing::debug!(frame_index, "skipping duplicate frame");
|
||||
self.last_frame_hash = current_hash;
|
||||
return Ok(EncodeOutcome::SkippedDuplicate);
|
||||
}
|
||||
self.last_frame_hash = current_hash;
|
||||
|
||||
let sws_start = Instant::now();
|
||||
// SAFETY: yuv_frame is an owned reusable YUV420P frame at the same dimensions as sw_nv12;
|
||||
// sws_ctx was created for NV12 -> YUV420P with no resize, so sws_scale only converts format.
|
||||
unsafe {
|
||||
let ret = ffi::av_frame_make_writable(self.yuv_frame);
|
||||
if ret < 0 {
|
||||
bail!("av_frame_make_writable failed: {}", ff_err(ret));
|
||||
}
|
||||
let src_slices = [
|
||||
frame.y_data.as_ptr(),
|
||||
frame.uv_data.as_ptr(),
|
||||
ptr::null(),
|
||||
ptr::null(),
|
||||
];
|
||||
let src_strides = [frame.y_stride as i32, frame.uv_stride as i32, 0, 0];
|
||||
let scaled = ffi::sws_scale(
|
||||
self.sws_ctx,
|
||||
src_slices.as_ptr(),
|
||||
src_strides.as_ptr(),
|
||||
0,
|
||||
self.enc_height as i32,
|
||||
(*self.yuv_frame).data.as_ptr() as *mut *mut u8,
|
||||
(*self.yuv_frame).linesize.as_ptr() as *const i32,
|
||||
);
|
||||
if scaled < 0 {
|
||||
bail!("sws_scale failed for software encoder: {scaled}");
|
||||
}
|
||||
}
|
||||
let sws_us = sws_start.elapsed().as_micros() as u64;
|
||||
|
||||
let pts = frame.pts;
|
||||
if self.starting_timestamp.is_none() {
|
||||
self.starting_timestamp = Some(pts);
|
||||
}
|
||||
let start_ts = self.starting_timestamp.unwrap_or(0);
|
||||
|
||||
let enc_start = Instant::now();
|
||||
// SAFETY: yuv_frame is initialized, writable, and matches the opened encoder format.
|
||||
// pict_type is reset every frame: the AVFrame is reused, so without resetting to NONE
|
||||
// a previously-forced I-type would leak into subsequent P-frames. With forced-idr=1
|
||||
// set on the encoder, AV_PICTURE_TYPE_I produces a true IDR NALU.
|
||||
unsafe {
|
||||
(*self.yuv_frame).pts = pts;
|
||||
(*self.yuv_frame).pict_type = if force_this_frame {
|
||||
ffi::AVPictureType::AV_PICTURE_TYPE_I
|
||||
} else {
|
||||
ffi::AVPictureType::AV_PICTURE_TYPE_NONE
|
||||
};
|
||||
let ret = ffi::avcodec_send_frame(self.enc_video.as_mut_ptr(), self.yuv_frame);
|
||||
if ret < 0 {
|
||||
bail!(
|
||||
"avcodec_send_frame failed for software encoder: {}",
|
||||
ff_err(ret)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if force_this_frame {
|
||||
self.force_keyframe_pending = false;
|
||||
}
|
||||
|
||||
let output_bytes = self.drain_encoder(start_ts)?;
|
||||
let encode_us = enc_start.elapsed().as_micros() as u64;
|
||||
|
||||
self.last_timing = SwEncodeTiming {
|
||||
sws_us,
|
||||
encode_us,
|
||||
output_bytes,
|
||||
};
|
||||
|
||||
Ok(EncodeOutcome::Encoded)
|
||||
}
|
||||
|
||||
pub(super) fn write_trailer_if_needed(&mut self) -> Result<()> {
|
||||
if self.frames_written {
|
||||
if let Some(FrameOutput::Muxer(ref mut octx)) = self.output {
|
||||
octx.write_trailer()
|
||||
.map_err(|e| anyhow::anyhow!("Failed to write trailer: {e}"))?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn drain_encoder(&mut self, start_ts: i64) -> Result<usize> {
|
||||
let mut total_bytes = 0usize;
|
||||
loop {
|
||||
let mut pkt = ff::Packet::empty();
|
||||
// SAFETY: enc_video is an open encoder; pkt is writable packet storage.
|
||||
let ret = unsafe {
|
||||
ffi::avcodec_receive_packet(self.enc_video.as_mut_ptr(), pkt.as_mut_ptr())
|
||||
};
|
||||
if ret < 0 {
|
||||
if ret == ffi::AVERROR(ffi::EAGAIN) || ret == ffi::AVERROR_EOF {
|
||||
break;
|
||||
}
|
||||
bail!("avcodec_receive_packet failed: {}", ff_err(ret));
|
||||
}
|
||||
|
||||
// Count encoded bytes produced before the Muxer/Channel match to
|
||||
// avoid branch duplication and handle multi-packet drain correctly.
|
||||
// SAFETY: pkt was just filled by a successful avcodec_receive_packet;
|
||||
// the size field is valid and initialized.
|
||||
let pkt_size = unsafe { (*pkt.as_mut_ptr()).size };
|
||||
if pkt_size > 0 {
|
||||
total_bytes += pkt_size as usize;
|
||||
}
|
||||
|
||||
match self.output {
|
||||
Some(FrameOutput::Muxer(ref mut octx)) => {
|
||||
encode_output::write_muxer_packet(
|
||||
&mut pkt,
|
||||
octx,
|
||||
self.enc_video.time_base(),
|
||||
start_ts,
|
||||
)?;
|
||||
self.frames_written = true;
|
||||
}
|
||||
Some(FrameOutput::Channel(ref tx))
|
||||
if encode_output::send_channel_packet(
|
||||
&mut pkt,
|
||||
tx,
|
||||
start_ts,
|
||||
self.last_capture_time.unwrap_or_else(Instant::now),
|
||||
)? == PacketOutput::Disconnected =>
|
||||
{
|
||||
self.webrtc_disconnected = true;
|
||||
break;
|
||||
}
|
||||
Some(FrameOutput::Channel(_)) => {}
|
||||
None => {}
|
||||
}
|
||||
}
|
||||
Ok(total_bytes)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for SwEncEncode {
|
||||
fn drop(&mut self) {
|
||||
if !self.sws_ctx.is_null() {
|
||||
// SAFETY: sws_ctx is owned by this state and was returned by sws_getContext.
|
||||
unsafe { ffi::sws_freeContext(self.sws_ctx) };
|
||||
self.sws_ctx = ptr::null_mut();
|
||||
}
|
||||
if !self.yuv_frame.is_null() {
|
||||
// SAFETY: yuv_frame is owned by this state and was allocated by av_frame_alloc.
|
||||
unsafe { ffi::av_frame_free(&mut self.yuv_frame) };
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user