diff options
Diffstat (limited to 'crates/asist-audio/src')
| -rw-r--r-- | crates/asist-audio/src/capture.rs | 153 | ||||
| -rw-r--r-- | crates/asist-audio/src/lib.rs | 148 | ||||
| -rw-r--r-- | crates/asist-audio/src/playback.rs | 287 | ||||
| -rw-r--r-- | crates/asist-audio/src/vad.rs | 403 |
4 files changed, 991 insertions, 0 deletions
diff --git a/crates/asist-audio/src/capture.rs b/crates/asist-audio/src/capture.rs new file mode 100644 index 0000000..05da992 --- /dev/null +++ b/crates/asist-audio/src/capture.rs @@ -0,0 +1,153 @@ +//! Captura desde el micrófono. + +use std::time::Instant; + +use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; +use cpal::{Sample, SampleFormat, SupportedStreamConfig}; +use crossbeam_channel::Sender; + +use asist_core::config::AudioConfig; +use asist_core::error::{Error, Result}; + +use crate::{describe, ASR_SAMPLE_RATE}; + +/// Formato con el que se abrió realmente el dispositivo. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct InputFormat { + pub sample_rate: usize, + pub channels: usize, +} + +/// Un bloque tal y como sale de la retrollamada, con la marca de tiempo que +/// permite luego medir cuánto se ha retrasado el pipeline respecto de la voz. +#[derive(Debug)] +pub struct CaptureBlock { + pub samples: Vec<f32>, + pub at: Instant, +} + +pub struct Capture { + stream: cpal::Stream, + pub format: InputFormat, + pub device_name: String, +} + +impl Capture { + /// Abre la entrada y empieza a empujar bloques por `tx`. + /// + /// Prefiere 16 kHz mono porque es justo lo que quiere el modelo: si el + /// dispositivo lo acepta, no hay remuestreo en ningún punto del camino. + pub fn open(config: &AudioConfig, tx: Sender<CaptureBlock>) -> Result<Self> { + let host = cpal::default_host(); + let device = select_device(&host, &config.input_device)?; + let device_name = describe(&device); + + let supported = preferred_config(&device)?; + let format = InputFormat { + sample_rate: supported.sample_rate() as usize, + channels: supported.channels() as usize, + }; + let stream_config: cpal::StreamConfig = supported.clone().into(); + let on_error = |err| tracing::error!(target: "audio", %err, "flujo de entrada"); + + // Dentro de la retrollamada: convertir a f32, enviar y salir. `send` + // sobre un canal sin límite no bloquea, que es la única propiedad que + // aquí importa. + macro_rules! build { + ($sample:ty) => { + device + .build_input_stream( + &stream_config, + move |data: &[$sample], _: &_| { + let samples = data.iter().map(|s| f32::from_sample(*s)).collect(); + let _ = tx.send(CaptureBlock { + samples, + at: Instant::now(), + }); + }, + on_error, + None, + ) + .map_err(|e| Error::Audio(format!("no se pudo abrir la entrada: {e}")))? + }; + } + + let stream = match supported.sample_format() { + SampleFormat::F32 => build!(f32), + SampleFormat::I16 => build!(i16), + SampleFormat::U16 => build!(u16), + other => { + return Err(Error::Audio(format!( + "formato de muestra no soportado: {other:?}" + ))) + } + }; + stream + .play() + .map_err(|e| Error::Audio(format!("no se pudo arrancar la entrada: {e}")))?; + + tracing::info!( + target: "audio", + dispositivo = %device_name, + hz = format.sample_rate, + canales = format.channels, + remuestreo = format.sample_rate != ASR_SAMPLE_RATE as usize || format.channels != 1, + "entrada abierta" + ); + + Ok(Self { + stream, + format, + device_name, + }) + } + + /// Cierra el dispositivo. Al soltar el emisor, la cadena de hilos se + /// desmonta sola de arriba abajo. + pub fn stop(self) { + drop(self.stream); + } +} + +fn select_device(host: &cpal::Host, wanted: &str) -> Result<cpal::Device> { + if wanted.is_empty() { + return host + .default_input_device() + .ok_or_else(|| Error::Audio("no hay dispositivo de entrada".into())); + } + let wanted_lower = wanted.to_lowercase(); + let devices = host + .input_devices() + .map_err(|e| Error::Audio(format!("no se pudieron listar las entradas: {e}")))?; + let mut seen = Vec::new(); + for device in devices { + let name = describe(&device); + if name.to_lowercase().contains(&wanted_lower) { + return Ok(device); + } + seen.push(name); + } + Err(Error::Audio(format!( + "ninguna entrada coincide con «{wanted}». Disponibles: {}", + seen.join(", ") + ))) +} + +fn preferred_config(device: &cpal::Device) -> Result<SupportedStreamConfig> { + let native = device + .supported_input_configs() + .map_err(|e| Error::Audio(format!("no se pudo consultar la entrada: {e}")))? + .filter(|range| range.channels() == 1) + .filter(|range| { + range.min_sample_rate() <= ASR_SAMPLE_RATE && ASR_SAMPLE_RATE <= range.max_sample_rate() + }) + .find(|range| range.sample_format() == SampleFormat::F32) + .map(|range| range.with_sample_rate(ASR_SAMPLE_RATE)); + + match native { + Some(config) => Ok(config), + None => device + .default_input_config() + .map_err(|e| Error::Audio(format!("sin configuración de entrada: {e}"))), + } +} diff --git a/crates/asist-audio/src/lib.rs b/crates/asist-audio/src/lib.rs new file mode 100644 index 0000000..f6979ca --- /dev/null +++ b/crates/asist-audio/src/lib.rs @@ -0,0 +1,148 @@ +//! Entrada y salida de audio, y el detector de voz que las separa en turnos. +//! +//! La regla que gobierna este crate: **la retrollamada de audio no bloquea +//! nunca**. cpal la ejecuta en un hilo de tiempo real y cualquier espera ahí +//! se oye como un chasquido, así que se limita a copiar muestras a un canal +//! (entrada) o a vaciar un anillo ya rellenado (salida). Todo el trabajo real +//! —remuestreo, VAD, HTTP— ocurre en hilos normales al otro lado. + +pub mod capture; +pub mod playback; +pub mod vad; + +/// Se reexporta para que el binario pueda listar dispositivos sin volver a +/// declarar cpal y arriesgarse a resolver otra versión. +pub use cpal; + +pub use capture::{Capture, CaptureBlock, InputFormat}; +pub use playback::{Playback, PlaybackHandle}; +pub use vad::{Gate, Segmenter, Utterance, VoiceEvent}; + +/// Frecuencia a la que trabaja Canary. La captura se abre directamente aquí +/// cuando el dispositivo lo permite, lo que quita el remuestreo del camino. +pub const ASR_SAMPLE_RATE: u32 = 16_000; + +/// Frecuencia a la que sintetiza qwentts. +pub const TTS_SAMPLE_RATE: u32 = 24_000; + +/// Nombre legible de un dispositivo. +/// +/// `DeviceTrait::name` está obsoleto en cpal 0.17 a favor de `description`, +/// que devuelve una ficha entera; aquí sólo interesa el nombre, y tener un +/// único sitio donde extraerlo evita repetir el desempaquetado. +pub fn describe(device: &impl cpal::traits::DeviceTrait) -> String { + device + .description() + .map(|d| d.name().to_string()) + .unwrap_or_else(|_| "desconocido".into()) +} + +/// Nivel RMS de un bloque, la medida con la que el VAD decide. +pub fn rms(samples: &[f32]) -> f32 { + if samples.is_empty() { + return 0.0; + } + let sum: f32 = samples.iter().map(|v| v * v).sum(); + (sum / samples.len() as f32).sqrt() +} + +/// Mezcla a mono y remuestrea linealmente a `target`. +/// +/// La interpolación lineal basta: el dispositivo ya entrega la señal limitada +/// en banda, y un remuestreador decente costaría más que la decodificación a +/// la que alimenta. +pub fn to_mono_at(samples: &[f32], format: InputFormat, target: u32) -> Vec<f32> { + let mono: Vec<f32> = if format.channels > 1 { + samples + .chunks(format.channels) + .map(|frame| frame.iter().sum::<f32>() / format.channels as f32) + .collect() + } else { + samples.to_vec() + }; + + if format.sample_rate == target as usize { + return mono; + } + let ratio = target as f64 / format.sample_rate as f64; + let out_len = (mono.len() as f64 * ratio) as usize; + (0..out_len) + .map(|i| { + let pos = i as f64 / ratio; + let idx = pos as usize; + let frac = (pos - idx as f64) as f32; + let a = mono.get(idx).copied().unwrap_or(0.0); + let b = mono.get(idx + 1).copied().unwrap_or(a); + a + (b - a) * frac + }) + .collect() +} + +/// Convierte s16le a f32 en [-1, 1]. Es el formato en el que el TTS entrega. +pub fn s16le_to_f32(bytes: &[u8], out: &mut Vec<f32>) { + for pair in bytes.chunks_exact(2) { + let sample = i16::from_le_bytes([pair[0], pair[1]]); + out.push(sample as f32 / 32768.0); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn el_rms_de_una_senal_constante_es_su_amplitud() { + assert!((rms(&[0.5; 100]) - 0.5).abs() < 1e-6); + assert_eq!(rms(&[]), 0.0); + } + + #[test] + fn el_estereo_se_mezcla_a_mono_promediando() { + let format = InputFormat { + sample_rate: 16_000, + channels: 2, + }; + let out = to_mono_at(&[1.0, 0.0, 0.5, 0.5], format, 16_000); + assert_eq!(out, vec![0.5, 0.5]); + } + + #[test] + fn el_remuestreo_ajusta_la_duracion() { + let format = InputFormat { + sample_rate: 48_000, + channels: 1, + }; + let out = to_mono_at(&vec![0.0; 4800], format, 16_000); + assert_eq!(out.len(), 1600, "48 kHz -> 16 kHz debe dividir por tres"); + } + + #[test] + fn a_la_misma_frecuencia_el_remuestreo_no_toca_nada() { + let format = InputFormat { + sample_rate: 16_000, + channels: 1, + }; + let input = vec![0.1, -0.2, 0.3]; + assert_eq!(to_mono_at(&input, format, 16_000), input); + } + + #[test] + fn s16le_recorre_el_rango_completo() { + let mut out = Vec::new(); + s16le_to_f32(&[0x00, 0x00, 0xff, 0x7f, 0x00, 0x80], &mut out); + assert_eq!(out[0], 0.0); + assert!((out[1] - 1.0).abs() < 1e-4); + assert!((out[2] + 1.0).abs() < 1e-6); + } + + #[test] + fn un_byte_suelto_no_produce_una_muestra_a_medias() { + let mut out = Vec::new(); + s16le_to_f32(&[0x00, 0x00, 0x11], &mut out); + assert_eq!( + out.len(), + 1, + "el byte impar se ignora en vez de corromper la muestra" + ); + } +} diff --git a/crates/asist-audio/src/playback.rs b/crates/asist-audio/src/playback.rs new file mode 100644 index 0000000..fcdfe7d --- /dev/null +++ b/crates/asist-audio/src/playback.rs @@ -0,0 +1,287 @@ +//! Reproducción del audio sintetizado. +//! +//! El TTS produce a ráfagas —un bloque de codec cada vez— y el altavoz consume +//! a ritmo constante, así que entre ambos hay un anillo. La retrollamada sólo +//! vacía el anillo; quien sintetiza sólo lo llena. Cortar la respuesta es +//! entonces una operación trivial y sin condiciones de carrera: se vacía el +//! anillo y la voz calla en el siguiente bloque de audio. + +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, Condvar, Mutex}; +use std::time::Duration; + +use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; +use cpal::SampleFormat; + +use asist_core::config::AudioConfig; +use asist_core::error::{Error, Result}; + +use crate::{describe, to_mono_at, InputFormat, TTS_SAMPLE_RATE}; + +/// Estado compartido entre quien sintetiza y la retrollamada del altavoz. +struct Ring { + samples: Mutex<std::collections::VecDeque<f32>>, + /// Despierta a quien espera a que se vacíe la cola. + drained: Condvar, + /// Muestras servidas al altavoz desde que arrancó. Sirve para saber si ya + /// ha sonado algo de verdad, que es la latencia que percibe el usuario. + played: AtomicU64, + /// Hay audio en curso (queda cola o se está alimentando). + active: AtomicBool, + gain: Mutex<f32>, +} + +impl Ring { + fn new(gain: f32) -> Self { + Self { + samples: Mutex::new(std::collections::VecDeque::new()), + drained: Condvar::new(), + played: AtomicU64::new(0), + active: AtomicBool::new(false), + gain: Mutex::new(gain), + } + } + + fn lock(&self) -> std::sync::MutexGuard<'_, std::collections::VecDeque<f32>> { + self.samples.lock().unwrap_or_else(|e| e.into_inner()) + } +} + +/// Mando a distancia del altavoz: se puede clonar y repartir por los hilos. +#[derive(Clone)] +pub struct PlaybackHandle { + ring: Arc<Ring>, + sample_rate: u32, +} + +impl PlaybackHandle { + /// Encola muestras mono ya a la frecuencia del dispositivo. + pub fn push(&self, samples: &[f32]) { + if samples.is_empty() { + return; + } + let mut queue = self.ring.lock(); + queue.extend(samples.iter().copied()); + self.ring.active.store(true, Ordering::SeqCst); + } + + /// Encola audio del TTS (24 kHz mono), remuestreando si hace falta. + pub fn push_tts(&self, samples: &[f32]) { + if self.sample_rate == TTS_SAMPLE_RATE { + self.push(samples); + return; + } + let format = InputFormat { + sample_rate: TTS_SAMPLE_RATE as usize, + channels: 1, + }; + self.push(&to_mono_at(samples, format, self.sample_rate)); + } + + /// Calla ahora mismo y tira lo que quedaba por sonar. + pub fn stop(&self) { + let mut queue = self.ring.lock(); + queue.clear(); + self.ring.active.store(false, Ordering::SeqCst); + self.ring.drained.notify_all(); + } + + /// Segundos de audio pendientes de sonar. + pub fn queued_secs(&self) -> f32 { + self.ring.lock().len() as f32 / self.sample_rate as f32 + } + + /// `true` si ya ha salido audio por el altavoz. + pub fn has_played(&self) -> bool { + self.ring.played.load(Ordering::SeqCst) > 0 + } + + pub fn reset_played(&self) { + self.ring.played.store(0, Ordering::SeqCst); + } + + pub fn is_active(&self) -> bool { + self.ring.active.load(Ordering::SeqCst) + } + + pub fn set_gain(&self, gain: f32) { + *self.ring.gain.lock().unwrap_or_else(|e| e.into_inner()) = gain; + } + + /// Espera a que se vacíe la cola, o a que venza el plazo. + /// + /// Devuelve `true` si terminó de sonar todo y `false` si se agotó el + /// tiempo, para que quien llama distinga «ya está» de «sigue sonando». + pub fn wait_drained(&self, timeout: Duration) -> bool { + let deadline = std::time::Instant::now() + timeout; + let mut queue = self.ring.lock(); + while !queue.is_empty() { + let left = deadline.saturating_duration_since(std::time::Instant::now()); + if left.is_zero() { + return false; + } + let (guard, result) = self + .ring + .drained + .wait_timeout(queue, left.min(Duration::from_millis(50))) + .unwrap_or_else(|e| e.into_inner()); + queue = guard; + if result.timed_out() && queue.is_empty() { + break; + } + } + self.ring.active.store(false, Ordering::SeqCst); + true + } +} + +pub struct Playback { + stream: cpal::Stream, + handle: PlaybackHandle, + pub device_name: String, + pub sample_rate: u32, +} + +impl Playback { + pub fn open(config: &AudioConfig) -> Result<Self> { + let host = cpal::default_host(); + let device = select_device(&host, &config.output_device)?; + let device_name = describe(&device); + + let supported = preferred_config(&device)?; + let sample_rate = supported.sample_rate(); + let channels = supported.channels() as usize; + let stream_config: cpal::StreamConfig = supported.clone().into(); + + let ring = Arc::new(Ring::new(config.output_gain)); + let on_error = |err| tracing::error!(target: "audio", %err, "flujo de salida"); + + // La retrollamada: coger lo que haya, rellenar con silencio lo que + // falte y salir. Nunca espera a que llegue audio, porque quedarse + // esperando aquí se oye como un corte. + macro_rules! build { + ($sample:ty, $silence:expr, $convert:expr) => {{ + let ring = Arc::clone(&ring); + device + .build_output_stream( + &stream_config, + move |data: &mut [$sample], _: &_| { + let gain = *ring.gain.lock().unwrap_or_else(|e| e.into_inner()); + let mut queue = ring.lock(); + let mut served = 0u64; + for frame in data.chunks_mut(channels) { + match queue.pop_front() { + Some(sample) => { + let value = (sample * gain).clamp(-1.0, 1.0); + // Mono a N canales: la misma muestra en todos. + for out in frame.iter_mut() { + *out = $convert(value); + } + served += 1; + } + None => { + for out in frame.iter_mut() { + *out = $silence; + } + } + } + } + if served > 0 { + ring.played.fetch_add(served, Ordering::SeqCst); + } + if queue.is_empty() { + ring.drained.notify_all(); + } + }, + on_error, + None, + ) + .map_err(|e| Error::Audio(format!("no se pudo abrir la salida: {e}")))? + }}; + } + + let stream = match supported.sample_format() { + SampleFormat::F32 => build!(f32, 0.0f32, |v: f32| v), + SampleFormat::I16 => build!(i16, 0i16, |v: f32| (v * 32767.0) as i16), + SampleFormat::U16 => { + build!(u16, u16::MAX / 2, |v: f32| ((v * 32767.0) as i32 + 32768) + as u16) + } + other => { + return Err(Error::Audio(format!( + "formato de salida no soportado: {other:?}" + ))) + } + }; + stream + .play() + .map_err(|e| Error::Audio(format!("no se pudo arrancar la salida: {e}")))?; + + tracing::info!( + target: "audio", + dispositivo = %device_name, + hz = sample_rate, + canales = channels, + "salida abierta" + ); + + Ok(Self { + stream, + handle: PlaybackHandle { ring, sample_rate }, + device_name, + sample_rate, + }) + } + + pub fn handle(&self) -> PlaybackHandle { + self.handle.clone() + } + + pub fn stop(self) { + self.handle.stop(); + drop(self.stream); + } +} + +fn select_device(host: &cpal::Host, wanted: &str) -> Result<cpal::Device> { + if wanted.is_empty() { + return host + .default_output_device() + .ok_or_else(|| Error::Audio("no hay dispositivo de salida".into())); + } + let wanted_lower = wanted.to_lowercase(); + let devices = host + .output_devices() + .map_err(|e| Error::Audio(format!("no se pudieron listar las salidas: {e}")))?; + let mut seen = Vec::new(); + for device in devices { + let name = describe(&device); + if name.to_lowercase().contains(&wanted_lower) { + return Ok(device); + } + seen.push(name); + } + Err(Error::Audio(format!( + "ninguna salida coincide con «{wanted}». Disponibles: {}", + seen.join(", ") + ))) +} + +fn preferred_config(device: &cpal::Device) -> Result<cpal::SupportedStreamConfig> { + // 24 kHz nativo evita remuestrear la síntesis; si no, se coge lo de fábrica. + let native = device + .supported_output_configs() + .map_err(|e| Error::Audio(format!("no se pudo consultar la salida: {e}")))? + .filter(|range| { + range.min_sample_rate() <= TTS_SAMPLE_RATE && TTS_SAMPLE_RATE <= range.max_sample_rate() + }) + .find(|range| range.sample_format() == SampleFormat::F32) + .map(|range| range.with_sample_rate(TTS_SAMPLE_RATE)); + + match native { + Some(config) => Ok(config), + None => device + .default_output_config() + .map_err(|e| Error::Audio(format!("sin configuración de salida: {e}"))), + } +} diff --git a/crates/asist-audio/src/vad.rs b/crates/asist-audio/src/vad.rs new file mode 100644 index 0000000..b30f833 --- /dev/null +++ b/crates/asist-audio/src/vad.rs @@ -0,0 +1,403 @@ +//! Detección de voz y troceado en intervenciones. +//! +//! Está escrito como una máquina de estados pura: se le dan muestras y +//! devuelve eventos, sin hilos ni canales dentro. Así el comportamiento que +//! más cuesta depurar a oído —cuándo arranca un turno, cuándo lo corta el +//! silencio, cuándo se ignora el eco del propio altavoz— se puede probar +//! entero con audio sintético y sin micrófono. + +use std::collections::VecDeque; + +use asist_core::config::VadConfig; + +use crate::{rms, ASR_SAMPLE_RATE}; + +/// Una intervención cerrada, lista para transcribir. +#[derive(Debug, Clone)] +pub struct Utterance { + pub samples: Vec<f32>, + pub sample_rate: u32, +} + +impl Utterance { + pub fn duration_secs(&self) -> f32 { + self.samples.len() as f32 / self.sample_rate as f32 + } +} + +/// Lo que el segmentador tiene que contar hacia fuera. +#[derive(Debug, Clone)] +pub enum VoiceEvent { + /// Ha empezado a hablarse. + Started, + /// Audio nuevo dentro de la intervención en curso, para las + /// transcripciones provisionales. + Audio(Vec<f32>), + /// Intervención terminada y lo bastante larga como para transcribirla. + Ended(Utterance), + /// Terminada pero demasiado corta: un golpe en la mesa, una tos. + Discarded, + /// Se ha detectado voz mientras el asistente hablaba, con el barge-in + /// activo. El orquestador corta la reproducción al recibirlo. + BargeIn, +} + +/// Qué hace el segmentador con el micrófono mientras suena el altavoz. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Gate { + /// Nadie está hablando por el altavoz: se escucha con normalidad. + Open, + /// El asistente habla. Según la configuración, o se ignora la entrada + /// (media dúplex) o se exige más volumen para interrumpir (barge-in). + Speaking, +} + +pub struct Segmenter { + config: VadConfig, + frame_samples: usize, + preroll_samples: usize, + silence_hold_samples: usize, + min_utterance_samples: usize, + max_utterance_samples: usize, + + pending: Vec<f32>, + preroll: VecDeque<f32>, + utterance: Vec<f32>, + /// Muestras de la intervención que estaban de verdad por encima del + /// umbral. El mínimo se mide sobre esto y no sobre `utterance`, que + /// arrastra el preroll: si no, 0,1 s de golpe en la mesa más 0,2 s de + /// preroll pasan por una intervención válida y disparan un turno entero. + voiced: usize, + noise_floor: f32, + silence_run: usize, + speaking: bool, + gate: Gate, +} + +impl Segmenter { + pub fn new(config: &VadConfig) -> Self { + let rate = ASR_SAMPLE_RATE as f32; + Self { + frame_samples: (rate * config.frame_seconds).max(1.0) as usize, + preroll_samples: (rate * config.preroll_seconds) as usize, + silence_hold_samples: (rate * config.silence_hold) as usize, + min_utterance_samples: (rate * config.min_utterance) as usize, + max_utterance_samples: (rate * config.max_utterance) as usize, + config: config.clone(), + pending: Vec::new(), + preroll: VecDeque::new(), + utterance: Vec::new(), + voiced: 0, + noise_floor: 0.0, + silence_run: 0, + speaking: false, + gate: Gate::Open, + } + } + + /// Abre o cierra el micrófono según hable o no el asistente. + pub fn set_gate(&mut self, gate: Gate) { + if self.gate == gate { + return; + } + self.gate = gate; + // Al volver a abrir tras una respuesta, lo acumulado es la cola del + // propio altavoz: arrancar un turno con eso daría un turno fantasma. + if gate == Gate::Open { + self.pending.clear(); + self.preroll.clear(); + self.utterance.clear(); + self.voiced = 0; + self.silence_run = 0; + self.speaking = false; + } + } + + pub fn is_speaking(&self) -> bool { + self.speaking + } + + pub fn noise_floor(&self) -> f32 { + self.noise_floor + } + + /// Umbral que separa voz de silencio ahora mismo. + pub fn threshold(&self) -> f32 { + let base = (self.noise_floor * self.config.threshold_factor) + .clamp(self.config.min_threshold, self.config.max_threshold); + // Con el altavoz sonando hay que hablar más alto para colarse: el + // micrófono se está oyendo a sí mismo. + if self.gate == Gate::Speaking { + base * self.config.barge_in_factor + } else { + base + } + } + + /// Cierra a la fuerza la intervención en curso (cierre del programa). + pub fn flush(&mut self) -> Option<Utterance> { + if !self.speaking || self.voiced < self.min_utterance_samples { + return None; + } + self.speaking = false; + self.voiced = 0; + Some(Utterance { + samples: std::mem::take(&mut self.utterance), + sample_rate: ASR_SAMPLE_RATE, + }) + } + + /// Alimenta audio mono a 16 kHz y recoge lo que haya que hacer. + pub fn push(&mut self, samples: &[f32]) -> Vec<VoiceEvent> { + let mut events = Vec::new(); + // En media dúplex el micrófono está apagado de hecho: sin esto, el + // asistente se transcribe a sí mismo y se responde solo. + if self.gate == Gate::Speaking && !self.config.barge_in { + return events; + } + self.pending.extend_from_slice(samples); + + while self.pending.len() >= self.frame_samples { + let frame: Vec<f32> = self.pending.drain(..self.frame_samples).collect(); + let level = rms(&frame); + self.track_noise_floor(level); + let threshold = self.threshold(); + + if level < threshold && !self.speaking { + self.preroll.extend(frame.iter().copied()); + while self.preroll.len() > self.preroll_samples { + self.preroll.pop_front(); + } + continue; + } + + if !self.speaking { + self.speaking = true; + self.utterance.clear(); + self.voiced = 0; + self.utterance.extend(self.preroll.drain(..)); + if self.gate == Gate::Speaking { + events.push(VoiceEvent::BargeIn); + } + events.push(VoiceEvent::Started); + } + + if level < threshold { + self.silence_run += frame.len(); + } else { + self.silence_run = 0; + self.voiced += frame.len(); + } + self.utterance.extend_from_slice(&frame); + + let ended = self.silence_run >= self.silence_hold_samples; + let too_long = self.utterance.len() >= self.max_utterance_samples; + if !ended && !too_long { + events.push(VoiceEvent::Audio(frame)); + continue; + } + + let long_enough = self.voiced >= self.min_utterance_samples; + events.push(if long_enough { + VoiceEvent::Ended(Utterance { + samples: std::mem::take(&mut self.utterance), + sample_rate: ASR_SAMPLE_RATE, + }) + } else { + VoiceEvent::Discarded + }); + + self.utterance.clear(); + self.voiced = 0; + self.preroll.clear(); + self.silence_run = 0; + // Un corte por longitud cae a mitad de frase: se sigue escuchando + // como si el usuario no hubiera dejado de hablar, que es la verdad. + self.speaking = too_long && !ended; + if self.speaking { + events.push(VoiceEvent::Started); + } + } + events + } + + /// El suelo de ruido sólo baja: una voz sostenida no debe poder arrastrar + /// el umbral por encima de sí misma y dejar de detectarse. + fn track_noise_floor(&mut self, level: f32) { + if self.noise_floor == 0.0 { + self.noise_floor = level; + } else if level < self.noise_floor * 1.5 { + self.noise_floor = self.noise_floor * 0.95 + level * 0.05; + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn config() -> VadConfig { + VadConfig { + frame_seconds: 0.1, + preroll_seconds: 0.2, + silence_hold: 0.3, + min_utterance: 0.3, + max_utterance: 2.0, + threshold_factor: 3.0, + min_threshold: 0.001, + max_threshold: 0.02, + barge_in: false, + barge_in_factor: 4.0, + } + } + + fn samples(secs: f32, amplitude: f32) -> Vec<f32> { + let n = (ASR_SAMPLE_RATE as f32 * secs) as usize; + // Alterna de signo para que el RMS sea la amplitud y no un continuo. + (0..n) + .map(|i| if i % 2 == 0 { amplitude } else { -amplitude }) + .collect() + } + + fn feed(seg: &mut Segmenter, secs: f32, amplitude: f32) -> Vec<VoiceEvent> { + seg.push(&samples(secs, amplitude)) + } + + #[test] + fn el_silencio_no_arranca_ningun_turno() { + let mut seg = Segmenter::new(&config()); + let events = feed(&mut seg, 2.0, 0.0001); + assert!( + events.is_empty(), + "el silencio no debe producir eventos: {events:?}" + ); + } + + #[test] + fn la_voz_seguida_de_silencio_cierra_una_intervencion() { + let mut seg = Segmenter::new(&config()); + feed(&mut seg, 1.0, 0.0002); // deja que el suelo de ruido se asiente + let mut events = feed(&mut seg, 0.8, 0.3); + events.extend(feed(&mut seg, 0.6, 0.0002)); + + assert!(matches!(events.first(), Some(VoiceEvent::Started))); + let ended = events.iter().find_map(|e| match e { + VoiceEvent::Ended(u) => Some(u), + _ => None, + }); + let utterance = ended.expect("la intervención debió cerrarse"); + assert!( + utterance.duration_secs() > 0.8, + "el preroll debe ir incluido, duró {}", + utterance.duration_secs() + ); + } + + #[test] + fn un_ruido_corto_se_descarta() { + let mut seg = Segmenter::new(&config()); + feed(&mut seg, 1.0, 0.0002); + let mut events = feed(&mut seg, 0.1, 0.3); + events.extend(feed(&mut seg, 0.6, 0.0002)); + // 0,1 s de golpe no llegan al mínimo de 0,3 s de voz, por mucho que el + // preroll haga que la intervención dure más. + assert!( + events.iter().any(|e| matches!(e, VoiceEvent::Discarded)), + "esperaba un descarte: {events:?}" + ); + } + + #[test] + fn una_intervencion_interminable_se_corta_y_se_sigue_escuchando() { + let mut seg = Segmenter::new(&config()); + feed(&mut seg, 1.0, 0.0002); + let events = feed(&mut seg, 3.0, 0.3); + assert!( + events.iter().any(|e| matches!(e, VoiceEvent::Ended(_))), + "a los 2 s debe cortarse: {events:?}" + ); + assert!( + seg.is_speaking(), + "tras un corte forzado se sigue en mitad de la frase" + ); + } + + #[test] + fn en_media_duplex_el_altavoz_no_se_transcribe_a_si_mismo() { + let mut seg = Segmenter::new(&config()); + feed(&mut seg, 1.0, 0.0002); + seg.set_gate(Gate::Speaking); + let events = feed(&mut seg, 2.0, 0.5); + assert!( + events.is_empty(), + "con barge_in apagado no debe entrar nada mientras habla el asistente: {events:?}" + ); + } + + #[test] + fn con_barge_in_hace_falta_hablar_mas_alto_para_cortar() { + let mut config = config(); + config.barge_in = true; + config.barge_in_factor = 4.0; + let mut seg = Segmenter::new(&config); + feed(&mut seg, 1.0, 0.0002); + seg.set_gate(Gate::Speaking); + + // Justo por encima del umbral normal pero por debajo del elevado. + let eco = seg.threshold() / config.barge_in_factor * 1.5; + let events = feed(&mut seg, 0.5, eco); + assert!( + events.is_empty(), + "el eco del altavoz no debe cortar: {events:?}" + ); + + let voz = seg.threshold() * 2.0; + let events = feed(&mut seg, 0.5, voz); + assert!( + events.iter().any(|e| matches!(e, VoiceEvent::BargeIn)), + "una voz clara sí debe cortar: {events:?}" + ); + } + + #[test] + fn al_reabrir_el_microfono_se_tira_la_cola_del_altavoz() { + let mut config = config(); + config.barge_in = true; + let mut seg = Segmenter::new(&config); + feed(&mut seg, 1.0, 0.0002); + seg.set_gate(Gate::Speaking); + feed(&mut seg, 0.5, 0.9); + assert!(seg.is_speaking()); + + seg.set_gate(Gate::Open); + assert!( + !seg.is_speaking(), + "reabrir debe descartar lo acumulado; si no, el turno siguiente arranca con eco" + ); + } + + #[test] + fn el_suelo_de_ruido_no_sube_con_la_voz() { + let mut seg = Segmenter::new(&config()); + feed(&mut seg, 1.0, 0.0002); + let quieto = seg.noise_floor(); + feed(&mut seg, 2.0, 0.5); + assert!( + seg.noise_floor() <= quieto * 1.5, + "hablar no debe elevar el suelo de ruido ({quieto} -> {})", + seg.noise_floor() + ); + } + + #[test] + fn el_cierre_entrega_la_intervencion_a_medias() { + let mut seg = Segmenter::new(&config()); |