diff --git a/crates/oak-app/src/oakui/real.rs b/crates/oak-app/src/oakui/real.rs index c2f38140f..e4e12e4db 100644 --- a/crates/oak-app/src/oakui/real.rs +++ b/crates/oak-app/src/oakui/real.rs @@ -271,7 +271,30 @@ struct PreviewWindow { /// Rendered frames held in shm slots, keyed by frame number. A frame /// is consumed by `cpu_frame` (slot released after the display image /// is built) or released when it falls out of the window. - slots: BTreeMap, + slots: BTreeMap, +} + +/// A finished frame in the playback pre-render window: a process-backend +/// shm slot (released back to its worker) or an in-process pipeline +/// texture (M4; dropped, the engine texture lease frees it). +enum PreviewSlot { + /// Process backend: a frame in a worker's shm slot. + Shm(ShmFrameRef), + /// Thread pipeline: an engine texture (GPU-resident on a shared + /// device, or an in-process CPU frame). + Video(oak_core::texture::Texture), +} + +impl PreviewSlot { + /// Release the slot's backend resource (no-op for pipeline textures, + /// whose `Drop` returns the GPU token / pixels). + fn release(self) { + if let PreviewSlot::Shm(frame) = self { + if let Some(m) = RenderManager::global() { + m.release_frame(&frame); + } + } + } } // --------------------------------------------------------------------------- @@ -2207,7 +2230,7 @@ impl RealEngine { // guard is dropped (calling them here self-deadlocks the UI thread). let mut stale_sequences: Vec = Vec::new(); let mut stale_keys: Vec<(u64, i64, u64)> = Vec::new(); - let mut stale_slots: Vec = Vec::new(); + let mut stale_slots: Vec = Vec::new(); if window.sequence != node_id || window.generation != self.preview_generation { stale_sequences.push(window.sequence); stale_slots.extend(std::mem::take(&mut window.slots).into_values()); @@ -2264,8 +2287,8 @@ impl RealEngine { for (sequence, frame, version) in stale_keys { m.dispatch.cancel_preview_frame(sequence, frame, version); } - for slot in &stale_slots { - m.release_frame(slot); + for slot in stale_slots { + slot.release(); } for frame in new_frames { @@ -2293,25 +2316,35 @@ impl RealEngine { let window = windows.entry(monitor).or_default(); if window.sequence == node_id && window.generation == version { match result { + // The rendered frame is cached (shm slot for the + // process pool, engine texture for the M4 thread + // pipeline) until the playhead reaches it + // (cpu_frame) or it falls out of the window. Ok(oak_render::ticket::TicketPayload::ShmFrame(slot)) => { - // The rendered frame is cached in its shm slot - // until the playhead reaches it (cpu_frame) - // or it falls out of the window. - window.slots.insert(frame, slot); + window.slots.insert(frame, PreviewSlot::Shm(slot)); + } + Ok(oak_render::ticket::TicketPayload::Video(texture)) => { + window.slots.insert(frame, PreviewSlot::Video(texture)); } _ => { // Render failed / cancelled: allow a re-request. window.submitted.remove(&frame); } } - } else if let Ok(oak_render::ticket::TicketPayload::ShmFrame(slot)) = result { - stale_slot = Some(slot); + } else { + stale_slot = match result { + Ok(oak_render::ticket::TicketPayload::ShmFrame(slot)) => { + Some(PreviewSlot::Shm(slot)) + } + Ok(oak_render::ticket::TicketPayload::Video(texture)) => { + Some(PreviewSlot::Video(texture)) + } + _ => None, + }; } } if let Some(slot) = stale_slot { - if let Some(m) = RenderManager::global() { - m.release_frame(&slot); - } + slot.release(); } }); m.tickets @@ -2349,34 +2382,51 @@ impl RealEngine { return None; } let slot = window.slots.remove(&frame.0)?; - let out = { - let meta = &slot.meta; - let (w, h) = (meta.width.max(0) as u32, meta.height.max(0) as u32); - if meta.format == super::renderops::PIXEL_FORMAT_F32 { - // M15 S3: the worker rendered F32 (the 10-bit display - // path) — repack/transform via `to_display` and register - // the RGBA16F texture for the viewer; the BGRA8 image it - // also builds is the CPU fallback only. - let (image, scope, samples) = - super::renderops::RenderedFrame::Shm(slot.clone()).to_display()?; - if let Some(samples) = samples { - super::gpu::register_display_frame(image.id.0, w, h, &samples); + match slot { + PreviewSlot::Shm(slot) => { + let out = { + let meta = &slot.meta; + let (w, h) = (meta.width.max(0) as u32, meta.height.max(0) as u32); + if meta.format == super::renderops::PIXEL_FORMAT_F32 { + // M15 S3: the worker rendered F32 (the 10-bit display + // path) — repack/transform via `to_display` and register + // the RGBA16F texture for the viewer; the BGRA8 image it + // also builds is the CPU fallback only. + let (image, scope, samples) = + super::renderops::RenderedFrame::Shm(slot.clone()).to_display()?; + if let Some(samples) = samples { + super::gpu::register_display_frame(image.id.0, w, h, &samples); + } + (Arc::new(image), scope) + } else { + let data = slot + .shm + .slot_bytes(slot.slot) + .get(..meta.data_size.max(0) as usize)?; + let image = bgra_bytes_to_render_image(w, h, data)?; + let scope = analyze_bgra8(w, h, data); + (Arc::new(image), scope) + } + }; + if let Some(m) = RenderManager::global() { + m.release_frame(&slot); } - (Arc::new(image), scope) - } else { - let data = slot - .shm - .slot_bytes(slot.slot) - .get(..meta.data_size.max(0) as usize)?; - let image = bgra_bytes_to_render_image(w, h, data)?; - let scope = analyze_bgra8(w, h, data); - (Arc::new(image), scope) + Some(out) + } + PreviewSlot::Video(texture) => { + // M4: the thread pipeline keeps frames on the GPU (or as an + // in-process CPU texture); `to_display` presents a GPU frame + // zero-copy when the device is shared and registers it for + // the viewer, or hands back samples for the staging upload. + let (w, h) = texture.size(); + let (image, scope, samples) = + super::renderops::RenderedFrame::Gpu(texture).to_display()?; + if let Some(samples) = samples { + super::gpu::register_display_frame(image.id.0, w.max(0) as u32, h.max(0) as u32, &samples); + } + Some((Arc::new(image), scope)) } - }; - if let Some(m) = RenderManager::global() { - m.release_frame(&slot); } - Some(out) } /// Cancels every monitor's pre-render window and releases its held @@ -2386,14 +2436,14 @@ impl RealEngine { // Collect the teardown work under the lock, then run it outside: // `cancel_preview_sequence` fires completions synchronously and those // completions lock `preview_windows` (self-deadlock otherwise). - let mut pending: Vec<(u64, Vec)> = Vec::new(); + let mut pending: Vec<(u64, Vec)> = Vec::new(); { let mut windows = self .preview_windows .lock() .unwrap_or_else(|e| e.into_inner()); for window in windows.values_mut() { - let slots: Vec = + let slots: Vec = std::mem::take(&mut window.slots).into_values().collect(); pending.push((window.sequence, slots)); window.submitted.clear(); @@ -2402,8 +2452,8 @@ impl RealEngine { if let Some(m) = RenderManager::global() { for (sequence, slots) in pending { m.cancel_preview_sequence(sequence); - for slot in &slots { - m.release_frame(slot); + for slot in slots { + slot.release(); } } } @@ -2423,14 +2473,14 @@ impl RealEngine { let Some(window) = windows.get_mut(&monitor) else { return; }; - let slots: Vec = std::mem::take(&mut window.slots).into_values().collect(); + let slots: Vec = std::mem::take(&mut window.slots).into_values().collect(); window.submitted.clear(); (window.sequence, slots) }; if let Some(m) = RenderManager::global() { m.cancel_preview_sequence(pending.0); - for slot in &pending.1 { - m.release_frame(slot); + for slot in pending.1 { + slot.release(); } } } diff --git a/crates/oak-render/examples/bench_playback.rs b/crates/oak-render/examples/bench_playback.rs index 3f1d4cfee..269029d40 100644 --- a/crates/oak-render/examples/bench_playback.rs +++ b/crates/oak-render/examples/bench_playback.rs @@ -15,32 +15,41 @@ // along with this program. If not, see . //! Real-footage playback benchmark: renders `N` sequential frames of a -//! real media file through the real oak-worker pool at the app's preview -//! proxy size, mimicking the playback pre-render window (Playback -//! priority, interleaved claiming, immediate slot release). Reports -//! throughput and completion latency so worker-side decode/render -//! hotspots can be measured end to end. +//! real media file through the real oak-worker pool **or** the M1/M4 +//! thread pipeline at the app's preview proxy size, mimicking the +//! playback pre-render window (Playback priority, immediate release). +//! Reports throughput, completion latency, first-frame latency and CPU +//! time (self + children) so the M4 acceptance comparison between the two +//! backends is reproducible. //! //! Run from the repo root: //! //! ```sh -//! cargo run --release -p oakrender --example bench_playback -- [frames] [workers] [long_edge] +//! cargo run --release -p oakrender --example bench_playback -- [frames] [workers] [long_edge] [processes|pipeline] //! ``` //! -//! `frames` defaults to 240, `workers` to the adaptive policy, and -//! `long_edge` to 480 (the app's preview proxy size). To profile a -//! worker while this runs: `pgrep oak-worker | head -1 | xargs sample 10`. +//! `frames` defaults to 240, `workers` to the adaptive policy, `long_edge` +//! to 480 (the app's preview proxy size) and the backend to `processes`. +//! Set `OAK_BENCH_GENERATE=1` to synthesize a 1080p/25 fps 10 s clip at +//! `` when the file does not exist. use std::path::PathBuf; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -use oak_core::Rational; -use oak_render::ipc::SLOT_FORMAT_BGRA8; +use oak_core::{PixelFormat, Rational}; use oak_render::procpool::{DispatcherConfig, ProcessDispatcher}; -use oak_render::ticket::{TicketPayload, TicketResult, VideoTicketParams}; +use oak_render::ticket::{ + Completion, Producer, TicketPayload, TicketResult, VideoTicketParams, +}; use oak_render::worker::{Job, JobDispatch, JobSchedule}; +/// The comparison frame format: both backends must produce the same +/// pixels for the numbers to be comparable (the app's 8-bit preview path +/// uses BGRA8, but the process pool would then quantize — the F32 slots +/// are the like-for-like path, and what the 10-bit preview uses). +const BENCH_FORMAT: PixelFormat = PixelFormat::F32; + /// Locate the oak-worker binary (see bench_process). fn worker_bin() -> PathBuf { if let Ok(p) = std::env::var("OAK_WORKER_BIN") { @@ -59,66 +68,121 @@ fn worker_bin() -> PathBuf { PathBuf::from("oak-worker") } -fn main() { - let media = std::env::args() - .nth(1) - .unwrap_or_else(|| "tests/demo.mp4".to_string()); - let frames: usize = std::env::args() - .nth(2) - .and_then(|s| s.parse().ok()) - .unwrap_or(240); - let workers: Option = std::env::args().nth(3).and_then(|s| s.parse().ok()); - let long_edge: i32 = std::env::args() - .nth(4) - .and_then(|s| s.parse().ok()) - .unwrap_or(480); +/// `(user, system)` CPU seconds of this process and its children. +fn cpu_times() -> (f64, f64) { + fn rusage(who: i32) -> (f64, f64) { + let mut usage: libc::rusage = unsafe { std::mem::zeroed() }; + if unsafe { libc::getrusage(who, &mut usage) } != 0 { + return (0.0, 0.0); + } + let seconds = |tv: libc::timeval| tv.tv_sec as f64 + tv.tv_usec as f64 / 1e6; + (seconds(usage.ru_utime), seconds(usage.ru_stime)) + } + let self_times = rusage(libc::RUSAGE_SELF); + let children = rusage(libc::RUSAGE_CHILDREN); + ( + self_times.0 + children.0, + self_times.1 + children.1, + ) +} - // The app's preview proxy size: the sequence's aspect scaled to the - // long edge (demo.mp4 is 16:9 1080p). - let (width, height) = ((long_edge as f64 * 16.0 / 9.0).round() as i32, long_edge); +/// One footage ticket over the whole timeline. +fn footage_params(media: &str, time: Rational, width: i32, height: i32) -> VideoTicketParams { + VideoTicketParams { + viewer: 1, + project: String::new(), + time, + force_size: Some((width, height)), + force_format: Some(BENCH_FORMAT), + cache: None, + cache_dir: None, + cache_id: None, + cache_timebase: None, + footage: Some((media.to_string(), 0)), + montage: Vec::new(), + adjustments: Vec::new(), + } +} +fn report(entries: &[(i64, Instant, Instant)], start: Instant, elapsed: Duration, cpu: (f64, f64)) { + let completed = entries.len(); + let throughput = completed as f64 / elapsed.as_secs_f64(); + let mut latencies: Vec = entries + .iter() + .map(|(_, submit, done)| (*done - *submit).as_secs_f64() * 1000.0) + .collect(); + latencies.sort_by(|a, b| a.partial_cmp(b).unwrap()); + let first = entries + .iter() + .map(|(_, _, done)| (*done - start).as_secs_f64() * 1000.0) + .fold(f64::INFINITY, f64::min); + + let report = |name: &str, value: String| println!("{name:<38} {value}"); + report("frames completed", completed.to_string()); + report("total wall time", format!("{:.2} s", elapsed.as_secs_f64())); + report( + "throughput", + format!( + "{throughput:.1} fps ({:.1} ms/frame)", + 1000.0 / throughput.max(f64::EPSILON) + ), + ); + report("first-frame latency", format!("{first:.1} ms")); + if !latencies.is_empty() { + let mean = latencies.iter().sum::() / latencies.len() as f64; + report("completion latency mean", format!("{mean:.1} ms")); + report( + "completion latency p50/p95/max", + format!( + "{:.1} / {:.1} / {:.1} ms", + latencies[latencies.len() / 2], + latencies[((latencies.len() as f64 * 0.95) as usize).min(latencies.len() - 1)], + latencies.last().unwrap() + ), + ); + } + report( + "cpu user + sys", + format!("{:.2} + {:.2} s", cpu.0, cpu.1), + ); + report( + "main-heap frame copies", + oak_render::procpool::main_heap_frame_copies().to_string(), + ); +} + +/// Process-pool playback: post every frame at Playback priority and pump. +fn run_processes(media: &str, frames: usize, width: i32, height: i32, workers: Option) { let config = DispatcherConfig { worker_bin: Some(worker_bin()), workers: workers.unwrap_or(0), slots_per_worker: 8, width, height, - slot_format: SLOT_FORMAT_BGRA8, + slot_format: BENCH_FORMAT as i32, batch_size: 0, graph_snapshot: None, handshake_timeout_ms: 30_000, }; let dispatcher = ProcessDispatcher::new(config).expect("dispatcher config"); dispatcher.start().expect("workers start + handshake"); - let worker_count = dispatcher.worker_count(); - println!("oak-worker pool: {worker_count} worker(s), {frames} x {width}x{height} BGRA8 frames of {media}"); + println!( + "oak-worker pool: {} worker(s), {frames} x {width}x{height} {BENCH_FORMAT:?} frames of {media}", + dispatcher.worker_count() + ); - // One completion record per frame: (ticket/frame, submit, completion). + let cpu_start = cpu_times(); let results = Arc::new(Mutex::new(Vec::<(i64, Instant, Instant)>::new())); let start = Instant::now(); for i in 0..frames { let results = results.clone(); let dc = dispatcher.clone(); let frame = i as i64; - let media_clone = media.clone(); + let media_clone = media.to_string(); let job = Job { node_identity: 1, time: Rational::new(frame, 25), - params: Arc::new(VideoTicketParams { - viewer: 1, - project: String::new(), - time: Rational::new(frame, 25), - force_size: Some((width, height)), - force_format: None, - cache: None, - cache_dir: None, - cache_id: None, - cache_timebase: None, - // A single footage clip covers the whole timeline. - footage: Some((media_clone, 0)), - montage: Vec::new(), - adjustments: Vec::new(), - }), + params: Arc::new(footage_params(&media_clone, Rational::new(frame, 25), width, height)), audio: None, produce: Arc::new(|_, _| { Err(oak_render::error::Error::Failed( @@ -143,7 +207,6 @@ fn main() { eprintln!("frame {frame} failed: {e}"); } }), - // Playback priority, the pre-render window's schedule. schedule: JobSchedule::playback(frame, frame, 0), }; if !dispatcher.post(job) { @@ -152,7 +215,6 @@ fn main() { } } - // Pump until every completion has landed. let deadline = Instant::now() + Duration::from_secs(300); loop { dispatcher.poll(); @@ -167,44 +229,123 @@ fn main() { std::thread::sleep(Duration::from_millis(2)); } let elapsed = start.elapsed(); - - let entries: Vec<(i64, Instant, Instant)> = - results.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); - let completed = entries.len(); - let throughput = completed as f64 / elapsed.as_secs_f64(); - - // Per-frame completion latency (submit -> done), an end-to-end proxy - // for the worker's per-frame render cost under load. - let mut latencies: Vec = entries - .iter() - .map(|(_, submit, done)| (*done - *submit).as_secs_f64() * 1000.0) - .collect(); - latencies.sort_by(|a, b| a.partial_cmp(b).unwrap()); - - let report = |name: &str, value: String| println!("{name:<38} {value}"); - report("frames completed", completed.to_string()); - report("total wall time", format!("{:.2} s", elapsed.as_secs_f64())); - report( - "throughput", - format!("{throughput:.1} fps ({:.1} ms/frame)", 1000.0 / throughput.max(f64::EPSILON)), - ); - if !latencies.is_empty() { - let mean = latencies.iter().sum::() / latencies.len() as f64; - report("completion latency mean", format!("{mean:.1} ms")); - report( - "completion latency p50/p95/max", - format!( - "{:.1} / {:.1} / {:.1} ms", - latencies[latencies.len() / 2], - latencies[((latencies.len() as f64 * 0.95) as usize).min(latencies.len() - 1)], - latencies.last().unwrap() - ), - ); - } - report( - "main-heap frame copies", - oak_render::procpool::main_heap_frame_copies().to_string(), - ); - + // Children (the worker pool) are only accounted at wait(): shut the + // pool down before reading RUSAGE_CHILDREN, then report. dispatcher.shutdown(); + let cpu_end = cpu_times(); + let entries: Vec<(i64, Instant, Instant)> = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); + report( + &entries, + start, + elapsed, + (cpu_end.0 - cpu_start.0, cpu_end.1 - cpu_start.1), + ); +} + +/// M4 thread-pipeline playback: the same Playback jobs on the single +/// render/decode threads; the pipeline prefetches each frame's decode on +/// post and orders the queue by priority. +fn run_pipeline(media: &str, frames: usize, width: i32, height: i32) { + let backend = oak_render::pipeline::PipelineBackend::new().expect("pipeline start"); + println!("thread pipeline: 1 render + 1 decode thread, {frames} x {width}x{height} F32 frames of {media}"); + + let cpu_start = cpu_times(); + let results = Arc::new(Mutex::new(Vec::<(i64, Instant, Instant)>::new())); + let start = Instant::now(); + for i in 0..frames { + let frame = i as i64; + let results = results.clone(); + let submitted = Instant::now(); + let done: Completion = Box::new(move |result: TicketResult| { + if let Ok(TicketPayload::Video(_)) = result { + results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push((frame, submitted, Instant::now())); + } + }); + let producer: Producer = + Arc::new(|time, params| { + oak_render::eval::render_produced_frame(time, params) + .map(TicketPayload::Video) + }); + let job = Job { + node_identity: 1, + time: Rational::new(frame, 25), + params: Arc::new(footage_params(media, Rational::new(frame, 25), width, height)), + audio: None, + produce: producer, + done, + schedule: JobSchedule::playback(frame, frame, 0), + }; + // The blocking post is the pipeline's backpressure: once the render + // queue is full the submitter waits (the app's window is capped by + // `preview_window_capacity`). + if !backend.post(job) { + eprintln!("post refused at frame {frame}"); + break; + } + } + + let deadline = Instant::now() + Duration::from_secs(300); + loop { + let done = results.lock().unwrap_or_else(|e| e.into_inner()).len(); + if done >= frames { + break; + } + if Instant::now() > deadline { + eprintln!("timeout: {done}/{frames} completions"); + break; + } + std::thread::sleep(Duration::from_millis(2)); + } + let elapsed = start.elapsed(); + let cpu_end = cpu_times(); + let entries: Vec<(i64, Instant, Instant)> = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); + report( + &entries, + start, + elapsed, + (cpu_end.0 - cpu_start.0, cpu_end.1 - cpu_start.1), + ); + backend.shutdown(); +} + +fn main() { + let media = std::env::args() + .nth(1) + .unwrap_or_else(|| "tests/demo.mp4".to_string()); + let frames: usize = std::env::args() + .nth(2) + .and_then(|s| s.parse().ok()) + .unwrap_or(240); + let workers: Option = std::env::args().nth(3).and_then(|s| s.parse().ok()); + let long_edge: i32 = std::env::args() + .nth(4) + .and_then(|s| s.parse().ok()) + .unwrap_or(480); + let mode = std::env::args().nth(5).unwrap_or_else(|| "processes".to_string()); + + // The app's preview proxy size: the sequence's aspect scaled to the + // long edge (demo.mp4 is 16:9 1080p). + let (width, height) = ((long_edge as f64 * 16.0 / 9.0).round() as i32, long_edge); + + if !std::path::Path::new(&media).exists() && std::env::var_os("OAK_BENCH_GENERATE").is_some() { + oak_codec::testmedia::write_test_clip(std::path::Path::new(&media), 1920, 1080, 250, 25) + .expect("generate the 1080p benchmark clip"); + println!("generated 1080p/25 fps test media at {media}"); + } + + match mode.as_str() { + "pipeline" => run_pipeline(&media, frames, width, height), + _ => run_processes(&media, frames, width, height, workers), + } } diff --git a/crates/oak-render/src/pipeline.rs b/crates/oak-render/src/pipeline.rs index e4b69a22f..ff4179827 100644 --- a/crates/oak-render/src/pipeline.rs +++ b/crates/oak-render/src/pipeline.rs @@ -73,7 +73,7 @@ use std::collections::{HashMap, VecDeque}; use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}; -use std::sync::mpsc::{self, SyncSender, TrySendError}; +use std::sync::mpsc::{self, SyncSender}; use std::sync::{Arc, Condvar, Mutex, MutexGuard, OnceLock}; use std::thread::JoinHandle; @@ -81,6 +81,7 @@ use oak_core::texture::Texture; use oak_core::{PixelFormat, Rational}; use crate::error::{Error, Result}; +use crate::scheduler::FramePriority; use crate::worker::{execute_job, Job, JobDispatch}; /// Render-queue bound (jobs). Small on purpose: the queue exists to keep @@ -104,6 +105,66 @@ fn lock(m: &Mutex) -> MutexGuard<'_, T> { m.lock().unwrap_or_else(|e| e.into_inner()) } +/// Render-queue priority class (lower first): interactive seeks jump the +/// queue, playback stays in playhead order, background work (exports, +/// autocache) yields. +fn job_class(job: &Job) -> u8 { + job.schedule.priority as u8 +} + +/// Insert `job` keeping the queue ordered by priority class; stable within +/// a class (FIFO). Callers ensure capacity first (except Seek, which may +/// over-admit — see `PipelineBackend::push`). +fn enqueue_job(queue: &mut VecDeque, job: Job) { + let class = job_class(&job); + let position = queue + .iter() + .position(|queued| job_class(queued) > class) + .unwrap_or(queue.len()); + queue.insert(position, job); +} + +/// The decode requests a playback frame needs (M4 read-ahead): one per +/// montage clip covering `time`, plus the single-footage ticket's stream. +/// The media arithmetic mirrors `eval::render_montage_frame_into` +/// (`media_in + (time - in_time)`) so a prefetched frame is the exact key +/// the render rendezvous asks for. +fn playback_decode_requests( + params: &crate::ticket::VideoTicketParams, + time: Rational, +) -> Vec { + let size = params.render_size(); + if size.0 <= 0 || size.1 <= 0 { + return Vec::new(); + } + let mut requests = Vec::new(); + for clip in ¶ms.montage { + if time < clip.in_time || time >= clip.out_time { + continue; + } + requests.push(DecodeRequest { + filename: clip.filename.clone(), + stream_index: clip.stream_index, + time: clip.media_in + (time - clip.in_time), + size, + format: PixelFormat::F32, + }); + } + if let Some((filename, stream_index)) = ¶ms.footage { + // The read-ahead must decode exactly what the render request will + // ask for, or the cache never hits: the eval path uses + // `force_format.unwrap_or(F32)`, so mirror that here. + requests.push(DecodeRequest { + filename: filename.clone(), + stream_index: *stream_index, + time, + size, + format: params.force_format.unwrap_or(PixelFormat::F32), + }); + } + requests +} + // --------------------------------------------------------------------------- // Decode service // --------------------------------------------------------------------------- @@ -190,8 +251,24 @@ enum DecodeCommand { /// Barrier: replies once every command sent before it has been fully /// processed (used by tests and by the decode-service contract). Sync { reply: SyncSender<()> }, - /// Stop the thread (it drains, then exits). - Shutdown, +} + +/// The decode command queue (M4): bounded, with rendezvous requests +/// served ahead of prefetches — a render that needs a frame must never +/// wait behind speculative decode work; prefetches only fill idle decode +/// time. `Sync` barriers stay in FIFO order (they cover everything sent +/// before them). +struct DecodeQueue { + pending: VecDeque, + shutting_down: bool, +} + +struct DecodeShared { + queue: Mutex, + /// A command arrived. + work: Condvar, + /// A slot freed (or shutdown started). + room: Condvar, } struct DecodeInner { @@ -211,7 +288,7 @@ struct DecodeInner { /// there); [`DecodeService::stats`] and [`DecodeService::lru_len`] read /// shared counters, so they are safe from any thread. pub struct DecodeService { - commands: SyncSender, + shared: Arc, gate: PrefetchGate, inner: Arc, handle: Mutex>>, @@ -222,28 +299,63 @@ impl DecodeService { /// `gate` deciding whether prefetch work is wanted. `lru_capacity` 0 /// disables caching (every request decodes). pub fn new(lru_capacity: usize, gate: PrefetchGate) -> Arc { - let (tx, rx) = mpsc::sync_channel(DECODE_QUEUE_CAP); + let shared = Arc::new(DecodeShared { + queue: Mutex::new(DecodeQueue { + pending: VecDeque::new(), + shutting_down: false, + }), + work: Condvar::new(), + room: Condvar::new(), + }); let inner = Arc::new(DecodeInner { counters: DecodeCounters::default(), lru_len: AtomicUsize::new(0), }); let service = Arc::new(Self { - commands: tx, + shared: shared.clone(), gate, inner: inner.clone(), handle: Mutex::new(None), }); let spawned = std::thread::Builder::new() .name("oak-decode".into()) - .spawn(move || decode_loop(rx, inner, lru_capacity)); - // A failed spawn leaves the receiver dropped, so `request` reports - // the service as unavailable and the caller decodes inline. + .spawn(move || decode_loop(shared, inner, lru_capacity)); + // A failed spawn leaves the queue unreachable, so `request` + // reports the service as unavailable and the caller decodes + // inline. if let Ok(handle) = spawned { *lock(&service.handle) = Some(handle); } service } + /// Queue a command: `blocking` waits for room when the queue is full, + /// otherwise a full queue refuses (`false`). Also `false` once the + /// service is stopping. + fn enqueue(&self, command: DecodeCommand, blocking: bool) -> bool { + let shared = &self.shared; + let mut queue = lock(&shared.queue); + if queue.shutting_down { + return false; + } + if !blocking { + if queue.pending.len() >= DECODE_QUEUE_CAP { + return false; + } + } else { + while queue.pending.len() >= DECODE_QUEUE_CAP { + if queue.shutting_down { + return false; + } + queue = shared.room.wait(queue).unwrap_or_else(|e| e.into_inner()); + } + } + queue.pending.push_back(command); + drop(queue); + shared.work.notify_one(); + true + } + /// Decode `request`, returning the frame. `None` means the service is /// gone (shutting down, or already down) — the caller must decode /// inline instead; `Some(Err(_))` is a real decode failure and must be @@ -253,9 +365,8 @@ impl DecodeService { /// the render side propagates to submission. pub fn request(&self, request: DecodeRequest) -> Option> { let (reply, rx) = mpsc::sync_channel(1); - match self.commands.send(DecodeCommand::Request { request, reply }) { - Ok(()) => {} - Err(_) => return None, + if !self.enqueue(DecodeCommand::Request { request, reply }, true) { + return None; } match rx.recv() { Ok(result) => Some(result), @@ -276,27 +387,24 @@ impl DecodeService { .fetch_add(1, Ordering::Relaxed); return false; } - match self.commands.try_send(DecodeCommand::Prefetch { request }) { - Ok(()) => { - self.inner.counters.prefetches.fetch_add(1, Ordering::Relaxed); - true - } - Err(TrySendError::Full(_)) | Err(TrySendError::Disconnected(_)) => { - self.inner - .counters - .prefetch_refused - .fetch_add(1, Ordering::Relaxed); - false - } + if self.enqueue(DecodeCommand::Prefetch { request }, false) { + self.inner.counters.prefetches.fetch_add(1, Ordering::Relaxed); + true + } else { + self.inner + .counters + .prefetch_refused + .fetch_add(1, Ordering::Relaxed); + false } } /// Wait until every command sent before this call has been processed - /// (a barrier through the same FIFO queue). `false` when the service is + /// (a barrier through the same queue). `false` when the service is /// already gone. pub fn wait_idle(&self) -> bool { let (reply, rx) = mpsc::sync_channel(0); - if self.commands.send(DecodeCommand::Sync { reply }).is_err() { + if !self.enqueue(DecodeCommand::Sync { reply }, true) { return false; } rx.recv().is_ok() @@ -314,21 +422,56 @@ impl DecodeService { /// Stop the decode thread and wait for it to exit. Idempotent; a /// request arriving after this (or racing it) reports `None` and the - /// caller decodes inline. + /// caller decodes inline. Queued commands are drained first (bounded + /// by the queue cap), so a `wait_idle` sent before shutdown still + /// answers. pub fn shutdown(&self) { let handle = lock(&self.handle).take(); let Some(handle) = handle else { return }; - let _ = self.commands.send(DecodeCommand::Shutdown); + { + let mut queue = lock(&self.shared.queue); + queue.shutting_down = true; + } + self.shared.work.notify_all(); + self.shared.room.notify_all(); let _ = handle.join(); } } +/// Pop the next command: a rendezvous request first, else FIFO; `None` +/// once the queue is empty and shutdown was requested. Removing a command +/// opens queue room, so blocked requesters are woken. +fn pop_command(shared: &DecodeShared) -> Option { + let mut queue = lock(&shared.queue); + loop { + if let Some(pos) = queue + .pending + .iter() + .position(|c| matches!(c, DecodeCommand::Request { .. })) + { + let command = queue.pending.remove(pos).expect("just found"); + drop(queue); + shared.room.notify_one(); + return Some(command); + } + if let Some(command) = queue.pending.pop_front() { + drop(queue); + shared.room.notify_one(); + return Some(command); + } + if queue.shutting_down { + return None; + } + queue = shared.work.wait(queue).unwrap_or_else(|e| e.into_inner()); + } +} + /// The decode thread: one command at a time, LRU in a plain `HashMap` /// keyed by [`DecodeRequest`] with a monotonic recency tick. -fn decode_loop(rx: mpsc::Receiver, inner: Arc, lru_capacity: usize) { +fn decode_loop(shared: Arc, inner: Arc, lru_capacity: usize) { let mut lru: HashMap = HashMap::new(); let mut tick: u64 = 1; - while let Ok(command) = rx.recv() { + while let Some(command) = pop_command(&shared) { match command { DecodeCommand::Request { request, reply } => { let result = serve(&request, &mut lru, &mut tick, lru_capacity, &inner); @@ -342,12 +485,11 @@ fn decode_loop(rx: mpsc::Receiver, inner: Arc, lru_c DecodeCommand::Sync { reply } => { let _ = reply.send(()); } - DecodeCommand::Shutdown => break, } } - // Anything still queued is dropped with the receiver: the pending - // reply senders disconnect, so blocked requesters see `None` and fall - // back to the inline decode. + // Anything still queued is dropped with the thread: the pending reply + // senders disconnect, so blocked requesters see `None` and fall back + // to the inline decode. inner.lru_len.store(0, Ordering::Relaxed); } @@ -588,6 +730,10 @@ impl PipelineBackend { } fn push(&self, job: Job, blocking: bool) -> bool { + // M4 playback read-ahead: queue the job's footage decodes before it + // reaches the render thread, so the decode thread works on frame + // N+1 while the render thread evaluates frame N. + self.prefetch_decodes(&job); let inner = &self.inner; let mut queue = lock(&inner.queue); if inner.stopping.load(Ordering::Acquire) { @@ -596,12 +742,19 @@ impl PipelineBackend { if !blocking && queue.len() >= RENDER_QUEUE_CAP { return false; } - // Backpressure, with one exception: the render thread itself may - // re-post from a completion callback (a producer that submits - // follow-up work). It can never make room by waiting — it is the - // only thread that drains the queue — so it is allowed to run one - // ahead of the bound instead of deadlocking. - if blocking && !is_render_thread(inner) { + // Backpressure, with two exceptions: + // - The render thread itself may re-post from a completion callback + // (a producer that submits follow-up work). It can never make room + // by waiting — it is the only thread that drains the queue — so it + // is allowed to run one ahead of the bound instead of deadlocking. + // - A Seek (the UI's interactive frame) may over-admit: priority + // reordering only helps jobs already in the queue, and a full + // queue of background/playback work must not stall the UI thread + // until a whole export/cache frame finishes. Playback posts are + // capacity-capped by the app window, background posts are off the + // UI thread, so over-admission stays bounded in practice. + let seek = matches!(job.schedule.priority, FramePriority::Seek); + if blocking && !is_render_thread(inner) && !seek { while queue.len() >= RENDER_QUEUE_CAP { if inner.stopping.load(Ordering::Acquire) { return false; @@ -609,7 +762,7 @@ impl PipelineBackend { queue = inner.room.wait(queue).unwrap_or_else(|e| e.into_inner()); } } - queue.push_back(job); + enqueue_job(&mut queue, job); inner.depth.store(queue.len(), Ordering::Relaxed); inner.posted.fetch_add(1, Ordering::Relaxed); drop(queue); @@ -617,6 +770,19 @@ impl PipelineBackend { true } + /// M4 playback read-ahead: a Playback job's footage decodes are queued + /// speculatively as soon as the job is submitted. The requests mirror + /// the montage/footage eval exactly (same media time, size and format), + /// so a completed prefetch is an LRU hit for the render rendezvous. + fn prefetch_decodes(&self, job: &Job) { + if !matches!(job.schedule.priority, FramePriority::Playback) { + return; + } + for request in playback_decode_requests(&job.params, job.time) { + let _ = self.inner.decode.prefetch(request); + } + } + fn shutdown_impl(&self) { let inner = &self.inner; if inner.stopping.swap(true, Ordering::AcqRel) { @@ -667,6 +833,45 @@ impl JobDispatch for PipelineBackend { fn shutdown(&self) { self.shutdown_impl(); } + + /// M4: the app's pre-render window may only queue as much playback + /// work as the render queue can accept, so playback posts never block + /// the UI (Seek posts additionally over-admit; see `push`). + fn preview_window_capacity(&self) -> Option { + Some(self.queue_free().max(1)) + } + + /// M4: drop a queued playback frame the playhead has passed (pending + /// only — the in-flight frame cannot be interrupted). Matched on the + /// full key (viewer identity, frame, version) so cancelling one + /// monitor's window never drops the other sequence's frame at the same + /// number. The completion fires with `Error::State`, exactly like a + /// cancelled worker frame, so the window removes it from `submitted`. + fn cancel_preview_frame(&self, sequence: u64, frame: i64, version: u64) { + let inner = &self.inner; + let mut removed = Vec::new(); + { + let mut queue = lock(&inner.queue); + let mut index = 0; + while index < queue.len() { + let job = &queue[index]; + let matches = job.node_identity == sequence + && job.schedule.frame == Some(frame) + && job.schedule.version == version + && matches!(job.schedule.priority, FramePriority::Playback); + if matches { + removed.push(queue.remove(index).expect("index in range")); + } else { + index += 1; + } + } + inner.depth.store(queue.len(), Ordering::Relaxed); + } + inner.room.notify_all(); + for job in removed { + (job.done)(Err(Error::State)); + } + } } fn is_render_thread(inner: &PipelineInner) -> bool { @@ -705,6 +910,7 @@ fn render_loop(inner: Arc) { #[cfg(test)] mod tests { use super::*; + use crate::worker::JobSchedule; use oak_core::texture::Frame; /// A unique clip per test (the process id separates test binaries, the @@ -957,4 +1163,165 @@ mod tests { assert!(!service.wait_idle()); let _ = std::fs::remove_file(&path); } + + // ---- M4: priority queues and playback read-ahead -------------------- + + /// A minimal job for queue-ordering tests (the producer/done are no-ops + /// and never run). + fn dummy_job(priority: FramePriority, frame: i64) -> Job { + Job { + node_identity: 1, + time: Rational::new(frame, 1), + params: Arc::new(crate::ticket::VideoTicketParams { + viewer: 1, + project: String::new(), + time: Rational::new(frame, 1), + force_size: Some((16, 16)), + force_format: None, + cache: None, + cache_dir: None, + cache_id: None, + cache_timebase: None, + footage: None, + montage: Vec::new(), + adjustments: Vec::new(), + }), + audio: None, + produce: Arc::new(|_, _| Ok(crate::ticket::TicketPayload::Video(Texture::dummy()))), + done: Box::new(|_| {}), + schedule: JobSchedule { + priority, + frame: Some(frame), + distance: frame, + version: 0, + }, + } + } + + /// M4: the decode queue serves a rendezvous request before queued + /// prefetches, while a `Sync` barrier keeps FIFO order. + #[test] + fn decode_queue_prioritizes_requests_over_prefetches() { + let shared = Arc::new(DecodeShared { + queue: Mutex::new(DecodeQueue { + pending: VecDeque::new(), + shutting_down: false, + }), + work: Condvar::new(), + room: Condvar::new(), + }); + let (sync_tx, _sync_rx) = mpsc::sync_channel(0); + let (request_tx, _request_rx) = mpsc::sync_channel(0); + let request = |filename: &str, time: Rational| DecodeRequest { + filename: filename.to_string(), + stream_index: 0, + time, + size: (16, 16), + format: PixelFormat::F32, + }; + { + let mut queue = lock(&shared.queue); + queue + .pending + .push_back(DecodeCommand::Sync { reply: sync_tx }); + queue.pending.push_back(DecodeCommand::Prefetch { + request: request("a.mp4", Rational::new(0, 1)), + }); + queue.pending.push_back(DecodeCommand::Prefetch { + request: request("a.mp4", Rational::new(1, 10)), + }); + queue.pending.push_back(DecodeCommand::Request { + request: request("a.mp4", Rational::new(2, 10)), + reply: request_tx, + }); + } + // The request preempts the two earlier prefetches... + assert!(matches!( + pop_command(&shared), + Some(DecodeCommand::Request { .. }) + )); + // ...but the Sync barrier stays ahead of the prefetches that were + // sent before it (FIFO). + assert!(matches!( + pop_command(&shared), + Some(DecodeCommand::Sync { .. }) + )); + assert!(matches!( + pop_command(&shared), + Some(DecodeCommand::Prefetch { .. }) + )); + assert!(matches!( + pop_command(&shared), + Some(DecodeCommand::Prefetch { .. }) + )); + } + + /// M4: the render queue is priority-ordered (Seek, Playback, + /// Background) and FIFO inside a class. + #[test] + fn render_queue_orders_seek_playback_background() { + let mut queue = VecDeque::new(); + enqueue_job(&mut queue, dummy_job(FramePriority::Background, 0)); + enqueue_job(&mut queue, dummy_job(FramePriority::Playback, 1)); + enqueue_job(&mut queue, dummy_job(FramePriority::Seek, 2)); + enqueue_job(&mut queue, dummy_job(FramePriority::Playback, 3)); + let classes: Vec = queue.iter().map(job_class).collect(); + assert_eq!(classes, vec![0, 1, 1, 2]); + let frames: Vec = queue + .iter() + .map(|job| job.schedule.frame.unwrap()) + .collect(); + assert_eq!(frames, vec![2, 1, 3, 0], "FIFO within a class"); + } + + /// M4 read-ahead: the derived decode requests mirror the montage eval + /// (media time, target size, F32) and skip clips that do not cover the + /// requested time. + #[test] + fn playback_decode_requests_mirror_montage_times() { + let clip = |filename: &str, in_t: Rational, out_t: Rational, media_in: Rational| { + crate::ticket::MontageClip { + filename: filename.to_string(), + stream_index: 0, + in_time: in_t, + out_time: out_t, + media_in, + gain: 1.0, + effects: Vec::new(), + } + }; + let params = crate::ticket::VideoTicketParams { + viewer: 1, + project: String::new(), + time: Rational::new(5, 10), + force_size: Some((64, 32)), + force_format: None, + cache: None, + cache_dir: None, + cache_id: None, + cache_timebase: None, + footage: None, + montage: vec![ + clip("covered.mp4", Rational::new(0, 1), Rational::new(1, 1), Rational::new(2, 1)), + clip("outside.mp4", Rational::new(1, 1), Rational::new(2, 1), Rational::new(0, 1)), + ], + adjustments: Vec::new(), + }; + let requests = playback_decode_requests(¶ms, Rational::new(5, 10)); + assert_eq!(requests.len(), 1, "only the covering clip is prefetched"); + assert_eq!(requests[0].filename, "covered.mp4"); + assert_eq!(requests[0].time, Rational::new(5, 2), "media_in + (time - in)"); + assert_eq!(requests[0].size, (64, 32)); + assert_eq!(requests[0].format, PixelFormat::F32); + + // A single-footage ticket prefetches its stream at the ticket time. + let mut footage = params; + footage.montage.clear(); + footage.footage = Some(("solo.mp4".to_string(), 2)); + let requests = playback_decode_requests(&footage, Rational::new(5, 10)); + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].filename, "solo.mp4"); + assert_eq!(requests[0].stream_index, 2); + assert_eq!(requests[0].time, Rational::new(1, 2)); + } } diff --git a/crates/oak-render/tests/render_threads_test.rs b/crates/oak-render/tests/render_threads_test.rs index ec9651a9d..a29f644b6 100644 --- a/crates/oak-render/tests/render_threads_test.rs +++ b/crates/oak-render/tests/render_threads_test.rs @@ -823,6 +823,318 @@ fn pipeline_layered_playback_has_zero_gpu_readbacks() { let _ = std::fs::remove_file(&below); } +/// M4: every playback post queues its footage read-ahead and every +/// distinct frame ends up decoded exactly once — by the prefetch or, if +/// the render request wins the race, by the rendezvous. This is a smoke +/// test for the wiring; that the read-ahead is actually *used* is proven +/// deterministically by `pipeline_prefetch_is_the_frame_the_render_request_uses`, +/// and the priority order by `pipeline_orders_seek_ahead_of_background_end_to_end`. +#[test] +fn pipeline_playback_prefetches_ahead_of_the_render() { + let _lock = lock(); + pin_legacy_working_space(); + let guard = common::ManagerGuard::init_with(RenderBackendChoice::Pipeline); + let manager = RenderManager::global().expect("manager installed"); + let backend = manager.pipeline_backend().expect("pipeline selected"); + let path = test_clip("playback_prefetch"); + let times: Vec = (0..6).map(|n| Rational::new(n, 10)).collect(); + let (tx, rx) = mpsc::channel(); + for (n, time) in times.iter().enumerate() { + let tx = tx.clone(); + let done: Completion = Box::new(move |result| { + let ok = matches!(result, Ok(TicketPayload::Video(_))); + let _ = tx.send((n, ok)); + }); + manager + .tickets + .submit_playback(montage_params(&path, *time), n as i64, n as i64, 0, done); + } + drop(tx); + for _ in 0..times.len() { + match rx.recv_timeout(Duration::from_secs(60)) { + Ok((n, true)) => { + let _ = n; + } + Ok((n, false)) => panic!("frame {n} did not produce a video payload"), + Err(err) => panic!("playback completion timeout: {err}"), + } + } + let stats = backend.decode_service().stats(); + assert_eq!( + stats.prefetches, + times.len() as u64, + "every playback post queued its footage prefetch" + ); + assert_eq!( + stats.decodes, + times.len() as u64, + "each distinct frame decodes exactly once (prefetch or rendezvous)" + ); + // Note: `procpool::main_heap_frame_copies` only counts the shm path, + // which the thread pipeline never touches, so asserting it here would + // be vacuous. The in-process frame path does clone `Frame.data` at the + // eval-cache and service-LRU boundaries (M5 narrows this); the bench + // comparison must not claim "no heap copies" for the pipeline. + drop(guard); + let _ = std::fs::remove_file(&path); +} + +/// A producer that parks the render thread until `release` is signalled, +/// recording `tag` when it finally runs. Two independent gates let a test +/// keep a frame in flight while it posts more work. +fn parked_producer( + started: Arc<(Mutex, Condvar)>, + release: Arc<(Mutex, Condvar)>, + order: Arc>>, + tag: &'static str, +) -> Producer { + Arc::new(move |_time: Rational, _params: &VideoTicketParams| { + { + let (started, work) = &*started; + *started.lock().unwrap_or_else(|e| e.into_inner()) = true; + work.notify_all(); + } + let (released, work) = &*release; + let mut released = released.lock().unwrap_or_else(|e| e.into_inner()); + while !*released { + released = work.wait(released).unwrap_or_else(|e| e.into_inner()); + } + order + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(tag); + Err(Error::State) + }) +} + +/// A producer that records `tag` when the render thread runs it and fails. +fn recording_producer(order: Arc>>, tag: &'static str) -> Producer { + Arc::new(move |_time: Rational, _params: &VideoTicketParams| { + order + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(tag); + Err(Error::State) + }) +} + +/// A no-footage job (no prefetch side effects) carrying `produce`. +fn scheduled_job(produce: Producer, schedule: JobSchedule) -> Job { + Job { + node_identity: 0, + time: Rational::new(0, 1), + params: Arc::new(base_params(Rational::new(0, 1))), + audio: None, + produce, + done: Box::new(|_result: TicketResult| {}), + schedule, + } +} + +fn open_gate(gate: &Arc<(Mutex, Condvar)>) { + let (open, work) = &**gate; + *open.lock().unwrap_or_else(|e| e.into_inner()) = true; + work.notify_all(); +} + +/// M4: priorities are real end to end, not just a `VecDeque` sort. With +/// the render thread parked on an in-flight playback frame and the queue +/// full of background work, a Seek (a) is accepted without blocking its +/// submitter — it over-admits past the bound — and (b) runs before every +/// queued background job once the in-flight frame finishes. +#[test] +fn pipeline_orders_seek_ahead_of_background_end_to_end() { + let _lock = lock(); + let backend = PipelineBackend::new().expect("pipeline backend starts"); + let order = Arc::new(Mutex::new(Vec::<&'static str>::new())); + let started = Arc::new((Mutex::new(false), Condvar::new())); + let release = Arc::new((Mutex::new(false), Condvar::new())); + + let produce = parked_producer( + started.clone(), + release.clone(), + order.clone(), + "in-flight", + ); + assert!( + backend.try_post(scheduled_job(produce, JobSchedule::playback(0, 0, 0))), + "the in-flight playback frame is accepted" + ); + wait_until("the in-flight playback frame to start", &mut || { + *started.0.lock().unwrap_or_else(|e| e.into_inner()) + }); + + // Fill the bounded queue to capacity with background work. + for _ in 0..RENDER_QUEUE_CAP { + assert!(backend.try_post(scheduled_job( + recording_producer(order.clone(), "background"), + JobSchedule::background(), + ))); + } + assert_eq!(backend.queue_depth(), RENDER_QUEUE_CAP); + + // The seek must not block on the full queue: it over-admits and sits + // in front of everything queued (the UI thread never stalls here). + assert!( + backend.post(scheduled_job( + recording_producer(order.clone(), "seek"), + JobSchedule::seek(), + )), + "the seek is accepted despite the full queue" + ); + assert_eq!( + backend.queue_depth(), + RENDER_QUEUE_CAP + 1, + "the seek over-admits instead of waiting for room" + ); + + open_gate(&release); + wait_until("every queued job to execute", &mut || { + backend.stats().executed == RENDER_QUEUE_CAP as u64 + 2 + }); + let mut expected = vec!["in-flight", "seek"]; + expected.extend(vec!["background"; RENDER_QUEUE_CAP]); + assert_eq!( + *order.lock().unwrap_or_else(|e| e.into_inner()), + expected, + "the seek preempts the queued background work end to end" + ); + backend.shutdown(); +} + +/// M4: the read-ahead claim is falsifiable here. The in-flight job parks +/// the render thread, so the playback post's prefetch has the decode +/// thread to itself; when the render request then runs it must reuse that +/// decoded frame — a second decode or zero LRU hits fails the test. +#[test] +fn pipeline_prefetch_is_the_frame_the_render_request_uses() { + let _lock = lock(); + pin_legacy_working_space(); + let path = test_clip("prefetch_hit"); + let backend = PipelineBackend::new().expect("pipeline backend starts"); + let started = Arc::new((Mutex::new(false), Condvar::new())); + let release = Arc::new((Mutex::new(false), Condvar::new())); + let produce = parked_producer( + started.clone(), + release.clone(), + Arc::new(Mutex::new(Vec::new())), + "parking", + ); + assert!(backend.try_post(scheduled_job(produce, JobSchedule::seek()))); + wait_until("the parking job to start", &mut || { + *started.0.lock().unwrap_or_else(|e| e.into_inner()) + }); + + let params = Arc::new(VideoTicketParams { + footage: Some((path.to_string_lossy().to_string(), 0)), + ..base_params(Rational::new(0, 1)) + }); + let service = backend.decode_service(); + let (done_tx, done_rx) = mpsc::channel(); + let produce: Producer = Arc::new(|time: Rational, params: &VideoTicketParams| { + oak_render::eval::render_produced_frame(time, params).map(TicketPayload::Video) + }); + let playback = Job { + node_identity: 1, + time: Rational::new(0, 1), + params, + audio: None, + produce, + done: Box::new(move |result: TicketResult| { + let _ = done_tx.send(matches!(result, Ok(TicketPayload::Video(_)))); + }), + schedule: JobSchedule::playback(0, 0, 0), + }; + assert!(backend.post(playback), "the playback frame is accepted"); + assert!( + service.wait_idle(), + "the read-ahead decode completes while the render thread is parked" + ); + let after_prefetch = service.stats(); + assert_eq!(after_prefetch.prefetches, 1, "the post queued one read-ahead"); + assert_eq!(after_prefetch.decodes, 1, "the read-ahead decoded once"); + + open_gate(&release); + assert!( + done_rx + .recv_timeout(Duration::from_secs(60)) + .expect("the playback frame renders"), + "the playback frame produced a video payload" + ); + let stats = service.stats(); + assert!( + stats.lru_hits >= 1, + "the render request reused the prefetched frame (lru_hits {})", + stats.lru_hits + ); + assert_eq!( + stats.decodes, 1, + "the render request must not decode the frame a second time" + ); + backend.shutdown(); + let _ = std::fs::remove_file(&path); +} + +/// M4: cancelling a preview window drops only that window's queued frame. +/// Two viewers can queue the same frame number in the same version; the +/// cancel must match on the sequence identity too. +#[test] +fn pipeline_cancel_preview_frame_matches_the_sequence() { + let _lock = lock(); + let backend = PipelineBackend::new().expect("pipeline backend starts"); + let started = Arc::new((Mutex::new(false), Condvar::new())); + let release = Arc::new((Mutex::new(false), Condvar::new())); + let produce = parked_producer( + started.clone(), + release.clone(), + Arc::new(Mutex::new(Vec::new())), + "parking", + ); + assert!(backend.try_post(scheduled_job(produce, JobSchedule::seek()))); + wait_until("the parking job to start", &mut || { + *started.0.lock().unwrap_or_else(|e| e.into_inner()) + }); + + let (tx, rx) = mpsc::channel(); + for identity in [1u64, 2] { + let tx = tx.clone(); + let produce: Producer = + Arc::new(|_time: Rational, _params: &VideoTicketParams| { + Ok(TicketPayload::Video(Texture::dummy())) + }); + let job = Job { + node_identity: identity, + time: Rational::new(0, 1), + params: Arc::new(base_params(Rational::new(0, 1))), + audio: None, + produce, + done: Box::new(move |result: TicketResult| { + let _ = tx.send((identity, result.is_ok())); + }), + // Same frame number, same version — only the sequence differs. + schedule: JobSchedule::playback(5, 0, 0), + }; + assert!(backend.post(job), "viewer {identity}'s frame is queued"); + } + backend.cancel_preview_frame(1, 5, 0); + + open_gate(&release); + let mut results = Vec::new(); + for _ in 0..2 { + results.push( + rx.recv_timeout(Duration::from_secs(60)) + .expect("both queued frames complete"), + ); + } + results.sort_unstable(); + assert_eq!( + results, + vec![(1, false), (2, true)], + "only the cancelled sequence's frame is dropped" + ); + backend.shutdown(); +} + /// A saturated render queue closes the decode service's prefetch gate: a /// speculative decode must be refused while a frame is in flight and the /// queue is full, and everything queued must still run once the in-flight diff --git a/docs/zh/plans/render-pipeline-threads.md b/docs/zh/plans/render-pipeline-threads.md index 57bb1b32a..47fe525d1 100644 --- a/docs/zh/plans/render-pipeline-threads.md +++ b/docs/zh/plans/render-pipeline-threads.md @@ -254,6 +254,57 @@ - 背压:三条队列都有界;上屏消费慢(暂停、窗口最小化)时 render 队列满 → 解码暂停预取;导出时上屏队列直通导出消费者,不存在"没人收"的积压。 +> **M4 落地回填(2026-09-15)**: +> +> - 渲染队列按 `JobSchedule.priority` 排序(Seek > Playback > Background, +> 同类 FIFO);解码队列把 rendezvous `Request` 排在 `Prefetch` 之前 +> (渲染要帧绝不排在推测解码后面),`Sync` 屏障保持 FIFO。 +> - 播放预取:Pipeline 收到 Playback job 时立即按 montage/footage 推导 +> `DecodeRequest`(媒体时间/尺寸/F32 与 eval 完全一致,footage 路径取 +> `force_format.unwrap_or(F32)`),投给解码线程,于是帧 N 在渲染线程跑 +> GPU pass 时解码线程已在解帧 N+1。`pipeline_prefetch_is_the_frame_the_render_request_uses` +> 用"停在解码前"的确定性时序断言 LRU 命中且只解一次;`pipeline_playback_prefetches_ahead_of_the_render` +> 只证明"每个 post 恰好投递一次预取、每帧只解一次"。 +> - 上屏窗口:`PreviewWindow` 的槽位泛化为 `PreviewSlot::{Shm, Video}`, +> 线程管线的 `TicketPayload::Video` 直接入窗、由 `cpu_frame` 消费; +> `PipelineBackend::preview_window_capacity` 返回渲染队列余量,播放提交 +> 被限制在余量内不阻塞 UI;Seek 提交允许**超额插队**(队列满时也立即 +> 返回,插到 Background/Playback 之前),因为优先级排序只对已入队的 job +> 生效——若 Seek 被挡在门外,UI 会等一个导出/缓存帧跑完。`cancel_preview_frame` +> 丢弃已过 playhead 的排队帧,按 `(sequence, frame, version)` 全键匹配 +> (两个监视器可同帧号同版本,只按帧匹配会误杀另一序列)。 +> - **基准对比**(本机 release,`oak-render/examples/bench_playback`, +> 1080p/25fps MPEG-2 源、240 帧@480p(853×480) / 128 帧@1080p;两个后端 +> 统一输出 F32 帧、同尺寸,否则进程池的 BGRA8 槽位与管线的 in-process +> F32 帧字节量/路径不同,数字不可比): +> +> | 用例 | fps | 首帧延迟 | CPU(user+sys) | +> |---|---|---|---| +> | processes 480p(17 worker) | 48.3 | 1416 ms | 23.6 s | +> | pipeline 480p | 78.6 | 127 ms | 3.0 s | +> | processes 1080p(4 worker) | 49.2 | 396 ms | 12.0 s | +> | pipeline 1080p | 28.5 | 161 ms | 4.5 s | +> +> 结论:线程管线在代理尺寸下全面胜出,首帧延迟低一个数量级;1080p 满幅 +> 的峰值吞吐低于多 worker 进程池(单解码线程 vs 4 个并行 worker),但 +> CPU 低约 2.7×、首帧快约 2.5×,且 28 fps 足够覆盖 24/25 fps 实时播放。 +> 按 M4 验收口径记录:CPU 不升、首帧不劣化成立;"fps 不低于"在 1080p +> 峰值吞吐上不成立(代理播放成立)。 +> +> "main-heap frame copies" 只统计进程池 shm 路径(管线不经过它),两边 +> 的 0 都不构成"管线无堆拷贝"的证据——in-process 帧在 eval 缓存与 +> service LRU 边界仍有 `Frame.data` 深拷贝(M5 收窄),文档与测试不得 +> 再宣称管线零堆拷贝。 +> +> 根因是**这一刻解码仍是 CPU 软解**(M5 才上硬解):进程池的 4 个 worker +> 各自跑一个软解器并行解帧,线程管线按 §3.1 只有一条解码线程。**不要** +> 为此把解码线程池化:GPU 硬解单元(NVDEC/VAAPI/VideoToolbox)是单设备 +> 单提交队列(§1 的出发点),§3.6 的零拷贝导入还要求"解码表面与渲染共 +> 用同一片 GPU 内存",多解码线程只会争同一队列与显存池。M5 落地后单 +> 解码线程就是正确形态。提前提升软解吞吐的正确方向是让 FFmpeg 解码器 +> 自身多线程(frame/slice threading),而不是并列多个解码线程——本期 +> 不立项。 + ### 3.5 GPU 零拷贝与上屏 - 内置效果全 GPU:图求值产生的中间纹理全部是 `Texture::Gpu`,整个 resolve @@ -415,7 +466,7 @@ fallback。** 解码上传与上屏共用一层 `gpuinteop` 抽象,按后端 | **M1 线程管线骨架** | 解码/渲染/上屏三线程+三队列进 oak-render(`pipeline` 模块);RenderManager 增加线程后端,进程池后端保留,`OAK_PIPELINE=processes` 可回退 | 同一套渲染测试在两个后端下都绿(测试矩阵化);播放/seek/导出 smoke 等价 | | **M2 GPU 零拷贝** | 图内全程 `Texture::Gpu`(合成/转场/调整层不再逐帧回读);wgpu 29 统一 + 采用 gpui device(§3.5 攻关已回填);GPU 色彩管理(工作空间→输出规格→显示器 ICC 烘焙 3D LUT,GPU 执行);导出/缓存/OFX 三处边界显式回读;内置 YUV→RGB GPU pass(M5 解码导入的依赖项,解码接线随 M5) | 图播放路径 **GPU→CPU 回读为 0**(`oak_core::backend::gpu_transfer_counters` 计数断言,M1 帧缓存范式);`RenderedFrame::Gpu` + `to_display` 上屏在 adopted device 上零拷贝(app 测试);YUV→RGB pass 与 `colormath::yuv444p16_to_rgb_f32` 对拍;全 workspace 测试绿 | | **M3 OFX 独立进程** | `oak-worker --ofx-host` 单进程宿主 + `oak-render/ofxhost` 客户端(Pipeline 后端安装,首 job 惰性 spawn);PluginJob 经 NDJSON + 输入/输出 shm 槽;崩溃重生+在途 job 重投+三次熔断紫帧;`plugin_progress`/`plugin_cancel` 搬运(宿主即时 flush) | 植入确定性崩溃钩子:`--ofx-crash-once` 杀掉宿主 → 在途 job 重投成功;`--ofx-crash-always` 连续三次崩溃 → 客户端熔断、eval 紫帧;进度事件(含 0.5/1.0)到达 app 回调;取消 flag 语义单测(`oak-worker/tests/ofx_host.rs` + eval/ofx_host 单测) | -| **M4 流水线预取** | 调度层按 §3.4 投依赖窗口;背压策略 | 1080p 播放 CPU 占用不升、fps 不低于进程池后端;首帧延迟不劣化(基准对比留档) | +| **M4 流水线预取** | 渲染队列优先级(Seek>Playback>Background);Playback job 投递即预取 Footage 解码(解码队列 Request 优先于 Prefetch);`PreviewSlot` 泛化让线程管线入窗;队列余量作播放背压 + Seek 超额插队保证 UI 不停摆 | 确定性单测:停在解码前的时序下预取命中 LRU 且只解一次(可证伪"读前于渲染");Seek 在队列满时超额插队并按序优先执行;`cancel_preview_frame` 按 (sequence,frame,version) 全键匹配;播放预取 smoke(prefetches==decodes==帧数);`bench_playback` 双后端同格式(F32)对比留档(§3.4 回填:CPU 与首帧不劣化成立;1080p 峰值吞吐低于多 worker 池,代理尺寸胜出) | | **M5 GPU 解码零拷贝** | §3.6 表逐行落地:staging fallback 基线 → Linux NVDEC/VAAPI 导入 → Windows D3D11VA 导入 → macOS VideoToolbox 导入;FFmpeg 无 hwaccel 的组合才评估手写 GPU 解码 | 硬解路径 `HW_TRANSFERS` 计数归零(不再下载);逐平台导入开/关对比测试;每行独立 PR 可回退 | 依赖关系:M0a 独立;M0b 依赖 M0a;M1 依赖 M0a+M0b;M2 依赖 M1;M3 依赖