diff --git a/src-tauri/src/commands/audio.rs b/src-tauri/src/commands/audio.rs index 1e3bc7e..62367da 100644 --- a/src-tauri/src/commands/audio.rs +++ b/src-tauri/src/commands/audio.rs @@ -1,12 +1,12 @@ use std::{ - collections::{BTreeMap, VecDeque}, + collections::BTreeMap, net::UdpSocket, sync::{ atomic::{AtomicBool, Ordering}, Arc, Mutex, }, thread, - time::{Duration, Instant}, + time::Duration, }; use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; @@ -14,55 +14,12 @@ use tauri::State; use crate::commands::config::ConfigState; -// ============================================================================ -// AUDIO FORMAT -// ============================================================================ - -const NETWORK_SAMPLE_RATE: u32 = 44_100; - -// 20ms packets. -// -// 44100 * 0.020 = 882 samples -const PACKET_DURATION_MS: usize = 20; -const NETWORK_PACKET_SAMPLES: usize = NETWORK_SAMPLE_RATE as usize * PACKET_DURATION_MS / 1000; - -// PCM is i16 => 2 bytes/sample. -const NETWORK_PACKET_BYTES: usize = NETWORK_PACKET_SAMPLES * 2; - -// Sequence number is a u32. -const AUDIO_HEADER_BYTES: usize = 4; - -// Maximum UDP packet we accept. -const MAX_UDP_PACKET_SIZE: usize = 4096; +const PACKET_SAMPLES: usize = 882; // 20ms @ 44.1kHz +const HEADER_SIZE: usize = 4; +const MAX_PACKET_SIZE: usize = 4096; // ============================================================================ -// JITTER BUFFER -// ============================================================================ -// -// We intentionally keep this relatively small. -// -// Increasing these values reduces underruns but increases latency. -// -// Current target: -// -// startup: 40ms -// normal target: 40ms -// maximum: 100ms -// -// This is appropriate for low-latency voice. -// - -const JITTER_TARGET_MS: usize = 40; -const JITTER_MAX_MS: usize = 100; - -// If a packet is missing, wait this long before considering it lost. -// -// Since packets are 20ms, 30ms gives us enough room for modest -// out-of-order delivery without making latency enormous. -const PACKET_LOSS_WAIT_MS: u64 = 30; - -// ============================================================================ -// VOICE STATE +// STATE // ============================================================================ pub struct VoiceState { @@ -71,13 +28,9 @@ pub struct VoiceState { pub struct VoiceSession { pub input_stream: cpal::Stream, - pub output_stream: Arc>, - pub socket: Arc, - pub pin: u64, - pub shutdown: Arc, } @@ -90,19 +43,17 @@ impl Default for VoiceState { } // ============================================================================ -// DEVICE LISTING +// DEVICES // ============================================================================ #[tauri::command] pub fn list_input_devices() -> Result, String> { let host = cpal::default_host(); - let devices = host + Ok(host .input_devices() - .map_err(|e| format!("Failed to enumerate input devices: {e}"))?; - - Ok(devices - .filter_map(|device| Some(device.description().ok()?.name().to_string())) + .map_err(|e| e.to_string())? + .filter_map(|d| d.description().ok().map(|x| x.name().to_string())) .collect()) } @@ -110,12 +61,10 @@ pub fn list_input_devices() -> Result, String> { pub fn list_output_devices() -> Result, String> { let host = cpal::default_host(); - let devices = host + Ok(host .output_devices() - .map_err(|e| format!("Failed to enumerate output devices: {e}"))?; - - Ok(devices - .filter_map(|device| Some(device.description().ok()?.name().to_string())) + .map_err(|e| e.to_string())? + .filter_map(|d| d.description().ok().map(|x| x.name().to_string())) .collect()) } @@ -147,35 +96,29 @@ pub fn connect_to_vc( let config = config_state.0.lock().unwrap().clone(); - // ======================================================================== + // ------------------------------------------------------------------------ // UDP - // ======================================================================== + // ------------------------------------------------------------------------ - let socket = - UdpSocket::bind("0.0.0.0:0").map_err(|e| format!("Failed to bind UDP socket: {e}"))?; + let socket = Arc::new( + UdpSocket::bind("0.0.0.0:0").map_err(|e| format!("Failed to bind UDP socket: {e}"))?, + ); socket .connect(&hostname) .map_err(|e| format!("Failed to connect UDP socket: {e}"))?; - // Required so disconnect can terminate the receiver thread. socket .set_read_timeout(Some(Duration::from_millis(100))) - .map_err(|e| format!("Failed to set UDP read timeout: {e}"))?; - - let socket = Arc::new(socket); - - // ======================================================================== - // AUTHENTICATION - // ======================================================================== + .map_err(|e| e.to_string())?; socket .send(&pin.to_be_bytes()) - .map_err(|e| format!("Failed to send pincode: {e}"))?; + .map_err(|e| format!("Failed to send pin: {e}"))?; - // ======================================================================== + // ------------------------------------------------------------------------ // DEVICES - // ======================================================================== + // ------------------------------------------------------------------------ let host = cpal::default_host(); @@ -183,110 +126,92 @@ pub fn connect_to_vc( Some(name) => host .input_devices() .map_err(|e| e.to_string())? - .find(|device| { - device - .description() + .find(|d| { + d.description() .ok() - .map(|description| description.name() == name) + .map(|x| x.name() == *name) .unwrap_or(false) }) - .ok_or_else(|| "Input device not found".to_string())?, + .ok_or("Input device not found")?, None => host .default_input_device() - .ok_or_else(|| "No default input device".to_string())?, + .ok_or("No default input device")?, }; let output_device = match &config.output_device_name { Some(name) => host .output_devices() .map_err(|e| e.to_string())? - .find(|device| { - device - .description() + .find(|d| { + d.description() .ok() - .map(|description| description.name() == name) + .map(|x| x.name() == *name) .unwrap_or(false) }) - .ok_or_else(|| "Output device not found".to_string())?, + .ok_or("Output device not found")?, None => host .default_output_device() - .ok_or_else(|| "No default output device".to_string())?, + .ok_or("No default output device")?, }; - // ======================================================================== - // SHUTDOWN - // ======================================================================== - let shutdown = Arc::new(AtomicBool::new(false)); - // ======================================================================== + // ------------------------------------------------------------------------ // PLAYBACK BUFFER - // ======================================================================== - // - // This contains samples already converted to the output device's - // channel layout. - // + // ------------------------------------------------------------------------ - let playback_buffer = Arc::new(Mutex::new(VecDeque::::new())); + let playback_buffer = Arc::new(Mutex::new(std::collections::VecDeque::::new())); - let playback_started = Arc::new(AtomicBool::new(false)); - - // ======================================================================== - // OUTPUT CONFIG - // ======================================================================== + // ------------------------------------------------------------------------ + // OUTPUT + // ------------------------------------------------------------------------ let output_config = build_output_config(&output_device)?; - - eprintln!( - "[vc] output: {} Hz / {} channels", - output_config.sample_rate, output_config.channels - ); - - let output_config_cell = Arc::new(Mutex::new(output_config.clone())); + let output_config = Arc::new(Mutex::new(output_config)); let needs_output_rebuild = Arc::new(AtomicBool::new(false)); - // ======================================================================== - // INPUT - // ======================================================================== + let output_stream = build_output_stream( + &output_device, + &output_config.lock().unwrap(), + playback_buffer.clone(), + needs_output_rebuild.clone(), + )?; - let input_socket = socket.clone(); + output_stream + .play() + .map_err(|e| format!("Failed to start output: {e}"))?; + + let output_stream = Arc::new(Mutex::new(output_stream)); + + // ------------------------------------------------------------------------ + // INPUT + // ------------------------------------------------------------------------ let mut input_config: cpal::StreamConfig = input_device .default_input_config() - .map_err(|e| format!("Failed to get input config: {e}"))? + .map_err(|e| e.to_string())? .into(); - // Network audio is mono. input_config.channels = 1; let input_sample_rate = input_config.sample_rate; - eprintln!( - "[vc] input: {} Hz -> {} Hz", - input_sample_rate, NETWORK_SAMPLE_RATE - ); + let input_socket = socket.clone(); let mut input_resampler = resampler::ResamplerFft::new( 1, - input_sample_rate.try_into().map_err(|e| { - format!( - "Invalid input sample rate: \ - {e:?}" - ) - })?, + input_sample_rate + .try_into() + .map_err(|e| format!("Invalid input sample rate: {e:?}"))?, resampler::SampleRate::Hz44100, ); let mut input_buffer = Vec::::new(); - - // Resampled samples waiting to form 20ms packets. - let mut packet_samples = Vec::::with_capacity(NETWORK_PACKET_SAMPLES * 2); - - // Sequence number for every audio packet. - let mut sequence: u32 = 0; + let mut packet_buffer = Vec::::new(); + let mut sequence = 0u32; let input_stream = input_device .build_input_stream( @@ -294,43 +219,24 @@ pub fn connect_to_vc( move |data: &[f32], _| { input_buffer.extend_from_slice(data); - let frame_size = input_resampler.chunk_size_input(); + let input_size = input_resampler.chunk_size_input(); + let output_size = input_resampler.chunk_size_output(); - while input_buffer.len() >= frame_size { - let input: Vec = input_buffer.drain(..frame_size).collect(); + while input_buffer.len() >= input_size { + let input: Vec = input_buffer.drain(..input_size).collect(); - let output_size = input_resampler.chunk_size_output(); - - let mut output = vec![0.0f32; output_size]; - - if let Err(e) = input_resampler.resample(&input, &mut output) { - eprintln!( - "[vc] input resample \ - failed: {e}" - ); + let mut output = vec![0.0; output_size]; + if input_resampler.resample(&input, &mut output).is_err() { continue; } - packet_samples.extend(output.into_iter().map(|sample| sample.clamp(-1.0, 1.0))); + packet_buffer.extend(output); - // ======================================================== - // CREATE EXACT 20ms PACKETS - // ======================================================== + while packet_buffer.len() >= PACKET_SAMPLES { + let samples: Vec = packet_buffer.drain(..PACKET_SAMPLES).collect(); - while packet_samples.len() >= NETWORK_PACKET_SAMPLES { - let samples: Vec = - packet_samples.drain(..NETWORK_PACKET_SAMPLES).collect(); - - // Header: - // - // [u32 sequence] - // - // Then: - // - // [i16 PCM...] - let mut packet = - Vec::with_capacity(AUDIO_HEADER_BYTES + NETWORK_PACKET_BYTES); + let mut packet = Vec::with_capacity(HEADER_SIZE + PACKET_SAMPLES * 2); packet.extend_from_slice(&sequence.to_be_bytes()); @@ -340,206 +246,103 @@ pub fn connect_to_vc( packet.extend_from_slice(&pcm.to_be_bytes()); } - if let Err(e) = input_socket.send(&packet) { - eprintln!( - "[vc] UDP send failed: \ - {e}" - ); - } + let _ = input_socket.send(&packet); sequence = sequence.wrapping_add(1); } } }, |err| { - eprintln!("[vc] input stream error: {err}"); + eprintln!("[vc] input error: {err}"); }, None, ) .map_err(|e| e.to_string())?; - // ======================================================================== - // OUTPUT STREAM - // ======================================================================== - - let initial_output_stream = build_output_stream( - &output_device, - &output_config, - playback_buffer.clone(), - needs_output_rebuild.clone(), - playback_started.clone(), - )?; - - initial_output_stream - .play() - .map_err(|e| format!("Failed to start output stream: {e}"))?; - - let output_stream = Arc::new(Mutex::new(initial_output_stream)); - - // ======================================================================== - // OUTPUT DEVICE WATCHER - // ======================================================================== + // ------------------------------------------------------------------------ + // OUTPUT DEVICE REBUILD WATCHER + // ------------------------------------------------------------------------ { let output_device = output_device.clone(); - - let needs_output_rebuild = needs_output_rebuild.clone(); - let output_stream = output_stream.clone(); - - let output_config_cell = output_config_cell.clone(); - + let output_config = output_config.clone(); let playback_buffer = playback_buffer.clone(); - - let playback_started = playback_started.clone(); - + let needs_rebuild = needs_output_rebuild.clone(); let shutdown = shutdown.clone(); thread::spawn(move || { while !shutdown.load(Ordering::SeqCst) { thread::sleep(Duration::from_millis(100)); - if shutdown.load(Ordering::SeqCst) { - return; - } - - if !needs_output_rebuild.swap(false, Ordering::SeqCst) { + if !needs_rebuild.swap(false, Ordering::SeqCst) { continue; } - eprintln!("[vc] rebuilding output stream..."); + let Ok(new_config) = build_output_config(&output_device) else { + needs_rebuild.store(true, Ordering::SeqCst); + continue; + }; - let result = build_output_config(&output_device).and_then(|new_config| { - build_output_stream( - &output_device, - &new_config, - playback_buffer.clone(), - needs_output_rebuild.clone(), - playback_started.clone(), - ) - .map(|stream| (new_config, stream)) - }); + let Ok(new_stream) = build_output_stream( + &output_device, + &new_config, + playback_buffer.clone(), + needs_rebuild.clone(), + ) else { + needs_rebuild.store(true, Ordering::SeqCst); + continue; + }; - match result { - Ok((new_config, new_stream)) => { - if let Err(e) = new_stream.play() { - eprintln!( - "[vc] failed to start \ - rebuilt output: {e}" - ); - - needs_output_rebuild.store(true, Ordering::SeqCst); - - continue; - } - - // The old output format is no longer valid. - // - // Throw away queued samples rather than playing - // them using the new device timing/layout. - playback_buffer.lock().unwrap().clear(); - - playback_started.store(false, Ordering::SeqCst); - - *output_config_cell.lock().unwrap() = new_config.clone(); - - *output_stream.lock().unwrap() = new_stream; - - eprintln!( - "[vc] output rebuilt: \ - {} Hz / {} channels", - new_config.sample_rate, new_config.channels - ); - } - - Err(e) => { - eprintln!( - "[vc] failed to rebuild \ - output: {e}" - ); - - needs_output_rebuild.store(true, Ordering::SeqCst); - } + if new_stream.play().is_err() { + needs_rebuild.store(true, Ordering::SeqCst); + continue; } + + playback_buffer.lock().unwrap().clear(); + + *output_config.lock().unwrap() = new_config; + *output_stream.lock().unwrap() = new_stream; } }); } - // ======================================================================== - // UDP RECEIVE + JITTER BUFFER - // ======================================================================== + // ------------------------------------------------------------------------ + // UDP RECEIVE + // ------------------------------------------------------------------------ { - let recv_socket = socket.clone(); - - let output_config_cell = output_config_cell.clone(); - + let socket = socket.clone(); let playback_buffer = playback_buffer.clone(); - - let playback_started = playback_started.clone(); - + let output_config = output_config.clone(); let shutdown = shutdown.clone(); thread::spawn(move || { - // ================================================================ - // PACKET REORDER BUFFER - // ================================================================ - - let mut packets: BTreeMap> = BTreeMap::new(); - - let mut expected_sequence: Option = None; - - // When we first notice a missing packet, remember when. - let mut missing_since: Option = None; - - // ================================================================ - // RESAMPLER - // ================================================================ - - let initial_config = output_config_cell.lock().unwrap().clone(); - - let mut last_sample_rate = initial_config.sample_rate; + let initial_config = output_config.lock().unwrap().clone(); let mut output_resampler = match create_output_resampler(initial_config.sample_rate) { - Ok(resampler) => resampler, - + Ok(r) => r, Err(e) => { - eprintln!( - "[vc] failed to create \ - output resampler: {e}" - ); - + eprintln!("[vc] output resampler: {e}"); return; } }; - let mut resample_input = Vec::::new(); + let mut last_sample_rate = initial_config.sample_rate; - // ================================================================ - // UDP BUFFER - // ================================================================ + let mut packets: BTreeMap> = BTreeMap::new(); + let mut expected: Option = None; - let mut udp_buffer = [0u8; MAX_UDP_PACKET_SIZE]; - - // ================================================================ - // MAIN LOOP - // ================================================================ + let mut resample_buffer = Vec::::new(); + let mut udp_buffer = [0u8; MAX_PACKET_SIZE]; while !shutdown.load(Ordering::SeqCst) { - // ============================================================ - // RECEIVE PACKET - // ============================================================ - - match recv_socket.recv(&mut udp_buffer) { + match socket.recv(&mut udp_buffer) { Ok(len) => { - if len <= AUDIO_HEADER_BYTES { + if len <= HEADER_SIZE { continue; } - // ---------------------------------------------------- - // Read sequence number. - // ---------------------------------------------------- - let sequence = u32::from_be_bytes([ udp_buffer[0], udp_buffer[1], @@ -547,342 +350,131 @@ pub fn connect_to_vc( udp_buffer[3], ]); - let pcm_bytes = &udp_buffer[AUDIO_HEADER_BYTES..len]; + let pcm = &udp_buffer[HEADER_SIZE..len]; - // Must contain complete i16 samples. - let pcm_len = pcm_bytes.len() & !1; + let mut samples = Vec::with_capacity(pcm.len() / 2); - if pcm_len == 0 { - continue; + for chunk in pcm.chunks_exact(2) { + let value = i16::from_be_bytes([chunk[0], chunk[1]]); + + samples.push(value as f32 / i16::MAX as f32); } - let mut samples = Vec::with_capacity(pcm_len / 2); - - for chunk in pcm_bytes[..pcm_len].chunks_exact(2) { - let pcm = i16::from_be_bytes([chunk[0], chunk[1]]); - - samples.push(pcm as f32 / i16::MAX as f32); - } - - // ---------------------------------------------------- - // Ignore packets that are already too old. - // ---------------------------------------------------- - - if let Some(expected) = expected_sequence { - if sequence_before(sequence, expected) { - continue; - } - } - - // ---------------------------------------------------- - // Insert into jitter buffer. - // - // BTreeMap automatically keeps sequence numbers - // ordered. - // ---------------------------------------------------- - packets.entry(sequence).or_insert(samples); - - // ---------------------------------------------------- - // Don't allow the packet jitter buffer itself to - // become a source of latency. - // ---------------------------------------------------- - - while packets.len() > 8 { - if let Some((&oldest, _)) = packets.iter().next() { - if let Some(expected) = expected_sequence { - if sequence_before(oldest, expected) { - packets.remove(&oldest); - } else { - break; - } - } else { - break; - } - } - } } - Err(e) - if e.kind() == std::io::ErrorKind::TimedOut - || e.kind() == std::io::ErrorKind::WouldBlock => - { - // Timeout is expected. It gives us a chance to - // process jitter-buffer timeouts and shutdown. + Err(e) if e.kind() == std::io::ErrorKind::TimedOut => { + continue; } Err(e) => { if !shutdown.load(Ordering::SeqCst) { - eprintln!( - "[vc] UDP receive failed: \ - {e}" - ); + eprintln!("[vc] UDP receive error: {e}"); } - - continue; } } - // ============================================================ - // INITIALIZE EXPECTED SEQUENCE - // ============================================================ + // ---------------------------------------------------------------- + // OUTPUT CONFIG + // ---------------------------------------------------------------- - if expected_sequence.is_none() { - if let Some((&first_sequence, _)) = packets.iter().next() { - expected_sequence = Some(first_sequence); + let config = output_config.lock().unwrap().clone(); - eprintln!( - "[vc] jitter buffer \ - synchronized at packet {}", - first_sequence - ); - } - } - - // ============================================================ - // CURRENT OUTPUT CONFIG - // ============================================================ - - let current_config = output_config_cell.lock().unwrap().clone(); - - // ============================================================ - // OUTPUT DEVICE SAMPLE RATE CHANGE - // ============================================================ - - if current_config.sample_rate != last_sample_rate { - eprintln!( - "[vc] output rate changed: \ - {} -> {}", - last_sample_rate, current_config.sample_rate - ); - - match create_output_resampler(current_config.sample_rate) { + if config.sample_rate != last_sample_rate { + match create_output_resampler(config.sample_rate) { Ok(new_resampler) => { output_resampler = new_resampler; - - last_sample_rate = current_config.sample_rate; + last_sample_rate = config.sample_rate; packets.clear(); - - resample_input.clear(); - + resample_buffer.clear(); playback_buffer.lock().unwrap().clear(); - - playback_started.store(false, Ordering::SeqCst); - - expected_sequence = None; - - missing_since = None; + expected = None; } - Err(e) => { - eprintln!( - "[vc] failed to create \ - output resampler: {e}" - ); - - packets.clear(); - resample_input.clear(); - - continue; - } + Err(_) => continue, } } - // ============================================================ - // MOVE READY PACKETS INTO RESAMPLER - // ============================================================ + // ---------------------------------------------------------------- + // INITIAL SEQUENCE + // ---------------------------------------------------------------- - loop { - let Some(expected) = expected_sequence else { - break; - }; - - // -------------------------------------------------------- - // Expected packet exists. - // -------------------------------------------------------- - - if let Some(samples) = packets.remove(&expected) { - resample_input.extend_from_slice(&samples); - - expected_sequence = Some(expected.wrapping_add(1)); - - missing_since = None; - - continue; - } - - // -------------------------------------------------------- - // Expected packet doesn't exist. - // - // If we don't have anything newer, there is nothing - // to do yet. - // -------------------------------------------------------- - - let has_newer_packet = packets.keys().any(|&seq| sequence_after(seq, expected)); - - if !has_newer_packet { - break; - } - - // -------------------------------------------------------- - // We have a later packet. - // - // Therefore the expected packet is either delayed or - // lost. - // - // Wait a short period before declaring it lost. - // -------------------------------------------------------- - - let now = Instant::now(); - - let since = missing_since.get_or_insert(now); - - if since.elapsed() < Duration::from_millis(PACKET_LOSS_WAIT_MS) { - break; - } - - // -------------------------------------------------------- - // Packet is considered lost. - // - // We insert a zero packet here. - // - // Because the missing packet is exactly 20ms, this - // results in a controlled 20ms gap rather than the - // playback clock getting permanently stuck. - // -------------------------------------------------------- - - resample_input.extend(std::iter::repeat(0.0f32).take(NETWORK_PACKET_SAMPLES)); - - expected_sequence = Some(expected.wrapping_add(1)); - - missing_since = None; - - eprintln!("[vc] lost UDP packet {}", expected); + if expected.is_none() { + expected = packets.keys().next().copied(); } - // ============================================================ + // ---------------------------------------------------------------- + // READ PACKETS IN ORDER + // ---------------------------------------------------------------- + + while let Some(seq) = expected { + if let Some(samples) = packets.remove(&seq) { + resample_buffer.extend(samples); + expected = Some(seq.wrapping_add(1)); + } else { + // Missing packet. + // + // Don't wait for it. + // Just insert 20ms of silence. + if packets.keys().any(|&x| x > seq) { + resample_buffer.extend(std::iter::repeat(0.0).take(PACKET_SAMPLES)); + + expected = Some(seq.wrapping_add(1)); + } + + break; + } + } + + // ---------------------------------------------------------------- // RESAMPLE - // ============================================================ + // ---------------------------------------------------------------- - let input_frame_size = output_resampler.chunk_size_input(); + let input_size = output_resampler.chunk_size_input(); + let output_size = output_resampler.chunk_size_output(); - while resample_input.len() >= input_frame_size { - let input: Vec = resample_input.drain(..input_frame_size).collect(); + while resample_buffer.len() >= input_size { + let input: Vec = resample_buffer.drain(..input_size).collect(); - let output_frame_size = output_resampler.chunk_size_output(); - - let mut resampled = vec![0.0f32; output_frame_size]; - - if let Err(e) = output_resampler.resample(&input, &mut resampled) { - eprintln!( - "[vc] output resample \ - failed: {e}" - ); + let mut output = vec![0.0; output_size]; + if output_resampler.resample(&input, &mut output).is_err() { continue; } - // ======================================================== - // MONO -> DEVICE CHANNELS - // ======================================================== - - let output = mono_to_output_channels(&resampled, current_config.channels); - - // ======================================================== - // PLAYBACK QUEUE - // ======================================================== - - let channels = current_config.channels.max(1) as usize; + let channels = config.channels.max(1) as usize; let mut queue = playback_buffer.lock().unwrap(); - queue.extend(output); + for sample in output { + for _ in 0..channels { + queue.push_back(sample); + } + } - // -------------------------------------------------------- - // Maximum playback queue. - // - // If this gets exceeded, THROW AWAY OLD AUDIO. - // - // This prevents latency from continuously growing. - // -------------------------------------------------------- - - let max_frames = current_config.sample_rate as usize * JITTER_MAX_MS / 1000; - - let max_samples = max_frames * channels; + // Don't let latency grow forever. + let max_samples = config.sample_rate as usize * channels / 5; while queue.len() > max_samples { queue.pop_front(); } - - // -------------------------------------------------------- - // Start playback once enough audio is available. - // -------------------------------------------------------- - - if !playback_started.load(Ordering::Acquire) { - let target_frames = - current_config.sample_rate as usize * JITTER_TARGET_MS / 1000; - - let target_samples = target_frames * channels; - - if queue.len() >= target_samples { - playback_started.store(true, Ordering::Release); - - eprintln!( - "[vc] playback started \ - with ~{}ms buffered", - JITTER_TARGET_MS - ); - } - } - } - - // ============================================================ - // RECOVER FROM PLAYBACK UNDERRUN - // ============================================================ - // - // If CPAL consumes everything, it will output silence. - // - // Once enough audio has accumulated again, playback can - // resume. - // - // We intentionally don't constantly toggle this state. - // That was one of the sources of the previous flicker. - // ============================================================ - - if playback_started.load(Ordering::Acquire) { - let channels = current_config.channels.max(1) as usize; - - let queue_len = playback_buffer.lock().unwrap().len(); - - let low_frames = current_config.sample_rate as usize * 10 / 1000; - - let low_samples = low_frames * channels; - - // We don't stop playback at 0ms. - // - // Only stop if the queue has actually become empty. - if queue_len == 0 { - playback_started.store(false, Ordering::Release); - } - - let _ = low_samples; } } }); } - // ======================================================================== + // ------------------------------------------------------------------------ // START INPUT - // ======================================================================== + // ------------------------------------------------------------------------ input_stream .play() - .map_err(|e| format!("Failed to start input stream: {e}"))?; + .map_err(|e| format!("Failed to start input: {e}"))?; - // ======================================================================== - // SAVE SESSION - // ======================================================================== + // ------------------------------------------------------------------------ + // SAVE + // ------------------------------------------------------------------------ *voice_state.session.lock().unwrap() = Some(VoiceSession { input_stream, @@ -899,26 +491,23 @@ pub fn connect_to_vc( // OUTPUT CONFIG // ============================================================================ -fn build_output_config(output_device: &cpal::Device) -> Result { - output_device +fn build_output_config(device: &cpal::Device) -> Result { + device .default_output_config() - .map_err(|e| e.to_string()) .map(Into::into) + .map_err(|e| e.to_string()) } // ============================================================================ -// OUTPUT RESAMPLER +// RESAMPLER // ============================================================================ fn create_output_resampler( - output_sample_rate: cpal::SampleRate, + sample_rate: cpal::SampleRate, ) -> Result { - let rate = output_sample_rate.try_into().map_err(|e| { - format!( - "Invalid output sample rate: \ - {e:?}" - ) - })?; + let rate = sample_rate + .try_into() + .map_err(|e| format!("Invalid sample rate: {e:?}"))?; Ok(resampler::ResamplerFft::new( 1, @@ -932,55 +521,26 @@ fn create_output_resampler( // ============================================================================ fn build_output_stream( - output_device: &cpal::Device, - output_config: &cpal::StreamConfig, - playback_buffer: Arc>>, + device: &cpal::Device, + config: &cpal::StreamConfig, + playback_buffer: Arc>>, needs_rebuild: Arc, - playback_started: Arc, ) -> Result { - output_device + device .build_output_stream( - output_config.clone(), - // ================================================================= - // AUDIO CALLBACK - // ================================================================= + config.clone(), move |data: &mut [f32], _| { - // ------------------------------------------------------------- - // Don't consume the queue until the jitter buffer has enough - // audio. - // ------------------------------------------------------------- - - if !playback_started.load(Ordering::Acquire) { - data.fill(0.0); - return; - } - let mut queue = playback_buffer.lock().unwrap(); - // ------------------------------------------------------------- - // CRITICAL: - // - // VecDeque::pop_front() is O(1). - // - // NEVER use Vec::remove(0) here. - // ------------------------------------------------------------- - - for output in data.iter_mut() { - *output = queue.pop_front().unwrap_or(0.0); + for sample in data { + *sample = queue.pop_front().unwrap_or(0.0); } - - // If the callback consumed the entire queue, the receiver - // thread will refill it. We don't modify playback_started - // here because the CPAL callback should stay extremely cheap. }, - // ================================================================= - // ERROR CALLBACK - // ================================================================= { let needs_rebuild = needs_rebuild.clone(); move |err| { - eprintln!("[vc] output stream error: {err}"); + eprintln!("[vc] output error: {err}"); let message = err.to_string(); @@ -996,60 +556,3 @@ fn build_output_stream( ) .map_err(|e| e.to_string()) } - -// ============================================================================ -// CHANNEL CONVERSION -// ============================================================================ - -fn mono_to_output_channels(mono: &[f32], channels: u16) -> Vec { - let channels = channels.max(1) as usize; - - if channels == 1 { - return mono.to_vec(); - } - - let mut output = Vec::with_capacity(mono.len() * channels); - - for &sample in mono { - for _ in 0..channels { - output.push(sample); - } - } - - output -} - -// ============================================================================ -// SEQUENCE NUMBER HELPERS -// ============================================================================ -// -// UDP sequence numbers eventually wrap around u32::MAX. -// -// These helpers make comparisons work correctly across the wrap. -// - -fn sequence_after(a: u32, b: u32) -> bool { - let diff = a.wrapping_sub(b); - - diff != 0 && diff < 0x8000_0000 -} - -fn sequence_before(a: u32, b: u32) -> bool { - let diff = a.wrapping_sub(b); - - diff != 0 && diff >= 0x8000_0000 -} - -// ============================================================================ -// OPTIONAL UTILITY -// ============================================================================ - -fn stereo_to_mono(stereo_data: &[f32]) -> Vec { - let mut mono = Vec::with_capacity(stereo_data.len() / 2); - - for chunk in stereo_data.chunks_exact(2) { - mono.push((chunk[0] + chunk[1]) * 0.5); - } - - mono -}