diff --git a/.jules/bolt.md b/.jules/bolt.md new file mode 100644 index 0000000..163408a --- /dev/null +++ b/.jules/bolt.md @@ -0,0 +1,3 @@ +## 2025-05-15 - [Optimization Pattern: Zero-Allocation Hot Loops in Audio Processing] +**Learning:** Performance in audio processing is heavily impacted by per-block heap allocations. The `livi` crate's API allows passing iterators for port connections, which can be leveraged to eliminate `Vec` allocations in the hot loop. Additionally, double-buffering can be avoided by copying directly from the decoder's buffer into a planar plugin buffer. +**Action:** Always check if library APIs support iterators or pre-allocated slices for high-frequency operations. When dealing with multi-channel audio, ensure that channel-copying logic is robust against channel mismatches to avoid panics. diff --git a/src/audio.rs b/src/audio.rs index 0861f7b..099addb 100644 --- a/src/audio.rs +++ b/src/audio.rs @@ -44,6 +44,8 @@ pub fn setup_input_audio(input_path: &Path) -> Result { use symphonia::core::conv::IntoSample; +/// This function is deprecated and will be removed in a future version. +/// Use `copy_to_planar_f32` for better performance. pub fn copy_to_f32_buffer(buf: AudioBufferRef<'_>, result: &mut [Vec]) { macro_rules! convert_buf { ($buf:expr) => { @@ -72,3 +74,46 @@ pub fn copy_to_f32_buffer(buf: AudioBufferRef<'_>, result: &mut [Vec]) { AudioBufferRef::F64(b) => convert_buf!(b), } } + +pub fn copy_to_planar_f32( + buf: &AudioBufferRef<'_>, + dst: &mut [f32], + src_offset: usize, + dst_offset: usize, + num_frames: usize, + block_size: usize, +) { + let dst_channels = dst.len() / block_size; + macro_rules! convert_buf { + ($buf:expr) => { + for ch in 0..std::cmp::min($buf.spec().channels.count(), dst_channels) { + let dst_start = ch * block_size + dst_offset; + let src_ch = $buf.chan(ch); + let src_slice = &src_ch[src_offset..src_offset + num_frames]; + for (i, &s) in src_slice.iter().enumerate() { + dst[dst_start + i] = s.into_sample(); + } + } + }; + } + + match buf { + AudioBufferRef::U8(b) => convert_buf!(b), + AudioBufferRef::U16(b) => convert_buf!(b), + AudioBufferRef::U24(b) => convert_buf!(b), + AudioBufferRef::U32(b) => convert_buf!(b), + AudioBufferRef::S8(b) => convert_buf!(b), + AudioBufferRef::S16(b) => convert_buf!(b), + AudioBufferRef::S24(b) => convert_buf!(b), + AudioBufferRef::S32(b) => convert_buf!(b), + AudioBufferRef::F32(b) => { + for ch in 0..std::cmp::min(b.spec().channels.count(), dst_channels) { + let dst_start = ch * block_size + dst_offset; + let src_ch = b.chan(ch); + let src_slice = &src_ch[src_offset..src_offset + num_frames]; + dst[dst_start..dst_start + num_frames].copy_from_slice(src_slice); + } + } + AudioBufferRef::F64(b) => convert_buf!(b), + } +} diff --git a/src/processing.rs b/src/processing.rs index fc3d27a..c1bfe9e 100644 --- a/src/processing.rs +++ b/src/processing.rs @@ -9,7 +9,7 @@ use symphonia::core::errors::Error as SymphoniaError; use symphonia::core::formats::{FormatReader, Track}; use crate::plugin::PluginInstance; -use crate::audio::copy_to_f32_buffer; +use crate::audio::copy_to_planar_f32; pub struct ProcessingContext<'a> { pub instances: &'a mut [PluginInstance], @@ -55,15 +55,51 @@ impl ProgressReporter { } } +fn run_instance<'a>( + control_inputs: &'a [f32], + control_outputs: &'a mut [f32], + audio_inputs: impl ExactSizeIterator, + audio_outputs: impl ExactSizeIterator, + atom_sequence_inputs: &'a mut [livi::event::LV2AtomSequence], + atom_sequence_outputs: &'a mut [livi::event::LV2AtomSequence], + instance: &mut livi::Instance, + block_size: usize, + worker_manager: &WorkerManager, +) -> Result<()> { + for seq in atom_sequence_inputs.iter_mut() { + seq.clear(); + } + for seq in atom_sequence_outputs.iter_mut() { + seq.clear(); + } + + let ports = EmptyPortConnections::new(block_size) + .with_control_inputs(control_inputs.iter()) + .with_control_outputs(control_outputs.iter_mut()) + .with_audio_inputs(audio_inputs) + .with_audio_outputs(audio_outputs) + .with_atom_sequence_inputs(atom_sequence_inputs.iter()) + .with_atom_sequence_outputs(atom_sequence_outputs.iter_mut()); + + unsafe { + instance + .run(ports) + .map_err(|e| anyhow!("Plugin run failed: {:?}", e))?; + } + worker_manager.run_workers(); + Ok(()) +} + pub fn process_audio( format: &mut dyn FormatReader, track: Track, mut ctx: ProcessingContext<'_>, ) -> Result<()> { - let mut decoder = symphonia::default::get_codecs().make(&track.codec_params, &DecoderOptions::default())?; + let mut decoder = symphonia::default::get_codecs() + .make(&track.codec_params, &DecoderOptions::default())?; let sample_rate = track.codec_params.sample_rate.unwrap_or(44100); let total_frames = track.codec_params.n_frames.unwrap_or(0); - + let first_plugin_input_count = ctx.instances[0].plugin.port_counts().audio_inputs; let mut input_buffer = vec![0.0f32; ctx.block_size * first_plugin_input_count]; let needs_upmix = ctx.input_channels == 1 && first_plugin_input_count == 2; @@ -73,116 +109,151 @@ pub fn process_audio( let mut is_draining = false; let mut frames_processed = 0u64; let mut frames_decoded = 0u64; - + let reporter = ProgressReporter::new(sample_rate, total_frames); - let mut decode_buffer: Option>> = None; loop { - let frames_in_chunk; if is_draining { - if drain_remaining == 0 { break; } - let chunk = min(ctx.block_size, drain_remaining); - frames_in_chunk = chunk; - let buf = decode_buffer.get_or_insert_with(|| vec![vec![0.0f32; ctx.block_size]; ctx.input_channels]); - for ch_buf in buf.iter_mut().take(ctx.input_channels) { - ch_buf.resize(chunk, 0.0); - ch_buf.fill(0.0); + if drain_remaining == 0 { + break; } - drain_remaining = drain_remaining.saturating_sub(chunk); - } else { - match format.next_packet() { - Ok(packet) => match decoder.decode(&packet) { - Ok(decoded) => { - let num_frames = decoded.frames(); - let num_chans = decoded.spec().channels.count(); - let buf = decode_buffer.get_or_insert_with(|| vec![vec![0.0f32; 0]; num_chans]); - for ch_buf in buf.iter_mut().take(num_chans) { - ch_buf.resize(num_frames, 0.0); + let chunk_size = min(ctx.block_size, drain_remaining); + input_buffer.fill(0.0); + process_chunk( + &mut ctx, + &mut input_buffer, + chunk_size, + needs_upmix, + first_plugin_input_count, + )?; + drain_remaining = drain_remaining.saturating_sub(chunk_size); + frames_processed += chunk_size as u64; + reporter.tick( + frames_processed, + frames_decoded, + is_draining, + drain_frames, + drain_remaining, + ctx.block_size, + ); + continue; + } + + match format.next_packet() { + Ok(packet) => match decoder.decode(&packet) { + Ok(decoded) => { + let num_frames = decoded.frames(); + let mut src_offset = 0; + while src_offset < num_frames { + let chunk_size = min(num_frames - src_offset, ctx.block_size); + copy_to_planar_f32( + &decoded, + &mut input_buffer, + src_offset, + 0, + chunk_size, + ctx.block_size, + ); + if chunk_size < ctx.block_size { + for ch in 0..first_plugin_input_count { + if !needs_upmix || ch == 0 { + input_buffer[ch * ctx.block_size + chunk_size..(ch + 1) * ctx.block_size].fill(0.0); + } + } } - copy_to_f32_buffer(decoded, buf); - frames_in_chunk = num_frames; - frames_decoded += num_frames as u64; - } - Err(e) => { - eprintln!("\nWarning: decode error: {}", e); - continue; + + process_chunk( + &mut ctx, + &mut input_buffer, + chunk_size, + needs_upmix, + first_plugin_input_count, + )?; + src_offset += chunk_size; + frames_processed += chunk_size as u64; + frames_decoded += chunk_size as u64; + reporter.tick( + frames_processed, + frames_decoded, + is_draining, + drain_frames, + drain_remaining, + ctx.block_size, + ); } - }, - Err(SymphoniaError::IoError(ref e)) if e.kind() == std::io::ErrorKind::UnexpectedEof => { - is_draining = true; - eprintln!("\nInput complete, draining plugin tails..."); - continue; } - Err(e) => return Err(e).context("Error reading packet"), + Err(e) => { + eprintln!("\nWarning: decode error: {}", e); + } + }, + Err(SymphoniaError::IoError(ref e)) if e.kind() == std::io::ErrorKind::UnexpectedEof => { + is_draining = true; + eprintln!("\nInput complete, draining plugin tails..."); } + Err(e) => return Err(e).context("Error reading packet"), } + } - let audio_buf = decode_buffer.as_ref().unwrap(); - let mut offset = 0; - while offset < frames_in_chunk { - let chunk_size = min(frames_in_chunk - offset, ctx.block_size); - - for ch in 0..first_plugin_input_count { - let dst_start = ch * ctx.block_size; - if ch < audio_buf.len() { - let src = &audio_buf[ch]; - input_buffer[dst_start..dst_start + chunk_size].copy_from_slice(&src[offset..offset + chunk_size]); - if chunk_size < ctx.block_size { - input_buffer[dst_start + chunk_size..dst_start + ctx.block_size].fill(0.0); - } - } else if !needs_upmix { - input_buffer[dst_start..dst_start + ctx.block_size].fill(0.0); - } - } - if needs_upmix { - input_buffer.copy_within(0..ctx.block_size, ctx.block_size); - } + reporter.finish(frames_processed); + ctx.writer.finalize()?; + Ok(()) +} - for i in 0..ctx.instances.len() { - let (prev, curr) = ctx.instances.split_at_mut(i); - let inst = &mut curr[0]; - - let (audio_inputs, audio_outputs) = if i == 0 { - let inputs: Vec<&[f32]> = (0..first_plugin_input_count).map(|ch| &input_buffer[ch * ctx.block_size..(ch + 1) * ctx.block_size]).collect(); - let outputs: Vec<&mut [f32]> = inst.audio_outputs.iter_mut().map(|b| b.as_mut_slice()).collect(); - (inputs, outputs) - } else { - let prev_inst = &prev[i-1]; - let inputs: Vec<&[f32]> = prev_inst.audio_outputs.iter().map(|b| b.as_slice()).collect(); - let outputs: Vec<&mut [f32]> = inst.audio_outputs.iter_mut().map(|b| b.as_mut_slice()).collect(); - (inputs, outputs) - }; - - for seq in &mut inst.atom_sequence_inputs { seq.clear(); } - for seq in &mut inst.atom_sequence_outputs { seq.clear(); } - - let ports = EmptyPortConnections::new(ctx.block_size) - .with_control_inputs(inst.control_inputs.iter()) - .with_control_outputs(inst.control_outputs.iter_mut()) - .with_audio_inputs(audio_inputs.into_iter()) - .with_audio_outputs(audio_outputs.into_iter()) - .with_atom_sequence_inputs(inst.atom_sequence_inputs.iter()) - .with_atom_sequence_outputs(inst.atom_sequence_outputs.iter_mut()); - - unsafe { inst.instance.run(ports).map_err(|e| anyhow!("Plugin run failed: {:?}", e))?; } - ctx.worker_manager.run_workers(); - } +fn process_chunk( + ctx: &mut ProcessingContext<'_>, + input_buffer: &mut [f32], + chunk_size: usize, + needs_upmix: bool, + first_plugin_input_count: usize, +) -> Result<()> { + if needs_upmix { + input_buffer.copy_within(0..ctx.block_size, ctx.block_size); + } - let final_outputs = &ctx.instances.last().unwrap().audio_outputs; - for frame_idx in 0..chunk_size { - for ch_buf in final_outputs { - ctx.writer.write_sample(ch_buf[frame_idx]).context("Failed to write sample")?; - } - } + for i in 0..ctx.instances.len() { + let (prev, curr) = ctx.instances.split_at_mut(i); + let inst = &mut curr[0]; - frames_processed += chunk_size as u64; - offset += chunk_size; + if i == 0 { + let inputs = (0..first_plugin_input_count) + .map(|ch| &input_buffer[ch * ctx.block_size..(ch + 1) * ctx.block_size]); + let outputs = inst.audio_outputs.iter_mut().map(|b| b.as_mut_slice()); + run_instance( + &inst.control_inputs, + &mut inst.control_outputs, + inputs, + outputs, + &mut inst.atom_sequence_inputs, + &mut inst.atom_sequence_outputs, + &mut inst.instance, + ctx.block_size, + ctx.worker_manager, + )?; + } else { + let prev_inst = &prev[i - 1]; + let inputs = prev_inst.audio_outputs.iter().map(|b| b.as_slice()); + let outputs = inst.audio_outputs.iter_mut().map(|b| b.as_mut_slice()); + run_instance( + &inst.control_inputs, + &mut inst.control_outputs, + inputs, + outputs, + &mut inst.atom_sequence_inputs, + &mut inst.atom_sequence_outputs, + &mut inst.instance, + ctx.block_size, + ctx.worker_manager, + )?; + }; + } - reporter.tick(frames_processed, frames_decoded, is_draining, drain_frames, drain_remaining, ctx.block_size); + let final_outputs = &ctx.instances.last().unwrap().audio_outputs; + for frame_idx in 0..chunk_size { + for ch_buf in final_outputs { + ctx.writer + .write_sample(ch_buf[frame_idx]) + .context("Failed to write sample")?; } } - - reporter.finish(frames_processed); - ctx.writer.finalize()?; Ok(()) } diff --git a/tests/audio/output/Sound-Bite-Space-Shuttle-Launch-Comm_UHF-Capcom-Wx-Update.wav.wav b/tests/audio/output/Sound-Bite-Space-Shuttle-Launch-Comm_UHF-Capcom-Wx-Update.wav.wav new file mode 100644 index 0000000..0fe8ee7 Binary files /dev/null and b/tests/audio/output/Sound-Bite-Space-Shuttle-Launch-Comm_UHF-Capcom-Wx-Update.wav.wav differ diff --git a/tests/audio/output/kassoulet-headhache-ringtone-2.mp3.wav b/tests/audio/output/kassoulet-headhache-ringtone-2.mp3.wav new file mode 100644 index 0000000..6bd1cc5 Binary files /dev/null and b/tests/audio/output/kassoulet-headhache-ringtone-2.mp3.wav differ diff --git a/tests/audio/output/original_Acid.wav.wav b/tests/audio/output/original_Acid.wav.wav new file mode 100644 index 0000000..b95db7a Binary files /dev/null and b/tests/audio/output/original_Acid.wav.wav differ diff --git a/tests/audio/output/original_Beat2.ogg.wav b/tests/audio/output/original_Beat2.ogg.wav new file mode 100644 index 0000000..254a5de Binary files /dev/null and b/tests/audio/output/original_Beat2.ogg.wav differ