diff --git a/crates/oakrender/examples/bench_playback.rs b/crates/oakrender/examples/bench_playback.rs new file mode 100644 index 000000000..6b3a898b4 --- /dev/null +++ b/crates/oakrender/examples/bench_playback.rs @@ -0,0 +1,208 @@ +// Oak Video Editor - Non-Linear Video Editor +// Copyright (C) 2026 Oak Team +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU General Public License for more details. +// +// You should have received a copy of the GNU General Public License +// 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. +//! +//! Run from the repo root: +//! +//! ```sh +//! cargo run --release -p oakrender --example bench_playback -- [frames] [workers] [long_edge] +//! ``` +//! +//! `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`. + +use std::path::PathBuf; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use oakcore_rs::Rational; +use oakrender::ipc::SLOT_FORMAT_BGRA8; +use oakrender::procpool::{DispatcherConfig, ProcessDispatcher}; +use oakrender::ticket::{TicketPayload, TicketResult, VideoTicketParams}; +use oakrender::worker::{Job, JobDispatch, JobSchedule}; + +/// Locate the oak-worker binary (see bench_process). +fn worker_bin() -> PathBuf { + if let Ok(p) = std::env::var("OAK_WORKER_BIN") { + return PathBuf::from(p); + } + if let Ok(exe) = std::env::current_exe() { + if let Some(examples) = exe.parent() { + if let Some(profile) = examples.parent() { + let candidate = profile.join(format!("oak-worker{}", std::env::consts::EXE_SUFFIX)); + if candidate.exists() { + return candidate; + } + } + } + } + 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); + + // 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); + + let config = DispatcherConfig { + worker_bin: Some(worker_bin()), + workers: workers.unwrap_or(0), + slots_per_worker: 8, + width, + height, + slot_format: SLOT_FORMAT_BGRA8, + 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}"); + + // One completion record per frame: (ticket/frame, submit, completion). + 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 job = Job { + node_identity: 1, + time: Rational::new(frame, 25), + params: Arc::new(VideoTicketParams { + viewer: 1, + 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(), + }), + audio: None, + produce: Arc::new(|_, _| { + Err(oakrender::error::Error::Failed( + "process backend does not use the in-process producer".into(), + )) + }), + done: Box::new(move |result: TicketResult| match result { + Ok(TicketPayload::ShmFrame(f)) => { + results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push((frame, start, Instant::now())); + dc.release_frame(&f); + } + Ok(TicketPayload::ShmAudio(a)) => { + dc.release_audio_frame(&a); + } + Ok(other) => { + eprintln!("unexpected payload: {other:?}"); + } + Err(e) => { + eprintln!("frame {frame} failed: {e}"); + } + }), + // Playback priority, the pre-render window's schedule. + schedule: JobSchedule::playback(frame, frame, 0), + }; + if !dispatcher.post(job) { + eprintln!("post refused at frame {frame}"); + break; + } + } + + // Pump until every completion has landed. + let deadline = Instant::now() + Duration::from_secs(300); + loop { + dispatcher.poll(); + 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 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", + oakrender::procpool::main_heap_frame_copies().to_string(), + ); + + dispatcher.shutdown(); +} diff --git a/src/oakui/real.rs b/src/oakui/real.rs index c5c2aef07..09bb8e220 100644 --- a/src/oakui/real.rs +++ b/src/oakui/real.rs @@ -3069,6 +3069,22 @@ impl AppEngine for RealEngine { }); return image; } + // During playback a cache miss must NOT block the UI thread on a + // synchronous render: the wait starves the tick loop that feeds + // the pre-render window, and the seek-priority ticket steals + // worker capacity from it, so the window never catches up (every + // painted frame blocked in `TicketArena::wait` — the choppy + // playback regression). Show the last displayed frame while the + // window warms up; the workers catch up within a few frames and + // the window then serves every playhead frame. + if self.clock(monitor).read(cx).transport.is_playing() { + if let Some(image) = cache + .get(&monitor) + .and_then(|e| e.proxy.as_ref().map(|p| p.image.clone())) + { + return image; + } + } // Both monitors render through the oakrender ticket arena (falling // back to the synthetic pattern when rendering is unavailable): the // program monitor renders the current sequence, the source monitor @@ -6008,6 +6024,80 @@ mod tests { } } + /// Playback pre-render window (M15 S2): during playback the playhead + /// frame must come from the worker-rendered shm slot cache, NOT the + /// synchronous render path (the main thread blocking in + /// `TicketArena::wait` on every painted frame is the "playback is + /// unusably choppy" regression — the UI must never sync-wait during + /// playback once the window has warmed up). + #[gpui::test] + async fn playback_window_supplies_playhead_frames(cx: &mut gpui::TestAppContext) { + let _media = media_lock(); + let _worker = WorkerBinGuard::set(); + let engine = cx.update(|cx| cx.new(|cx| RealEngine::create(cx))); + cx.update(|app| engine.update(app, |engine, cx| engine.new_project(cx))); + + let media = std::env::temp_dir().join(format!( + "oakapp_playback_window_{}.mp4", + std::process::id() + )); + oakcodec::testmedia::write_test_clip(&media, 64, 64, 250, 25) + .expect("generate playback test media"); + cx.update(|app| { + engine.update(app, |engine, cx| { + engine.import_footage(media.clone(), cx).expect("import") + }) + }); + let name = media.file_name().unwrap().to_string_lossy().into_owned(); + let entry = cx.read(|app| { + engine + .read(app) + .roots() + .into_iter() + .find(|e| e.name.as_ref() == name) + }) + .expect("imported footage is listed"); + cx.update(|app| { + engine.update(app, |engine, cx| { + engine.drop_footage(entry.id, TrackKind::Video, 0, Frame(0), cx) + }) + }); + + // Start playback and drive the tick loop: the window must fill. + cx.update(|app| engine.update(app, |engine, cx| engine.play(Monitor::Program, cx))); + let deadline = std::time::Instant::now() + Duration::from_secs(30); + let mut filled = 0usize; + let mut hit = false; + loop { + cx.update(|app| engine.update(app, |engine, cx| engine.tick(cx))); + let (slots, submitted) = cx.read(|app| { + let engine = engine.read(app); + let windows = engine.preview_windows.lock().unwrap(); + let window = windows.get(&Monitor::Program); + ( + window.map(|w| w.slots.len()).unwrap_or(0), + window.map(|w| w.submitted.len()).unwrap_or(0), + ) + }); + filled = filled.max(slots); + let playhead = cx.read(|app| engine.read(app).clock_frame(Monitor::Program, app)); + hit = hit + || cx + .update(|app| engine.update(app, |engine, _cx| engine.preview_slot_frame(Monitor::Program, playhead))) + .is_some(); + if hit { + break; + } + assert!( + std::time::Instant::now() < deadline, + "the playback window must supply playhead frames (peak cached {filled}, submitted {submitted})" + ); + std::thread::sleep(Duration::from_millis(10)); + } + assert!(filled > 0, "the window cached worker frames"); + let _ = std::fs::remove_file(&media); + } + // ---- M15 S3 audio prefetch ------------------------------------------ /// A lightweight `RenderedAudio` stand-in (the prefetch logic only