use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use rtrb::{Consumer, Producer, RingBuffer};
use crate::{
audio::{AudioBackend, AudioBuffers, AudioConfig, AudioLevels, AudioStream, ChannelLevel},
error::{Error, Result},
midi::MidiEvent,
plugin::Plugin,
realtime::{RealtimePluginRunner, RtControl, TransportCommand},
};
const SIDE_CHANNEL_CAPACITY: usize = 4096;
enum HybridCommand {
Midi { event: MidiEvent, offset: i32 },
Param { id: u32, value: f64 },
Transport(TransportCommand),
Panic,
}
fn channel_peak(buf: &[f32]) -> f32 {
buf.iter()
.map(|&x| if x.is_finite() { x.abs() } else { 0.0 })
.fold(0.0_f32, f32::max)
}
struct AudioSideChannels {
control_rx: Consumer<HybridCommand>,
out_midi_tx: Producer<MidiEvent>,
param_tx: Producer<(u32, f64)>,
levels: Arc<[AtomicU32]>,
}
impl AudioSideChannels {
fn apply_control(&mut self, plugin: &mut Plugin) {
crate::realtime::drain_commands(&mut self.control_rx, |command| match command {
HybridCommand::Midi { event, offset } => {
let _ = plugin.send_midi_event_at(event, offset);
}
HybridCommand::Param { id, value } => {
let _ = plugin.queue_processor_parameter_at(id, value, 0);
}
HybridCommand::Transport(change) => {
change.apply(plugin);
}
HybridCommand::Panic => {
let _ = plugin.midi_panic();
}
});
}
fn publish_levels(&mut self, outputs: &[Vec<f32>]) {
for (ch, atomic) in self.levels.iter().enumerate() {
let peak = outputs.get(ch).map(|b| channel_peak(b)).unwrap_or(0.0);
atomic.fetch_max(peak.to_bits(), Ordering::Relaxed);
}
}
fn publish_feedback(&mut self, plugin: &Plugin) {
for event in plugin.take_output_midi() {
let _ = self.out_midi_tx.push(event);
}
for change in plugin.get_parameter_changes() {
let _ = self.param_tx.push(change);
}
}
}
struct UiSideChannels {
control_tx: Arc<Mutex<Producer<HybridCommand>>>,
out_midi_rx: Mutex<Consumer<MidiEvent>>,
param_rx: Mutex<Consumer<(u32, f64)>>,
levels: Arc<[AtomicU32]>,
}
fn queue_command(tx: &Mutex<Producer<HybridCommand>>, command: HybridCommand) -> bool {
tx.lock()
.map(|mut tx| tx.push(command).is_ok())
.unwrap_or(false)
}
fn make_side_channels(channels: usize) -> (AudioSideChannels, UiSideChannels) {
let (control_tx, control_rx) = RingBuffer::<HybridCommand>::new(SIDE_CHANNEL_CAPACITY);
let (out_midi_tx, out_midi_rx) = RingBuffer::<MidiEvent>::new(SIDE_CHANNEL_CAPACITY);
let (param_tx, param_rx) = RingBuffer::<(u32, f64)>::new(SIDE_CHANNEL_CAPACITY);
let levels: Arc<[AtomicU32]> = (0..channels).map(|_| AtomicU32::new(0)).collect();
let audio = AudioSideChannels {
control_rx,
out_midi_tx,
param_tx,
levels: Arc::clone(&levels),
};
let ui = UiSideChannels {
control_tx: Arc::new(Mutex::new(control_tx)),
out_midi_rx: Mutex::new(out_midi_rx),
param_rx: Mutex::new(param_rx),
levels,
};
(audio, ui)
}
pub struct AudioHandle {
_stream: Box<dyn AudioStream>,
_input_stream: Option<Box<dyn AudioStream>>,
plugin: Arc<Mutex<Plugin>>,
ui: UiSideChannels,
}
#[derive(Clone)]
pub struct MidiSink {
control_tx: Arc<Mutex<Producer<HybridCommand>>>,
}
impl MidiSink {
pub fn send_midi(&self, event: MidiEvent) -> bool {
self.send_midi_at(event, 0)
}
pub fn send_midi_at(&self, event: MidiEvent, sample_offset: i32) -> bool {
queue_command(
&self.control_tx,
HybridCommand::Midi {
event,
offset: sample_offset.max(0),
},
)
}
}
impl AudioHandle {
pub fn lock(&self) -> MutexGuard<'_, Plugin> {
self.plugin
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn try_lock(&self) -> Option<MutexGuard<'_, Plugin>> {
match self.plugin.try_lock() {
Ok(guard) => Some(guard),
Err(std::sync::TryLockError::Poisoned(p)) => Some(p.into_inner()),
Err(std::sync::TryLockError::WouldBlock) => None,
}
}
pub fn send_midi(&self, event: MidiEvent) -> bool {
self.send_midi_at(event, 0)
}
pub fn send_midi_at(&self, event: MidiEvent, sample_offset: i32) -> bool {
queue_command(
&self.ui.control_tx,
HybridCommand::Midi {
event,
offset: sample_offset.max(0),
},
)
}
pub fn midi_sink(&self) -> MidiSink {
MidiSink {
control_tx: Arc::clone(&self.ui.control_tx),
}
}
pub fn set_parameter(&self, id: u32, value: f64) -> bool {
crate::realtime::is_normalized(value)
&& queue_command(&self.ui.control_tx, HybridCommand::Param { id, value })
}
pub fn set_tempo(&self, bpm: f64) -> bool {
if !(bpm.is_finite() && bpm > 0.0) {
return false;
}
queue_command(
&self.ui.control_tx,
HybridCommand::Transport(TransportCommand::Tempo(bpm)),
)
}
pub fn set_time_signature(&self, numerator: i32, denominator: i32) -> bool {
if numerator <= 0 || !matches!(denominator, 1 | 2 | 4 | 8 | 16) {
return false;
}
queue_command(
&self.ui.control_tx,
HybridCommand::Transport(TransportCommand::TimeSignature(numerator, denominator)),
)
}
pub fn set_playing(&self, playing: bool) -> bool {
queue_command(
&self.ui.control_tx,
HybridCommand::Transport(TransportCommand::Playing(playing)),
)
}
pub fn midi_panic(&self) -> bool {
queue_command(&self.ui.control_tx, HybridCommand::Panic)
}
pub fn output_levels(&self) -> AudioLevels {
let channels = self
.ui
.levels
.iter()
.map(|atomic| {
let peak = f32::from_bits(atomic.swap(0, Ordering::Relaxed));
ChannelLevel {
peak,
rms: 0.0,
peak_hold: peak,
}
})
.collect();
AudioLevels { channels }
}
pub fn drain_output_midi(&self) -> Vec<MidiEvent> {
let mut out = Vec::new();
if let Ok(mut rx) = self.ui.out_midi_rx.lock() {
while let Ok(event) = rx.pop() {
out.push(event);
}
}
out
}
pub fn drain_parameter_changes(&self) -> Vec<(u32, f64)> {
let mut out = Vec::new();
if let Ok(mut rx) = self.ui.param_rx.lock() {
while let Ok(change) = rx.pop() {
out.push(change);
}
}
out
}
pub fn plugin(&self) -> Arc<Mutex<Plugin>> {
Arc::clone(&self.plugin)
}
pub fn stop(self) {}
}
pub(crate) fn interleave_outputs(outputs: &[Vec<f32>], out: &mut [f32], channels: usize) {
if channels == 0 {
return;
}
let frames = out.len() / channels;
for ch in 0..channels.min(outputs.len()) {
let src = &outputs[ch];
for frame in 0..frames.min(src.len()) {
out[frame * channels + ch] = src[frame];
}
}
}
fn push_capture_frames(producer: &mut Producer<f32>, data: &[f32], channels: usize) -> usize {
if channels == 0 {
return 0;
}
let room_frames = producer.slots() / channels;
let mut pushed = 0;
for frame in data.chunks_exact(channels).take(room_frames) {
for &sample in frame {
let _ = producer.push(sample);
}
pushed += 1;
}
pushed
}
fn pop_capture_frames(
consumer: &mut Consumer<f32>,
inputs: &mut [Vec<f32>],
frames: usize,
) -> usize {
let channels = inputs.len();
if channels == 0 {
return 0;
}
let available = (consumer.slots() / channels).min(frames);
for f in 0..available {
for ch in inputs.iter_mut() {
let sample = consumer.pop().unwrap_or(0.0);
if let Some(slot) = ch.get_mut(f) {
*slot = sample;
}
}
}
available
}
fn bridge_ring_capacity(block_size: usize, channels: usize) -> usize {
const MIN_BRIDGE_FRAMES: usize = 1024;
const MAX_BRIDGE_FRAMES: usize = 1 << 18; let frames = block_size
.saturating_mul(8)
.clamp(MIN_BRIDGE_FRAMES, MAX_BRIDGE_FRAMES);
frames.saturating_mul(channels.max(1))
}
fn prepare_scratch(scratch: &mut AudioBuffers, frames: usize) {
for ch in &mut scratch.outputs {
if ch.len() != frames {
ch.resize(frames, 0.0);
}
ch.fill(0.0);
}
for ch in &mut scratch.inputs {
if ch.len() != frames {
ch.resize(frames, 0.0);
}
ch.fill(0.0);
}
scratch.block_size = frames;
}
pub fn play_with_backend<B: AudioBackend>(
backend: &B,
plugin: Plugin,
config: AudioConfig,
) -> Result<AudioHandle> {
let device = backend
.default_output_device()
.ok_or_else(|| Error::AudioBackendError("No default output device available".into()))?;
let channels = config.output_channels;
let sample_rate = config.sample_rate;
let plugin = Arc::new(Mutex::new(plugin));
plugin
.lock()
.unwrap_or_else(|p| p.into_inner())
.start_processing()?;
let plugin_cb = Arc::clone(&plugin);
let (mut side, ui) = make_side_channels(channels);
let mut scratch = AudioBuffers::new(0, channels, config.block_size, sample_rate);
let data_cb = Box::new(move |data: &mut [f32]| {
data.fill(0.0);
if channels == 0 {
return;
}
let frames = data.len() / channels;
prepare_scratch(&mut scratch, frames);
let mut p = match plugin_cb.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
side.apply_control(&mut p);
if p.process_audio(&mut scratch).is_ok() {
interleave_outputs(&scratch.outputs, data, channels);
side.publish_levels(&scratch.outputs);
}
side.publish_feedback(&p);
});
let err_cb = Box::new(|e: B::Error| {
log::error!("audio stream error: {}", e);
});
let stream = backend
.create_output_stream(&device, config, data_cb, err_cb)
.map_err(|e| Error::AudioBackendError(format!("Failed to create output stream: {}", e)))?;
stream
.play()
.map_err(|e| Error::AudioBackendError(format!("Failed to start stream: {}", e)))?;
Ok(AudioHandle {
_stream: Box::new(stream),
_input_stream: None,
plugin,
ui,
})
}
pub fn play_with_input_backend<B: AudioBackend>(
backend: &B,
plugin: Plugin,
config: AudioConfig,
) -> Result<AudioHandle> {
let in_device = backend
.default_input_device()
.ok_or_else(|| Error::AudioBackendError("No default input device available".into()))?;
let out_device = backend
.default_output_device()
.ok_or_else(|| Error::AudioBackendError("No default output device available".into()))?;
let in_channels = config.input_channels.max(1);
let out_channels = config.output_channels;
let sample_rate = config.sample_rate;
let plugin = Arc::new(Mutex::new(plugin));
plugin
.lock()
.unwrap_or_else(|p| p.into_inner())
.start_processing()?;
let (mut producer, mut consumer) =
rtrb::RingBuffer::<f32>::new(bridge_ring_capacity(config.block_size, in_channels));
let in_data_cb = Box::new(move |data: &[f32]| {
push_capture_frames(&mut producer, data, in_channels);
});
let in_err_cb = Box::new(|e: B::Error| log::error!("input stream error: {}", e));
let input_stream = backend
.create_input_stream(&in_device, config, in_data_cb, in_err_cb)
.map_err(|e| Error::AudioBackendError(format!("Failed to create input stream: {}", e)))?;
let plugin_cb = Arc::clone(&plugin);
let (mut side, ui) = make_side_channels(out_channels);
let mut scratch = AudioBuffers::new(in_channels, out_channels, config.block_size, sample_rate);
let out_data_cb = Box::new(move |data: &mut [f32]| {
data.fill(0.0);
if out_channels == 0 {
return;
}
let frames = data.len() / out_channels;
prepare_scratch(&mut scratch, frames);
pop_capture_frames(&mut consumer, &mut scratch.inputs, frames);
let mut p = match plugin_cb.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
side.apply_control(&mut p);
if p.process_audio(&mut scratch).is_ok() {
interleave_outputs(&scratch.outputs, data, out_channels);
side.publish_levels(&scratch.outputs);
}
side.publish_feedback(&p);
});
let out_err_cb = Box::new(|e: B::Error| log::error!("output stream error: {}", e));
let output_stream = backend
.create_output_stream(&out_device, config, out_data_cb, out_err_cb)
.map_err(|e| Error::AudioBackendError(format!("Failed to create output stream: {}", e)))?;
input_stream
.play()
.map_err(|e| Error::AudioBackendError(format!("Failed to start input stream: {}", e)))?;
output_stream
.play()
.map_err(|e| Error::AudioBackendError(format!("Failed to start output stream: {}", e)))?;
Ok(AudioHandle {
_stream: Box::new(output_stream),
_input_stream: Some(Box::new(input_stream)),
plugin,
ui,
})
}
pub struct RtAudioHandle {
_stream: Box<dyn AudioStream>,
control: RtControl,
}
impl RtAudioHandle {
pub fn control(&mut self) -> &mut RtControl {
&mut self.control
}
pub fn stop(self) {}
}
pub fn play_realtime_with_backend<B: AudioBackend>(
backend: &B,
plugin: Plugin,
config: AudioConfig,
command_capacity: usize,
) -> Result<RtAudioHandle> {
let device = backend
.default_output_device()
.ok_or_else(|| Error::AudioBackendError("No default output device available".into()))?;
let channels = config.output_channels;
let sample_rate = config.sample_rate;
let (mut runner, control) = RealtimePluginRunner::new(plugin, command_capacity);
runner.start()?;
let mut scratch = AudioBuffers::new(0, channels, config.block_size, sample_rate);
let data_cb = Box::new(move |data: &mut [f32]| {
data.fill(0.0);
if channels == 0 {
return;
}
let frames = data.len() / channels;
prepare_scratch(&mut scratch, frames);
if runner.process(&mut scratch).is_ok() {
interleave_outputs(&scratch.outputs, data, channels);
}
});
let err_cb = Box::new(|e: B::Error| {
log::error!("audio stream error: {}", e);
});
let stream = backend
.create_output_stream(&device, config, data_cb, err_cb)
.map_err(|e| Error::AudioBackendError(format!("Failed to create output stream: {}", e)))?;
stream
.play()
.map_err(|e| Error::AudioBackendError(format!("Failed to start stream: {}", e)))?;
Ok(RtAudioHandle {
_stream: Box::new(stream),
control,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn interleaves_two_channels() {
let outputs = vec![vec![1.0, 2.0, 3.0], vec![-1.0, -2.0, -3.0]];
let mut out = vec![0.0; 6]; interleave_outputs(&outputs, &mut out, 2);
assert_eq!(out, vec![1.0, -1.0, 2.0, -2.0, 3.0, -3.0]);
}
#[test]
fn hybrid_midi_carries_sample_offset() {
let (mut tx, mut rx) = RingBuffer::<HybridCommand>::new(4);
tx.push(HybridCommand::Midi {
event: MidiEvent::NoteOn {
channel: crate::midi::MidiChannel::Ch1,
note: 60,
velocity: 100,
},
offset: 200,
})
.expect("queue scheduled note");
match rx.pop().expect("note queued") {
HybridCommand::Midi { offset, .. } => assert_eq!(offset, 200),
_ => panic!("expected a MIDI command"),
}
}
#[test]
fn channel_peak_is_max_abs_and_sanitizes_non_finite() {
assert_eq!(channel_peak(&[0.1, -0.5, 0.3]), 0.5);
assert_eq!(channel_peak(&[]), 0.0);
assert_eq!(channel_peak(&[f32::NAN, 0.2, f32::INFINITY]), 0.2);
}
#[test]
fn nonneg_f32_bits_are_monotonic_so_fetch_max_is_float_max() {
let peaks = [0.0_f32, 1e-6, 0.01, 0.25, 0.5, 0.999, 1.0];
for w in peaks.windows(2) {
assert!(w[0].to_bits() < w[1].to_bits(), "{} vs {}", w[0], w[1]);
}
}
#[test]
fn ignores_extra_plugin_channels() {
let outputs = vec![vec![1.0, 2.0], vec![3.0, 4.0], vec![9.0, 9.0]];
let mut out = vec![0.0; 4];
interleave_outputs(&outputs, &mut out, 2);
assert_eq!(out, vec![1.0, 3.0, 2.0, 4.0]);
}
#[test]
fn leaves_missing_channels_as_silence() {
let outputs = vec![vec![0.5, 0.6]];
let mut out = vec![0.0; 4];
interleave_outputs(&outputs, &mut out, 2);
assert_eq!(out, vec![0.5, 0.0, 0.6, 0.0]);
}
#[test]
fn zero_channels_is_a_noop() {
let outputs = vec![vec![1.0, 2.0]];
let mut out = vec![7.0, 7.0];
interleave_outputs(&outputs, &mut out, 0);
assert_eq!(out, vec![7.0, 7.0]);
}
#[test]
fn bridge_overflow_drops_whole_frames_never_splits_one() {
const CH: usize = 2;
let (mut producer, mut consumer) = RingBuffer::<f32>::new(4 * CH);
let data: Vec<f32> = (0..8).flat_map(|f| [f as f32, -(f as f32)]).collect();
let pushed = push_capture_frames(&mut producer, &data, CH);
assert_eq!(pushed, 4, "only the frames that fit are pushed");
let mut inputs = vec![vec![0.0f32; 8], vec![0.0f32; 8]];
let popped = pop_capture_frames(&mut consumer, &mut inputs, 8);
assert_eq!(popped, 4);
assert_eq!(inputs[0], vec![0.0, 1.0, 2.0, 3.0, 0.0, 0.0, 0.0, 0.0]);
assert_eq!(inputs[1], vec![-0.0, -1.0, -2.0, -3.0, 0.0, 0.0, 0.0, 0.0]);
}
#[test]
fn bridge_stays_frame_aligned_across_overflow_rounds() {
const CH: usize = 2;
let (mut producer, mut consumer) = RingBuffer::<f32>::new(4 * CH);
for round in 0..5 {
let base = round * 100;
let data: Vec<f32> = (0..6)
.flat_map(|f| [(base + f) as f32, -((base + f) as f32)])
.collect();
push_capture_frames(&mut producer, &data, CH);
let mut inputs = vec![vec![0.0f32; 3], vec![0.0f32; 3]];
let popped = pop_capture_frames(&mut consumer, &mut inputs, 3);
for f in 0..popped {
assert_eq!(
inputs[1][f], -inputs[0][f],
"round {round} frame {f} lost channel alignment: {:?}",
inputs
);
}
}
}
#[test]
fn bridge_ignores_a_partial_trailing_frame() {
const CH: usize = 3;
let (mut producer, mut consumer) = RingBuffer::<f32>::new(8 * CH);
let data = [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0];
assert_eq!(push_capture_frames(&mut producer, &data, CH), 2);
assert_eq!(consumer.slots() % CH, 0, "ring holds whole frames only");
let mut inputs = vec![vec![0.0f32; 2], vec![0.0f32; 2], vec![0.0f32; 2]];
assert_eq!(pop_capture_frames(&mut consumer, &mut inputs, 2), 2);
assert_eq!(inputs[0], vec![1.0, 4.0]);
assert_eq!(inputs[1], vec![2.0, 5.0]);
assert_eq!(inputs[2], vec![3.0, 6.0]);
}
#[test]
fn bridge_consumer_leaves_an_incomplete_frame_alone() {
const CH: usize = 2;
let (mut producer, mut consumer) = RingBuffer::<f32>::new(4 * CH);
producer.push(1.0).expect("room for the first sample");
let mut inputs = vec![vec![0.0f32; 2], vec![0.0f32; 2]];
assert_eq!(pop_capture_frames(&mut consumer, &mut inputs, 2), 0);
assert_eq!(consumer.slots(), 1, "the half frame is still queued");
producer.push(-1.0).expect("room for the second sample");
assert_eq!(pop_capture_frames(&mut consumer, &mut inputs, 2), 1);
assert_eq!(inputs[0][0], 1.0);
assert_eq!(inputs[1][0], -1.0);
}
#[test]
fn bridge_capacity_is_a_whole_number_of_frames() {
for block_size in [0usize, 1, 32, 64, 128, 512, 1024, usize::MAX] {
for channels in 1..=8usize {
let cap = bridge_ring_capacity(block_size, channels);
assert_eq!(
cap % channels,
0,
"block {block_size} x {channels} ch: capacity {cap} splits a frame"
);
assert!(
cap >= channels,
"block {block_size} x {channels} ch: capacity {cap} holds no frame"
);
}
}
assert_ne!(2048 % 3, 0);
}
#[test]
fn bridge_capacity_survives_zero_channels() {
assert!(bridge_ring_capacity(512, 0) > 0);
}
#[test]
fn queue_command_drops_on_a_full_ring() {
let (tx, mut rx) = RingBuffer::<HybridCommand>::new(2);
let tx = Mutex::new(tx);
assert!(queue_command(&tx, HybridCommand::Panic));
assert!(queue_command(&tx, HybridCommand::Panic));
assert!(!queue_command(&tx, HybridCommand::Panic), "ring is full");
assert_eq!(rx.slots(), 2);
assert!(rx.pop().is_ok());
}
#[test]
fn prepare_scratch_resizes_and_clears() {
let mut scratch = AudioBuffers::new(1, 2, 4, 48000.0);
scratch.outputs[0][0] = 9.0;
prepare_scratch(&mut scratch, 8);
assert_eq!(scratch.block_size, 8);
assert!(scratch.outputs.iter().all(|c| c.len() == 8));
assert!(scratch.inputs.iter().all(|c| c.len() == 8));
assert!(scratch.outputs.iter().flatten().all(|&s| s == 0.0));
}
}