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; use anyhow::{bail, Result}; use ffmpeg_next as ff; use ffmpeg_next::ffi; use ffmpeg_next::packet::Mut as _; use crate::cap_portal::PwDmaBufFrame; use crate::transform::{transpose_if_transform_transposed, Transform}; // --------------------------------------------------------------------------- // Bitrate feedback command (WebRTC BWE → SW encoder) // --------------------------------------------------------------------------- /// 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, } #[derive(Clone, Copy, Debug)] pub struct ResolutionChange { pub width: u32, pub height: u32, } /// 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, } /// 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 // --------------------------------------------------------------------------- pub struct AvHwDevCtx { ptr: *mut ffi::AVBufferRef, } // 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 { pub fn new_vaapi(drm_device: &Path) -> Result { let device_cstr = CString::new(drm_device.to_str().unwrap())?; let mut p: *mut ffi::AVBufferRef = ptr::null_mut(); // 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 }) } pub fn as_ptr(&self) -> *mut ffi::AVBufferRef { self.ptr } pub fn ref_clone(&self) -> *mut ffi::AVBufferRef { // SAFETY: av_buffer_ref atomically increments refcount and returns a new ref. unsafe { ffi::av_buffer_ref(self.ptr) } } } impl Drop for AvHwDevCtx { fn drop(&mut self) { if !self.ptr.is_null() { // SAFETY: av_buffer_unref decrements refcount; frees the buffer when it hits zero. unsafe { ffi::av_buffer_unref(&mut self.ptr) }; } } } // --------------------------------------------------------------------------- // AvHwFrameCtx // --------------------------------------------------------------------------- pub struct AvHwFrameCtx { ptr: *mut ffi::AVBufferRef, } // 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 { fn new_inner(hw_dev: &AvHwDevCtx, w: u32, h: u32, sw_fmt: ff::format::Pixel) -> Result { // 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"); } // 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; } // 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 { // 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 }) } pub fn for_capture( hw_dev: &AvHwDevCtx, w: u32, h: u32, sw_fmt: ff::format::Pixel, ) -> Result { 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 { // SAFETY: av_buffer_ref atomically increments refcount and returns a new ref. unsafe { ffi::av_buffer_ref(self.ptr) } } } impl Drop for AvHwFrameCtx { fn drop(&mut self) { if !self.ptr.is_null() { // SAFETY: av_buffer_unref decrements refcount; frees when zero. unsafe { ffi::av_buffer_unref(&mut self.ptr) }; } } } /// 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)?; // 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(()) } /// 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 { let duped_fd = libc::dup(raw_fd); if duped_fd < 0 { bail!("dup(fd) failed: {}", std::io::Error::last_os_error()); } 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; let desc_box = Box::new(desc); let desc_ptr = Box::into_raw(desc_box); let buf_ref = ffi::av_buffer_create( desc_ptr as *mut u8, std::mem::size_of::(), 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"); } 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; } let mut dst = ff::frame::Video::empty(); // 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"); } } // 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) } 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); } /// Convert an FFmpeg error code to a human-readable string. pub(crate) fn av_err_to_string(err: i32) -> String { let mut buf = vec![0u8; 128]; // 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()); } String::from_utf8_lossy(&buf) .trim_end_matches('\0') .to_string() } /// 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 // --------------------------------------------------------------------------- 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, frames_written: bool, } unsafe impl Send for EncState {} impl EncState { #[allow(clippy::too_many_arguments)] 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, ) -> Result { tracing::info!( "EncState::new: {width}x{height} enc={enc_width}x{enc_height} transform={transform:?}" ); // 1. VAAPI device — reuse existing context if provided 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(|| anyhow::anyhow!("filter 'out' not found"))?; // 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 }; // 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 let codec = ff::encoder::find_by_name("h264_vaapi") .ok_or_else(|| anyhow::anyhow!("h264_vaapi encoder not found"))?; 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. 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; } // 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; } // 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. 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) let output_cstr = CString::new(output_path.to_str().unwrap())?; let mut fmt_ctx_ptr: *mut ffi::AVFormatContext = ptr::null_mut(); // 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)); } // 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"); } // 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"); } // 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) ); } // SAFETY: Copy encoder time_base to stream. unsafe { (*stream_ptr).time_base = (*enc_video.as_ptr()).time_base; } // 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) ); } // 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)); } // SAFETY: We created fmt_ctx_ptr above and it's valid. let octx = unsafe { ff::format::context::Output::wrap(fmt_ctx_ptr) }; Ok(Self { enc_video, frames_rgb, video_filter, hw_device_ctx, octx, starting_timestamp: None, frames_written: false, }) } pub fn frames_rgb(&self) -> &AvHwFrameCtx { &self.frames_rgb } 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}"))?; loop { let mut filtered = ff::frame::Video::empty(); match filter_sink.frame(&mut filtered) { Ok(()) => { 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(); // 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)); } self.drain_encoder(start_ts)?; } Ok(()) } 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) = 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"))?; let mut filter_sink = filter_sink_ctx.sink(); loop { let mut filtered = ff::frame::Video::empty(); match filter_sink.frame(&mut filtered) { Ok(()) => { let start_ts = self.starting_timestamp.unwrap_or(0); // 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, } } // 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); self.drain_encoder(start_ts)?; // Write trailer only if at least one frame was encoded. if self.frames_written { self.octx .write_trailer() .map_err(|e| anyhow::anyhow!("Failed to write trailer: {e}"))?; } Ok(()) } fn drain_encoder(&mut self, start_ts: i64) -> Result<()> { loop { let mut pkt = ff::Packet::empty(); // 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 == 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 let enc_tb = self.enc_video.time_base(); // 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"); } 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 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(&mut self.octx) .map_err(|e| anyhow::anyhow!("Failed to write packet: {e}"))?; self.frames_written = true; } Ok(()) } } // --------------------------------------------------------------------------- // SwEncState - VAAPI GPU downscale + software H.264 encode // --------------------------------------------------------------------------- pub enum FrameOutput { Muxer(ff::format::context::Output), Channel(crossbeam_channel::Sender>), } /// 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, pub uv_data: Vec, pub y_stride: usize, pub uv_stride: usize, pub pts: i64, } 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>, encoder_resolution_tx: Option>, } 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 { let hw_dev = AvHwDevCtx::new_vaapi(drm_device)?; let frames_rgb = AvHwFrameCtx::for_capture(&hw_dev, width, height, ff::format::Pixel::BGRA)?; let filter_graph = build_swenc_filter_graph( &hw_dev, &frames_rgb, width, height, enc_width, enc_height, fps, )?; 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, encoder_resolution_tx: crossbeam_channel::Sender, ) -> Result { 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 { self.poll_resolution_commands()?; 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(); filter_src .add(hw_frame) .map_err(|e| anyhow::anyhow!("software pipeline filter source add failed: {e}"))?; let mut first = None; let mut extra_count = 0usize; loop { let mut filtered = ff::frame::Video::empty(); match filter_sink.frame(&mut filtered) { Ok(()) => { 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; } } 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.ok_or_else(|| anyhow::anyhow!("software pipeline produced no scaled frame")) } pub fn flush_import(&mut self) -> Result> { 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 { // 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"); } // 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 { // 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) ); } // 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); 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, } }; Ok(frame) } } pub struct SwEncEncode { sws_ctx: *mut ffi::SwsContext, enc_video: ff::codec::encoder::video::Video, output: Option, yuv_frame: *mut ffi::AVFrame, last_frame_hash: u64, frame_count: u64, starting_timestamp: Option, frames_written: bool, webrtc_disconnected: bool, webrtc_paused: Option>, bitrate_rx: crossbeam_channel::Receiver, resolution_rx: crossbeam_channel::Receiver, 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, } const FNV1A_OFFSET_BASIS: u64 = 0xcbf29ce484222325; const FNV1A_PRIME: u64 = 0x100000001b3; const Y_PLANE_HASH_ROW_STEP: usize = 8; 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 } // 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 { #[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 { 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)?; 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(), }) } #[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>, webrtc_paused: Arc, bitrate_rx: crossbeam_channel::Receiver, resolution_rx: crossbeam_channel::Receiver, ) -> Result { 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(), }) } 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 { self.last_timing = SwEncodeTiming::default(); 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 } => { 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 % 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(); // 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) } 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() { // 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() { // 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 { 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)) => { let enc_tb = self.enc_video.time_base(); // 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) }; pkt.rescale_ts(enc_tb, stream_tb); 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)) => { // 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() { // 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) }; match tx.try_send(data.to_vec()) { Ok(()) => {} Err(crossbeam_channel::TrySendError::Full(frame)) => { tracing::warn!( "WebRTC channel full, dropping frame: {} bytes lost", frame.len() ); } Err(crossbeam_channel::TrySendError::Disconnected(frame)) => { tracing::warn!( "WebRTC channel disconnected: {} bytes lost", frame.len() ); self.webrtc_disconnected = true; break; } } } } 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) }; } } } 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 { 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>, webrtc_paused: Arc, ) -> Result { 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, gop_size: Option, existing_hw_ctx: Option, ) -> Result { 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 { 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::::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 { 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); enc.set_time_base(ff::Rational::new(1, fps as i32)); 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(); let vbv_maxrate = bitrate; let vbv_bufsize = bitrate / 4; let val = CString::new(format!( "repeat_headers=1:vbv-maxrate={vbv_maxrate}:vbv-bufsize={vbv_bufsize}" )) .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 { 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::::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::::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; let vbv_maxrate = bitrate; let vbv_bufsize = bitrate / 4; let opts = format!("repeat_headers=1:vbv-maxrate={vbv_maxrate}:vbv-bufsize={vbv_bufsize}"); assert!(opts.contains("vbv-maxrate=5000000")); assert!(opts.contains("vbv-bufsize=1250000")); } #[test] fn vbv_bufsize_is_quarter_of_maxrate() { for bitrate in [1_000_000, 5_000_000, 10_000_000] { let maxrate = bitrate; let bufsize = bitrate / 4; assert_eq!(bufsize * 4, maxrate, "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" ); } }