aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-app/src
diff options
context:
space:
mode:
Diffstat (limited to 'crates/asist-app/src')
-rw-r--r--crates/asist-app/src/main.rs474
-rw-r--r--crates/asist-app/src/pipeline.rs634
-rw-r--r--crates/asist-app/src/render.rs258
-rw-r--r--crates/asist-app/src/session.rs160
-rw-r--r--crates/asist-app/src/supervisor.rs228
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,
+ });
+ }
+ timer.sentences = splitter.emitted();
+ chat.push(Message::assistant(outcome.text.clone()));
+ // El turno siguiente vuelve a empezar decidiendo.
+ if prompts.two_phase {
+ chat.set_system(&prompts.deciding);
+ }
+ let _ = events.send(Event::ReplyDone {
+ turn,
+ text: outcome.text,
+ });
+ return;
+ }
+
+ // El modelo ha pedido herramientas: se ejecutan, se le devuelve el
+ // resultado y se le vuelve a preguntar.
+ if round == config.llm.max_tool_rounds {
+ let _ = events.send(Event::Warning {
+ turn,
+ message: format!(
+ "el modelo siguió pidiendo herramientas tras {} vueltas; se corta",
+ config.llm.max_tool_rounds
+ ),
+ });
+ return;
+ }
+
+ let tools = tools.expect("no puede haber llamadas sin herramientas registradas");
+ chat.push(Message::tool_request(
+ outcome.text.clone(),
+ outcome.tool_calls.clone(),
+ ));
+ let tools_started = Instant::now();
+ for call in &outcome.tool_calls {
+ let _ = events.send(Event::ToolRequested {
+ turn,
+ name: call.name.clone(),
+ arguments: call.arguments.clone(),
+ });
+ let result = tools.dispatch(call);
+ let _ = events.send(Event::ToolFinished {
+ turn,
+ name: result.name.clone(),
+ ok: result.ok,
+ output: result.output.clone(),
+ took: result.took,
+ });
+ chat.push(Message::tool_result(&result));
+ timer.tool_calls += 1;
+ }
+ timer.record(Stage::Tools, tools_started.elapsed());
+
+ // Con el resultado ya en la conversación, la siguiente vuelta sólo
+ // tiene que redactar: es el momento de recuperar la instrucción de
+ // voz, que en la pasada anterior habría impedido la llamada.
+ if prompts.two_phase {
+ chat.set_system(&prompts.speaking);
+ }</