aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-audio/src
diff options
context:
space:
mode:
Diffstat (limited to 'crates/asist-audio/src')
-rw-r--r--crates/asist-audio/src/capture.rs153
-rw-r--r--crates/asist-audio/src/lib.rs148
-rw-r--r--crates/asist-audio/src/playback.rs287
-rw-r--r--crates/asist-audio/src/vad.rs403
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());