render: the M4 playback prefetch — dependency window, priorities and backpressure
docs/zh/plans/render-pipeline-threads.md M4: the thread pipeline now
keeps its decode thread ahead of the render thread and the app's
playback window consumes in-process frames.
- Render queue: priority-ordered by JobSchedule.priority (Seek >
Playback > Background, FIFO within a class), so interactive frames
jump playback exports/autocache. Seek posts may over-admit the bound:
priority only reorders queued jobs, so a full queue of background work
must not park the UI thread until an export frame finishes.
- Decode queue: rendezvous Requests are served ahead of queued
Prefetches (a frame the renderer needs never waits behind speculative
decodes); Sync barriers stay FIFO. The queue is a bounded
Mutex+Condvar structure, preserving the request backpressure and the
wait_idle contract.
- Playback read-ahead: a Playback job's footage decode requests are
derived from its montage/footage spec on post (same media time, size
and force_format.unwrap_or(F32) as the eval) and queued immediately,
so frame N+1 decodes while frame N runs its GPU passes.
- App window: PreviewWindow slots are generalized to
PreviewSlot::{Shm, Video}; the pipeline's in-process TicketPayload is
cached and consumed by cpu_frame exactly like a worker slot.
PipelineBackend::preview_window_capacity reports the render-queue
headroom, so playback posts are capped to what the queue can take;
cancel_preview_frame drops queued frames the playhead has passed,
matched on the full (sequence, frame, version) key so one monitor's
window never drops the other sequence's same-numbered frame.
- Tests: decode-queue preemption/FIFO, render-queue ordering, request
derivation, and deterministic end-to-end M4 tests: a prefetch that
must be reused by the render request (LRU hit, single decode — the
read-ahead claim is falsifiable), a parked-render-thread priority test
where a full queue of background work still lets a Seek over-admit and
run first, and a sequence-aware cancel test. The playback prefetch
smoke asserts prefetches == distinct decodes == frames; it does not
claim zero heap copies (Frame.data is deep-copied at the eval-cache
and service-LRU boundaries today).
- bench_playback gains a pipeline mode with CPU (self+children) and
first-frame latency; both backends now produce F32 frames so the
comparison is like-for-like. The §3.4 backfill records the numbers:
at the proxy size the pipeline is faster with a lower first frame; at
1080p peak throughput is below the multi-worker pool, but that is an
artifact of the decode still being CPU software (M5), not a case for
pooling decode threads — GPU decode is a single device/queue and the
zero-copy import shares one GPU memory pool, so the single decode
thread stays the target shape.
This commit is contained in:
@@ -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<i64, ShmFrameRef>,
|
||||
slots: BTreeMap<i64, PreviewSlot>,
|
||||
}
|
||||
|
||||
/// 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<u64> = Vec::new();
|
||||
let mut stale_keys: Vec<(u64, i64, u64)> = Vec::new();
|
||||
let mut stale_slots: Vec<ShmFrameRef> = Vec::new();
|
||||
let mut stale_slots: Vec<PreviewSlot> = 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<ShmFrameRef>)> = Vec::new();
|
||||
let mut pending: Vec<(u64, Vec<PreviewSlot>)> = 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<ShmFrameRef> =
|
||||
let slots: Vec<PreviewSlot> =
|
||||
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<ShmFrameRef> = std::mem::take(&mut window.slots).into_values().collect();
|
||||
let slots: Vec<PreviewSlot> = 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,32 +15,41 @@
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
//! 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 -- <media> [frames] [workers] [long_edge]
|
||||
//! cargo run --release -p oakrender --example bench_playback -- <media> [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
|
||||
//! `<media>` 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<usize> = 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<f64> = 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::<f64>() / 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<usize>) {
|
||||
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<f64> = 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::<f64>() / 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<usize> = 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),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<T>(m: &Mutex<T>) -> 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: 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<DecodeRequest> {
|
||||
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<DecodeCommand>,
|
||||
shutting_down: bool,
|
||||
}
|
||||
|
||||
struct DecodeShared {
|
||||
queue: Mutex<DecodeQueue>,
|
||||
/// 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<DecodeCommand>,
|
||||
shared: Arc<DecodeShared>,
|
||||
gate: PrefetchGate,
|
||||
inner: Arc<DecodeInner>,
|
||||
handle: Mutex<Option<JoinHandle<()>>>,
|
||||
@@ -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<Self> {
|
||||
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<Result<Texture>> {
|
||||
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<DecodeCommand> {
|
||||
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<DecodeCommand>, inner: Arc<DecodeInner>, lru_capacity: usize) {
|
||||
fn decode_loop(shared: Arc<DecodeShared>, inner: Arc<DecodeInner>, lru_capacity: usize) {
|
||||
let mut lru: HashMap<DecodeRequest, (Texture, u64)> = 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<DecodeCommand>, inner: Arc<DecodeInner>, 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<usize> {
|
||||
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<PipelineInner>) {
|
||||
#[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<u8> = queue.iter().map(job_class).collect();
|
||||
assert_eq!(classes, vec![0, 1, 1, 2]);
|
||||
let frames: Vec<i64> = 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));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<Rational> = (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<bool>, Condvar)>,
|
||||
release: Arc<(Mutex<bool>, Condvar)>,
|
||||
order: Arc<Mutex<Vec<&'static str>>>,
|
||||
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<Mutex<Vec<&'static str>>>, 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<bool>, 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
|
||||
|
||||
@@ -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 依赖
|
||||
|
||||
Reference in New Issue
Block a user