diff --git a/crates/pf-capture/src/linux/pipewire.rs b/crates/pf-capture/src/linux/pipewire.rs index ade8861c..0f80219d 100644 --- a/crates/pf-capture/src/linux/pipewire.rs +++ b/crates/pf-capture/src/linux/pipewire.rs @@ -96,6 +96,10 @@ struct UserData { /// into the first-frame-timeout retry loop; the promised renegotiation normally lands /// within a frame or two). gate_since: Option, + /// Deferred requeue of raw-passthrough buffers (see [`DeferredRequeue`]): the encode thread + /// reads the dmabuf long after `.process` returns, so the buffer must not rejoin the + /// producer's pool until the frame's [`BufferHold`] drops. + defer: std::sync::Arc, } impl UserData { @@ -113,6 +117,46 @@ impl UserData { } let _ = self.wake.try_send(()); } + + /// Withhold the raw-passthrough buffer from the producer's pool until the returned hold + /// drops — the deferred requeue that closes the rewrite-while-the-encoder-reads race. + /// `None` (pool too shallow, or `PUNKTFUNK_ZEROCOPY_HOLD=0`) falls back to the immediate + /// `.process`-epilogue requeue, i.e. the old racy contract; said once per session. + fn try_defer(&mut self, pw_buf: *mut pw::sys::pw_buffer) -> Option { + if !zerocopy_hold_enabled() { + return None; + } + let buf = pw_buf as usize; + let pool_live = self.pool.live; + let generation = self.defer.book.lock().ok()?.try_hold(buf, pool_live); + let Some(generation) = generation else { + if !self.defer.logged_shallow.swap(true, Ordering::Relaxed) { + tracing::warn!( + pool_depth = pool_live, + reserve = HOLD_POOL_RESERVE, + "zero-copy: the producer's buffer pool cannot spare a buffer to hold across \ + the encode — falling back to the immediate requeue, which the producer may \ + rewrite mid-encode (torn/discolored frames under load); PUNKTFUNK_FORCE_SHM=1 \ + trades CPU for a race-free capture if artifacts appear" + ); + } + return None; + }; + if !self.defer.logged_active.swap(true, Ordering::Relaxed) { + tracing::info!( + pool_depth = pool_live, + reserve = HOLD_POOL_RESERVE, + "zero-copy: withholding each published buffer from the producer until the \ + encoder releases it (deferred requeue — the producer can no longer rewrite a \ + frame mid-encode); PUNKTFUNK_ZEROCOPY_HOLD=0 restores the immediate requeue" + ); + } + Some(std::sync::Arc::new(BufferHold { + defer: self.defer.clone(), + buf, + generation, + })) + } } /// Everything the zero-copy negotiation decision depends on, gathered at ONE point in time. @@ -510,11 +554,12 @@ impl FenceWaitStats { /// PW5 stage 1: how many buffers the producer actually allocated for this stream. /// -/// **Nothing in this codebase had ever counted them.** The zero-copy path dups the dmabuf fd and -/// publishes the frame while the SPA buffer is handed straight back to the producer at `.process` -/// return — so the only thing keeping capture untorn is that the producer round-robins a pool -/// deeper than our import+encode window. That depth was an unmeasured assumption; this makes it a -/// logged number, on every producer, before anything is built on it. +/// **Nothing in this codebase had ever counted them.** The zero-copy path used to hand the SPA +/// buffer straight back to the producer at `.process` return, leaving pool depth as the only +/// thing keeping capture untorn. The deferred requeue ([`DeferredRequeue`]) now withholds +/// published buffers until the consumer is done, but the depth still matters twice over: it is +/// the budget `HoldBook::try_hold` spends (a pool of ≤ [`HOLD_POOL_RESERVE`] cannot defer at +/// all and runs the old race), and for un-deferred frames it remains the race window. /// /// `live` is maintained by the `add_buffer`/`remove_buffer` stream callbacks, which PipeWire fires /// on the loop thread as the pool is allocated (and again, remove-then-add, on a renegotiation that @@ -586,6 +631,104 @@ impl PassthroughFallbacks { /// short streak of dropped frames the capturer fails loudly and the session renegotiates. const IMPORT_FAIL_POISON: u32 = 3; +/// Buffers the deferred requeue always leaves in the producer's pool. One for the frame the +/// producer is rendering right now, one in transit — withholding past that would make the +/// producer skip frames whenever our holds are at their worst (host frame + up to two encoder +/// ring slots), which is a pacing hiccup, not corruption, but there is no reason to court it. +const HOLD_POOL_RESERVE: u32 = 2; + +/// `PUNKTFUNK_ZEROCOPY_HOLD=0` restores the immediate `.process`-return requeue (the racy +/// pre-hold behavior) — a field bisect lever, not a tuning knob. `env_on` grammar like every +/// other capture knob (a bare `== "0"` compare is the trap `PUNKTFUNK_FORCE_SHM` already fell in). +fn zerocopy_hold_enabled() -> bool { + static ON: std::sync::OnceLock = std::sync::OnceLock::new(); + *ON.get_or_init(|| pf_host_config::env_on("PUNKTFUNK_ZEROCOPY_HOLD").unwrap_or(true)) +} + +/// Pure bookkeeping for the deferred requeue: which buffers are currently withheld from the +/// producer, each under a per-hold generation so a pointer-value reuse across a pool +/// renegotiation can never satisfy a stale hold's release (see `complete`). +/// +/// Threading contract (what makes the single-requeue invariant hold with no atomics): entries are +/// INSERTED (`try_hold`) and REMOVED (`complete` via the requeue channel's callback, `purge` via +/// `remove_buffer`) only on the PipeWire loop thread; a dropping [`BufferHold`] on any other +/// thread only *sends* the release message. So between a hold's creation and the loop servicing +/// its release, `contains` is stable — which is exactly what the `.process` epilogue relies on to +/// decide "requeue now" vs "the hold owns the requeue". +#[derive(Default)] +struct HoldBook { + /// Withheld buffers: `*mut pw_buffer` as usize → the generation of the hold that owns it. + out: std::collections::HashMap, + /// Last issued hold generation (monotonic per stream). + last_gen: u64, +} + +impl HoldBook { + /// Withhold `buf` if the pool can spare it: at most `pool_live - HOLD_POOL_RESERVE` buffers + /// out at once. Returns the generation to release with, or `None` (pool too shallow / buffer + /// somehow already out — the caller falls back to the immediate requeue). + fn try_hold(&mut self, buf: usize, pool_live: u32) -> Option { + let cap = pool_live.saturating_sub(HOLD_POOL_RESERVE) as usize; + if self.out.len() >= cap || self.out.contains_key(&buf) { + return None; + } + self.last_gen += 1; + self.out.insert(buf, self.last_gen); + Some(self.last_gen) + } + + /// A hold released: take `buf` out of the book iff this generation still owns it. `true` ⇒ + /// the caller must requeue the buffer; `false` ⇒ the entry was purged (pool renegotiated — + /// the pointer may even be a NEW buffer under a reused address) and the buffer must NOT be + /// touched. + fn complete(&mut self, buf: usize, generation: u64) -> bool { + match self.out.get(&buf) { + Some(&g) if g == generation => { + self.out.remove(&buf); + true + } + _ => false, + } + } + + /// `remove_buffer`: the buffer is being freed under us (renegotiation/teardown) — forget it. + /// Its hold's later release finds the generation gone and becomes a no-op. + fn purge(&mut self, buf: usize) { + self.out.remove(&buf); + } + + fn contains(&self, buf: usize) -> bool { + self.out.contains_key(&buf) + } +} + +/// Shared between the loop thread ([`HoldBook`] ops) and the [`BufferHold`] guards riding +/// published frames to the encode thread. +struct DeferredRequeue { + book: std::sync::Mutex, + /// Wakes the loop to requeue `(buffer, generation)`. Send failure = the loop (and with it + /// the stream and every buffer) is gone — nothing to release. + tx: pw::channel::Sender<(usize, u64)>, + /// One-per-session lines: the first successful defer, and the shallow-pool fallback. + logged_active: std::sync::atomic::AtomicBool, + logged_shallow: std::sync::atomic::AtomicBool, +} + +/// The concrete [`pf_frame::FrameHold`]: releases its buffer back to the producer when the last +/// clone drops. Send-only from the dropping thread — the actual `pw_stream_queue_buffer` runs in +/// the requeue channel's loop-thread callback. +struct BufferHold { + defer: std::sync::Arc, + buf: usize, + generation: u64, +} + +impl Drop for BufferHold { + fn drop(&mut self) { + let _ = self.defer.tx.send((self.buf, self.generation)); + } +} + /// Log a frame-drop reason once per process (the process callback runs per frame; a stuck /// pipeline must say why without flooding). fn warn_once(msg: &'static str) { @@ -644,7 +787,14 @@ impl Drop for DmabufMap { /// `.process` callback with the NEWEST drained buffer (latest-frame-only). `datas` is sourced /// via the same transparent cast libspa's `Buffer::datas_mut` performs, so the safe `Data` /// accessors (`.type_()`, `.chunk()`, `.data()`, `.fd()`, `.as_raw()`) keep working. -fn consume_frame(ud: &mut UserData, spa_buf: *mut spa::sys::spa_buffer) { +/// +/// `pw_buf` is the buffer's `pw_buffer` handle (`spa_buf`'s owner), used only as the identity a +/// raw-passthrough publish withholds via [`UserData::try_defer`] — never dereferenced here. +fn consume_frame( + ud: &mut UserData, + spa_buf: *mut spa::sys::spa_buffer, + pw_buf: *mut pw::sys::pw_buffer, +) { // No active stream: release the buffer without the (expensive at 5K) de-pad. if !ud.signals.active.load(Ordering::Relaxed) { return; @@ -822,8 +972,11 @@ fn consume_frame(ud: &mut UserData, spa_buf: *mut spa::sys::spa_buffer) { None }; // dup the fd so it survives the SPA buffer recycle — the encode thread - // imports it. Content stability across the brief import/encode window relies - // on the compositor's buffer-pool depth, like any zero-copy capture. + // imports it. Content stability across the read window comes from the deferred + // requeue below (`try_defer` — the producer does not get this buffer back until + // the frame's hold drops); with no hold (shallow pool / PUNKTFUNK_ZEROCOPY_HOLD=0) + // it falls back to the compositor's pool depth outrunning the encode, the old + // racy contract. // SAFETY: `datas[0].fd()` is the dmabuf fd owned by the live PipeWire buffer (valid // for this callback). `fcntl(fd, F_DUPFD_CLOEXEC, 0)` reads only the integer fd, // touches no Rust memory, and returns a fresh independent CLOEXEC duplicate (or -1). @@ -836,6 +989,7 @@ fn consume_frame(ud: &mut UserData, spa_buf: *mut spa::sys::spa_buffer) { .duration_since(UNIX_EPOCH) .map(|d| d.as_nanos() as u64) .unwrap_or(0); + let hold = ud.try_defer(pw_buf); ud.publish(CapturedFrame { width: w as u32, height: h as u32, @@ -852,6 +1006,7 @@ fn consume_frame(ud: &mut UserData, spa_buf: *mut spa::sys::spa_buffer) { offset, stride, plane1, + hold, }), // Cursor-as-metadata is blended only by RGB→NV12 backends. Gamescope // embeds its pointer in the produced pixels, so native NV12 has none. @@ -1434,6 +1589,18 @@ pub fn pipewire_thread( ); } + // Deferred requeue (the rewrite-while-encoding fix): holds riding published frames release + // their buffers through this channel from whatever thread drops them last; the receiver — + // attached to the loop below, after the stream exists — is the single place a withheld + // buffer rejoins the producer's pool. + let (requeue_tx, requeue_rx) = pw::channel::channel::<(usize, u64)>(); + let defer = std::sync::Arc::new(DeferredRequeue { + book: std::sync::Mutex::new(HoldBook::default()), + tx: requeue_tx, + logged_active: std::sync::atomic::AtomicBool::new(false), + logged_shallow: std::sync::atomic::AtomicBool::new(false), + }); + let data = UserData { info: VideoInfoRaw::default(), format: None, @@ -1459,6 +1626,7 @@ pub fn pipewire_thread( }, gate_skips: 0, gate_since: None, + defer: defer.clone(), }; let stream = pw::stream::StreamBox::new( @@ -1562,10 +1730,18 @@ pub fn pipewire_thread( } }) // PW5 stage 1 — the pool census. PipeWire fires these on the loop thread as it allocates - // (and, on a renegotiation, frees then re-allocates) the stream's buffers. Counting only: - // the buffer pointer is not touched, so no lifetime question arises here. + // (and, on a renegotiation, frees then re-allocates) the stream's buffers. The census only + // counts; `remove_buffer` additionally purges the buffer from the deferred-requeue book — + // the buffer is being freed under any hold still riding a frame, so that hold's later + // release must become a no-op (the generation check in `HoldBook::complete` also covers + // the freed address being reused by a new pool's buffer). .add_buffer(|_stream, ud, _buf| ud.pool.add()) - .remove_buffer(|_stream, ud, _buf| ud.pool.remove()) + .remove_buffer(|_stream, ud, buf| { + ud.pool.remove(); + if let Ok(mut book) = ud.defer.book.lock() { + book.purge(buf as usize); + } + }) .process(|stream, ud| { // Latest-frame-only (OBS pattern): Mutter delivers buffers in bursts and recycles its // pool; an older queued buffer carries a STALE frame. Drain all queued buffers, requeue @@ -1598,19 +1774,19 @@ pub fn pipewire_thread( // value. MEASURED, not requested: `build_dmabuf_buffers` asks for a range and the // producer picks — this line is the only place the picked number is visible. // - // Why it matters beyond curiosity: `stream.queue_raw_buffer(newest)` at the end of this - // callback hands the buffer back while the encode thread may still be importing and - // reading its dmabuf, so content stability rests entirely on the producer not cycling - // back to this buffer before we are done with it. That window is `pool_depth` buffer - // periods wide. A pool of 2 has essentially none. + // Why it matters beyond curiosity: the depth is the budget the deferred requeue + // (`HoldBook::try_hold`) spends withholding published buffers from the producer + // while the encoder reads them. A pool of ≤ HOLD_POOL_RESERVE cannot defer at all — + // those sessions run the old contract, where a requeued buffer may be rewritten + // mid-encode and only pool depth keeps frames untorn. if let Some(depth) = ud.pool.note_frame() { tracing::info!( pool_depth = depth, high_water = ud.pool.high_water, drained, - "pipewire buffer pool negotiated — this is the producer's ACTUAL count \ - (add_buffer/remove_buffer), the window in which a buffer we handed back may \ - be rewritten while the encoder still reads it" + "pipewire buffer pool negotiated — the producer's ACTUAL count \ + (add_buffer/remove_buffer): the deferred-requeue budget, and the rewrite \ + window for any frame published without a hold" ); } // Sacrificial-mode gate (kwin.rs `create`): until the producer renegotiates to the @@ -1766,14 +1942,30 @@ pub fn pipewire_thread( return; } - consume_frame(ud, spa_buf); + consume_frame(ud, spa_buf, newest); })); // Hand `newest` back to the stream exactly once, on EVERY path — normal, corrupted-skip, - // or a caught panic in the closure above. This single requeue is what keeps the fixed - // buffer pool from draining. - // SAFETY: all reads of `spa_buf`/`newest` (update_cursor_meta, consume_frame) completed - // inside the closure above; `newest` was dequeued from this stream and not yet requeued. - unsafe { stream.queue_raw_buffer(newest) }; + // or a caught panic in the closure above — UNLESS a raw-passthrough publish withheld it + // (`try_defer` put it in the hold book): then the requeue duty belongs to the frame's + // `BufferHold`, and requeueing here too would hand the producer the same buffer twice. + // The book is stable across this check: only this thread removes entries (the requeue + // channel's callback / `remove_buffer`), and neither can run inside `.process` — a + // consumer racing the frame to its drop merely queues the release message. A panic + // AFTER the publish leaves the hold live on the published frame, so skipping the + // immediate requeue remains correct on that path too. + let withheld = ud + .defer + .book + .lock() + .map(|b| b.contains(newest as usize)) + .unwrap_or(false); + if !withheld { + // SAFETY: all reads of `spa_buf`/`newest` (update_cursor_meta, consume_frame) + // completed inside the closure above; `newest` was dequeued from this stream, + // not yet requeued, and — per the `withheld` check — carries no hold that would + // requeue it a second time. + unsafe { stream.queue_raw_buffer(newest) }; + } if outcome.is_err() { // In the per-frame `.process` callback: a deterministic panic (e.g. a bad // format) would fire this every frame, so power-of-two throttle it — enough to @@ -1789,6 +1981,34 @@ pub fn pipewire_thread( .register() .context("register stream listener")?; + // The deferred-requeue service. A `BufferHold` dropping on any thread only *sends* + // `(buffer, generation)`; this callback — on the loop thread, like every other stream op — + // is where a withheld buffer actually rejoins the producer's pool. `HoldBook::complete` + // makes a release for a renegotiated-away buffer (or a freed address reused by a new + // pool's buffer) a no-op, so a stale hold can never queue somebody else's buffer. + let defer_cb = defer.clone(); + let stream_ptr = stream.as_raw_ptr() as usize; + let _requeue_attach = requeue_rx.attach(mainloop.loop_(), move |(buf, generation)| { + let requeue = defer_cb + .book + .lock() + .map(|mut b| b.complete(buf, generation)) + .unwrap_or(false); + if requeue { + // SAFETY: `complete` returned true ⇒ this buffer was withheld by exactly this hold + // and no `remove_buffer` has freed it since (that purges the book), so the pointer + // is a live buffer of this stream that we own (dequeued, never requeued). The + // stream outlives this attached receiver (declared after it, dropped before it), + // and the loop stops dispatching once `run()` returns. + let _ = unsafe { + pw::sys::pw_stream_queue_buffer( + stream_ptr as *mut pw::sys::pw_stream, + buf as *mut pw::sys::pw_buffer, + ) + }; + } + }); + // Debug knob: offer a single fixed format (PUNKTFUNK_PW_FIXED_POD="WxH") to bisect // negotiation failures against a producer's exact EnumFormat (e.g. gamescope). let fixed_pod: Option<(u32, u32)> = std::env::var("PUNKTFUNK_PW_FIXED_POD") @@ -2479,4 +2699,77 @@ mod tests { assert_eq!(p.note_frame(), Some(0)); assert_eq!(p.high_water, 0); } + + use super::{HoldBook, HOLD_POOL_RESERVE}; + + /// The book must always leave [`HOLD_POOL_RESERVE`] buffers with the producer: an 8-pool + /// spares 6, and the pools at or below the reserve spare NOTHING — those sessions must fall + /// back to the immediate requeue rather than starve the compositor of render targets. + #[test] + fn hold_book_spends_at_most_pool_minus_reserve() { + let mut b = HoldBook::default(); + for i in 0..6 { + assert!( + b.try_hold(0x1000 + i, 8).is_some(), + "hold {i} within budget" + ); + } + assert!( + b.try_hold(0x2000, 8).is_none(), + "7th of 8 exceeds the budget" + ); + assert!( + HoldBook::default() + .try_hold(0x1000, HOLD_POOL_RESERVE) + .is_none(), + "a pool of exactly the reserve cannot spare a buffer" + ); + assert!( + HoldBook::default() + .try_hold(0x1000, HOLD_POOL_RESERVE + 1) + .is_some(), + "one past the reserve spares exactly one" + ); + } + + /// One hold ⇒ one requeue: the first `complete` releases, a duplicate release (a bug shape, + /// but also the benign stale-message case) must NOT requeue a second time — handing the + /// producer the same buffer twice corrupts its pool. + #[test] + fn hold_book_releases_exactly_once() { + let mut b = HoldBook::default(); + let g = b.try_hold(0x1000, 8).unwrap(); + assert!(b.complete(0x1000, g), "first release requeues"); + assert!(!b.complete(0x1000, g), "second release is a no-op"); + assert!(!b.contains(0x1000)); + } + + /// The renegotiation hazard the generation exists for: the pool is replaced (`remove_buffer` + /// purges), a NEW buffer lands on the SAME address and is withheld, and only then does the + /// OLD hold's release arrive. Matching by pointer alone would requeue the new tenant while + /// its own hold is still out — the mid-encode rewrite race, reintroduced by the fix itself. + #[test] + fn hold_book_generation_outlives_an_address_reuse() { + let mut b = HoldBook::default(); + let old = b.try_hold(0x1000, 8).unwrap(); + b.purge(0x1000); // remove_buffer: pool renegotiated away under the hold + assert!(!b.complete(0x1000, old), "purged hold releases nothing"); + let new = b.try_hold(0x1000, 8).unwrap(); // new pool's buffer, same address + assert!( + !b.complete(0x1000, old), + "the OLD hold cannot release the NEW tenant" + ); + assert!(b.contains(0x1000), "new tenant still withheld"); + assert!(b.complete(0x1000, new), "its own hold releases it"); + } + + /// A buffer already out cannot be withheld again (one requeue duty per buffer): `.process` + /// can only re-see an address after its requeue, so a duplicate try_hold means state + /// confusion — refuse it and let the epilogue requeue immediately. + #[test] + fn hold_book_refuses_a_buffer_already_out() { + let mut b = HoldBook::default(); + b.try_hold(0x1000, 8).unwrap(); + assert!(b.try_hold(0x1000, 8).is_none()); + } } diff --git a/crates/pf-capture/src/linux/pw_pods.rs b/crates/pf-capture/src/linux/pw_pods.rs index f22b27cc..ab852d6e 100644 --- a/crates/pf-capture/src/linux/pw_pods.rs +++ b/crates/pf-capture/src/linux/pw_pods.rs @@ -322,12 +322,14 @@ pub(super) fn build_shm_only_buffers() -> Result> { /// PW5 stage 2: the buffer-pool depth we ASK for on the zero-copy path, as a Choice range. /// -/// The zero-copy path hands the SPA buffer back to the producer at `.process` return, while the -/// encode thread still holds a dup of its dmabuf fd and has not yet imported, let alone read, the -/// contents. Nothing bounds that window — see the `queue_raw_buffer` comment in `pipewire.rs` — so -/// the only thing that keeps capture untorn is the producer round-robining a pool deeper than our -/// import+encode latency. Until PW5 stage 1 nobody had ever counted what that pool was; we never -/// even asked for a size (`build_dmabuf_buffers` set `dataType` and nothing else). +/// The raw-passthrough arm now WITHHOLDS each published buffer from the producer until the +/// consumer's hold drops (`DeferredRequeue` in `pipewire.rs` — the fix for the producer +/// rewriting a buffer mid-encode), spending up to `pool - HOLD_POOL_RESERVE` buffers of this +/// depth. A pool at the old floor of 2 has nothing to spend and falls back to the racy +/// immediate requeue, where only the producer round-robining a pool deeper than our +/// import+encode latency keeps capture untorn. Until PW5 stage 1 nobody had ever counted what +/// that pool was; we never even asked for a size (`build_dmabuf_buffers` set `dataType` and +/// nothing else). /// /// A **range**, deliberately, not a fixed count: SPA intersects the consumer's and producer's /// Buffers params, so a fixed 8 against a producer that can only afford 4 empties the intersection diff --git a/crates/pf-encode/src/enc/linux/pyrowave.rs b/crates/pf-encode/src/enc/linux/pyrowave.rs index 8e3b6c08..d382a588 100644 --- a/crates/pf-encode/src/enc/linux/pyrowave.rs +++ b/crates/pf-encode/src/enc/linux/pyrowave.rs @@ -2803,6 +2803,7 @@ mod tests { plane1: None, offset: 0, stride: 64 * 4, + hold: None, } }; let fd_count = || std::fs::read_dir("/proc/self/fd").expect("procfs").count(); diff --git a/crates/pf-encode/src/enc/linux/pyrowave_remote.rs b/crates/pf-encode/src/enc/linux/pyrowave_remote.rs index 4b79525f..7ca4eb06 100644 --- a/crates/pf-encode/src/enc/linux/pyrowave_remote.rs +++ b/crates/pf-encode/src/enc/linux/pyrowave_remote.rs @@ -936,6 +936,7 @@ mod tests { plane1: None, offset: 0, stride: 1920 * 4, + hold: None, }), cursor, } diff --git a/crates/pf-encode/src/enc/linux/vulkan_video.rs b/crates/pf-encode/src/enc/linux/vulkan_video.rs index 7545f624..36d48ad2 100644 --- a/crates/pf-encode/src/enc/linux/vulkan_video.rs +++ b/crates/pf-encode/src/enc/linux/vulkan_video.rs @@ -586,6 +586,12 @@ struct Frame { pts_ns: u64, keyframe: bool, recovery_anchor: bool, + /// The captured dmabuf's deferred-requeue hold ([`pf_frame::FrameHold`]), cloned at submit and + /// dropped when this slot retires (fence signaled — `poll`/backpressure/`reset`). This is what + /// extends "the producer must not rewrite the buffer" across the whole asynchronous GPU read: + /// the host's own clone only lives until it takes the NEXT frame, which with a ring of 2 is + /// before this slot's encode finished. `None` for non-dmabuf sources or un-held frames. + src_hold: Option, } pub struct VulkanVideoEncoder { @@ -2274,7 +2280,9 @@ impl VulkanVideoEncoder { // First import: acquire from the foreign producer (UNDEFINED preserves the modifier-tiled // bytes). Cached re-read: we still own it, so no queue-family transfer — just a visibility // barrier so the shader read sees the content the producer wrote out-of-band this frame - // (single-GPU coherent; the capture layer guarantees the buffer is ready at hand-off). + // (single-GPU coherent). The barrier orders nothing against the PRODUCER — content + // stability across this read is the frame's deferred-requeue hold (`Frame::src_hold`): + // the producer does not get the buffer back to rewrite until this slot's fence retires. let (old, src_qf, dst_qf) = if fresh { ( vk::ImageLayout::UNDEFINED, @@ -3875,11 +3883,24 @@ impl VulkanVideoEncoder { ), Err(e) => return Err(e.into()), } + // Fence signaled ⟹ the GPU is done reading this slot's captured dmabuf — release + // its hold so the capture layer requeues the producer's buffer. + self.frames[slot].src_hold = None; let done = self.read_slot(slot)?; self.pending.push_back(done); } let slot = self.ring; self.ring = (self.ring + 1) % self.frames.len(); + // Take over the frame's deferred-requeue hold for this occupancy BEFORE recording: the + // producer must not get the buffer back until this slot's fence retires (poll / + // backpressure / reset), because the encode reads the imported dmabuf for its whole + // duration — the host's own clone drops as soon as it takes the next frame. Assigned + // even if `record_submit` then fails: an over-hold until the slot's next tenant is + // harmless, a released-while-referenced buffer is the exact race this closes. + self.frames[slot].src_hold = match &frame.payload { + FramePayload::Dmabuf(d) => d.hold.clone(), + _ => None, + }; self.record_submit(slot, frame, wire)?; self.in_flight.push_back(slot); Ok(()) @@ -4002,6 +4023,9 @@ impl Encoder for VulkanVideoEncoder { Err(e) => return Err(e.into()), } self.in_flight.pop_front(); + // Fence signaled ⟹ the GPU is done reading this slot's captured dmabuf — release its + // hold so the capture layer requeues the producer's buffer. + self.frames[slot].src_hold = None; // SAFETY: fence signaled ⟹ this slot's CSC+encode is complete; read its bitstream. Ok(Some(unsafe { self.read_slot(slot)? })) } @@ -4064,6 +4088,11 @@ impl Encoder for VulkanVideoEncoder { } self.in_flight.clear(); self.pending.clear(); + // The waits above proved every slot's GPU read is done — release the captured-dmabuf + // holds so the capture layer (possibly mid-rebuild itself) gets its buffers back. + for f in &mut self.frames { + f.src_hold = None; + } self.ring = 0; self.first_frame = true; self.force_kf = false; diff --git a/crates/pf-encode/src/enc/linux/worker.rs b/crates/pf-encode/src/enc/linux/worker.rs index e31c8a51..1f4a0f66 100644 --- a/crates/pf-encode/src/enc/linux/worker.rs +++ b/crates/pf-encode/src/enc/linux/worker.rs @@ -692,6 +692,9 @@ fn encode_one( plane1: req.plane1, offset: req.offset, stride: req.stride, + // The deferred-requeue hold stays host-side: this backend is synchronous at depth 1 + // (see below), so the host's frame — hold and all — outlives the whole encode. + hold: None, }), cursor, }; diff --git a/crates/pf-frame/src/lib.rs b/crates/pf-frame/src/lib.rs index 38cc96f7..36bb3748 100644 --- a/crates/pf-frame/src/lib.rs +++ b/crates/pf-frame/src/lib.rs @@ -231,6 +231,23 @@ pub struct CapturedFrame { pub cursor: Option, } +/// Keeps the producer's buffer behind a zero-copy frame OUT of the producer's pool. +/// +/// The fd on a [`DmabufFrame`] only keeps the buffer object from being *freed*; nothing stops the +/// compositor from *re-rendering into it* once the capture layer hands the buffer back — which it +/// used to do at `.process` return, before the encoder had even imported the dmabuf (the +/// gamescope-at-120fps torn-frame race). This handle is the fix: the PipeWire capture attaches one +/// to every raw-passthrough frame (pool depth permitting) and defers the requeue until the LAST +/// clone drops. A consumer that reads the dmabuf asynchronously (the Vulkan encoder's ring) clones +/// it into whatever tracks the read (its ring slot) and drops it when the GPU is provably done +/// (the slot's fence), so content stability covers exactly the read window. Consumers that finish +/// their read while the frame is alive need to do nothing — the frame's own clone is enough. +/// +/// Opaque on purpose: the concrete guard lives in the capture crate; everyone else only clones and +/// drops. +#[cfg(target_os = "linux")] +pub type FrameHold = std::sync::Arc; + /// A captured frame still living in a DMA-BUF. Packed RGB uses one plane. Native Linux NV12 /// (gamescope PipeWire) travels in ONE fd: Y starts at `offset`, and the interleaved UV plane /// lives at `plane1`'s offset/stride when the producer reported them — else at the contiguous @@ -238,8 +255,9 @@ pub struct CapturedFrame { /// /// Owns a *dup* of the PipeWire buffer's fd, so the frame can travel to the encode thread and be /// imported there without the compositor's buffer being closed underneath it. Content stability -/// across the brief import window relies on the compositor's buffer pool depth, like any zero-copy -/// capture. +/// across the read window comes from [`hold`](Self::hold) when present (the producer does not get +/// the buffer back until the hold drops); a `None` hold falls back to the old contract — the +/// compositor's pool depth outrunning the import+encode window. #[cfg(target_os = "linux")] pub struct DmabufFrame { pub fd: std::os::fd::OwnedFd, @@ -253,6 +271,9 @@ pub struct DmabufFrame { pub plane1: Option<(u32, u32)>, pub offset: u32, pub stride: u32, + /// Deferred-requeue hold on the producer's buffer (see [`FrameHold`]); `None` when the + /// capture could not spare a buffer from the pool (shallow pool, or `PUNKTFUNK_ZEROCOPY_HOLD=0`). + pub hold: Option, } /// Where a captured frame's pixels live.