1use std::sync::atomic::{AtomicU32, Ordering};
12use std::sync::{Arc, Mutex, MutexGuard};
13
14use rtrb::{Consumer, Producer, RingBuffer};
15
16use crate::{
17 audio::{AudioBackend, AudioBuffers, AudioConfig, AudioLevels, AudioStream, ChannelLevel},
18 error::{Error, Result},
19 midi::MidiEvent,
20 plugin::Plugin,
21 realtime::{RealtimePluginRunner, RtControl},
22};
23
24const SIDE_CHANNEL_CAPACITY: usize = 4096;
28
29enum HybridCommand {
32 Midi(MidiEvent),
33 Param { id: u32, value: f64 },
34 Panic,
35}
36
37fn channel_peak(buf: &[f32]) -> f32 {
39 buf.iter()
40 .map(|&x| if x.is_finite() { x.abs() } else { 0.0 })
41 .fold(0.0_f32, f32::max)
42}
43
44struct AudioSideChannels {
51 control_rx: Consumer<HybridCommand>,
52 out_midi_tx: Producer<MidiEvent>,
53 param_tx: Producer<(u32, f64)>,
54 levels: Arc<[AtomicU32]>,
58}
59
60impl AudioSideChannels {
61 fn apply_control(&mut self, plugin: &mut Plugin) {
63 while let Ok(cmd) = self.control_rx.pop() {
64 match cmd {
65 HybridCommand::Midi(event) => {
66 let _ = plugin.send_midi_event(event);
67 }
68 HybridCommand::Param { id, value } => {
69 let _ = plugin.set_parameter(id, value);
70 }
71 HybridCommand::Panic => {
72 let _ = plugin.midi_panic();
73 }
74 }
75 }
76 }
77
78 fn publish_levels(&mut self, outputs: &[Vec<f32>]) {
81 for (ch, atomic) in self.levels.iter().enumerate() {
82 let peak = outputs.get(ch).map(|b| channel_peak(b)).unwrap_or(0.0);
83 atomic.fetch_max(peak.to_bits(), Ordering::Relaxed);
84 }
85 }
86
87 fn publish_feedback(&mut self, plugin: &Plugin) {
92 for event in plugin.take_output_midi() {
93 let _ = self.out_midi_tx.push(event);
94 }
95 for change in plugin.get_parameter_changes() {
96 let _ = self.param_tx.push(change);
97 }
98 }
99}
100
101struct UiSideChannels {
105 control_tx: Mutex<Producer<HybridCommand>>,
106 out_midi_rx: Mutex<Consumer<MidiEvent>>,
107 param_rx: Mutex<Consumer<(u32, f64)>>,
108 levels: Arc<[AtomicU32]>,
109}
110
111fn make_side_channels(channels: usize) -> (AudioSideChannels, UiSideChannels) {
114 let (control_tx, control_rx) = RingBuffer::<HybridCommand>::new(SIDE_CHANNEL_CAPACITY);
115 let (out_midi_tx, out_midi_rx) = RingBuffer::<MidiEvent>::new(SIDE_CHANNEL_CAPACITY);
116 let (param_tx, param_rx) = RingBuffer::<(u32, f64)>::new(SIDE_CHANNEL_CAPACITY);
117 let levels: Arc<[AtomicU32]> = (0..channels).map(|_| AtomicU32::new(0)).collect();
118
119 let audio = AudioSideChannels {
120 control_rx,
121 out_midi_tx,
122 param_tx,
123 levels: Arc::clone(&levels),
124 };
125 let ui = UiSideChannels {
126 control_tx: Mutex::new(control_tx),
127 out_midi_rx: Mutex::new(out_midi_rx),
128 param_rx: Mutex::new(param_rx),
129 levels,
130 };
131 (audio, ui)
132}
133
134pub struct AudioHandle {
140 _stream: Box<dyn AudioStream>,
143 _input_stream: Option<Box<dyn AudioStream>>,
146 plugin: Arc<Mutex<Plugin>>,
147 ui: UiSideChannels,
150}
151
152impl AudioHandle {
153 pub fn lock(&self) -> MutexGuard<'_, Plugin> {
158 self.plugin
159 .lock()
160 .unwrap_or_else(|poisoned| poisoned.into_inner())
161 }
162
163 pub fn try_lock(&self) -> Option<MutexGuard<'_, Plugin>> {
172 match self.plugin.try_lock() {
173 Ok(guard) => Some(guard),
174 Err(std::sync::TryLockError::Poisoned(p)) => Some(p.into_inner()),
175 Err(std::sync::TryLockError::WouldBlock) => None,
176 }
177 }
178
179 pub fn send_midi(&self, event: MidiEvent) -> bool {
185 self.ui
186 .control_tx
187 .lock()
188 .map(|mut tx| tx.push(HybridCommand::Midi(event)).is_ok())
189 .unwrap_or(false)
190 }
191
192 pub fn set_parameter(&self, id: u32, value: f64) -> bool {
195 self.ui
196 .control_tx
197 .lock()
198 .map(|mut tx| tx.push(HybridCommand::Param { id, value }).is_ok())
199 .unwrap_or(false)
200 }
201
202 pub fn midi_panic(&self) -> bool {
205 self.ui
206 .control_tx
207 .lock()
208 .map(|mut tx| tx.push(HybridCommand::Panic).is_ok())
209 .unwrap_or(false)
210 }
211
212 pub fn output_levels(&self) -> AudioLevels {
219 let channels = self
220 .ui
221 .levels
222 .iter()
223 .map(|atomic| {
224 let peak = f32::from_bits(atomic.swap(0, Ordering::Relaxed));
225 ChannelLevel {
226 peak,
227 rms: 0.0,
228 peak_hold: peak,
229 }
230 })
231 .collect();
232 AudioLevels { channels }
233 }
234
235 pub fn drain_output_midi(&self) -> Vec<MidiEvent> {
238 let mut out = Vec::new();
239 if let Ok(mut rx) = self.ui.out_midi_rx.lock() {
240 while let Ok(event) = rx.pop() {
241 out.push(event);
242 }
243 }
244 out
245 }
246
247 pub fn drain_parameter_changes(&self) -> Vec<(u32, f64)> {
250 let mut out = Vec::new();
251 if let Ok(mut rx) = self.ui.param_rx.lock() {
252 while let Ok(change) = rx.pop() {
253 out.push(change);
254 }
255 }
256 out
257 }
258
259 pub fn plugin(&self) -> Arc<Mutex<Plugin>> {
261 Arc::clone(&self.plugin)
262 }
263
264 pub fn stop(self) {}
266}
267
268pub(crate) fn interleave_outputs(outputs: &[Vec<f32>], out: &mut [f32], channels: usize) {
275 if channels == 0 {
276 return;
277 }
278 let frames = out.len() / channels;
279 for ch in 0..channels.min(outputs.len()) {
280 let src = &outputs[ch];
281 for frame in 0..frames.min(src.len()) {
282 out[frame * channels + ch] = src[frame];
283 }
284 }
285}
286
287fn prepare_scratch(scratch: &mut AudioBuffers, frames: usize) {
289 for ch in &mut scratch.outputs {
290 if ch.len() != frames {
291 ch.resize(frames, 0.0);
292 }
293 ch.fill(0.0);
294 }
295 for ch in &mut scratch.inputs {
296 if ch.len() != frames {
297 ch.resize(frames, 0.0);
298 }
299 ch.fill(0.0);
300 }
301 scratch.block_size = frames;
302}
303
304pub fn play_with_backend<B: AudioBackend>(
313 backend: &B,
314 plugin: Plugin,
315 config: AudioConfig,
316) -> Result<AudioHandle> {
317 let device = backend
318 .default_output_device()
319 .ok_or_else(|| Error::AudioBackendError("No default output device available".into()))?;
320
321 let channels = config.output_channels;
322 let sample_rate = config.sample_rate;
323
324 let plugin = Arc::new(Mutex::new(plugin));
325 plugin
327 .lock()
328 .unwrap_or_else(|p| p.into_inner())
329 .start_processing()?;
330
331 let plugin_cb = Arc::clone(&plugin);
332 let (mut side, ui) = make_side_channels(channels);
334 let mut scratch = AudioBuffers::new(0, channels, config.block_size, sample_rate);
336
337 let data_cb = Box::new(move |data: &mut [f32]| {
338 data.fill(0.0);
340 if channels == 0 {
341 return;
342 }
343 let frames = data.len() / channels;
344 prepare_scratch(&mut scratch, frames);
345
346 let mut p = match plugin_cb.lock() {
350 Ok(guard) => guard,
351 Err(poisoned) => poisoned.into_inner(),
352 };
353 side.apply_control(&mut p);
357 if p.process_audio(&mut scratch).is_ok() {
358 interleave_outputs(&scratch.outputs, data, channels);
359 side.publish_levels(&scratch.outputs);
360 }
361 side.publish_feedback(&p);
362 });
363
364 let err_cb = Box::new(|e: B::Error| {
365 log::error!("audio stream error: {}", e);
366 });
367
368 let stream = backend
369 .create_output_stream(&device, config, data_cb, err_cb)
370 .map_err(|e| Error::AudioBackendError(format!("Failed to create output stream: {}", e)))?;
371
372 stream
373 .play()
374 .map_err(|e| Error::AudioBackendError(format!("Failed to start stream: {}", e)))?;
375
376 Ok(AudioHandle {
377 _stream: Box::new(stream),
378 _input_stream: None,
379 plugin,
380 ui,
381 })
382}
383
384pub fn play_with_input_backend<B: AudioBackend>(
396 backend: &B,
397 plugin: Plugin,
398 config: AudioConfig,
399) -> Result<AudioHandle> {
400 let in_device = backend
401 .default_input_device()
402 .ok_or_else(|| Error::AudioBackendError("No default input device available".into()))?;
403 let out_device = backend
404 .default_output_device()
405 .ok_or_else(|| Error::AudioBackendError("No default output device available".into()))?;
406
407 let in_channels = config.input_channels.max(1);
408 let out_channels = config.output_channels;
409 let sample_rate = config.sample_rate;
410
411 let plugin = Arc::new(Mutex::new(plugin));
412 plugin
413 .lock()
414 .unwrap_or_else(|p| p.into_inner())
415 .start_processing()?;
416
417 let ring_cap = (config.block_size * in_channels * 8).max(2048);
420 let (mut producer, mut consumer) = rtrb::RingBuffer::<f32>::new(ring_cap);
421
422 let in_data_cb = Box::new(move |data: &[f32]| {
423 for &s in data {
425 let _ = producer.push(s);
426 }
427 });
428 let in_err_cb = Box::new(|e: B::Error| log::error!("input stream error: {}", e));
429 let input_stream = backend
430 .create_input_stream(&in_device, config, in_data_cb, in_err_cb)
431 .map_err(|e| Error::AudioBackendError(format!("Failed to create input stream: {}", e)))?;
432
433 let plugin_cb = Arc::clone(&plugin);
434 let (mut side, ui) = make_side_channels(out_channels);
437 let mut scratch = AudioBuffers::new(in_channels, out_channels, config.block_size, sample_rate);
438 let out_data_cb = Box::new(move |data: &mut [f32]| {
439 data.fill(0.0);
440 if out_channels == 0 {
441 return;
442 }
443 let frames = data.len() / out_channels;
444 prepare_scratch(&mut scratch, frames);
445 for f in 0..frames {
448 for ch in scratch.inputs.iter_mut() {
449 ch[f] = consumer.pop().unwrap_or(0.0);
450 }
451 }
452 let mut p = match plugin_cb.lock() {
453 Ok(guard) => guard,
454 Err(poisoned) => poisoned.into_inner(),
455 };
456 side.apply_control(&mut p);
457 if p.process_audio(&mut scratch).is_ok() {
458 interleave_outputs(&scratch.outputs, data, out_channels);
459 side.publish_levels(&scratch.outputs);
460 }
461 side.publish_feedback(&p);
462 });
463 let out_err_cb = Box::new(|e: B::Error| log::error!("output stream error: {}", e));
464 let output_stream = backend
465 .create_output_stream(&out_device, config, out_data_cb, out_err_cb)
466 .map_err(|e| Error::AudioBackendError(format!("Failed to create output stream: {}", e)))?;
467
468 input_stream
469 .play()
470 .map_err(|e| Error::AudioBackendError(format!("Failed to start input stream: {}", e)))?;
471 output_stream
472 .play()
473 .map_err(|e| Error::AudioBackendError(format!("Failed to start output stream: {}", e)))?;
474
475 Ok(AudioHandle {
476 _stream: Box::new(output_stream),
477 _input_stream: Some(Box::new(input_stream)),
478 plugin,
479 ui,
480 })
481}
482
483pub struct RtAudioHandle {
487 _stream: Box<dyn AudioStream>,
488 control: RtControl,
489}
490
491impl RtAudioHandle {
492 pub fn control(&mut self) -> &mut RtControl {
495 &mut self.control
496 }
497
498 pub fn stop(self) {}
500}
501
502pub fn play_realtime_with_backend<B: AudioBackend>(
508 backend: &B,
509 plugin: Plugin,
510 config: AudioConfig,
511 command_capacity: usize,
512) -> Result<RtAudioHandle> {
513 let device = backend
514 .default_output_device()
515 .ok_or_else(|| Error::AudioBackendError("No default output device available".into()))?;
516
517 let channels = config.output_channels;
518 let sample_rate = config.sample_rate;
519
520 let (mut runner, control) = RealtimePluginRunner::new(plugin, command_capacity);
521 runner.start()?;
522
523 let mut scratch = AudioBuffers::new(0, channels, config.block_size, sample_rate);
525
526 let data_cb = Box::new(move |data: &mut [f32]| {
527 data.fill(0.0);
528 if channels == 0 {
529 return;
530 }
531 let frames = data.len() / channels;
532 prepare_scratch(&mut scratch, frames);
533
534 if runner.process(&mut scratch).is_ok() {
536 interleave_outputs(&scratch.outputs, data, channels);
537 }
538 });
539
540 let err_cb = Box::new(|e: B::Error| {
541 log::error!("audio stream error: {}", e);
542 });
543
544 let stream = backend
545 .create_output_stream(&device, config, data_cb, err_cb)
546 .map_err(|e| Error::AudioBackendError(format!("Failed to create output stream: {}", e)))?;
547
548 stream
549 .play()
550 .map_err(|e| Error::AudioBackendError(format!("Failed to start stream: {}", e)))?;
551
552 Ok(RtAudioHandle {
553 _stream: Box::new(stream),
554 control,
555 })
556}
557
558#[cfg(test)]
559mod tests {
560 use super::*;
561
562 #[test]
563 fn interleaves_two_channels() {
564 let outputs = vec![vec![1.0, 2.0, 3.0], vec![-1.0, -2.0, -3.0]];
566 let mut out = vec![0.0; 6]; interleave_outputs(&outputs, &mut out, 2);
568 assert_eq!(out, vec![1.0, -1.0, 2.0, -2.0, 3.0, -3.0]);
569 }
570
571 #[test]
572 fn channel_peak_is_max_abs_and_sanitizes_non_finite() {
573 assert_eq!(channel_peak(&[0.1, -0.5, 0.3]), 0.5);
574 assert_eq!(channel_peak(&[]), 0.0);
575 assert_eq!(channel_peak(&[f32::NAN, 0.2, f32::INFINITY]), 0.2);
577 }
578
579 #[test]
580 fn nonneg_f32_bits_are_monotonic_so_fetch_max_is_float_max() {
581 let peaks = [0.0_f32, 1e-6, 0.01, 0.25, 0.5, 0.999, 1.0];
584 for w in peaks.windows(2) {
585 assert!(w[0].to_bits() < w[1].to_bits(), "{} vs {}", w[0], w[1]);
586 }
587 }
588
589 #[test]
590 fn ignores_extra_plugin_channels() {
591 let outputs = vec![vec![1.0, 2.0], vec![3.0, 4.0], vec![9.0, 9.0]];
593 let mut out = vec![0.0; 4];
594 interleave_outputs(&outputs, &mut out, 2);
595 assert_eq!(out, vec![1.0, 3.0, 2.0, 4.0]);
596 }
597
598 #[test]
599 fn leaves_missing_channels_as_silence() {
600 let outputs = vec![vec![0.5, 0.6]];
602 let mut out = vec![0.0; 4];
603 interleave_outputs(&outputs, &mut out, 2);
604 assert_eq!(out, vec![0.5, 0.0, 0.6, 0.0]);
606 }
607
608 #[test]
609 fn zero_channels_is_a_noop() {
610 let outputs = vec![vec![1.0, 2.0]];
611 let mut out = vec![7.0, 7.0];
612 interleave_outputs(&outputs, &mut out, 0);
613 assert_eq!(out, vec![7.0, 7.0]);
614 }
615
616 #[test]
617 fn prepare_scratch_resizes_and_clears() {
618 let mut scratch = AudioBuffers::new(1, 2, 4, 48000.0);
619 scratch.outputs[0][0] = 9.0;
620 prepare_scratch(&mut scratch, 8);
621 assert_eq!(scratch.block_size, 8);
622 assert!(scratch.outputs.iter().all(|c| c.len() == 8));
623 assert!(scratch.inputs.iter().all(|c| c.len() == 8));
624 assert!(scratch.outputs.iter().flatten().all(|&s| s == 0.0));
625 }
626}