From f8f98e83481e7376a235fb305a095552433f79ab Mon Sep 17 00:00:00 2001 From: elvis Date: Sun, 6 Sep 2026 19:21:16 -0300 Subject: Local voice assistant on top of Canary, llama.cpp and qwentts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 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 --- crates/asist-app/src/main.rs | 474 +++++++++++++++++++++++++++ crates/asist-app/src/pipeline.rs | 634 +++++++++++++++++++++++++++++++++++++ crates/asist-app/src/render.rs | 258 +++++++++++++++ crates/asist-app/src/session.rs | 160 ++++++++++ crates/asist-app/src/supervisor.rs | 228 +++++++++++++ 5 files changed, 1754 insertions(+) create mode 100644 crates/asist-app/src/main.rs create mode 100644 crates/asist-app/src/pipeline.rs create mode 100644 crates/asist-app/src/render.rs create mode 100644 crates/asist-app/src/session.rs create mode 100644 crates/asist-app/src/supervisor.rs (limited to 'crates/asist-app/src') 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::(); + 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 Fichero de configuración [config/asistente.toml] + --log-dir Carpeta de los registros de los servidores [logs/] + --voice 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, + no_manage: bool, + no_partials: bool, + barge_in: bool, + shell: bool, + verbose: bool, + help: bool, +} + +impl Args { + fn parse() -> Result { + 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 = 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>, + pub events: Receiver, + pub metrics: Arc, + session: Session, + capture_tx: Option>, +} + +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, + capture_tx: Sender, + input_format: InputFormat, + ) -> Result { + let (event_tx, event_rx) = unbounded::(); + let (asr_tx, asr_rx) = unbounded::(); + let (result_tx, result_rx) = unbounded::(); + // 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::(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, + asr: Sender, + events: Sender, + 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, + speak: Sender, + events: Sender, + metrics: Arc, +) { + 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, + events: &Sender, + 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); + } + } +} + +// --------------------------------------------------------------------------- +// Síntesis: frases -> altavoz +// --------------------------------------------------------------------------- + +fn run_speech( + session: Session, + tts: TtsClient, + jobs: Receiver, + playback: PlaybackHandle, + events: Sender, +) { + // Un `Select` en vez de `recv()` a secas para poder despertar + // periódicamente y bajar la bandera de «hablando» cuando el anillo se + // vacía: si no, el micrófono seguiría cerrado tras la última frase. + let mut select = Select::new(); + let job_index = select.recv(&jobs); + let mut speaking_turn: Option = None; + + loop { + let job = match select.select_timeout(Duration::from_millis(100)) { + Ok(op) if op.index() == job_index => match op.recv(&jobs) { + Ok(job) => Some(job), + Err(_) => break, + }, + Ok(_) => None, + Err(_) => None, + }; + + let Some(job) = job else { + // Sin trabajo: si ya no queda audio, el turno ha terminado de sonar. + if let Some(turn) = speaking_turn { + if !playback.is_active() || playback.queued_secs() <= 0.0 { + session.set_speaking(false); + speaking_turn = None; + let _ = events.send(Event::AudioFinished { turn }); + } + } + if session.is_stopping() { + break; + } + continue; + }; + + // Una frase de un turno ya superado no debe llegar a oírse. + if !session.is_current(job.turn) { + tracing::debug!(target: "tts", turno = job.turn.0, "frase obsoleta, descartada"); + continue; + } + + session.set_speaking(true); + speaking_turn = Some(job.turn); + if job.index == 0 { + playback.reset_played(); + } + + let turn = job.turn; + let first_of_turn = job.index == 0; + let mut announced = false; + let result = tts.speak(&job.text, session.tts_cancel(), |samples| { + if !session.is_current(turn) { + return false; + } + playback.push_tts(samples); + if first_of_turn && !announced && playback.has_played() { + announced = true; + let _ = events.send(Event::AudioStarted { + turn, + latency: job.turn_started.elapsed(), + }); + } + true + }); + + match result { + Ok(outcome) => { + if outcome.cancelled { + let _ = events.send(Event::Interrupted { + turn, + reason: InterruptReason::UserSpoke, + }); + continue; + } + // La cuenta del primer audio puede no haberse dado si el + // anillo aún no había servido nada cuando llegó el bloque. + if first_of_turn && !announced { + let _ = events.send(Event::AudioStarted { + turn, + latency: job.turn_started.elapsed(), + }); + } + if outcome.rtf() > 1.0 { + tracing::warn!( + target: "tts", + rtf = outcome.rtf(), + "la síntesis va por detrás del tiempo real; la voz se cortará a trozos" + ); + } + } + Err(err) => { + let _ = events.send(Event::Failed { + turn, + message: err.to_string(), + }); + } + } + } + + playback.stop(); + session.set_speaking(false); +} diff --git a/crates/asist-app/src/render.rs b/crates/asist-app/src/render.rs new file mode 100644 index 0000000..7fd1e86 --- /dev/null +++ b/crates/asist-app/src/render.rs @@ -0,0 +1,258 @@ +//! Presentación en el terminal. +//! +//! Una sola línea viva que se reescribe con la transcripción provisional, y +//! líneas fijas para lo que ya es definitivo. La distinción importa: ver que +//! el asistente te está oyendo *mientras* hablas es lo que hace que la espera +//! se note menos de lo que dura. + +use std::io::Write; +use std::time::Duration; + +use asist_core::event::{Event, InterruptReason}; + +const RESET: &str = "\x1b[0m"; +const DIM: &str = "\x1b[2m"; +const BOLD: &str = "\x1b[1m"; +const CYAN: &str = "\x1b[36m"; +const GREEN: &str = "\x1b[32m"; +const YELLOW: &str = "\x1b[33m"; +const RED: &str = "\x1b[31m"; + +pub struct Renderer { + width: usize, + line_open: bool, + color: bool, + show_partials: bool, + reply: String, +} + +impl Renderer { + pub fn new(show_partials: bool) -> Self { + Self { + width: terminal_width(), + line_open: false, + // Sin terminal interactivo, los códigos de color sólo ensucian + // el fichero de registro. + color: std::env::var_os("NO_COLOR").is_none() && is_tty(), + show_partials, + reply: String::new(), + } + } + + fn paint(&self, code: &str, text: &str) -> String { + if self.color { + format!("{code}{text}{RESET}") + } else { + text.to_string() + } + } + + /// Borra la línea viva para poder escribir encima algo definitivo. + fn close_line(&mut self) { + if self.line_open { + print!("\r\x1b[K"); + self.line_open = false; + } + } + + fn live(&mut self, text: &str) { + let mut shown: String = text.chars().collect(); + if shown.chars().count() > self.width.saturating_sub(4) { + let keep = self.width.saturating_sub(7); + let skip = shown.chars().count() - keep; + shown = format!("…{}", shown.chars().skip(skip).collect::()); + } + print!("\r\x1b[K{shown}"); + let _ = std::io::stdout().flush(); + self.line_open = true; + } + + pub fn apply(&mut self, event: &Event) { + match event { + Event::SpeechStarted { .. } => { + self.reply.clear(); + if self.show_partials { + self.live(&self.paint(DIM, "escuchando…")); + } + } + Event::Partial { + committed, + volatile, + .. + } => { + if !self.show_partials { + return; + } + let line = format!( + "{} {}", + self.paint(DIM, committed), + self.paint(DIM, volatile) + ); + self.live(&line); + } + Event::Transcript { + text, + audio_secs, + decode, + .. + } => { + self.close_line(); + println!( + "{} {} {}", + self.paint(BOLD, "tú >"), + text, + self.paint( + DIM, + &format!("({audio_secs:.1} s, asr {} ms)", decode.as_millis()) + ) + ); + } + Event::Discarded { .. } => { + self.close_line(); + } + Event::ReplyStarted { ttft, .. } => { + self.close_line(); + print!("{} ", self.paint(CYAN, "asistente >")); + let _ = std::io::stdout().flush(); + tracing::debug!(target: "render", ttft_ms = ttft.as_millis(), "primer token"); + } + Event::ReplyDelta { text, .. } => { + self.reply.push_str(text); + print!("{text}"); + let _ = std::io::stdout().flush(); + } + Event::ReplyDone { .. } => { + println!(); + } + Event::ToolRequested { + name, arguments, .. + } => { + self.close_line(); + println!( + "{} {name} {}", + self.paint(YELLOW, "herramienta >"), + self.paint(DIM, &short(arguments, 120)) + ); + } + Event::ToolFinished { + name, + ok, + output, + took, + .. + } => { + let mark = if *ok { + self.paint(GREEN, "ok") + } else { + self.paint(RED, "error") + }; + println!( + "{} {name} {mark} {}", + self.paint(YELLOW, "herramienta <"), + self.paint( + DIM, + &format!("{} · {} ms", short(output, 120), took.as_millis()) + ) + ); + } + Event::AudioStarted { latency, .. } => { + tracing::info!( + target: "latencia", + ms = latency.as_millis(), + "primer audio (latencia percibida)" + ); + if self.color { + println!( + "{}", + self.paint(DIM, &format!(" ▸ voz en {}", human(*latency))) + ); + } + } + Event::Sentence { index, text, .. } => { + tracing::debug!(target: "render", frase = index, %text, "a sintetizar"); + } + Event::AudioFinished { .. } => {} + Event::Interrupted { reason, .. } => { + self.close_line(); + let why = match reason { + InterruptReason::UserSpoke => "te has adelantado", + InterruptReason::Requested => "cortado", + }; + println!("{}", self.paint(YELLOW, &format!(" ▪ {why}"))); + } + Event::Warning { message, .. } => { + self.close_line(); + println!("{} {message}", self.paint(YELLOW, "aviso >")); + } + Event::Failed { message, .. } => { + self.close_line(); + println!("{} {message}", self.paint(RED, "error >")); + } + Event::Shutdown => { + self.close_line(); + } + } + } + + pub fn finish(&mut self) { + self.close_line(); + let _ = std::io::stdout().flush(); + } +} + +fn short(text: &str, max: usize) -> String { + let text = text.replace('\n', " "); + if text.chars().count() <= max { + return text; + } + text.chars().take(max).collect::() + "…" +} + +fn human(d: Duration) -> String { + if d.as_millis() < 1000 { + format!("{} ms", d.as_millis()) + } else { + format!("{:.1} s", d.as_secs_f32()) + } +} + +fn terminal_width() -> usize { + std::env::var("COLUMNS") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(100) +} + +fn is_tty() -> bool { + #[cfg(unix)] + { + unsafe extern "C" { + fn isatty(fd: i32) -> i32; + } + unsafe { isatty(1) == 1 } + } + #[cfg(not(unix))] + true +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn el_texto_largo_se_recorta_con_puntos_suspensivos() { + assert_eq!(short("abcdefghij", 5), "abcde…"); + assert_eq!(short("abc", 5), "abc"); + } + + #[test] + fn los_saltos_de_linea_no_rompen_la_linea_viva() { + assert_eq!(short("a\nb", 10), "a b"); + } + + #[test] + fn las_duraciones_se_leen_en_la_unidad_adecuada() { + assert_eq!(human(Duration::from_millis(430)), "430 ms"); + assert_eq!(human(Duration::from_millis(2500)), "2.5 s"); + } +} diff --git a/crates/asist-app/src/session.rs b/crates/asist-app/src/session.rs new file mode 100644 index 0000000..b74a4f3 --- /dev/null +++ b/crates/asist-app/src/session.rs @@ -0,0 +1,160 @@ +//! Estado compartido entre los hilos del pipeline. +//! +//! Sólo hay dos cosas que de verdad tienen que compartirse: qué turno es el +//! vigente, y las banderas que permiten cortar lo que está en marcha. Todo lo +//! demás viaja por canales. Mantener esta superficie pequeña es lo que hace +//! que la interrupción sea razonable de seguir: cortar es subir el turno y +//! levantar dos banderas. + +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::Arc; + +use asist_core::event::TurnId; +use asist_core::http::Cancel; + +#[derive(Clone)] +pub struct Session { + inner: Arc, +} + +struct Inner { + /// Turno que se está atendiendo. Cualquier trabajo de un turno anterior + /// que aparezca por un canal es basura y se tira. + current: AtomicU64, + /// El asistente está hablando (o a punto de hacerlo). + speaking: AtomicBool, + /// Cierre solicitado. + stopping: AtomicBool, + /// Corta la respuesta del modelo en curso. + llm: Cancel, + /// Corta la síntesis en curso. + tts: Cancel, +} + +impl Default for Session { + fn default() -> Self { + Self::new() + } +} + +impl Session { + pub fn new() -> Self { + Self { + inner: Arc::new(Inner { + current: AtomicU64::new(0), + speaking: AtomicBool::new(false), + stopping: AtomicBool::new(false), + llm: Cancel::new(), + tts: Cancel::new(), + }), + } + } + + pub fn current(&self) -> TurnId { + TurnId(self.inner.current.load(Ordering::SeqCst)) + } + + /// Abre un turno nuevo y devuelve su identificador. + pub fn begin_turn(&self) -> TurnId { + let id = TurnId(self.inner.current.fetch_add(1, Ordering::SeqCst) + 1); + self.inner.llm.reset(); + self.inner.tts.reset(); + id + } + + /// ¿Sigue siendo `turn` el turno vigente? + /// + /// Lo consultan los hilos antes de gastar trabajo: una frase que + /// pertenece a un turno ya superado no debe sintetizarse ni oírse. + pub fn is_current(&self, turn: TurnId) -> bool { + turn == self.current() + } + + /// Corta todo lo que esté en marcha para el turno actual. + pub fn interrupt(&self) { + self.inner.llm.cancel(); + self.inner.tts.cancel(); + self.inner.speaking.store(false, Ordering::SeqCst); + } + + pub fn llm_cancel(&self) -> &Cancel { + &self.inner.llm + } + + pub fn tts_cancel(&self) -> &Cancel { + &self.inner.tts + } + + pub fn set_speaking(&self, speaking: bool) { + self.inner.speaking.store(speaking, Ordering::SeqCst); + } + + pub fn is_speaking(&self) -> bool { + self.inner.speaking.load(Ordering::SeqCst) + } + + pub fn request_stop(&self) { + self.inner.stopping.store(true, Ordering::SeqCst); + self.interrupt(); + } + + pub fn is_stopping(&self) -> bool { + self.inner.stopping.load(Ordering::SeqCst) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cada_turno_recibe_un_identificador_creciente() { + let session = Session::new(); + assert_eq!(session.begin_turn(), TurnId(1)); + assert_eq!(session.begin_turn(), TurnId(2)); + assert_eq!(session.current(), TurnId(2)); + } + + #[test] + fn el_trabajo_de_un_turno_viejo_deja_de_ser_vigente() { + let session = Session::new(); + let primero = session.begin_turn(); + session.begin_turn(); + assert!( + !session.is_current(primero), + "el turno viejo debe descartarse" + ); + } + + #[test] + fn abrir_turno_rearma_las_cancelaciones() { + let session = Session::new(); + session.begin_turn(); + session.interrupt(); + assert!(session.llm_cancel().is_cancelled()); + + session.begin_turn(); + assert!( + !session.llm_cancel().is_cancelled(), + "un turno nuevo no puede heredar la cancelación del anterior" + ); + } + + #[test] + fn interrumpir_calla_y_cancela_las_dos_etapas() { + let session = Session::new(); + session.begin_turn(); + session.set_speaking(true); + session.interrupt(); + assert!(session.tts_cancel().is_cancelled()); + assert!(!session.is_speaking()); + } + + #[test] + fn pedir_el_cierre_tambien_interrumpe() { + let session = Session::new(); + session.request_stop(); + assert!(session.is_stopping()); + assert!(session.llm_cancel().is_cancelled()); + } +} diff --git a/crates/asist-app/src/supervisor.rs b/crates/asist-app/src/supervisor.rs new file mode 100644 index 0000000..17360da --- /dev/null +++ b/crates/asist-app/src/supervisor.rs @@ -0,0 +1,228 @@ +//! Arranque y parada de los servidores locales. +//! +//! El asistente puede levantar `llama-server` y `tts-server` él mismo, para +//! que ponerlo en marcha sea un solo comando. Los procesos se lanzan en su +//! propio grupo y se paran con SIGTERM antes de recurrir a SIGKILL: matar en +//! seco a llama-server deja la GPU ocupada hasta que el driver la recupera. + +use std::path::Path; +use std::process::{Child, Command, Stdio}; +use std::time::{Duration, Instant}; + +use asist_core::config::{Config, LlamaProcess, TtsProcess}; +use asist_core::error::{Error, Result}; + +pub struct Supervisor { + children: Vec, + log_dir: std::path::PathBuf, +} + +struct Managed { + name: &'static str, + child: Child, +} + +impl Supervisor { + pub fn new(log_dir: impl AsRef) -> Result { + let log_dir = log_dir.as_ref().to_path_buf(); + std::fs::create_dir_all(&log_dir)?; + Ok(Self { + children: Vec::new(), + log_dir, + }) + } + + /// Lanza los servidores que la configuración pida y que no estén ya arriba. + /// + /// Reutilizar uno que ya escucha es deliberado: durante el desarrollo se + /// reinicia el asistente muchas veces y volver a cargar los modelos + /// cuesta más de un minuto. + pub fn start(&mut self, config: &Config, llm_up: bool, tts_up: bool) -> Result<()> { + if !config.supervisor.manage { + return Ok(()); + } + if llm_up { + tracing::info!(target: "supervisor", "llama-server ya está escuchando, se reutiliza"); + } else { + let command = llama_command(&config.supervisor.llama, config)?; + self.spawn("llama-server", command)?; + } + if tts_up { + tracing::info!(target: "supervisor", "tts-server ya está escuchando, se reutiliza"); + } else { + let command = tts_command(&config.supervisor.tts, config)?; + self.spawn("tts-server", command)?; + } + Ok(()) + } + + fn spawn(&mut self, name: &'static str, mut command: Command) -> Result<()> { + let log_path = self.log_dir.join(format!("{name}.log")); + let log = std::fs::File::create(&log_path)?; + let errors = log.try_clone()?; + + // La salida va al fichero: los modelos escupen cientos de líneas y + // taparían la transcripción en pantalla. Cuando algo falla, el + // mensaje de error apunta aquí. + command + .stdin(Stdio::null()) + .stdout(Stdio::from(log)) + .stderr(Stdio::from(errors)); + + let child = command.spawn().map_err(|e| { + Error::Config(format!( + "no se pudo lanzar {name} ({:?}): {e}", + command.get_program() + )) + })?; + tracing::info!( + target: "supervisor", + proceso = name, + pid = child.id(), + registro = %log_path.display(), + "lanzado" + ); + self.children.push(Managed { name, child }); + Ok(()) + } + + /// Comprueba si alguno se ha muerto solo, y devuelve su nombre. + pub fn crashed(&mut self) -> Option<(&'static str, Option)> { + for managed in &mut self.children { + if let Ok(Some(status)) = managed.child.try_wait() { + return Some((managed.name, status.code())); + } + } + None + } + + pub fn log_path(&self, name: &str) -> std::path::PathBuf { + self.log_dir.join(format!("{name}.log")) + } + + /// Para todo lo lanzado. SIGTERM primero para que liberen la GPU. + pub fn shutdown(&mut self) { + for managed in &mut self.children { + terminate(&mut managed.child, managed.name); + } + self.children.clear(); + } +} + +impl Drop for Supervisor { + fn drop(&mut self) { + self.shutdown(); + } +} + +fn terminate(child: &mut Child, name: &str) { + if matches!(child.try_wait(), Ok(Some(_))) { + return; + } + // SIGTERM directo: `Child::kill` manda SIGKILL, que no da al servidor + // ocasión de soltar la memoria de la GPU. + #[cfg(unix)] + unsafe { + libc_kill(child.id() as i32, 15); + } + #[cfg(not(unix))] + let _ = child.kill(); + + let deadline = Instant::now() + Duration::from_secs(10); + while Instant::now() < deadline { + match child.try_wait() { + Ok(Some(_)) => { + tracing::info!(target: "supervisor", proceso = name, "detenido"); + return; + } + Ok(None) => std::thread::sleep(Duration::from_millis(100)), + Err(_) => break, + } + } + tracing::warn!(target: "supervisor", proceso = name, "no respondió a SIGTERM, se fuerza"); + let _ = child.kill(); + let _ = child.wait(); +} + +#[cfg(unix)] +unsafe fn libc_kill(pid: i32, signal: i32) { + // Se declara aquí el único símbolo de libc que hace falta, en vez de + // arrastrar el crate entero por una llamada. + unsafe extern "C" { + fn kill(pid: i32, sig: i32) -> i32; + } + unsafe { + kill(pid, signal); + } +} + +fn require(path: &Path, what: &str) -> Result<()> { + if path.exists() { + return Ok(()); + } + Err(Error::Config(format!( + "falta {what}: {}. Ejecuta scripts/bootstrap.sh para compilar los \ + motores y enlazar los modelos", + path.display() + ))) +} + +fn llama_command(process: &LlamaProcess, config: &Config) -> Result { + require(&process.binary, "el binario de llama-server")?; + require(&process.model, "el modelo del LLM")?; + + let mut command = Command::new(&process.binary); + command + .arg("--model") + .arg(&process.model) + .arg("--host") + .arg(&config.llm.host) + .arg("--port") + .arg(config.llm.port.to_string()); + + if process.mmproj.exists() { + command.arg("--mmproj").arg(&process.mmproj); + } + // Sin esto, la plantilla del modelo abre y no lo cierra nunca: el + // asistente se pasa entre 7 y 9 s razonando antes de la primera palabra. + if process.chat_template.exists() { + command + .arg("--jinja") + .arg("--chat-template-file") + .arg(&process.chat_template); + } else { + tracing::warn!( + target: "supervisor", + ruta = %process.chat_template.display(), + "no está la plantilla sin razonamiento: el modelo tardará varios \ + segundos en empezar a hablar" + ); + } + command.args(&process.extra_args); + Ok(command) +} + +fn tts_command(process: &TtsProcess, config: &Config) -> Result { + require(&process.binary, "el binario de tts-server")?; + require(&process.model, "el modelo hablante del TTS")?; + require(&process.codec, "el codec del TTS")?; + + let mut command = Command::new(&process.binary); + command + .arg("--model") + .arg(&process.model) + .arg("--codec") + .arg(&process.codec) + .arg("--host") + .arg(&config.tts.host) + .arg("--port") + .arg(config.tts.port.to_string()) + .arg("--lang") + .arg(&config.tts.language) + // El ajuste con más efecto de todo el sistema: de fábrica son 24 s, que + // en la práctica significa no devolver nada hasta terminar la frase. + .arg("--codec-chunk-dur") + .arg(process.codec_chunk_dur.to_string()); + command.args(&process.extra_args); + Ok(command) +} -- cgit v1.2.3