diff --git a/crates/oak-render/examples/bench_playback.rs b/crates/oak-render/examples/bench_playback.rs index 269029d40..8b7dcb841 100644 --- a/crates/oak-render/examples/bench_playback.rs +++ b/crates/oak-render/examples/bench_playback.rs @@ -39,9 +39,7 @@ use std::time::{Duration, Instant}; use oak_core::{PixelFormat, Rational}; use oak_render::procpool::{DispatcherConfig, ProcessDispatcher}; -use oak_render::ticket::{ - Completion, Producer, 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 @@ -80,10 +78,7 @@ fn cpu_times() -> (f64, f64) { } let self_times = rusage(libc::RUSAGE_SELF); let children = rusage(libc::RUSAGE_CHILDREN); - ( - self_times.0 + children.0, - self_times.1 + children.1, - ) + (self_times.0 + children.0, self_times.1 + children.1) } /// One footage ticket over the whole timeline. @@ -141,10 +136,7 @@ fn report(entries: &[(i64, Instant, Instant)], start: Instant, elapsed: Duration ), ); } - report( - "cpu user + sys", - format!("{:.2} + {:.2} s", cpu.0, cpu.1), - ); + 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(), @@ -182,7 +174,12 @@ fn run_processes(media: &str, frames: usize, width: i32, height: i32, workers: O let job = Job { node_identity: 1, time: Rational::new(frame, 25), - params: Arc::new(footage_params(&media_clone, Rational::new(frame, 25), width, height)), + params: Arc::new(footage_params( + &media_clone, + Rational::new(frame, 25), + width, + height, + )), audio: None, produce: Arc::new(|_, _| { Err(oak_render::error::Error::Failed( @@ -191,10 +188,11 @@ fn run_processes(media: &str, frames: usize, width: i32, height: i32, workers: O }), done: Box::new(move |result: TicketResult| match result { Ok(TicketPayload::ShmFrame(f)) => { - results - .lock() - .unwrap_or_else(|e| e.into_inner()) - .push((frame, start, Instant::now())); + results.lock().unwrap_or_else(|e| e.into_inner()).push(( + frame, + start, + Instant::now(), + )); dc.release_frame(&f); } Ok(TicketPayload::ShmAudio(a)) => { @@ -208,6 +206,7 @@ fn run_processes(media: &str, frames: usize, width: i32, height: i32, workers: O } }), schedule: JobSchedule::playback(frame, frame, 0), + cancelled: None, }; if !dispatcher.post(job) { eprintln!("post refused at frame {frame}"); @@ -262,25 +261,30 @@ fn run_pipeline(media: &str, frames: usize, width: i32, height: i32) { 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())); + 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 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)), + params: Arc::new(footage_params( + media, + Rational::new(frame, 25), + width, + height, + )), audio: None, produce: producer, done, schedule: JobSchedule::playback(frame, frame, 0), + cancelled: None, }; // The blocking post is the pipeline's backpressure: once the render // queue is full the submitter waits (the app's window is capped by @@ -332,7 +336,9 @@ fn main() { .nth(4) .and_then(|s| s.parse().ok()) .unwrap_or(480); - let mode = std::env::args().nth(5).unwrap_or_else(|| "processes".to_string()); + 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). diff --git a/crates/oak-render/examples/bench_process.rs b/crates/oak-render/examples/bench_process.rs index c71f36be1..9ca61dccb 100644 --- a/crates/oak-render/examples/bench_process.rs +++ b/crates/oak-render/examples/bench_process.rs @@ -57,8 +57,7 @@ fn worker_bin() -> PathBuf { // target//examples/bench_process -> target//oak-worker if let Some(examples) = exe.parent() { if let Some(profile) = examples.parent() { - let candidate = - profile.join(format!("oak-worker{}", std::env::consts::EXE_SUFFIX)); + let candidate = profile.join(format!("oak-worker{}", std::env::consts::EXE_SUFFIX)); if candidate.exists() { return candidate; } @@ -91,9 +90,7 @@ fn main() { 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" - ); + println!("oak-worker pool: {worker_count} worker(s), {frames} x {width}x{height} BGRA8 frames"); // One completion record per frame: (frame number, wall-clock completion). let results = Arc::new(Mutex::new(Vec::<(i64, Instant)>::new())); @@ -152,6 +149,7 @@ fn main() { // Playback priority: the pre-render window schedule (seek/playback // prioritization is what the scheduler benchmark measures). schedule: JobSchedule::playback(frame, 0, 0), + cancelled: None, }; if !dispatcher.post(job) { eprintln!("post refused at frame {frame}"); @@ -175,7 +173,11 @@ fn main() { } let elapsed = start.elapsed(); - let mut entries: Vec<(i64, Instant)> = results.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); + let mut entries: Vec<(i64, Instant)> = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); entries.sort_by_key(|(id, _)| *id); let completed = entries.len(); let throughput = completed as f64 / elapsed.as_secs_f64(); @@ -195,18 +197,30 @@ fn main() { }; 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)); + report( + "throughput", + format!("{throughput:.1} fps ({:.1} ms/frame)", 1000.0 / throughput), + ); if !deltas.is_empty() { let mean = deltas.iter().sum::() / deltas.len() as f64; let p95 = deltas[((deltas.len() as f64 * 0.95) as usize).min(deltas.len() - 1)]; report("adjacent-frame delta (pairs)", deltas.len().to_string()); - report(" max", format!("{:.3} ms", deltas.last().unwrap() * 1000.0)); + report( + " max", + format!("{:.3} ms", deltas.last().unwrap() * 1000.0), + ); report(" mean", format!("{:.3} ms", mean * 1000.0)); report(" p95", format!("{:.3} ms", p95 * 1000.0)); } else { - report("adjacent-frame delta", "no adjacent pairs completed".to_string()); + report( + "adjacent-frame delta", + "no adjacent pairs completed".to_string(), + ); } - report("main-heap frame copies", oak_render::procpool::main_heap_frame_copies().to_string()); + report( + "main-heap frame copies", + oak_render::procpool::main_heap_frame_copies().to_string(), + ); dispatcher.shutdown(); } diff --git a/crates/oak-render/src/autocacher.rs b/crates/oak-render/src/autocacher.rs index a007bb126..fb81eaa1c 100644 --- a/crates/oak-render/src/autocacher.rs +++ b/crates/oak-render/src/autocacher.rs @@ -139,8 +139,15 @@ impl PreviewAutoCacher { } /// Start a caching job for `range` of `owner`. + /// + /// M4 priority fix: the range job is Background (§3.3 交互 > 播放 > + /// 导出 > 自动缓存). It used to go through `submit_video` (Seek), which + /// after the M4 seek over-admission let every autocache frame jump + /// ahead of playback *and* past the render-queue bound. The + /// interactive single-frame preview ([`PreviewAutoCacher::single_frame`]) + /// keeps Seek. fn start_range_job(&mut self, owner: u64, range: TimeRange) { - let id = self.arena.submit_video( + let id = self.arena.submit_video_background( VideoTicketParams { viewer: owner, project: String::new(), @@ -324,6 +331,44 @@ mod tests { (PreviewAutoCacher::new(arena), d) } + /// A dispatcher that records the posted schedules instead of running + /// them (the priority test only needs the posted `JobSchedule`). + #[derive(Default)] + struct RecordingDispatcher { + schedules: Mutex>, + } + + impl JobDispatch for RecordingDispatcher { + fn post(&self, job: crate::worker::Job) -> bool { + lock(&self.schedules).push(job.schedule); + true + } + + fn shutdown(&self) {} + } + + #[test] + fn range_jobs_are_background_and_single_frames_stay_seek() { + let d = Arc::new(RecordingDispatcher::default()); + let arena = Arc::new(TicketArena::new(d.clone(), frame_producer())); + let mut c = PreviewAutoCacher::new(arena); + c.attach(7).unwrap(); + c.on_cache_request(7, TimeRange::new(Rational::new(0, 1), Rational::new(10, 1))); + c.single_frame(Rational::new(0, 1)); + let schedules = lock(&d.schedules); + assert_eq!(schedules.len(), 2, "one range job + one single frame"); + assert_eq!( + schedules[0].priority, + crate::scheduler::FramePriority::Background, + "autocache must yield to playback (§3.3)" + ); + assert_eq!( + schedules[1].priority, + crate::scheduler::FramePriority::Seek, + "the interactive single-frame preview stays Seek" + ); + } + #[test] fn attach_detach_lifecycle() { let (mut c, d) = new_cacher(); @@ -418,7 +463,6 @@ mod tests { stop: AtomicU32, } - #[test] fn events_deliver_progress() { let (mut c, d) = new_cacher(); diff --git a/crates/oak-render/src/pipeline.rs b/crates/oak-render/src/pipeline.rs index ff4179827..5fbbccdab 100644 --- a/crates/oak-render/src/pipeline.rs +++ b/crates/oak-render/src/pipeline.rs @@ -96,10 +96,15 @@ pub const RENDER_QUEUE_CAP: usize = 8; /// 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; +/// Decoded-frame LRU capacity, in frames (~33 MB each at 1080p F32, so a +/// small hand-off buffer, not a cache of record). The eval-side +/// `decoded_frames` LRU (24 frames) is the cache of record; this one only +/// bridges the decode thread and the render request. Keeping it at the +/// read-ahead window bounds the double-cache overhead: at 1080p F32 the +/// service copy stays ~2 frames instead of 8 (~200 MB saved). A request +/// the hand-off missed is served from the eval cache without a decode — +/// only the deep copy is paid again. +pub const DECODE_LRU_CAP: usize = 2; fn lock(m: &Mutex) -> MutexGuard<'_, T> { m.lock().unwrap_or_else(|e| e.into_inner()) @@ -388,7 +393,10 @@ impl DecodeService { return false; } if self.enqueue(DecodeCommand::Prefetch { request }, false) { - self.inner.counters.prefetches.fetch_add(1, Ordering::Relaxed); + self.inner + .counters + .prefetches + .fetch_add(1, Ordering::Relaxed); true } else { self.inner @@ -506,7 +514,14 @@ fn serve( return Ok(texture); } let texture = decode(request, inner)?; - lru_insert(lru, request.clone(), texture.clone(), lru_capacity, inner, tick); + lru_insert( + lru, + request.clone(), + texture.clone(), + lru_capacity, + inner, + tick, + ); inner.lru_len.store(lru.len(), Ordering::Relaxed); Ok(texture) } @@ -670,9 +685,8 @@ impl PipelineBackend { // 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 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()), @@ -687,7 +701,9 @@ impl PipelineBackend { executed: AtomicU64::new(0), drained: AtomicU64::new(0), }); - let backend = Arc::new(Self { inner: inner.clone() }); + let backend = Arc::new(Self { + inner: inner.clone(), + }); let handle = std::thread::Builder::new() .name("oak-render".into()) .spawn(move || render_loop(inner)) @@ -797,7 +813,9 @@ impl PipelineBackend { inner.depth.store(0, Ordering::Relaxed); jobs }; - inner.drained.fetch_add(jobs.len() as u64, Ordering::Relaxed); + inner + .drained + .fetch_add(jobs.len() as u64, Ordering::Relaxed); inner.work.notify_all(); inner.room.notify_all(); for job in jobs { @@ -921,8 +939,7 @@ mod tests { "oakrender_pipeline_{tag}_{}.mp4", std::process::id() )); - oak_codec::testmedia::write_test_clip(&path, 64, 64, 10, 10) - .expect("test clip generation"); + oak_codec::testmedia::write_test_clip(&path, 64, 64, 10, 10).expect("test clip generation"); path } @@ -971,7 +988,9 @@ mod tests { 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[i] = f32::from_le_bytes( + frame.data[off + i * 4..off + i * 4 + 4].try_into().unwrap(), + ); } out }; @@ -981,7 +1000,10 @@ mod tests { 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!( + b > 0.5 && r < 0.4 && g < 0.4, + "{tag}: blue half {r},{g},{b}" + ); assert!(a > 0.9, "{tag}: opaque {a}"); } @@ -1056,8 +1078,7 @@ mod tests { let err = service .request(req) .expect("service available") - .err() - .expect("decoding a missing file must fail"); + .expect_err("decoding a missing file must fail"); let _ = err.code(); // an explainable error, not a panic let stats = service.stats(); assert_eq!(stats.errors, 1); @@ -1108,9 +1129,10 @@ mod tests { 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 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"); @@ -1136,7 +1158,10 @@ mod tests { fn wait_idle_barrier_covers_queued_commands() { pin_legacy_working_space(); let path = test_clip("barrier"); - let service = DecodeService::new(DECODE_LRU_CAP, always()); + // The barrier test needs four live entries; the production + // hand-off capacity is deliberately tiny (see DECODE_LRU_CAP), so + // this test sizes its own service. + let service = DecodeService::new(4, always()); for n in 0..4 { assert!(service.prefetch(request(&path, Rational::new(n, 10)))); } @@ -1158,7 +1183,9 @@ mod tests { 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 + .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); @@ -1195,6 +1222,7 @@ mod tests { distance: frame, version: 0, }, + cancelled: None, } } @@ -1302,15 +1330,29 @@ mod tests { 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)), + 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].time, + Rational::new(5, 2), + "media_in + (time - in)" + ); assert_eq!(requests[0].size, (64, 32)); assert_eq!(requests[0].format, PixelFormat::F32); diff --git a/crates/oak-render/src/ticket.rs b/crates/oak-render/src/ticket.rs index 8406e0b7f..422e422ff 100644 --- a/crates/oak-render/src/ticket.rs +++ b/crates/oak-render/src/ticket.rs @@ -465,6 +465,7 @@ impl TicketArena { let params = Arc::new(params); let producer = self.producer.clone(); let slot_done = slot.clone(); + let slot_cancelled = slot.clone(); let job = crate::worker::Job { node_identity: params.viewer, time: params.time, @@ -473,6 +474,13 @@ impl TicketArena { produce: producer, done: Box::new(move |result| slot_done.finish(result)), schedule, + // Audit B: let the dispatcher skip producing a frame whose + // ticket was cancelled after posting (exactly-once delivery is + // unchanged — `finish` still fires, here with the same + // `Error::State` a cancelled result would get). + cancelled: Some(Arc::new(move || { + slot_cancelled.cancel.load(Ordering::Acquire) + })), }; if !self.dispatch.post(job) { // Backend is gone (shutdown raced the submit): deliver now. @@ -507,11 +515,7 @@ impl TicketArena { /// Submit a Background-priority frame (M15 S2 exports/precache): the /// scheduler renders it whenever no Seek/Playback work is pending. - pub fn submit_video_background( - &self, - params: VideoTicketParams, - done: Completion, - ) -> TicketId { + pub fn submit_video_background(&self, params: VideoTicketParams, done: Completion) -> TicketId { let id = self.next_id(); self.submit_video_background_with_id(id, params, done) } @@ -590,6 +594,7 @@ impl TicketArena { let make_job = |slot_done: Arc| { let ap_job = ap.clone(); let ap_prod = ap.clone(); + let slot_cancelled = slot_done.clone(); let producer: Producer = Arc::new(move |_, _| eval::render_audio_samples(&ap_prod)); crate::worker::Job { node_identity: viewer, @@ -612,6 +617,9 @@ impl TicketArena { produce: producer, done: Box::new(move |result| slot_done.finish(result)), schedule: JobSchedule::seek(), + cancelled: Some(Arc::new(move || { + slot_cancelled.cancel.load(Ordering::Acquire) + })), } }; let job = make_job(slot.clone()); diff --git a/crates/oak-render/src/worker.rs b/crates/oak-render/src/worker.rs index c8d210be3..d0a392fc9 100644 --- a/crates/oak-render/src/worker.rs +++ b/crates/oak-render/src/worker.rs @@ -69,6 +69,11 @@ pub struct Job { pub done: Completion, /// Scheduler hints (M15 S2). Defaults to a Seek single-frame request. pub schedule: JobSchedule, + /// Mid-flight cancellation probe (audit B): when it returns true the + /// dispatcher must finish the job with `Error::State` without running + /// the producer. The ticket arena installs the slot's cancel atom; + /// hand-built jobs (tests, non-ticket callers) pass `None`. + pub cancelled: Option bool + Send + Sync>>, } /// Scheduler hints a posted job carries (M15 S2). The process dispatcher @@ -234,8 +239,14 @@ impl InlineDispatcher { /// 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). +/// which runs the same producer on its render thread). A job the arena +/// already cancelled (audit B) is short-circuited with `Error::State` — +/// a cancelled frame must not burn render/GPU work only to be discarded. pub(crate) fn execute_job(job: Job) { + if job.cancelled.as_ref().is_some_and(|cancelled| cancelled()) { + (job.done)(Err(Error::State)); + return; + } let result = catch_unwind(AssertUnwindSafe(|| (job.produce)(job.time, &job.params))) .unwrap_or_else(|_| Err(Error::Failed("frame producer panicked".into()))); (job.done)(result); @@ -302,8 +313,8 @@ impl GraphSnapshotStore { "oakrender-snapshots-{}-{:x}", std::process::id(), { - use std::sync::atomic::{AtomicU64, Ordering}; - static SEQ: AtomicU64 = AtomicU64::new(0); + use std::sync::atomic::{AtomicU64, Ordering}; + static SEQ: AtomicU64 = AtomicU64::new(0); SEQ.fetch_add(1, Ordering::Relaxed) } )); @@ -350,9 +361,10 @@ impl GraphSnapshotStore { // Atomic staging: temp file + rename. The rename is a single // directory entry swap, so a concurrent worker load observes // either the old file or the complete new one. - let tmp = self - .dir - .join(format!("graph-{uuid}-{revision}.{}.tmp", std::process::id())); + let tmp = self.dir.join(format!( + "graph-{uuid}-{revision}.{}.tmp", + std::process::id() + )); std::fs::write(&tmp, &xml) .map_err(|e| Error::Failed(format!("write snapshot temp: {e}")))?; if let Err(e) = std::fs::rename(&tmp, &path) { @@ -390,9 +402,10 @@ impl GraphSnapshotStore { let path = self.dir.join(format!("graph-{uuid}-{revision}.xml")); let path_str = path.to_string_lossy().into_owned(); // Atomic staging: temp file + rename (see [`acquire`]). - let tmp = self - .dir - .join(format!("graph-{uuid}-{revision}.{}.tmp", std::process::id())); + let tmp = self.dir.join(format!( + "graph-{uuid}-{revision}.{}.tmp", + std::process::id() + )); std::fs::write(&tmp, &xml) .map_err(|e| Error::Failed(format!("write snapshot temp: {e}")))?; if let Err(e) = std::fs::rename(&tmp, &path) { @@ -466,13 +479,13 @@ impl Default for GraphSnapshotStore { #[cfg(test)] mod tests { - use super::*; - use std::sync::mpsc; - use std::time::Duration; + use super::*; + use std::sync::mpsc; + use std::time::Duration; - use oak_core::texture::Texture; + use oak_core::texture::Texture; - fn job(tag: u64, tx: mpsc::Sender, gate: Option>) -> Job { + fn job(tag: u64, tx: mpsc::Sender, gate: Option>) -> Job { let produce: Producer = Arc::new(move |_, _| { if let Some(g) = &gate { if g.load(Ordering::Acquire) { @@ -505,6 +518,7 @@ mod tests { let _ = tx.send(tag); }), schedule: JobSchedule::seek(), + cancelled: None, } } @@ -552,7 +566,8 @@ mod tests { let (tx, rx) = mpsc::channel(); for _ in 0..4 { let tx = tx.clone(); - let p: Producer = Arc::new(|_, _| Ok(crate::ticket::TicketPayload::Video(Texture::dummy()))); + let p: Producer = + Arc::new(|_, _| Ok(crate::ticket::TicketPayload::Video(Texture::dummy()))); d.post(Job { node_identity: 1, time: Rational::new(0, 1), @@ -576,6 +591,7 @@ mod tests { let _ = tx.send(r.is_err()); }), schedule: JobSchedule::seek(), + cancelled: None, }); } drop(tx); @@ -585,7 +601,10 @@ mod tests { delivered.push(err); } assert_eq!(delivered.len(), 4, "all queued completions fire"); - assert!(delivered.iter().all(|&e| e), "queued jobs cancel at shutdown"); + assert!( + delivered.iter().all(|&e| e), + "queued jobs cancel at shutdown" + ); } #[test] @@ -594,7 +613,8 @@ mod tests { let (tx, rx) = mpsc::channel(); let tx1 = tx.clone(); let boom: Producer = Arc::new(|_, _| panic!("boom")); - let ok: Producer = Arc::new(|_, _| Ok(crate::ticket::TicketPayload::Video(Texture::dummy()))); + let ok: Producer = + Arc::new(|_, _| Ok(crate::ticket::TicketPayload::Video(Texture::dummy()))); let params = Arc::new(VideoTicketParams { viewer: 0, project: String::new(), @@ -620,6 +640,7 @@ mod tests { let _ = tx1.send(1u64); }), schedule: JobSchedule::seek(), + cancelled: None, }); d.post(Job { node_identity: 1, @@ -632,6 +653,7 @@ mod tests { let _ = tx.send(2u64); }), schedule: JobSchedule::seek(), + cancelled: None, }); let mut got = Vec::new(); while let Ok(v) = rx.recv_timeout(Duration::from_secs(5)) { @@ -681,10 +703,7 @@ mod tests { // First snapshot: the default project (working space ACEScg). let p1 = store.acquire(&project, 1).unwrap(); let before = std::fs::read_to_string(&p1).unwrap(); - assert!( - std::path::Path::new(&p1).exists(), - "snapshot file written" - ); + assert!(std::path::Path::new(&p1).exists(), "snapshot file written"); assert_eq!(store.refs(&p1), 1); // A settings change with no revision bump: plain acquire must NOT @@ -695,7 +714,11 @@ mod tests { } let p2 = store.acquire(&project, 1).unwrap(); assert_eq!(p1, p2, "same (uuid, revision) key reuses the file"); - assert_eq!(std::fs::read_to_string(&p1).unwrap(), before, "acquire never rewrites"); + assert_eq!( + std::fs::read_to_string(&p1).unwrap(), + before, + "acquire never rewrites" + ); store.release(&p2); // ...while acquire_rewrite rewrites it in place at the same key. diff --git a/crates/oak-render/tests/render_threads_test.rs b/crates/oak-render/tests/render_threads_test.rs index a29f644b6..a86790cba 100644 --- a/crates/oak-render/tests/render_threads_test.rs +++ b/crates/oak-render/tests/render_threads_test.rs @@ -25,6 +25,7 @@ //! and each creates and tears its manager down explicitly. use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{mpsc, Arc, Condvar, Mutex, MutexGuard}; use std::time::{Duration, Instant}; @@ -44,8 +45,7 @@ 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, + decode_service, DecodeRequest, DecodeStats, PipelineBackend, PipelineStats, RENDER_QUEUE_CAP, }; use oak_render::ticket::{ Completion, MontageClip, Producer, TicketPayload, TicketResult, VideoTicketParams, @@ -71,8 +71,7 @@ fn test_clip(tag: &str) -> PathBuf { "oakrender_threads_{tag}_{}.mp4", std::process::id() )); - oak_codec::testmedia::write_test_clip(&path, 64, 64, 10, 10) - .expect("test clip generation"); + oak_codec::testmedia::write_test_clip(&path, 64, 64, 10, 10).expect("test clip generation"); path } @@ -211,7 +210,8 @@ fn assert_known_pattern(frame: &Frame, index: i32, tag: &str) { 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[i] = + f32::from_le_bytes(frame.data[off + i * 4..off + i * 4 + 4].try_into().unwrap()); } out }; @@ -221,7 +221,10 @@ fn assert_known_pattern(frame: &Frame, index: i32, tag: &str) { 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!( + b > 0.5 && r < 0.4 && g < 0.4, + "{tag}: blue half {r},{g},{b}" + ); assert!(a > 0.9, "{tag}: opaque {a}"); } @@ -306,10 +309,9 @@ fn filler_job() -> Job { Path::new("/definitely/not/here-filler.mp4"), Rational::new(0, 1), )); - let produce: Producer = - Arc::new(|_time: Rational, _params: &VideoTicketParams| -> TicketResult { - Err(Error::State) - }); + let produce: Producer = Arc::new( + |_time: Rational, _params: &VideoTicketParams| -> TicketResult { Err(Error::State) }, + ); Job { node_identity: 0, time: Rational::new(0, 1), @@ -318,6 +320,7 @@ fn filler_job() -> Job { produce, done: Box::new(|_result: TicketResult| {}), schedule: JobSchedule::seek(), + cancelled: None, } } @@ -468,8 +471,18 @@ fn build_layered_project( // V1: A [0,1) + transition [0.5,1.5) + B [1,2). let (v1_core, v1_beh) = TrackBehavior::create(); let v1 = p.graph.add_node(v1_core, v1_beh); - let a = add_clip(&mut p.graph, first, Rational::new(0, 1), Rational::new(1, 1)); - let b = add_clip(&mut p.graph, second, Rational::new(1, 1), Rational::new(2, 1)); + let a = add_clip( + &mut p.graph, + first, + Rational::new(0, 1), + Rational::new(1, 1), + ); + let b = add_clip( + &mut p.graph, + second, + Rational::new(1, 1), + Rational::new(2, 1), + ); let (tcore, tbehavior) = oak_node::block::transition_create(); let transition = p.graph.add_node(tcore, tbehavior); { @@ -507,7 +520,12 @@ fn build_layered_project( // V2: C [0,2), overlapping the transition track. let (v2_core, v2_beh) = TrackBehavior::create(); let v2 = p.graph.add_node(v2_core, v2_beh); - let c = add_clip(&mut p.graph, below, Rational::new(0, 1), Rational::new(2, 1)); + let c = add_clip( + &mut p.graph, + below, + Rational::new(0, 1), + Rational::new(2, 1), + ); track_mut(&mut p.graph, v2).append_block(c); // V3: an adjustment layer [0,2) with an Opacity(0.75) chain. @@ -774,7 +792,8 @@ fn pipeline_graph_playback_has_zero_gpu_readbacks() { #[test] fn pipeline_layered_playback_has_zero_gpu_readbacks() { let _lock = lock(); - if oak_core::backend::shared_gpu_or_skip("the layered playback zero-readback assertion").is_none() + if oak_core::backend::shared_gpu_or_skip("the layered playback zero-readback assertion") + .is_none() { return; } @@ -899,10 +918,7 @@ fn parked_producer( while !*released { released = work.wait(released).unwrap_or_else(|e| e.into_inner()); } - order - .lock() - .unwrap_or_else(|e| e.into_inner()) - .push(tag); + order.lock().unwrap_or_else(|e| e.into_inner()).push(tag); Err(Error::State) }) } @@ -910,10 +926,7 @@ fn parked_producer( /// A producer that records `tag` when the render thread runs it and fails. fn recording_producer(order: Arc>>, tag: &'static str) -> Producer { Arc::new(move |_time: Rational, _params: &VideoTicketParams| { - order - .lock() - .unwrap_or_else(|e| e.into_inner()) - .push(tag); + order.lock().unwrap_or_else(|e| e.into_inner()).push(tag); Err(Error::State) }) } @@ -928,6 +941,7 @@ fn scheduled_job(produce: Producer, schedule: JobSchedule) -> Job { produce, done: Box::new(|_result: TicketResult| {}), schedule, + cancelled: None, } } @@ -950,12 +964,7 @@ fn pipeline_orders_seek_ahead_of_background_end_to_end() { 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", - ); + 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" @@ -1044,6 +1053,7 @@ fn pipeline_prefetch_is_the_frame_the_render_request_uses() { let _ = done_tx.send(matches!(result, Ok(TicketPayload::Video(_)))); }), schedule: JobSchedule::playback(0, 0, 0), + cancelled: None, }; assert!(backend.post(playback), "the playback frame is accepted"); assert!( @@ -1051,7 +1061,10 @@ fn pipeline_prefetch_is_the_frame_the_render_request_uses() { "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.prefetches, 1, + "the post queued one read-ahead" + ); assert_eq!(after_prefetch.decodes, 1, "the read-ahead decoded once"); open_gate(&release); @@ -1098,10 +1111,9 @@ fn pipeline_cancel_preview_frame_matches_the_sequence() { 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 produce: Producer = Arc::new(|_time: Rational, _params: &VideoTicketParams| { + Ok(TicketPayload::Video(Texture::dummy())) + }); let job = Job { node_identity: identity, time: Rational::new(0, 1), @@ -1113,6 +1125,7 @@ fn pipeline_cancel_preview_frame_matches_the_sequence() { }), // Same frame number, same version — only the sequence differs. schedule: JobSchedule::playback(5, 0, 0), + cancelled: None, }; assert!(backend.post(job), "viewer {identity}'s frame is queued"); } @@ -1135,6 +1148,48 @@ fn pipeline_cancel_preview_frame_matches_the_sequence() { backend.shutdown(); } +/// Audit B: a job whose ticket was cancelled after posting must not run +/// its producer. The arena installs `Job.cancelled` from the slot's cancel +/// atom; this hand-built job pins the dispatcher behaviour and the +/// exactly-once completion (`Error::State`). +#[test] +fn pipeline_skips_a_cancelled_job() { + let _lock = lock(); + let backend = PipelineBackend::new().expect("pipeline backend starts"); + let ran = Arc::new(AtomicBool::new(false)); + let flag = Arc::new(AtomicBool::new(true)); + let ran_producer = ran.clone(); + let produce: Producer = Arc::new(move |_time: Rational, _params: &VideoTicketParams| { + ran_producer.store(true, Ordering::Release); + Ok(TicketPayload::Video(Texture::dummy())) + }); + let probe = flag.clone(); + let (tx, rx) = mpsc::channel(); + let 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(move |result: TicketResult| { + let _ = tx.send(result.is_err()); + }), + schedule: JobSchedule::seek(), + cancelled: Some(Arc::new(move || probe.load(Ordering::Acquire))), + }; + assert!(backend.try_post(job), "the cancelled job is still accepted"); + assert!( + rx.recv_timeout(Duration::from_secs(10)) + .expect("the completion fires exactly once"), + "a cancelled job completes with Error::State" + ); + assert!( + !ran.load(Ordering::Acquire), + "the producer must not run for a cancelled job" + ); + 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 @@ -1158,8 +1213,12 @@ fn pipeline_queue_backpressure_closes_the_prefetch_gate() { 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()); + *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; @@ -1181,21 +1240,21 @@ fn pipeline_queue_backpressure_closes_the_prefetch_gate() { produce, done: Box::new(|_result: TicketResult| {}), schedule: JobSchedule::seek(), + cancelled: None, }; 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(); + 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"), @@ -1238,8 +1297,7 @@ fn pipeline_queue_backpressure_closes_the_prefetch_gate() { work.notify_all(); } wait_until("every queued job to execute", &mut || { - backend.stats().executed == RENDER_QUEUE_CAP as u64 + 1 - && backend.queue_depth() == 0 + backend.stats().executed == RENDER_QUEUE_CAP as u64 + 1 && backend.queue_depth() == 0 }); let stats = backend.stats(); assert_eq!( @@ -1249,11 +1307,11 @@ fn pipeline_queue_backpressure_closes_the_prefetch_gate() { ); backend.shutdown(); + assert!(!backend.try_post(filler_job()), "shutdown rejects new work"); assert!( - !backend.try_post(filler_job()), - "shutdown rejects new work" + decode_service().is_none(), + "shutdown uninstalls the service" ); - assert!(decode_service().is_none(), "shutdown uninstalls the service"); } /// The backend owns the process-wide decode service slot: it is installed @@ -1274,8 +1332,5 @@ fn pipeline_installs_and_uninstalls_the_decode_service() { decode_service().is_none(), "shutdown uninstalls the decode service" ); - assert!( - !backend.try_post(filler_job()), - "shutdown rejects new work" - ); + assert!(!backend.try_post(filler_job()), "shutdown rejects new work"); } diff --git a/crates/oak-worker/tests/procpool_integration.rs b/crates/oak-worker/tests/procpool_integration.rs index a432f26fd..dfedf82b9 100644 --- a/crates/oak-worker/tests/procpool_integration.rs +++ b/crates/oak-worker/tests/procpool_integration.rs @@ -42,9 +42,7 @@ use oak_render::ipc::SLOT_FORMAT_BGRA8; use oak_render::procpool::{ main_heap_frame_copies, reset_main_heap_frame_copies, DispatcherConfig, ProcessDispatcher, }; -use oak_render::ticket::{ - AudioTicketParams, TicketPayload, TicketResult, VideoTicketParams, -}; +use oak_render::ticket::{AudioTicketParams, TicketPayload, TicketResult, VideoTicketParams}; use oak_render::worker::{Job, JobDispatch, JobSchedule}; /// Serialize every test in this file (shared process environment + @@ -117,9 +115,13 @@ fn submit( )) }), done: Box::new(move |result| { - results.lock().unwrap_or_else(|e| e.into_inner()).push(result); + results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(result); }), schedule: JobSchedule::seek(), + cancelled: None, }; assert!(dispatcher.post(job), "post accepted while alive"); } @@ -166,8 +168,11 @@ fn pool_resize_drains_shrunk_workers_and_regrows() { let deadline = Instant::now() + Duration::from_secs(60); while completed < 12 { dispatcher.poll(); - let drained: Vec = - results.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); + let drained: Vec = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); for result in drained { let payload = result.expect("frame rendered"); let TicketPayload::ShmFrame(frame) = payload else { @@ -181,7 +186,10 @@ fn pool_resize_drains_shrunk_workers_and_regrows() { } std::thread::sleep(Duration::from_millis(2)); } - assert_eq!(completed, 12, "all tickets completed before the resize settles"); + assert_eq!( + completed, 12, + "all tickets completed before the resize settles" + ); // The pool reached the target (the shrink took effect). assert_eq!(dispatcher.worker_count(), 1); // Let the retired children exit naturally (their EOF is reaped on @@ -200,8 +208,11 @@ fn pool_resize_drains_shrunk_workers_and_regrows() { let results2 = Arc::new(Mutex::new(Vec::new())); submit(&dispatcher, &results2, 6, None); pump_until(&dispatcher, &results2, 6); - let frames: Vec = - results2.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); + let frames: Vec = results2 + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); assert_eq!(frames.len(), 6, "the regrown pool renders a new wave"); for result in frames { let payload = result.expect("regrown wave rendered"); @@ -237,8 +248,11 @@ fn two_workers_render_two_waves_zero_copy() { let deadline = Instant::now() + Duration::from_secs(60); while completed < 12 { dispatcher.poll(); - let drained: Vec = - results.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect(); + let drained: Vec = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); for result in drained { let payload = result.expect("frame rendered"); let TicketPayload::ShmFrame(frame) = payload else { @@ -304,10 +318,8 @@ fn crash_isolation_restarts_worker_and_frame_still_renders() { // One-shot crash hook: the worker dies with SIGSEGV while rendering // ticket 1; the marker file it leaves behind makes the restarted // worker render the re-queued frame for real. - let marker = std::env::temp_dir().join(format!( - "oak-procpool-crash-marker-{}", - std::process::id() - )); + let marker = + std::env::temp_dir().join(format!("oak-procpool-crash-marker-{}", std::process::id())); let _ = std::fs::remove_file(&marker); std::env::set_var("OAK_WORKER_CRASH_ON_TICKET", "1"); std::env::set_var("OAK_WORKER_CRASH_MARKER", &marker); @@ -387,7 +399,9 @@ fn worker_decodes_real_footage_into_slot() { // Decoded video is opaque: every BGRA alpha byte is 255. let pixels = &frame.shm.slot_bytes(frame.slot)[..frame.meta.data_size as usize]; let alpha_ok = pixels - .chunks_exact(4) + .as_chunks::<4>() + .0 + .iter() .filter(|px| px[3] == 255) .count(); assert!( @@ -431,10 +445,7 @@ fn montage_effects_render_through_the_worker() { type_id: "org.olivevideoeditor.Olive.opacity".into(), enabled: true, effect_input_id: Some("tex_in".into()), - params: vec![( - "opacity_in".into(), - oak_node::value::NodeValue::Float(0.5), - )], + params: vec![("opacity_in".into(), oak_node::value::NodeValue::Float(0.5))], }])], ]; @@ -455,9 +466,13 @@ fn montage_effects_render_through_the_worker() { )) }), done: Box::new(move |result| { - results.lock().unwrap_or_else(|e| e.into_inner()).push(result); + results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(result); }), schedule: JobSchedule::seek(), + cancelled: None, }; assert!(dispatcher.post(job), "post accepted while alive"); } @@ -480,7 +495,9 @@ fn montage_effects_render_through_the_worker() { }; // Pick an opaque, non-black pixel in the plain render as the probe. let probe = plain - .chunks_exact(4) + .as_chunks::<4>() + .0 + .iter() .position(|px| px[3] == 255 && px[0] > 40) .expect("the fixture frame has an opaque non-black pixel"); let p = &plain[probe * 4..probe * 4 + 4]; @@ -549,9 +566,13 @@ fn submit_audio( )) }), done: Box::new(move |result| { - results.lock().unwrap_or_else(|e| e.into_inner()).push(result); + results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(result); }), schedule: JobSchedule::seek(), + cancelled: None, }; assert!(dispatcher.post(job), "post accepted while alive"); } @@ -572,7 +593,13 @@ fn audio_tickets_roundtrip_through_shm_slots() { // 1/24 s at 48 kHz stereo = 2000 sample frames x 2 ch = 16000 bytes — // fits the 64x64 BGRA8 slot (16384 bytes). for i in 0..4 { - submit_audio(&dispatcher, &results, 1, Rational::new(i, 24), Rational::new(1, 24)); + submit_audio( + &dispatcher, + &results, + 1, + Rational::new(i, 24), + Rational::new(1, 24), + ); } pump_until(&dispatcher, &results, 4); @@ -590,7 +617,10 @@ fn audio_tickets_roundtrip_through_shm_slots() { // Empty montage: total silence, parsed back as f32. let samples = audio.samples(); assert_eq!(samples.len(), 2000 * 2); - assert!(samples.iter().all(|&v| v == 0.0), "empty montage is silence"); + assert!( + samples.iter().all(|&v| v == 0.0), + "empty montage is silence" + ); // Audio tickets are Seek priority — claimable by ANY worker (the // seek-starvation fix), so there is no shard-spread assertion; what // matters is that every ticket rendered on a live worker. @@ -692,7 +722,13 @@ fn audio_crash_isolation_restarts_worker_and_audio_still_renders() { // The dispatcher's first ticket id is 1 — exactly the crash ticket // (an audio ticket this time). for i in 0..4 { - submit_audio(&dispatcher, &results, 1, Rational::new(i, 24), Rational::new(1, 24)); + submit_audio( + &dispatcher, + &results, + 1, + Rational::new(i, 24), + Rational::new(1, 24), + ); } pump_until(&dispatcher, &results, 4); @@ -769,6 +805,7 @@ fn oversized_audio_ticket_is_refused_by_process_backend() { results2.lock().unwrap_or_else(|e| e.into_inner()).push(()); }), schedule: JobSchedule::seek(), + cancelled: None, }; assert!( !dispatcher.post(job), @@ -792,7 +829,12 @@ fn f32_ticket_gets_f32_slot_and_bgra8_stays_bgra8() { let results = Arc::new(Mutex::new(Vec::new())); submit(&dispatcher, &results, 1, None); pump_until(&dispatcher, &results, 1); - let bg = results.lock().unwrap().pop().unwrap().expect("frame rendered"); + let bg = results + .lock() + .unwrap() + .pop() + .unwrap() + .expect("frame rendered"); let TicketPayload::ShmFrame(frame) = bg else { panic!("ShmFrame payload"); }; @@ -829,14 +871,23 @@ fn f32_ticket_gets_f32_slot_and_bgra8_stays_bgra8() { )) }), done: Box::new(move |result| { - results2.lock().unwrap_or_else(|e| e.into_inner()).push(result); + results2 + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(result); }), schedule: JobSchedule::seek(), + cancelled: None, }; assert!(dispatcher.post(job), "f32 post accepted"); } pump_until(&dispatcher, &results2, 1); - let f32res = results2.lock().unwrap().pop().unwrap().expect("f32 frame rendered"); + let f32res = results2 + .lock() + .unwrap() + .pop() + .unwrap() + .expect("f32 frame rendered"); let TicketPayload::ShmFrame(frame) = f32res else { panic!("ShmFrame payload"); }; @@ -845,7 +896,9 @@ fn f32_ticket_gets_f32_slot_and_bgra8_stays_bgra8() { // The F32 bytes are zero (transparent black pipeline output). let pixels = frame.shm.slot_bytes(frame.slot); assert!( - pixels[..frame.meta.data_size as usize].iter().all(|&b| b == 0), + pixels[..frame.meta.data_size as usize] + .iter() + .all(|&b| b == 0), "generated F32 frame is transparent black" ); dispatcher.release_frame(&frame); @@ -854,7 +907,12 @@ fn f32_ticket_gets_f32_slot_and_bgra8_stays_bgra8() { // segment (the wire format is per ticket). submit(&dispatcher, &results, 1, None); pump_until(&dispatcher, &results, 1); - let bg2 = results.lock().unwrap().pop().unwrap().expect("frame rendered"); + let bg2 = results + .lock() + .unwrap() + .pop() + .unwrap() + .expect("frame rendered"); let TicketPayload::ShmFrame(frame) = bg2 else { panic!("ShmFrame payload"); }; @@ -892,7 +950,12 @@ fn build_graph_project(clip: &std::path::Path) -> (Arc>, NodeId) let (ccore, cbehavior) = oak_node::block::clip_create(); let clip_node = p.graph.add_node(ccore, cbehavior); p.graph - .connect(footage, clip_node, oak_node::block::clip_input::TEXTURE_INPUT, -1) + .connect( + footage, + clip_node, + oak_node::block::clip_input::TEXTURE_INPUT, + -1, + ) .expect("connect footage to clip"); let clip_behavior = p @@ -974,9 +1037,13 @@ fn post_graph_job( )) }), done: Box::new(move |result| { - results.lock().unwrap_or_else(|e| e.into_inner()).push(result); + results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(result); }), schedule: JobSchedule::seek(), + cancelled: None, }; assert!(dispatcher.post(job), "post accepted while alive"); } @@ -995,7 +1062,12 @@ fn assert_graph_frame_opaque(dispatcher: &ProcessDispatcher, payload: TicketResu assert_eq!(frame.meta.format, SLOT_FORMAT_BGRA8); assert_eq!(frame.meta.data_size, 64 * 64 * 4); let pixels = &frame.shm.slot_bytes(frame.slot)[..frame.meta.data_size as usize]; - let alpha_ok = pixels.chunks_exact(4).filter(|px| px[3] == 255).count(); + let alpha_ok = pixels + .as_chunks::<4>() + .0 + .iter() + .filter(|px| px[3] == 255) + .count(); assert!( alpha_ok as f64 >= 0.99 * (64 * 64) as f64, "graph-rendered frame must be opaque ({alpha_ok}/4096)" @@ -1042,11 +1114,19 @@ fn graph_mode_renders_sequence_viewer_from_snapshot() { dispatcher.start().expect("worker starts"); let results = Arc::new(Mutex::new(Vec::new())); - let project_uuid = project.lock().unwrap_or_else(|e| e.into_inner()).uuid.clone(); + let project_uuid = project + .lock() + .unwrap_or_else(|e| e.into_inner()) + .uuid + .clone(); post_graph_job(&dispatcher, &results, viewer, &project_uuid); pump_until(&dispatcher, &results, 1); - let result = results.lock().unwrap_or_else(|e| e.into_inner()).pop().unwrap(); + let result = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .pop() + .unwrap(); assert_graph_frame_opaque(&dispatcher, result); dispatcher.shutdown(); @@ -1088,11 +1168,19 @@ fn set_graph_snapshot_after_start_reroutes_tickets() { dispatcher.set_graph_snapshot(Some(snapshot.display().to_string())); let results = Arc::new(Mutex::new(Vec::new())); - let project_uuid = project.lock().unwrap_or_else(|e| e.into_inner()).uuid.clone(); + let project_uuid = project + .lock() + .unwrap_or_else(|e| e.into_inner()) + .uuid + .clone(); post_graph_job(&dispatcher, &results, viewer, &project_uuid); pump_until(&dispatcher, &results, 1); - let result = results.lock().unwrap_or_else(|e| e.into_inner()).pop().unwrap(); + let result = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .pop() + .unwrap(); assert_graph_frame_opaque(&dispatcher, result); dispatcher.shutdown(); diff --git a/docs/zh/plans/render-pipeline-threads.md b/docs/zh/plans/render-pipeline-threads.md index 47fe525d1..f15b5dea9 100644 --- a/docs/zh/plans/render-pipeline-threads.md +++ b/docs/zh/plans/render-pipeline-threads.md @@ -273,6 +273,20 @@ > 生效——若 Seek 被挡在门外,UI 会等一个导出/缓存帧跑完。`cancel_preview_frame` > 丢弃已过 playhead 的排队帧,按 `(sequence, frame, version)` 全键匹配 > (两个监视器可同帧号同版本,只按帧匹配会误杀另一序列)。 +> - 自动缓存优先级修正:`PreviewAutoCacher::start_range_job` 改走 +> `submit_video_background`(Background),不再以 Seek 插到播放之前—— +> 否则 Seek 超额插队会让逐张 autocache 帧绕过队列容量约束、倒挂 +> §3.3 的"交互 > 播放 > 导出 > 自动缓存";`single_frame` 交互帧保持 +> Seek(`range_jobs_are_background_and_single_frames_stay_seek`)。 +> - 取消即停(审计 B):`Job.cancelled` 钩子由 arena 从 slot 的 cancel 原子 +> 安装,`execute_job` 在生产前检查——`TicketArena::cancel` 之后仍未开跑的 +> 帧直接以 `Error::State` 完成,不再烧一遍 GPU/CPU 再丢弃;交付仍是 +> exactly-once(`pipeline_skips_a_cancelled_job` 断言 producer 未运行)。 +> - 双份帧缓存(审计 C):eval 内层 `decoded_frames`(24 帧)才是缓存主拷贝, +> `DecodeService` LRU 只是解码线程→渲染请求的交接缓冲,已由 8 帧降到 +> `DECODE_LRU_CAP=2`;1080p F32 下副缓存常驻 ~8×33MB → ~2×33MB,交接 +> 未命中由 eval 缓存兜底(不再解码,只多一次拷贝)。彻底合并两份缓存需 +> 让 Texture/缓存持有共享 Arc 帧,另行立项。 > - **基准对比**(本机 release,`oak-render/examples/bench_playback`, > 1080p/25fps MPEG-2 源、240 帧@480p(853×480) / 128 帧@1080p;两个后端 > 统一输出 F32 帧、同尺寸,否则进程池的 BGRA8 槽位与管线的 in-process