diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 19af2119..92e53e93 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -26,6 +26,22 @@ jobs: apt-get update apt-get install -y --no-install-recommends libgtk-4-dev libadwaita-1-dev libsdl3-dev + # The committed pf-zerocopy SPIR-V blobs are pulled in with include_bytes! and rebuilt only + # by hand — edit a .comp, forget the rebuild, and the OLD kernel ships with no compile error + # or failing test. Recompile each shader and diff the disassembly. Filtering OpSourceExtension + # (+ --no-header) is exactly what absorbs the shaderc-vs-glslang generator difference; every + # instruction, ID and constant must match. + - name: Shader SPIR-V drift gate (pf-zerocopy) + run: | + apt-get install -y --no-install-recommends glslang-tools spirv-tools + for s in rgb2nv12_buf cursor_blend; do + d=crates/pf-zerocopy/src/imp + glslangValidator -V "$d/$s.comp" -o "/tmp/$s.spv" >/dev/null + diff <(spirv-dis --no-header "$d/$s.spv" | grep -v OpSourceExtension) \ + <(spirv-dis --no-header "/tmp/$s.spv" | grep -v OpSourceExtension) \ + || { echo "::error::$d/$s.spv is stale — rebuild it from $s.comp"; exit 1; } + done + # Best-effort caches (act_runner's built-in cache server). Keyed on Cargo.lock: # registry/git are download caches, target/ the incremental build. The target key # carries the rustc version — resolved via `rustc --version` (below) rather than parsed diff --git a/crates/pf-capture/src/linux/pipewire.rs b/crates/pf-capture/src/linux/pipewire.rs index 70e10810..4fbd3a53 100644 --- a/crates/pf-capture/src/linux/pipewire.rs +++ b/crates/pf-capture/src/linux/pipewire.rs @@ -321,14 +321,15 @@ fn consume_frame(ud: &mut UserData, spa_buf: *mut spa::sys::spa_buffer) { // attach no fence. Covers both the GPU import and the CPU mmap read below. if datas[0].type_() == pw::spa::buffer::DataType::DmaBuf { match pf_zerocopy::dmabuf_fence::wait_read_ready(datas[0].fd(), 100) { - Ok(waited) => { + Ok(outcome) => { static F1: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(true); if F1.swap(false, Ordering::Relaxed) { tracing::info!( - waited, - "dmabuf implicit-fence sync active (waited=true → driver fences \ - the render, race closed; false → no implicit fence, zero-copy \ - may still show stale frames)" + ?outcome, + "dmabuf implicit-fence sync active (Signaled → driver fences the \ + render, race closed; NoFence → no implicit fence, zero-copy may \ + still show stale frames; TimedOut → fence pending past 100ms, \ + proceeded anyway)" ); } } diff --git a/crates/pf-zerocopy/src/dmabuf_fence.rs b/crates/pf-zerocopy/src/dmabuf_fence.rs index cb5afb24..406d99cc 100644 --- a/crates/pf-zerocopy/src/dmabuf_fence.rs +++ b/crates/pf-zerocopy/src/dmabuf_fence.rs @@ -11,12 +11,14 @@ //! fd we can `poll()` — readable once the producer's writes complete. This makes zero-copy capture //! race-free WITHOUT the producer doing anything, *iff* the driver actually attaches the fence. If it //! attaches none, the export yields an already-signaled sync_file (poll returns immediately) — no -//! wait, no harm, and `waited=false` tells us the driver doesn't fence (so zero-copy would still race). +//! wait, no harm, and [`WaitOutcome::NoFence`] tells us the driver doesn't fence (so zero-copy +//! would still race). // Every `unsafe` block in this file carries a `// SAFETY:` proof; enforce it (unsafe-proof program). #![deny(clippy::undocumented_unsafe_blocks)] use std::os::fd::RawFd; +use std::time::{Duration, Instant}; // linux/dma-buf.h ioctls on the DMA_BUF_BASE ('b' = 0x62) magic. _IOWR = dir(3)<<30 | size<<16 | base<<8 | nr. const DMA_BUF_BASE: u64 = 0x62; @@ -34,11 +36,25 @@ const DMA_BUF_IOCTL_EXPORT_SYNC_FILE: u64 = iowr(2, std::mem::size_of:: std::io::Result { +/// What the implicit-fence wait actually observed — the operator-facing diagnostic for "does +/// implicit fencing work on this box" must not conflate these (a timeout or an interrupted wait +/// proceeds with a possibly mid-render buffer; a signaled fence closed the race for real). +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +pub enum WaitOutcome { + /// No sync_file / already signaled — the driver attaches no implicit fence (or the render + /// finished before we looked); zero-copy can still race. + NoFence, + /// A render was in flight and we waited until its fence signaled — the race was real, now closed. + Signaled, + /// A render was in flight and `timeout_ms` elapsed first — we proceed (fail-open: blocking + /// longer would stall capture), possibly encoding a mid-render buffer. + TimedOut, +} + +/// Wait until the producer's writes to `dmabuf_fd` complete (or `timeout_ms` elapses; negative = +/// no timeout). `Err` means the ioctl or the poll itself failed (e.g. the kernel/driver lacks +/// `EXPORT_SYNC_FILE`); see [`WaitOutcome`] for the success cases. +pub fn wait_read_ready(dmabuf_fd: RawFd, timeout_ms: i32) -> std::io::Result { let mut req = DmaBufExportSyncFile { flags: DMA_BUF_SYNC_READ, fd: -1, @@ -54,31 +70,90 @@ pub fn wait_read_ready(dmabuf_fd: RawFd, timeout_ms: i32) -> std::io::Result= 0`). - // `poll` reads `fd`/`events` and writes `revents` for this non-blocking (timeout 0) probe, then - // returns — `pfd` outlives the call and aliases nothing. - let pending = unsafe { libc::poll(&mut pfd, 1, 0) } == 0; - if pending { - pfd.revents = 0; - // SAFETY: same live single-element `pfd` (its `revents` reset to 0 just above), `nfds == 1`, - // and `sync_fd` still open. This blocking `poll` (up to `timeout_ms`) waits for the render - // fence to signal; it reads `fd`/`events`, writes `revents`, and returns before `pfd` ends. - unsafe { libc::poll(&mut pfd, 1, timeout_ms) }; // block until the render fence signals + return Ok(WaitOutcome::NoFence); // no sync_file exported } + let outcome = poll_readable(sync_fd, timeout_ms); // SAFETY: `sync_fd` is the sync_file fd the EXPORT_SYNC_FILE ioctl created and handed us to own; // this point is reached only when `sync_fd >= 0`, this `close` runs exactly once on it, and it is // never used afterward — no double-close or use-after-close. unsafe { libc::close(sync_fd) }; - Ok(pending) + outcome +} + +/// Poll `fd` for `POLLIN`: a non-blocking probe first (already-readable ⇒ [`WaitOutcome::NoFence`] +/// — the fence was signaled before we looked), then a blocking wait up to `timeout_ms` (negative = +/// no timeout). `EINTR` retries with the remaining budget instead of silently skipping the wait — +/// the host spawns subprocesses, so a `SIGCHLD` mid-poll is a real occurrence, and skipping here +/// would reopen the exact stale-frame race this file exists to close. +fn poll_readable(fd: RawFd, timeout_ms: i32) -> std::io::Result { + let mut pfd = libc::pollfd { + fd, + events: libc::POLLIN, + revents: 0, + }; + // Non-blocking probe: not-yet-readable (poll == 0) means the producer is still rendering. + let probed = loop { + // SAFETY: `&mut pfd` points at a single live `libc::pollfd` and `nfds == 1` matches that + // one element; `fd` is the caller's live sync_file fd. `poll` reads `fd`/`events` and + // writes `revents` for this non-blocking (timeout 0) probe, then returns — `pfd` outlives + // the call and aliases nothing. + let r = unsafe { libc::poll(&mut pfd, 1, 0) }; + if r >= 0 { + break r; + } + let e = std::io::Error::last_os_error(); + if e.raw_os_error() != Some(libc::EINTR) { + return Err(e); + } + }; + if probed > 0 { + if pfd.revents & libc::POLLIN != 0 { + return Ok(WaitOutcome::NoFence); // signaled before we looked + } + // POLLERR/POLLNVAL without POLLIN — the fd is broken, not signaled. + return Err(std::io::Error::other(format!( + "poll(sync_file) revents {:#x} without POLLIN", + pfd.revents + ))); + } + let deadline = + (timeout_ms >= 0).then(|| Instant::now() + Duration::from_millis(timeout_ms as u64)); + loop { + let remaining = match deadline { + None => -1, // poll's "no timeout" + Some(d) => match d.checked_duration_since(Instant::now()) { + None => return Ok(WaitOutcome::TimedOut), + // +1: round up so a sub-millisecond remainder still waits instead of busy-polling. + Some(rem) => (rem.as_millis() as i32).saturating_add(1), + }, + }; + pfd.revents = 0; + // SAFETY: same live single-element `pfd` (its `revents` reset to 0 just above), `nfds == 1`, + // and `fd` still open (closed by the caller only after this function returns). This blocking + // `poll` (up to `remaining` ms) waits for the fence to signal; it reads `fd`/`events`, + // writes `revents`, and returns before `pfd` ends. + let r = unsafe { libc::poll(&mut pfd, 1, remaining) }; + match r { + 0 => return Ok(WaitOutcome::TimedOut), + r if r > 0 => { + if pfd.revents & libc::POLLIN != 0 { + return Ok(WaitOutcome::Signaled); + } + // POLLERR/POLLNVAL without POLLIN — the fd is broken, not signaled. + return Err(std::io::Error::other(format!( + "poll(sync_file) revents {:#x} without POLLIN", + pfd.revents + ))); + } + _ => { + let e = std::io::Error::last_os_error(); + if e.raw_os_error() != Some(libc::EINTR) { + return Err(e); + } + // Interrupted — recompute the remaining budget and keep waiting. + } + } + } } #[cfg(test)] @@ -90,4 +165,24 @@ mod tests { fn ioctl_number_matches_dma_buf_h() { assert_eq!(DMA_BUF_IOCTL_EXPORT_SYNC_FILE, 0xC008_6202); } + + /// The poll state machine, driven by a pipe (readable ⇔ signaled): a not-yet-readable fd that + /// stays quiet is a truthful `TimedOut` (not the old `Ok(true)`), and one already readable at + /// the probe is `NoFence`. + #[test] + fn poll_readable_reports_the_truth() { + use std::io::Write; + use std::os::fd::AsRawFd; + + let (r, mut w) = std::io::pipe().unwrap(); + assert_eq!( + poll_readable(r.as_raw_fd(), 10).unwrap(), + WaitOutcome::TimedOut + ); + w.write_all(b"x").unwrap(); + assert_eq!( + poll_readable(r.as_raw_fd(), 10).unwrap(), + WaitOutcome::NoFence + ); + } } diff --git a/crates/pf-zerocopy/src/imp/client.rs b/crates/pf-zerocopy/src/imp/client.rs index e260dc6a..303ef213 100644 --- a/crates/pf-zerocopy/src/imp/client.rs +++ b/crates/pf-zerocopy/src/imp/client.rs @@ -1,5 +1,5 @@ //! Host side of the isolated zero-copy GPU import (design: -//! [`design/zerocopy-worker-isolation.md`]): spawns the `zerocopy-worker` subprocess, mirrors the +//! `design/zerocopy-worker-isolation.md`): spawns the `zerocopy-worker` subprocess, mirrors the //! [`super::egl::EglImporter`] entry points over the [`super::proto`] socket, and materializes //! the worker's pooled CUDA buffers in this process via CUDA IPC (each buffer's handles are //! opened exactly once and reused as the pool recycles). A worker death — the whole point of the @@ -151,7 +151,7 @@ pub struct RemoteImporter { impl RemoteImporter { /// Spawn the worker from this host binary and complete the readiness handshake. The worker - /// is exec'd through the pinned [`SELF_EXE`] fd, so it is always the exact image this + /// is exec'd through the pinned `SELF_EXE` fd, so it is always the exact image this /// process runs — even after the installed binary was replaced mid-flight. An `Err` here /// means "no isolated zero-copy available" — callers fall back to the CPU path, exactly like /// an in-process `EglImporter::new()` failure. diff --git a/crates/pf-zerocopy/src/imp/cuda.rs b/crates/pf-zerocopy/src/imp/cuda.rs index 7a5a98f1..0526e558 100644 --- a/crates/pf-zerocopy/src/imp/cuda.rs +++ b/crates/pf-zerocopy/src/imp/cuda.rs @@ -1,4 +1,4 @@ -//! CUDA driver-side state for the zero-copy path, layered over the raw driver-API FFI in [`ffi`] +//! CUDA driver-side state for the zero-copy path, layered over the raw driver-API FFI in `ffi` //! (the `dlopen`'d `libcuda.so.1` symbol table — hand-rolled because no Rust crate exposes the //! GL-interop calls, and runtime-loaded so one binary runs on NVIDIA *and* on AMD/Intel where //! `libcuda` is absent). This facade owns the higher-level pieces on top of that layer: @@ -356,7 +356,10 @@ fn alloc_pitched_nv12( // SAFETY: two independent `cuMemAllocPitch_v2` calls (wrapper → live table). `&mut y_ptr`/ // `&mut y_pitch` and `&mut uv_ptr`/`&mut uv_pitch` are live, distinct stack out-params the // driver writes each plane's pointer and pitch into; all outlive their synchronous calls. The - // dimension/element-size args are by-value ints. No aliasing — four separate locals. + // dimension/element-size args are by-value ints. No aliasing — four separate locals. If the UV + // allocation fails, the just-created Y allocation is freed before the error propagates — this + // runs per frame under VRAM pressure (`BufferPool::get`'s pool-miss fallback), so a leak here + // would compound exactly when memory is already scarce. unsafe { ck( cuMemAllocPitch_v2( @@ -369,7 +372,7 @@ fn alloc_pitched_nv12( "cuMemAllocPitch_v2(Y)", )?; // Chroma is W/2 samples wide at 2 bytes each = W bytes; H/2 rows. - ck( + if let Err(e) = ck( cuMemAllocPitch_v2( &mut uv_ptr, &mut uv_pitch, @@ -378,7 +381,10 @@ fn alloc_pitched_nv12( 16, ), "cuMemAllocPitch_v2(UV)", - )?; + ) { + let _ = cuMemFree_v2(y_ptr); + return Err(e); + } } Ok(((y_ptr, y_pitch), (uv_ptr, uv_pitch))) } @@ -432,7 +438,7 @@ impl InputSurface { } /// Planar YUV444 (8-bit 4:4:4): one allocation, Y|U|V full-res planes stacked (see - /// [`alloc_pitched_yuv444`]). + /// `alloc_pitched_yuv444`). pub fn alloc_yuv444(width: u32, height: u32) -> Result { let (ptr, pitch) = alloc_pitched_yuv444(width, height)?; Ok(InputSurface { ptr, pitch, height }) @@ -509,7 +515,7 @@ pub struct BufferPool { pitch: usize, /// NV12 pools carry a second (chroma) pitch; `Some` ⇒ buffers from this pool have a UV plane. uv_pitch: Option, - /// YUV444 pools: one allocation of 3·`height` stacked 1-byte planes (see [`alloc_pitched_yuv444`]). + /// YUV444 pools: one allocation of 3·`height` stacked 1-byte planes (see `alloc_pitched_yuv444`). yuv444: bool, } @@ -629,10 +635,10 @@ pub struct DeviceBuffer { pub pitch: usize, pub width: u32, pub height: u32, - /// NV12 only: the interleaved chroma plane `(ptr, pitch)` paired with the Y plane in [`ptr`]. - /// `None` for the default 4-byte RGB/BGRx path. When `Some`, [`ptr`] is the Y plane (1 byte/px). + /// NV12 only: the interleaved chroma plane `(ptr, pitch)` paired with the Y plane in [`ptr`](Self::ptr). + /// `None` for the default 4-byte RGB/BGRx path. When `Some`, [`ptr`](Self::ptr) is the Y plane (1 byte/px). pub uv: Option<(CUdeviceptr, usize)>, - /// Planar YUV444: [`ptr`] is ONE allocation of 3·[`height`](Self::height) rows at + /// Planar YUV444: [`ptr`](Self::ptr) is ONE allocation of 3·[`height`](Self::height) rows at /// [`pitch`](Self::pitch) — the full-res 1-byte Y, U, V planes stacked in that order /// (`uv` stays `None`; the single-plane wire/IPC path carries it unchanged). pub yuv444: bool, @@ -1005,8 +1011,15 @@ pub fn copy_nv12_to_device( // `sync: false` shifts the source-lifetime obligation to the caller (documented above). // Wrappers → live table. unsafe { - copy_async(&y, "cuMemcpy2DAsync_v2(nv12 Y dev->dev)")?; - copy_async(&uv, "cuMemcpy2DAsync_v2(nv12 UV dev->dev)")?; + // On a failed enqueue, drain the stream before propagating: the caller drops `src` on + // `Err`, which recycles it into its pool — a copy still in flight would then race the + // next frame written into the same allocation. + let r = copy_async(&y, "cuMemcpy2DAsync_v2(nv12 Y dev->dev)") + .and_then(|()| copy_async(&uv, "cuMemcpy2DAsync_v2(nv12 UV dev->dev)")); + if r.is_err() { + let _ = sync_copy_stream(); + return r; + } if sync { sync_copy_stream()?; } @@ -1045,8 +1058,15 @@ pub fn copy_yuv444_to_device( // `src.ptr + pitch·h·i` stays within the live 3·H-row stacked allocation (`yuv444` // checked above), `dst_ptr`/`dst_pitch` is the caller's live NVENC plane; `w`×`h` fits // both. Completion is the trailing stream sync below (`sync`) or the caller's - // stream-ordering obligation (`sync: false`, documented above). Wrapper → live table. - unsafe { copy_async(©, "cuMemcpy2DAsync_v2(yuv444 plane dev->dev)")? }; + // stream-ordering obligation (`sync: false`, documented above). A failed enqueue drains + // the stream first — earlier planes are already queued, and the caller drops (recycles) + // `src` on `Err`. Wrapper → live table. + unsafe { + if let Err(e) = copy_async(©, "cuMemcpy2DAsync_v2(yuv444 plane dev->dev)") { + let _ = sync_copy_stream(); + return Err(e); + } + } } if sync { // SAFETY: one stream sync after the last enqueue covers all three planes (FIFO) — the diff --git a/crates/pf-zerocopy/src/imp/cursor_blend.comp b/crates/pf-zerocopy/src/imp/cursor_blend.comp index 0d940361..41aa8479 100644 --- a/crates/pf-zerocopy/src/imp/cursor_blend.comp +++ b/crates/pf-zerocopy/src/imp/cursor_blend.comp @@ -17,7 +17,8 @@ // rows for NV12, whose 2 chroma bytes-pairs land in one exclusive word). Spans are aligned to the // SURFACE, not the cursor, so neighbouring invocations never share a word even at odd `ox`. // -// Rebuild: glslc cursor_blend.comp -o cursor_blend.spv (vendored beside this file) +// Rebuild: glslangValidator -V cursor_blend.comp -o cursor_blend.spv (vendored beside this +// file; or glslc — CI diffs the disassembly against this source) layout(local_size_x = 8, local_size_y = 8, local_size_z = 1) in; diff --git a/crates/pf-zerocopy/src/imp/egl.rs b/crates/pf-zerocopy/src/imp/egl.rs index 4f4c482c..41f95733 100644 --- a/crates/pf-zerocopy/src/imp/egl.rs +++ b/crates/pf-zerocopy/src/imp/egl.rs @@ -8,7 +8,7 @@ //! format is opaque). So we follow OBS/Sunshine: bind the `EGLImage` to a GL texture //! (`glEGLImageTargetTexture2DOES`), render it through a fullscreen-triangle shader into a plain //! immutable `GL_RGBA8` texture (de-tiling and swizzling to the BGRx the encoder wants), then -//! register *that* texture with CUDA ([`MappedTexture`]) and copy it device-to-device into an +//! register *that* texture with CUDA (`cuda::RegisteredTexture`) and copy it device-to-device into an //! owned [`DeviceBuffer`] so the dmabuf can be returned to the compositor immediately. #![allow(non_upper_case_globals)] diff --git a/crates/pf-zerocopy/src/imp/mod.rs b/crates/pf-zerocopy/src/imp/mod.rs index e51d1278..2e2893ff 100644 --- a/crates/pf-zerocopy/src/imp/mod.rs +++ b/crates/pf-zerocopy/src/imp/mod.rs @@ -24,11 +24,33 @@ use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; pub use cuda::DeviceBuffer; pub use egl::{DmabufPlane, EglImporter}; -/// Whether a `PUNKTFUNK_*` flag is truthy (`1`/`true`/`yes`/`on`), or `None` when unset. +/// Whether a `PUNKTFUNK_*` flag is truthy (`1`/`true`/`yes`/`on`, any case), falsy (`0`/`false`/ +/// `no`/`off`, any case), or `None` when unset — and, crucially, `None` for an *unrecognised* +/// spelling too. These flags default ON (`PUNKTFUNK_ZEROCOPY`, `PUNKTFUNK_NV12`), so treating a +/// typo'd `TRUE` as "off" silently inverted the operator's intent host-wide (a systemd drop-in or +/// Nix module is exactly where such spellings come from). Unrecognised values are logged once so +/// they are visible instead of silent. fn flag_opt(name: &str) -> Option { - std::env::var(name) - .ok() - .map(|v| matches!(v.trim(), "1" | "true" | "yes" | "on")) + let v = std::env::var(name).ok()?; + match v.trim().to_ascii_lowercase().as_str() { + "1" | "true" | "yes" | "on" => Some(true), + "0" | "false" | "no" | "off" | "" => Some(false), + _ => { + use std::collections::HashSet; + use std::sync::Mutex; + static WARNED: Mutex>> = Mutex::new(None); + let mut g = WARNED.lock().unwrap(); + if g.get_or_insert_with(HashSet::new).insert(name.to_string()) { + tracing::warn!( + flag = name, + value = %v, + "unrecognised boolean value for this PUNKTFUNK_* flag — expected \ + 1/true/yes/on or 0/false/no/off (any case); using the flag's default" + ); + } + None + } + } } /// Whether a `PUNKTFUNK_*` flag is truthy (`1`/`true`/`yes`/`on`); unset ⇒ false. @@ -49,7 +71,7 @@ pub fn vaapi_dmabuf_forced() -> bool { /// (`design/zerocopy-worker-isolation.md`), so a driver fault on a producer-invalidated dmabuf kills /// the worker and the host degrades to its capture-loss rebuild instead of dying — the reason the /// NVENC path stayed opt-in is gone. Fallbacks stay in place: VAAPI has a one-shot CPU downgrade if -/// the LINEAR-dmabuf offer never negotiates ([`note_vaapi_dmabuf_failed`]); NVENC falls back per +/// the LINEAR-dmabuf offer never negotiates ([`note_raw_dmabuf_negotiation_failed`]); NVENC falls back per /// capture when no importer/importable modifier is available and latches the import off after /// repeated worker deaths. `PUNKTFUNK_ZEROCOPY=0` opts out; `PUNKTFUNK_FORCE_SHM` forces the /// race-free SHM path. @@ -204,7 +226,7 @@ static GPU_IMPORT_DISABLED: AtomicBool = AtomicBool::new(false); const GPU_IMPORT_DEATH_LATCH: u32 = 3; /// Record a worker death (transport-level failure). Latches the process-wide disable after -/// [`GPU_IMPORT_DEATH_LATCH`] consecutive deaths. +/// `GPU_IMPORT_DEATH_LATCH` consecutive deaths. pub fn note_gpu_import_death() { let streak = GPU_IMPORT_DEATH_STREAK.fetch_add(1, Ordering::Relaxed) + 1; if streak >= GPU_IMPORT_DEATH_LATCH && !GPU_IMPORT_DISABLED.swap(true, Ordering::Relaxed) { @@ -243,7 +265,7 @@ static RAW_DMABUF_DISABLED: AtomicBool = AtomicBool::new(false); const RAW_DMABUF_FAILURE_LATCH: u32 = 3; /// Record an encoder-side raw-dmabuf import failure. Latches the process-wide disable after -/// [`RAW_DMABUF_FAILURE_LATCH`] consecutive failures. +/// `RAW_DMABUF_FAILURE_LATCH` consecutive failures. pub fn note_raw_dmabuf_import_failure(reason: &str) { let streak = RAW_DMABUF_FAILURE_STREAK.fetch_add(1, Ordering::Relaxed) + 1; if streak >= RAW_DMABUF_FAILURE_LATCH && !RAW_DMABUF_DISABLED.swap(true, Ordering::Relaxed) { @@ -465,6 +487,38 @@ pub fn nv12_selftest() -> anyhow::Result<()> { mod tests { use super::*; + /// [`bt709_limited`] is the sole oracle for the GPU colour self-test — a coefficient typo here + /// would fail the self-test on a correct GPU (operator blames the driver), or, mirrored into + /// the shaders, pass it on genuinely wrong output. Pin it to externally-known BT.709 + /// limited-range anchors instead of trusting it to check itself. + #[test] + fn bt709_limited_reference_matches_known_anchors() { + let close = |a: f64, b: f64| (a - b).abs() < 1e-9; + + let (y, u, v) = bt709_limited(0, 0, 0); // black + assert!( + close(y, 16.0) && close(u, 128.0) && close(v, 128.0), + "black → ({y}, {u}, {v})" + ); + + let (y, u, v) = bt709_limited(255, 255, 255); // white + assert!( + close(y, 235.0) && close(u, 128.0) && close(v, 128.0), + "white → ({y}, {u}, {v})" + ); + + // Pure red saturates V (Kr row sums to exactly +0.5), pure blue saturates U. + let (_, _, v) = bt709_limited(255, 0, 0); + assert!(close(v, 240.0), "red V → {v}"); + let (_, u, _) = bt709_limited(0, 0, 255); + assert!(close(u, 240.0), "blue U → {u}"); + + // One mid-scale anchor so a swapped Kr/Kb pair can't cancel out: BT.709 Y of pure green + // is 16 + 219·0.7152. + let (y, _, _) = bt709_limited(0, 255, 0); + assert!(close(y, 16.0 + 219.0 * 0.7152), "green Y → {y}"); + } + /// Single test owning the process-global latch statics (they are never reset by design). #[test] fn gpu_import_death_latch() { diff --git a/crates/pf-zerocopy/src/imp/proto.rs b/crates/pf-zerocopy/src/imp/proto.rs index 01635347..ccc774db 100644 --- a/crates/pf-zerocopy/src/imp/proto.rs +++ b/crates/pf-zerocopy/src/imp/proto.rs @@ -1,6 +1,6 @@ //! Wire protocol between the PipeWire capture thread and the isolated zero-copy GPU-import //! worker process (`punktfunk-host zerocopy-worker`; design: -//! [`design/zerocopy-worker-isolation.md`]). Transport is a `SOCK_SEQPACKET` unix socketpair — +//! `design/zerocopy-worker-isolation.md`). Transport is a `SOCK_SEQPACKET` unix socketpair — //! reliable, ordered, message-framed (one `sendmsg` = one message) — with dmabuf fds riding as //! `SCM_RIGHTS` control data. Bodies are small serde_json blobs (~200 B/frame); pixels never //! cross the socket (they move GPU-side via CUDA IPC, see [`super::cuda::ipc_export`]). diff --git a/crates/pf-zerocopy/src/imp/rgb2nv12_buf.comp b/crates/pf-zerocopy/src/imp/rgb2nv12_buf.comp index fe809206..8e25cbf3 100644 --- a/crates/pf-zerocopy/src/imp/rgb2nv12_buf.comp +++ b/crates/pf-zerocopy/src/imp/rgb2nv12_buf.comp @@ -11,7 +11,8 @@ // touch the same word (an 8-bit-storage-free way to write byte planes race-free). All pitches // and the UV offset are in WORDS and must be word-aligned (the Rust side sizes them so). // -// Rebuild: glslc rgb2nv12_buf.comp -o rgb2nv12_buf.spv +// Rebuild: glslangValidator -V rgb2nv12_buf.comp -o rgb2nv12_buf.spv (or glslc; CI diffs the +// disassembly against this source, tolerating only the generator header) layout(local_size_x = 8, local_size_y = 8) in; layout(std430, binding = 0) readonly buffer Src { diff --git a/crates/pf-zerocopy/src/imp/vkslot.rs b/crates/pf-zerocopy/src/imp/vkslot.rs index 1a3dec0e..a61426a9 100644 --- a/crates/pf-zerocopy/src/imp/vkslot.rs +++ b/crates/pf-zerocopy/src/imp/vkslot.rs @@ -45,7 +45,7 @@ use ash::vk; pub const CURSOR_MAX: u32 = cuda::CURSOR_MAX; /// The vendored SPIR-V for `cursor_blend.comp` (beside this file; rebuild with -/// `glslc cursor_blend.comp -o cursor_blend.spv`). +/// `glslangValidator -V cursor_blend.comp -o cursor_blend.spv`; CI gates drift). const CURSOR_SPV: &[u8] = include_bytes!("cursor_blend.spv"); /// NVENC input-surface layout — selects the spec-constant `MODE` pipeline and the allocation @@ -84,9 +84,10 @@ impl SlotFormat { } /// What the encoder holds per ring slot: the CUDA view it registers with NVENC plus the id it -/// hands back to [`VkSlotBlend::blend`]. The backing Vulkan objects + CUDA mapping live in the -/// [`VkSlotBlend`] (freed by [`free_slots`](VkSlotBlend::free_slots) / drop), so this is Copy — -/// the encoder's ring keeps its existing shape. +/// hands back to [`VkSlotBlend::blend_ref`] / [`VkSlotBlend::blend_ref_ordered`]. The backing +/// Vulkan objects + CUDA mapping live in the [`VkSlotBlend`] (freed by +/// [`free_slots`](VkSlotBlend::free_slots) / drop), so this is Copy — the encoder's ring keeps +/// its existing shape. #[derive(Clone, Copy)] pub struct VkSlotRef { /// Device pointer NVENC registers (CUDA's mapping of the Vulkan memory). @@ -654,7 +655,7 @@ impl VkSlotBlend { } /// Free every allocated slot (encoder teardown, alongside its ring clear). CUDA mappings drop - /// first (field order in [`SlotAlloc`] frees `cuda` via its own `Drop` before we free the VK + /// first (field order in `SlotAlloc` frees `cuda` via its own `Drop` before we free the VK /// objects explicitly here). pub fn free_slots(&mut self) { // Ordered blends return with their work still on the queue — quiesce before freeing the @@ -1062,3 +1063,62 @@ impl Drop for VkSlotBlend { } } } + +#[cfg(test)] +mod tests { + use super::*; + + fn slot() -> VkSlotRef { + VkSlotRef { + ptr: 0, + pitch: 2048, + height: 1080, + id: 0, + } + } + + fn geo(fmt: SlotFormat, cw: u32, ch: u32, ox: i32, oy: i32) -> Option<(Push, u32, u32)> { + VkSlotBlend::blend_geometry(&slot(), fmt, 1920, cw, ch, ox, oy) + } + + /// An empty (clamped-away) cursor rect dispatches nothing. + #[test] + fn empty_rect_is_none() { + assert!(geo(SlotFormat::Argb, 0, 32, 10, 10).is_none()); + assert!(geo(SlotFormat::Nv12, 32, 0, 10, 10).is_none()); + } + + /// Oversized bitmaps clamp to `CURSOR_MAX` — the push constants must agree with the staging + /// buffer's capacity, or the shader reads past the uploaded bitmap. + #[test] + fn cursor_dims_clamp_to_max() { + let (push, _, _) = geo(SlotFormat::Argb, CURSOR_MAX + 100, CURSOR_MAX + 1, 0, 0).unwrap(); + assert_eq!(push.cur_w, CURSOR_MAX); + assert_eq!(push.cur_h, CURSOR_MAX); + } + + /// ARGB dispatches per cursor pixel; NV12/YUV444 per word-aligned 4-px span, NV12 walking + /// 2-row blocks and YUV444 single rows. A 32×32 cursor at ox=13: spans cover the aligned + /// x∈[12,48) → 9 spans → 2 groups of 8. + #[test] + fn group_counts_per_format() { + let (_, gx, gy) = geo(SlotFormat::Argb, 32, 32, 13, 0).unwrap(); + assert_eq!((gx, gy), (4, 4)); // 32/8 in both axes + + let (_, gx, gy) = geo(SlotFormat::Nv12, 32, 32, 13, 0).unwrap(); + assert_eq!((gx, gy), (2, 2)); // 9 spans → 2 groups; 16 2-row blocks → 2 groups + + let (_, gx, gy) = geo(SlotFormat::Yuv444, 32, 32, 13, 0).unwrap(); + assert_eq!((gx, gy), (2, 4)); // 9 spans → 2 groups; 32 rows → 4 groups + } + + /// Negative `ox` must anchor spans with FLOOR alignment (`>>` on the signed value), not + /// truncating division: at ox=-5 the aligned start is -8, giving 9 spans over a 32-px cursor + /// (truncation would start at -4 and dispatch only 8 — silently dropping the right edge). + #[test] + fn negative_ox_floor_aligns_the_span_origin() { + let (push, gx, _) = geo(SlotFormat::Nv12, 32, 32, -5, 0).unwrap(); + assert_eq!(push.ox, -5, "push constants carry the true origin"); + assert_eq!(gx, 2, "9 spans from the floor-aligned start → 2 groups"); + } +} diff --git a/crates/pf-zerocopy/src/imp/vulkan.rs b/crates/pf-zerocopy/src/imp/vulkan.rs index cd8806fa..b0a9c846 100644 --- a/crates/pf-zerocopy/src/imp/vulkan.rs +++ b/crates/pf-zerocopy/src/imp/vulkan.rs @@ -51,7 +51,7 @@ struct Csc { } /// The buffer-to-buffer RGB→NV12 compute shader (see `rgb2nv12_buf.comp` beside this file; -/// rebuild with `glslc rgb2nv12_buf.comp -o rgb2nv12_buf.spv`). +/// rebuild with `glslangValidator -V rgb2nv12_buf.comp -o rgb2nv12_buf.spv`; CI gates drift). const CSC_SPV: &[u8] = include_bytes!("rgb2nv12_buf.spv"); pub struct VkBridge { diff --git a/crates/pf-zerocopy/src/imp/worker.rs b/crates/pf-zerocopy/src/imp/worker.rs index d6ce0a3a..8fcf3304 100644 --- a/crates/pf-zerocopy/src/imp/worker.rs +++ b/crates/pf-zerocopy/src/imp/worker.rs @@ -1,5 +1,5 @@ //! The isolated zero-copy GPU-import worker (`punktfunk-host zerocopy-worker`; design: -//! [`design/zerocopy-worker-isolation.md`]). It owns the fragile driver stack — the headless +//! `design/zerocopy-worker-isolation.md`). It owns the fragile driver stack — the headless //! EGLDisplay + GL context, the CUDA context, and the Vulkan bridge — so that a driver fault on a //! producer-invalidated dmabuf (the `cuGraphicsMapResources` SIGSEGV the F44 Game→Desktop switch //! reproduced) kills THIS process, not the streaming host. The host observes the dead socket, @@ -43,9 +43,24 @@ pub fn run_from_args(args: &[String]) -> Result<()> { .transpose() .context("parse --fd")? .unwrap_or(3); + // Refuse to adopt anything that cannot be the spawning host's socket: a negative fd is + // outright UB inside `OwnedFd` (its niche), and 0–2 would make this worker close one of its + // own stdio streams on exit. Then confirm the number actually holds an open socket — the + // subcommand is hidden but runnable by hand, and adopting an arbitrary inherited fd would + // close it behind its real owner. + anyhow::ensure!(fd >= 3, "--fd must be >= 3 (got {fd})"); + // SAFETY: `libc::stat` is plain-old-data for which all-zero is a valid value, so + // `mem::zeroed()` is a sound initializer; `fstat` writes into the live, correctly-sized + // `&mut st` and only reads `fd`. `st_mode` is read only after the return value is checked. + let is_socket = unsafe { + let mut st: libc::stat = std::mem::zeroed(); + libc::fstat(fd, &mut st) == 0 && (st.st_mode & libc::S_IFMT) == libc::S_IFSOCK + }; + anyhow::ensure!(is_socket, "--fd {fd} is not an open socket"); // SAFETY: the spawning host `dup2`'d its socketpair end onto exactly this fd number before - // exec (the subcommand's contract) and nothing else in this fresh process owns it, so - // `OwnedFd` takes sole ownership and closes it exactly once at exit. + // exec (the subcommand's contract, just verified to be an open socket ≥ 3) and nothing else + // in this fresh process owns it, so `OwnedFd` takes sole ownership and closes it exactly + // once at exit. let sock = unsafe { OwnedFd::from_raw_fd(fd) }; run(sock) }