Files
wl-webrtc/src/avhw.rs
T

2617 lines
121 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! # avhw — FFmpeg / VAAPI 硬件加速 FFI 绑定(FFI 最密集的文件)
//!
//! 本文件是整个 crate 中 `unsafe` 块密度最高的模块:直接调用 FFmpeg C API。
//! `ffmpeg-next` 是 ffmpeg-sys-nextC FFI 绑定)的薄包装;硬件加速相关 API
//! `AVBufferRef`/`av_hwframe_*`)仍是裸 C 接口。
//!
//! ## FFmpeg 引用计数对象的生命周期
//! - `AVBufferRef`FFmpeg 引用计数句柄,`av_buffer_ref` 加 1、`av_buffer_unref` 减 1。
//! `AvHwDevCtx`/`AvHwFrameCtx` 用 RAII 包装:`new` alloc`Drop` unref。
//! - `AVFrame`:硬件帧(`AV_PIX_FMT_VAAPI`)的 `data[0]` 是 GPU 表面指针。
//!
//! ## `unsafe impl Send` 的存在原因(AGENTS.md 明确警告)
//! `*mut AVBufferRef` 默认 `!Send`wrapper 显式声明 Send 的前提是**外部调用方
//! 保证 `&mut self` 独占访问**。AGENTS.md:不要跨线程移动这些 wrapper 而不重新
//! 检查 exclusivity 假设。
//!
//! ## 与 Go cgo 的类比
//! - `extern "C" fn` ≈ Go `//export` C 回调;
//! - `unsafe { ffi::av_*() }` ≈ Go `C.av_xxx()`
//! - RAII `Drop` ≈ Go `runtime.SetFinalizer`(但 Rust Drop 是确定性的)。
// std 导入:CStringC 字符串)、mem::zeroedFFI 零初始化)、RawFd/AsRawFdfd 桥),
// c_voidC void 跨语言)、PathDRM 设备路径)、ptrnull_mut 等裸指针工具),
// slice(从裸指针构造切片)、AtomicBool/Ordering/Arc(跨线程暂停标志,T10b 用),
// Instant(编码计时,T10b 用)。本 sub-todolines 1-330)仅部分使用。
use std::ffi::CString;
use std::mem;
use std::os::fd::{AsRawFd, RawFd};
use std::os::raw::c_void;
use std::path::Path;
use std::ptr;
use std::slice;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Instant;
// anyhow 错误处理:bail!(提前返回 Err)、Result<T>(错误传播)。
// ffmpeg-next:稍安全的 APIff::frame::Video、ff::format::Pixel、ff::codec 等)。
// ffmpeg_next::ffiFFmpeg C 头绑定的裸 APIAVBufferRef、av_hwdevice_ctx_create 等),
// 所有 `unsafe` FFI 调用的入口。`packet::Mut as _`:导入 packet 的 Mut trait 但匿名,
// 仅用于 trait method 解析,不污染命名空间。
use anyhow::{bail, Result};
use ffmpeg_next as ff;
use ffmpeg_next::ffi;
use ffmpeg_next::packet::Mut as _;
// crate 内部模块:PwDmaBufFramePipeWire DMA-BUF 帧元数据:fd/width/height/stride/modifier),
// TransformWayland wl_output 变换,决定 ROI 是否需要转置处理)。
use crate::cap_portal::PwDmaBufFrame;
use crate::transform::{transpose_if_transform_transposed, Transform};
// ---------------------------------------------------------------------------
// Bitrate feedback command (WebRTC BWE → SW encoder)
// ---------------------------------------------------------------------------
// 跨线程控制信令:WebRTC 线程根据 BWEbandwidth estimation)估算的可用带宽,
// 通过 MPSC channel 把 `BitrateCommand` 发给软件编码线程,由后者调整 x264 参数。
// 注意:硬件(VAAPI)路径目前不接受动态码率调整,只有软件(x264)路径消费这些命令。
// 类比 Go`chan BitrateCommand` 单向 channel。
/// Commands sent from the WebRTC thread to the SW encoder when the
/// bandwidth estimate changes significantly.
pub enum BitrateCommand {
UpdateBitrate { target_bps: u64 },
UpdateResolution { width: u32, height: u32 },
/// Force the next encoded frame to be an IDR. Sent by the WebRTC thread
/// in response to str0m `Event::KeyframeRequest` or a resolution change.
ForceKeyframe,
}
// 分辨率变更事件:WebRTC 线程根据对端能力(SVC/layered encoding 协商)建议切换分辨率,
// 编码线程收到后重建 encoder(x264 不支持运行时改变 width/height,必须销毁重建)。
// `#[derive(Clone, Copy, Debug)]`Copy 表示按位拷贝(小值类型,无需 clone 调用),
// 类比 Go 的 `type ResolutionChange struct{...}`Go 默认值语义即 Copy)。
#[derive(Clone, Copy, Debug)]
pub struct ResolutionChange {
pub width: u32,
pub height: u32,
}
// 每帧编码耗时快照,由编码线程填好后通过 channel 发给统计线程(stats thread)。
// `#[derive(Default, Clone, Copy, Debug)]`Default 允许 `SwEncodeTiming::default()` 全零初始化;
// Copy/Clone 表示小值类型按位拷贝(无堆字段)。
//
// 字段语义:
// - sws_uslibswscale 把 NV12VAAPI 输出)转 YUV420Px264 输入)的耗时(微秒)
// - encode_us`avcodec_send_frame` + drain 接收所有 packet 的总耗时(微秒)
// - output_bytes:本帧产出的 H.264 字节数(即使下游 WebRTC 因暂停丢帧也计入)
//
// 类比 Go`type SwEncodeTiming struct{...}` + atomic store/load 跨 goroutine 传递。
/// Per-frame timing snapshot for the software encoder, consumed by the stats
/// thread. `sws_us` measures NV12→YUV420P conversion, `encode_us` measures
/// `avcodec_send_frame` + drain, and `output_bytes` counts encoded bytes
/// produced by libavcodec (even if downstream delivery later drops them).
#[derive(Default, Clone, Copy, Debug)]
pub struct SwEncodeTiming {
pub sws_us: u64,
pub encode_us: u64,
pub output_bytes: usize,
}
// `encode_cpu_frame` 的返回值,编码线程据此判断是否计 encoded_fps(只有 `Encoded`
// 才算一次真实编码)。`#[derive(Debug, Clone, Copy, PartialEq, Eq)]`PartialEq/Eq
// 允许 `==` 比较(测试和 stats 聚合用),Copy 表示小值类型按位拷贝。
// 类比 Go`type EncodeOutcome int` + const 枚举值。
/// Outcome of a single `encode_cpu_frame` call. Used by the encode thread
/// to decide whether to report timing stats (only real encodes tick encoded_fps).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EncodeOutcome {
/// Frame was actually encoded and produced output bytes.
Encoded,
/// Frame was dropped because WebRTC is paused (no client connected).
SkippedPaused,
/// Frame was dropped because the encoder is in disconnected state.
SkippedDisconnected,
/// Frame was dropped because its Y-plane hash matched the previous frame.
SkippedDuplicate,
}
// ---------------------------------------------------------------------------
// AvHwDevCtx
// ---------------------------------------------------------------------------
// VAAPI 设备上下文的 RAII 包装。`ptr` 是 FFmpeg 引用计数的 `AVBufferRef`
// 包装 `AVHWDeviceContext`VAAPI 内部的 `VADisplay`)。`new_vaapi` 时 alloc
// `Drop` 时 `av_buffer_unref``ref_clone` 用于把内部 ref 传递给编码器的
// `AVCodecContext::hw_device_ctx`FFmpeg 会再 `av_buffer_ref` 一次)。
// **独占访问假设**`unsafe impl Send` 的前提是外部调用方保证 `&mut self`
// 不允许多线程共享 `&AvHwDevCtx` 同时调用 mut method(详见 AGENTS.md 警告)。
pub struct AvHwDevCtx {
// FFmpeg 内部引用计数句柄,类型擦除为 `*mut AVBufferRef`。
// 永不为 null(除非在 Drop 中被置 null)。
ptr: *mut ffi::AVBufferRef,
}
// 中文概述:声明 `AvHwDevCtx` 可以跨线程 `Send`。前提是调用方用 `&mut self`
// 独占访问(无并发可变借用)。VAAPI 的 `VADisplay` 在硬件驱动层线程安全,
// FFmpeg 的 `AVHWDeviceContext` 内部状态在 send/receive 编码模式下也线程安全。
// 不修复:保留 AGENTS.md 文档的 exclusivity 假设(这是本 crate 的并发模型基础)。
// SAFETY: AvHwDevCtx wraps an FFmpeg AVBufferRef which is not Send by default,
// but we guarantee exclusive access through &mut self. The underlying VAAPI
// device context is thread-safe for the operations we perform.
unsafe impl Send for AvHwDevCtx {}
impl AvHwDevCtx {
// 创建 VAAPI 设备上下文。`drm_device` 通常是 `/dev/dri/renderD128`。
// 失败原因:设备不存在、无 VAAPI 驱动、权限不足(用户不在 `video` 组)。
// 成功时返回 `AvHwDevCtx`ptr 持有一个 refrefcount=1)。
pub fn new_vaapi(drm_device: &Path) -> Result<Self> {
// CString::new 在路径含内部 NUL 时返回 ErrFFmpeg C API 要求 NUL 结尾)。
// `to_str().unwrap()`:路径非 UTF-8 时 panicLinux DRM 设备路径通常 ASCII)。
let device_cstr = CString::new(drm_device.to_str().unwrap())?;
// 初始化为 null_mutFFmpeg 的 out-pointer 约定(调用者置 null,被调者赋值)。
let mut p: *mut ffi::AVBufferRef = ptr::null_mut();
// 中文概述:调用 FFmpeg C API 创建 VAAPI `AVHWDeviceContext`,写入 `*p`。
// 失败返回负的 AVERROR;成功返回 0 且 `p` 指向新分配的 `AVBufferRef`refcount=1)。
// SAFETY: device_cstr is a valid C string for the duration of the call;
// p is a valid out-pointer that FFmpeg initializes on success.
let ret = unsafe {
ffi::av_hwdevice_ctx_create(
&mut p,
ffi::AVHWDeviceType::AV_HWDEVICE_TYPE_VAAPI,
device_cstr.as_ptr(),
ptr::null_mut(),
0,
)
};
if ret < 0 {
bail!(
"Failed to create VAAPI device context from {}: {}",
drm_device.display(),
ff_err(ret)
);
}
Ok(Self { ptr: p })
}
// 返回裸 `*mut AVBufferRef`,调用方负责管理引用计数。
// 通常用于传给 FFmpeg API 的 `hw_device_ctx` 字段(FFmpeg 内部会 av_buffer_ref)。
pub fn as_ptr(&self) -> *mut ffi::AVBufferRef {
self.ptr
}
// 增加 refcount,返回新的 `*mut AVBufferRef`(不复制底层缓冲)。
// 用于把设备 ref 传给编码器的 `AVCodecContext::hw_device_ctx`,编码器 Drop 时
// 调用 `av_buffer_unref` 释放自己持有的 ref,不影响本 wrapper。
pub fn ref_clone(&self) -> *mut ffi::AVBufferRef {
// 中文概述:FFmpeg 引用计数 +1,返回新 ref;线程安全(内部原子操作)。
// SAFETY: av_buffer_ref atomically increments refcount and returns a new ref.
unsafe { ffi::av_buffer_ref(self.ptr) }
}
}
// RAII 析构:减少 refcount。如果是最后一个 refrefcount→0),FFmpeg 释放
// 底层的 `AVHWDeviceContext` 和 `VADisplay`。`is_null()` 守卫是为了应对
// `Drop` 被多次调用或 `new_vaapi` 失败后 ptr 仍为 null 的边缘情况(实际上
// Rust 的 Drop 不会被调用两次,但 FFmpeg 的 av_buffer_unref 接受 null 入参)。
impl Drop for AvHwDevCtx {
fn drop(&mut self) {
if !self.ptr.is_null() {
// 中文概述:FFmpeg 引用计数 -1,refcount=0 时释放底层设备上下文。
// SAFETY: av_buffer_unref decrements refcount; frees the buffer when it hits zero.
unsafe { ffi::av_buffer_unref(&mut self.ptr) };
}
}
}
// ---------------------------------------------------------------------------
// AvHwFrameCtx
// ---------------------------------------------------------------------------
// VAAPI 硬件帧池的 RAII 包装。`AVHWFramesContext` 是 FFmpeg 的硬件帧分配器,
// 编码器通过它分配 `AV_PIX_FMT_VAAPI` 帧(GPU 表面)。`initial_pool_size=4`
// 表示池大小为 4(足够 send/receive pipeline1 帧 encode 中 + 1 帧 queue + 2 缓冲)。
// 字段配置(width/height/sw_format)必须在 `av_hwframe_ctx_init` 前设置;
// `sw_fmt` 是 CPU 可读的软件格式(VAAPI 内部从该格式上传/下载)。
pub struct AvHwFrameCtx {
ptr: *mut ffi::AVBufferRef,
}
// 中文概述:声明 `AvHwFrameCtx` 可以跨线程 `Send`。前提是 `&mut self` 独占访问,
// 且底层硬件帧池在 send/receive 模式下线程安全。与 `AvHwDevCtx` 同样的 exclusivity 假设。
// 不修复:保留 AGENTS.md 文档的并发约束。
// SAFETY: AvHwFrameCtx wraps an FFmpeg AVBufferRef to an AVHWFramesContext.
// It is only accessed through &mut self, ensuring no concurrent mutation.
// The underlying hardware frames pool is thread-safe for the send/receive pattern.
unsafe impl Send for AvHwFrameCtx {}
impl AvHwFrameCtx {
// 内部构造函数:分配 `AVHWFramesContext`、配置字段、调用 `av_hwframe_ctx_init`。
// 失败时(alloc 或 init 失败)返回 Err,确保 ptr 不泄漏。
fn new_inner(hw_dev: &AvHwDevCtx, w: u32, h: u32, sw_fmt: ff::format::Pixel) -> Result<Self> {
// 中文概述:分配 `AVHWFramesContext`refcount=1,关联到 hw_dev 的设备上下文。
// 返回的 ref 的 `data` 字段指向未初始化的 `AVHWFramesContext`,需手动配置。
// SAFETY: hw_dev is a live AVHWDeviceContext; FFmpeg returns either a valid
// frames context ref or null (checked below).
let mut p = unsafe { ffi::av_hwframe_ctx_alloc(hw_dev.as_ptr()) };
if p.is_null() {
bail!("av_hwframe_ctx_alloc returned null");
}
// 中文概述:通过 `(*p).data as *mut AVHWFramesContext` 访问 FFmpeg 内部类型擦除的
// 字段(`AVBufferRef::data` 是 `*mut u8`,实际指向具体 context 类型)。
// 直接赋值 `format`/`sw_format`/`width`/`height`/`initial_pool_size` 5 个字段。
// SAFETY: p is a valid AVBufferRef from av_hwframe_ctx_alloc.
// Its .data field points to an AVHWFramesContext that we must configure.
unsafe {
let fc = (*p).data as *mut ffi::AVHWFramesContext;
(*fc).format = ff::format::Pixel::VAAPI.into();
(*fc).sw_format = sw_fmt.into();
(*fc).width = w as i32;
(*fc).height = h as i32;
(*fc).initial_pool_size = 4;
}
// 中文概述:初始化帧池(实际分配 GPU 表面)。失败时必须手动 unref p 否则泄漏。
// SAFETY: p is a valid AVHWFramesContext ref configured above and not yet
// transferred or freed.
let ret = unsafe { ffi::av_hwframe_ctx_init(p) };
if ret < 0 {
// 中文概述:init 失败的清理路径——p 仍有效但不可用,unref 释放。
// SAFETY: p is valid but init failed; clean up.
unsafe { ffi::av_buffer_unref(&mut p) };
bail!("av_hwframe_ctx_init failed: {}", ff_err(ret));
}
Ok(Self { ptr: p })
}
// 为采集路径创建硬件帧池。`sw_fmt` 通常是 `BGRA`PipeWire DMA-BUF 像素格式)。
pub fn for_capture(
hw_dev: &AvHwDevCtx,
w: u32,
h: u32,
sw_fmt: ff::format::Pixel,
) -> Result<Self> {
Self::new_inner(hw_dev, w, h, sw_fmt)
}
pub fn as_ptr(&self) -> *mut ffi::AVBufferRef {
self.ptr
}
pub fn ref_clone(&self) -> *mut ffi::AVBufferRef {
// 中文概述:FFmpeg 引用计数 +1,返回新 ref;用于把帧池 ref 传给编码器。
// SAFETY: av_buffer_ref atomically increments refcount and returns a new ref.
unsafe { ffi::av_buffer_ref(self.ptr) }
}
}
// RAII 析构:与 AvHwDevCtx::Drop 同样的模式(unref 释放 ref)。
impl Drop for AvHwFrameCtx {
fn drop(&mut self) {
if !self.ptr.is_null() {
// 中文概述:FFmpeg 引用计数 -1refcount=0 时释放帧池和 GPU 表面。
// SAFETY: av_buffer_unref decrements refcount; frees when zero.
unsafe { ffi::av_buffer_unref(&mut self.ptr) };
}
}
}
// 启动时探针:测试 `drm_device` 能否通过 VAAPI 导入 PipeWire 给的 DMA-BUF 帧。
// 用于 backend_detect 决定走硬件(VAAPI)还是软件(x264)编码路径。
// 失败原因:硬件不支持给定格式/修饰符组合、DRM 设备无 VAAPI、内核 driver 限制。
// 成功只是"探针通过",不代表实际编码时一直可用(运行时仍可能因 OOM 等失败)。
/// Test whether `drm_device` can import the PipeWire DMA-BUF frame via VAAPI.
pub fn test_dma_buf_import(drm_device: &Path, frame: &PwDmaBufFrame) -> Result<()> {
let hw_dev = AvHwDevCtx::new_vaapi(drm_device)?;
let frames =
AvHwFrameCtx::for_capture(&hw_dev, frame.width, frame.height, ff::format::Pixel::BGRA)?;
// 中文概述:调用下面的 `import_dma_buf_to_vaapi`,用 `av_hwframe_map` 把 DMA-BUF
// 映射到 VAAPI 表面(零拷贝)。`as_raw_fd()` 把 `OwnedFd` 转为裸 fd 传给 FFI。
// SAFETY: frames is a live VAAPI frames context; frame carries valid DMA-BUF metadata.
unsafe {
import_dma_buf_to_vaapi(
frames.as_ptr(),
frame.fd.as_raw_fd(),
frame.width,
frame.height,
frame.format,
frame.modifier,
frame.stride,
frame.offset,
)
}?;
Ok(())
}
// 中文概述:把 PipeWire 的 DMA-BUF 帧零拷贝映射为 VAAPI `AVFrame`。
// DMA-BUF 是 Linux 内核的 GPU 缓冲共享机制(fd 指向 GPU 内存);`av_hwframe_map`
// 在 VAAPI 内部通过 `vaCreateSurfaces` + `vaMapBuffer` 直接 GPU 映射,无需 CPU 拷贝。
// 类比 GoGo 没有等价物——cgo + libavutil 才能实现同样的零拷贝路径。
/// Import a DMA-BUF into a VAAPI hardware frame via zero-copy `av_hwframe_map`.
///
/// # Safety
/// - `frames_ctx` must point to an initialized AVHWCramesContext for VAAPI
/// - `raw_fd` must be a valid DMA-BUF file descriptor
pub unsafe fn import_dma_buf_to_vaapi(
frames_ctx: *mut ffi::AVBufferRef,
raw_fd: RawFd,
width: u32,
height: u32,
drm_format: u32,
modifier: u64,
stride: u32,
offset: u64,
) -> Result<ff::frame::Video> {
// `libc::dup` 复制 fd——`AVDRMFrameDescriptor` 持有 dup 后的 fd
// 当 descriptor 释放时 `cleanup_drm_descriptor` 会 `libc::close` 它;
// 原始 `raw_fd` 不被本函数消费(仍由调用方拥有,通常是 PipeWire)。
let duped_fd = libc::dup(raw_fd);
if duped_fd < 0 {
bail!("dup(fd) failed: {}", std::io::Error::last_os_error());
}
// `mem::zeroed()`FFI 零初始化(所有字段置 0),类似 Go 的 `&drmDesc{}`。
// SAFETY 后续不再赘述——这是 `unsafe fn` 内部,调用契约由函数签名保证。
let mut desc: ffi::AVDRMFrameDescriptor = mem::zeroed();
desc.nb_objects = 1;
desc.objects[0].fd = duped_fd;
desc.objects[0].size = (height as usize) * (stride as usize);
desc.objects[0].format_modifier = modifier;
desc.nb_layers = 1;
desc.layers[0].format = drm_format;
desc.layers[0].nb_planes = 1;
desc.layers[0].planes[0].object_index = 0;
desc.layers[0].planes[0].offset = offset as isize;
desc.layers[0].planes[0].pitch = stride as isize;
// `Box::new(desc)` 把栈上的 descriptor 移到堆上;`Box::into_raw` 交出所有权给 FFmpeg。
// `av_buffer_create` 会接管这块内存:当 AVBufferRef refcount 归零时调用
// `cleanup_drm_descriptor`(自定义 free 回调)→ close fd + Box::from_raw 释放堆内存。
let desc_box = Box::new(desc);
let desc_ptr = Box::into_raw(desc_box);
// 创建 FFmpeg 自定义 AVBufferRefdata 指向 descriptorfree 回调为 cleanup_drm_descriptor。
// 失败时(极少)需要手动恢复 Box 并 close fd,否则泄漏。
let buf_ref = ffi::av_buffer_create(
desc_ptr as *mut u8,
std::mem::size_of::<ffi::AVDRMFrameDescriptor>(),
Some(cleanup_drm_descriptor),
ptr::null_mut(),
0,
);
if buf_ref.is_null() {
let desc_box = Box::from_raw(desc_ptr);
libc::close(desc_box.objects[0].fd);
bail!("av_buffer_create returned null for DRM descriptor");
}
// 构造 src AVFrame:格式为 `AV_PIX_FMT_DRM_PRIME`FFmpeg 的 DMA-BUF 包装格式),
// `data[0]` 指向 AVBufferRef 的 data(即 descriptor),`buf[0]` 持有 ref 防止提前释放。
let mut src = ff::frame::Video::empty();
{
let sp = src.as_mut_ptr();
(*sp).format = ffi::AVPixelFormat::AV_PIX_FMT_DRM_PRIME as i32;
(*sp).width = width as i32;
(*sp).height = height as i32;
(*sp).data[0] = (*buf_ref).data;
(*sp).buf[0] = buf_ref;
}
// 构造 dst AVFrame:格式为 `AV_PIX_FMT_VAAPI`,挂上 `hw_frames_ctx` 让 FFmpeg 知道用哪个池。
let mut dst = ff::frame::Video::empty();
// 中文概述:把 dst 的 format 设为 VAAPI、`hw_frames_ctx` 设为新的 ref。
// `av_buffer_ref(frames_ctx)` 必须成功,否则 av_hwframe_map 无法定位帧池。
// SAFETY: frames_ctx is guaranteed by this unsafe function's contract to be a
// valid initialized VAAPI frames context; we set format/hw_frames_ctx on a
// freshly allocated dst frame.
unsafe {
let dp = dst.as_mut_ptr();
(*dp).format = ffi::AVPixelFormat::AV_PIX_FMT_VAAPI as i32;
(*dp).hw_frames_ctx = ffi::av_buffer_ref(frames_ctx);
if (*dp).hw_frames_ctx.is_null() {
bail!("av_buffer_ref(frames_ctx) returned null");
}
}
// 中文概述:核心 FFI 调用——FFmpeg 在内部调用 VAAPI `vaCreateSurfaces`+`vaMapBuffer`
// 把 DMA-BUF 映射到 GPU 表面,写入 dst。零 CPU 拷贝。
// SAFETY: src and dst are initialized AVFrames; dst has a valid hw_frames_ctx
// ref and av_hwframe_map fills dst from src.
let ret = unsafe {
ffi::av_hwframe_map(
dst.as_mut_ptr(),
src.as_ptr(),
ffi::AV_HWFRAME_MAP_READ as i32,
)
};
if ret < 0 {
bail!("av_hwframe_map failed: {}", ff_err(ret));
}
Ok(dst)
}
// `av_buffer_create` 的 free 回调:当自定义 AVBufferRef 的 refcount 归零时被 FFmpeg 调用。
// `unsafe extern "C" fn` 是因为 FFmpeg 通过 C 函数指针调用它(C ABI 兼容)。
// 职责:(1) close DMA-BUF 的 dup'd fd(2) `Box::from_raw` 恢复堆内存所有权并立即 drop。
// `_opaque`FFmpeg 的 AVBufferRef 允许附带 opaque 数据(本例未用,故 `_` 前缀)。
// 类比 Go:相当于 `runtime.SetFinalizer` 的回调,但是 FFmpeg 主动调用的(确定性)。
unsafe extern "C" fn cleanup_drm_descriptor(_opaque: *mut c_void, data: *mut u8) {
let desc = data as *mut ffi::AVDRMFrameDescriptor;
if !desc.is_null() && (*desc).nb_objects > 0 && (*desc).objects[0].fd >= 0 {
libc::close((*desc).objects[0].fd);
}
let _ = Box::from_raw(data as *mut ffi::AVDRMFrameDescriptor);
}
// 把 FFmpeg 的负数错误码(AVERROR)转成人类可读字符串。
// `av_strerror` 是 FFmpeg 的线程安全 strerror 等价物,写入调用方提供的缓冲区。
// `pub(crate)` 可见性:仅供本 crate 内部错误格式化使用(不暴露给外部)。
/// Convert an FFmpeg error code to a human-readable string.
pub(crate) fn av_err_to_string(err: i32) -> String {
// 128 字节足够 AVERROR 描述(FFmpeg 内部字符串都很短)。
let mut buf = vec![0u8; 128];
// 中文概述:调用 FFmpeg `av_strerror`,把错误码描述写入 buf。
// 失败(无效错误码)时 buf 仍为全 0,trim 后为空字符串。
// SAFETY: buf points to 128 writable bytes and lives for the duration of
// av_strerror.
unsafe {
ffi::av_strerror(err, buf.as_mut_ptr() as *mut i8, buf.len());
}
// `from_utf8_lossy` 容忍非 UTF-8 字节;`trim_end_matches('\0')` 去掉 C 字符串结尾的 NUL。
String::from_utf8_lossy(&buf)
.trim_end_matches('\0')
.to_string()
}
// 带数字码的错误格式化:`"error -22 (Invalid argument)"`。
// 数字码很重要——AVERROR 是负数,但有些 FFmpeg 错误是宏(如 AVERROR(EAGAIN)),
// 看数字能区分底层 errno 还是 FFmpeg 自定义错误。
/// Format an FFmpeg error code with both numeric value and description.
/// Example output: "error -22 (Invalid argument)"
pub(crate) fn ff_err(ret: i32) -> String {
format!("error {ret} ({})", av_err_to_string(ret))
}
// ---------------------------------------------------------------------------
// EncState
// ---------------------------------------------------------------------------
// 编码状态机:把 VAAPI 硬件帧 → H.264 ES 流 → MP4 文件的"主循环编排者"。
//
// 字段角色(按数据流顺序):
// - enc_video : 已打开的 h264_vaapi 编码器句柄(FFmpeg AVCodecContext 包装)
// - video_filter : scale/crop/format 过滤图(BGRA→NV12 + VAAPI 像素格式上传)
// - frames_rgb : 捕获侧 hw frames 池(与 import_dma_buf_to_vaapi 共享)
// - hw_device_ctx : VAAPI 设备上下文(DRM render node 句柄)
// - octx : 输出 muxer 上下文(MP4 文件)
// - starting_timestamp : 首帧 PTS 锚点,用于把绝对时间戳归零(避免 MP4 起始时间漂移)
// - frames_written : 是否已写出过帧(用于决定 flush 时是否写 trailer
//
// 跨线程约束:state.rs/state_portal.rs 的设计是"单线程驱动"——`&mut EncState`
// 由外层 main 串行化访问,因此 `enc_video` / `octx` 等 FFmpeg 上下文不需要锁。
// `unsafe impl Send` 仅表达"类型上可以跨线程移动",不代表"可以并发访问"。
pub struct EncState {
enc_video: ff::codec::encoder::video::Video,
frames_rgb: AvHwFrameCtx,
video_filter: ff::filter::Graph,
hw_device_ctx: AvHwDevCtx,
octx: ff::format::context::Output,
starting_timestamp: Option<i64>,
frames_written: bool,
}
// 安全说明:FFmpeg 的 AVCodecContext/AVFormatContext 等 C 对象不是自动 Sync 的,
// 但本工程的并发模型保证 EncState 只在单一编码线程内被 `&mut` 访问,因此标记
// Send(跨线程移动)是安全的;若未来引入并发编码必须重新审视(参见 AGENTS.md
// 关于 avhw.rs 显式 `unsafe impl Send` 与独占性假设的说明)。
unsafe impl Send for EncState {}
// impl EncState:编码主循环对外暴露的方法集合。
//
// 方法职责(按数据流顺序):
// - new : 构造 VAAPI 设备/帧池/filter 图/编码器/muxer 6 步流水线(FFmpeg 严格顺序)
// - frames_rgb : 暴露内部 hw frames 池给捕获侧 import_dma_buf_to_vaapi 共享
// - encode_frame : 单帧驱动:送入 filter → 拉过滤后帧 → avcodec_send_frame → drain_encoder
// - flush : EOF 处理:filter drain + 空帧 avcodec_send_frame 触发 encoder flush + write_trailer
// - drain_encoder: 内部辅助:循环 avcodec_receive_packet 直到 EAGAIN/EOF,重缩放 PTS 并写 muxer
impl EncState {
// 抑制 clippy::too_many_arguments:构造函数需要 11 个参数(捕获尺寸/编码尺寸/码率/GOP/帧率/变换/共享 hw ctx),
// 拆分为 builder 模式反而会增加 FFI 调用顺序出错的风险(FFmpeg 各步骤有严格依赖关系)。
#[allow(clippy::too_many_arguments)]
// 构造编码器:类比 Go `func NewEncState(...) (*EncState, error)`11 个参数对应 6 步流水线配置。
// 返回 `Result<Self>` 用 `?` 把任何 FFI/IO 错误传播给调用方(state.rs/state_portal.rs)。
pub fn new(
drm_device: &Path,
output_path: &Path,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
bitrate: u64,
gop_size: u32,
fps: u32,
transform: Transform,
existing_hw_ctx: Option<AvHwDevCtx>,
) -> Result<Self> {
tracing::info!(
"EncState::new: {width}x{height} enc={enc_width}x{enc_height} transform={transform:?}"
);
// 1. VAAPI device — reuse existing context if provided
// `match` 是 Rust 的 exhaustive 模式匹配(类似 Go `switch` 但更强:编译器强制覆盖所有分支)。
// `Option<AvHwDevCtx>` 等价于 Go `*AvHwDevCtx`Some=valueNone=nil),但类型系统强制处理 nil 情况。
let hw_device_ctx = match existing_hw_ctx {
Some(ctx) => ctx,
None => AvHwDevCtx::new_vaapi(drm_device)?,
};
let frames_rgb =
AvHwFrameCtx::for_capture(&hw_device_ctx, width, height, ff::format::Pixel::BGRA)?;
// 3. Filter graph — must be built BEFORE encoder config so we can derive
// hw_frames_ctx from the buffersink output (correct surface pool dimensions).
let mut video_filter = build_filter_graph(
&hw_device_ctx,
&frames_rgb,
width,
height,
enc_width,
enc_height,
fps,
transform,
)?;
let mut sink_ctx = video_filter
.get("out")
// `ok_or_else(|| ...)` 把 `Option<T>` 转为 `Result<T, E>`None 时执行闭包构造错误。
// 闭包 `|| ...` 延迟构造(类比 Go `if x == nil { return fmt.Errorf(...) }`);
// 末尾 `?` 把 Err 传播给调用方,Ok(t) 则解包继续。
.ok_or_else(|| anyhow::anyhow!("filter 'out' not found"))?;
// 中文概述:从 buffersink 取出硬件帧上下文并 `av_buffer_ref` 提升为拥有引用。
// SAFETY: sink_ctx is a live buffersink; the returned hw_frames_ctx is
// borrowed, so av_buffer_ref creates an owned reference.
let sink_hw_frames = unsafe {
let raw = ffi::av_buffersink_get_hw_frames_ctx(sink_ctx.as_mut_ptr());
if raw.is_null() {
bail!("buffersink has no hw_frames_ctx — filter graph may not be configured for hardware output");
}
let hw_ref = ffi::av_buffer_ref(raw);
if hw_ref.is_null() {
bail!("av_buffer_ref failed for buffersink hw_frames_ctx — likely out of memory");
}
hw_ref
};
// 中文概述:解引用 sink_hw_frames 检查 filter graph 输出尺寸与编码器期望尺寸是否一致;
// 不一致仅 warnfilter 可能做隐式 scale),不视为硬错误。
// SAFETY: sink_hw_frames is an owned AVBufferRef to an AVHWFramesContext
// returned by the validated filter graph.
unsafe {
let fc = (*sink_hw_frames).data as *mut ffi::AVHWFramesContext;
let actual_w = (*fc).width as u32;
let actual_h = (*fc).height as u32;
if actual_w != enc_width || actual_h != enc_height {
tracing::warn!(
"Filter output dimensions {actual_w}x{actual_h} differ from encoder dimensions {enc_width}x{enc_height}"
);
}
}
// 4. Find h264_vaapi encoder
// `ff::encoder::find_by_name("h264_vaapi")` 返回 `Option<Codec>`FFmpeg 未编译 VAAPI 时为 None。
// `.ok_or_else(...)?` 链式:Option → Result → 自动传播。
let codec = ff::encoder::find_by_name("h264_vaapi")
.ok_or_else(|| anyhow::anyhow!("h264_vaapi encoder not found"))?;
// 块表达式 `{ ... }` 求值为最后一个表达式(无分号)的值:这里 `ctx.encoder().video()?`
// 返回 `Result<Encoder<Video>>``?` 解开为 `Encoder<Video>`,作为整个块的值赋给 `enc`。
// 类比 Go 的 `enc := func() *Encoder { ... return x }()` 但 Rust 的块是表达式级。
let mut enc = {
let ctx = ff::codec::Context::new_with_codec(codec);
ctx.encoder().video()?
};
enc.set_width(enc_width);
enc.set_height(enc_height);
enc.set_format(ff::format::Pixel::VAAPI);
enc.set_bit_rate(bitrate as usize);
enc.set_gop(gop_size);
enc.set_time_base(ff::Rational::new(1, fps as i32));
enc.set_max_b_frames(0);
// VBV rate limiting: caps IDR burst size for WebRTC. Without this a 4K
// scene change can produce a 256KB keyframe that overflows the UDP send
// buffer. bufsize=bitrate/4 ≈ 250ms of video at the target bitrate.
// 中文概述:直接写裸指针字段 rc_max_rate/rc_buffer_size 设置 VBV 上限,避免 4K IDR 帧击爆 UDP 缓冲。
unsafe {
let ctx_ptr = enc.as_mut_ptr();
(*ctx_ptr).rc_max_rate = bitrate as i64;
(*ctx_ptr).rc_buffer_size = (bitrate / 4) as i32;
}
// 中文概述:AV_CODEC_FLAG_GLOBAL_HEADER 必须在 open 之前设置,触发 SPS/PPS extradata 生成(muxer 做 Annex B → AVCC 转换时需要)。
// SAFETY: AV_CODEC_FLAG_GLOBAL_HEADER must be set BEFORE opening the encoder.
// It triggers SPS/PPS extradata generation needed by the muxer for
// Annex B to AVCC conversion.
unsafe {
(*enc.as_mut_ptr()).flags |= ffi::AV_CODEC_FLAG_GLOBAL_HEADER as i32;
}
// 中文概述:把 hw_device_ctx 与 sink_hw_framesfilter 输出帧池)引用赋给编码器,编码器从该池分配 VAAPI surface。
// SAFETY: Assign hw device and frames ctx to the encoder.
unsafe {
(*enc.as_mut_ptr()).hw_device_ctx = hw_device_ctx.ref_clone();
(*enc.as_mut_ptr()).hw_frames_ctx = sink_hw_frames;
}
// SAFETY: Set repeat_pps=1 on the encoder so PPS is inserted in every encoded frame.
// This ensures decoders can start decoding from any frame (important for WebRTC).
// Note: repeat_pps is only available in FFmpeg 7.0+ (not in 6.x). On older FFmpeg,
// IDR frames carry SPS by default; PPS repetition depends on the driver.
// For SPS repetition: IDR frames carry SPS by default, controlled by gop_size/idr_interval.
{
let key = CString::new("repeat_pps").unwrap();
let val = CString::new("1").unwrap();
let ret = unsafe {
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0)
};
if ret < 0 {
tracing::warn!("av_opt_set repeat_pps failed ({}), likely FFmpeg < 7.0; continuing without per-frame PPS", ff_err(ret));
}
}
// 5. Open encoder. Video::open() returns Encoder(Video); .0 extracts the Video.
// `.map_err(|e| ...)?` 用闭包包装错误上下文(保留原始 ffmpeg::Error,附加上下文消息),
// `?` 自动传播。元组解构 `opened.0``Encoder(Video)` 是 newtype 包装,`.0` 取内部 Video。
let opened = enc
.open()
.map_err(|e| anyhow::anyhow!("Failed to open h264_vaapi encoder: {e}"))?;
let enc_video = opened.0;
// 6. Muxer setup (strict order)
// `CString` 是 Rust ↔ C FFI 的桥梁:FFmpeg C API 需要以 `\0` 结尾的字符串,
// Rust `String`/`&str` 不带终止符,必须用 `CString::new` 包装。
// `output_path.to_str().unwrap()` 把 `Path` 转 UTF-8 字符串切片(非 UTF-8 路径会 panic)。
let output_cstr = CString::new(output_path.to_str().unwrap())?;
// 裸指针声明:`*mut ffi::AVFormatContext` 是 C 指针(FFI 类型),Rust 中默认不可解引用。
// `ptr::null_mut()` 等价于 C 的 `NULL`FFmpeg C API 接收 `**` 时会重新分配并写回。
let mut fmt_ctx_ptr: *mut ffi::AVFormatContext = ptr::null_mut();
// 中文概述:根据文件扩展名(.mp4/.mkv/...)分配输出 format context,但不打开文件。
// SAFETY: avformat_alloc_output_context2 creates format context from
// the file extension. Does NOT open the file.
let ret = unsafe {
ffi::avformat_alloc_output_context2(
&mut fmt_ctx_ptr,
ptr::null_mut(),
ptr::null(),
output_cstr.as_ptr(),
)
};
if ret < 0 || fmt_ctx_ptr.is_null() {
bail!("Failed to allocate output format context: {}", ff_err(ret));
}
// 中文概述:检查编码器 codec_id 与 muxer oformat 是否兼容(H.264 + MP4 应返回 ≥0)。
// SAFETY: avformat_query_codec checks codec+format compatibility.
let codec_id = unsafe { (*enc_video.as_ptr()).codec_id };
let oformat = unsafe { (*fmt_ctx_ptr).oformat };
let compat = unsafe {
ffi::avformat_query_codec(oformat, codec_id, ffi::FF_COMPLIANCE_NORMAL as i32)
};
if compat < 0 {
bail!("H.264 codec not supported by output container format");
}
// 中文概述:在 format context 中新建一条 stream(视频流),返回的指针由 fmt_ctx 拥有,无需手动释放。
// SAFETY: avformat_new_stream creates a new stream in the format context.
let stream_ptr = unsafe { ffi::avformat_new_stream(fmt_ctx_ptr, ptr::null()) };
if stream_ptr.is_null() {
bail!("Failed to create new stream in output context");
}
// 中文概述:从 encoder 拷贝参数(含 SPS/PPS extradata)到 stream->codecpar,供 demuxer 解析。
// SAFETY: avcodec_parameters_from_context copies encoder params + extradata.
let ret = unsafe {
ffi::avcodec_parameters_from_context((*stream_ptr).codecpar, enc_video.as_ptr())
};
if ret < 0 {
bail!(
"Failed to copy encoder parameters to stream: {}",
ff_err(ret)
);
}
// 中文概述:把 encoder 的 time_base1/fps)拷贝到 stream,muxer 写头部时按此时基打时间戳。
// SAFETY: Copy encoder time_base to stream.
unsafe {
(*stream_ptr).time_base = (*enc_video.as_ptr()).time_base;
}
// 中文概述:打开输出文件 IO 上下文(AVIO),把 pbAVIOContext*)挂到 fmt_ctx。
// SAFETY: avio_open opens the output file for writing.
let ret = unsafe {
ffi::avio_open(
&mut (*fmt_ctx_ptr).pb,
output_cstr.as_ptr(),
ffi::AVIO_FLAG_WRITE,
)
};
if ret < 0 {
bail!(
"Failed to open output file '{}': {}",
output_path.display(),
ff_err(ret)
);
}
// 中文概述:写入容器头部(MP4 ftyp box 等),失败需关闭 IO 释放资源(此处简化为直接 bail)。
// SAFETY: avformat_write_header writes the container header.
let ret = unsafe { ffi::avformat_write_header(fmt_ctx_ptr, ptr::null_mut()) };
if ret < 0 {
bail!("Failed to write output header: {}", ff_err(ret));
}
// 中文概述:把裸 fmt_ctx_ptr 包装回 safe Rust 的 Output 类型,由 octx 的 Drop 负责最终释放。
// SAFETY: We created fmt_ctx_ptr above and it's valid.
let octx = unsafe { ff::format::context::Output::wrap(fmt_ctx_ptr) };
// `Ok(Self { ... })` 是 Rust 构造语法:struct literal 字段简写(`enc_video` 等价于 `enc_video: enc_video`)。
// 类比 Go `return &EncState{enc_video: ..., ...}, nil`,但 Rust 的 `Self` 是类型别名。
Ok(Self {
enc_video,
frames_rgb,
video_filter,
hw_device_ctx,
octx,
starting_timestamp: None,
frames_written: false,
})
}
// 共享访问器:返回内部 hw frames 池的引用,捕获侧(state.rs/state_portal.rs)用它
// 给 `import_dma_buf_to_vaapi` 分配目标 VAAPI surface。借用检查器保证调用方不能
// 在持有 `&AvHwFrameCtx` 期间 `&mut self`(例如调用 `encode_frame`)。
pub fn frames_rgb(&self) -> &AvHwFrameCtx {
&self.frames_rgb
}
// 单帧驱动主循环:`&mut self` 表示独占借用(类比 Go `*receiver` + 单线程保证)。
// 接收一个 VAAPI 硬件帧,通过 filter 图(BGRA→NV12 + scale/crop)后送给编码器。
pub fn encode_frame(&mut self, hw_frame: &ff::frame::Video) -> Result<()> {
let mut filter_src_ctx = self
.video_filter
.get("in")
.ok_or_else(|| anyhow::anyhow!("filter 'in' not found"))?;
let mut filter_src = filter_src_ctx.source();
let mut filter_sink_ctx = self
.video_filter
.get("out")
.ok_or_else(|| anyhow::anyhow!("filter 'out' not found"))?;
let mut filter_sink = filter_sink_ctx.sink();
// SAFETY: hw_frame is a valid VAAPI hardware frame from capture.
filter_src
.add(hw_frame)
.map_err(|e| anyhow::anyhow!("Filter source add failed: {e}"))?;
// 持续从 sink 拉过滤后帧,直到 EAGAINfilter 缓冲空)。
// `loop {}` 是 Rust 无限循环(类比 Go `for {}`),靠 `break` 退出。
loop {
// `ff::frame::Video::empty()` 分配零值帧(无数据),filter_sink.frame() 填充它。
// `&mut filtered` 把可变借用传给 ffmpeg-next,类比 Go `&filtered`。
let mut filtered = ff::frame::Video::empty();
// 多路 matchOk 表示成功;带 guard 的 `Err(ff::Error::Other { errno }) if errno == ffi::EAGAIN`
// 是 pattern guard(类比 Go `switch v := err.(type); ...`);通配 `Err(e)` 兜底。
match filter_sink.frame(&mut filtered) {
Ok(()) => {
// 若 filter 没设 PTSpts()==None),用输入 hw_frame 的 PTS 兜底。
if filtered.pts().is_none() {
filtered.set_pts(hw_frame.pts());
}
}
Err(ff::Error::Other { errno }) if errno == ffi::EAGAIN => break,
Err(e) => bail!("Filter sink get frame failed: {e}"),
}
let pts = filtered.pts().unwrap_or(0);
if self.starting_timestamp.is_none() {
self.starting_timestamp = Some(pts);
}
let start_ts = self.starting_timestamp.unwrap();
// 中文概述:把过滤后的 NV12 VAAPI surface 发送给编码器输入队列;<0 表示 FFI 失败。
// SAFETY: avcodec_send_frame sends a valid NV12 VAAPI surface to the encoder.
let ret =
unsafe { ffi::avcodec_send_frame(self.enc_video.as_mut_ptr(), filtered.as_ptr()) };
if ret < 0 {
bail!("avcodec_send_frame failed: {}", ff_err(ret));
}
// 立即 drain 编码器输出队列:FFmpeg 编码可能延迟一帧或多帧才产出 packet,
// send_frame 后必须 receive_packet 直到 EAGAIN,否则编码器内部缓冲溢出。
self.drain_encoder(start_ts)?;
}
Ok(())
}
// EOF 处理:在调用方(state.rs/state_portal.rs)完成所有 `encode_frame` 后调用,
// 负责把 filter 图残余帧排空 + 触发编码器 EOS + 写 muxer trailermoov box 等)。
pub fn flush(&mut self) -> Result<()> {
// Flush filter graph
let mut filter_src_ctx = self
.video_filter
.get("in")
.ok_or_else(|| anyhow::anyhow!("filter 'in' not found"))?;
let mut filter_src = filter_src_ctx.source();
// `if let Err(e) = ...` 模式匹配:只关心 Err 分支,Ok 丢弃。类比 Go `if err := ...; err != nil {}`。
if let Err(e) = filter_src.flush() {
tracing::debug!("filter source flush error: {e}");
}
// Drain filter
let mut filter_sink_ctx = self
.video_filter
.get("out")
.ok_or_else(|| anyhow::anyhow!("filter 'out' not found"))?;
// Drain filter:循环拉 sink,把 filter 图中残余的帧全部送入编码器。
// 与 encode_frame 不同:这里 Err 通配 `Err(_) => break`(任何错误都视为 drain 完毕),
// 因为 flush 是收尾阶段,不再传播错误。
let mut filter_sink = filter_sink_ctx.sink();
loop {
let mut filtered = ff::frame::Video::empty();
match filter_sink.frame(&mut filtered) {
Ok(()) => {
// `unwrap_or(0)` 是 Option 的兜底方法:None 返回 0Some(v) 返回 v。
// flush 时若从未编码过帧(starting_timestamp=None),用 0 作为 PTS 基线。
let start_ts = self.starting_timestamp.unwrap_or(0);
// 中文概述:把 drain 出的帧送编码器;flush 路径失败也 bail(停止 flush)。
// SAFETY: filtered is a valid VAAPI frame drained from the
// filter graph; enc_video is an opened encoder.
let ret = unsafe {
ffi::avcodec_send_frame(self.enc_video.as_mut_ptr(), filtered.as_ptr())
};
if ret < 0 {
bail!("avcodec_send_frame failed during flush: {}", ff_err(ret));
}
self.drain_encoder(start_ts)?;
}
Err(_) => break,
}
}
// 中文概述:发送 NULL 帧给编码器,触发 EOSEnd Of Stream),编码器内部缓冲的帧会被强制输出。
// SAFETY: Sending null frame signals end of stream to encoder.
unsafe {
ffi::avcodec_send_frame(self.enc_video.as_mut_ptr(), ptr::null());
}
let start_ts = self.starting_timestamp.unwrap_or(0);
// 最后一次 drain:取出 EOS 触发后所有剩余 packet。
self.drain_encoder(start_ts)?;
// Write trailer only if at least one frame was encoded.
// 写 trailer 条件:若全程 0 帧编码,写 trailer 会让 muxer 生成空 moov(损坏文件),所以用 frames_written 守卫。
if self.frames_written {
self.octx
.write_trailer()
.map_err(|e| anyhow::anyhow!("Failed to write trailer: {e}"))?;
}
Ok(())
}
// 内部辅助方法(无 pub):每次 send_frame 后调用,循环 receive_packet 直到 EAGAIN/EOF。
// `start_ts: i64` 是首帧 PTS,用于把所有 packet 的 PTS/DTS 减去首帧 PTS(归零起点)。
fn drain_encoder(&mut self, start_ts: i64) -> Result<()> {
loop {
let mut pkt = ff::Packet::empty();
// 中文概述:从编码器拉一个已编码 packet;返回 <0 时区分 EAGAIN(暂时无)和 EOF(流结束)。
// SAFETY: avcodec_receive_packet retrieves an encoded packet.
let ret = unsafe {
ffi::avcodec_receive_packet(self.enc_video.as_mut_ptr(), pkt.as_mut_ptr())
};
if ret < 0 {
// 嵌套 if:先判 ret<0FFI 错误),再判具体码(EAGAIN/EOF)。
// `ffi::AVERROR(ffi::EAGAIN)` 是 FFmpeg 的 errno 包装宏(负数)。
if ret == ffi::AVERROR(ffi::EAGAIN) || ret == ffi::AVERROR_EOF {
break;
}
bail!("avcodec_receive_packet failed: {}", ff_err(ret));
}
// Rescale timestamps from encoder time_base to stream time_base
// 时间基重缩放:encoder 用 1/fpsstream 可能用 1/1000 等不同基。
// rescale_ts 用有理数运算,避免浮点误差累积。
let enc_tb = self.enc_video.time_base();
// 中文概述:通过裸指针读 octx.streams[0].time_base,得到 stream 的时间基(由 muxer 决定)。
// SAFETY: octx was created with stream 0 during muxer setup; streams
// is non-null and stream 0 remains owned by the format context.
let stream_tb = unsafe {
let fmt = *self.octx.as_ptr();
if fmt.nb_streams == 0 || fmt.streams.is_null() {
bail!("no streams in output context");
}
// `fmt.streams.add(0)` 是裸指针算术: streams 是 `*mut *mut AVStream`
// `.add(0)` 取首元素地址,外层 `*` 解引用得到 `*mut AVStream`。
let st = *fmt.streams.add(0);
ff::Rational::from((*st).time_base)
};
pkt.rescale_ts(enc_tb, stream_tb);
// Offset timestamps so first frame starts at 0
// PTS/DTS 归零:把首帧 PTS 减去自身得到 0,后续帧相对首帧的偏移。`if let Some(pts) = ...`
// 模式匹配:只有当 PTS 存在时才改写(部分 packet 可能没有 PTS)。
if let Some(pts) = pkt.pts() {
pkt.set_pts(Some(pts - start_ts));
}
if let Some(dts) = pkt.dts() {
pkt.set_dts(Some(dts - start_ts));
}
// 标记 packet 属于 stream 0(视频流)。
pkt.set_stream(0);
// 交错写入:muxer 内部会按 DTS 排序后写出,类比 Go `mp4Writer.WriteInterleaved(pkt)`。
pkt.write_interleaved(&mut self.octx)
.map_err(|e| anyhow::anyhow!("Failed to write packet: {e}"))?;
// 标记:至少一帧已写入 muxer(flush 时用此守卫决定是否写 trailer)。
self.frames_written = true;
}
Ok(())
}
}
// ---------------------------------------------------------------------------
// SwEncState - VAAPI GPU downscale + software H.264 encode
// ---------------------------------------------------------------------------
// 编码完成的 H.264 帧载体,是软件编码路径(SwEncEncode → WebRTC 通道)的输出单元。
// 与硬件路径 EncState 不同:硬件路径直接对接 avformat muxerPTS 由 muxer 内部维护;
// 软件路径通过 crossbeam channel 把帧送到 webrtc 线程,必须显式携带 PTS 与捕获时间戳。
//
// 字段语义:
// - dataH.264 NAL 字节流(Annex B 或 AVCC,取决于 encoder 配置;本工程用 libx264 默认 Annex B
// - pts_ticks:编码器 time_base 单位下的 PTSWebRTC 模式下 time_base=1/90000
// 直接对应 RTP 时间戳,90kHz 精度)
// - capture_time:墙上时钟捕获时刻,从 CpuNv12Frame 透传,用于 stats 的 frame_age 指标
//
// issue #24 背景:Wayland damage-driven 帧率波动大,若 PTS 用帧计数 * (1/fps) 估算,
// 浏览器 jitter buffer 会膨胀到秒级;必须用真实捕获时刻换算 PTS。
//
// 与 Go 对照:类似 Go 视频包结构体 `type Frame struct { Data []byte; PTSTicks int64; CaptureTime time.Time }`。
/// Encoded H.264 frame with timing metadata for WebRTC output.
///
/// MP4 file output (FrameOutput::Muxer) does NOT use this - it writes via
/// avformat which preserves PTS internally. WebRTC output (FrameOutput::Channel)
/// requires explicit PTS propagation so RTP timestamps reflect real capture time.
/// Without this, WebRTC clients' jitter buffers grow to seconds under
/// damage-driven variable frame rate. See issue #24.
#[derive(Debug)]
pub struct EncodedH264Frame {
/// H.264 NAL byte stream (Annex B or AVCC depending on encoder configuration)
pub data: Vec<u8>,
/// PTS in encoder time_base units (1/fps seconds), normalized so first frame = 0.
/// Derived from real capture time, NOT frame counter.
pub pts_ticks: i64,
/// Wall-clock capture time, propagated from CpuNv12Frame for frame_age stat.
pub capture_time: std::time::Instant,
}
// 编码器输出目的地二选一枚举:
// - Muxer:直接对接 FFmpeg avformat 写 MP4 文件(自包含 PTS,硬编码路径 EncState 风格)
// - Channelcrossbeam Sender<EncodedH264Frame>,把帧送 WebRTC 信令线程(必须显式带 PTS)
// 类比 Go`type FrameOutput interface { Write(Frame) error }` 的两种实现。
// 该枚举让 SwEncEncode 在 new_muxer / new_webrtc 两个构造函数间共享同一 drain_encoder 主循环。
pub enum FrameOutput {
Muxer(ff::format::context::Output),
Channel(crossbeam_channel::Sender<EncodedH264Frame>),
}
// 主线程持有:完成 VAAPI GPU 下采样 + av_hwframe_transfer_data 回到 CPU 的 NV12 帧缓冲。
// 数据流:PipeWire/wlr 帧通过 import_and_scale 进入 → VAAPI scale_vaapi 转 NV12 →
// transfer_filtered_to_cpu 拷贝到 Vec<u8> → 通过 channel 跨线程送 SwEncEncode 编码线程。
// NV12 格式:Y 平面满分辨率,UV 交错平面半高(chroma 高度减半);stride 可能 > width
// (对齐到 GPU 要求)。字段不带 GPU 资源(纯 CPU 内存),自动 Send。
/// Owned CPU NV12 frame data for cross-thread transfer.
/// Produced by main thread (VAAPI import + GPU scale + transfer), consumed by encode thread.
pub struct CpuNv12Frame {
pub y_data: Vec<u8>,
pub uv_data: Vec<u8>,
pub y_stride: usize,
pub uv_stride: usize,
pub pts: i64,
/// Wall-clock time when this frame was captured (PipeWire delivery).
/// Used for frame_age stat: time from capture to WebRTC send.
pub capture_time: std::time::Instant,
}
// 软件编码路径的"导入 + GPU 缩放"半部:负责把外部 DMA-BUF 帧(PipeWire/wlr)经
// VAAPI 硬件零拷贝导入 → scale_vaapi GPU 下采样 → av_hwframe_transfer_data 回传 CPU NV12
// 最后封装为 CpuNv12Frame 喂给 SwEncEncodelibx264 软编)。
//
// 与 EncState(硬编直连)的区别:本结构仅做 GPU 辅助的导入/缩放,H.264 编码由 CPU 完成;
// 因此自始至终只在导入线程访问,跨线程边界是 channelCpuNv12Frame)。
//
// 字段:
// - hw_dev / frames_rgbVAAPI 设备 + BGRA 硬件帧上下文(DMA-BUF 导入目的地)
// - filter_graphFFmpeg 滤镜图(scale_vaapi=M,w=enc:h=enc,format=nv12
// - width/height:源帧尺寸;enc_width/enc_height:目标缩放后尺寸
// - resolution_rx / encoder_resolution_tx:可选的动态分辨率控制通道对
pub struct SwEncImport {
hw_dev: AvHwDevCtx,
frames_rgb: AvHwFrameCtx,
filter_graph: ff::filter::Graph,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
fps: u32,
resolution_rx: Option<crossbeam_channel::Receiver<BitrateCommand>>,
encoder_resolution_tx: Option<crossbeam_channel::Sender<ResolutionChange>>,
}
// impl SwEncImport5 个方法按数据流顺序
// (1) new / new_with_resolution_control:构造函数(建立 VAAPI + filter graph
// (2) frames_rgb:访问器(暴露硬件帧上下文给外部 DMA-BUF 导入)
// (3) import_and_scale:单帧主流程(filter_src.add → filter_sink.frame 循环 → transfer_to_cpu
// (4) flush_importEOS 时排空 filter graph 缓冲
// (5) poll_resolution_commands + transfer_filtered_to_cpu:私有辅助
impl SwEncImport {
#[allow(clippy::too_many_arguments)]
pub fn new(
drm_device: &Path,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
fps: u32,
) -> Result<Self> {
// drm_device 形如 /dev/dri/renderD128AvHwDevCtx::new_vaapi 内部打开 DRM fd 并
// 调 vaGetDisplayDRM / vaInitialize 建立 VAAPI 连接(FFmpeg hwcontext_vaapi)。
// ? 自动传播 Result(类比 Go `if err != nil { return err }`)。
let hw_dev = AvHwDevCtx::new_vaapi(drm_device)?;
// 硬件帧上下文:调用 vaCreateSurfaces 申请 BGRA surface 池供 DMA-BUF 导入使用;
// AV_PIX_FMT_BGRA 对应 PipeWire/wlr 帧的常见 RGB 格式。返回的 AvHwFrameCtx 持有
// AVHWFramesContextsurface 池大小由 FFmpeg 默认(通常 4-8 帧)。
let frames_rgb =
AvHwFrameCtx::for_capture(&hw_dev, width, height, ff::format::Pixel::BGRA)?;
// 建立 scale_vaapi 滤镜图:源=BGRA/VAAPIscale=enc_width×enc_heightformat=nv12。
// 滤镜图绑定的 hw_dev 与 frames_rgb 必须同属一个 VAAPI 设备,否则 vaCreateSurfaces
// 在不同 display 上会失败。
let filter_graph = build_swenc_filter_graph(
&hw_dev,
&frames_rgb,
width,
height,
enc_width,
enc_height,
fps,
)?;
// Self 是返回类型的别名(Rust 构造器惯用法,类比 Go `return &Foo{...}` 中的 Foo)。
// 字段简写:当局部变量名与字段名相同时可省略 `field: value`,等价于 Go 结构体字面量简写。
// resolution_rx / encoder_resolution_tx 留 None,由 new_with_resolution_control 后续填充。
Ok(Self {
hw_dev,
frames_rgb,
filter_graph,
width,
height,
enc_width,
enc_height,
fps,
resolution_rx: None,
encoder_resolution_tx: None,
})
}
#[allow(clippy::too_many_arguments)]
pub fn new_with_resolution_control(
drm_device: &Path,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
fps: u32,
resolution_rx: crossbeam_channel::Receiver<BitrateCommand>,
encoder_resolution_tx: crossbeam_channel::Sender<ResolutionChange>,
) -> Result<Self> {
let mut this = Self::new(drm_device, width, height, enc_width, enc_height, fps)?;
this.resolution_rx = Some(resolution_rx);
this.encoder_resolution_tx = Some(encoder_resolution_tx);
Ok(this)
}
pub fn frames_rgb(&self) -> &AvHwFrameCtx {
let _ = self.hw_dev.as_ptr();
&self.frames_rgb
}
pub fn import_and_scale(&mut self, hw_frame: &ff::frame::Video) -> Result<CpuNv12Frame> {
// 先消费 ResolutionChange 命令:若分辨率改变,本函数后续的 filter graph 已是新图。
self.poll_resolution_commands()?;
// 拿到滤镜图的输入/输出端点("in"/"out" 是 build_swenc_filter_graph 中注册的名称)。
// ok_or_else 把 Option<Result> 转 Result(延迟构造错误,类比 Go `if x == nil { return err }`)。
// ? 自动传播 Resultmut 因为 source()/sink() 返回 &mut 借用。
let mut filter_src_ctx = self
.filter_graph
.get("in")
.ok_or_else(|| anyhow::anyhow!("filter 'in' not found"))?;
let mut filter_src = filter_src_ctx.source();
let mut filter_sink_ctx = self
.filter_graph
.get("out")
.ok_or_else(|| anyhow::anyhow!("filter 'out' not found"))?;
let mut filter_sink = filter_sink_ctx.sink();
// 把外部导入的硬件帧(hw_frame 是 AV_PIX_FMT_DRM_PRIME/VAAPIVASurfaceID 在 data[3]
// 送入滤镜图入口。滤镜图内部 scale_vaapi 执行 GPU 下采样 + 格式转换。
// map_err 把 ffmpeg::Error 包装成 anyhow::Error 加上下文。
filter_src
.add(hw_frame)
.map_err(|e| anyhow::anyhow!("software pipeline filter source add failed: {e}"))?;
// filter graph 可能产生多帧(理论上 scale_vaapi 1:1,但 deinterlace/ fps 滤镜可能 1:N)。
// 这里只取第一帧;额外帧计数后丢并 warn。EAGAIN 表示滤镜图已 drain 完毕。
let mut first = None;
let mut extra_count = 0usize;
loop {
let mut filtered = ff::frame::Video::empty();
match filter_sink.frame(&mut filtered) {
Ok(()) => {
// 若 filter 没有显式设置 PTS(罕见),用源帧 PTS 兜底。
if filtered.pts().is_none() {
filtered.set_pts(hw_frame.pts());
}
let cpu_frame = self.transfer_filtered_to_cpu(&filtered)?;
if first.is_none() {
first = Some(cpu_frame);
} else {
extra_count += 1;
}
}
// ffi::EAGAIN:滤镜图当前没有输出可取;这是正常的"等下一帧"信号,break 主循环。
// 类比 Go 中 io.Read 返回 EAGAIN。
Err(ff::Error::Other { errno }) if errno == ffi::EAGAIN => break,
Err(e) => bail!("software pipeline filter sink get frame failed: {e}"),
}
}
if extra_count > 0 {
tracing::warn!(
"software import filter produced {extra_count} extra frame(s); dropping extras"
);
}
// first 一定存在:filter_src.add 成功后 scale_vaapi 至少产出一帧;
// 若不产出说明上游 filter graph 配置错误或 hw_frame 无效,返回错误。
first.ok_or_else(|| anyhow::anyhow!("software pipeline produced no scaled frame"))
}
pub fn flush_import(&mut self) -> Result<Vec<CpuNv12Frame>> {
let mut filter_src_ctx = self
.filter_graph
.get("in")
.ok_or_else(|| anyhow::anyhow!("filter 'in' not found"))?;
let mut filter_src = filter_src_ctx.source();
if let Err(e) = filter_src.flush() {
tracing::debug!("filter source flush error: {e}");
}
let mut filter_sink_ctx = self
.filter_graph
.get("out")
.ok_or_else(|| anyhow::anyhow!("filter 'out' not found"))?;
let mut filter_sink = filter_sink_ctx.sink();
let mut frames = Vec::new();
loop {
let mut filtered = ff::frame::Video::empty();
match filter_sink.frame(&mut filtered) {
Ok(()) => frames.push(self.transfer_filtered_to_cpu(&filtered)?),
Err(_) => break,
}
}
Ok(frames)
}
fn poll_resolution_commands(&mut self) -> Result<()> {
let Some(rx) = self.resolution_rx.as_ref().cloned() else {
return Ok(());
};
let mut requested = None;
while let Ok(cmd) = rx.try_recv() {
match cmd {
BitrateCommand::UpdateResolution { width, height } => {
requested = Some((width & !1, height & !1));
}
BitrateCommand::UpdateBitrate { .. } => {}
BitrateCommand::ForceKeyframe => {}
}
}
let Some((width, height)) = requested else {
return Ok(());
};
if width == self.enc_width && height == self.enc_height {
return Ok(());
}
tracing::info!(
from = format_args!("{}x{}", self.enc_width, self.enc_height),
to = format_args!("{}x{}", width, height),
"rebuilding software import filter graph for resolution change"
);
let _ = self.flush_import();
self.filter_graph = build_swenc_filter_graph(
&self.hw_dev,
&self.frames_rgb,
self.width,
self.height,
width,
height,
self.fps,
)?;
self.enc_width = width;
self.enc_height = height;
if let Some(tx) = &self.encoder_resolution_tx {
tx.send(ResolutionChange { width, height })
.map_err(|_| anyhow::anyhow!("encoder resolution channel disconnected"))?;
}
Ok(())
}
fn transfer_filtered_to_cpu(&self, filtered: &ff::frame::Video) -> Result<CpuNv12Frame> {
// 中文概述:分配一个 CPU 侧的 AVFrame 作为 av_hwframe_transfer_data 的目标。
// 英文 SAFETY 详细论证见下方。
// SAFETY: av_frame_alloc returns a newly allocated AVFrame or null,
// which is checked below.
let mut sw_nv12 = unsafe { ffi::av_frame_alloc() };
if sw_nv12.is_null() {
bail!("av_frame_alloc failed for NV12 transfer frame");
}
// 中文概述:GPU→CPU 数据回传。filtered 指向 VAAPI NV12 surface
// av_hwframe_transfer_data 内部调用 vaGetImage 等同步原语把 GPU 显存拷贝到 sw_nv12。
// 英文 SAFETY 论证下方。
// SAFETY: sw_nv12 is an allocated destination frame; filtered is a valid VAAPI NV12
// surface produced by scale_vaapi at encoder dimensions.
let transfer_ret = unsafe { ffi::av_hwframe_transfer_data(sw_nv12, filtered.as_ptr(), 0) };
if transfer_ret < 0 {
// 中文概述:transfer 失败时必须释放 sw_nv12 防止泄漏;FFmpeg C API 无 RAII。
// SAFETY: sw_nv12 was allocated above and has not been freed yet.
unsafe { ffi::av_frame_free(&mut sw_nv12) };
bail!(
"av_hwframe_transfer_data failed for GPU-downscaled frame: {}",
ff_err(transfer_ret)
);
}
// 中文概述:从 AVFrame 字段读 Y/UV 平面指针、linesize、尺寸,校验后拷贝到 Vec<u8>。
// 错误路径同样调用 av_frame_free 释放。最后再次 av_frame_free 释放 sw_nv12 自身。
// 英文 SAFETY 详细论证见下方。
// SAFETY: sw_nv12 was filled by av_hwframe_transfer_data. NV12 planes 0 and 1 are
// initialized for enc_width x enc_height; linesize values define each row's byte span.
let frame = unsafe {
let y_ptr = (*sw_nv12).data[0];
let uv_ptr = (*sw_nv12).data[1];
if y_ptr.is_null() || uv_ptr.is_null() {
ffi::av_frame_free(&mut sw_nv12);
bail!("NV12 transfer frame missing Y/UV plane data");
}
let y_stride = (*sw_nv12).linesize[0] as usize;
let uv_stride = (*sw_nv12).linesize[1] as usize;
if (*sw_nv12).width != self.enc_width as i32
|| (*sw_nv12).height != self.enc_height as i32
{
ffi::av_frame_free(&mut sw_nv12);
bail!("NV12 transfer frame has unexpected dimensions");
}
let y_len = y_stride * self.enc_height as usize;
let uv_len = uv_stride * (self.enc_height as usize / 2);
// slice::from_raw_parts:把 C 的裸指针 + 长度变成 Rust 切片(零拷贝)。
// 紧接着 .to_vec() 立即深拷贝,避免 sw_nv12 释放后切片悬空。
let y_data = slice::from_raw_parts(y_ptr, y_len).to_vec();
let uv_data = slice::from_raw_parts(uv_ptr, uv_len).to_vec();
let pts = filtered.pts().unwrap_or(0);
ffi::av_frame_free(&mut sw_nv12);
CpuNv12Frame {
y_data,
uv_data,
y_stride,
uv_stride,
pts,
capture_time: std::time::Instant::now(),
}
};
Ok(frame)
}
}
// 软件编码路径的"CPU 编码"半部:消费 CpuNv12FrameNV12),用 sws_scale 转 YUV420P
// 喂 libx264 软件 H.264 编码器,输出 EncodedH264Frame。与 SwEncImport 配对使用:
// SwEncImport 负责 GPU 辅助导入/缩放,SwEncEncode 负责 CPU 编码;两者通过 channel 通信。
//
// 持有的 C 资源(裸指针,非 Send/Sync 自动,需手写 unsafe impl Send 见下):
// - sws_ctxFFmpeg SwsContextNV12 → YUV420P 颜色空间转换器,无 resize)
// - yuv_frame:可复用的 AVFrameYUV420P,尺寸 = enc_width × enc_height
// - enc_video:打开的 libx264 AVCodecContext
//
// 其他字段是控制状态(动态码率/分辨率/关键帧请求、去重哈希、计时统计、WebRTC 暂停/断开)。
// 详见每个字段上方英文 ///(保留)。
pub struct SwEncEncode {
sws_ctx: *mut ffi::SwsContext,
enc_video: ff::codec::encoder::video::Video,
output: Option<FrameOutput>,
yuv_frame: *mut ffi::AVFrame,
last_frame_hash: u64,
frame_count: u64,
starting_timestamp: Option<i64>,
frames_written: bool,
webrtc_disconnected: bool,
webrtc_paused: Option<Arc<AtomicBool>>,
bitrate_rx: crossbeam_channel::Receiver<BitrateCommand>,
resolution_rx: crossbeam_channel::Receiver<ResolutionChange>,
enc_width: u32,
enc_height: u32,
fps: u32,
bitrate: u64,
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.
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.
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).
last_capture_time: Option<Instant>,
}
// FNV-1a 64 位哈希常量,用于帧去重(见 encode_cpu_frame 调用 hash_sampled_y_plane):
// - FNV1A_OFFSET_BASIS:初始哈希值(FNV 论文推荐的 64 位 offset basis
// - FNV1A_PRIME:每字节乘法因子(64 位素数)
// - Y_PLANE_HASH_ROW_STEPY 平面采样步长(仅哈希每 8 行,性能与碰撞率权衡)
// 算法:hash = (hash XOR byte) * PRIME,每字节迭代;XOR + 乘法双混淆,分布均匀。
// 选择 FNV-1a 而非 CRC32:性能相当但无查表依赖,纯整数运算适合帧级 hot path。
// 64 位宽:1080p60 一帧 1920*1080*0.125 字节 ≈ 260KB 采样数据,32 位易碰撞。
const FNV1A_OFFSET_BASIS: u64 = 0xcbf29ce484222325;
const FNV1A_PRIME: u64 = 0x100000001b3;
const Y_PLANE_HASH_ROW_STEP: usize = 8;
// WebRTC 模式下的编码器 time_base 分母(90kHz)。匹配 RTP 视频时钟频率(RFC 3551)。
// 这样编码器输出的 PTS 直接就是 RTP 时间戳,无需二次换算(microsecond 精度足够)。
// MP4 模式保留 1/fps time_base 以简化文件写入。
/// 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;
// 计算 NV12 Y 平面的 FNV-1a 采样哈希,用于检测连续帧是否像素相同(静帧去重)。
// 步长 Y_PLANE_HASH_ROW_STEP=8:仅哈希每 8 行的全部像素,跳过中间 7 行。
// 这是性能优化:1080p60 全量哈希 ~10MB/帧 * 60fps = 600MB/s 内存带宽;
// 采样到 ~1.3MB/帧仍能可靠检测桌面静止(多数静止帧大面积像素相同)。
//
// 等价 Go
// hash := uint64(0xcbf29ce484222325)
// for row := 0; row < height; row += 8 {
// off := row * stride
// for _, b := range y_data[off:off+width] {
// hash ^= uint64(b); hash *= 0x100000001b3
// }
// }
// wrapping_mul:故意允许溢出回绕(u64 取模 2^64),等价于 Go 无符号溢出语义。
fn hash_sampled_y_plane(y_data: &[u8], width: usize, height: usize, stride: usize) -> u64 {
let mut hash = FNV1A_OFFSET_BASIS;
for row in (0..height).step_by(Y_PLANE_HASH_ROW_STEP) {
let row_start = row * stride;
let row_end = row_start + width;
for &byte in &y_data[row_start..row_end] {
hash ^= u64::from(byte);
hash = hash.wrapping_mul(FNV1A_PRIME);
}
}
hash
}
// unsafe impl Send 中文概述:SwEncEncode 含裸指针(sws_ctx / yuv_frame / enc_video),
// Rust 默认认为裸指针非 Send(防止跨线程共享 C 资源)。这里手写 unsafe impl Send 表示:
// 本类型在构造完成后只被 move 到一个编码线程,所有访问通过 &mut self 独占借用,
// 满足"单线程独占"模型。AGENTS.md 明确警告:跨线程移动这些 wrapper 前必须重审 exclusivity 假设。
// 英文 SAFETY 详细论证下方保留不动。
// 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 SwEncEncode6 个方法按数据流顺序
// (1) new_muxer:构造函数,输出到 MP4 文件(new_with_resolution_control 调用)
// (2) new_webrtc:构造函数,输出到 crossbeam Sender<EncodedH264Frame>
// (3) flushEOS 时排空编码器(send null frame + drain
// (4) take_timing:取出最近一帧的 sws/encode 耗时统计
// (5) encode_cpu_frame:单帧主循环(NV12 → YUV420P → libx264 → drain → channel/muxer
// (6) recreate_encoder / write_trailer_if_needed / drain_encoder:私有辅助
impl SwEncEncode {
#[allow(clippy::too_many_arguments)]
fn new_muxer(
output_path: &Path,
enc_width: u32,
enc_height: u32,
fps: u32,
bitrate: u64,
gop_size: u32,
) -> Result<Self> {
// MP4 muxer 模式构造函数:用 avformat 写文件,不接 WebRTC。
// 调用专用工厂函数(create_software_h264_muxer 内部完成 avformat_alloc_output +
// libx264 encoder + avio_open + write_header)。
let sws_ctx = create_nv12_to_yuv420p_sws(enc_width, enc_height)?;
let (enc_video, octx) =
create_software_h264_muxer(output_path, enc_width, enc_height, fps, bitrate, gop_size)?;
let yuv_frame = alloc_yuv420p_frame(enc_width, enc_height)?;
// 创建后立即 drop 发送端的 dummy channel idiom:构造一个容量 1 的 channel
// 拿到接收端后立刻丢弃发送端(drop(dummy_tx)),使接收端永远 recv 到 Disconnected。
// encode_cpu_frame 中的 try_recv 循环因此立即结束,等效"无分辨率控制"。
// 比用 Option<Receiver> 更优雅:避免每个调用点都 match Option。
let (dummy_tx, bitrate_rx) = crossbeam_channel::bounded(1);
drop(dummy_tx);
let (dummy_resolution_tx, resolution_rx) = crossbeam_channel::bounded(1);
drop(dummy_resolution_tx);
Ok(Self {
sws_ctx,
enc_video,
output: Some(FrameOutput::Muxer(octx)),
yuv_frame,
last_frame_hash: 0,
frame_count: 0,
starting_timestamp: None,
frames_written: false,
webrtc_disconnected: false,
webrtc_paused: None,
bitrate_rx,
resolution_rx,
enc_width,
enc_height,
fps,
bitrate,
gop_size,
force_keyframe_pending: false,
last_timing: SwEncodeTiming::default(),
last_capture_time: None,
})
}
#[allow(clippy::too_many_arguments)]
pub fn new_webrtc(
enc_width: u32,
enc_height: u32,
fps: u32,
bitrate: u64,
gop_size: u32,
tx: crossbeam_channel::Sender<EncodedH264Frame>,
webrtc_paused: Arc<AtomicBool>,
bitrate_rx: crossbeam_channel::Receiver<BitrateCommand>,
resolution_rx: crossbeam_channel::Receiver<ResolutionChange>,
) -> Result<Self> {
let sws_ctx = create_nv12_to_yuv420p_sws(enc_width, enc_height)?;
let enc_video =
create_software_h264_encoder(enc_width, enc_height, fps, bitrate, gop_size)?;
let yuv_frame = alloc_yuv420p_frame(enc_width, enc_height)?;
Ok(Self {
sws_ctx,
enc_video,
output: Some(FrameOutput::Channel(tx)),
yuv_frame,
last_frame_hash: 0,
frame_count: 0,
starting_timestamp: None,
frames_written: false,
webrtc_disconnected: false,
webrtc_paused: Some(webrtc_paused),
bitrate_rx,
resolution_rx,
enc_width,
enc_height,
fps,
bitrate,
gop_size,
force_keyframe_pending: false,
last_timing: SwEncodeTiming::default(),
last_capture_time: None,
})
}
pub fn flush(&mut self) -> Result<()> {
// 中文概述:flush 阶段,向打开的 libx264 编码器发送 NULL frame 触发 drain 模式。
// FFmpeg 契约:send_frame(NULL) 表示输入 EOS,后续 receive_packet 直到返回 AVERROR_EOF。
// AVERROR_EOF 在内部被忽略(已 drain 完毕),其他错误传播。
// 英文 SAFETY 见下方。
// 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> {
// 每帧开始:把计时统计置零,避免上一帧的 stale 值泄漏(last_timing 字段注释见 struct)。
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);
// WebRTC 通道已断开(之前 drain_encoder 检测到 Disconnected)→ 后续帧直接跳过。
if self.webrtc_disconnected {
return Ok(EncodeOutcome::SkippedDisconnected);
}
// 必须先 drain bitrate 命令再做 stride 检查:导入线程在 emit ResolutionChange 后
// 才送来新的(更小 stride 的)帧;若先检查 stride 会因旧 stride 与新帧不匹配误报。
// 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;
// 中文概述:直接写 libx264 AVCodecContext 的 bit_rate 字段(i64),
// 实时调整编码器输出码率(BWE 反馈驱动)。英文 SAFETY 见下方。
// 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");
}
}
}
// 暂存 keyframe 请求状态:实际使用在本函数后段(pict_type 设置处),中间可能被覆盖。
let force_this_frame = self.force_keyframe_pending;
// drain ResolutionChange 通道:若分辨率变化,调用 recreate_encoder 重建 libx264 +
// sws_ctx + yuv_frame(避免在帧处理中途重建,确保后续 unsafe 块看到的字段一致)。
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");
}
// 暂停检查:webrtc_paused 是 Arc<AtomicBool>(跨线程共享),Relaxed 内存序足够
// (单字段读,不参与 happens-before 协调;类比 Go atomic.LoadInt32)。
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;
// required_y_lenY 平面所需最小字节数 = (height-1) 行完整 stride + 末行 width 字节。
// 用于校验源帧尺寸合法(防御式编程,避免 from_raw_parts 越界)。
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);
// 帧去重:FNV-1a 哈希 Y 平面采样点,与上一帧相同则丢帧(静帧优化)。
let current_hash = hash_sampled_y_plane(&frame.y_data, width, height, frame.y_stride);
// GOP 强制关键帧:每 gop_size 帧一个 I-framelibx264 也可由 keyint 控制,但这里
// 用 force_gop_frame 显式触发)。
let force_gop_frame = self.gop_size > 0 && frame_index % u64::from(self.gop_size) == 0;
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();
// 中文概述:sws_scale 把 NV12(含交错的 UV 半高平面)转 YUV420PY + U + V 三平面)。
// yuv_frame 是预分配的可复用 AVFrameav_frame_make_writable 确保独占写权限。
// src_slices/src_strides 是 4 元素数组(FFmpeg 固定 API),后两个填 null/0。
// 英文 SAFETY 见下方。
// 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();
// 中文概述:把 YUV420P 帧送入 libx264。pict_type 重置为 NONE(除非 force_this_frame
// 强制 I-frame)—— yuv_frame 复用,若不重置上一帧的 I-type 会泄漏到当前 P-frame。
// forced-idr=1 编码器选项让 AV_PICTURE_TYPE_I 触发真正的 IDR NALU(不是普通 I-frame)。
// 英文 SAFETY 详细论证下方。
// 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;
}
// drain_encoder 内部循环 avcodec_receive_packet,把所有编码好的 packet 写 muxer 或 channel。
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)
}
fn recreate_encoder(&mut self, width: u32, height: u32) -> Result<()> {
if width == self.enc_width && height == self.enc_height {
return Ok(());
}
tracing::info!(
from = format_args!("{}x{}", self.enc_width, self.enc_height),
to = format_args!("{}x{}", width, height),
"recreating WebRTC software encoder for resolution change"
);
if !self.sws_ctx.is_null() {
// 中文概述:分辨率变化前先释放旧 sws_ctx(NV12→YUV420P 转换器),
// 紧接着把字段置 null 防止二次释放。英文 SAFETY 见下方。
// SAFETY: sws_ctx is owned exclusively by self and will be replaced below.
unsafe { ffi::sws_freeContext(self.sws_ctx) };
self.sws_ctx = ptr::null_mut();
}
if !self.yuv_frame.is_null() {
// 中文概述:同步释放旧 yuv_frameav_frame_free 内部清空 data/extended_data。
// SAFETY: yuv_frame is owned exclusively by self and will be replaced below.
unsafe { ffi::av_frame_free(&mut self.yuv_frame) };
}
self.sws_ctx = create_nv12_to_yuv420p_sws(width, height)?;
self.enc_video =
create_software_h264_encoder(width, height, self.fps, self.bitrate, self.gop_size)?;
self.yuv_frame = alloc_yuv420p_frame(width, height)?;
self.enc_width = width;
self.enc_height = height;
self.last_frame_hash = 0;
self.frame_count = 0;
self.force_keyframe_pending = true;
Ok(())
}
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;
// drain 循环:avcodec_send_frame 之后,必须反复调用 avcodec_receive_packet 直到 EAGAIN/EOF。
// FFmpeg 编码 API 是异步的:send_frame 立即返回,packet 通过 receive_packet 取出。
// 类比 Go 的 chan: send_frame=send, receive_packet=recv, EAGAIN=chan 空。
loop {
let mut pkt = ff::Packet::empty();
// 中文概述:从打开的 libx264 编码器取出一个已编码 packet。
// 英文 SAFETY 见下方。
// 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.
// 中文概述:从 AVPacket.size 读字节数累加,避免在 Muxer/Channel 分支里重复统计。
// 英文 SAFETY 见下方。
// 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;
}
// 输出目的地二选一:Muxer 写文件,Channel 送 WebRTC 线程。
match self.output {
Some(FrameOutput::Muxer(ref mut octx)) => {
let enc_tb = self.enc_video.time_base();
// 中文概述:从 AVFormatContext 读 stream 0 的 time_basemuxer 写入时用)。
// 必须校验 nb_streams > 0 && streams 非空,防止空 muxer 触发 UB。
// 英文 SAFETY 见下方。
// SAFETY: muxer output was created with stream 0 during setup;
// streams is non-null and stream 0 remains owned by the format context.
let stream_tb = unsafe {
let fmt = *octx.as_ptr();
if fmt.nb_streams == 0 || fmt.streams.is_null() {
bail!("no streams in output context");
}
let st = *fmt.streams.add(0);
ff::Rational::from((*st).time_base)
};
// rescale_ts:编码器 time_base → 容器 time_base 的有理数换算。
pkt.rescale_ts(enc_tb, stream_tb);
// 减去起始 PTS 让首帧 PTS=0Muxer 模式下 AVFormatContext 自己管 DTS。
if let Some(pts) = pkt.pts() {
pkt.set_pts(Some(pts - start_ts));
}
if let Some(dts) = pkt.dts() {
pkt.set_dts(Some(dts - start_ts));
}
pkt.set_stream(0);
pkt.write_interleaved(octx)
.map_err(|e| anyhow::anyhow!("Failed to write packet: {e}"))?;
self.frames_written = true;
}
Some(FrameOutput::Channel(ref tx)) => {
// 中文概述:把 AVPacket 字段读出到栈上 raw 结构体,用于检查 size/data。
// 英文 SAFETY 见下方。
// SAFETY: pkt is a valid AVPacket just filled by
// avcodec_receive_packet; this copies fields for
// read-only inspection before pkt is dropped.
let raw = unsafe { *pkt.as_mut_ptr() };
if raw.size > 0 && !raw.data.is_null() {
// 中文概述:从 raw.data 裸指针构造临时切片(零拷贝),随即 .to_vec() 深拷贝。
// 类比 Go 的 unsafe Slice((*byte)(data), size) → append(nil, slice...)。
// SAFETY: `pkt` is a valid AVPacket just filled by a successful
// `avcodec_receive_packet` call. We checked `size > 0` and
// `data` is non-null, so `data` points to `size` initialized
// bytes owned by the packet. `u8` has alignment 1, and the
// slice is copied into a Vec before the packet is unreffed.
let data: &[u8] =
unsafe { std::slice::from_raw_parts(raw.data, raw.size as usize) };
// Normalize PTS: subtract starting_timestamp so first frame = 0.
// Mirrors the Muxer branch normalization above. `start_ts` is
// self.starting_timestamp.unwrap_or(0) (passed by caller), so
// when no origin is recorded yet the subtraction is a no-op.
let pts_ticks = match pkt.pts() {
Some(p) => p - start_ts,
None => {
// libx264 should always set PTS; emitting RTP ts=0
// here would recreate issue #24. Drop the packet.
tracing::warn!(
"encoder produced packet without PTS, dropping"
);
continue;
}
};
// try_send 非阻塞:通道满或断开都立即返回(不阻塞 encode 线程)。
// 类比 Go 的 `select { case ch <- frame: default: drop }`。
match tx.try_send(EncodedH264Frame {
data: data.to_vec(),
pts_ticks,
capture_time: self
.last_capture_time
.unwrap_or_else(Instant::now),
}) {
Ok(()) => {}
Err(crossbeam_channel::TrySendError::Full(frame)) => {
tracing::warn!(
"WebRTC channel full, dropping frame: {} bytes lost",
frame.data.len()
);
}
Err(crossbeam_channel::TrySendError::Disconnected(frame)) => {
tracing::warn!(
"WebRTC channel disconnected: {} bytes lost",
frame.data.len()
);
// 设置断开标志,后续 encode_cpu_frame 调用立即 SkippedDisconnected
// 避免每帧都重复走 sws_scale + encode + drain 的浪费路径。
self.webrtc_disconnected = true;
break;
}
}
}
}
None => {}
}
}
Ok(total_bytes)
}
}
// impl Drop for SwEncEncodeC 资源释放。enc_video 由 ff::codec::encoder::video::Video
// 包装,Drop 会自动调用 avcodec_closesws_ctx 和 yuv_frame 是裸指针,必须在此手动释放。
// 类比 Go 的 `defer ffi.sws_freeContext(sws_ctx)`——Rust 没有原生 defer,但 Drop trait 等价。
impl Drop for SwEncEncode {
fn drop(&mut self) {
if !self.sws_ctx.is_null() {
// 中文概述:释放 sws_ctx。释放后置 null 防止二次释放(虽然 drop 后 self 立即销毁,
// 但这是防御式编程惯例,与 recreate_encoder 中的释放路径保持一致)。
// 英文 SAFETY 见下方。
// 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() {
// 中文概述:释放 yuv_frame。av_frame_free 接收 **mut 并内部置 null,无需手动置空。
// 英文 SAFETY 见下方。
// SAFETY: yuv_frame is owned by this state and was allocated by av_frame_alloc.
unsafe { ffi::av_frame_free(&mut self.yuv_frame) };
}
}
}
pub struct SwEncState {
import: SwEncImport,
encode: SwEncEncode,
}
// SAFETY: SwEncState owns import and encode state exclusively and existing sync callers move it
// between threads only with external serialization; all FFI handles are accessed through &mut self.
unsafe impl Send for SwEncState {}
impl SwEncState {
#[allow(clippy::too_many_arguments)]
pub fn new(
drm_device: &Path,
output_path: &Path,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
fps: u32,
bitrate: u64,
gop_size: u32,
) -> Result<Self> {
tracing::info!(
"SwEncState::new: GPU downscale {width}x{height} BGRA -> {enc_width}x{enc_height} NV12, software H.264"
);
let import = SwEncImport::new(drm_device, width, height, enc_width, enc_height, fps)?;
let encode =
SwEncEncode::new_muxer(output_path, enc_width, enc_height, fps, bitrate, gop_size)?;
Ok(Self { import, encode })
}
#[allow(clippy::too_many_arguments)]
pub fn new_webrtc(
drm_device: &Path,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
fps: u32,
bitrate: u64,
gop_size: u32,
tx: crossbeam_channel::Sender<EncodedH264Frame>,
webrtc_paused: Arc<AtomicBool>,
) -> Result<Self> {
tracing::info!(
"SwEncState::new_webrtc: GPU downscale {width}x{height} BGRA -> {enc_width}x{enc_height} NV12, software H.264 -> WebRTC"
);
let import = SwEncImport::new(drm_device, width, height, enc_width, enc_height, fps)?;
let (dummy_tx, bitrate_rx) = crossbeam_channel::bounded(1);
drop(dummy_tx);
let (dummy_resolution_tx, resolution_rx) = crossbeam_channel::bounded(1);
drop(dummy_resolution_tx);
let encode = SwEncEncode::new_webrtc(
enc_width,
enc_height,
fps,
bitrate,
gop_size,
tx,
webrtc_paused,
bitrate_rx,
resolution_rx,
)?;
Ok(Self { import, encode })
}
pub fn frames_rgb(&self) -> &AvHwFrameCtx {
self.import.frames_rgb()
}
pub fn encode_frame(&mut self, hw_frame: &ff::frame::Video) -> Result<()> {
let cpu_frame = self.import.import_and_scale(hw_frame)?;
self.encode.encode_cpu_frame(&cpu_frame).map(|_| ())
}
pub fn flush(&mut self) -> Result<()> {
for frame in self.import.flush_import()? {
self.encode.encode_cpu_frame(&frame)?;
}
self.encode.flush()?;
self.encode.write_trailer_if_needed()
}
}
// ---------------------------------------------------------------------------
// Shared encoder creation (used by both wlr-screencopy and portal paths)
// ---------------------------------------------------------------------------
/// Create a fully configured encoder with VAAPI hardware acceleration.
///
/// Convenience wrapper around [`EncState::new`] that computes default values
/// for `bitrate` and `gop_size` when not provided, and handles encoder dimension
/// transposition for rotated/transformed outputs.
#[allow(clippy::too_many_arguments)]
pub fn create_encoder(
drm_device: &Path,
output_path: &Path,
width: u32,
height: u32,
fps: u32,
transform: Transform,
bitrate: Option<u64>,
gop_size: Option<u32>,
existing_hw_ctx: Option<AvHwDevCtx>,
) -> Result<EncState> {
let (enc_w, enc_h) = transpose_if_transform_transposed(transform, width as i32, height as i32);
let actual_bitrate =
bitrate.unwrap_or_else(|| 2 * (width as u64) * (height as u64) * (fps as u64) / 100);
let actual_gop_size = gop_size.unwrap_or(fps);
EncState::new(
drm_device,
output_path,
width,
height,
enc_w as u32,
enc_h as u32,
actual_bitrate,
actual_gop_size,
fps,
transform,
existing_hw_ctx,
)
}
// ---------------------------------------------------------------------------
// Software-encode GPU-downscale helpers
// ---------------------------------------------------------------------------
#[allow(clippy::too_many_arguments)]
fn build_swenc_filter_graph(
hw_dev: &AvHwDevCtx,
frames_rgb: &AvHwFrameCtx,
width: u32,
height: u32,
enc_width: u32,
enc_height: u32,
fps: u32,
) -> Result<ff::filter::Graph> {
let mut graph = ff::filter::Graph::new();
let buffersrc =
ff::filter::find("buffer").ok_or_else(|| anyhow::anyhow!("filter 'buffer' not found"))?;
let buffersink = ff::filter::find("buffersink")
.ok_or_else(|| anyhow::anyhow!("filter 'buffersink' not found"))?;
let scale_vaapi = ff::filter::find("scale_vaapi")
.ok_or_else(|| anyhow::anyhow!("filter 'scale_vaapi' not found"))?;
// FFmpeg 8.0+ rejects VAAPI pix_fmt in buffer args before hw_frames_ctx is attached.
// Use a SW placeholder, then override format/hw_frames_ctx with av_buffersrc_parameters_set.
let args = format!(
"video_size={}x{}:pix_fmt=bgra:time_base=1/{fps}:pixel_aspect=1/1",
width, height,
);
let mut src_ctx = graph.add(&buffersrc, "in", &args)?;
// SAFETY: av_buffersrc_parameters_alloc returns newly allocated parameters
// or null, which is checked below.
let par = unsafe { ffi::av_buffersrc_parameters_alloc() };
if par.is_null() {
bail!("av_buffersrc_parameters_alloc returned null");
}
// SAFETY: par and src_ctx are valid; frames_rgb.ref_clone returns an owned hw_frames_ctx ref
// that buffersrc consumes on successful parameter set.
unsafe {
(*par).format = Into::<ffi::AVPixelFormat>::into(ff::format::Pixel::VAAPI) as i32;
(*par).width = width as i32;
(*par).height = height as i32;
(*par).time_base = ffi::AVRational {
num: 1,
den: fps as i32,
};
(*par).hw_frames_ctx = frames_rgb.ref_clone();
let ret = ffi::av_buffersrc_parameters_set(src_ctx.as_mut_ptr(), par);
ffi::av_free(par as *mut _);
if ret < 0 {
bail!("av_buffersrc_parameters_set failed: {}", ff_err(ret));
}
}
let mut scale_ctx = graph.add(
&scale_vaapi,
"scale",
&format!("{enc_width}:{enc_height}:format=nv12"),
)?;
// SAFETY: scale_vaapi keeps a ref-counted device context while the graph is alive.
unsafe {
(*scale_ctx.as_mut_ptr()).hw_device_ctx = hw_dev.ref_clone();
}
let mut sink_ctx = graph.add(&buffersink, "out", "")?;
src_ctx.link(0, &mut scale_ctx, 0);
scale_ctx.link(0, &mut sink_ctx, 0);
graph
.validate()
.map_err(|e| anyhow::anyhow!("software GPU filter graph validation failed: {e}"))?;
Ok(graph)
}
fn create_nv12_to_yuv420p_sws(width: u32, height: u32) -> Result<*mut ffi::SwsContext> {
// SAFETY: sws_getContext creates an owned scaler context for same-size NV12 -> YUV420P.
let ctx = unsafe {
ffi::sws_getContext(
width as i32,
height as i32,
ffi::AVPixelFormat::AV_PIX_FMT_NV12,
width as i32,
height as i32,
ffi::AVPixelFormat::AV_PIX_FMT_YUV420P,
2,
ptr::null_mut(),
ptr::null_mut(),
ptr::null_mut(),
)
};
if ctx.is_null() {
bail!("Failed to create NV12 -> YUV420P sws_scale context");
}
Ok(ctx)
}
fn alloc_yuv420p_frame(width: u32, height: u32) -> Result<*mut ffi::AVFrame> {
// SAFETY: Allocate an AVFrame, configure format/dimensions, then allocate writable buffers.
unsafe {
let mut frame = ffi::av_frame_alloc();
if frame.is_null() {
bail!("av_frame_alloc failed");
}
(*frame).width = width as i32;
(*frame).height = height as i32;
(*frame).format = ffi::AVPixelFormat::AV_PIX_FMT_YUV420P as i32;
let ret = ffi::av_frame_get_buffer(frame, 0);
if ret < 0 {
ffi::av_frame_free(&mut frame);
bail!("av_frame_get_buffer failed: {}", ff_err(ret));
}
Ok(frame)
}
}
fn create_software_h264_muxer(
output_path: &Path,
width: u32,
height: u32,
fps: u32,
bitrate: u64,
gop_size: u32,
) -> Result<(
ff::codec::encoder::video::Video,
ff::format::context::Output,
)> {
let output_cstr = CString::new(output_path.to_str().unwrap())?;
let codec = ff::encoder::find_by_name("libx264")
.or_else(|| ff::encoder::find_by_name("libopenh264"))
.ok_or_else(|| {
anyhow::anyhow!("No H.264 software encoder found (tried libx264, libopenh264)")
})?;
let codec_name = codec.name().to_string();
let mut enc = {
let ctx = ff::codec::Context::new_with_codec(codec);
ctx.encoder().video()?
};
enc.set_width(width);
enc.set_height(height);
enc.set_format(ff::format::Pixel::YUV420P);
enc.set_bit_rate(bitrate as usize);
enc.set_gop(gop_size);
enc.set_time_base(ff::Rational::new(1, fps as i32));
enc.set_max_b_frames(3);
// SAFETY: global headers are needed by MP4 and harmless for other common muxers.
unsafe {
(*enc.as_mut_ptr()).flags |= ffi::AV_CODEC_FLAG_GLOBAL_HEADER as i32;
}
if codec_name == "libx264" {
// SAFETY: priv_data and codec context belong to the unopened encoder;
// strings live for each av_opt_set call.
unsafe {
let key = CString::new("preset").unwrap();
let val = CString::new("fast").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
let key = CString::new("threads").unwrap();
let val = CString::new("6").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
(*enc.as_mut_ptr()).profile = ffi::AV_PROFILE_H264_HIGH as i32;
// SAFETY: enc is a valid, initialized AVCodecContext from
// avcodec_alloc_context3. Setting level is a simple i32 field
// assignment on a properly aligned struct.
(*enc.as_mut_ptr()).level = 40; // H.264 Level 4.0 (up to 1080p@30)
}
}
let opened = enc
.open()
.map_err(|e| anyhow::anyhow!("Failed to open {codec_name} encoder: {e}"))?;
let enc_video = opened.0;
let use_null = output_path
.to_str()
.map(|s| s.contains("null"))
.unwrap_or(false);
let fmt_name = if use_null {
CString::new("null").unwrap()
} else {
CString::new("").unwrap()
};
let fmt_name_ptr = if use_null {
fmt_name.as_ptr()
} else {
ptr::null()
};
let mut fmt_ctx_ptr: *mut ffi::AVFormatContext = ptr::null_mut();
// SAFETY: fmt_ctx_ptr is initialized by FFmpeg; C strings live across the call.
let ret = unsafe {
ffi::avformat_alloc_output_context2(
&mut fmt_ctx_ptr,
ptr::null_mut(),
fmt_name_ptr,
output_cstr.as_ptr(),
)
};
if ret < 0 || fmt_ctx_ptr.is_null() {
bail!("Failed to allocate output format context: {}", ff_err(ret));
}
// SAFETY: fmt_ctx_ptr is valid; stream and codec parameters are owned by the format context.
let stream_ptr = unsafe { ffi::avformat_new_stream(fmt_ctx_ptr, ptr::null()) };
if stream_ptr.is_null() {
bail!("Failed to create output stream");
}
// SAFETY: stream_ptr and encoder context are valid; parameters are copied into stream.
let ret =
unsafe { ffi::avcodec_parameters_from_context((*stream_ptr).codecpar, enc_video.as_ptr()) };
if ret < 0 {
bail!("Failed to copy codec parameters to stream: {}", ff_err(ret));
}
// SAFETY: stream_ptr is valid and writable during muxer setup.
unsafe {
(*stream_ptr).time_base = (*enc_video.as_ptr()).time_base;
}
// SAFETY: open an AVIO only for muxers that require files; null muxer advertises NOFILE.
unsafe {
if (*(*fmt_ctx_ptr).oformat).flags & ffi::AVFMT_NOFILE == 0 {
let ret = ffi::avio_open(
&mut (*fmt_ctx_ptr).pb,
output_cstr.as_ptr(),
ffi::AVIO_FLAG_WRITE,
);
if ret < 0 {
bail!(
"Failed to open output file '{}': {}",
output_path.display(),
ff_err(ret)
);
}
}
}
// SAFETY: fmt_ctx_ptr is fully configured.
let ret = unsafe { ffi::avformat_write_header(fmt_ctx_ptr, ptr::null_mut()) };
if ret < 0 {
bail!("Failed to write output header: {}", ff_err(ret));
}
// SAFETY: ownership of fmt_ctx_ptr transfers to ffmpeg-next Output wrapper.
let octx = unsafe { ff::format::context::Output::wrap(fmt_ctx_ptr) };
tracing::info!("Using software H.264 encoder: {codec_name}");
Ok((enc_video, octx))
}
fn create_software_h264_encoder(
width: u32,
height: u32,
fps: u32,
bitrate: u64,
gop_size: u32,
) -> Result<ff::codec::encoder::video::Video> {
let codec = ff::encoder::find_by_name("libx264")
.or_else(|| ff::encoder::find_by_name("libopenh264"))
.ok_or_else(|| anyhow::anyhow!("No H.264 software encoder found"))?;
let codec_name = codec.name().to_string();
let mut enc = {
let ctx = ff::codec::Context::new_with_codec(codec);
ctx.encoder().video()?
};
enc.set_width(width);
enc.set_height(height);
enc.set_format(ff::format::Pixel::YUV420P);
enc.set_bit_rate(bitrate as usize);
enc.set_gop(gop_size);
// 90kHz media clock matches RTP directly. Eliminates 1/fps quantization
// that previously caused sequential RTP timestamps during 60fps capture,
// leading to 2x RTP time inflation and 10s+ browser jitter buffer growth.
// See issue #25.
enc.set_time_base(ff::Rational::new(1, 90_000));
// Explicit framerate is REQUIRED when time_base is not 1/fps, otherwise
// libx264 infers wrong fps from the 90kHz time_base and VBV rate control
// breaks. Per Oracle review round for #25.
enc.set_frame_rate(Some(ff::Rational::new(fps as i32, 1)));
enc.set_max_b_frames(0);
if codec_name == "libx264" {
// SAFETY: priv_data and codec context belong to the unopened encoder;
// each CString lives for the duration of its av_opt_set call.
unsafe {
let key = CString::new("preset").unwrap();
let val = CString::new("veryfast").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
let key = CString::new("tune").unwrap();
let val = CString::new("zerolatency").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
let key = CString::new("threads").unwrap();
let val = CString::new("6").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
// High profile via AVCodecContext.profile (not x264opts — x264 rejects it there).
// High enables CABAC + 8x8dct automatically.
(*enc.as_mut_ptr()).profile = ffi::AV_PROFILE_H264_HIGH as i32;
// SAFETY: enc is a valid, initialized AVCodecContext from
// avcodec_alloc_context3. Setting level is a simple i32 field
// assignment on a properly aligned struct.
(*enc.as_mut_ptr()).level = 42; // H.264 Level 4.2 (up to 1440p@30)
// SAFETY: priv_data belongs to the unopened libx264 encoder context.
// `forced-idr` is an FFmpeg-level private option (not x264-native),
// so it must be set via av_opt_set, NOT via the x264opts string.
// With forced-idr=1, setting AV_PICTURE_TYPE_I on an input frame
// produces a true IDR NALU with inline SPS/PPS (repeat_headers=1).
let key = CString::new("forced-idr").unwrap();
let val = CString::new("1").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
let key = CString::new("x264opts").unwrap();
// x264's vbv-maxrate unit is kbit/s and vbv-bufsize is kbit (NOT bps).
// Confirmed via x264 source encoder/ratecontrol.c:658-661 which multiplies
// these values by 1000 to convert kbit → bit at use site. Passing bps makes
// VBV effectively unbounded (5.5 Mbps becomes 5.5 Gbps, clipped to 2 Gbps).
// See https://github.com/mirror/x264/blob/c24e06c2e184345ceb33eb20a15d1024d9fd3497/encoder/ratecontrol.c#L658-L661
let vbv_maxrate_kbps = bitrate / 1000;
let vbv_bufsize_kbps = (bitrate / 4) / 1000;
let val = CString::new(format!(
"repeat_headers=1:vbv-maxrate={vbv_maxrate_kbps}:vbv-bufsize={vbv_bufsize_kbps}"
))
.unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
}
}
let opened = enc
.open()
.map_err(|e| anyhow::anyhow!("Failed to open {codec_name} encoder: {e}"))?;
tracing::info!("WebRTC encoder: {codec_name} {width}x{height} @ {fps}fps {bitrate}bps (profile High, preset veryfast)");
Ok(opened.0)
}
// ---------------------------------------------------------------------------
// Filter graph (inline)
// ---------------------------------------------------------------------------
fn build_filter_graph(
hw_dev: &AvHwDevCtx,
frames_rgb: &AvHwFrameCtx,
width: u32,
height: u32,
_enc_width: u32,
_enc_height: u32,
fps: u32,
transform: Transform,
) -> Result<ff::filter::Graph> {
let mut graph = ff::filter::Graph::new();
let buffersrc =
ff::filter::find("buffer").ok_or_else(|| anyhow::anyhow!("filter 'buffer' not found"))?;
let buffersink = ff::filter::find("buffersink")
.ok_or_else(|| anyhow::anyhow!("filter 'buffersink' not found"))?;
let scale_vaapi = ff::filter::find("scale_vaapi")
.ok_or_else(|| anyhow::anyhow!("filter 'scale_vaapi' not found"))?;
// buffersrc — use AVBufferSrcParameters to set hw_frames_ctx properly
let args = format!(
"video_size={}x{}:pix_fmt={}:time_base=1/{fps}:pixel_aspect=1/1",
width,
height,
Into::<ffi::AVPixelFormat>::into(ff::format::Pixel::VAAPI) as i32,
);
let mut src_ctx = graph.add(&buffersrc, "in", &args)?;
// SAFETY: av_buffersrc_parameters_alloc allocates params for the buffersrc.
let par = unsafe { ffi::av_buffersrc_parameters_alloc() };
if par.is_null() {
bail!("av_buffersrc_parameters_alloc returned null");
}
// SAFETY: Set hw_frames_ctx on the buffersrc parameters, then apply.
unsafe {
(*par).format = Into::<ffi::AVPixelFormat>::into(ff::format::Pixel::VAAPI) as i32;
(*par).width = width as i32;
(*par).height = height as i32;
(*par).time_base = ffi::AVRational {
num: 1,
den: fps as i32,
};
(*par).hw_frames_ctx = frames_rgb.ref_clone();
let ret = ffi::av_buffersrc_parameters_set(src_ctx.as_mut_ptr(), par);
ffi::av_free(par as *mut _);
if ret < 0 {
bail!("av_buffersrc_parameters_set failed: {}", ff_err(ret));
}
}
// scale_vaapi: hardware scaling and colourspace conversion (keeps original dimensions)
let mut scale_ctx = graph.add(
&scale_vaapi,
"scale",
&format!("{width}:{height}:format=nv12"),
)?;
// SAFETY: scale_vaapi needs hw_device_ctx for VAAPI device access.
unsafe {
(*scale_ctx.as_mut_ptr()).hw_device_ctx = hw_dev.ref_clone();
}
// buffersink
let mut sink_ctx = graph.add(&buffersink, "out", "")?;
// Build filter chain: src -> scale -> [transpose] -> sink
src_ctx.link(0, &mut scale_ctx, 0);
match transform {
Transform::Normal => {
scale_ctx.link(0, &mut sink_ctx, 0);
}
other => {
let transpose = ff::filter::find("transpose_vaapi")
.ok_or_else(|| anyhow::anyhow!("filter 'transpose_vaapi' not found"))?;
let dir_val = match other {
Transform::Normal90 => "1",
Transform::Normal180 => "4",
Transform::Normal270 => "2",
Transform::Flipped => "5",
Transform::Flipped90 => "3",
Transform::Flipped180 => "6",
Transform::Flipped270 => "0",
Transform::Normal => unreachable!(),
};
let mut trans_ctx = graph.add(&transpose, "transpose", &format!("dir={dir_val}"))?;
// SAFETY: trans_ctx is a live transpose_vaapi filter context;
// scale_vaapi/transpose_vaapi keep a ref-counted device context.
unsafe {
(*trans_ctx.as_mut_ptr()).hw_device_ctx = hw_dev.ref_clone();
}
scale_ctx.link(0, &mut trans_ctx, 0);
trans_ctx.link(0, &mut sink_ctx, 0);
}
}
graph
.validate()
.map_err(|e| anyhow::anyhow!("Filter graph validation failed: {e}"))?;
Ok(graph)
}
#[cfg(test)]
mod tests {
use super::*;
// ── Task 1: VBV x264opts formatting ──
#[test]
fn vbv_x264opts_format() {
let bitrate: u64 = 5_000_000;
// x264 expects kbit/s and kbit, not bps
let vbv_maxrate_kbps = bitrate / 1000;
let vbv_bufsize_kbps = (bitrate / 4) / 1000;
let opts = format!("repeat_headers=1:vbv-maxrate={vbv_maxrate_kbps}:vbv-bufsize={vbv_bufsize_kbps}");
assert_eq!(vbv_maxrate_kbps, 5000);
assert_eq!(vbv_bufsize_kbps, 1250);
assert!(opts.contains("vbv-maxrate=5000"));
assert!(opts.contains("vbv-bufsize=1250"));
}
#[test]
fn vbv_bufsize_is_quarter_of_maxrate() {
for bitrate in [1_000_000, 5_000_000, 10_000_000] {
// x264 expects kbit/s and kbit; both scaled by /1000, ratio preserved
let maxrate_kbps = bitrate / 1000;
let bufsize_kbps = (bitrate / 4) / 1000;
assert_eq!(bufsize_kbps * 4, maxrate_kbps, "bufsize should be maxrate/4");
}
}
// ── Task 3: GOP formula ──
#[test]
fn webrtc_gop_formula() {
assert_eq!((15u32 * 2).max(20), 30); // 15fps -> 30
assert_eq!((30u32 * 2).max(20), 60); // 30fps -> 60
assert_eq!((60u32 * 2).max(20), 120); // 60fps -> 120
assert_eq!((5u32 * 2).max(20), 20); // 5fps -> 20 (floor)
}
#[test]
fn h264_level_values() {
// Level 4.0 supports up to 1080p@30fps (used for file muxer)
assert_eq!(40i32, 40);
// Level 4.2 supports up to 1440p@30fps (used for WebRTC low-latency encoder)
assert_eq!(42i32, 42);
}
// ── Task 4: Duplicate frame hash detection ──
#[test]
fn hash_sampled_y_plane_first_frame_consistent() {
let width = 64;
let height = 64;
let stride = 64;
let y_data = vec![0u8; stride * height];
let hash1 = hash_sampled_y_plane(&y_data, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data, width, height, stride);
assert_eq!(hash1, hash2, "same input should produce same hash");
}
#[test]
fn hash_sampled_y_plane_detects_different_frames() {
let width = 64;
let height = 64;
let stride = 64;
let y_data1 = vec![0u8; stride * height];
let y_data2 = vec![128u8; stride * height];
let hash1 = hash_sampled_y_plane(&y_data1, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data2, width, height, stride);
assert_ne!(
hash1, hash2,
"different frame data should produce different hashes"
);
}
#[test]
fn hash_sampled_y_plane_samples_every_8th_row() {
// Changing a non-sampled row (e.g., row 1) should NOT change the hash
let width = 64;
let height = 64;
let stride = 64;
let y_data1 = vec![0u8; stride * height];
let mut y_data2 = vec![0u8; stride * height];
// Row 1 is NOT sampled (sampling is every 8th row: 0, 8, 16, ...)
y_data2[stride * 1..stride * 1 + width].fill(255);
let hash1 = hash_sampled_y_plane(&y_data1, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data2, width, height, stride);
assert_eq!(
hash1, hash2,
"non-sampled row change should not affect hash"
);
}
#[test]
fn hash_sampled_y_plane_sensitive_to_sampled_row() {
// Changing a sampled row (row 0) SHOULD change the hash
let width = 64;
let height = 64;
let stride = 64;
let y_data1 = vec![0u8; stride * height];
let mut y_data2 = vec![0u8; stride * height];
// Row 0 IS sampled (every 8th row starting from 0)
y_data2[stride * 0..stride * 0 + width].fill(255);
let hash1 = hash_sampled_y_plane(&y_data1, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data2, width, height, stride);
assert_ne!(
hash1, hash2,
"sampled row change should produce different hash"
);
}
#[test]
fn hash_sampled_y_plane_handles_stride_greater_than_width() {
// Stride can be larger than width due to alignment; unused padding should not affect hash
let width = 32;
let height = 16;
let stride = 64; // padded stride
let y_data1 = vec![0u8; stride * height];
let mut y_data2 = vec![0u8; stride * height];
// Fill the padding area (columns 32..63) of row 0 with garbage
y_data2[width..stride].fill(0xFF);
let hash1 = hash_sampled_y_plane(&y_data1, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data2, width, height, stride);
assert_eq!(
hash1, hash2,
"padding bytes beyond width should not affect hash"
);
}
}