diff options
| author | elvis <elvis@claros.ar> | 2026-09-06 19:21:16 -0300 |
|---|---|---|
| committer | elvis <elvis@claros.ar> | 2026-09-06 19:23:37 -0300 |
| commit | f8f98e83481e7376a235fb305a095552433f79ab (patch) | |
| tree | 6d0ecddb26bef8430a086431de40f8ce9b1039e6 /crates/asist-app/src | |
| download | asist-p-f8f98e83481e7376a235fb305a095552433f79ab.tar.gz asist-p-f8f98e83481e7376a235fb305a095552433f79ab.zip | |
Local voice assistant on top of Canary, llama.cpp and qwentts
Pipeline en Rust de hilos y canales que une los tres motores: el micrófono
alimenta un segmentador con VAD, las intervenciones cerradas van al
reconocedor, la transcripción al modelo y cada frase que este cierra sale
hacia el sintetizador sin esperar al resto de la respuesta.
Dos reglas sostienen el diseño: ninguna etapa bloquea a la anterior —quien va
sobrado descarta trabajo en lugar de acumular retraso— y todo lo que viaja por
los canales lleva el turno al que pertenece, así que interrumpir es subir el
contador y levantar dos banderas de cancelación.
Seis crates: core (configuración, eventos, HTTP, telemetría, herramientas),
audio (cpal, VAD, anillo de reproducción), asr, llm, tts y app (supervisor de
procesos y orquestador). Los motores van como submódulos fijados a un commit,
con los cambios locales en vendor/patches.
Midiendo el pipeline aparecieron tres cuellos de botella de configuración que
valieron más que cualquier cambio de código, todos documentados en
docs/RENDIMIENTO.md:
- tts-server decodificaba el audio en bloques de 24 s, de modo que el modo
«streaming» llegaba de una pieza: 4948 ms -> 585 ms hasta el primer audio.
- La plantilla de chat del modelo abre <think> y no lo cierra nunca, sin
variable que lo apague: 8630 ms -> 413 ms hasta el primer token, con una
copia de la plantilla que deja el bloque cerrado de entrada.
- Cualquier indicación de estilo junto a la guía de herramientas hace que este
modelo de 2B deje de llamarlas y se invente el dato (8/8 aciertos con la
guía sola, 0/8 con la persona de asistente de voz). El turno alterna ahora
entre dos instrucciones de sistema.
La ejecución de órdenes del sistema queda implementada y apagada, tras cuatro
barreras: lista blanca sobre el ejecutable, rutas rechazadas, sin shell que
interprete metacaracteres y plazo máximo.
68 pruebas unitarias sin modelos, más seis de integración que se saltan solas
si no hay servidores y se turnan la GPU: en paralelo, los dos servidores no
caben en 4 GB y miden contención en vez de latencia.
Claude-Session: https://claude.ai/code/session_01FNxz5cSdQSscJH9H7b8uGU
Diffstat (limited to 'crates/asist-app/src')
| -rw-r--r-- | crates/asist-app/src/main.rs | 474 | ||||
| -rw-r--r-- | crates/asist-app/src/pipeline.rs | 634 | ||||
| -rw-r--r-- | crates/asist-app/src/render.rs | 258 | ||||
| -rw-r--r-- | crates/asist-app/src/session.rs | 160 | ||||
| -rw-r--r-- | crates/asist-app/src/supervisor.rs | 228 |
5 files changed, 1754 insertions, 0 deletions
diff --git a/crates/asist-app/src/main.rs b/crates/asist-app/src/main.rs new file mode 100644 index 0000000..2788d97 --- /dev/null +++ b/crates/asist-app/src/main.rs @@ -0,0 +1,474 @@ +//! Asistente de voz local: escucha, piensa y responde hablando. +//! +//! Une tres motores que ya existen —Canary para oír, llama.cpp para pensar y +//! qwentts para hablar— en un pipeline de hilos y canales donde ninguna etapa +//! espera a la siguiente. + +mod pipeline; +mod render; +mod session; +mod supervisor; + +use std::path::PathBuf; +use std::time::Duration; + +use anyhow::{bail, Context, Result}; +use crossbeam_channel::unbounded; + +use asist_asr::Recognizer; +use asist_audio::{Capture, CaptureBlock, Playback}; +use asist_core::config::Config; +use asist_core::event::Event; +use asist_core::tools::ToolRegistry; +use asist_llm::LlmClient; +use asist_tts::TtsClient; + +use pipeline::Pipeline; +use render::Renderer; +use session::Session; +use supervisor::Supervisor; + +fn main() -> Result<()> { + let args = Args::parse()?; + if args.help { + println!("{USAGE}"); + return Ok(()); + } + init_logging(&args); + + let mut config = Config::load(&args.config) + .with_context(|| format!("no se pudo cargar {}", args.config.display()))?; + args.apply(&mut config); + + match args.command { + Command::Run => run(config, &args), + Command::Devices => devices(), + Command::Check => check(&config), + } +} + +fn run(config: Config, args: &Args) -> Result<()> { + let session = Session::new(); + + // Los clientes se crean antes de arrancar nada: así se detecta lo que ya + // esté escuchando y no se levanta un servidor por duplicado. + let llm = LlmClient::new(config.llm_authority(), &config.llm); + let tts = TtsClient::new(config.tts_authority(), &config.tts); + let (llm_up, tts_up) = (llm.healthy(), tts.healthy()); + + let mut supervisor = Supervisor::new(&args.log_dir)?; + supervisor.start(&config, llm_up, tts_up)?; + + let timeout = Duration::from_secs(config.supervisor.startup_timeout_secs); + eprintln!("Esperando a los motores (hasta {} s)…", timeout.as_secs()); + llm.wait_ready(timeout).map_err(|e| { + anyhow::anyhow!( + "{e}\nRevisa {}", + supervisor.log_path("llama-server").display() + ) + })?; + tts.wait_ready(timeout).map_err(|e| { + anyhow::anyhow!( + "{e}\nRevisa {}", + supervisor.log_path("tts-server").display() + ) + })?; + llm.warn_if_thinking_template(); + + if let Some(reference) = &config.tts.reference { + tts.register_voice(reference) + .context("no se pudo registrar la voz clonada")?; + } + // Se paga aquí la construcción de los grafos del sintetizador, unos 3,5 s + // que si no se los comería el primer turno de verdad. + if let Err(err) = tts.warmup() { + tracing::warn!(target: "tts", %err, "falló el precalentado; el primer turno irá lento"); + } + + eprintln!("Cargando el modelo de voz a texto…"); + let recognizer = Recognizer::load(&config.asr)?; + + let playback = Playback::open(&config.audio)?; + let (capture_tx, capture_rx) = unbounded::<CaptureBlock>(); + let capture = Capture::open(&config.audio, capture_tx.clone())?; + + let tools = ToolRegistry::from_config(&config.tools); + if tools.is_empty() { + eprintln!("Herramientas: ninguna"); + } else { + eprintln!("Herramientas: {}", tools.names().join(", ")); + } + if config.tools.shell { + eprintln!( + "Ejecución de órdenes ACTIVA — permitidas: {}{}", + config.tools.shell_allowlist.join(", "), + if config.tools.shell_dry_run { + " (simulación)" + } else { + "" + } + ); + } + + let pipeline = Pipeline::spawn( + &config, + session.clone(), + recognizer, + llm, + tts, + tools, + playback.handle(), + capture_rx, + capture_tx, + capture.format, + )?; + + eprintln!( + "\nListo. Micrófono: {} · Altavoz: {} ({} Hz)\nHabla cuando quieras. Enter para salir.\n", + capture.device_name, playback.device_name, playback.sample_rate + ); + + // Enter cierra, pero sólo si hay alguien delante para pulsarlo. Con la + // entrada redirigida —bajo systemd, en un contenedor, con `< /dev/null`— + // `read_line` devuelve EOF al instante y el asistente se cerraría nada + // más arrancar. En ese caso se espera a una señal. + if stdin_is_tty() { + let session = session.clone(); + std::thread::spawn(move || { + let mut line = String::new(); + let _ = std::io::stdin().read_line(&mut line); + session.request_stop(); + }); + } else { + eprintln!("(entrada no interactiva: para cerrar, manda SIGINT o SIGTERM)"); + install_signal_handler(session.clone()); + } + + let mut renderer = Renderer::new(config.asr.partials); + let metrics = std::sync::Arc::clone(&pipeline.metrics); + let mut next_health_check = std::time::Instant::now() + Duration::from_secs(5); + loop { + match pipeline.events.recv_timeout(Duration::from_millis(150)) { + Ok(event) => { + let shutdown = matches!(event, Event::Shutdown); + renderer.apply(&event); + if shutdown { + break; + } + } + Err(crossbeam_channel::RecvTimeoutError::Timeout) => {} + Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break, + } + if session.is_stopping() { + break; + } + // Si un servidor se muere a mitad de sesión, todos los turnos + // siguientes fallarían con un error de transporte sin explicar por + // qué. Mejor decirlo una vez, con el registro a mano, y cerrar. + if std::time::Instant::now() >= next_health_check { + next_health_check = std::time::Instant::now() + Duration::from_secs(5); + if let Some((name, code)) = supervisor.crashed() { + renderer.finish(); + eprintln!( + "\n{name} ha terminado inesperadamente (código {}). Revisa {}", + code.map(|c| c.to_string()) + .unwrap_or_else(|| "desconocido".into()), + supervisor.log_path(name).display() + ); + break; + } + } + } + + renderer.finish(); + capture.stop(); + pipeline.shutdown(); + playback.stop(); + supervisor.shutdown(); + + if metrics.turns() > 0 { + eprintln!("\n{}", metrics.summary()); + } + Ok(()) +} + +/// Lista los dispositivos de audio, para poder nombrarlos en la configuración. +fn devices() -> Result<()> { + use asist_audio::describe; + use cpal_reexport::traits::HostTrait; + let host = cpal_reexport::default_host(); + + let default_in = host.default_input_device().map(|d| describe(&d)); + println!("Entradas:"); + for device in host.input_devices()? { + let name = describe(&device); + let mark = if Some(&name) == default_in.as_ref() { + " (por defecto)" + } else { + "" + }; + println!(" {name}{mark}"); + } + + let default_out = host.default_output_device().map(|d| describe(&d)); + println!("\nSalidas:"); + for device in host.output_devices()? { + let name = describe(&device); + let mark = if Some(&name) == default_out.as_ref() { + " (por defecto)" + } else { + "" + }; + println!(" {name}{mark}"); + } + Ok(()) +} + +/// Comprueba que está todo en su sitio, sin abrir el micrófono. +fn check(config: &Config) -> Result<()> { + let mut problems = Vec::new(); + let ok = |label: &str, detail: String| println!(" ok {label:<22} {detail}"); + + let paths: [(&str, &std::path::Path); 6] = [ + ("modelo asr", &config.asr.model_dir), + ("binario llama", &config.supervisor.llama.binary), + ("modelo llm", &config.supervisor.llama.model), + ("plantilla chat", &config.supervisor.llama.chat_template), + ("binario tts", &config.supervisor.tts.binary), + ("modelo tts", &config.supervisor.tts.model), + ]; + println!("Ficheros:"); + for (label, path) in paths { + if path.exists() { + ok(label, path.display().to_string()); + } else { + println!(" FALTA {label:<22} {}", path.display()); + problems.push(format!("falta {label}: {}", path.display())); + } + } + + if let Some(reference) = &config.tts.reference { + println!("Voz de referencia «{}»:", reference.name); + for (label, path) in [ + ("hablante (.spk)", &reference.speaker), + ("códigos (.rvq)", &reference.codes), + ("transcripción", &reference.transcript), + ] { + if path.exists() { + ok(label, path.display().to_string()); + } else { + println!(" FALTA {label:<22} {}", path.display()); + problems.push(format!("falta {label}")); + } + } + } + + println!("Servidores:"); + let llm = LlmClient::new(config.llm_authority(), &config.llm); + let tts = TtsClient::new(config.tts_authority(), &config.tts); + for (label, up, authority) in [ + ("llama-server", llm.healthy(), config.llm_authority()), + ("tts-server", tts.healthy(), config.tts_authority()), + ] { + println!( + " {:<5} {label:<22} {authority}", + if up { "ok" } else { "-" } + ); + } + if llm.healthy() { + llm.warn_if_thinking_template(); + } + if tts.healthy() { + match tts.voices() { + Ok(voices) if voices.contains(&config.tts.voice) => { + println!(" ok voz «{}» registrada", config.tts.voice); + } + Ok(voices) => println!( + " - voz «{}» aún no registrada (hay: {})", + config.tts.voice, + voices.join(", ") + ), + Err(err) => println!(" - no se pudieron listar las voces: {err}"), + } + } + + if problems.is_empty() { + println!("\nTodo listo."); + Ok(()) + } else { + bail!( + "{} problema(s):\n - {}", + problems.len(), + problems.join("\n - ") + ) + } +} + +// cpal se usa aquí sólo para listar dispositivos; llega a través de asist-audio. +use asist_audio::cpal as cpal_reexport; + +const USAGE: &str = "\ +asistente — asistente de voz local + +USO: + asistente [ORDEN] [OPCIONES] + +ÓRDENES: + run Escucha y responde (por defecto) + check Comprueba ficheros, modelos y servidores, y sale + devices Lista los dispositivos de audio + +OPCIONES: + -c, --config <RUTA> Fichero de configuración [config/asistente.toml] + --log-dir <RUTA> Carpeta de los registros de los servidores [logs/] + --voice <NOMBRE> Voz del sintetizador + --no-manage No lanzar los servidores; suponerlos ya arriba + --no-partials Sin transcripción provisional (ahorra CPU) + --barge-in Permitir cortar al asistente hablando encima + --shell Activar la ejecución de órdenes del sistema + -v, --verbose Registro de depuración (equivale a RUST_LOG=debug) + -h, --help Esta ayuda + +VARIABLES: + RUST_LOG Filtro de registro, p. ej. «info,tts=debug,latencia=info» +"; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Command { + Run, + Check, + Devices, +} + +struct Args { + command: Command, + config: PathBuf, + log_dir: PathBuf, + voice: Option<String>, + no_manage: bool, + no_partials: bool, + barge_in: bool, + shell: bool, + verbose: bool, + help: bool, +} + +impl Args { + fn parse() -> Result<Self> { + let mut args = Self { + command: Command::Run, + config: PathBuf::from("config/asistente.toml"), + log_dir: PathBuf::from("logs"), + voice: None, + no_manage: false, + no_partials: false, + barge_in: false, + shell: false, + verbose: false, + help: false, + }; + let mut raw = std::env::args().skip(1); + while let Some(arg) = raw.next() { + match arg.as_str() { + "run" => args.command = Command::Run, + "check" => args.command = Command::Check, + "devices" => args.command = Command::Devices, + "-c" | "--config" => { + args.config = raw.next().context("--config necesita una ruta")?.into() + } + "--log-dir" => { + args.log_dir = raw.next().context("--log-dir necesita una ruta")?.into() + } + "--voice" => args.voice = Some(raw.next().context("--voice necesita un nombre")?), + "--no-manage" => args.no_manage = true, + "--no-partials" => args.no_partials = true, + "--barge-in" => args.barge_in = true, + "--shell" => args.shell = true, + "-v" | "--verbose" => args.verbose = true, + "-h" | "--help" => args.help = true, + other => bail!("opción desconocida: {other}\n\n{USAGE}"), + } + } + Ok(args) + } + + /// La línea de órdenes manda sobre el fichero. + fn apply(&self, config: &mut Config) { + if let Some(voice) = &self.voice { + config.tts.voice = voice.clone(); + } + if self.no_manage { + config.supervisor.manage = false; + } + if self.no_partials { + config.asr.partials = false; + } + if self.barge_in { + config.vad.barge_in = true; + } + if self.shell { + config.tools.shell = true; + } + } +} + +/// ¿Hay una persona al otro lado de la entrada estándar? +fn stdin_is_tty() -> bool { + #[cfg(unix)] + { + unsafe extern "C" { + fn isatty(fd: i32) -> i32; + } + unsafe { isatty(0) == 1 } + } + #[cfg(not(unix))] + true +} + +/// Cierre ordenado ante SIGINT/SIGTERM. +/// +/// Importa más de lo que parece: matar el proceso en seco deja a los +/// servidores hijos con la memoria de la GPU tomada. +#[cfg(unix)] +fn install_signal_handler(session: Session) { + use std::sync::OnceLock; + static SESSION: OnceLock<Session> = OnceLock::new(); + let _ = SESSION.set(session); + + extern "C" fn on_signal(_: i32) { + if let Some(session) = SESSION.get() { + session.request_stop(); + } + } + unsafe extern "C" { + fn signal(sig: i32, handler: extern "C" fn(i32)) -> usize; + } + unsafe { + signal(2, on_signal); // SIGINT + signal(15, on_signal); // SIGTERM + } +} + +#[cfg(not(unix))] +fn install_signal_handler(_session: Session) {} + +fn init_logging(args: &Args) { + use tracing_subscriber::{fmt, EnvFilter}; + // ONNX Runtime informa a nivel INFO de cada transformación del grafo: + // varios cientos de líneas que tapan la transcripción. Se silencia salvo + // que se pida expresamente por RUST_LOG. + let default = if args.verbose { + "debug,ort=warn" + } else { + "info,ort=warn" + }; + let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(default)); + fmt() + .with_env_filter(filter) + // Al terminal de texto van la transcripción y la respuesta; el registro + // va a stderr para que se puedan separar con una redirección. + .with_writer(std::io::stderr) + .with_target(true) + .without_time() + .init(); +} diff --git a/crates/asist-app/src/pipeline.rs b/crates/asist-app/src/pipeline.rs new file mode 100644 index 0000000..9e9037e --- /dev/null +++ b/crates/asist-app/src/pipeline.rs @@ -0,0 +1,634 @@ +//! El orquestador: los hilos del asistente y los canales que los unen. +//! +//! ```text +//! micrófono ──muestras──> segmentador ──intervención──> ASR ──texto──┐ +//! (cpal, tiempo real) (VAD, turnos) (Canary) │ +//! v +//! altavoz <──muestras── síntesis <──frases── conversación <──────────┘ +//! (cpal, anillo) (qwentts) (llama.cpp + herramientas) +//! ``` +//! +//! Cada caja es un hilo y cada flecha un canal. Dos consecuencias que valen +//! por todo el diseño: +//! +//! * **Nada bloquea al que va delante.** El micrófono nunca espera al ASR, y +//! el modelo nunca espera al sintetizador; quien va sobrado descarta trabajo +//! en lugar de acumular retraso. +//! * **Se responde por frases, no por respuestas.** La conversación entrega +//! cada frase al sintetizador en cuanto está cerrada, así que el asistente +//! empieza a hablar mientras el modelo sigue escribiendo. Es lo que separa +//! una latencia de medio segundo de una de cinco. + +use std::sync::Arc; +use std::thread::{self, JoinHandle}; +use std::time::{Duration, Instant}; + +use crossbeam_channel::{bounded, unbounded, Receiver, Select, Sender}; + +use asist_asr::{AsrJob, AsrResult, Recognizer}; +use asist_audio::{ + CaptureBlock, Gate, InputFormat, PlaybackHandle, Segmenter, VoiceEvent, ASR_SAMPLE_RATE, +}; +use asist_core::config::Config; +use asist_core::error::Result; +use asist_core::event::{Event, InterruptReason, TurnId}; +use asist_core::telemetry::{Metrics, Stage, TurnTimer}; +use asist_core::text::SentenceSplitter; +use asist_core::tools::ToolRegistry; +use asist_llm::chat::Conversation; +use asist_llm::{Delta, LlmClient, Message}; +use asist_tts::TtsClient; + +use crate::session::Session; + +/// Una frase pendiente de sintetizar. +#[derive(Debug)] +pub struct SpeakJob { + pub turn: TurnId, + pub index: usize, + pub text: String, + /// Origen desde el que se mide la latencia percibida del turno. + pub turn_started: Instant, +} + +/// Los hilos en marcha, para poder esperarlos al cerrar. +pub struct Pipeline { + handles: Vec<JoinHandle<()>>, + pub events: Receiver<Event>, + pub metrics: Arc<Metrics>, + session: Session, + capture_tx: Option<Sender<CaptureBlock>>, +} + +impl Pipeline { + /// Monta el pipeline entero. `capture_rx` viene del micrófono. + #[allow(clippy::too_many_arguments)] + pub fn spawn( + config: &Config, + session: Session, + recognizer: Recognizer, + llm: LlmClient, + tts: TtsClient, + tools: ToolRegistry, + playback: PlaybackHandle, + capture_rx: Receiver<CaptureBlock>, + capture_tx: Sender<CaptureBlock>, + input_format: InputFormat, + ) -> Result<Self> { + let (event_tx, event_rx) = unbounded::<Event>(); + let (asr_tx, asr_rx) = unbounded::<AsrJob>(); + let (result_tx, result_rx) = unbounded::<AsrResult>(); + // Acotado a propósito: si la síntesis se atasca, el hilo de + // conversación debe notarlo y frenar en lugar de acumular frases que + // llegarán tarde a un turno que quizá ya se ha interrumpido. + let (speak_tx, speak_rx) = bounded::<SpeakJob>(8); + + let metrics = Arc::new(Metrics::default()); + let mut handles = Vec::new(); + + handles.push(spawn_named("segmentador", { + let config = config.clone(); + let session = session.clone(); + let events = event_tx.clone(); + let playback = playback.clone(); + move || { + run_segmenter( + &config, + session, + capture_rx, + asr_tx, + events, + playback, + input_format, + ) + } + })); + + handles.push(spawn_named("asr", move || { + recognizer.run(asr_rx, result_tx) + })); + + handles.push(spawn_named("conversación", { + let config = config.clone(); + let session = session.clone(); + let events = event_tx.clone(); + let metrics = Arc::clone(&metrics); + move || { + run_brain( + &config, session, llm, tools, result_rx, speak_tx, events, metrics, + ) + } + })); + + handles.push(spawn_named("síntesis", { + let session = session.clone(); + let events = event_tx; + move || run_speech(session, tts, speak_rx, playback, events) + })); + + Ok(Self { + handles, + events: event_rx, + metrics, + session, + capture_tx: Some(capture_tx), + }) + } + + /// Cierra en orden: se corta lo que suena, se suelta el emisor del + /// micrófono y con eso la cadena de hilos se desmonta sola. + pub fn shutdown(mut self) { + self.session.request_stop(); + self.capture_tx.take(); + for handle in self.handles.drain(..) { + let _ = handle.join(); + } + } +} + +fn spawn_named(name: &str, body: impl FnOnce() + Send + 'static) -> JoinHandle<()> { + thread::Builder::new() + .name(name.to_string()) + .spawn(body) + .expect("no se pudo crear un hilo del pipeline") +} + +// --------------------------------------------------------------------------- +// Segmentador: micrófono -> intervenciones +// --------------------------------------------------------------------------- + +fn run_segmenter( + config: &Config, + session: Session, + capture: Receiver<CaptureBlock>, + asr: Sender<AsrJob>, + events: Sender<Event>, + playback: PlaybackHandle, + input_format: InputFormat, +) { + let mut segmenter = Segmenter::new(&config.vad); + let mut turn = session.current(); + + while let Ok(block) = capture.recv() { + if session.is_stopping() { + break; + } + // El micrófono se cierra mientras suena el altavoz. Con barge-in + // activo no se cierra, sólo sube el listón de volumen. + segmenter.set_gate(if session.is_speaking() || playback.is_active() { + Gate::Speaking + } else { + Gate::Open + }); + + // El dispositivo se abre a 16 kHz mono siempre que puede, y entonces + // esto no hace nada; cuando no puede, es aquí donde se convierte. + let mono = asist_audio::to_mono_at(&block.samples, input_format, ASR_SAMPLE_RATE); + + for event in segmenter.push(&mono) { + match event { + VoiceEvent::BargeIn => { + // Se corta antes de abrir el turno nuevo: así lo que + // quedaba en el anillo no se cuela por encima. + playback.stop(); + session.interrupt(); + let _ = events.send(Event::Interrupted { + turn: session.current(), + reason: InterruptReason::UserSpoke, + }); + } + VoiceEvent::Started => { + turn = session.begin_turn(); + let _ = events.send(Event::SpeechStarted { turn, at: block.at }); + } + VoiceEvent::Audio(samples) => { + if config.asr.partials + && asr + .send(AsrJob::Window { + turn, + samples, + at: block.at, + }) + .is_err() + { + return; + } + } + VoiceEvent::Ended(utterance) => { + let at = Instant::now(); + if asr + .send(AsrJob::Utterance { + turn, + samples: utterance.samples, + at, + }) + .is_err() + { + return; + } + } + VoiceEvent::Discarded => { + if asr.send(AsrJob::Reset { turn }).is_err() { + return; + } + } + } + } + } + // Lo que quedara a medias al cerrar todavía merece transcribirse. + if let Some(utterance) = segmenter.flush() { + let _ = asr.send(AsrJob::Utterance { + turn, + samples: utterance.samples, + at: Instant::now(), + }); + } +} + +// --------------------------------------------------------------------------- +// Conversación: transcripción -> frases +// --------------------------------------------------------------------------- + +#[allow(clippy::too_many_arguments)] +fn run_brain( + config: &Config, + session: Session, + llm: LlmClient, + tools: ToolRegistry, + results: Receiver<AsrResult>, + speak: Sender<SpeakJob>, + events: Sender<Event>, + metrics: Arc<Metrics>, +) { + let prompts = Prompts::new(config, &tools); + let mut chat = Conversation::new(prompts.deciding.clone(), config.general.history_turns); + + while let Ok(result) = results.recv() { + if session.is_stopping() { + break; + } + match result { + AsrResult::Partial { + turn, + committed, + volatile, + .. + } => { + let _ = events.send(Event::Partial { + turn, + committed, + volatile, + }); + } + AsrResult::Empty { turn } => { + let _ = events.send(Event::Discarded { turn }); + } + AsrResult::Error { turn, message } => { + let _ = events.send(Event::Warning { turn, message }); + } + AsrResult::Final { + turn, + text, + audio_secs, + decode, + spoken_at, + } => { + // Una transcripción de un turno superado llega tarde: el + // usuario ya ha dicho otra cosa. + if !session.is_current(turn) { + tracing::debug!(target: "brain", turno = turn.0, "transcripción obsoleta"); + continue; + } + let mut timer = TurnTimer::new(turn); + timer.input_secs = audio_secs; + timer.record(Stage::Asr, decode); + + let _ = events.send(Event::Transcript { + turn, + text: text.clone(), + audio_secs, + decode, + }); + + chat.push(Message::user(text)); + answer( + config, &session, &llm, &tools, &prompts, &mut chat, &speak, &events, + &mut timer, spoken_at, + ); + + if config.general.report_latency { + tracing::info!(target: "latencia", "\n{}", timer.report()); + } + metrics.record_turn(&timer); + } + } + } +} + +/// Las dos instrucciones de sistema entre las que alterna un turno. +/// +/// Que sean dos no es un capricho de diseño sino un resultado medido: con +/// Qwen3.5-2B, cualquier indicación de estilo junto a la guía de herramientas +/// hace que el modelo deje de llamarlas y se invente el dato (8/8 aciertos con +/// la guía sola, 0/8 con la persona de asistente de voz añadida). Así que la +/// pasada en que el modelo *decide* si actuar lleva la guía a solas, y la +/// instrucción de voz se reserva para redactar lo que se va a pronunciar. +struct Prompts { + /// Para la pasada en que el modelo decide si llamar a una herramienta. + deciding: String, + /// Para redactar la respuesta hablada, ya con los resultados en la mano. + speaking: String, + /// `false` cuando no hay herramientas: entonces ambas son la misma. + two_phase: bool, +} + +impl Prompts { + fn new(config: &Config, tools: &ToolRegistry) -> Self { + let speaking = config.general.system_prompt.trim().to_string(); + let two_phase = config.tools.enabled && config.tools.dedicated_prompt && !tools.is_empty(); + Self { + deciding: if two_phase { + config.general.tools_prompt.trim().to_string() + } else { + speaking.clone() + }, + speaking, + two_phase, + } + } +} + +/// Genera la respuesta de un turno, resolviendo herramientas si el modelo las +/// pide, y va soltando frases al sintetizador según se cierran. +#[allow(clippy::too_many_arguments)] +fn answer( + config: &Config, + session: &Session, + llm: &LlmClient, + tools: &ToolRegistry, + prompts: &Prompts, + chat: &mut Conversation, + speak: &Sender<SpeakJob>, + events: &Sender<Event>, + timer: &mut TurnTimer, + turn_started: Instant, +) { + let turn = timer.turn; + let tools = (!tools.is_empty() && config.tools.enabled).then_some(tools); + let mut splitter = SentenceSplitter::new(); + let mut spoken = String::new(); + + for round in 0..=config.llm.max_tool_rounds { + let llm_started = Instant::now(); + let mut first_token = None; + + let outcome = llm.stream(chat, tools, session.llm_cancel(), |delta| { + if !session.is_current(turn) { + return false; + } + match delta { + Delta::Text(text) => { + if first_token.is_none() { + first_token = Some(llm_started.elapsed()); + let _ = events.send(Event::ReplyStarted { + turn, + ttft: llm_started.elapsed(), + }); + } + let _ = events.send(Event::ReplyDelta { + turn, + text: text.clone(), + }); + spoken.push_str(text); + // Aquí está el solape: cada frase cerrada sale hacia el + // sintetizador sin esperar al resto de la respuesta. + for sentence in splitter.push(text) { + let index = splitter.emitted() - 1; + let _ = events.send(Event::Sentence { + turn, + index, + text: sentence.clone(), + }); + if speak + .send(SpeakJob { + turn, + index, + text: sentence, + turn_started, + }) + .is_err() + { + return false; + } + } + true + } + Delta::ToolCalls(_) => true, + } + }); + + if let Some(ttft) = first_token { + timer.record(Stage::LlmFirstToken, ttft); + timer.record(Stage::LlmRest, llm_started.elapsed().saturating_sub(ttft)); + } + + let outcome = match outcome { + Ok(outcome) => outcome, + Err(err) => { + let _ = events.send(Event::Failed { + turn, + message: err.to_string(), + }); + return; + } + }; + + if !session.is_current(turn) { + return; + } + + if outcome.tool_calls.is_empty() { + if let Some(rest) = splitter.flush() { + let index = splitter.emitted() - 1; + let _ = events.send(Event::Sentence { + turn, + index, + text: rest.clone(), + }); + let _ = speak.send(SpeakJob { + turn, + index, + text: rest, + turn_started, + }); |