diff --git a/crates/pf-encode/src/enc/windows/nvenc.rs b/crates/pf-encode/src/enc/windows/nvenc.rs index e84227a7..cc9cec1c 100644 --- a/crates/pf-encode/src/enc/windows/nvenc.rs +++ b/crates/pf-encode/src/enc/windows/nvenc.rs @@ -48,7 +48,7 @@ use super::nvenc_core::{ CeilingKey, LowLatencyConfig, NvStatusExt, RangePlan, }; use super::nvenc_status; -use super::{ChromaFormat, Codec, EncodedFrame, Encoder, EncoderCaps}; +use super::{AuChunk, ChromaFormat, Codec, EncodedFrame, Encoder, EncoderCaps}; use anyhow::{anyhow, bail, Context, Result}; use pf_frame::{CapturedFrame, FramePayload, PixelFormat}; use std::collections::{HashMap, VecDeque}; @@ -592,7 +592,7 @@ pub struct NvencD3d11Encoder { /// in-flight encode. The fourth field tags the first frame encoded after a successful /// [`invalidate_ref_frames`](Encoder::invalidate_ref_frames) — the clean re-anchor P-frame the /// client lifts its post-loss freeze on (see [`EncodedFrame::recovery_anchor`]). - pending: VecDeque<(nv::NV_ENC_OUTPUT_PTR, nv::NV_ENC_INPUT_PTR, u64, bool)>, + pending: VecDeque<(nv::NV_ENC_OUTPUT_PTR, nv::NV_ENC_INPUT_PTR, u64, bool, bool)>, /// The frame number of the NEXT submission (also its `inputTimeStamp`). Pinned per frame by /// [`Encoder::submit_indexed`] to the WIRE frame index the AU will carry, so the DPB timestamps /// `invalidate_ref_frames` compares client frame numbers against stay 1:1 with the wire across @@ -622,6 +622,18 @@ pub struct NvencD3d11Encoder { /// there could flip `enableSubFrameWrite` mid-session, and would re-log the arbitration on /// every ABR retarget). subframe_on: bool, + /// Slice count the live session was configured with ([`resolve_slices`] over the + /// negotiated ceiling, latched by `init_session` and consumed by `build_config` so an + /// in-place reconfigure presents the same slicing). Chunked poll needs ≥ 2. + slices: u32, + /// Client-decoder slice ceiling from negotiation (see [`NvencD3d11Encoder::open`]). + max_slices: u32, + /// Sub-frame chunked poll armed for the live session (`slices ≥ 2` ∧ sub-frame write ∧ + /// SYNC retrieve — the async retrieve owns the bitstream from its thread, a doNotWait + /// sampler here would race it). + subframe_chunks: bool, + /// In-progress chunked readback of the FRONT `pending` AU (see [`ChunkState`]). + chunk: Option, session_async: bool, /// The last reference-frame range we invalidated — dedupes repeated RFI requests for the same /// loss event (the client resends until it sees recovery). @@ -655,6 +667,41 @@ pub struct NvencD3d11Encoder { // during the move — so `Send` introduces no data race on the non-`Send` fields. unsafe impl Send for NvencD3d11Encoder {} +/// Sampling cadence of the sub-frame chunked poll's doNotWait lock — mirrors the Linux twin +/// (slice completions land ~0.5-1 ms apart; 50 µs keeps the added per-chunk delivery delay +/// well under one slice time without hammering the driver). +const CHUNK_SAMPLE_INTERVAL: std::time::Duration = std::time::Duration::from_micros(50); + +/// Progress of a sub-frame chunked readback (P2f — the Windows twin of the Linux LN1 state) +/// for the FRONT in-flight AU: how much of the bitstream has already been handed out as +/// chunks. `Some` from an AU's first emitted chunk until its `last` — [`Encoder::poll`] +/// refuses to run while it exists (a plain poll would re-emit the already-shipped prefix). +struct ChunkState { + /// Bytes emitted so far — also the next chunk's start offset (always a slice boundary). + emitted: usize, + /// Completed slices already covered by emitted chunks. + slices_out: u32, + /// The AU-opening chunk (`AuChunk::first`) has been handed out. + opened: bool, + /// Debug-build shadow of every emitted byte, cross-checked against the finishing blocking + /// lock's full AU — a mis-cut chunk fails loudly in debug instead of silently corrupting + /// the wire. Compiled out of release builds. + #[cfg(debug_assertions)] + shadow: Vec, +} + +impl ChunkState { + fn new() -> Self { + ChunkState { + emitted: 0, + slices_out: 0, + opened: false, + #[cfg(debug_assertions)] + shadow: Vec::new(), + } + } +} + impl NvencD3d11Encoder { #[allow(clippy::too_many_arguments)] pub fn open( @@ -666,6 +713,10 @@ impl NvencD3d11Encoder { bitrate_bps: u64, bit_depth: u8, chroma: ChromaFormat, + // Client-decoder slice ceiling from negotiation (`VIDEO_CAP_MULTI_SLICE`, or + // GameStream's `videoEncoderSlicesPerFrame`); 1 = single-slice only — the safe shape + // toward decoders that never asked (Amlogic TV SoCs wedge on multi-slice AUs). + max_slices: u32, ) -> Result { // The runtime DLL load is the real "is NVENC possible here" gate: fail the open with a // clear reason (backend misdetect / forced PUNKTFUNK_ENCODER=nvenc on a non-NVIDIA box) @@ -706,6 +757,10 @@ impl NvencD3d11Encoder { custom_vbv: false, split_mode: nv::NV_ENC_SPLIT_ENCODE_MODE::NV_ENC_SPLIT_DISABLE_MODE as u32, subframe_on: false, + slices: 1, + max_slices: max_slices.max(1), + subframe_chunks: false, + chunk: None, session_async: false, last_rfi_range: None, init_device: ptr::null_mut(), @@ -734,7 +789,7 @@ impl NvencD3d11Encoder { while rt.done_rx.try_recv().is_ok() {} } // Unmap any in-flight inputs, then unregister every cached texture and destroy the bitstreams. - for (_, map, _, _) in &self.pending { + for (_, map, _, _, _) in &self.pending { if !map.is_null() { let _ = (api().unmap_input_resource)(self.encoder, *map); } @@ -795,6 +850,8 @@ impl NvencD3d11Encoder { self.regs.clear(); // drops the texture clones, releasing our refs self.bitstreams.clear(); self.pending.clear(); + self.chunk = None; + self.subframe_chunks = false; self.encoder = ptr::null_mut(); self.inited = false; self.next = 0; @@ -961,10 +1018,9 @@ impl NvencD3d11Encoder { av1_input_depth_minus8: if ten_bit_in { 2 } else { 0 }, hdr: self.hdr, rfi_supported: self.rfi_supported, - // Env-only on Windows (default single slice): the Phase-3 default-on is a - // Linux-direct-NVENC decision — the Windows async path stays untouched, and a - // Windows operator opting in must choose slices+sync over async retrieve. - slices: resolve_slices(self.codec, 1), + // Latched by `init_session` from the negotiated ceiling (P2f) — a later + // reconfigure re-presents the same slicing. + slices: self.slices, }, ); Ok(cfg) @@ -1091,6 +1147,11 @@ impl NvencD3d11Encoder { // The init-failure fallback below disables it if a codec/config rejects it. let pixel_rate = self.width as u64 * self.height as u64 * self.fps.max(1) as u64; let split_mode: u32 = resolve_split_mode(self.bit_depth, pixel_rate); + // Negotiated multi-slice (P2f): the direct-NVENC default of 4, clamped by the + // client's ceiling — a single-slice client keeps today's shape, a + // VIDEO_CAP_MULTI_SLICE / Moonlight slices-per-frame client gets real slices. + // `PUNKTFUNK_NVENC_SLICES` stays the operator override in both directions. + self.slices = resolve_slices(self.codec, 4.min(self.max_slices)); // Split × sub-frame arbitration (Phase 8) before the ladder/ceiling key. On Windows // sub-frame is env-opt-in only, so resolved == forced by construction. let subframe_req = resolve_subframe(false); @@ -1255,6 +1316,17 @@ impl NvencD3d11Encoder { self.split_mode = used_split; self.subframe_on = subframe_req; self.session_async = use_async; + // Sub-frame chunked poll (P2f, the Windows leg of the slice pipeline): sync + // retrieve only — chunked poll is a depth-1 sync feature; the async retrieve's + // thread owns the bitstream lock. Sub-frame write itself stays env-gated + // (`PUNKTFUNK_NVENC_SUBFRAME=1`) until the Windows on-glass A/B validates it. + self.subframe_chunks = self.slices >= 2 && subframe_req && !use_async; + if self.subframe_chunks { + tracing::info!( + slices = self.slices, + "NVENC sub-frame chunked poll armed (poll_chunk emits slice-boundary AU chunks)" + ); + } // Session-budget accounting (Stage W3): record what this open holds so admission can // decline a parallel display the hardware can't afford. Weighted by the FINAL split // mode (a split session occupies one hardware session per engine). @@ -1337,7 +1409,7 @@ impl NvencD3d11Encoder { /// error surfaces AFTER the unmap (the resource is retired either way) so the session glue's /// rebuild path starts from clean state. fn absorb_done(&mut self, done: RetrieveDone) -> Result<()> { - let Some((bs, map, pts_ns, anchor)) = self.pending.pop_front() else { + let Some((bs, map, pts_ns, anchor, _idr_hint)) = self.pending.pop_front() else { bail!("NVENC async: completion with no in-flight frame (pairing bug)"); }; if bs as usize != done.bs { @@ -1585,6 +1657,10 @@ impl Encoder for NvencD3d11Encoder { // first one encoded after the invalidation — the clean re-anchor. A simultaneous // forced IDR is itself the re-anchor, so the tag is dropped in that case. let anchor = std::mem::take(&mut self.pending_anchor) && flags == 0; + // Submit-time IDR intent: chunked poll must flag an AU's EARLY chunks before the + // driver reports `pictureType` (only the finishing lock sees it). Exact under + // P-only + infinite GOP: IDRs happen only when forced. + let idr_hint = flags != 0; let mut pic = nv::NV_ENC_PIC_PARAMS { version: nv::NV_ENC_PIC_PARAMS_VER, inputWidth: self.width, @@ -1661,6 +1737,7 @@ impl Encoder for NvencD3d11Encoder { mp.mappedResource, captured.pts_ns, anchor, + idr_hint, )); // Async: hand the in-flight encode to the retrieve thread (channel capacity = POOL ≥ // in-flight, so this send never blocks). The pending entry above pairs with its @@ -1783,6 +1860,11 @@ impl Encoder for NvencD3d11Encoder { } fn poll(&mut self) -> Result> { + // A partially-chunked AU must be finished through `poll_chunk`: its emitted prefix is + // already with the caller, so a whole-AU poll here would double-emit those bytes. + if self.chunk.is_some() { + bail!("NVENC poll() called mid-chunked-AU — drain it via poll_chunk (caller bug)"); + } // Async mode: drain whatever the retrieve thread has finished (non-blocking) and hand out // the oldest ready AU. `None` = nothing completed yet — the session loop keeps the frame // in flight and re-polls next tick, capture never blocks on the WDDM scheduling wait. @@ -1803,7 +1885,7 @@ impl Encoder for NvencD3d11Encoder { .ready .pop_front()); } - let Some((bs, map, pts_ns, anchor)) = self.pending.pop_front() else { + let Some((bs, map, pts_ns, anchor, _idr_hint)) = self.pending.pop_front() else { return Ok(None); }; // SAFETY: a non-empty `pending` implies `submit` ran, so `self.encoder` is the live session @@ -1850,6 +1932,172 @@ impl Encoder for NvencD3d11Encoder { } } + fn supports_chunked_poll(&self) -> bool { + // Dynamic on purpose: a rebuild can land in a different mode (async opt-in, slice + // fallback) and `teardown` drops the latch — a caller re-querying per AU sees the mode + // fall away instead of chunk-polling a session that can't serve it. + self.subframe_chunks && self.async_rt.is_none() + } + + fn poll_chunk(&mut self) -> Result> { + // Not a chunked session (knobs off, escalated to async): degrade to a single whole-AU + // chunk so a chunk-driven caller works against every session shape. The + // `chunk.is_none()` arm is defensive — a mid-AU state must always finish below. + if !self.supports_chunked_poll() && self.chunk.is_none() { + return Ok(self.poll()?.map(AuChunk::whole)); + } + let Some(&(bs, _, pts_ns, anchor, idr_hint)) = self.pending.front() else { + return Ok(None); + }; + // Sampling budget: if this driver branch never publishes intermediate slices, stop + // burning CPU after ~2 frame intervals and finish through the blocking lock — worst + // case poll_chunk behaves like sync `poll` plus a few failed doNotWait attempts. + let budget = std::time::Duration::from_micros(2_000_000 / self.fps.max(1) as u64); + let t0 = std::time::Instant::now(); + let mut offsets = [0u32; 32]; + loop { + let emitted = self.chunk.as_ref().map_or(0, |c| c.emitted); + let slices_out = self.chunk.as_ref().map_or(0, |c| c.slices_out); + // SAFETY: `bs` is the front `pending` entry's pool bitstream (a prior + // `encode_picture` targeted it, and `teardown` clears `pending` whenever it nulls + // the session), the session is live, and this runs on the encode thread. `lock` + // (version set, doNotWait) and `offsets` are live stack locals across the + // synchronous call; `reportSliceOffsets` was armed at init so the driver may write + // up to `numSlices` ≤ 32 offsets (`sliceModeData` is clamped 2..=32 by + // `resolve_slices`). On a successful sub-frame lock `bitstreamBufferPtr` holds + // `bitstreamSizeInBytes` readable bytes of COMPLETED slices (enableSubFrameWrite + // publishes them mid-encode) valid until the matching unlock; the emitted range is + // copied out BEFORE the unlock. Every successful lock is unlocked exactly once on + // all paths through the body. + unsafe { + let mut lock = nv::NV_ENC_LOCK_BITSTREAM { + version: nv::NV_ENC_LOCK_BITSTREAM_VER, + outputBitstream: bs, + sliceOffsets: offsets.as_mut_ptr(), + ..Default::default() + }; + lock.set_doNotWait(1); + if (api().lock_bitstream)(self.encoder, &mut lock) + .nv_ok() + .is_ok() + { + let n = lock.numSlices; + let bytes = lock.bitstreamSizeInBytes as usize; + if n >= self.slices { + // Every slice is readable — fall through to the finishing blocking + // lock (the completion authority; `numSlices` alone is not trusted + // across driver branches). + let _ = (api().unlock_bitstream)(self.encoder, bs); + break; + } + if n > slices_out && bytes > emitted { + // New completed slice(s): cut `[emitted..bytes)`. `bytes` with `n` + // reported slices is the end of slice n (slices are contiguous + // Annex-B), so the cut lands on a NAL boundary. + let data = + std::slice::from_raw_parts(lock.bitstreamBufferPtr as *const u8, bytes) + [emitted..] + .to_vec(); + (api().unlock_bitstream)(self.encoder, bs) + .nv_ok() + .map_err(|e| nvenc_status::call_err("unlock_bitstream (chunk)", e))?; + let cs = self.chunk.get_or_insert_with(ChunkState::new); + #[cfg(debug_assertions)] + cs.shadow.extend_from_slice(&data); + let first = !cs.opened; + cs.opened = true; + cs.emitted = bytes; + cs.slices_out = n; + return Ok(Some(AuChunk { + data, + pts_ns, + keyframe: idr_hint, + recovery_anchor: anchor, + chunk_aligned: false, + first, + last: false, + })); + } + let _ = (api().unlock_bitstream)(self.encoder, bs); + } + // Non-SUCCESS (LOCK_BUSY) = not ready — never an error here; the finishing + // blocking lock below owns real failures. + } + if t0.elapsed() > budget { + break; + } + std::thread::sleep(CHUNK_SAMPLE_INTERVAL); + } + + // Finish: ONE blocking lock — the completion authority, exactly like sync `poll` (so + // the final chunk blocks and the AU tail never rides a +1 tick — the depth-1 pump + // contract). Emits whatever the sampler hadn't handed out. + let (bs, map, pts_ns, anchor, idr_hint) = + self.pending.pop_front().expect("front() checked above"); + // SAFETY: same contract as `poll`'s blocking lock: `bs` is the popped in-flight pool + // bitstream on the live session (encode thread); the blocking `lock_bitstream` + // (version set) returns when the encode finished, yielding `bitstreamSizeInBytes` + // CPU-readable bytes at `bitstreamBufferPtr` valid until `unlock_bitstream` — every + // read (tail copy + debug prefix check) happens BEFORE the unlock. `map` (paired with + // `bs` in `pending`) is unmapped here, after completion, exactly once. + unsafe { + let mut lock = nv::NV_ENC_LOCK_BITSTREAM { + version: nv::NV_ENC_LOCK_BITSTREAM_VER, + outputBitstream: bs, + ..Default::default() + }; + (api().lock_bitstream)(self.encoder, &mut lock) + .nv_ok() + .map_err(|e| nvenc_status::call_err("lock_bitstream (chunk finish)", e))?; + let total = lock.bitstreamSizeInBytes as usize; + let full = std::slice::from_raw_parts(lock.bitstreamBufferPtr as *const u8, total); + let cs = self.chunk.take().unwrap_or_else(ChunkState::new); + if cs.emitted > total { + let _ = (api().unlock_bitstream)(self.encoder, bs); + bail!( + "NVENC chunked poll: {} bytes already emitted but the finished AU is only \ + {} — sub-frame readback reported bytes the final lock disowns", + cs.emitted, + total + ); + } + #[cfg(debug_assertions)] + if cs.shadow.as_slice() != &full[..cs.emitted] { + let _ = (api().unlock_bitstream)(self.encoder, bs); + bail!("NVENC chunked poll: emitted chunks diverge from the finished AU prefix"); + } + let data = full[cs.emitted..].to_vec(); + let keyframe = matches!( + lock.pictureType, + nv::NV_ENC_PIC_TYPE::NV_ENC_PIC_TYPE_IDR | nv::NV_ENC_PIC_TYPE::NV_ENC_PIC_TYPE_I + ); + (api().unlock_bitstream)(self.encoder, bs) + .nv_ok() + .map_err(|e| nvenc_status::call_err("unlock_bitstream (chunk finish)", e))?; + if !map.is_null() { + let _ = (api().unmap_input_resource)(self.encoder, map); + } + if cs.opened && keyframe != idr_hint { + // Can't happen under P-only + infinite GOP; if a driver branch ever proves + // otherwise, the earlier chunks carried the wrong flag — make it visible. + tracing::warn!( + predicted = idr_hint, + actual = keyframe, + "NVENC chunked poll: picture type diverged from the submit-time prediction" + ); + } + Ok(Some(AuChunk { + data, + pts_ns, + keyframe, + recovery_anchor: anchor, + chunk_aligned: false, + first: !cs.opened, + last: true, + })) + } + } + /// Encode-stall recovery: tear the whole session down (the same teardown a capture-device /// change uses) and let the next `submit` rebuild it lazily on the current device — the owed /// AUs are forfeited and the fresh session opens on an IDR. Gives the encode-stall watchdog a @@ -2260,6 +2508,7 @@ mod tests { 100_000_000, // high rate: the 1-px stripes must survive quantization 8, chroma, + 1, ) .expect("NVENC open"); let mut out = Vec::new(); @@ -2360,6 +2609,7 @@ mod tests { 20_000_000, 8, ChromaFormat::Yuv420, + 1, ) .expect("NVENC open"); diff --git a/crates/pf-encode/src/lib.rs b/crates/pf-encode/src/lib.rs index 81f6bfba..66c71ab3 100644 --- a/crates/pf-encode/src/lib.rs +++ b/crates/pf-encode/src/lib.rs @@ -679,6 +679,7 @@ fn open_video_backend( bitrate_bps, bit_depth, chroma, + max_slices, ) .map(|e| (Box::new(e) as Box, "nvenc")) }