1use js_sys::{Array, Atomics, SharedArrayBuffer, Uint32Array}; 2use std::cmp; 3use wasm_bindgen::prelude::*; 4use wasm_bindgen_futures::JsFuture; 5use web_sys::{AudioContext, AudioWorkletNode, AudioWorkletNodeOptions, ChannelCountMode}; 6 7pub const SAMPLE_RATE: u32 = 48000; 8 9const HEADER_LEN: u32 = 2; 10const HEADER_LEN_BYTES: u32 = HEADER_LEN * 4; 11const START_INDEX: u32 = 0; 12const END_INDEX: u32 = 1; 13 14pub const BUFFER_LEN_SAMPLES: u32 = 8192; 15const BUFFER_LEN_BYTES: u32 = BUFFER_LEN_SAMPLES * 4; 16const BUFFER_INDEX_MASK: u32 = BUFFER_LEN_SAMPLES - 1; 17 18#[derive(Debug, Clone, Copy, PartialEq, Eq)] 19pub enum EnqueueResult { 20 Successful, 21 BufferFull, 22}
A very simple lock-free queue implemented using a circular buffer. The header contains two 32-bit integers containing the current start and exclusive end indices.
34impl Default for AudioQueue { 35 fn default() -> Self { 36 Self::new() 37 } 38} 39 40impl AudioQueue { 41 pub fn new() -> Self { 42 let header = SharedArrayBuffer::new(HEADER_LEN_BYTES); 43 let buffer = SharedArrayBuffer::new(BUFFER_LEN_BYTES); 44 Self::from_buffers(header, buffer) 45 } 46 47 pub fn try_from_js_value(value: JsValue) -> Result<Self, JsValue> { 48 let array = value.dyn_into::<Array>()?; 49 let header = array.get(0).dyn_into::<SharedArrayBuffer>()?; 50 let buffer = array.get(1).dyn_into::<SharedArrayBuffer>()?; 51 Ok(Self::from_buffers(header, buffer)) 52 } 53 54 pub fn from_buffers(header: SharedArrayBuffer, buffer: SharedArrayBuffer) -> Self { 55 let header_typed = Uint32Array::new(&header); 56 let buffer_typed = Uint32Array::new(&buffer); 57 Self { header, header_typed, buffer, buffer_typed } 58 } 59 60 pub fn push_if_space( 61 &self, 62 (sample_l, sample_r): (f32, f32), 63 ) -> Result<EnqueueResult, JsValue> { 64 let mut end = Atomics::load(&self.header_typed, END_INDEX)? as u32; 65 let start = Atomics::load(&self.header_typed, START_INDEX)? as u32; 66 67 for sample in [sample_l, sample_r] { 68 if end == start.wrapping_sub(1) & BUFFER_INDEX_MASK { 69 return Ok(EnqueueResult::BufferFull); 70 } 71 72 Atomics::store(&self.buffer_typed, end, sample.to_bits() as i32)?; 73 end = (end + 1) & BUFFER_INDEX_MASK; 74 } 75 76 Atomics::store(&self.header_typed, END_INDEX, end as i32)?; 77 78 Ok(EnqueueResult::Successful) 79 } 80 81 pub fn drain_into(&self, out: &mut Vec<f32>, limit: u32) -> Result<(), JsValue> { 82 let loaded_start = Atomics::load(&self.header_typed, START_INDEX)? as u32; 83 let end = Atomics::load(&self.header_typed, END_INDEX)? as u32; 84 85 let queue_len = if loaded_start <= end { 86 end - loaded_start 87 } else { 88 end + BUFFER_LEN_SAMPLES - loaded_start 89 }; 90 let drain_len = cmp::min(queue_len as usize, limit as usize); 91 92 let mut start = loaded_start; 93 for _ in 0..drain_len { 94 let value = Atomics::load(&self.buffer_typed, start)?; 95 let sample = f32::from_bits(value as u32); 96 out.push(sample); 97 98 start = (start + 1) & BUFFER_INDEX_MASK; 99 } 100 101 if start != loaded_start { 102 Atomics::store(&self.header_typed, START_INDEX, start as i32)?; 103 } 104 105 Ok(()) 106 } 107 108 pub fn len(&self) -> Result<u32, JsValue> { 109 let end = Atomics::load(&self.header_typed, END_INDEX)? as u32; 110 let start = Atomics::load(&self.header_typed, START_INDEX)? as u32; 111 112 if start <= end { Ok(end - start) } else { Ok(end + BUFFER_LEN_SAMPLES - start) } 113 } 114 115 fn to_js_value(&self) -> JsValue { 116 Array::of2(&self.header, &self.buffer).into() 117 } 118} 119 120#[wasm_bindgen] 121pub struct AudioProcessor { 122 audio_queue: AudioQueue, 123 buffer: Vec<f32>, 124} 125 126#[wasm_bindgen] 127impl AudioProcessor { 128 #[wasm_bindgen(constructor)] 129 pub fn new(audio_queue: JsValue) -> AudioProcessor { 130 let audio_queue = AudioQueue::try_from_js_value(audio_queue) 131 .expect("Unable to initialize audio queue in audio worklet processor"); 132 133 AudioProcessor { audio_queue, buffer: Vec::with_capacity(BUFFER_LEN_SAMPLES as usize) } 134 } 135 136 pub fn process(&mut self, output_l: &mut [f32], output_r: &mut [f32]) { 137 self.buffer.clear(); 138 self.audio_queue 139 .drain_into(&mut self.buffer, 2 * output_l.len() as u32) 140 .expect("Unable to drain audio queue"); 141 142 for (chunk, (out_l, out_r)) in 143 self.buffer.chunks_exact(2).zip(output_l.iter_mut().zip(output_r.iter_mut())) 144 { 145 let &[sample_l, sample_r] = chunk else { unreachable!("chunks_exact(2)") }; 146 *out_l = sample_l; 147 *out_r = sample_r; 148 } 149 } 150} 151 152pub async fn initialize_audio_worklet( 153 audio_ctx: &AudioContext, 154 audio_queue: &AudioQueue, 155) -> Result<AudioWorkletNode, JsValue> { 156 // Polyfill for TextDecoder and TextEncoder; needs to run before loading audio processor module 157 run_polyfill(); 158 159 // Append a random query parameter because Firefox caches this file way too aggressively and 160 // Ctrl+Shift+R doesn't force a reload because it's not loaded on page load. The file itself is 161 // less than 1KB and is only loaded at most once per page load, so not a big deal to not cache it. 162 let module_url = format!("./js/audio-processor.js?r={}", rand::random::<u32>()); 163 JsFuture::from(audio_ctx.audio_worklet()?.add_module(&module_url)?).await?; 164 165 let node_options = AudioWorkletNodeOptions::new(); 166 node_options.set_channel_count_mode(ChannelCountMode::Explicit); 167 node_options.set_output_channel_count(&Array::of1(&JsValue::from(2))); 168 node_options.set_processor_options(Some(&Array::of3( 169 &wasm_bindgen::module(), 170 &wasm_bindgen::memory(), 171 &audio_queue.to_js_value(), 172 ))); 173 174 let worklet_node = 175 AudioWorkletNode::new_with_options(audio_ctx, "audio-processor", &node_options)?; 176 worklet_node.connect_with_audio_node(&audio_ctx.destination())?; 177 178 Ok(worklet_node) 179} 180 181#[wasm_bindgen(module = "/js/polyfill.js")] 182extern "C" { 183 fn run_polyfill(); 184}