1use crate::{
2 executor::NodeJob,
3 message::{
4 Action, Message, OfflineAutomationLane, OfflineAutomationPoint, OfflineAutomationTarget,
5 OfflineBounceWork, ProcessTask,
6 },
7 midi::io::MidiEvent,
8 render_plan::Op,
9};
10#[cfg(unix)]
11use nix::libc;
12use std::sync::Arc;
13use std::time::{Duration, Instant};
14use tokio::sync::mpsc::{Receiver, Sender};
15
16pub(crate) struct NodeJobResult {
17 pub(crate) worker_id: usize,
18 pub(crate) epoch: u64,
19 pub(crate) node: u32,
20 pub(crate) output_linear: Vec<f32>,
21 pub(crate) parameter_updates: Vec<Action>,
22 pub(crate) latency_changed: bool,
23}
24
25#[derive(Debug)]
26pub struct Worker {
27 id: usize,
28 rx: Receiver<Message>,
29 tx: Sender<Message>,
30 realtime_priority: i32,
31}
32
33impl Worker {
34 fn automation_lane_value_at(points: &[OfflineAutomationPoint], sample: usize) -> Option<f32> {
35 if points.is_empty() {
36 return None;
37 }
38 if sample <= points[0].sample {
39 return Some(points[0].value.clamp(0.0, 1.0));
40 }
41 if sample >= points[points.len().saturating_sub(1)].sample {
42 return Some(points[points.len().saturating_sub(1)].value.clamp(0.0, 1.0));
43 }
44 for segment in points.windows(2) {
45 let left = &segment[0];
46 let right = &segment[1];
47 if sample < left.sample || sample > right.sample {
48 continue;
49 }
50 let span = right.sample.saturating_sub(left.sample).max(1) as f32;
51 let t = (sample.saturating_sub(left.sample) as f32 / span).clamp(0.0, 1.0);
52 return Some((left.value + (right.value - left.value) * t).clamp(0.0, 1.0));
53 }
54 None
55 }
56
57 fn apply_freeze_automation_at_sample(
58 track: &mut crate::track::TrackData,
59 sample: usize,
60 lanes: &[OfflineAutomationLane],
61 ) {
62 for lane in lanes {
63 if matches!(
64 lane.target,
65 OfflineAutomationTarget::Volume | OfflineAutomationTarget::Balance
66 ) {
67 continue;
68 }
69 let Some(value) = Self::automation_lane_value_at(&lane.points, sample) else {
70 continue;
71 };
72 match lane.target {
73 OfflineAutomationTarget::Volume | OfflineAutomationTarget::Balance => {}
74 OfflineAutomationTarget::MidiCc { channel, cc } => {
75 let cc_value = (value * 127.0).round() as u8;
76 track.rt.pending_automation_midi_events.push(MidiEvent::new(
77 0,
78 vec![0xB0 | channel.min(15), cc.min(127), cc_value],
79 ));
80 }
81 #[cfg(unix)]
82 OfflineAutomationTarget::Lv2Parameter {
83 instance_id,
84 index,
85 min,
86 max,
87 } => {
88 let lo = min.min(max);
89 let hi = max.max(min);
90 let param_value = (lo + value * (hi - lo)).clamp(lo, hi);
91 let _ = track.set_lv2_control_value(
92 instance_id,
93 index as usize,
94 param_value as f64,
95 );
96 }
97 OfflineAutomationTarget::Vst3Parameter {
98 instance_id,
99 param_id,
100 } => {
101 let _ = track.set_vst3_parameter(instance_id, param_id, value.clamp(0.0, 1.0));
102 }
103 OfflineAutomationTarget::ClapParameter {
104 instance_id,
105 param_id,
106 min,
107 max,
108 } => {
109 let lo = min.min(max);
110 let hi = max.max(min);
111 let param_value = (lo + value as f64 * (hi - lo)).clamp(lo, hi);
112 let _ = track.set_clap_parameter_at(instance_id, param_id, param_value, 0);
113 }
114 }
115 }
116 }
117
118 fn prepare_track_for_freeze_render(track: &mut crate::track::TrackData) -> (f32, f32) {
119 let original_level = track.level();
120 let original_balance = track.balance();
121 track.set_level(0.0);
122 track.set_balance(0.0);
123 (original_level, original_balance)
124 }
125
126 fn restore_track_after_freeze_render(
127 track: &mut crate::track::TrackData,
128 original_level: f32,
129 original_balance: f32,
130 ) {
131 track.set_level(original_level);
132 track.set_balance(original_balance);
133 }
134
135 async fn process_offline_bounce(&self, job: OfflineBounceWork) {
136 let track_handle = job.state.tracks.get(&job.track_name).cloned();
137 let Some(target_track) = track_handle else {
138 let _ = self
139 .tx
140 .send(Message::OfflineBounceFinished {
141 result: Err(format!("Track not found: {}", job.track_name)),
142 })
143 .await;
144 let _ = self.tx.send(Message::Ready(self.id)).await;
145 return;
146 };
147 let (channels, block_size, sample_rate) = {
148 let t = target_track.lock();
149 let block_size = t
150 .audio
151 .outs
152 .first()
153 .map(|io| io.buffer_size())
154 .or_else(|| t.audio.ins.first().map(|io| io.buffer_size()))
155 .unwrap_or(0)
156 .max(1);
157 (
158 t.audio.outs.len().max(1),
159 block_size,
160 t.sample_rate.round().max(1.0) as i32,
161 )
162 };
163 let freeze_state = if job.apply_fader {
164 None
165 } else {
166 let mut t = target_track.lock();
167 Some(Self::prepare_track_for_freeze_render(&mut t))
168 };
169
170 let all_tracks: Vec<_> = job.state.tracks.values().cloned().collect();
171 let plan_collector = basedrop::Collector::new();
172 let render_plan = crate::render_plan::RenderPlan::compile(&job.state, &[], &[], block_size);
173 let render_plan = Arc::new(basedrop::Owned::new(&plan_collector.handle(), render_plan));
174
175 let mut output_samples =
176 Vec::<f32>::with_capacity(job.length_samples.saturating_mul(channels.max(1)));
177
178 let mut cursor = 0usize;
179 let mut last_reported_progress = 0.0_f32;
180 let mut total_process_time = Duration::ZERO;
181 let mut total_write_time = Duration::ZERO;
182 while cursor < job.length_samples {
183 if job.cancel.load(std::sync::atomic::Ordering::Relaxed) {
184 let _ = std::fs::remove_file(&job.output_path);
185 if let Some((original_level, original_balance)) = freeze_state {
186 let mut t = target_track.lock();
187 Self::restore_track_after_freeze_render(
188 &mut t,
189 original_level,
190 original_balance,
191 );
192 }
193 let _ = self
194 .tx
195 .send(Message::OfflineBounceFinished {
196 result: Ok(Action::TrackOfflineBounceCanceled {
197 track_name: job.track_name.clone(),
198 }),
199 })
200 .await;
201 let _ = self.tx.send(Message::Ready(self.id)).await;
202 return;
203 }
204
205 let step = (job.length_samples - cursor).min(block_size);
206 for handle in &all_tracks {
207 let mut t = handle.lock();
208 t.audio.set_finished(false);
209 t.audio.set_processing(false);
210 t.set_transport_sample(job.start_sample.saturating_add(cursor));
211 t.set_loop_config(false, None);
212 t.set_transport_timing(job.tempo_bpm, job.tsig_num, job.tsig_denom);
213 t.set_clip_playback_enabled(true);
214 t.set_record_tap_enabled(false);
215 }
216 {
217 let mut t = target_track.lock();
218 Self::apply_freeze_automation_at_sample(
219 &mut t,
220 job.start_sample.saturating_add(cursor),
221 &job.automation_lanes,
222 );
223 }
224
225 let p_start = Instant::now();
226 for node in 0..render_plan.nodes.len() as crate::render_plan::NodeId {
227 let _ = Self::process_node_job_result(
228 self.id,
229 NodeJob {
230 epoch: 0,
231 plan: render_plan.clone(),
232 node,
233 },
234 );
235 }
236 total_process_time += p_start.elapsed();
237
238 let write_start = Instant::now();
239 {
240 let t = target_track.lock();
241 let outs = t.last_audio_outputs();
242 for i in 0..step {
243 for ch in 0..channels {
244 let sample = outs
245 .get(ch)
246 .and_then(|out| out.get(i))
247 .copied()
248 .unwrap_or(0.0);
249 output_samples.push(sample);
250 }
251 }
252 }
253 total_write_time += write_start.elapsed();
254
255 cursor = cursor.saturating_add(step);
256 let progress = (cursor as f32 / job.length_samples as f32).clamp(0.0, 1.0);
257
258 if progress - last_reported_progress >= 0.01 || cursor >= job.length_samples {
259 last_reported_progress = progress;
260 let _ = self
261 .tx
262 .send(Message::OfflineBounceFinished {
263 result: Ok(Action::TrackOfflineBounceProgress {
264 track_name: job.track_name.clone(),
265 progress,
266 operation: Some("Rendering freeze".to_string()),
267 }),
268 })
269 .await;
270 }
271 }
272
273 if let Err(e) = crate::audio_codec::write_wav_f32(
274 std::path::Path::new(&job.output_path),
275 &output_samples,
276 channels,
277 sample_rate as u32,
278 ) {
279 let _ = std::fs::remove_file(&job.output_path);
280 if let Some((original_level, original_balance)) = freeze_state {
281 let mut t = target_track.lock();
282 Self::restore_track_after_freeze_render(&mut t, original_level, original_balance);
283 }
284 let _ = self
285 .tx
286 .send(Message::OfflineBounceFinished {
287 result: Err(format!(
288 "Failed to write offline bounce '{}': {e}",
289 job.output_path
290 )),
291 })
292 .await;
293 let _ = self.tx.send(Message::Ready(self.id)).await;
294 return;
295 }
296
297 if let Some((original_level, original_balance)) = freeze_state {
298 let mut t = target_track.lock();
299 Self::restore_track_after_freeze_render(&mut t, original_level, original_balance);
300 }
301
302 let _ = self
303 .tx
304 .send(Message::OfflineBounceFinished {
305 result: Ok(Action::TrackOfflineBounce {
306 track_name: job.track_name,
307 output_path: job.output_path,
308 start_sample: job.start_sample,
309 length_samples: job.length_samples,
310 automation_lanes: vec![],
311 apply_fader: job.apply_fader,
312 }),
313 })
314 .await;
315 let _ = self.tx.send(Message::Ready(self.id)).await;
316 }
317
318 #[cfg(unix)]
319 pub(crate) fn try_enable_realtime(priority: i32) -> Result<(), String> {
320 let thread = unsafe { libc::pthread_self() };
321 let policy = libc::SCHED_FIFO;
322 let param = unsafe {
323 let mut p = std::mem::zeroed::<libc::sched_param>();
324 p.sched_priority = priority;
325 p
326 };
327 let rc = unsafe { libc::pthread_setschedparam(thread, policy, ¶m) };
328 if rc == 0 {
329 Ok(())
330 } else {
331 Err(format!("pthread_setschedparam failed with errno {}", rc))
332 }
333 }
334
335 #[cfg(target_os = "windows")]
336 pub(crate) fn try_enable_realtime(_priority: i32) -> Result<(), String> {
337 use std::{cell::Cell, ffi::OsStr, os::windows::ffi::OsStrExt};
338
339 #[link(name = "avrt")]
340 unsafe extern "system" {
341 fn AvSetMmThreadCharacteristicsW(task_name: *const u16, task_index: *mut u32) -> isize;
342 }
343
344 thread_local! {
345 static MMCSS_TASK_HANDLE: Cell<isize> = const { Cell::new(0) };
346 }
347
348 MMCSS_TASK_HANDLE.with(|handle| {
349 if handle.get() != 0 {
350 return Ok(());
351 }
352
353 let task_name: Vec<u16> = OsStr::new("Pro Audio")
354 .encode_wide()
355 .chain(Some(0))
356 .collect();
357 let mut task_index = 0_u32;
358 let mmcss_handle =
359 unsafe { AvSetMmThreadCharacteristicsW(task_name.as_ptr(), &mut task_index) };
360 if mmcss_handle == 0 {
361 Err(format!(
362 "AvSetMmThreadCharacteristicsW(Pro Audio) failed: {}",
363 std::io::Error::last_os_error()
364 ))
365 } else {
366 handle.set(mmcss_handle);
367 Ok(())
368 }
369 })
370 }
371
372 #[cfg(all(not(unix), not(target_os = "windows")))]
373 pub(crate) fn try_enable_realtime(_priority: i32) -> Result<(), String> {
374 Err("Realtime thread priority is not supported on this platform".to_string())
375 }
376
377 pub async fn new(
378 id: usize,
379 rx: Receiver<Message>,
380 tx: Sender<Message>,
381 realtime_priority: i32,
382 ) -> Worker {
383 let worker = Worker {
384 id,
385 rx,
386 tx,
387 realtime_priority,
388 };
389 worker.send(Message::Ready(id)).await;
390 worker
391 }
392
393 pub async fn send(&self, message: Message) {
394 self.tx
395 .send(message)
396 .await
397 .expect("Failed to send message from worker");
398 }
399
400 fn arena_input_slices<'a>(
401 plan: &'a crate::render_plan::RenderPlan,
402 ins: &[crate::render_plan::BufferId],
403 ) -> Vec<&'a [f32]> {
404 ins.iter()
405 .map(|&buf| {
406 unsafe { plan.buffer(buf) }
409 })
410 .collect()
411 }
412
413 fn arena_source_slices<'a>(
414 plan: &'a crate::render_plan::RenderPlan,
415 writable: &[crate::render_plan::BufferId],
416 ) -> Vec<(usize, &'a [f32], usize)> {
417 plan.port_map
418 .iter()
419 .filter_map(|(&key, &buf)| {
420 if writable.contains(&buf) {
421 return None;
422 }
423 Some((key, unsafe { plan.buffer(buf) }, plan.buffer_latency(buf)))
428 })
429 .collect()
430 }
431
432 fn metronome_output_buffer(
433 plan: &crate::render_plan::RenderPlan,
434 t: &crate::track::TrackData,
435 outs: &[crate::render_plan::BufferId],
436 ) -> Option<crate::render_plan::BufferId> {
437 let source = t.metronome_source()?;
438 let key = Arc::as_ptr(&source) as usize;
439 let &buf = plan.port_map.get(&key)?;
440 outs.contains(&buf).then_some(buf)
441 }
442
443 pub(crate) fn process_node_job_result(worker_id: usize, job: NodeJob) -> NodeJobResult {
450 let NodeJob { epoch, plan, node } = job;
451 let (output_linear, parameter_updates, latency_changed) = match &plan.nodes[node as usize] {
452 Op::Zero { output } => {
453 unsafe { &mut *plan.buffer_ptr(*output) }.fill(0.0);
457 plan.set_buffer_latency(*output, 0);
458 (Vec::new(), Vec::new(), false)
459 }
460 Op::Sum {
461 inputs,
462 delays,
463 output,
464 } => {
465 let out = unsafe { &mut *plan.buffer_ptr(*output) };
468 out.fill(0.0);
469 let max_latency = inputs
470 .iter()
471 .map(|&input| plan.buffer_latency(input))
472 .max()
473 .unwrap_or(0);
474 for (idx, &input) in inputs.iter().enumerate() {
475 let src = unsafe { plan.buffer(input) };
476 let delay = max_latency.saturating_sub(plan.buffer_latency(input));
477 let line = unsafe { &mut *delays[idx].get() };
481 line.process(src, delay, out, idx != 0);
482 }
483 plan.set_buffer_latency(*output, max_latency);
484 (Vec::new(), Vec::new(), false)
485 }
486 Op::HwInput { output, .. } => {
487 plan.set_buffer_latency(*output, 0);
490 (Vec::new(), Vec::new(), false)
491 }
492 Op::Task { task, ins, outs } => {
493 let track = match task {
494 ProcessTask::Track(t)
495 | ProcessTask::FolderInput(t)
496 | ProcessTask::FolderOutput(t) => t,
497 ProcessTask::Plugin { track, .. } => track,
498 };
499 let mut t = track.lock();
500 match task {
501 ProcessTask::Track(_) => {
502 let audio_out_count = t.audio.outs.len();
503 let metronome_output = Self::metronome_output_buffer(&plan, &t, outs);
504 let input_ptrs = ins
505 .iter()
506 .map(|&buf| {
507 unsafe { plan.buffer_ptr(buf) }
510 })
511 .collect::<Vec<_>>();
512 let mut inputs = input_ptrs
513 .iter()
514 .map(|&ptr| {
515 unsafe { (&mut *ptr).as_mut_slice() }
518 })
519 .collect::<Vec<_>>();
520 let source_buffers = Self::arena_source_slices(&plan, outs);
521 let output_ptrs = outs
522 .iter()
523 .take(audio_out_count)
524 .map(|&buf| {
525 unsafe { plan.buffer_ptr(buf) }
528 })
529 .collect::<Vec<_>>();
530 let mut outputs = output_ptrs
531 .iter()
532 .map(|&ptr| {
533 unsafe { (&mut *ptr).as_mut_slice() }
536 })
537 .collect::<Vec<_>>();
538 let metronome_output_ptr = metronome_output.map(|buf| {
539 unsafe { plan.buffer_ptr(buf) }
542 });
543 let metronome_output = metronome_output_ptr.map(|ptr| {
544 unsafe { (&mut *ptr).as_mut_slice() }
547 });
548 t.process_render_block_with_audio_buffers_and_metronome(
549 &mut inputs,
550 &mut outputs,
551 &source_buffers,
552 metronome_output,
553 );
554 for &out in outs.iter().take(audio_out_count) {
555 plan.set_buffer_latency(out, t.plugin_graph_latency_samples());
556 }
557 }
558 ProcessTask::FolderInput(_) => {
559 let metronome_output = Self::metronome_output_buffer(&plan, &t, outs);
560 let input_ptrs = ins
561 .iter()
562 .map(|&buf| {
563 unsafe { plan.buffer_ptr(buf) }
566 })
567 .collect::<Vec<_>>();
568 let mut inputs = input_ptrs
569 .iter()
570 .map(|&ptr| {
571 unsafe { (&mut *ptr).as_mut_slice() }
574 })
575 .collect::<Vec<_>>();
576 let metronome_output_ptr = metronome_output.map(|buf| {
577 unsafe { plan.buffer_ptr(buf) }
580 });
581 let metronome_output = metronome_output_ptr.map(|ptr| {
582 unsafe { (&mut *ptr).as_mut_slice() }
585 });
586 t.process_folder_input_with_audio_buffers_and_metronome(
587 &mut inputs,
588 metronome_output,
589 );
590 for &input in ins {
591 plan.set_buffer_latency(input, 0);
592 }
593 }
594 ProcessTask::FolderOutput(_) => {
595 let source_buffers = Self::arena_source_slices(&plan, outs);
596 let output_ptrs = outs
597 .iter()
598 .map(|&buf| {
599 unsafe { plan.buffer_ptr(buf) }
602 })
603 .collect::<Vec<_>>();
604 let mut outputs = output_ptrs
605 .iter()
606 .map(|&ptr| {
607 unsafe { (&mut *ptr).as_mut_slice() }
610 })
611 .collect::<Vec<_>>();
612 t.process_folder_output_with_audio_buffers(&mut outputs, &source_buffers);
613 for &out in outs {
614 plan.set_buffer_latency(out, t.plugin_graph_latency_samples());
615 }
616 }
617 ProcessTask::Plugin { kind, index, .. } => {
618 let input_latency = ins
619 .iter()
620 .map(|&input| plan.buffer_latency(input))
621 .max()
622 .unwrap_or(0);
623 let inputs = Self::arena_input_slices(&plan, ins);
624 let output_ptrs = outs
625 .iter()
626 .map(|&buf| {
627 unsafe { plan.buffer_ptr(buf) }
630 })
631 .collect::<Vec<_>>();
632 let mut outputs = output_ptrs
633 .iter()
634 .map(|&ptr| {
635 unsafe { (&mut *ptr).as_mut_slice() }
638 })
639 .collect::<Vec<_>>();
640 t.process_plugin_with_audio_buffers(*kind, *index, &inputs, &mut outputs);
641 let latency =
642 input_latency.saturating_add(t.plugin_latency_samples(*kind, *index));
643 for &out in outs {
644 plan.set_buffer_latency(out, latency);
645 }
646 }
647 }
648 t.audio.set_processing(false);
649 let latency_changed = t.take_plugin_latency_changed();
650 let updates = std::mem::take(&mut t.rt.echoed_parameter_updates);
651 let meter = t.output_meter_linear();
652 {
653 static TRACK_OUTPUT_PEAK_LOG_COUNT: std::sync::atomic::AtomicUsize =
654 std::sync::atomic::AtomicUsize::new(0);
655 let peak = meter.iter().copied().fold(0.0_f32, f32::max);
656 let count = TRACK_OUTPUT_PEAK_LOG_COUNT
657 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
658 if count < 64 || peak > 0.0 {
659 let task_kind = match task {
660 ProcessTask::Track(_) => "track",
661 ProcessTask::FolderInput(_) => "folder_input",
662 ProcessTask::FolderOutput(_) => "folder_output",
663 ProcessTask::Plugin { .. } => "plugin",
664 };
665 tracing::debug!(
666 worker_id,
667 node,
668 track = %t.name,
669 task_kind,
670 peak,
671 meter = ?meter,
672 "track output meter peak"
673 );
674 }
675 }
676 (meter, updates, latency_changed)
677 }
678 };
679 NodeJobResult {
680 worker_id,
681 epoch,
682 node,
683 output_linear,
684 parameter_updates,
685 latency_changed,
686 }
687 }
688
689 async fn process_node_job(&self, job: NodeJob) {
690 let result = Self::process_node_job_result(self.id, job);
691 let _ = self.tx.send(result.into()).await;
692 }
693
694 pub async fn work(&mut self) {
695 crate::enable_flush_denormals_to_zero();
696 if let Err(e) = Self::try_enable_realtime(self.realtime_priority) {
697 tracing::warn!(
698 "Worker {} realtime priority {} not enabled: {}",
699 self.id,
700 self.realtime_priority,
701 e
702 );
703 }
704 while let Some(message) = self.rx.recv().await {
705 match message {
706 Message::Request(Action::Quit) => {
707 return;
708 }
709 Message::ProcessOfflineBounce(job) => {
710 self.process_offline_bounce(job).await;
711 }
712 Message::NodeJob(job) => {
713 self.process_node_job(job).await;
714 }
715 _ => {}
716 }
717 }
718 }
719}
720
721impl From<NodeJobResult> for Message {
722 fn from(result: NodeJobResult) -> Self {
723 Message::NodeDone {
724 worker_id: result.worker_id,
725 epoch: result.epoch,
726 node: result.node,
727 output_linear: result.output_linear,
728 parameter_updates: result.parameter_updates,
729 latency_changed: result.latency_changed,
730 }
731 }
732}
733
734#[cfg(test)]
735mod tests {
736 use super::Worker;
737 use crate::message::{
738 Action, Message, OfflineAutomationLane, OfflineAutomationPoint, OfflineAutomationTarget,
739 OfflineBounceWork,
740 };
741 use crate::state::State;
742 use crate::track::Track;
743 use std::path::PathBuf;
744 use std::sync::{Arc, atomic::AtomicBool};
745 use std::time::{SystemTime, UNIX_EPOCH};
746 use tokio::sync::mpsc::channel;
747
748 fn make_state_with_track(track: Track) -> State {
749 let mut state = State::default();
750 state.tracks.insert(track.name.clone(), Arc::new(track));
751 state
752 }
753
754 fn unique_temp_wav(name: &str) -> PathBuf {
755 let nanos = SystemTime::now()
756 .duration_since(UNIX_EPOCH)
757 .expect("clock")
758 .as_nanos();
759 std::env::temp_dir().join(format!("maolan_{name}_{nanos}.wav"))
760 }
761
762 #[test]
763 fn prepare_track_for_freeze_render_neutralizes_level_and_balance() {
764 let mut track = Track::new("track".to_string(), 1, 2, 0, 0, 64, 48_000.0);
765 track.set_level(-6.0);
766 track.set_balance(0.35);
767
768 let (level, balance) = Worker::prepare_track_for_freeze_render(&mut track);
769
770 assert_eq!(level, -6.0);
771 assert_eq!(balance, 0.35);
772 assert_eq!(track.level(), 0.0);
773 assert_eq!(track.balance(), 0.0);
774
775 Worker::restore_track_after_freeze_render(&mut track, level, balance);
776 assert_eq!(track.level(), -6.0);
777 assert_eq!(track.balance(), 0.35);
778 }
779
780 #[test]
781 fn freeze_automation_ignores_volume_and_balance_lanes() {
782 let mut track = Track::new("track".to_string(), 1, 2, 0, 1, 64, 48_000.0);
783 let lanes = vec![
784 OfflineAutomationLane {
785 target: OfflineAutomationTarget::Volume,
786 visible: true,
787 points: vec![OfflineAutomationPoint {
788 sample: 0,
789 value: 0.0,
790 }],
791 },
792 OfflineAutomationLane {
793 target: OfflineAutomationTarget::Balance,
794 visible: true,
795 points: vec![OfflineAutomationPoint {
796 sample: 0,
797 value: 1.0,
798 }],
799 },
800 OfflineAutomationLane {
801 target: OfflineAutomationTarget::MidiCc { channel: 0, cc: 7 },
802 visible: true,
803 points: vec![OfflineAutomationPoint {
804 sample: 0,
805 value: 1.0,
806 }],
807 },
808 ];
809
810 Worker::apply_freeze_automation_at_sample(&mut track, 0, &lanes);
811
812 assert_eq!(track.level(), 0.0);
813 assert_eq!(track.balance(), 0.0);
814 assert_eq!(track.rt.pending_automation_midi_events.len(), 1);
815 assert_eq!(
816 track.rt.pending_automation_midi_events[0].data,
817 vec![0xB0, 7, 127]
818 );
819 }
820
821 #[test]
822 fn automation_lane_value_at_interpolates_between_points() {
823 let value = Worker::automation_lane_value_at(
824 &[
825 OfflineAutomationPoint {
826 sample: 10,
827 value: 0.25,
828 },
829 OfflineAutomationPoint {
830 sample: 20,
831 value: 0.75,
832 },
833 ],
834 15,
835 )
836 .expect("value");
837
838 assert!((value - 0.5).abs() < 1.0e-6);
839 }
840
841 #[test]
842 fn freeze_automation_applies_interpolated_midi_cc_lane() {
843 let mut track = Track::new("track".to_string(), 1, 1, 0, 1, 64, 48_000.0);
844 let lanes = vec![OfflineAutomationLane {
845 target: OfflineAutomationTarget::MidiCc { channel: 0, cc: 7 },
846 visible: true,
847 points: vec![
848 OfflineAutomationPoint {
849 sample: 0,
850 value: 0.0,
851 },
852 OfflineAutomationPoint {
853 sample: 10,
854 value: 1.0,
855 },
856 ],
857 }];
858
859 Worker::apply_freeze_automation_at_sample(&mut track, 5, &lanes);
860 assert_eq!(track.rt.pending_automation_midi_events.len(), 1);
861 assert_eq!(track.rt.pending_automation_midi_events[0].data[2], 64);
862
863 track.rt.pending_automation_midi_events.clear();
864 Worker::apply_freeze_automation_at_sample(&mut track, 2, &lanes);
865 assert_eq!(track.rt.pending_automation_midi_events.len(), 1);
866 assert_eq!(track.rt.pending_automation_midi_events[0].data[2], 25);
867 }
868
869 #[cfg_attr(
870 all(miri, target_os = "freebsd"),
871 ignore = "Tokio runtime uses kqueue, which Miri does not support on FreeBSD"
872 )]
873 #[tokio::test]
874 async fn process_node_job_sums_arena_buffers() {
875 use crate::render_plan::{Op, RenderPlan};
876 use std::cell::UnsafeCell;
877 use std::collections::HashMap;
878
879 let (_rx_unused_tx, rx_unused) = channel(1);
880 let (tx, mut out_rx) = channel(8);
881 let worker = Worker {
882 id: 4,
883 rx: rx_unused,
884 tx,
885 realtime_priority: 0,
886 };
887 let collector = basedrop::Collector::new();
888 let plan = RenderPlan {
889 buffer_size: 4,
890 buffers: (0..3).map(|_| UnsafeCell::new(vec![0.0; 4])).collect(),
891 buffer_latencies: (0..3)
892 .map(|_| std::sync::atomic::AtomicUsize::new(0))
893 .collect(),
894 nodes: vec![Op::Sum {
895 inputs: vec![0, 1],
896 delays: vec![
897 UnsafeCell::new(crate::render_plan::DelayLine::new()),
898 UnsafeCell::new(crate::render_plan::DelayLine::new()),
899 ],
900 output: 2,
901 }],
902 indegree: vec![0],
903 dependents: vec![vec![]],
904 sources: vec![0],
905 hw_in_map: vec![],
906 hw_out_map: vec![],
907 port_map: HashMap::new(),
908 midi_edges: vec![],
909 forced: vec![],
910 };
911 unsafe {
913 (&mut *plan.buffer_ptr(0)).copy_from_slice(&[0.25, 0.5, 0.75, 1.0]);
914 (&mut *plan.buffer_ptr(1)).copy_from_slice(&[0.75, 0.5, 0.25, f32::NAN]);
915 }
916 let shared = std::sync::Arc::new(basedrop::Owned::new(&collector.handle(), plan));
917
918 worker
919 .process_node_job(crate::executor::NodeJob {
920 epoch: 1,
921 plan: shared.clone(),
922 node: 0,
923 })
924 .await;
925
926 unsafe {
929 assert_eq!(
930 &*shared.buffer_ptr(2),
931 &vec![1.0, 1.0, 1.0, 1.0],
932 "sanitized sum in the arena"
933 );
934 }
935 match out_rx.recv().await.expect("message") {
936 Message::NodeDone {
937 worker_id,
938 epoch,
939 node,
940 output_linear,
941 ..
942 } => {
943 assert_eq!(worker_id, 4);
944 assert_eq!(epoch, 1);
945 assert_eq!(node, 0);
946 assert!(output_linear.is_empty());
947 }
948 other => panic!("unexpected message: {other:?}"),
949 }
950 }
951
952 #[cfg_attr(
953 all(miri, target_os = "freebsd"),
954 ignore = "Tokio runtime uses kqueue, which Miri does not support on FreeBSD"
955 )]
956 #[tokio::test]
957 async fn process_offline_bounce_errors_when_track_is_missing() {
958 let (_rx_unused_tx, rx_unused) = channel(1);
959 let (tx, mut out_rx) = channel(8);
960 let worker = Worker {
961 id: 7,
962 rx: rx_unused,
963 tx,
964 realtime_priority: 0,
965 };
966 let job = OfflineBounceWork {
967 state: Arc::new(State::default().snapshot()),
968 track_name: "missing".to_string(),
969 output_path: unique_temp_wav("missing").to_string_lossy().to_string(),
970 start_sample: 0,
971 length_samples: 8,
972 tempo_bpm: 120.0,
973 tsig_num: 4,
974 tsig_denom: 4,
975 automation_lanes: vec![],
976 cancel: Arc::new(AtomicBool::new(false)),
977 apply_fader: false,
978 };
979
980 worker.process_offline_bounce(job).await;
981
982 match out_rx.recv().await.expect("message") {
983 Message::OfflineBounceFinished { result: Err(err) } => {
984 assert!(err.contains("Track not found: missing"));
985 }
986 other => panic!("unexpected message: {other:?}"),
987 }
988 }
989
990 #[cfg_attr(
991 all(miri, target_os = "freebsd"),
992 ignore = "Tokio runtime uses kqueue, which Miri does not support on FreeBSD"
993 )]
994 #[tokio::test]
995 async fn process_offline_bounce_cancels_and_restores_track_state() {
996 let (_rx_unused_tx, rx_unused) = channel(1);
997 let (tx, mut out_rx) = channel(8);
998 let worker = Worker {
999 id: 5,
1000 rx: rx_unused,
1001 tx,
1002 realtime_priority: 0,
1003 };
1004 let track = Track::new("track".to_string(), 1, 2, 0, 0, 4, 48_000.0);
1005 track.set_level(-9.0);
1006 track.set_balance(-0.3);
1007 let state = make_state_with_track(track);
1008 let job = OfflineBounceWork {
1009 state: Arc::new(state.lock().snapshot()),
1010 track_name: "track".to_string(),
1011 output_path: unique_temp_wav("cancel").to_string_lossy().to_string(),
1012 start_sample: 0,
1013 length_samples: 8,
1014 tempo_bpm: 120.0,
1015 tsig_num: 4,
1016 tsig_denom: 4,
1017 automation_lanes: vec![],
1018 cancel: Arc::new(AtomicBool::new(true)),
1019 apply_fader: false,
1020 };
1021
1022 worker.process_offline_bounce(job).await;
1023
1024 match out_rx.recv().await.expect("message") {
1025 Message::OfflineBounceFinished {
1026 result: Ok(Action::TrackOfflineBounceCanceled { track_name }),
1027 } => assert_eq!(track_name, "track"),
1028 other => panic!("unexpected message: {other:?}"),
1029 }
1030 assert!(matches!(out_rx.recv().await, Some(Message::Ready(5))));
1031 let state_guard = state.lock();
1032 let track = state_guard.tracks.get("track").expect("track").lock();
1033 assert_eq!(track.level(), -9.0);
1034 assert_eq!(track.balance(), -0.3);
1035 }
1036
1037 #[cfg_attr(
1038 all(miri, target_os = "freebsd"),
1039 ignore = "Tokio runtime uses kqueue, which Miri does not support on FreeBSD"
1040 )]
1041 #[tokio::test]
1042 async fn process_offline_bounce_restores_track_state_on_write_failure() {
1043 let (_rx_unused_tx, rx_unused) = channel(1);
1044 let (tx, mut out_rx) = channel(8);
1045 let worker = Worker {
1046 id: 3,
1047 rx: rx_unused,
1048 tx,
1049 realtime_priority: 0,
1050 };
1051 let track = Track::new("track".to_string(), 1, 2, 0, 0, 4, 48_000.0);
1052 track.set_level(-4.0);
1053 track.set_balance(0.25);
1054 let state = make_state_with_track(track);
1055 let output_path = std::env::temp_dir().to_string_lossy().to_string();
1056 let job = OfflineBounceWork {
1057 state: Arc::new(state.lock().snapshot()),
1058 track_name: "track".to_string(),
1059 output_path,
1060 start_sample: 0,
1061 length_samples: 4,
1062 tempo_bpm: 120.0,
1063 tsig_num: 4,
1064 tsig_denom: 4,
1065 automation_lanes: vec![],
1066 cancel: Arc::new(AtomicBool::new(false)),
1067 apply_fader: false,
1068 };
1069
1070 worker.process_offline_bounce(job).await;
1071
1072 let mut saw_error = false;
1073 while let Some(message) = out_rx.recv().await {
1074 match message {
1075 Message::OfflineBounceFinished {
1076 result: Ok(Action::TrackOfflineBounceProgress { .. }),
1077 } => {}
1078 Message::OfflineBounceFinished { result: Err(err) } => {
1079 assert!(
1080 err.contains("Failed to create offline bounce")
1081 || err.contains("Failed to write offline bounce")
1082 || err.contains("Failed to finalize offline bounce")
1083 );
1084 saw_error = true;
1085 }
1086 Message::Ready(3) => break,
1087 other => panic!("unexpected message: {other:?}"),
1088 }
1089 }
1090 assert!(saw_error);
1091 let state_guard = state.lock();
1092 let track = state_guard.tracks.get("track").expect("track").lock();
1093 assert_eq!(track.level(), -4.0);
1094 assert_eq!(track.balance(), 0.25);
1095 }
1096
1097 #[cfg_attr(
1098 all(miri, target_os = "freebsd"),
1099 ignore = "Tokio runtime uses kqueue, which Miri does not support on FreeBSD"
1100 )]
1101 #[tokio::test]
1102 async fn process_offline_bounce_emits_progress_and_completion() {
1103 let (_rx_unused_tx, rx_unused) = channel(1);
1104 let (tx, mut out_rx) = channel(16);
1105 let worker = Worker {
1106 id: 2,
1107 rx: rx_unused,
1108 tx,
1109 realtime_priority: 0,
1110 };
1111 let track = Track::new("track".to_string(), 1, 1, 0, 0, 4, 48_000.0);
1112 track.set_level(-3.0);
1113 track.set_balance(0.4);
1114 let state = make_state_with_track(track);
1115 let output = unique_temp_wav("success");
1116 let job = OfflineBounceWork {
1117 state: Arc::new(state.lock().snapshot()),
1118 track_name: "track".to_string(),
1119 output_path: output.to_string_lossy().to_string(),
1120 start_sample: 0,
1121 length_samples: 8,
1122 tempo_bpm: 120.0,
1123 tsig_num: 4,
1124 tsig_denom: 4,
1125 automation_lanes: vec![],
1126 cancel: Arc::new(AtomicBool::new(false)),
1127 apply_fader: false,
1128 };
1129
1130 worker.process_offline_bounce(job).await;
1131
1132 let mut saw_progress = false;
1133 let mut saw_complete = false;
1134 let mut saw_ready = false;
1135 while let Some(message) = out_rx.recv().await {
1136 match message {
1137 Message::OfflineBounceFinished {
1138 result:
1139 Ok(Action::TrackOfflineBounceProgress {
1140 track_name,
1141 progress,
1142 ..
1143 }),
1144 } => {
1145 assert_eq!(track_name, "track");
1146 assert!(progress > 0.0);
1147 saw_progress = true;
1148 }
1149 Message::OfflineBounceFinished {
1150 result:
1151 Ok(Action::TrackOfflineBounce {
1152 track_name,
1153 output_path,
1154 ..
1155 }),
1156 } => {
1157 assert_eq!(track_name, "track");
1158 assert_eq!(output_path, output.to_string_lossy());
1159 saw_complete = true;
1160 }
1161 Message::Ready(2) => {
1162 saw_ready = true;
1163 break;
1164 }
1165 other => panic!("unexpected message: {other:?}"),
1166 }
1167 }
1168
1169 assert!(saw_progress);
1170 assert!(saw_complete);
1171 assert!(saw_ready);
1172 assert!(output.exists());
1173 std::fs::remove_file(&output).expect("remove temp wav");
1174 let state_guard = state.lock();
1175 let track = state_guard.tracks.get("track").expect("track").lock();
1176 assert_eq!(track.level(), -3.0);
1177 assert_eq!(track.balance(), 0.4);
1178 assert!(!track.muted());
1179 }
1180}