diff --git a/crates/oak-worker/tests/procpool_integration.rs b/crates/oak-worker/tests/procpool_integration.rs index 2c628654a..c1f8185a8 100644 --- a/crates/oak-worker/tests/procpool_integration.rs +++ b/crates/oak-worker/tests/procpool_integration.rs @@ -418,6 +418,64 @@ fn audio_tickets_roundtrip_through_shm_slots() { dispatcher.shutdown(); } +/// A claim batch that mixes audio and video tickets must assign slots in +/// the worker's acquisition order (the video message is processed first, +/// the audio message second). Interleaved assignment scrambled the +/// worker's free ring and flooded "slot assignment mismatch" failures +/// during playback — this is the regression guard. +#[test] +fn mixed_audio_video_batch_keeps_slot_assignment_order() { + let _guard = lock_test(); + // One worker, four slots: the first four posts dispatch singly (post + // pumps once itself), the remaining posts queue behind the busy + // slots. As completions are released below, claims gather up to + // `batch_size` queued tickets — mixed audio/video batches. + let dispatcher = ProcessDispatcher::new(config(1, 4)).expect("dispatcher config"); + dispatcher.start().expect("worker starts"); + + let results = Arc::new(Mutex::new(Vec::new())); + for i in 0..8 { + submit(&dispatcher, &results, 1, None); + submit_audio( + &dispatcher, + &results, + 1, + Rational::new(i, 24), + Rational::new(1, 24), + ); + } + // Drain completions, releasing each slot immediately so the queued + // tickets keep flowing into new (mixed) claims. + let deadline = Instant::now() + Duration::from_secs(60); + let mut total = 0usize; + loop { + dispatcher.poll(); + let done: Vec = results + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect(); + total += done.len(); + for result in &done { + match result { + Ok(TicketPayload::ShmFrame(frame)) => dispatcher.release_frame(frame), + Ok(TicketPayload::ShmAudio(audio)) => dispatcher.release_audio_frame(audio), + Err(e) => panic!("mixed-batch ticket failed: {e}"), + _ => panic!("unexpected payload variant"), + } + } + if total >= 16 { + break; + } + if Instant::now() > deadline { + panic!("timeout waiting for the mixed batch ({total}/16)"); + } + std::thread::sleep(Duration::from_millis(2)); + } + + dispatcher.shutdown(); +} + /// An audio render crashing mid-mix (SIGSEGV hook) must not take down the /// main process: the audio ticket is re-queued, the worker restarted and /// the samples still arrive. diff --git a/crates/oakrender/src/procpool.rs b/crates/oakrender/src/procpool.rs index c3d388989..913b4b1e4 100644 --- a/crates/oakrender/src/procpool.rs +++ b/crates/oakrender/src/procpool.rs @@ -1199,13 +1199,27 @@ impl ProcessDispatcher { let Some(batch) = inner.scheduler.claim_batch(worker, credit, max_bytes) else { return; }; - let mut video_tickets = Vec::with_capacity(batch.frames.len()); - let mut audio_tickets: Vec = Vec::new(); - for req in &batch.frames { + // Slot assignment order MUST match the worker's acquisition + // order: the batch is delivered as the video message first and + // the audio message second, and the worker pops one slot per + // ticket in that message order, checking each pop against the + // assignment. Assigning in the scheduler's interleaved frame + // order scrambles the free ring (every audio ticket in a mixed + // batch mismatched, and each mismatch leaked a slot — the + // "slot assignment mismatch" flood). Two passes: video first. + let (video_reqs, audio_reqs): (Vec<_>, Vec<_>) = + batch.frames.iter().partition(|r| { + !inner + .tickets + .get(&r.payload) + .is_some_and(|pt| pt.audio.is_some()) + }); + let mut video_tickets = Vec::with_capacity(video_reqs.len()); + let mut audio_tickets: Vec = Vec::with_capacity(audio_reqs.len()); + for req in video_reqs.into_iter().chain(audio_reqs) { let ticket = req.payload; - let slot = match inner.workers[worker].free_slots.pop_front() { - Some(s) => s, - None => break, // credit accounting drifted; stop cleanly + let Some(slot) = inner.workers[worker].free_slots.pop_front() else { + break; // credit accounting drifted; stop cleanly }; inner.workers[worker].outstanding.insert(ticket, slot); let Some(pt) = inner.tickets.get(&ticket) else {