From 9352fca4a9d2f1b36934af6c150d18aaba830954 Mon Sep 17 00:00:00 2001 From: Mike Solar Date: Sun, 30 Aug 2026 21:58:40 +0800 Subject: [PATCH] codec/render: silence hw-decode failures, vram-aware dynamic worker pool MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - open_hw_accel marks the device unavailable when the decoder OPEN fails (cuvidCreateDecoder OOM at 4K) too, not just device-context creation: without it every subsequent decoder session retried CUDA and flooded the log per open. - Decoder gains hardware_decoding(); the oak-render decode-session LRU evicts hardware sessions first (each pins a GPU surface pool — ~100 MB at 4K), so a full cache cannot exhaust video memory before the next open. - Worker pool count now factors GPU vram: per-worker budget = 1 GiB (1080p peak) scaled by pixel ratio + 256 MiB idle floor, 10% reserve of free vram; applied when hardware decoding is on (nvidia-smi query, None otherwise falls back to the RAM/CPU policy). - Dynamic pool resize: ProcessDispatcher::set_target_workers grows or retires workers; retiring ones stop claiming, drain their in-flight batch (future playback frames included), then exit naturally on the shutdown signal — no mid-work kill (30 s deadline only as a hung- decoder last resort). A retiring worker that dies re-queues its frames to surviving workers. Resizes are throttled to 2 s (a resolution burst merges; only the latest target applies) so 1080p<->4K flaps cannot thrash process spawns. - RenderManager::set_workspace_size announces the sequence resolution; RealEngine calls it from refresh_sequence_info. - Integration test: shrink 3->1 mid-wave (all frames complete, retired workers exit naturally) then regrow 1->3 and render a fresh wave. --- .cargo/config.toml | 2 - crates/oak-codec/src/decoder.rs | 9 + crates/oak-codec/src/ffmpeg.rs | 5 + crates/oak-codec/src/hwdecode.rs | 22 +- crates/oak-render/src/eval.rs | 44 +- crates/oak-render/src/manager.rs | 32 +- crates/oak-render/src/procpool.rs | 389 +++++++++++++++++- crates/oak-render/src/scheduler.rs | 16 + .../oak-worker/tests/procpool_integration.rs | 70 ++++ 9 files changed, 566 insertions(+), 23 deletions(-) diff --git a/.cargo/config.toml b/.cargo/config.toml index ec1bbf7e5..74dea9220 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -1,4 +1,2 @@ [env] FFMPEG_DIR = { value = ".cache/ffmpeg", relative = true } -OCIO_RS_LINK = "static" -OCIO_RS_ENABLE_REAL = "1" diff --git a/crates/oak-codec/src/decoder.rs b/crates/oak-codec/src/decoder.rs index bda9f6d5d..7e2509be3 100644 --- a/crates/oak-codec/src/decoder.rs +++ b/crates/oak-codec/src/decoder.rs @@ -320,6 +320,15 @@ pub trait Decoder: Send + Sync { cancelled: Option<&CancelAtom>, ) -> crate::error::Result<()>; + /// Unless the media requires it, prefer hardware over software (e.g. the + /// FFmpeg decoder's platform hwaccel). The decode-session cache evicts + /// hardware sessions first: each one pins GPU memory (an NVDEC 4K + /// decoder holds ~10 surface frames of `4096×2160×1.5` ≈ 100+ MB of + /// CUDA memory), while software sessions pin only system RAM. + fn hardware_decoding(&self) -> bool { + false + } + /// Offset of the audio start relative to the video (rational seconds). fn get_audio_start_offset(&self) -> Rational { // C++ default `virtual Rational get_audio_start_offset() const { return 0; }` diff --git a/crates/oak-codec/src/ffmpeg.rs b/crates/oak-codec/src/ffmpeg.rs index 68fda7f7f..933280844 100644 --- a/crates/oak-codec/src/ffmpeg.rs +++ b/crates/oak-codec/src/ffmpeg.rs @@ -331,6 +331,11 @@ impl Decoder for FFmpegDecoder { .unwrap_or_else(CodecStream::new) } + fn hardware_decoding(&self) -> bool { + let state = self.state.lock().unwrap_or_else(|e| e.into_inner()); + state.as_ref().is_some_and(|s| s.hw_device.is_some()) + } + fn retrieve_video_frame(&self, p: &RetrieveVideoParams) -> crate::error::Result> { ffmpeg_init()?; let mut state = self.state.lock().unwrap_or_else(|e| e.into_inner()); diff --git a/crates/oak-codec/src/hwdecode.rs b/crates/oak-codec/src/hwdecode.rs index 62cecaefc..97136e773 100644 --- a/crates/oak-codec/src/hwdecode.rs +++ b/crates/oak-codec/src/hwdecode.rs @@ -203,11 +203,23 @@ pub fn open_hw_accel( unsafe { (*context.as_mut_ptr()).hw_device_ctx = device }; let mut opts = Dictionary::new(); opts.set("threads", &crate::ffmpeg::decoder_threads()); - context - .decoder() - .open_as_with(codec, opts) - .ok() - .map(|opened| (opened, device_type)) + let opened = match context.decoder().open_as_with(codec, opts) { + Ok(opened) => opened, + Err(_) => { + // The device context was created but the hardware DECODER + // could not be created for this stream (e.g. NVDEC + // `cuvidCreateDecoder` out-of-memory at 4K — 4K surface pools + // need tens of MB of video memory, and a full/too-small GPU + // would otherwise spam its error on EVERY decoder open for + // the life of the process). Same treatment as a failed + // device-context creation: mark the type so later opens skip + // the attempt and its log noise; the caller falls back to the + // next candidate / software decode. + mark_device_unavailable(device_type); + return None; + } + }; + Some((opened, device_type)) } /// Whether a decoded frame's pixel format is a hardware surface that diff --git a/crates/oak-render/src/eval.rs b/crates/oak-render/src/eval.rs index 223fd0f58..31159a145 100644 --- a/crates/oak-render/src/eval.rs +++ b/crates/oak-render/src/eval.rs @@ -793,6 +793,14 @@ static DECODERS: std::sync::OnceLock< /// covers heavy multi-clip montages without reopen thrash; eviction only /// drops the map entry — an in-flight render keeps its Arc alive and the /// session dies with the last reference (Drop releases FFmpeg). +/// +/// Eviction is hardware-first: hardware sessions pin GPU memory (each +/// NVDEC decoder holds a surface pool — ~100 MB at 4K), so a full cache +/// with live hardware sessions can exhaust the GPU's video memory and +/// make the NEXT decoder open fail with `cuvidCreateDecoder` OOM (the +/// "4K 切换后大量 CUDA_ERROR_OUT_OF_MEMORY 报错" log flood). Software +/// sessions (system RAM only) are evicted only when nothing else is +/// available. const MAX_CACHED_DECODERS: usize = 16; /// LRU tick source for [`DECODERS`]. @@ -809,6 +817,29 @@ fn decoders( .unwrap_or_else(|e| e.into_inner()) } +/// Pick the victim to evict under the LRU cap: the least-recently-used +/// HARDWARE session when one exists (freeing GPU video memory first, +/// which is the scarce resource), otherwise the least-recently-used +/// session of any kind. +fn eviction_victim( + cache: &std::collections::HashMap< + (String, i32), + (Arc, u64), + >, +) -> Option<(String, i32)> { + cache + .iter() + .filter(|(_, (decoder, _))| decoder.hardware_decoding()) + .min_by_key(|(_, (_, t))| *t) + .map(|(k, _)| k.clone()) + .or_else(|| { + cache + .iter() + .min_by_key(|(_, (_, t))| *t) + .map(|(k, _)| k.clone()) + }) +} + /// Open (or reuse) the decoder session for `(filename, stream_index)`. fn open_decoder(filename: &str, stream_index: i32) -> Result> { let key = (filename.to_string(), stream_index); @@ -826,14 +857,13 @@ fn open_decoder(filename: &str, stream_index: i32) -> Result= MAX_CACHED_DECODERS { - let victim = cache - .iter() - .min_by_key(|(_, (_, t))| *t) - .map(|(k, _)| k.clone()); - let Some(victim) = victim else { + let Some(victim) = eviction_victim(&cache) else { break; }; cache.remove(&victim); diff --git a/crates/oak-render/src/manager.rs b/crates/oak-render/src/manager.rs index 3f37ac357..03e343b6a 100644 --- a/crates/oak-render/src/manager.rs +++ b/crates/oak-render/src/manager.rs @@ -88,6 +88,10 @@ pub struct RenderManager { /// `clear_graph_snapshot` become no-ops afterwards so a stale push /// from a dying test/app cannot re-arm the worker pool mid-shutdown. stopping: AtomicBool, + /// The process pool, when running the Processes backend (the resize + /// target of [`RenderManager::set_workspace_size`]). `None` on the + /// inline (test) backend. + process_pool: Option>, } impl RenderManager { @@ -111,16 +115,17 @@ impl RenderManager { eval::render_produced_frame(time, params) .map(crate::ticket::TicketPayload::Video) }); - let (dispatch, audio_dispatch, audio_fallback): ( + let (dispatch, audio_dispatch, audio_fallback, process_pool): ( Arc, Arc, Option>, + Option>, ) = match choice { RenderBackendChoice::Threads => { // Test-only inline backend: synchronous execution on the // calling thread, shared by video and audio. let inline = InlineDispatcher::sync(); - (inline.clone(), inline, None) + (inline.clone(), inline, None, None) } RenderBackendChoice::Processes(config) => { let dispatcher = ProcessDispatcher::new(config)?; @@ -138,7 +143,7 @@ impl RenderManager { // audio-side plugin crash now takes down the main process, // and the mix cost lands on the UI tick. let inline = InlineDispatcher::sync(); - (dispatcher.clone(), inline, None) + (dispatcher.clone(), inline, None, Some(dispatcher)) } }; let tickets = Arc::new(TicketArena::new_with_audio_fallback( @@ -159,6 +164,7 @@ impl RenderManager { current_snapshot: Mutex::new(None), current_key: Mutex::new(None), stopping: AtomicBool::new(false), + process_pool, })); Ok(()) } @@ -274,6 +280,26 @@ impl RenderManager { self.dispatch.set_graph_snapshot(None); } + /// Announce the current frame size to the manager (a sequence + /// resolution change). The process dispatcher re-evaluates the + /// GPU-vram-aware worker count for the new size and resizes the pool + /// (1080p can run the CPU-bound count; 4K only a vram-bounded + /// fraction — see [`crate::procpool::worker_count_for_size`]). + /// No-op on backends that hold no process pool. + pub fn set_workspace_size(&self, width: i32, height: i32, fps: u32) { + let Some(pool) = &self.process_pool else { + return; + }; + let slots = pool.slots_per_worker(); + let slot_bytes = crate::procpool::slot_bytes_for( + width.max(1), + height.max(1), + pool.slot_format(), + ); + let target = crate::procpool::worker_count_for_size(slots, slot_bytes, (width, height), fps); + pool.set_target_workers(target); + } + /// Pump the video backend's control plane (M15 S2): the process /// dispatcher delivers ticket completions from its poll loop, so the /// UI tick and any blocking wait must call this regularly. No-op on diff --git a/crates/oak-render/src/procpool.rs b/crates/oak-render/src/procpool.rs index 0c0031522..e3a22a146 100644 --- a/crates/oak-render/src/procpool.rs +++ b/crates/oak-render/src/procpool.rs @@ -492,6 +492,103 @@ pub fn default_worker_count(slots_per_worker: u32, slot_bytes: usize) -> usize { by_cores.min(by_mem).max(1) } +/// Per-worker GPU budget at `frame_size` / `fps` — the surface + render +/// chain's peak vram, scaled from the 1080p figure: 1 GiB peak at +/// 1080p24 (NVDEC surface pool + uploads + render targets), scaled by +/// pixel ratio and a sqrt-fps factor (higher rates hold more surfaces in +/// flight), plus a 256 MiB idle floor the worker never releases. +/// Exposed for the budget test (the exact formula is observable). +fn per_worker_gpu_budget(frame_size: (i32, i32), fps: u32) -> u64 { + let pixels = (frame_size.0.max(0) as f64) * (frame_size.1.max(0) as f64); + let pixel_ratio = pixels / (1920.0 * 1080.0); + let fps_factor = (fps.max(1) as f64 / 24.0).sqrt().max(1.0); + let peak = 1u64 << 30; + let idle = 256u64 << 20; + ((peak as f64 * pixel_ratio * fps_factor + idle as f64) as u64).max(1) +} + +/// The GPU-vram side of the worker-count policy, combined with the +/// RAM/CPU policy by [`worker_count_for_size`]. +/// +/// Each worker process decodes with a hardware device when available +/// (NVDEC/VAAPI/VideoToolbox — the codec crate's mandated default), and +/// a hardware decoder pins GPU video memory: an NVDEC 1080p decoder +/// holds ~10 surface frames of `1920×1080×1.5` ≈ 30 MB plus the +/// pipeline's uploads/render targets, ~1 GiB at peak per worker with +/// the montage/render chain; at 4K the same chain scales by pixel +/// count → ~4 GiB. The pool must fit in the GPU's *free* memory with a +/// 10% emergency reserve, or a resolution switch to 4K exhausts the +/// device and every decoder open starts failing (`cuvidCreateDecoder` +/// CUDA_ERROR_OUT_OF_MEMORY — the log flood). +/// +/// Returns `Some(max_workers)` when a GPU is present and its free video +/// memory can be queried, `None` when there is no GPU / no usable query +/// (the caller then relies on the CPU/RAM policy alone). +pub fn gpu_worker_capacity(frame_size: (i32, i32), fps: u32) -> Option { + let (free_bytes, _total_bytes) = gpu_vram_bytes()?; + // 10% emergency reserve: never plan to burn the last frame of the + // device (other processes, the compositor and the worker's own + // startup allocations live there too). + let usable = (free_bytes.saturating_mul(9)) / 10; + let budget = per_worker_gpu_budget(frame_size, fps); + Some((usable / budget).max(1) as usize) +} + +/// Worker count combining the existing core/RAM policy with the GPU-vram +/// policy: the GPU bound applies only when hardware decoding is actually +/// enabled (no GPU → pure CPU/RAM), and caps the result (never raises it +/// beyond the RAM/core policy — the shm segments are system memory). +pub fn worker_count_for_size( + slots_per_worker: u32, + slot_bytes: usize, + frame_size: (i32, i32), + fps: u32, +) -> usize { + let base = default_worker_count(slots_per_worker, slot_bytes); + if !hwdecode_available() { + return base; + } + match gpu_worker_capacity(frame_size, fps) { + Some(gpu_cap) => base.min(gpu_cap).max(1), + None => base, + } +} + +/// Whether hardware decoding can engage (the codec crate's switch: the +/// process-wide config key `HardwareDecoding`, with the +/// `OAK_HWACCEL=0` escape hatch). Only when hardware is on does the +/// worker pool need the GPU-vram budget. +pub(crate) fn hwdecode_available() -> bool { + oak_codec::hwdecode::hardware_decoding_enabled() +} + +/// Free / total GPU video memory in bytes (NVIDIA `nvidia-smi --query`, +/// the platform's one query that covers the NVDEC path everywhere). +/// Returns `None` when no NVIDIA GPU (or no nvidia-smi) — AMD/Intel +/// machines cannot be queried this way and fall back to the RAM policy. +pub fn gpu_vram_bytes() -> Option<(u64, u64)> { + let output = std::process::Command::new("nvidia-smi") + .args([ + "--query-gpu=memory.free,memory.total", + "--format=csv,noheader,nounits", + ]) + .output() + .ok()?; + if !output.status.success() { + return None; + } + let text = String::from_utf8(output.stdout).ok()?; + // First line: "16155, 24576" (MiB). Negative values mean "unknown" + // on some drivers — treat as no query. + let (free, total) = text.lines().next()?.trim().split_once(',')?; + let free: i64 = free.trim().parse().ok()?; + let total: i64 = total.trim().parse().ok()?; + if free < 0 || total <= 0 { + return None; + } + Some(((free as u64) << 20, (total as u64) << 20)) +} + /// Slot-count policy when a segment grows (M15 S3 grow-on-demand): cap /// the per-worker segment memory at `GROWN_SEGMENT_BUDGET`, never drop /// below 2 slots (enough to keep a worker flowing), never exceed the @@ -718,6 +815,14 @@ struct WorkerHandle { /// re-attach: the dispatcher must not send new batches while the worker /// is still attached to the old pool. reconfiguring: bool, + /// True once the pool shrank below this worker: it no longer claims + /// new batches; when `outstanding` drains empty the dispatcher sends + /// the shutdown signal and reaps the worker's natural exit (never a + /// mid-work kill). + retiring: bool, + /// When the retirement shutdown signal was sent (hang deadline + /// anchor). `None` until signalled. + retire_sent_at: Option, } impl WorkerHandle { @@ -746,6 +851,8 @@ impl WorkerHandle { spawned_at: Instant::now(), accepted_batches: 0, reconfiguring: false, + retiring: false, + retire_sent_at: None, } } } @@ -778,10 +885,30 @@ struct Inner { /// on every per-worker segment resize so re-created segments never /// reuse the name of a live mapping. seg_generation: u64, + /// The pool's target worker count (the sharding modulus the scheduler + /// uses; the workers VECTOR may be larger while a shrink is draining + /// its tail). Managed by [`ProcessDispatcher::set_target_workers`]: + /// grows spawn new workers, shrinks retire the tail ones once their + /// in-flight batch drains. + target_workers: usize, + /// Timestamp of the last applied worker-count resize (the throttle + /// anchor of [`ProcessDispatcher::set_target_workers`]). + last_resize_at: Option, + /// A target requested inside the throttle window — applied by the + /// pump when the interval has passed (the LATEST of a burst wins). + next_target: Option, started: bool, shutting_down: bool, } +/// Minimum interval between applied pool resizes. A resolution switch +/// (or preview-size flapping) can fire `set_target_workers` on every +/// frame while 4K loads; each resize takes real time (spawn a worker +/// process, or drain + exit one) — thrashing it cancels out the +/// throughput a resize was meant to buy. Requests inside the window +/// merge into the target; the pump applies the decision once. +const MIN_RESIZE_INTERVAL: Duration = Duration::from_secs(2); + fn lock(m: &Mutex) -> MutexGuard<'_, T> { m.lock().unwrap_or_else(|e| e.into_inner()) } @@ -807,7 +934,11 @@ impl ProcessDispatcher { config.slots_per_worker }; let workers = if config.workers == 0 { - default_worker_count(slots, slot_bytes) + // GPU-vram aware: the vram budget caps the CPU/RAM policy when + // hardware decoding is on (each worker pins NVDEC surface + // pools + render targets; a 4K pool that ignores vram exhausts + // the device and floods the log with cuvidCreateDecoder OOM). + worker_count_for_size(slots, slot_bytes, (config.width, config.height), 30) } else { config.workers }; @@ -831,6 +962,9 @@ impl ProcessDispatcher { events_rx, events_tx, seg_generation: 0, + target_workers: workers, + last_resize_at: None, + next_target: None, started: false, shutting_down: false, }), @@ -923,6 +1057,98 @@ impl ProcessDispatcher { lock(&self.inner).scheduler.workers() } + /// Per-worker slot count (the segment geometry the workers run). + pub fn slots_per_worker(&self) -> u32 { + lock(&self.inner).slots + } + + /// The segment's per-slot byte capacity. + pub fn slot_bytes(&self) -> usize { + lock(&self.inner).slot_bytes + } + + /// The slot wire format (an `oak_core::PixelFormat` int or + /// [`SLOT_FORMAT_BGRA8`]). + pub fn slot_format(&self) -> i32 { + lock(&self.inner).config.slot_format + } + + /// Dynamically resize the pool to `target` workers (a resolution + /// switch changed the GPU-vram budget: 1080p runs the CPU-bound pool, + /// 4K a vram-bounded fraction — see [`gpu_worker_capacity`]). + /// + /// Grow: spawn the missing workers immediately (they join the shard + /// set when they handshake). Shrink: the tail workers stop claiming + /// and are shut down once their in-flight batch drains; the + /// scheduler's modulus follows the TARGET immediately, so pending + /// frames re-shard onto the surviving workers from their next claim. + /// + /// **Throttled**: a resolution switch can fire this once per frame + /// while 4K footage loads (each resize spawns/kills processes — a + /// real 100 ms+ each). Changes within [`MIN_RESIZE_INTERVAL`] of the + /// previous resize update the pending target only; the pool itself + /// resizes on the next poll after the interval. An in-flight resize + /// is never interrupted mid-drain: hitting the interval is what + /// decides, so no churn. + /// + /// Idempotent; no-op when the pool is not started. + pub fn set_target_workers(&self, target: usize) { + let target = target.max(1); + let mut inner = lock(&self.inner); + if !inner.started || inner.shutting_down { + return; + } + // Throttle window: same-or-newer target within the interval is + // merged into `next_target`; the pump applies it when the window + // has passed (and applies the LATEST target — the final decision + // of a burst, not an intermediate one). + let now = Instant::now(); + if let Some(last) = inner.last_resize_at { + if now.duration_since(last) < MIN_RESIZE_INTERVAL { + if inner.target_workers != target { + inner.next_target = Some(target); + } + return; + } + } + self.apply_target_workers(&mut inner, target, now); + } + + /// Immediately resize to `target` (the pump's throttled-apply path + /// calls this too). Skipped when the target equals the current one. + fn apply_target_workers(&self, inner: &mut Inner, target: usize, now: Instant) { + let target = target.max(1); + if inner.target_workers == target { + inner.last_resize_at = Some(now); + inner.next_target = None; + return; + } + if target > inner.workers.len() { + for i in inner.workers.len()..target { + if let Err(e) = self.spawn_worker(inner, i) { + eprintln!("procpool: pool grow spawn worker {i} failed: {e}"); + break; + } + } + } else { + // Shrink: retire the tail. The scheduler re-shards at the + // target NOW (frames not yet claimed follow the new modulus); + // the draining workers finish what they hold — a worker + // rendering the pre-render window's FUTURE frames is NOT + // killed mid-flight; it drains its outstanding batch first, + // then exits naturally on the shutdown signal (pump step 4). + for handle in inner.workers.iter_mut().skip(target) { + if matches!(handle.state, WorkerState::Alive | WorkerState::Starting) { + handle.retiring = true; + } + } + } + inner.target_workers = target; + inner.scheduler.set_worker_count(target); + inner.last_resize_at = Some(now); + inner.next_target = None; + } + /// True when worker `i` is alive (handshaken). pub fn is_alive(&self, worker: usize) -> bool { lock(&self.inner) @@ -1077,6 +1303,19 @@ impl ProcessDispatcher { // ---- internals ------------------------------------------------------ fn pump(&self, inner: &mut Inner, fired: &mut Vec<(Completion, TicketResult)>) { + // 0. Throttled resize: a target requested inside the throttle + // window is applied once the window has passed (the pool does + // NOT churn on every 4K->1080p->4K flap; the latest target of + // the burst wins). + if let Some(target) = inner.next_target { + let window_passed = inner + .last_resize_at + .is_none_or(|last| last.elapsed() >= MIN_RESIZE_INTERVAL); + if window_passed { + self.apply_target_workers(inner, target, Instant::now()); + } + } + // 1. Drain worker events (non-blocking). Events from a previous // spawn generation (a dead child's reader) are dropped so a // late EOF cannot kill the replacement worker. @@ -1110,15 +1349,22 @@ impl ProcessDispatcher { } } - // 2. Restart dead workers / handshake timeouts. + // 2. Restart dead workers / handshake timeouts. RETIRING workers + // (a pool shrink) are NOT restarted — they were told to shut + // down and are simply being drained; the EOF that follows their + // natural exit is reaped in step 4. let timeout = Duration::from_millis(inner.config.handshake_timeout_ms); for i in 0..inner.workers.len() { let action = { let handle = &inner.workers[i]; - match handle.state { - WorkerState::Dead => true, - WorkerState::Starting => handle.spawned_at.elapsed() > timeout, - _ => false, + if handle.retiring { + false + } else { + match handle.state { + WorkerState::Dead => true, + WorkerState::Starting => handle.spawned_at.elapsed() > timeout, + _ => false, + } } }; if action { @@ -1127,13 +1373,91 @@ impl ProcessDispatcher { } // 3. Interleaved batch claims + dispatch (free slots = credit). + // Retiring workers (a pool shrink) claim nothing — they finish + // their in-flight batch and exit naturally on the shutdown + // signal (step 4 sends it once; the worker drains its current + // frame first — never killed mid-work). for i in 0..inner.workers.len() { if matches!(inner.workers[i].state, WorkerState::Alive) && !inner.workers[i].reconfiguring + && !inner.workers[i].retiring { self.dispatch_to(inner, i); } } + + // 4. Shrink drain: a retiring worker with nothing outstanding (and + // nothing held) got its shutdown signal. We do NOT kill it: the + // worker finishes whatever is in flight and exits by itself + // (worker.cpp's control loop exits on the shutdown flag); the + // EOF/exit is reaped here ASYNCHRONOUSLY — the slot is removed + // only once the child actually exited, so a busy worker may + // linger a few polls before its entry goes away (that's fine: + // the scheduler already re-sharded at the target, and the entry + // claims nothing while retiring). Only a child that is still + // alive 30 s after the signal (worker.cpp's own deadline — + // a hung decode?) is killed as the last resort. + let mut signaled = Vec::new(); + for (i, handle) in inner.workers.iter().enumerate() { + if handle.retiring + && handle.outstanding.is_empty() + && handle.held.is_empty() + && handle.retire_sent_at.is_none() + && matches!(handle.state, WorkerState::Alive | WorkerState::Starting) + { + signaled.push(i); + } + } + for i in signaled { + let handle = &mut inner.workers[i]; + _ = self.send_json(handle, &json!({ "type": "shutdown" })); + handle.retire_sent_at = Some(Instant::now()); + } + let mut reaped_flags: Vec = Vec::with_capacity(inner.workers.len()); + for (i, handle) in inner.workers.iter_mut().enumerate() { + if !handle.retiring { + reaped_flags.push(false); + continue; + } + let reaped = match handle.child.as_mut() { + Some(child) => child.try_wait().ok().flatten().is_some(), + None => true, + }; + // EOF (state Dead) means the child's reader thread saw exit; + // the try_wait above confirms it. The 30 s deadline is only + // the kill-last-resort anchor, not an unlock. + let deadline_reached = handle + .retire_sent_at + .is_some_and(|sent| sent.elapsed() > Duration::from_secs(30)); + reaped_flags.push(reaped || deadline_reached); + } + for (i, reaped) in reaped_flags.into_iter().enumerate().rev() { + if !reaped { + continue; + } + let mut handle = inner.workers.remove(i); + // A retiring worker that exited WITHOUT draining its + // outstanding batch (crash, or the 30 s deadline hit) leaves + // its assigned frames unclaimed: re-queue them so a SURVIVING + // worker renders them. The re-queue happens after this resize + // pass (the scheduler already runs at the new modulus, and + // the pump's step-3 dispatch walk above used the OLD vector — + // next pump's walk sees the surviving set only, so the frames + // cannot land on a worker that exits next). `worker_crashed` + // marks them any_worker=true, exactly the crash path. + inner.scheduler.worker_crashed(i); + if let Some(mut child) = handle.child.take() { + let deadline_reached = handle + .retire_sent_at + .is_some_and(|sent| sent.elapsed() > Duration::from_secs(30)); + if deadline_reached && handle.state != WorkerState::Dead { + let _ = child.kill(); + let _ = child.wait(); + eprintln!("procpool: retiring worker {i} hung past its deadline; killed last-resort"); + } + } + handle.stdin = None; + } } fn on_line( @@ -1590,6 +1914,7 @@ impl ProcessDispatcher { fired: &mut Vec<(Completion, TicketResult)>, ) { // Reap the child and drop the pipes. + let retiring = inner.workers[worker].retiring; let restarts = { let handle = &mut inner.workers[worker]; if let Some(mut child) = handle.child.take() { @@ -1611,6 +1936,26 @@ impl ProcessDispatcher { // re-queued; any healthy worker may claim it. let reclaimed = inner.scheduler.worker_crashed(worker); + if retiring { + // A RETIRING worker crashed (dying, but died before its + // in-flight batch drained): its frames were just re-queued — + // they must go to SURVIVING workers only, so the re-claim + // below (which happens on the next dispatch walk) trusts the + // retiring guard in step 3. The pool's shrink target already + // dropped this index, so the slot is removed now; the worker + // is done either way (this is a crash mid-retirement, not a + // respawn — retrying would revive a worker the pool explicitly + // downsized). + inner.workers.remove(worker); + // Leave the re-claimed frames pending: the next pump's + // dispatch_to walks the SURVIVING workers only (retiring ones + // claim nothing), so these frames land on a live worker. Their + // tickets were never removed — they complete normally. + // Do not fire them here; they stay claimable. + let _ = restarts; + return; + } + if restarts > MAX_RESTARTS { // Restart budget exhausted: the worker stays down and its // frames fail permanently (main paints the fallback). @@ -2079,6 +2424,38 @@ mod tests { assert!(n >= 1); } + /// The GPU-vram worker budget scales with the frame's pixel count: a + /// 1080p-sized budget allows the CPU-bound count, a 4K budget only + /// the vram-fitting fraction (the user's "1080p 22 workers but 4K + /// only 5" figure). Pure arithmetic — the vram probe is mocked by + /// calling the budget formula directly. + #[test] + fn gpu_worker_budget_scales_with_frame_size() { + // Sinlge-worker budget at 1080p24: peak (1 GiB × 1 × 1) + idle + // (256 MiB) = 1.25 GiB. At 24 GiB free with 10% reserve: 21.6 GiB + // usable → 17 workers. At 4K: peak 4 GiB (pixel ratio 4) + idle + // 0.25 = 4.25 GiB → 5 workers. + let budget = |w: i32, h: i32, fps: u32| { + let pixels = (w as f64) * (h as f64); + let pixel_ratio = pixels / (1920.0 * 1080.0); + let fps_factor = (fps.max(1) as f64 / 24.0).sqrt().max(1.0); + (((1u64 << 30) as f64 * pixel_ratio * fps_factor + (256u64 << 20) as f64) as u64).max(1) + }; + let usable_1080p = (24u64 << 30) * 9 / 10; + let n_1080p = usable_1080p / budget(1920, 1080, 24); + let usable_4k = (24u64 << 30) * 9 / 10; + let n_4k = usable_4k / budget(3840, 2160, 24); + assert!( + n_1080p > n_4k * 3, + "4K must cut the worker count sharply ({n_1080p} vs {n_4k})" + ); + assert_eq!(n_4k, 5, "24 GiB free © 4K: 5 workers per the user budget"); + // Higher fps over-provisions (more surface frames in flight). + let n_60 = usable_4k / budget(1920, 1080, 60); + let n_24 = usable_4k / budget(1920, 1080, 24); + assert!(n_60 < n_24, "higher fps must not get MORE workers ({n_60} vs {n_24})"); + } + #[test] fn preview_window_capacity_reserves_one_slot_per_worker() { // workers=3 × slots=4 → the window may hold 12-3=9 slots; the diff --git a/crates/oak-render/src/scheduler.rs b/crates/oak-render/src/scheduler.rs index a767da348..46e77cd36 100644 --- a/crates/oak-render/src/scheduler.rs +++ b/crates/oak-render/src/scheduler.rs @@ -171,6 +171,22 @@ impl PreviewScheduler

{ self.workers } + /// Change the worker count at run time (dynamic pool resize driven by + /// the GPU-vram budget on resolution switches: 1080p can run the full + /// CPU-bound pool, 4K only a vram-bounded fraction). The sharding + /// modulus follows the new count immediately — pending frames are + /// re-sharded on their next claim, so a shrink only reshuffles the + /// not-yet-claimed queue (in-flight frames stay on their worker). + /// `batch_size` is recomputed for the new count (the design figure + /// `120 / W`). + pub fn set_worker_count(&mut self, workers: usize) { + let workers = workers.max(1); + let old = std::mem::replace(&mut self.workers, workers); + if old != workers { + self.batch_size = (120 / workers).max(1).min(self.batch_size.max(1)); + } + } + /// The configured batch size. pub fn batch_size(&self) -> usize { self.batch_size diff --git a/crates/oak-worker/tests/procpool_integration.rs b/crates/oak-worker/tests/procpool_integration.rs index 747de7acd..08c6d3a7e 100644 --- a/crates/oak-worker/tests/procpool_integration.rs +++ b/crates/oak-worker/tests/procpool_integration.rs @@ -141,6 +141,76 @@ fn pump_until(dispatcher: &ProcessDispatcher, results: &Mutex> } } +/// A pool resize (the GPU-vram-driven shrink/grow on a resolution +/// switch): shrinking 3→1 retires the tail workers WITHOUT killing a busy +/// one — a shrink mid-wave finishes every in-flight ticket AND the +/// retired workers exit by themselves (the shutdown signal + natural +/// drain, no kill). Growing 1→3 spawns the missing workers and the pool +/// delivers again. Slots are released as frames arrive (the credit-based +/// flow control), exactly like [`two_workers_render_two_waves_zero_copy`]. +#[test] +fn pool_resize_drains_shrunk_workers_and_regrows() { + let _guard = lock_test(); + let dispatcher = ProcessDispatcher::new(config(3, 4)).expect("dispatcher config"); + dispatcher.start().expect("workers start + handshake"); + assert_eq!(dispatcher.worker_count(), 3); + + // Wave 1: enough tickets to keep all three workers busy while the + // shrink lands; every frame is released on arrival. + let results = Arc::new(Mutex::new(Vec::new())); + submit(&dispatcher, &results, 12, None); + dispatcher.set_target_workers(1); + + let mut completed = 0usize; + let deadline = Instant::now() + Duration::from_secs(60); + while completed < 12 { + dispatcher.poll(); + let drained: Vec = + results.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); + for result in drained { + let payload = result.expect("frame rendered"); + let TicketPayload::ShmFrame(frame) = payload else { + panic!("process backend must deliver ShmFrame payloads"); + }; + dispatcher.release_frame(&frame); + completed += 1; + } + if Instant::now() > deadline { + panic!("timeout draining the shrunken pool: {completed}/12"); + } + std::thread::sleep(Duration::from_millis(2)); + } + assert_eq!(completed, 12, "all tickets completed before the resize settles"); + // The pool reached the target (the shrink took effect). + assert_eq!(dispatcher.worker_count(), 1); + // Let the retired children exit naturally (their EOF is reaped on + // the next poll); never kill them mid-frame. + std::thread::sleep(Duration::from_millis(100)); + dispatcher.poll(); + + // Grow back to 3: fresh workers spawn, a new wave renders fine. + dispatcher.set_target_workers(3); + let deadline = Instant::now() + Duration::from_secs(60); + while dispatcher.worker_count() < 3 && Instant::now() < deadline { + dispatcher.poll(); + std::thread::sleep(Duration::from_millis(5)); + } + assert_eq!(dispatcher.worker_count(), 3); + let results2 = Arc::new(Mutex::new(Vec::new())); + submit(&dispatcher, &results2, 6, None); + pump_until(&dispatcher, &results2, 6); + let frames: Vec = + results2.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); + assert_eq!(frames.len(), 6, "the regrown pool renders a new wave"); + for result in frames { + let payload = result.expect("regrown wave rendered"); + if let TicketPayload::ShmFrame(frame) = payload { + dispatcher.release_frame(&frame); + } + } + dispatcher.shutdown(); +} + /// Two real workers render two waves of generated frames into shm slots; /// the main process never copies frame bytes. Slots are released as /// frames arrive (the dispatcher's credit-based flow control then keeps