diff --git a/crates/oak-render/src/eval.rs b/crates/oak-render/src/eval.rs index 4fba38bb2..20454e43f 100644 --- a/crates/oak-render/src/eval.rs +++ b/crates/oak-render/src/eval.rs @@ -1059,6 +1059,23 @@ const MAX_CACHED_DECODERS: usize = 6; /// LRU tick source for [`DECODERS`]. static DECODER_TICK: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); +/// Count of [`render_footage_frame_inner`] entries, i.e. of real codec work +/// on the footage path (frame-cache hits and decode-service LRU hits do not +/// count). Visible for the decode-service tests, which use it to tell a +/// cached frame from a re-decode. +static DECODE_INVOCATIONS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + +/// The number of real footage decodes since the last +/// [`reset_decode_invocations`] (or process start). +pub fn decode_invocations() -> u64 { + DECODE_INVOCATIONS.load(std::sync::atomic::Ordering::Relaxed) +} + +/// Reset the [`decode_invocations`] counter to zero. +pub fn reset_decode_invocations() { + DECODE_INVOCATIONS.store(0, std::sync::atomic::Ordering::Relaxed); +} + /// Cap on decoded FRAMES cached per process (release builds of the graph /// renderer re-decode every frame the pre-render window pulls — each is a /// seek+decode on the shared decoder session, which serializes the graph @@ -1175,6 +1192,12 @@ fn open_decoder(filename: &str, stream_index: i32) -> Result Result { let started = std::time::Instant::now(); - let result = render_footage_frame_inner(filename, stream_index, time, size, format); + let result = match crate::pipeline::decode_service() { + Some(service) => { + let request = crate::pipeline::DecodeRequest { + filename: filename.to_string(), + stream_index, + time, + size, + format, + }; + match service.request(request) { + Some(result) => result, + // The service went away (shutdown) mid-flight: decode here + // rather than failing the frame. + None => render_footage_frame_inner(filename, stream_index, time, size, format), + } + } + None => render_footage_frame_inner(filename, stream_index, time, size, format), + }; if std::env::var_os("OAK_PERF").is_some() { eprintln!( "[decode] {:?} s{} time {}/{} ({:.3}s) size {:?} -> {:.3}s {:?}", @@ -1200,7 +1240,10 @@ pub fn render_footage_frame( result } -fn render_footage_frame_inner( +/// The synchronous decode behind [`render_footage_frame`] (frame-LRU → +/// decoder session → codec → F32 frame). `pub(crate)` because the decode +/// service runs exactly this on its own thread. +pub(crate) fn render_footage_frame_inner( filename: &str, stream_index: i32, time: Rational, @@ -1224,6 +1267,9 @@ fn render_footage_frame_inner( return Ok(Texture::wrap_frame(frame)); } } + // Past the frame cache: this call really goes to the codec (the + // decode-service tests read this counter to prove it). + DECODE_INVOCATIONS.fetch_add(1, std::sync::atomic::Ordering::Relaxed); let decoder = open_decoder(filename, stream_index)?; let params = RetrieveVideoParams { stream: CodecStream::with_block(filename.to_string(), stream_index, None), diff --git a/crates/oak-render/src/lib.rs b/crates/oak-render/src/lib.rs index c59a2570c..b790363cc 100644 --- a/crates/oak-render/src/lib.rs +++ b/crates/oak-render/src/lib.rs @@ -35,6 +35,8 @@ //! - `ipc` — render-worker NDJSON protocol + shm frame-slot transport //! - `scheduler` — preview frame scheduler (interleaved batch claims) //! - `procpool` — process-isolated render backend (M15) +//! - `pipeline` — the M1 thread pipeline backend (render thread + decode +//! thread; selected by `OAK_PIPELINE=threads`) //! - `error` — re-exports `oak_core::error` (the backend/color/texture/ //! frame value types moved to `oak-core` in the oak-common merge) @@ -51,6 +53,7 @@ pub mod frameio; pub mod handle; pub mod ipc; pub mod manager; +pub mod pipeline; pub mod procpool; pub mod scheduler; pub mod shaderfx; diff --git a/crates/oak-render/src/manager.rs b/crates/oak-render/src/manager.rs index adb92d1c5..31b9a6d06 100644 --- a/crates/oak-render/src/manager.rs +++ b/crates/oak-render/src/manager.rs @@ -51,6 +51,22 @@ pub enum RenderBackendChoice { /// Process-isolated oak-worker pool (crash isolation + shm frames). /// The M15 S2 default. Processes(DispatcherConfig), + /// The M1 thread pipeline ([`crate::pipeline::PipelineBackend`]): one + /// render thread running the same producer the other backends run, + /// plus the decode thread footage rendering rendezvouses with. + /// Selected by `OAK_PIPELINE=threads`; the process pool stays the + /// default until the M4 acceptance. + Pipeline, +} + +/// The backend `OAK_PIPELINE` asks for: `threads` selects the M1 thread +/// pipeline, anything else (including unset) the default process pool. +fn backend_choice_from_env() -> RenderBackendChoice { + if std::env::var("OAK_PIPELINE").as_deref() == Ok("threads") { + RenderBackendChoice::Pipeline + } else { + RenderBackendChoice::Processes(DispatcherConfig::default()) + } } /// The manager. Created by `oakrender_manager_init` (C ABI), accessed @@ -97,19 +113,25 @@ pub struct RenderManager { /// target of [`RenderManager::set_workspace_size`]). `None` on the /// inline (test) backend. process_pool: Option>, + /// The thread pipeline, when running the Pipeline backend (`None` on + /// the inline and process backends). Holding the `Arc` keeps the + /// render/decode threads alive for the manager's lifetime. + pipeline: Option>, } impl RenderManager { - /// Initialize the process-wide manager with the default backend — the - /// process-isolated oak-worker pool (M15 S2 mandate; idempotent; C++ - /// instance() semantics — only the main GUI process does this). + /// Initialize the process-wide manager with the backend `OAK_PIPELINE` + /// asks for — the process-isolated oak-worker pool (M15 S2 mandate) + /// unless the variable says `threads` (idempotent; C++ instance() + /// semantics — only the main GUI process does this). pub fn init() -> Result<()> { - Self::init_with_backend(RenderBackendChoice::Processes(DispatcherConfig::default())) + Self::init_with_backend(backend_choice_from_env()) } /// Initialize the process-wide manager with an explicit backend. /// `Threads` is the test-only inline backend (no worker threads, no - /// child processes); `Processes` spawns the oak-worker pool. + /// child processes); `Processes` spawns the oak-worker pool; + /// `Pipeline` starts the M1 render/decode threads. pub fn init_with_backend(choice: RenderBackendChoice) -> Result<()> { let mut guard = lock(&MANAGER); if guard.is_some() { @@ -149,17 +171,18 @@ impl RenderManager { eval::render_produced_frame(time, params) .map(crate::ticket::TicketPayload::Video) }); - let (dispatch, audio_dispatch, audio_fallback, process_pool): ( + let (dispatch, audio_dispatch, audio_fallback, process_pool, pipeline): ( Arc, Arc, Option>, Option>, + Option>, ) = match choice { RenderBackendChoice::Threads => { // Test-only inline backend: synchronous execution on the // calling thread, shared by video and audio. let inline = InlineDispatcher::sync(); - (inline.clone(), inline, None, None) + (inline.clone(), inline, None, None, None) } RenderBackendChoice::Processes(config) => { let dispatcher = ProcessDispatcher::new(config)?; @@ -177,7 +200,16 @@ impl RenderManager { // audio-side plugin crash now takes down the main process, // and the mix cost lands on the UI tick. let inline = InlineDispatcher::sync(); - (dispatcher.clone(), inline, None, Some(dispatcher)) + (dispatcher.clone(), inline, None, Some(dispatcher), None) + } + RenderBackendChoice::Pipeline => { + // M1 thread pipeline: the render thread executes the same + // producer (so graph mode included), and the decode thread + // the producer's footage decodes rendezvous with. Audio + // stays inline, as on the process backend. + let pipeline = crate::pipeline::PipelineBackend::new()?; + let inline = InlineDispatcher::sync(); + (pipeline.clone(), inline, None, None, Some(pipeline)) } }; let tickets = Arc::new(TicketArena::new_with_audio_fallback( @@ -200,6 +232,7 @@ impl RenderManager { inline_graph: Mutex::new(None), stopping: AtomicBool::new(false), process_pool, + pipeline, })); Ok(()) } @@ -213,6 +246,12 @@ impl RenderManager { lock(&self.inline_graph).clone() } + /// The thread pipeline, when the Pipeline backend is live (its queue + /// stats and decode service). `None` on the inline and process backends. + pub fn pipeline_backend(&self) -> Option> { + self.pipeline.clone() + } + /// Global access; `None` before init. pub fn global() -> Option> { lock(&MANAGER).clone() @@ -559,4 +598,28 @@ mod tests { assert_eq!(disk_cache_size().unwrap(), 0); std::env::remove_var("OAK_CONFIG_DIR"); } + + #[test] + fn backend_choice_env_selects_pipeline() { + let _guard = oak_core::commonutil::ENV_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); + std::env::remove_var("OAK_PIPELINE"); + assert!(matches!( + backend_choice_from_env(), + RenderBackendChoice::Processes(_) + )); + std::env::set_var("OAK_PIPELINE", "threads"); + assert!(matches!( + backend_choice_from_env(), + RenderBackendChoice::Pipeline + )); + // Any other value keeps the default process pool. + std::env::set_var("OAK_PIPELINE", "processes"); + assert!(matches!( + backend_choice_from_env(), + RenderBackendChoice::Processes(_) + )); + std::env::remove_var("OAK_PIPELINE"); + } } diff --git a/crates/oak-render/src/pipeline.rs b/crates/oak-render/src/pipeline.rs new file mode 100644 index 000000000..e4b69a22f --- /dev/null +++ b/crates/oak-render/src/pipeline.rs @@ -0,0 +1,960 @@ +// 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 . + +//! M1: the thread pipeline skeleton (design §3.1 / §3.3 / §3.4). +//! +//! Three threads, one queue pair, no more: +//! +//! - **UI / main thread** — the consumer. It submits through the ticket +//! arena and consumes finished frames. There is deliberately no +//! presentation thread (design §3.3): the arena's completion mechanism +//! *is* the present queue. A finished ticket (the `done` callback, or +//! `TicketArena::wait` returning) is the signal that a frame is ready to +//! be shown; the UI thread reads the payload on its next tick and +//! uploads/blits it. Nothing in this module uploads or presents anything. +//! - **Render thread** ([`PipelineBackend`], thread name `oak-render`) — +//! the pipeline backend's ONLY render executor: it takes jobs off the +//! render queue and runs the same producer the inline and process +//! backends run ([`crate::worker::execute_job`]), then completes the +//! ticket. Every GPU call the backend makes happens on this thread; +//! `GpuContext::shared()` is unchanged. A producer that submits follow-up +//! work re-posts it to this same thread. +//! - **Decode thread** ([`DecodeService`], thread name `oak-decode`) — the +//! optional process-wide decode service. [`crate::eval::render_footage_frame`] +//! renders through it by rendezvous while it is installed, which is what +//! moves real decoding (and the FFmpeg / NVDEC calls it makes) off the +//! render thread. With no service installed the decode path is exactly +//! the synchronous one it has always been. +//! +//! ### Queues and backpressure +//! +//! - **Render queue**: bounded at [`RENDER_QUEUE_CAP`] jobs, FIFO. A post +//! blocks once it is full (the arena submits from the UI thread, so a +//! stalled pipeline throttles submission instead of growing without +//! bound). The single exception is a post made *from the render thread*, +//! which can never wait for room it would have to make itself. +//! - **Decode queue**: bounded at [`DECODE_QUEUE_CAP`] commands. The +//! rendezvous request blocks when the queue is full — the same throttle, +//! one stage upstream of the render queue. +//! - **Prefetch** is best-effort: a prefetch is *dropped* (never queued) +//! when the pipeline is saturated. [`DecodeService`] carries a +//! [`PrefetchGate`] — the backend wires it to "the render queue has +//! room" — so speculative decoding never delays a foreground request. +//! Dropped prefetches are counted ([`DecodeStats::prefetch_refused`]). +//! +//! ### Present path +//! +//! Ticket completion → the UI thread's payload read → upload/blit, all on +//! the UI thread; the pipeline adds no thread and no queue for it. This is +//! the M1 reading of design §3.3 ("上屏队列 = 完成队列"). +//! +//! ### Relationship to the process pool +//! +//! [`PipelineBackend`] is an additional [`crate::worker::JobDispatch`] +//! implementation next to [`crate::procpool::ProcessDispatcher`]: the +//! process pool stays the default backend, and `OAK_PIPELINE=threads` (or +//! [`crate::manager::RenderBackendChoice::Pipeline`]) selects the thread +//! backend explicitly. Ticket arena, scheduler hints, snapshot push and +//! cancellation semantics are identical on both — only the executor +//! differs. + +use std::collections::{HashMap, VecDeque}; +use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}; +use std::sync::mpsc::{self, SyncSender, TrySendError}; +use std::sync::{Arc, Condvar, Mutex, MutexGuard, OnceLock}; +use std::thread::JoinHandle; + +use oak_core::texture::Texture; +use oak_core::{PixelFormat, Rational}; + +use crate::error::{Error, Result}; +use crate::worker::{execute_job, Job, JobDispatch}; + +/// Render-queue bound (jobs). Small on purpose: the queue exists to keep +/// the render thread fed across a UI tick, not to buffer a whole +/// pre-render window — a backlog this deep already means the pipeline is +/// slower than real time, and blocking the submitter is the honest +/// response (design §3.4). +pub const RENDER_QUEUE_CAP: usize = 8; + +/// Decode-command queue bound (requests + prefetch messages). Deeper than +/// the render queue because decodes are short and the service is the only +/// path to the media, but still small enough to bound latency. +pub const DECODE_QUEUE_CAP: usize = 16; + +/// Decoded-frame LRU capacity, in frames (~31 MB each at 1080p F32, so a +/// handful of frames is already hundreds of MB). Sized for a playback +/// window plus the montage baseline; the M2 decoder work lands here. +pub const DECODE_LRU_CAP: usize = 8; + +fn lock(m: &Mutex) -> MutexGuard<'_, T> { + m.lock().unwrap_or_else(|e| e.into_inner()) +} + +// --------------------------------------------------------------------------- +// Decode service +// --------------------------------------------------------------------------- + +/// One decode request: everything that identifies a decodable frame, and +/// the exact key of the service's frame cache. The size is part of the key +/// (the same media at a different target resolution is a different frame: +/// an interleaved viewer/proxy request must never reuse a wrongly-sized +/// buffer); `(0, 0)` means "native size". +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct DecodeRequest { + /// Footage filename. + pub filename: String, + /// Media stream index. + pub stream_index: i32, + /// Frame time (the frame containing this time is decoded). + pub time: Rational, + /// Target size, or `(0, 0)` for the media's native size. + pub size: (i32, i32), + /// Target pixel format. + pub format: PixelFormat, +} + +/// Decode-service counters (M1's evidence that the service really decodes +/// and really caches; the M2 work keeps them). +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct DecodeStats { + /// Rendezvous requests received. + pub requests: u64, + /// Requests served from the frame LRU (no decode). + pub lru_hits: u64, + /// Frames actually decoded (LRU misses and accepted prefetches). + pub decodes: u64, + /// Prefetch messages accepted for execution. + pub prefetches: u64, + /// Prefetch messages dropped by the gate or a full queue. + pub prefetch_refused: u64, + /// LRU entries evicted to stay under the capacity. + pub evictions: u64, + /// Failed decodes (reported to the requester, or swallowed for a + /// prefetch). + pub errors: u64, +} + +#[derive(Default)] +struct DecodeCounters { + requests: AtomicU64, + lru_hits: AtomicU64, + decodes: AtomicU64, + prefetches: AtomicU64, + prefetch_refused: AtomicU64, + evictions: AtomicU64, + errors: AtomicU64, +} + +impl DecodeCounters { + fn snapshot(&self) -> DecodeStats { + DecodeStats { + requests: self.requests.load(Ordering::Relaxed), + lru_hits: self.lru_hits.load(Ordering::Relaxed), + decodes: self.decodes.load(Ordering::Relaxed), + prefetches: self.prefetches.load(Ordering::Relaxed), + prefetch_refused: self.prefetch_refused.load(Ordering::Relaxed), + evictions: self.evictions.load(Ordering::Relaxed), + errors: self.errors.load(Ordering::Relaxed), + } + } +} + +/// Prefetch gate: `true` = prefetch work may be queued. The backend wires +/// this to "the render queue has room", so prefetch never competes with +/// foreground frames for the decode queue. +pub type PrefetchGate = Arc bool + Send + Sync>; + +/// A command for the decode thread. +enum DecodeCommand { + /// Rendezvous decode: the reply carries the frame or the decode error. + Request { + request: DecodeRequest, + reply: SyncSender>, + }, + /// Cache-fill only: the result is never delivered anywhere. + Prefetch { request: DecodeRequest }, + /// 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, +} + +struct DecodeInner { + counters: DecodeCounters, + lru_len: AtomicUsize, +} + +/// The optional process-wide decode service (design §3.3): one decode +/// thread, a bounded command queue, and a bounded decoded-frame LRU. +/// +/// Installed by [`PipelineBackend::new`] and removed by its shutdown; also +/// usable standalone ([`DecodeService::new`]), which is how the unit tests +/// drive it. [`crate::eval::render_footage_frame`] routes through it while +/// it is installed and falls back to the synchronous decode otherwise. +/// +/// The LRU lives on the decode thread (no lock: it is only ever touched +/// there); [`DecodeService::stats`] and [`DecodeService::lru_len`] read +/// shared counters, so they are safe from any thread. +pub struct DecodeService { + commands: SyncSender, + gate: PrefetchGate, + inner: Arc, + handle: Mutex>>, +} + +impl DecodeService { + /// Start a decode thread with an LRU of `lru_capacity` frames and + /// `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 inner = Arc::new(DecodeInner { + counters: DecodeCounters::default(), + lru_len: AtomicUsize::new(0), + }); + let service = Arc::new(Self { + commands: tx, + 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. + if let Ok(handle) = spawned { + *lock(&service.handle) = Some(handle); + } + service + } + + /// 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 + /// propagated as such. + /// + /// Blocks while the decode queue is full: that block is the backpressure + /// 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, + } + match rx.recv() { + Ok(result) => Some(result), + // The command was dropped (shutdown drained the queue) or the + // decode thread died mid-request: the caller falls back inline. + Err(_) => None, + } + } + + /// Queue a speculative decode of `request`; `false` when the gate says + /// no (or the queue is full / the service is down), in which case + /// nothing was decoded and nothing was queued. + pub fn prefetch(&self, request: DecodeRequest) -> bool { + if !(self.gate)() { + self.inner + .counters + .prefetch_refused + .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 + } + } + } + + /// Wait until every command sent before this call has been processed + /// (a barrier through the same FIFO 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() { + return false; + } + rx.recv().is_ok() + } + + /// The current counters. + pub fn stats(&self) -> DecodeStats { + self.inner.counters.snapshot() + } + + /// The number of frames currently in the LRU. + pub fn lru_len(&self) -> usize { + self.inner.lru_len.load(Ordering::Relaxed) + } + + /// 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. + pub fn shutdown(&self) { + let handle = lock(&self.handle).take(); + let Some(handle) = handle else { return }; + let _ = self.commands.send(DecodeCommand::Shutdown); + let _ = handle.join(); + } +} + +/// 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) { + let mut lru: HashMap = HashMap::new(); + let mut tick: u64 = 1; + while let Ok(command) = rx.recv() { + match command { + DecodeCommand::Request { request, reply } => { + let result = serve(&request, &mut lru, &mut tick, lru_capacity, &inner); + // The requester may have gone away; the frame is still + // cached, so the decode was not wasted. + let _ = reply.send(result); + } + DecodeCommand::Prefetch { request } => { + prefetch_into(&request, &mut lru, &mut tick, lru_capacity, &inner); + } + 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. + inner.lru_len.store(0, Ordering::Relaxed); +} + +fn serve( + request: &DecodeRequest, + lru: &mut HashMap, + tick: &mut u64, + lru_capacity: usize, + inner: &DecodeInner, +) -> Result { + inner.counters.requests.fetch_add(1, Ordering::Relaxed); + if let Some(texture) = lru_get(lru, request, tick) { + inner.counters.lru_hits.fetch_add(1, Ordering::Relaxed); + return Ok(texture); + } + let texture = decode(request, inner)?; + lru_insert(lru, request.clone(), texture.clone(), lru_capacity, inner, tick); + inner.lru_len.store(lru.len(), Ordering::Relaxed); + Ok(texture) +} + +fn prefetch_into( + request: &DecodeRequest, + lru: &mut HashMap, + tick: &mut u64, + lru_capacity: usize, + inner: &DecodeInner, +) { + if lru_get(lru, request, tick).is_some() { + // Already decoded: `lru_get` refreshed its recency, nothing to do. + return; + } + match decode(request, inner) { + Ok(texture) => { + lru_insert(lru, request.clone(), texture, lru_capacity, inner, tick); + inner.lru_len.store(lru.len(), Ordering::Relaxed); + } + // Best-effort: a prefetch failure is counted, never raised — no + // ticket is waiting on it, and the foreground request that follows + // will report the same error through its own channel. + Err(_) => {} + } +} + +/// The real decode: the same producer code the inline path runs, on the +/// decode thread. +fn decode(request: &DecodeRequest, inner: &DecodeInner) -> Result { + inner.counters.decodes.fetch_add(1, Ordering::Relaxed); + let result = crate::eval::render_footage_frame_inner( + &request.filename, + request.stream_index, + request.time, + request.size, + request.format, + ); + if result.is_err() { + inner.counters.errors.fetch_add(1, Ordering::Relaxed); + } + result +} + +/// LRU read + recency refresh. Returns the cached frame; the clone is the +/// price of handing a texture to the render thread while keeping the +/// service's copy (a GPU-handle texture makes this a token in M2 — the +/// texture type is already the seam for it). +fn lru_get( + lru: &mut HashMap, + key: &DecodeRequest, + tick: &mut u64, +) -> Option { + let entry = lru.get_mut(key)?; + let texture = entry.0.clone(); + entry.1 = *tick; + *tick += 1; + Some(texture) +} + +/// LRU insert under the capacity, evicting the least recently used entry. +fn lru_insert( + lru: &mut HashMap, + key: DecodeRequest, + texture: Texture, + capacity: usize, + inner: &DecodeInner, + tick: &mut u64, +) { + if capacity == 0 { + return; + } + while lru.len() >= capacity && !lru.contains_key(&key) { + let Some(victim) = lru + .iter() + .min_by_key(|(_, (_, t))| *t) + .map(|(k, _)| k.clone()) + else { + break; + }; + lru.remove(&victim); + inner.counters.evictions.fetch_add(1, Ordering::Relaxed); + } + lru.insert(key, (texture, *tick)); + *tick += 1; +} + +// --------------------------------------------------------------------------- +// Process-wide service slot (the `PLUGIN_EXECUTOR` pattern from eval.rs) +// --------------------------------------------------------------------------- + +static DECODE_SERVICE: OnceLock>>> = OnceLock::new(); + +fn decode_service_slot() -> &'static Mutex>> { + DECODE_SERVICE.get_or_init(|| Mutex::new(None)) +} + +/// Install (or clear) the process-wide decode service. Installed by +/// [`PipelineBackend::new`]; `None` restores the synchronous decode path +/// for every renderer in the process. +pub fn install_decode_service(service: Option>) { + *lock(decode_service_slot()) = service; +} + +/// The installed decode service, if any. `eval` consults this per footage +/// frame; `None` means "decode inline as before". +pub fn decode_service() -> Option> { + lock(decode_service_slot()).clone() +} + +// --------------------------------------------------------------------------- +// Render thread +// --------------------------------------------------------------------------- + +/// Pipeline counters (tests/reporting). +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct PipelineStats { + /// Jobs accepted into the render queue. + pub posted: u64, + /// Jobs executed on the render thread. + pub executed: u64, + /// Jobs drained unexecuted by shutdown (completed with + /// [`Error::State`]). + pub drained: u64, +} + +/// The thread-backed [`JobDispatch`]: one render thread draining a bounded +/// FIFO, and the decode service that feeds its footage decodes. +/// +/// Only one instance may be live per process at a time — [`PipelineBackend::new`] +/// installs its decode service into the process-wide slot, and a second +/// live backend would steal the first one's slot (the manager owns the +/// singleton in production; tests serialize on a lock). +pub struct PipelineBackend { + inner: Arc, +} + +struct PipelineInner { + queue: Mutex>, + /// A job arrived. + work: Condvar, + /// Room appeared (or shutdown started). + room: Condvar, + /// `queue` length, mirror of the queue for gate/stat reads without the + /// lock; shared with the decode service's prefetch gate. + depth: Arc, + stopping: AtomicBool, + /// The render thread's id, set by the thread itself on entry. + render_thread: OnceLock, + decode: Arc, + handle: Mutex>>, + posted: AtomicU64, + executed: AtomicU64, + drained: AtomicU64, +} + +impl PipelineBackend { + /// Start the render and decode threads and install the decode service. + pub fn new() -> Result> { + let depth = Arc::new(AtomicUsize::new(0)); + // Prefetch is allowed only while the render queue has room: it must + // never fill the decode queue behind a saturated pipeline. + let gate_depth = depth.clone(); + let gate: PrefetchGate = Arc::new(move || { + gate_depth.load(Ordering::Relaxed) < RENDER_QUEUE_CAP + }); + let decode = DecodeService::new(DECODE_LRU_CAP, gate); + let inner = Arc::new(PipelineInner { + queue: Mutex::new(VecDeque::new()), + work: Condvar::new(), + room: Condvar::new(), + depth, + stopping: AtomicBool::new(false), + render_thread: OnceLock::new(), + decode: decode.clone(), + handle: Mutex::new(None), + posted: AtomicU64::new(0), + executed: AtomicU64::new(0), + drained: AtomicU64::new(0), + }); + let backend = Arc::new(Self { inner: inner.clone() }); + let handle = std::thread::Builder::new() + .name("oak-render".into()) + .spawn(move || render_loop(inner)) + .map_err(|e| Error::Failed(format!("render thread spawn: {e}")))?; + *lock(&backend.inner.handle) = Some(handle); + install_decode_service(Some(decode)); + Ok(backend) + } + + /// The backend's decode service (for prefetch injection and stats). + pub fn decode_service(&self) -> Arc { + self.inner.decode.clone() + } + + /// Jobs waiting (not yet taken by the render thread). + pub fn queue_depth(&self) -> usize { + self.inner.depth.load(Ordering::Relaxed) + } + + /// Free render-queue slots — the decode service's prefetch gate. + pub fn queue_free(&self) -> usize { + RENDER_QUEUE_CAP.saturating_sub(self.queue_depth()) + } + + /// Counters. + pub fn stats(&self) -> PipelineStats { + PipelineStats { + posted: self.inner.posted.load(Ordering::Relaxed), + executed: self.inner.executed.load(Ordering::Relaxed), + drained: self.inner.drained.load(Ordering::Relaxed), + } + } + + /// Non-blocking enqueue: `false` when the queue is full or the backend + /// is stopping. The blocking path is [`JobDispatch::post`]; this exists + /// for callers that prefer to drop work over waiting (and for tests + /// that fill the queue deterministically). + pub fn try_post(&self, job: Job) -> bool { + self.push(job, false) + } + + fn push(&self, job: Job, blocking: bool) -> bool { + let inner = &self.inner; + let mut queue = lock(&inner.queue); + if inner.stopping.load(Ordering::Acquire) { + return false; + } + 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) { + while queue.len() >= RENDER_QUEUE_CAP { + if inner.stopping.load(Ordering::Acquire) { + return false; + } + queue = inner.room.wait(queue).unwrap_or_else(|e| e.into_inner()); + } + } + queue.push_back(job); + inner.depth.store(queue.len(), Ordering::Relaxed); + inner.posted.fetch_add(1, Ordering::Relaxed); + drop(queue); + inner.work.notify_one(); + true + } + + fn shutdown_impl(&self) { + let inner = &self.inner; + if inner.stopping.swap(true, Ordering::AcqRel) { + return; + } + // The job the render thread is running finishes (M1 does not + // interrupt an in-flight frame); everything still queued is + // delivered as cancelled, exactly like the inline dispatcher. + let jobs: Vec = { + let mut queue = lock(&inner.queue); + let jobs: Vec = queue.drain(..).collect(); + inner.depth.store(0, Ordering::Relaxed); + jobs + }; + inner.drained.fetch_add(jobs.len() as u64, Ordering::Relaxed); + inner.work.notify_all(); + inner.room.notify_all(); + for job in jobs { + (job.done)(Err(Error::State)); + } + if !is_render_thread(inner) { + if let Some(handle) = lock(&inner.handle).take() { + let _ = handle.join(); + } + } + // Uninstall first (no new request may be routed to a service that + // is about to stop), then stop the decode thread. + if decode_service().is_some_and(|s| Arc::ptr_eq(&s, &inner.decode)) { + install_decode_service(None); + } + inner.decode.shutdown(); + } +} + +impl Drop for PipelineBackend { + fn drop(&mut self) { + // Safety net: the last handle going away must not leave the render + // thread parked on an empty queue. + self.shutdown_impl(); + } +} + +impl JobDispatch for PipelineBackend { + fn post(&self, job: Job) -> bool { + self.push(job, true) + } + + fn shutdown(&self) { + self.shutdown_impl(); + } +} + +fn is_render_thread(inner: &PipelineInner) -> bool { + inner.render_thread.get() == Some(&std::thread::current().id()) +} + +/// The render thread body: take a job, run it, repeat. Exits when the +/// queue is empty *and* shutdown has been requested. +fn render_loop(inner: Arc) { + let _ = inner.render_thread.set(std::thread::current().id()); + loop { + let job = { + let mut queue = lock(&inner.queue); + loop { + if let Some(job) = queue.pop_front() { + inner.depth.store(queue.len(), Ordering::Relaxed); + inner.room.notify_all(); + break Some(job); + } + if inner.stopping.load(Ordering::Acquire) { + break None; + } + queue = inner.work.wait(queue).unwrap_or_else(|e| e.into_inner()); + } + }; + match job { + Some(job) => { + execute_job(job); + inner.executed.fetch_add(1, Ordering::Relaxed); + } + None => break, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use oak_core::texture::Frame; + + /// A unique clip per test (the process id separates test binaries, the + /// tag separates tests inside one binary — the decode caches are + /// process-wide). + fn test_clip(tag: &str) -> std::path::PathBuf { + let path = std::env::temp_dir().join(format!( + "oakrender_pipeline_{tag}_{}.mp4", + std::process::id() + )); + oak_codec::testmedia::write_test_clip(&path, 64, 64, 10, 10) + .expect("test clip generation"); + path + } + + /// Pin the working space to the legacy sRGB pass-through: these tests + /// assert the decoded pattern, not the color transform (the ACEScg + /// default would remap the values). + fn pin_legacy_working_space() { + oak_core::color::set_pipeline_color_settings( + oak_core::colormath::WorkingColorSpace::SrgbLegacy, + oak_core::colormath::OutputColorSpec::default(), + ); + } + + fn request(filename: &std::path::Path, time: Rational) -> DecodeRequest { + DecodeRequest { + filename: filename.to_string_lossy().to_string(), + stream_index: 0, + time, + size: (64, 64), + format: PixelFormat::F32, + } + } + + fn always() -> PrefetchGate { + Arc::new(|| true) + } + + fn frame_of(texture: &Texture) -> &Frame { + let Texture::Cpu(frame) = texture else { + panic!("decode produced a non-CPU texture"); + }; + frame + } + + /// The decoded test pattern: a red|blue split that steps with the frame + /// index — `oak_codec::testmedia` shifts it by `index * width / (2 * fps)` + /// columns (9 columns for frame 3 of this 64 px / 10 fps clip). The + /// sampled columns track that shift; MPEG-2 is lossy, so the assertions + /// use dominance with generous margins. + fn assert_known_pattern(frame: &Frame, index: i32, tag: &str) { + assert_eq!((frame.width, frame.height), (64, 64), "{tag}"); + assert_eq!(frame.format, PixelFormat::F32, "{tag}"); + let stride = frame.linesize_bytes(); + let shift = (index * frame.width / 20).rem_euclid(frame.width); + let read = |x: usize, y: usize| -> [f32; 4] { + let off = y * stride + x * 16; + let mut out = [0f32; 4]; + for i in 0..4 { + out[i] = f32::from_le_bytes(frame.data[off + i * 4..off + i * 4 + 4].try_into().unwrap()); + } + out + }; + // Sample the middle of each half: the generator's split has the red + // half where `(x + shift) % 64` is below 32. + let [r, g, b, a] = read((16 - shift).rem_euclid(64) as usize, 32); + assert!(r > 0.5 && g < 0.4 && b < 0.4, "{tag}: red half {r},{g},{b}"); + assert!(a > 0.9, "{tag}: opaque {a}"); + let [r, g, b, a] = read((48 - shift).rem_euclid(64) as usize, 32); + assert!(b > 0.5 && r < 0.4 && g < 0.4, "{tag}: blue half {r},{g},{b}"); + assert!(a > 0.9, "{tag}: opaque {a}"); + } + + /// The service decodes real media through the real codec path. + #[test] + fn request_decodes_real_media() { + pin_legacy_working_space(); + let path = test_clip("real"); + let service = DecodeService::new(DECODE_LRU_CAP, always()); + + let texture = service + .request(request(&path, Rational::new(0, 1))) + .expect("service available") + .expect("frame 0 decodes"); + assert_known_pattern(frame_of(&texture), 0, "frame 0"); + + let stats = service.stats(); + assert_eq!(stats.requests, 1); + assert_eq!(stats.decodes, 1, "one real decode"); + assert_eq!(stats.lru_hits, 0); + assert_eq!(stats.errors, 0); + assert_eq!(service.lru_len(), 1); + + service.shutdown(); + let _ = std::fs::remove_file(&path); + } + + /// Prefetch really decodes and really fills the LRU: the request that + /// follows is served from the cache with no decode at all. + #[test] + fn prefetch_then_request_hits_the_lru() { + pin_legacy_working_space(); + let path = test_clip("prefetch"); + let service = DecodeService::new(DECODE_LRU_CAP, always()); + + let req = request(&path, Rational::new(3, 10)); + assert!(service.prefetch(req.clone()), "prefetch accepted"); + assert!(service.wait_idle(), "decode thread processed the prefetch"); + let after_prefetch = service.stats(); + assert_eq!(after_prefetch.prefetches, 1); + assert_eq!(after_prefetch.decodes, 1, "the prefetch really decoded"); + assert_eq!(service.lru_len(), 1); + + // The rendezvous request for the same frame: served from the LRU, + // so the decode counter does NOT move (this is the assertion that + // distinguishes a real cache hit from a re-decode). + let texture = service + .request(req) + .expect("service available") + .expect("cached frame"); + assert_known_pattern(frame_of(&texture), 3, "prefetched frame"); + let after_request = service.stats(); + assert_eq!(after_request.lru_hits, 1); + assert_eq!(after_request.decodes, after_prefetch.decodes); + assert_eq!(after_request.requests, 1); + + service.shutdown(); + let _ = std::fs::remove_file(&path); + } + + /// Decode failures travel back to the requester unchanged. + #[test] + fn request_propagates_decode_errors() { + let service = DecodeService::new(DECODE_LRU_CAP, always()); + let req = DecodeRequest { + filename: "/definitely/not/here.mp4".into(), + stream_index: 0, + time: Rational::new(0, 1), + size: (64, 64), + format: PixelFormat::F32, + }; + let err = service + .request(req) + .expect("service available") + .err() + .expect("decoding a missing file must fail"); + let _ = err.code(); // an explainable error, not a panic + let stats = service.stats(); + assert_eq!(stats.errors, 1); + assert_eq!(stats.decodes, 1); + assert_eq!(service.lru_len(), 0, "failures are not cached"); + service.shutdown(); + } + + /// The LRU is bounded and evicts least-recently-used first (an evicted + /// frame must be decoded again). + #[test] + fn lru_evicts_bounded() { + pin_legacy_working_space(); + let path = test_clip("evict"); + let service = DecodeService::new(2, always()); + let time = |n: i64| request(&path, Rational::new(n, 10)); + + service.request(time(0)).unwrap().unwrap(); + service.request(time(1)).unwrap().unwrap(); + assert_eq!(service.lru_len(), 2); + assert_eq!(service.stats().evictions, 0); + + // Frame 2 evicts frame 0 (the least recently used). + service.request(time(2)).unwrap().unwrap(); + assert_eq!(service.lru_len(), 2, "capacity holds"); + assert_eq!(service.stats().evictions, 1); + let decodes = service.stats().decodes; + assert_eq!(decodes, 3); + + // Frame 0 was evicted: asking for it decodes again... + service.request(time(0)).unwrap().unwrap(); + assert_eq!(service.stats().decodes, decodes + 1); + assert_eq!(service.stats().lru_hits, 0); + // ...while frame 2 (still cached) does not. + service.request(time(2)).unwrap().unwrap(); + assert_eq!(service.stats().decodes, decodes + 1); + assert_eq!(service.stats().lru_hits, 1); + + service.shutdown(); + let _ = std::fs::remove_file(&path); + } + + /// The prefetch gate is consulted per message: a closed gate refuses + /// (and counts) without queueing anything. + #[test] + fn prefetch_gate_refuses_and_recovers() { + pin_legacy_working_space(); + let path = test_clip("gate"); + let open = Arc::new(AtomicBool::new(false)); + let gate_open = open.clone(); + let service = DecodeService::new(DECODE_LRU_CAP, Arc::new(move || { + gate_open.load(Ordering::Relaxed) + })); + + let req = request(&path, Rational::new(1, 10)); + assert!(!service.prefetch(req.clone()), "closed gate refuses"); + let stats = service.stats(); + assert_eq!(stats.prefetch_refused, 1); + assert_eq!(stats.prefetches, 0); + assert_eq!(stats.decodes, 0, "nothing was decoded behind the gate"); + + open.store(true, Ordering::Relaxed); + assert!(service.prefetch(req.clone()), "open gate accepts"); + assert!(service.wait_idle()); + let stats = service.stats(); + assert_eq!(stats.prefetches, 1); + assert_eq!(stats.decodes, 1); + + service.shutdown(); + let _ = std::fs::remove_file(&path); + } + + /// `wait_idle` is a real barrier through the command queue: every + /// prefetch sent before it has been decoded when it returns. + #[test] + fn wait_idle_barrier_covers_queued_commands() { + pin_legacy_working_space(); + let path = test_clip("barrier"); + let service = DecodeService::new(DECODE_LRU_CAP, always()); + for n in 0..4 { + assert!(service.prefetch(request(&path, Rational::new(n, 10)))); + } + assert!(service.wait_idle()); + let stats = service.stats(); + assert_eq!(stats.prefetches, 4); + assert_eq!(stats.decodes, 4); + assert_eq!(service.lru_len(), 4); + service.shutdown(); + let _ = std::fs::remove_file(&path); + } + + /// After shutdown the service reports itself gone (`None`), so the + /// eval path falls back to decoding inline instead of failing frames. + #[test] + fn shutdown_makes_the_service_unavailable() { + pin_legacy_working_space(); + let path = test_clip("shutdown"); + let service = DecodeService::new(DECODE_LRU_CAP, always()); + service.shutdown(); + service.shutdown(); // idempotent + assert!(service.request(request(&path, Rational::new(0, 1))).is_none()); + assert!(!service.prefetch(request(&path, Rational::new(0, 1)))); + assert!(!service.wait_idle()); + let _ = std::fs::remove_file(&path); + } +} diff --git a/crates/oak-render/src/worker.rs b/crates/oak-render/src/worker.rs index d05450597..c8d210be3 100644 --- a/crates/oak-render/src/worker.rs +++ b/crates/oak-render/src/worker.rs @@ -232,7 +232,10 @@ impl InlineDispatcher { } } -fn execute_job(job: Job) { +/// Run one job on the calling thread: the producer, then its completion +/// (used by the inline dispatcher and by [`crate::pipeline::PipelineBackend`], +/// which runs the same producer on its render thread). +pub(crate) fn execute_job(job: Job) { let result = catch_unwind(AssertUnwindSafe(|| (job.produce)(job.time, &job.params))) .unwrap_or_else(|_| Err(Error::Failed("frame producer panicked".into()))); (job.done)(result); diff --git a/crates/oak-render/tests/common/mod.rs b/crates/oak-render/tests/common/mod.rs index 0169e0d18..35f1843bd 100644 --- a/crates/oak-render/tests/common/mod.rs +++ b/crates/oak-render/tests/common/mod.rs @@ -38,12 +38,16 @@ impl ManagerGuard { /// Initialize the manager and hold the serialization lock. Uses the /// test-only inline backend so no oak-worker children are spawned. pub fn init() -> Self { + Self::init_with(oak_render::manager::RenderBackendChoice::Threads) + } + + /// Initialize the manager with an explicit backend choice and hold the + /// serialization lock. Callers that pass `Pipeline` get the M1 thread + /// pipeline; everything else behaves like [`ManagerGuard::init`]. + pub fn init_with(choice: oak_render::manager::RenderBackendChoice) -> Self { let guard = MANAGER_LOCK.lock().unwrap_or_else(|e| e.into_inner()); oak_render::manager::RenderManager::shutdown(); - oak_render::manager::RenderManager::init_with_backend( - oak_render::manager::RenderBackendChoice::Threads, - ) - .expect("manager init"); + oak_render::manager::RenderManager::init_with_backend(choice).expect("manager init"); Self { _guard: guard } } } diff --git a/crates/oak-render/tests/render_threads_test.rs b/crates/oak-render/tests/render_threads_test.rs new file mode 100644 index 000000000..3185ba969 --- /dev/null +++ b/crates/oak-render/tests/render_threads_test.rs @@ -0,0 +1,695 @@ +// 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 . + +//! The M1 thread pipeline: the same ticket stream on the inline (test) +//! backend and on the thread pipeline must produce byte-identical frames. +//! +//! Covers the pipeline backend selection (`OAK_PIPELINE=threads`), the +//! decode-service LRU under seek patterns, the prefetch gate under render +//! queue backpressure, and the decode-service install/uninstall lifecycle. +//! The manager singleton, the `OAK_PIPELINE` variable and the decode +//! service slot are process-wide, so every test here serializes on `LOCK` +//! and each creates and tears its manager down explicitly. + +use std::path::{Path, PathBuf}; +use std::sync::{mpsc, Arc, Condvar, Mutex, MutexGuard}; +use std::time::{Duration, Instant}; + +use oak_core::commonutil::ENV_TEST_LOCK; +use oak_core::texture::{Frame, Texture}; +use oak_core::{PixelFormat, Rational, TimeRange}; + +use oak_node::block::{clip_create, clip_input, ClipBlockBehavior}; +use oak_node::footage::FootageBehavior; +use oak_node::id::NodeId; +use oak_node::node::NodeCore; +use oak_node::project::Project; +use oak_node::sequence::SequenceBehavior; +use oak_node::track::{TrackBehavior, TrackListBehavior}; + +use oak_render::error::Error; +use oak_render::eval::{decode_invocations, reset_decode_invocations}; +use oak_render::manager::{RenderBackendChoice, RenderManager}; +use oak_render::pipeline::{ + decode_service, DecodeRequest, DecodeStats, PipelineBackend, PipelineStats, + RENDER_QUEUE_CAP, +}; +use oak_render::ticket::{ + Completion, MontageClip, Producer, TicketPayload, TicketResult, VideoTicketParams, +}; +use oak_render::worker::{Job, JobDispatch, JobSchedule}; + +mod common; + +/// Serializes this binary's tests: the manager singleton, the environment +/// variable, the decode-service slot and the eval decode caches are all +/// process-wide. +static LOCK: Mutex<()> = Mutex::new(()); + +fn lock() -> MutexGuard<'static, ()> { + LOCK.lock().unwrap_or_else(|e| e.into_inner()) +} + +/// A unique clip per test (the process id separates test binaries, the +/// tag separates tests inside one binary — the decode caches are +/// process-wide). +fn test_clip(tag: &str) -> PathBuf { + let path = std::env::temp_dir().join(format!( + "oakrender_threads_{tag}_{}.mp4", + std::process::id() + )); + oak_codec::testmedia::write_test_clip(&path, 64, 64, 10, 10) + .expect("test clip generation"); + path +} + +/// A second copy of `src` under a fresh name: the decode and eval frame +/// caches are keyed by filename, so the two backends must not share one +/// (otherwise the pipeline run would replay the inline run's cache instead +/// of decoding). +fn test_clip_copy(src: &Path, tag: &str) -> PathBuf { + let path = std::env::temp_dir().join(format!( + "oakrender_threads_{tag}_{}.mp4", + std::process::id() + )); + std::fs::copy(src, &path).expect("copy the test clip"); + path +} + +/// Pin the working space to the legacy sRGB pass-through: these tests +/// assert the decoded pattern, not the color transform (the ACEScg +/// default would remap the values). +fn pin_legacy_working_space() { + oak_core::color::set_pipeline_color_settings( + oak_core::colormath::WorkingColorSpace::SrgbLegacy, + oak_core::colormath::OutputColorSpec::default(), + ); +} + +fn base_params(time: Rational) -> VideoTicketParams { + VideoTicketParams { + viewer: 0, + project: String::new(), + time, + force_size: Some((64, 64)), + force_format: Some(PixelFormat::F32), + cache: None, + cache_dir: None, + cache_id: None, + cache_timebase: None, + footage: None, + montage: Vec::new(), + adjustments: Vec::new(), + } +} + +/// A one-clip montage ticket over `[0s, 1s)`. +fn montage_params(filename: &Path, time: Rational) -> VideoTicketParams { + VideoTicketParams { + montage: vec![MontageClip { + filename: filename.to_string_lossy().to_string(), + stream_index: 0, + in_time: Rational::new(0, 1), + out_time: Rational::new(1, 1), + media_in: Rational::new(0, 1), + gain: 1.0, + effects: Vec::new(), + }], + ..base_params(time) + } +} + +/// A viewer ticket: the manager's graph mode renders `viewer` of the +/// project whose uuid is `uuid` (armed via `set_inline_project`). +fn viewer_params(uuid: &str, viewer: u64, time: Rational) -> VideoTicketParams { + VideoTicketParams { + viewer, + project: uuid.to_string(), + ..base_params(time) + } +} + +/// Submit one video ticket and wait for its frame. The reserved id keeps +/// the arena slot alive past the completion; `result()` reaps it. +fn render_video(params: VideoTicketParams) -> Texture { + let manager = RenderManager::global().expect("manager installed"); + let id = manager.tickets.next_id(); + let (tx, rx) = mpsc::sync_channel::(1); + let done: Completion = Box::new(move |result: TicketResult| { + let _ = tx.send(result); + }); + manager.tickets.submit_video_with_id(id, params, done); + let result = rx + .recv_timeout(Duration::from_secs(60)) + .expect("ticket completed within 60s"); + let _ = manager.tickets.result(id); + match result { + Ok(TicketPayload::Video(texture)) => texture, + Ok(other) => panic!("unexpected ticket payload: {other:?}"), + Err(e) => panic!("ticket failed: {e}"), + } +} + +fn frame_of(texture: &Texture) -> &Frame { + let Texture::Cpu(frame) = texture else { + panic!("ticket produced a non-CPU texture"); + }; + frame +} + +/// Byte-for-byte frame equality with a first-difference report. +fn assert_same_frame(expected: &Frame, actual: &Frame, tag: &str) { + assert_eq!( + (expected.width, expected.height), + (actual.width, actual.height), + "{tag}: frame size" + ); + assert_eq!(expected.format, actual.format, "{tag}: pixel format"); + assert_eq!(expected.channels, actual.channels, "{tag}: channel count"); + assert_eq!(expected.data.len(), actual.data.len(), "{tag}: data length"); + if let Some(offset) = expected + .data + .iter() + .zip(actual.data.iter()) + .position(|(a, b)| a != b) + { + let stride = expected.linesize_bytes(); + let row = offset / stride; + let byte_column = offset % stride; + panic!( + "{tag}: pixel bytes differ at offset {offset} (row {row}, byte column {byte_column}): \ + expected {}, got {}", + expected.data[offset], actual.data[offset] + ); + } +} + +/// The decoded test pattern: a red|blue split that steps with the frame +/// index — `oak_codec::testmedia` shifts it by `index * width / (2 * fps)` +/// columns (9 columns for frame 3 of this 64 px / 10 fps clip). The +/// sampled columns track that shift; MPEG-2 is lossy, so the assertions +/// use dominance with generous margins. +fn assert_known_pattern(frame: &Frame, index: i32, tag: &str) { + assert_eq!((frame.width, frame.height), (64, 64), "{tag}"); + assert_eq!(frame.format, PixelFormat::F32, "{tag}"); + let stride = frame.linesize_bytes(); + let shift = (index * frame.width / 20).rem_euclid(frame.width); + let read = |x: usize, y: usize| -> [f32; 4] { + let off = y * stride + x * 16; + let mut out = [0f32; 4]; + for i in 0..4 { + out[i] = f32::from_le_bytes(frame.data[off + i * 4..off + i * 4 + 4].try_into().unwrap()); + } + out + }; + // Sample the middle of each half: the generator's split has the red + // half where `(x + shift) % 64` is below 32. + let [r, g, b, a] = read((16 - shift).rem_euclid(64) as usize, 32); + assert!(r > 0.5 && g < 0.4 && b < 0.4, "{tag}: red half {r},{g},{b}"); + assert!(a > 0.9, "{tag}: opaque {a}"); + let [r, g, b, a] = read((48 - shift).rem_euclid(64) as usize, 32); + assert!(b > 0.5 && r < 0.4 && g < 0.4, "{tag}: blue half {r},{g},{b}"); + assert!(a > 0.9, "{tag}: opaque {a}"); +} + +/// Poll `check` until it holds (10s cap) — the pipeline counters advance +/// on the render thread, so a fixed sleep is both slow and flaky. +fn wait_until(what: &str, check: &mut dyn FnMut() -> bool) { + let deadline = Instant::now() + Duration::from_secs(10); + if check() { + return; + } + while Instant::now() < deadline { + std::thread::sleep(Duration::from_millis(5)); + if check() { + return; + } + } + panic!("timed out waiting for {what}"); +} + +/// Render `times` on the inline backend (no extra threads, no children). +fn render_inline( + project: Option<&Arc>>, + times: &[Rational], + params: impl Fn(Rational) -> VideoTicketParams, +) -> Vec { + let _guard = common::ManagerGuard::init(); + if let Some(project) = project { + RenderManager::global() + .expect("manager installed") + .set_inline_project(project.clone()); + } + times + .iter() + .map(|&time| frame_of(&render_video(params(time))).clone()) + .collect() +} + +/// Render `times` on the thread pipeline, then drain and report its +/// counters. The manager is torn down before returning (the decode service +/// must be uninstalled with it). +fn render_pipeline( + project: Option<&Arc>>, + times: &[Rational], + params: impl Fn(Rational) -> VideoTicketParams, +) -> (Vec, PipelineStats, DecodeStats) { + let guard = common::ManagerGuard::init_with(RenderBackendChoice::Pipeline); + let manager = RenderManager::global().expect("manager installed"); + if let Some(project) = project { + manager.set_inline_project(project.clone()); + } + let backend = manager + .pipeline_backend() + .expect("thread pipeline selected"); + reset_decode_invocations(); + let frames: Vec = times + .iter() + .map(|&time| frame_of(&render_video(params(time))).clone()) + .collect(); + let rendered = frames.len() as u64; + wait_until("all pipeline jobs executed", &mut || { + backend.stats().executed == rendered && backend.queue_depth() == 0 + }); + let stats = backend.stats(); + let service = backend.decode_service(); + assert!(service.wait_idle(), "decode service drained"); + let decode = service.stats(); + drop(service); + drop(backend); + drop(manager); + drop(guard); + assert!( + decode_service().is_none(), + "manager shutdown uninstalls the decode service" + ); + (frames, stats, decode) +} + +/// A job whose producer always fails: fills the queue / proves that a +/// stopped backend refuses work. +fn filler_job() -> Job { + let params = Arc::new(montage_params( + Path::new("/definitely/not/here-filler.mp4"), + Rational::new(0, 1), + )); + let produce: Producer = + Arc::new(|_time: Rational, _params: &VideoTicketParams| -> TicketResult { + Err(Error::State) + }); + Job { + node_identity: 0, + time: Rational::new(0, 1), + params, + audio: None, + produce, + done: Box::new(|_result: TicketResult| {}), + schedule: JobSchedule::seek(), + } +} + +/// One sequence + one video track list with one track per clip +/// `(filename, [in, out))`. The LAST entry's track composites on top +/// (NLE stacking: the highest-numbered track is topmost). +/// +/// The project is initialized like a real one, so the root folder takes +/// the first arena slot: a ticket names its viewer by `NodeId::identity`, +/// and identity 0 is the ticket API's "no graph viewer" sentinel — a +/// sequence created into slot 0 would silently take the montage fall-back +/// instead of the graph. +fn build_project(clips: &[(&str, Rational, Rational)]) -> (Arc>, NodeId) { + pin_legacy_working_space(); + let project = Project::new(); + let seq; + { + let mut p = project.lock().unwrap(); + p.initialize().expect("initialize the project"); + let (score, sbehavior) = SequenceBehavior::create(); + seq = p.graph.add_node(score, sbehavior); + + let (tcore, tbehavior) = TrackListBehavior::create(); + let tl = p.graph.add_node(tcore, tbehavior); + + for &(path, in_, out) in clips { + let (tcore, tbehavior) = TrackBehavior::create(); + let track = p.graph.add_node(tcore, tbehavior); + + let mut footage = FootageBehavior::new(path); + footage.probe().expect("probe the generated clip"); + let footage = p.graph.add_node(NodeCore::new(), Box::new(footage)); + + let (ccore, cbehavior) = clip_create(); + let clip = p.graph.add_node(ccore, cbehavior); + p.graph + .connect(footage, clip, clip_input::TEXTURE_INPUT, -1) + .expect("connect footage to clip"); + + let clip_behavior = p + .graph + .get_mut(clip) + .unwrap() + .behavior + .as_any_mut() + .unwrap() + .downcast_mut::() + .expect("clip block"); + clip_behavior.core.range = TimeRange::new(in_, out); + + p.graph + .get_mut(track) + .unwrap() + .behavior + .as_any_mut() + .unwrap() + .downcast_mut::() + .expect("video track") + .append_block(clip); + p.graph + .get_mut(tl) + .unwrap() + .behavior + .as_any_mut() + .unwrap() + .downcast_mut::() + .expect("video track list") + .tracks + .push(track); + } + + p.graph + .get_mut(seq) + .unwrap() + .behavior + .as_any_mut() + .unwrap() + .downcast_mut::() + .expect("sequence") + .track_lists + .push(tl); + } + (project, seq) +} + +/// `OAK_PIPELINE=threads` selects the thread pipeline (the default stays +/// the process backend, which the manager-guard init below exercises as +/// the test-only inline choice). +#[test] +fn oak_pipeline_env_selects_the_thread_backend() { + let _lock = lock(); + let _env = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner()); + RenderManager::shutdown(); + std::env::set_var("OAK_PIPELINE", "threads"); + RenderManager::init().expect("manager init with OAK_PIPELINE=threads"); + { + let manager = RenderManager::global().expect("manager installed"); + let backend = manager + .pipeline_backend() + .expect("OAK_PIPELINE=threads selects the thread pipeline"); + assert_eq!(backend.queue_depth(), 0, "idle queue"); + assert_eq!(backend.queue_free(), RENDER_QUEUE_CAP); + let stats = backend.stats(); + assert_eq!( + (stats.posted, stats.executed, stats.drained), + (0, 0, 0), + "fresh pipeline counters" + ); + assert!( + decode_service().is_some(), + "the pipeline installs its decode service" + ); + } + RenderManager::shutdown(); + std::env::remove_var("OAK_PIPELINE"); + assert!( + decode_service().is_none(), + "shutdown uninstalls the decode service" + ); + let _guard = common::ManagerGuard::init(); + assert!( + RenderManager::global() + .unwrap() + .pipeline_backend() + .is_none(), + "the inline test backend runs no thread pipeline" + ); +} + +/// Six consecutive frames: the two backends must agree byte for byte and +/// every frame must be a real decode through the service. +#[test] +fn pipeline_matches_inline_pixels_across_consecutive_frames() { + let _lock = lock(); + pin_legacy_working_space(); + let inline_path = test_clip("consecutive_inline"); + let pipeline_path = test_clip_copy(&inline_path, "consecutive_pipeline"); + let times: Vec = (0..6).map(|n| Rational::new(n, 10)).collect(); + + let inline_frames = render_inline(None, ×, |time| montage_params(&inline_path, time)); + let (pipeline_frames, stats, decode) = + render_pipeline(None, ×, |time| montage_params(&pipeline_path, time)); + + assert_eq!((stats.posted, stats.executed, stats.drained), (6, 6, 0)); + assert_eq!(decode.requests, 6, "one decode request per frame"); + assert_eq!(decode.decodes, 6, "each frame decodes once"); + assert_eq!(decode.lru_hits, 0, "consecutive frames never repeat"); + assert_eq!(decode_invocations(), 6, "the service did the decoding"); + for (frame, time) in pipeline_frames.iter().zip(times.iter()) { + assert_eq!(frame.timestamp, *time, "the rendered frame keeps its time"); + } + for (index, (a, b)) in inline_frames.iter().zip(pipeline_frames.iter()).enumerate() { + assert_same_frame(a, b, &format!("frame {index}")); + } + assert_known_pattern(&inline_frames[0], 0, "inline frame 0"); + assert_known_pattern(&pipeline_frames[0], 0, "pipeline frame 0"); + assert!( + inline_frames.windows(2).any(|w| w[0].data != w[1].data), + "the generated clip really moves between frames" + ); + + let _ = std::fs::remove_file(&inline_path); + let _ = std::fs::remove_file(&pipeline_path); +} + +/// Out-of-order seeks with a repeat: the service's LRU must absorb the +/// repeated frame, and the pixels must still match the inline path. +#[test] +fn pipeline_seek_out_of_order_matches_inline() { + let _lock = lock(); + pin_legacy_working_space(); + let inline_path = test_clip("seek_inline"); + let pipeline_path = test_clip_copy(&inline_path, "seek_pipeline"); + let seeks: [i64; 5] = [7, 2, 5, 0, 5]; + let times: Vec = seeks.iter().map(|&n| Rational::new(n, 10)).collect(); + + let inline_frames = render_inline(None, ×, |time| montage_params(&inline_path, time)); + let (pipeline_frames, stats, decode) = + render_pipeline(None, ×, |time| montage_params(&pipeline_path, time)); + + assert_eq!((stats.posted, stats.executed, stats.drained), (5, 5, 0)); + assert_eq!(decode.requests, 5, "one decode request per seek"); + assert_eq!(decode.lru_hits, 1, "the repeated 5/10 seek hits the LRU"); + assert_eq!(decode.decodes, 4, "four distinct frames decode"); + assert_eq!(decode_invocations(), 4, "the LRU absorbed the repeat"); + for (index, (a, b)) in inline_frames.iter().zip(pipeline_frames.iter()).enumerate() { + assert_same_frame(a, b, &format!("seek {index}")); + } + assert_known_pattern(&inline_frames[3], 0, "inline seek to 0/10"); + assert_known_pattern(&pipeline_frames[3], 0, "pipeline seek to 0/10"); + + let _ = std::fs::remove_file(&inline_path); + let _ = std::fs::remove_file(&pipeline_path); +} + +/// The graph (viewer) path through the pipeline: the same node-graph +/// render as the inline backend — not a silent fall-back to a blank +/// generated frame (which the pattern assertions would catch). +#[test] +fn pipeline_viewer_ticket_matches_inline_pixels() { + let _lock = lock(); + let path = test_clip("viewer"); + let filename = path.to_string_lossy().to_string(); + let clip = (filename.as_str(), Rational::new(0, 1), Rational::new(1, 1)); + let (project, sequence) = build_project(&[clip]); + let uuid = project.lock().unwrap().uuid.clone(); + let viewer = sequence.identity(); + let times = [ + Rational::new(0, 1), + Rational::new(3, 10), + Rational::new(8, 10), + ]; + + let inline_frames = render_inline(Some(&project), ×, |time| { + viewer_params(&uuid, viewer, time) + }); + let (pipeline_frames, stats, decode) = render_pipeline(Some(&project), ×, |time| { + viewer_params(&uuid, viewer, time) + }); + + assert_eq!((stats.posted, stats.executed, stats.drained), (3, 3, 0)); + assert_eq!(decode.requests, 3, "the graph path decodes via the service"); + assert_known_pattern(&inline_frames[0], 0, "inline graph frame 0"); + assert_known_pattern(&pipeline_frames[0], 0, "pipeline graph frame 0"); + for (index, (a, b)) in inline_frames.iter().zip(pipeline_frames.iter()).enumerate() { + assert_same_frame(a, b, &format!("graph frame {index}")); + } + + let _ = std::fs::remove_file(&path); +} + +/// 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 +/// frame completes. +#[test] +fn pipeline_queue_backpressure_closes_the_prefetch_gate() { + let _lock = lock(); + assert!( + decode_service().is_none(), + "the decode service slot starts empty" + ); + let backend = PipelineBackend::new().expect("pipeline backend starts"); + + // The in-flight job parks in its producer until released; that is what + // lets this test fill the queue deterministically. + let started = Arc::new((Mutex::new(None::), Condvar::new())); + let release = Arc::new((Mutex::new(false), Condvar::new())); + let job_started = started.clone(); + let job_release = release.clone(); + let produce: Producer = Arc::new( + move |_time: Rational, _params: &VideoTicketParams| -> TicketResult { + { + let (name, work) = &*job_started; + *name.lock().unwrap_or_else(|e| e.into_inner()) = + Some(std::thread::current().name().unwrap_or_default().to_string()); + work.notify_all(); + } + let (released, work) = &*job_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()); + } + Err(Error::State) + }, + ); + let hold_job = Job { + node_identity: 0, + time: Rational::new(0, 1), + params: Arc::new(montage_params( + Path::new("/definitely/not/here-hold.mp4"), + Rational::new(0, 1), + )), + audio: None, + produce, + done: Box::new(|_result: TicketResult| {}), + schedule: JobSchedule::seek(), + }; + + assert!(backend.try_post(hold_job), "the in-flight job is accepted"); + wait_until("the render thread to pick up the in-flight job", &mut || { + started + .0 + .lock() + .unwrap_or_else(|e| e.into_inner()) + .is_some() + }); + let thread_name = started + .0 + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(); + assert_eq!( + thread_name.as_deref(), + Some("oak-render"), + "the render thread runs the producer" + ); + + let mut accepted = 0usize; + while accepted < RENDER_QUEUE_CAP && backend.try_post(filler_job()) { + accepted += 1; + } + assert_eq!(accepted, RENDER_QUEUE_CAP, "the bounded queue fills up"); + assert_eq!(backend.queue_depth(), RENDER_QUEUE_CAP); + assert_eq!(backend.queue_free(), 0, "the prefetch gate's depth reads 0"); + assert!( + !backend.try_post(filler_job()), + "the queue refuses overflow" + ); + + let service = backend.decode_service(); + let refused = service.prefetch(DecodeRequest { + filename: "/definitely/not/here-prefetch.mp4".to_string(), + stream_index: 0, + time: Rational::new(0, 1), + size: (64, 64), + format: PixelFormat::F32, + }); + assert!(!refused, "a saturated pipeline refuses prefetch"); + let decode = service.stats(); + assert_eq!(decode.prefetch_refused, 1, "the refusal is counted"); + assert_eq!(decode.prefetches, 0, "nothing was queued"); + assert_eq!(decode.decodes, 0, "nothing was decoded"); + assert!( + !backend.try_post(filler_job()), + "the refused prefetch made no room" + ); + + { + let (released, work) = &*release; + *released.lock().unwrap_or_else(|e| e.into_inner()) = true; + work.notify_all(); + } + wait_until("every queued job to execute", &mut || { + backend.stats().executed == RENDER_QUEUE_CAP as u64 + 1 + && backend.queue_depth() == 0 + }); + let stats = backend.stats(); + assert_eq!( + (stats.posted, stats.drained), + (RENDER_QUEUE_CAP as u64 + 1, 0), + "all posted jobs ran" + ); + + backend.shutdown(); + assert!( + !backend.try_post(filler_job()), + "shutdown rejects new work" + ); + assert!(decode_service().is_none(), "shutdown uninstalls the service"); +} + +/// The backend owns the process-wide decode service slot: it is installed +/// at startup and uninstalled on shutdown. +#[test] +fn pipeline_installs_and_uninstalls_the_decode_service() { + let _lock = lock(); + assert!( + decode_service().is_none(), + "the decode service slot starts empty" + ); + let backend = PipelineBackend::new().expect("pipeline backend starts"); + let installed = decode_service().expect("decode service installed"); + assert!(Arc::ptr_eq(&installed, &backend.decode_service())); + + backend.shutdown(); + assert!( + decode_service().is_none(), + "shutdown uninstalls the decode service" + ); + assert!( + !backend.try_post(filler_job()), + "shutdown rejects new work" + ); +}