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-core/Cargo.toml | 15 + crates/asist-core/src/config.rs | 508 +++++++++++++++++++++++++++++++++ crates/asist-core/src/error.rs | 82 ++++++ crates/asist-core/src/event.rs | 127 +++++++++ crates/asist-core/src/http.rs | 351 +++++++++++++++++++++++ crates/asist-core/src/lib.rs | 17 ++ crates/asist-core/src/telemetry.rs | 255 +++++++++++++++++ crates/asist-core/src/text.rs | 330 ++++++++++++++++++++++ crates/asist-core/src/tools.rs | 563 +++++++++++++++++++++++++++++++++++++ 9 files changed, 2248 insertions(+) create mode 100644 crates/asist-core/Cargo.toml create mode 100644 crates/asist-core/src/config.rs create mode 100644 crates/asist-core/src/error.rs create mode 100644 crates/asist-core/src/event.rs create mode 100644 crates/asist-core/src/http.rs create mode 100644 crates/asist-core/src/lib.rs create mode 100644 crates/asist-core/src/telemetry.rs create mode 100644 crates/asist-core/src/text.rs create mode 100644 crates/asist-core/src/tools.rs (limited to 'crates/asist-core') diff --git a/crates/asist-core/Cargo.toml b/crates/asist-core/Cargo.toml new file mode 100644 index 0000000..bd8c3ae --- /dev/null +++ b/crates/asist-core/Cargo.toml @@ -0,0 +1,15 @@ +[package] +name = "asist-core" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true + +[dependencies] +anyhow.workspace = true +crossbeam-channel.workspace = true +serde.workspace = true +serde_json.workspace = true +thiserror.workspace = true +toml.workspace = true +tracing.workspace = true diff --git a/crates/asist-core/src/config.rs b/crates/asist-core/src/config.rs new file mode 100644 index 0000000..4ca8d8c --- /dev/null +++ b/crates/asist-core/src/config.rs @@ -0,0 +1,508 @@ +//! Configuración del asistente: un TOML con valores por defecto que ya +//! incorporan lo aprendido midiendo el pipeline (ver `docs/RENDIMIENTO.md`). + +use std::path::{Path, PathBuf}; + +use serde::{Deserialize, Serialize}; + +use crate::error::{Error, Result}; + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct Config { + pub general: General, + pub audio: AudioConfig, + pub vad: VadConfig, + pub asr: AsrConfig, + pub llm: LlmConfig, + pub tts: TtsConfig, + pub tools: ToolsConfig, + pub supervisor: SupervisorConfig, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct General { + /// Idioma de la conversación; se propaga a ASR y TTS. + pub language: String, + /// Instrucción de sistema. Pide frases cortas y sin markdown a propósito: + /// el TTS lee literalmente los asteriscos y las viñetas. + pub system_prompt: String, + /// Instrucción de sistema de la pasada en que el modelo decide si usar + /// una herramienta. + /// + /// Va aparte, y **sustituye** a `system_prompt` en esa pasada, por una + /// razón medida y no por gusto: con Qwen3.5-2B, añadir cualquier + /// indicación de estilo a esta guía —en cualquier posición, incluso dos + /// palabras— hace que el modelo deje de llamar a las herramientas y se + /// invente el dato. Medido sobre 8 intentos: la guía sola acierta 8/8; + /// con «Responde breve.» detrás, 1/8; con la persona de asistente de voz, + /// 0/8. Ver docs/RENDIMIENTO.md. + pub tools_prompt: String, + /// Turnos de historial que se envían al modelo (0 = sin memoria). + pub history_turns: usize, + /// Imprime un resumen de latencias por turno al terminar cada respuesta. + pub report_latency: bool, +} + +impl Default for General { + fn default() -> Self { + Self { + language: "es".into(), + system_prompt: concat!( + "Eres un asistente de voz en español. Tus respuestas se leen en voz alta, ", + "así que responde en una o dos frases cortas, en texto plano corrido. ", + "No uses markdown, ni listas, ni asteriscos, ni emojis, ni encabezados. ", + "No escribas URLs ni código salvo que te lo pidan explícitamente. ", + "Si no sabes algo, dilo en una frase." + ) + .into(), + tools_prompt: concat!( + "Antes de responder, comprueba si alguna de tus herramientas te da el dato. ", + "Si es así, llámala primero y espera su resultado; no contestes de memoria. ", + "Sólo cuando tengas el resultado, resúmelo en una frase." + ) + .into(), + history_turns: 8, + report_latency: true, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct AudioConfig { + /// Nombre (o subcadena) del dispositivo de entrada; vacío = el de por defecto. + pub input_device: String, + /// Ídem para la salida. + pub output_device: String, + /// Segundos de audio que la reproducción mantiene en cola antes de + /// arrancar. Amortigua los baches del TTS sin añadir latencia perceptible. + pub playback_prebuffer: f32, + /// Ganancia aplicada a la reproducción. + pub output_gain: f32, +} + +impl Default for AudioConfig { + fn default() -> Self { + Self { + input_device: String::new(), + output_device: String::new(), + playback_prebuffer: 0.20, + output_gain: 1.0, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct VadConfig { + /// Ventana de audio por decisión del detector. + pub frame_seconds: f32, + /// Audio conservado antes del inicio de voz para no cortar la primera sílaba. + pub preroll_seconds: f32, + /// Silencio que cierra una intervención. + pub silence_hold: f32, + /// Intervención más corta que merezca una transcripción final. + pub min_utterance: f32, + /// Corte forzado, que acota el coste de la decodificación final. + pub max_utterance: f32, + /// Múltiplo del suelo de ruido a partir del cual se considera voz. + pub threshold_factor: f32, + pub min_threshold: f32, + pub max_threshold: f32, + /// Permite hablar encima del asistente para cortarlo. + /// + /// Desactivado por defecto: con altavoces abiertos el micrófono se oye a sí + /// mismo y el asistente se interrumpe solo. Actívalo con auriculares o con + /// cancelación de eco del sistema. + pub barge_in: bool, + /// Con barge-in activo, cuánto más fuerte que el umbral normal debe sonar + /// la voz para cortar. Sube el listón frente al eco del altavoz. + pub barge_in_factor: f32, +} + +impl Default for VadConfig { + fn default() -> Self { + Self { + frame_seconds: 0.1, + preroll_seconds: 0.3, + silence_hold: 0.8, + min_utterance: 0.3, + max_utterance: 20.0, + threshold_factor: 3.0, + min_threshold: 0.0008, + max_threshold: 0.02, + barge_in: false, + barge_in_factor: 4.0, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct AsrConfig { + /// Carpeta del modelo Canary en ONNX. + pub model_dir: PathBuf, + pub source_lang: String, + pub target_lang: String, + /// Ventana deslizante de las transcripciones provisionales. + pub window: f32, + /// Avance de la ventana entre decodificaciones. + pub step: f32, + /// Ventanas que deben coincidir para dar una palabra por estable. + pub stability: usize, + /// Emitir transcripciones provisionales. Cuestan CPU y sólo sirven para + /// verlas en pantalla: el turno se dispara con la final. + pub partials: bool, + /// Carga una segunda instancia del modelo para que las decodificaciones + /// finales no bloqueen a las provisionales. Duplica la memoria. + pub dedicated_final_model: bool, + /// Proveedor de ejecución de ONNX Runtime: `cpu`, `cuda`, ... + pub execution_provider: String, + pub inter_threads: usize, + pub intra_threads: usize, +} + +impl Default for AsrConfig { + fn default() -> Self { + Self { + model_dir: PathBuf::from("vendor/canary-rs/models/canary-180m-flash-onnx"), + source_lang: "es".into(), + target_lang: "es".into(), + window: 6.0, + step: 0.4, + stability: 2, + partials: true, + dedicated_final_model: false, + // CPU a propósito: la GPU de 4 GB está ocupada por el hablante del + // TTS, y disputársela sale más caro que decodificar en CPU. + execution_provider: "cpu".into(), + inter_threads: 2, + intra_threads: 4, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct LlmConfig { + pub host: String, + pub port: u16, + pub model: String, + pub temperature: f32, + pub top_p: f32, + pub top_k: i32, + pub min_p: f32, + pub repeat_penalty: f32, + pub max_tokens: u32, + /// Vueltas máximas del bucle de herramientas antes de rendirse. + pub max_tool_rounds: usize, + pub request_timeout_secs: u64, +} + +impl Default for LlmConfig { + fn default() -> Self { + Self { + host: "127.0.0.1".into(), + port: 8012, + model: String::new(), + temperature: 0.7, + top_p: 0.9, + top_k: 40, + min_p: 0.1, + repeat_penalty: 1.1, + // Una respuesta hablada larga cansa; el recorte también acota el + // coste de la síntesis, que es la etapa lenta. + max_tokens: 300, + max_tool_rounds: 4, + request_timeout_secs: 120, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct TtsConfig { + pub host: String, + pub port: u16, + /// Voz registrada en el servidor, o una del modelo. + pub voice: String, + /// Voz clonada que se registra al arrancar, si se define. + pub reference: Option, + pub language: String, + pub temperature: f32, + pub top_k: i32, + pub top_p: f32, + pub repetition_penalty: f32, + pub max_new_tokens: u32, + /// Sintetiza una frase corta al arrancar. La primera petición paga la + /// construcción de los grafos: ~3,5 s que conviene no gastar en el + /// primer turno real. + pub warmup: bool, + pub request_timeout_secs: u64, +} + +impl Default for TtsConfig { + fn default() -> Self { + Self { + host: "127.0.0.1".into(), + port: 8013, + voice: "asistente".into(), + reference: None, + language: "spanish".into(), + temperature: 0.9, + top_k: 50, + top_p: 1.0, + repetition_penalty: 1.05, + max_new_tokens: 2048, + warmup: true, + request_timeout_secs: 180, + } + } +} + +/// Latentes de una voz clonada, tal y como los produce `qwen-codec --talker`. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ReferenceVoice { + /// Nombre con el que se registra en el servidor. + pub name: String, + /// Embedding del hablante (`.spk`). + pub speaker: PathBuf, + /// Códigos de referencia (`.rvq`), que activan el clonado ICL. + pub codes: PathBuf, + /// Transcripción de la referencia (`.txt`). + pub transcript: PathBuf, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct ToolsConfig { + /// Deja que el modelo llame a herramientas. + pub enabled: bool, + /// Usa `general.tools_prompt` en solitario para la pasada en que el + /// modelo decide si llamar a una herramienta, y `general.system_prompt` + /// sólo para redactar la respuesta hablada. + /// + /// Es lo único que hace que las herramientas funcionen de verdad con este + /// modelo (ver `tools_prompt`), a cambio de que las respuestas que no + /// usan herramienta pierdan la guía de estilo. Se compensa en parte + /// porque el limpiador de texto quita el markdown antes de hablar. + /// Ponlo a `false` para priorizar el estilo sobre las herramientas. + pub dedicated_prompt: bool, + /// Habilita la herramienta de ejecución de órdenes del sistema. + /// + /// Apagada por defecto a conciencia: darle una shell a un modelo que + /// obedece a lo que oye por el micrófono es un cambio de postura de + /// seguridad, no una opción de comodidad. + pub shell: bool, + /// Órdenes admitidas, comparadas contra el ejecutable (argv[0]). + /// Una lista vacía deniega todo aunque `shell` esté activo. + pub shell_allowlist: Vec, + /// Segundos que puede durar una orden antes de que se la mate. + pub shell_timeout_secs: u64, + /// Registra la orden pero no la ejecuta. Útil para estrenar la lista blanca. + pub shell_dry_run: bool, + /// Directorio de trabajo de las órdenes; vacío = el del proceso. + pub shell_working_dir: String, +} + +impl Default for ToolsConfig { + fn default() -> Self { + Self { + enabled: true, + dedicated_prompt: true, + shell: false, + shell_allowlist: vec![ + "date".into(), + "uptime".into(), + "free".into(), + "df".into(), + "ls".into(), + ], + shell_timeout_secs: 10, + shell_dry_run: false, + shell_working_dir: String::new(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct SupervisorConfig { + /// Lanza los servidores en lugar de suponer que ya están arriba. + pub manage: bool, + /// Segundos de espera a que un servidor conteste a `/health`. + pub startup_timeout_secs: u64, + pub llama: LlamaProcess, + pub tts: TtsProcess, +} + +impl Default for SupervisorConfig { + fn default() -> Self { + Self { + manage: true, + startup_timeout_secs: 180, + llama: LlamaProcess::default(), + tts: TtsProcess::default(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct LlamaProcess { + pub binary: PathBuf, + pub model: PathBuf, + pub mmproj: PathBuf, + /// Plantilla de chat que se pasa con `--chat-template-file`. + /// + /// La del modelo abre `` sin cerrarlo nunca, y el asistente se pasa + /// entonces varios segundos razonando antes de decir la primera palabra. + /// Esta copia deja el bloque cerrado de entrada. + pub chat_template: PathBuf, + pub extra_args: Vec, +} + +impl Default for LlamaProcess { + fn default() -> Self { + Self { + binary: PathBuf::from("vendor/llama.cpp/build/bin/llama-server"), + model: PathBuf::from("models/Qwen3.5-2B.Q5_K_M.gguf"), + mmproj: PathBuf::from("models/mmproj-BF16.gguf"), + chat_template: PathBuf::from("config/qwen35-no-think.jinja"), + extra_args: [ + "--threads", + "10", + "--threads-batch", + "10", + "--batch-size", + "512", + "--ubatch-size", + "256", + "--gpu-layers", + "10", + "--split-mode", + "layer", + "--tensor-split", + "1", + "--main-gpu", + "0", + "--no-mmap", + "--ctx-size", + "8192", + "--parallel", + "2", + "--cache-ram", + "6144", + "--rope-freq-base", + "1000000", + "--rope-freq-scale", + "0.25", + "--jinja", + ] + .iter() + .map(|s| s.to_string()) + .collect(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct TtsProcess { + pub binary: PathBuf, + pub model: PathBuf, + pub codec: PathBuf, + /// Segundos de audio que el codec acumula antes de decodificar un bloque. + /// + /// El valor de fábrica son 24 s, que en la práctica significa «no + /// devuelvas nada hasta terminar la frase entera»: medido, baja el primer + /// audio de 4,9 s a 0,6 s. Es el ajuste con más efecto de todo el sistema. + pub codec_chunk_dur: f32, + pub extra_args: Vec, +} + +impl Default for TtsProcess { + fn default() -> Self { + Self { + binary: PathBuf::from("vendor/qwentts.cpp/build/tts-server"), + model: PathBuf::from("models/qwen-talker-1.7b-base-Q8_0.gguf"), + codec: PathBuf::from("models/qwen-tokenizer-12hz-Q8_0.gguf"), + codec_chunk_dur: 1.0, + extra_args: Vec::new(), + } + } +} + +impl Config { + pub fn load(path: impl AsRef) -> Result { + let path = path.as_ref(); + let raw = std::fs::read_to_string(path) + .map_err(|e| Error::Config(format!("no se pudo leer {}: {e}", path.display())))?; + let mut config: Config = + toml::from_str(&raw).map_err(|e| Error::Config(format!("{}: {e}", path.display())))?; + // Las rutas relativas se resuelven contra la carpeta del TOML, no + // contra el directorio desde el que se lanza el binario. + if let Some(base) = path.parent().filter(|p| !p.as_os_str().is_empty()) { + config.rebase(base); + } + config.validate()?; + Ok(config) + } + + /// Reinterpreta las rutas relativas respecto de `base`. + pub fn rebase(&mut self, base: &Path) { + let fix = |p: &mut PathBuf| { + if p.is_relative() { + *p = base.join(&*p); + } + }; + fix(&mut self.asr.model_dir); + fix(&mut self.supervisor.llama.binary); + fix(&mut self.supervisor.llama.model); + fix(&mut self.supervisor.llama.mmproj); + fix(&mut self.supervisor.llama.chat_template); + fix(&mut self.supervisor.tts.binary); + fix(&mut self.supervisor.tts.model); + fix(&mut self.supervisor.tts.codec); + if let Some(reference) = self.tts.reference.as_mut() { + fix(&mut reference.speaker); + fix(&mut reference.codes); + fix(&mut reference.transcript); + } + } + + fn validate(&self) -> Result<()> { + if self.vad.silence_hold <= 0.0 { + return Err(Error::Config("vad.silence_hold debe ser > 0".into())); + } + if self.vad.min_utterance >= self.vad.max_utterance { + return Err(Error::Config( + "vad.min_utterance debe ser menor que vad.max_utterance".into(), + )); + } + if self.asr.step <= 0.0 || self.asr.window <= self.asr.step { + return Err(Error::Config( + "asr.window debe ser mayor que asr.step, y ambos > 0".into(), + )); + } + if self.tools.shell && self.tools.shell_allowlist.is_empty() { + return Err(Error::Config( + "tools.shell está activo pero tools.shell_allowlist está vacía: \ + declara las órdenes permitidas o desactiva tools.shell" + .into(), + )); + } + Ok(()) + } + + pub fn llm_authority(&self) -> String { + format!("{}:{}", self.llm.host, self.llm.port) + } + + pub fn tts_authority(&self) -> String { + format!("{}:{}", self.tts.host, self.tts.port) + } +} diff --git a/crates/asist-core/src/error.rs b/crates/asist-core/src/error.rs new file mode 100644 index 0000000..b559bee --- /dev/null +++ b/crates/asist-core/src/error.rs @@ -0,0 +1,82 @@ +use std::fmt; + +pub type Result = std::result::Result; + +/// Un fallo de una etapa del pipeline. El orquestador decide si es fatal o si +/// basta con abortar el turno en curso, así que el error lleva esa distinción +/// en lugar de dejarla al criterio de quien lo recibe. +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("configuración inválida: {0}")] + Config(String), + + #[error("audio: {0}")] + Audio(String), + + #[error("ASR: {0}")] + Asr(String), + + #[error("LLM: {0}")] + Llm(String), + + #[error("TTS: {0}")] + Tts(String), + + #[error("HTTP {status} en {url}: {body}")] + Http { + status: u16, + url: String, + body: String, + }, + + #[error("transporte hacia {url}: {source}")] + Transport { + url: String, + #[source] + source: std::io::Error, + }, + + #[error("herramienta «{tool}»: {message}")] + Tool { tool: String, message: String }, + + #[error("operación cancelada")] + Cancelled, + + #[error(transparent)] + Json(#[from] serde_json::Error), + + #[error(transparent)] + Io(#[from] std::io::Error), +} + +impl Error { + /// Un error fatal tumba el asistente; el resto sólo aborta el turno. + /// + /// Los servidores locales se caen y se reinician, y el usuario prefiere + /// oír «no te he entendido» a que el proceso muera, así que sólo la + /// configuración y el audio (que no se pueden reintentar) son fatales. + pub fn is_fatal(&self) -> bool { + matches!(self, Error::Config(_) | Error::Audio(_)) + } + + /// `true` cuando reintentar tiene sentido: el servidor aún está cargando + /// el modelo, o se ha caído la conexión a mitad de una petición. + pub fn is_retryable(&self) -> bool { + match self { + Error::Transport { .. } => true, + Error::Http { status, .. } => *status == 503 || *status == 429 || *status >= 500, + _ => false, + } + } +} + +/// Contexto legible para los `Result` que cruzan una frontera de crate. +pub trait Context { + fn ctx(self, f: impl FnOnce() -> String) -> Result; +} + +impl Context for std::result::Result { + fn ctx(self, f: impl FnOnce() -> String) -> Result { + self.map_err(|e| Error::Config(format!("{}: {}", f(), e))) + } +} diff --git a/crates/asist-core/src/event.rs b/crates/asist-core/src/event.rs new file mode 100644 index 0000000..549bec0 --- /dev/null +++ b/crates/asist-core/src/event.rs @@ -0,0 +1,127 @@ +use std::time::{Duration, Instant}; + +/// Identifica un turno de conversación (una intervención del usuario y la +/// respuesta que provoca). Todo lo que viaja por el bus lo lleva, para que un +/// resultado que llega tarde no se confunda con el turno que ya está en curso +/// —el caso típico cuando el usuario interrumpe al asistente. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Default)] +pub struct TurnId(pub u64); + +impl std::fmt::Display for TurnId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "#{}", self.0) + } +} + +/// Motivo por el que se corta la reproducción en curso. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum InterruptReason { + /// El usuario ha empezado a hablar encima del asistente (barge-in). + UserSpoke, + /// Petición explícita: tecla, señal o cierre. + Requested, +} + +/// Todo lo que las etapas se cuentan entre sí. Un único enum mantiene el bus +/// observable: el renderizador y la telemetría ven exactamente los mismos +/// eventos que el orquestador, sin canales paralelos que se desincronicen. +#[derive(Debug, Clone)] +pub enum Event { + /// El VAD ha detectado el arranque de una intervención. + SpeechStarted { turn: TurnId, at: Instant }, + + /// Transcripción provisional de la ventana deslizante: `committed` ya es + /// estable, `volatile` todavía puede cambiar. + Partial { + turn: TurnId, + committed: String, + volatile: String, + }, + + /// Transcripción definitiva de la intervención completa. + Transcript { + turn: TurnId, + text: String, + audio_secs: f32, + decode: Duration, + }, + + /// La intervención no contenía nada transcribible. + Discarded { turn: TurnId }, + + /// Primer fragmento de texto que devuelve el modelo. + ReplyStarted { turn: TurnId, ttft: Duration }, + + /// Un trozo de respuesta según llega del modelo. + ReplyDelta { turn: TurnId, text: String }, + + /// Una frase completa lista para sintetizar. + Sentence { + turn: TurnId, + index: usize, + text: String, + }, + + /// Respuesta completa del modelo para este turno. + ReplyDone { turn: TurnId, text: String }, + + /// El modelo ha pedido ejecutar una herramienta. + ToolRequested { + turn: TurnId, + name: String, + arguments: String, + }, + + /// Resultado de esa ejecución. + ToolFinished { + turn: TurnId, + name: String, + ok: bool, + output: String, + took: Duration, + }, + + /// Primer audio audible del turno: la métrica que de verdad percibe quien habla. + AudioStarted { turn: TurnId, latency: Duration }, + + /// Se ha terminado de reproducir todo lo del turno. + AudioFinished { turn: TurnId }, + + /// Se ha cortado la reproducción. + Interrupted { + turn: TurnId, + reason: InterruptReason, + }, + + /// Aviso no fatal; el turno continúa. + Warning { turn: TurnId, message: String }, + + /// Fallo que aborta el turno. + Failed { turn: TurnId, message: String }, + + /// Cierre ordenado. + Shutdown, +} + +impl Event { + pub fn turn(&self) -> Option { + match self { + Event::SpeechStarted { turn, .. } + | Event::Partial { turn, .. } + | Event::Transcript { turn, .. } + | Event::Discarded { turn } + | Event::ReplyStarted { turn, .. } + | Event::ReplyDelta { turn, .. } + | Event::Sentence { turn, .. } + | Event::ReplyDone { turn, .. } + | Event::ToolRequested { turn, .. } + | Event::ToolFinished { turn, .. } + | Event::AudioStarted { turn, .. } + | Event::AudioFinished { turn } + | Event::Interrupted { turn, .. } + | Event::Warning { turn, .. } + | Event::Failed { turn, .. } => Some(*turn), + Event::Shutdown => None, + } + } +} diff --git a/crates/asist-core/src/http.rs b/crates/asist-core/src/http.rs new file mode 100644 index 0000000..32b6680 --- /dev/null +++ b/crates/asist-core/src/http.rs @@ -0,0 +1,351 @@ +//! Cliente HTTP/1.1 mínimo para hablar con los servidores locales. +//! +//! Los tres motores (llama-server, tts-server) escuchan en localhost sobre HTTP +//! plano, así que no hace falta TLS ni un runtime asíncrono. A cambio de las +//! ~300 líneas de aquí se gana lo único que ninguna librería genérica da +//! cómodo y que este pipeline necesita de verdad: **abortar una respuesta a +//! media descarga**. Cuando el usuario interrumpe al asistente hay que dejar de +//! leer audio ya generado y cerrar el socket en ese mismo instante; con un +//! cliente que sólo expone `read_to_end` eso no se puede hacer. + +use std::io::{BufRead, BufReader, Read, Write}; +use std::net::TcpStream; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +use crate::error::{Error, Result}; + +/// Bandera compartida que aborta una lectura en curso. +/// +/// Se comprueba entre bloques, así que el corte ocurre como muy tarde un +/// `read` después de activarla —del orden de milisegundos con los tamaños de +/// bloque que usan los servidores. +#[derive(Clone, Default)] +pub struct Cancel(Arc); + +impl Cancel { + pub fn new() -> Self { + Self::default() + } + + pub fn cancel(&self) { + self.0.store(true, Ordering::SeqCst); + } + + pub fn is_cancelled(&self) -> bool { + self.0.load(Ordering::SeqCst) + } + + pub fn reset(&self) { + self.0.store(false, Ordering::SeqCst); + } +} + +#[derive(Clone, Debug)] +pub struct HttpClient { + /// `host:puerto` del servidor local. + pub authority: String, + pub connect_timeout: Duration, + /// Tiempo máximo sin recibir un solo byte. No es el tiempo total: una + /// síntesis larga puede tardar minutos siempre que siga fluyendo. + pub read_timeout: Duration, +} + +impl HttpClient { + pub fn new(authority: impl Into) -> Self { + Self { + authority: authority.into(), + connect_timeout: Duration::from_secs(5), + read_timeout: Duration::from_secs(120), + } + } + + pub fn with_read_timeout(mut self, timeout: Duration) -> Self { + self.read_timeout = timeout; + self + } + + fn connect(&self, path: &str) -> Result { + let addr = resolve(&self.authority, path)?; + let stream = TcpStream::connect_timeout(&addr, self.connect_timeout).map_err(|e| { + Error::Transport { + url: format!("http://{}{}", self.authority, path), + source: e, + } + })?; + stream.set_read_timeout(Some(self.read_timeout))?; + // Los cuerpos son pequeños y la latencia manda: nada de Nagle. + let _ = stream.set_nodelay(true); + Ok(stream) + } + + fn send(&self, method: &str, path: &str, body: Option<&[u8]>) -> Result { + let mut stream = self.connect(path)?; + let mut head = format!( + "{method} {path} HTTP/1.1\r\nHost: {}\r\nConnection: close\r\nAccept: */*\r\n", + self.authority + ); + if let Some(body) = body { + head.push_str("Content-Type: application/json\r\n"); + head.push_str(&format!("Content-Length: {}\r\n", body.len())); + } + head.push_str("\r\n"); + + let url = format!("http://{}{}", self.authority, path); + let io = |e: std::io::Error| Error::Transport { + url: url.clone(), + source: e, + }; + stream.write_all(head.as_bytes()).map_err(io)?; + if let Some(body) = body { + stream.write_all(body).map_err(io)?; + } + stream.flush().map_err(io)?; + + Response::read_head(BufReader::new(stream), url) + } + + /// GET que devuelve el cuerpo completo ya decodificado. + pub fn get(&self, path: &str) -> Result> { + self.send("GET", path, None)?.into_body() + } + + /// POST de JSON que devuelve el cuerpo completo. + pub fn post_json(&self, path: &str, body: &serde_json::Value) -> Result> { + let payload = serde_json::to_vec(body)?; + self.send("POST", path, Some(&payload))?.into_body() + } + + /// POST de JSON cuya respuesta se consume a trozos según llega. + /// + /// `on_chunk` recibe cada bloque en cuanto está disponible y devuelve + /// `false` para cortar; junto a `cancel` son las dos vías por las que el + /// barge-in detiene una síntesis a medias. + pub fn post_json_streaming( + &self, + path: &str, + body: &serde_json::Value, + cancel: &Cancel, + mut on_chunk: impl FnMut(&[u8]) -> bool, + ) -> Result<()> { + let payload = serde_json::to_vec(body)?; + let mut response = self.send("POST", path, Some(&payload))?; + response.error_for_status()?; + response.stream(cancel, &mut on_chunk) + } + + /// POST de JSON que entrega la respuesta línea a línea (SSE de llama-server). + pub fn post_json_lines( + &self, + path: &str, + body: &serde_json::Value, + cancel: &Cancel, + mut on_line: impl FnMut(&str) -> bool, + ) -> Result<()> { + let mut pending = Vec::new(); + self.post_json_streaming(path, body, cancel, |chunk| { + pending.extend_from_slice(chunk); + while let Some(nl) = pending.iter().position(|b| *b == b'\n') { + let line: Vec = pending.drain(..=nl).collect(); + let line = String::from_utf8_lossy(&line); + if !on_line(line.trim_end_matches(['\r', '\n'])) { + return false; + } + } + true + }) + } + + /// Sondea `/health` y devuelve `true` si el servidor está listo. + pub fn healthy(&self, path: &str) -> bool { + self.get(path).is_ok() + } +} + +fn resolve(authority: &str, path: &str) -> Result { + use std::net::ToSocketAddrs; + authority + .to_socket_addrs() + .map_err(|e| Error::Transport { + url: format!("http://{authority}{path}"), + source: e, + })? + .next() + .ok_or_else(|| Error::Config(format!("no se pudo resolver «{authority}»"))) +} + +/// Codificación del cuerpo anunciada por el servidor. +enum Body { + Chunked, + Length(usize), + /// Sin longitud ni chunked: el cuerpo acaba cuando cierra el socket. + UntilClose, +} + +pub struct Response { + pub status: u16, + url: String, + reader: BufReader, + body: Body, +} + +impl Response { + fn read_head(mut reader: BufReader, url: String) -> Result { + let io = |e: std::io::Error| Error::Transport { + url: url.clone(), + source: e, + }; + + let mut line = String::new(); + reader.read_line(&mut line).map_err(io)?; + let status = line + .split_whitespace() + .nth(1) + .and_then(|s| s.parse::().ok()) + .ok_or_else(|| Error::Config(format!("respuesta HTTP ilegible de {url}: {line:?}")))?; + + let mut body = Body::UntilClose; + loop { + let mut header = String::new(); + if reader.read_line(&mut header).map_err(io)? == 0 { + break; + } + let header = header.trim_end(); + if header.is_empty() { + break; + } + let Some((name, value)) = header.split_once(':') else { + continue; + }; + let (name, value) = (name.trim().to_ascii_lowercase(), value.trim()); + match name.as_str() { + "transfer-encoding" if value.eq_ignore_ascii_case("chunked") => { + body = Body::Chunked; + } + "content-length" => { + if let Ok(n) = value.parse() { + body = Body::Length(n); + } + } + _ => {} + } + } + + Ok(Self { + status, + url, + reader, + body, + }) + } + + fn error_for_status(&mut self) -> Result<()> { + if (200..300).contains(&self.status) { + return Ok(()); + } + let status = self.status; + let url = self.url.clone(); + // Leer el cuerpo del error es lo que convierte un «500» en un mensaje útil. + let body = self.drain().unwrap_or_default(); + Err(Error::Http { + status, + url, + body: String::from_utf8_lossy(&body).chars().take(500).collect(), + }) + } + + fn into_body(mut self) -> Result> { + self.error_for_status()?; + self.drain() + } + + fn drain(&mut self) -> Result> { + let mut out = Vec::new(); + let cancel = Cancel::new(); + self.stream(&cancel, &mut |chunk| { + out.extend_from_slice(chunk); + true + })?; + Ok(out) + } + + fn stream(&mut self, cancel: &Cancel, on_chunk: &mut dyn FnMut(&[u8]) -> bool) -> Result<()> { + let url = self.url.clone(); + let io = |e: std::io::Error| Error::Transport { + url: url.clone(), + source: e, + }; + let mut buf = vec![0u8; 16 * 1024]; + + match self.body { + Body::Chunked => loop { + if cancel.is_cancelled() { + return Err(Error::Cancelled); + } + let mut size_line = String::new(); + if self.reader.read_line(&mut size_line).map_err(io)? == 0 { + return Ok(()); + } + let size_line = size_line.trim(); + if size_line.is_empty() { + continue; + } + // Las extensiones de chunk van tras un ';' y aquí no interesan. + let size = + usize::from_str_radix(size_line.split(';').next().unwrap_or("0").trim(), 16) + .map_err(|_| { + Error::Config(format!("tamaño de chunk ilegible: {size_line:?}")) + })?; + if size == 0 { + return Ok(()); + } + let mut left = size; + while left > 0 { + if cancel.is_cancelled() { + return Err(Error::Cancelled); + } + let want = left.min(buf.len()); + self.reader.read_exact(&mut buf[..want]).map_err(io)?; + if !on_chunk(&buf[..want]) { + return Err(Error::Cancelled); + } + left -= want; + } + // CRLF de cierre del chunk. + let mut crlf = [0u8; 2]; + self.reader.read_exact(&mut crlf).map_err(io)?; + }, + Body::Length(total) => { + let mut left = total; + while left > 0 { + if cancel.is_cancelled() { + return Err(Error::Cancelled); + } + let want = left.min(buf.len()); + let n = self.reader.read(&mut buf[..want]).map_err(io)?; + if n == 0 { + return Ok(()); + } + if !on_chunk(&buf[..n]) { + return Err(Error::Cancelled); + } + left -= n; + } + Ok(()) + } + Body::UntilClose => loop { + if cancel.is_cancelled() { + return Err(Error::Cancelled); + } + let n = self.reader.read(&mut buf).map_err(io)?; + if n == 0 { + return Ok(()); + } + if !on_chunk(&buf[..n]) { + return Err(Error::Cancelled); + } + }, + } + } +} diff --git a/crates/asist-core/src/lib.rs b/crates/asist-core/src/lib.rs new file mode 100644 index 0000000..db3ebfb --- /dev/null +++ b/crates/asist-core/src/lib.rs @@ -0,0 +1,17 @@ +//! Tipos, configuración y utilidades compartidas por todo el asistente. +//! +//! Este crate no habla con ningún motor: sólo define el vocabulario común +//! (eventos, errores, configuración, métricas) y las dos piezas de +//! infraestructura que los demás crates necesitan: un cliente HTTP mínimo +//! para hablar con los servidores locales, y el registro de herramientas. + +pub mod config; +pub mod error; +pub mod event; +pub mod http; +pub mod telemetry; +pub mod text; +pub mod tools; + +pub use error::{Error, Result}; +pub use event::{Event, TurnId}; diff --git a/crates/asist-core/src/telemetry.rs b/crates/asist-core/src/telemetry.rs new file mode 100644 index 0000000..964aaa6 --- /dev/null +++ b/crates/asist-core/src/telemetry.rs @@ -0,0 +1,255 @@ +//! Medición de latencia por turno. +//! +//! El objetivo no es un histograma bonito sino contestar a una pregunta +//! concreta cada vez que el asistente responde: *¿quién se ha comido el +//! tiempo?* Por eso se marcan los cinco instantes que separan las etapas y se +//! imprime el desglose, en lugar de un único total que no dice dónde mirar. + +use std::collections::BTreeMap; +use std::sync::Mutex; +use std::time::{Duration, Instant}; + +use crate::event::TurnId; + +/// Etapas del pipeline, en el orden en que ocurren. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum Stage { + /// Del final del habla a la transcripción definitiva. + Asr, + /// De la transcripción al primer fragmento de texto del modelo. + LlmFirstToken, + /// Del primer fragmento a la respuesta completa. + LlmRest, + /// Ejecución de herramientas. + Tools, + /// De la primera frase al primer audio recibido. + TtsFirstAudio, + /// Síntesis del resto del turno. + TtsRest, + /// Del final del habla al primer audio: lo que de verdad se percibe. + PerceivedLatency, +} + +impl Stage { + pub fn label(self) -> &'static str { + match self { + Stage::Asr => "asr.final", + Stage::LlmFirstToken => "llm.primer_token", + Stage::LlmRest => "llm.resto", + Stage::Tools => "herramientas", + Stage::TtsFirstAudio => "tts.primer_audio", + Stage::TtsRest => "tts.resto", + Stage::PerceivedLatency => "LATENCIA PERCIBIDA", + } + } +} + +/// Cronómetro de un turno. Se rellena a medida que el turno avanza y se +/// vuelca de una vez al terminar. +#[derive(Debug)] +pub struct TurnTimer { + pub turn: TurnId, + started: Instant, + marks: BTreeMap, + /// Segundos de audio hablados por el usuario, para calcular el RTF del ASR. + pub input_secs: f32, + /// Segundos de audio sintetizados, para el RTF del TTS. + pub output_secs: f32, + pub sentences: usize, + pub tool_calls: usize, +} + +impl TurnTimer { + pub fn new(turn: TurnId) -> Self { + Self { + turn, + started: Instant::now(), + marks: BTreeMap::new(), + input_secs: 0.0, + output_secs: 0.0, + sentences: 0, + tool_calls: 0, + } + } + + /// Instante de referencia del turno: el fin del habla del usuario. + pub fn started(&self) -> Instant { + self.started + } + + pub fn record(&mut self, stage: Stage, took: Duration) { + // Etapas que se repiten (una síntesis por frase) se acumulan; las que + // marcan un hito (primer audio) se quedan con la primera medida, que + // es la que describe la latencia de arranque. + match stage { + Stage::TtsFirstAudio | Stage::LlmFirstToken | Stage::PerceivedLatency => { + self.marks.entry(stage).or_insert(took); + } + _ => *self.marks.entry(stage).or_default() += took, + } + } + + /// Marca una etapa con el tiempo transcurrido desde el inicio del turno. + pub fn mark_since_start(&mut self, stage: Stage) { + let took = self.started.elapsed(); + self.record(stage, took); + } + + pub fn get(&self, stage: Stage) -> Option { + self.marks.get(&stage).copied() + } + + pub fn total(&self) -> Duration { + self.started.elapsed() + } + + /// Etapa que más ha tardado, excluyendo el total percibido (que las agrega). + pub fn bottleneck(&self) -> Option<(Stage, Duration)> { + self.marks + .iter() + .filter(|(stage, _)| **stage != Stage::PerceivedLatency) + .max_by_key(|(_, took)| **took) + .map(|(stage, took)| (*stage, *took)) + } + + /// Informe de una línea por etapa, con el cuello de botella señalado. + pub fn report(&self) -> String { + let mut out = format!("turno {} — desglose de latencia\n", self.turn); + let worst = self.bottleneck().map(|(stage, _)| stage); + for (stage, took) in &self.marks { + let flag = if Some(*stage) == worst { + " <== cuello de botella" + } else { + "" + }; + out.push_str(&format!( + " {:<20} {:>8.0} ms{}\n", + stage.label(), + took.as_secs_f64() * 1000.0, + flag + )); + } + if self.input_secs > 0.0 { + if let Some(asr) = self.get(Stage::Asr) { + out.push_str(&format!( + " {:<20} {:>8.2} x tiempo real ({:.1} s de voz)\n", + "asr.rtf", + asr.as_secs_f32() / self.input_secs, + self.input_secs + )); + } + } + if self.output_secs > 0.0 { + let tts: Duration = self.get(Stage::TtsFirstAudio).unwrap_or_default() + + self.get(Stage::TtsRest).unwrap_or_default(); + out.push_str(&format!( + " {:<20} {:>8.2} x tiempo real ({:.1} s de audio, {} frases)\n", + "tts.rtf", + tts.as_secs_f32() / self.output_secs, + self.output_secs, + self.sentences + )); + } + out.push_str(&format!( + " {:<20} {:>8.0} ms\n", + "total del turno", + self.total().as_secs_f64() * 1000.0 + )); + out + } +} + +/// Agregado de todos los turnos de la sesión, para el resumen del cierre. +#[derive(Debug, Default)] +pub struct Metrics { + inner: Mutex, +} + +#[derive(Debug, Default)] +struct MetricsInner { + turns: usize, + totals: BTreeMap, +} + +impl Metrics { + pub fn record_turn(&self, timer: &TurnTimer) { + let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner()); + inner.turns += 1; + for (stage, took) in &timer.marks { + let entry = inner.totals.entry(*stage).or_default(); + entry.0 += *took; + entry.1 += 1; + } + } + + pub fn turns(&self) -> usize { + self.inner.lock().unwrap_or_else(|e| e.into_inner()).turns + } + + /// Media por etapa sobre toda la sesión. + pub fn summary(&self) -> String { + let inner = self.inner.lock().unwrap_or_else(|e| e.into_inner()); + if inner.turns == 0 { + return "sin turnos completados\n".into(); + } + let mut out = format!("resumen de {} turno(s) — media por etapa\n", inner.turns); + for (stage, (total, count)) in &inner.totals { + out.push_str(&format!( + " {:<20} {:>8.0} ms\n", + stage.label(), + total.as_secs_f64() * 1000.0 / *count as f64 + )); + } + out + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn el_cuello_de_botella_es_la_etapa_mas_lenta() { + let mut timer = TurnTimer::new(TurnId(1)); + timer.record(Stage::Asr, Duration::from_millis(300)); + timer.record(Stage::LlmFirstToken, Duration::from_millis(500)); + timer.record(Stage::TtsFirstAudio, Duration::from_millis(900)); + assert_eq!(timer.bottleneck().unwrap().0, Stage::TtsFirstAudio); + } + + #[test] + fn la_latencia_percibida_no_compite_como_cuello_de_botella() { + let mut timer = TurnTimer::new(TurnId(1)); + timer.record(Stage::Asr, Duration::from_millis(300)); + timer.record(Stage::PerceivedLatency, Duration::from_secs(9)); + assert_eq!(timer.bottleneck().unwrap().0, Stage::Asr); + } + + #[test] + fn las_etapas_repetidas_se_acumulan_y_los_hitos_no() { + let mut timer = TurnTimer::new(TurnId(1)); + timer.record(Stage::TtsRest, Duration::from_millis(100)); + timer.record(Stage::TtsRest, Duration::from_millis(150)); + assert_eq!(timer.get(Stage::TtsRest), Some(Duration::from_millis(250))); + + timer.record(Stage::TtsFirstAudio, Duration::from_millis(600)); + timer.record(Stage::TtsFirstAudio, Duration::from_millis(50)); + assert_eq!( + timer.get(Stage::TtsFirstAudio), + Some(Duration::from_millis(600)), + "el primer audio describe el arranque; una frase posterior no debe rebajarlo" + ); + } + + #[test] + fn el_resumen_promedia_entre_turnos() { + let metrics = Metrics::default(); + for ms in [100, 300] { + let mut timer = TurnTimer::new(TurnId(1)); + timer.record(Stage::Asr, Duration::from_millis(ms)); + metrics.record_turn(&timer); + } + assert_eq!(metrics.turns(), 2); + assert!(metrics.summary().contains("200 ms")); + } +} diff --git a/crates/asist-core/src/text.rs b/crates/asist-core/src/text.rs new file mode 100644 index 0000000..e799561 --- /dev/null +++ b/crates/asist-core/src/text.rs @@ -0,0 +1,330 @@ +//! Preparación del texto que va a leerse en voz alta. +//! +//! Dos problemas distintos, ambos en el camino crítico de la latencia: +//! +//! 1. **Trocear en frases mientras el modelo escribe.** Esperar a la respuesta +//! completa antes de sintetizar suma el tiempo del LLM al del TTS. Cortando +//! por frases, el asistente empieza a hablar mientras sigue pensando. +//! 2. **Quitar lo que no se pronuncia.** El sintetizador lee los asteriscos y +//! las almohadillas tal cual, así que el markdown que se le escapa al modelo +//! hay que limpiarlo antes. + +/// Trocea un flujo de texto en frases pronunciables. +/// +/// La primera frase sale en cuanto es mínimamente decente y las siguientes +/// esperan a tener más cuerpo: el arranque manda en la latencia percibida, +/// pero una vez que suena la voz, las frases largas se entonan mejor. +#[derive(Debug)] +pub struct SentenceSplitter { + buffer: String, + emitted: usize, + /// Longitud mínima de la primera frase. + first_min: usize, + /// Longitud mínima de las siguientes. + rest_min: usize, + /// Longitud a partir de la cual se corta aunque no haya puntuación, para + /// que un modelo que no puntúa no deje al asistente mudo. + hard_max: usize, +} + +impl Default for SentenceSplitter { + fn default() -> Self { + Self { + buffer: String::new(), + emitted: 0, + first_min: 12, + rest_min: 40, + hard_max: 240, + } + } +} + +impl SentenceSplitter { + pub fn new() -> Self { + Self::default() + } + + pub fn with_limits(first_min: usize, rest_min: usize, hard_max: usize) -> Self { + Self { + first_min, + rest_min, + hard_max, + ..Self::default() + } + } + + pub fn emitted(&self) -> usize { + self.emitted + } + + /// Añade texto recién llegado y devuelve las frases que ya están completas. + pub fn push(&mut self, delta: &str) -> Vec { + self.buffer.push_str(delta); + let mut out = Vec::new(); + while let Some(sentence) = self.take_ready() { + out.push(sentence); + } + out + } + + /// Entrega lo que quede al terminar la respuesta. + pub fn flush(&mut self) -> Option { + let rest = clean_for_speech(&std::mem::take(&mut self.buffer)); + if !is_speakable(&rest) { + return None; + } + self.emitted += 1; + Some(rest) + } + + fn min_len(&self) -> usize { + if self.emitted == 0 { + self.first_min + } else { + self.rest_min + } + } + + fn take_ready(&mut self) -> Option { + let cut = self.boundary()?; + let head: String = self.buffer.drain(..cut).collect(); + let head = clean_for_speech(&head); + if !is_speakable(&head) { + // Sólo era puntuación o markdown: se descarta sin gastar un turno + // de síntesis, pero el corte ya se ha consumido. + return self.take_ready(); + } + self.emitted += 1; + Some(head) + } + + /// Índice de byte por el que cortar, si hay alguno. + fn boundary(&self) -> Option { + let min = self.min_len(); + let mut last_soft = None; + + for (i, c) in self.buffer.char_indices() { + let end = i + c.len_utf8(); + if end < min { + continue; + } + if is_terminator(c) { + // El punto de una abreviatura o de un decimal no cierra frase. + if c == '.' && !ends_sentence(&self.buffer, i) { + continue; + } + // Sólo cierra si viene un espacio detrás, o si ya no queda + // nada: si no, se estaría cortando a mitad de «3.14». + match self.buffer[end..].chars().next() { + Some(next) if next.is_whitespace() => return Some(end), + None => {} + Some(_) => continue, + } + } + // Una coma o un punto y coma valen como corte de emergencia si la + // frase ya se ha pasado de largo. + if matches!(c, ',' | ';' | ':') { + last_soft = Some(end); + } + if end >= self.hard_max { + return Some(last_soft.unwrap_or(end)); + } + } + None + } +} + +/// ¿Hay algo que pronunciar? +/// +/// Un fragmento que sólo tiene puntuación —lo que queda al limpiar un «**» o +/// unos puntos suspensivos sueltos— cuesta una síntesis entera y no suena a +/// nada, así que no llega a salir del troceador. +pub fn is_speakable(text: &str) -> bool { + text.chars().any(char::is_alphanumeric) +} + +fn is_terminator(c: char) -> bool { + matches!(c, '.' | '!' | '?' | '…' | '\n') +} + +/// ¿El punto en `idx` cierra frase de verdad? +/// +/// Descarta decimales («3.14»), elipsis a medio escribir y las abreviaturas +/// más comunes en español, que si no parten la frase justo antes del nombre. +fn ends_sentence(text: &str, idx: usize) -> bool { + let before = &text[..idx]; + let after = &text[idx + 1..]; + + if after.starts_with(|c: char| c.is_ascii_digit()) + && before.ends_with(|c: char| c.is_ascii_digit()) + { + return false; + } + if after.starts_with('.') || before.ends_with('.') { + return false; + } + let word = before + .rsplit(|c: char| c.is_whitespace()) + .next() + .unwrap_or("") + .to_lowercase(); + const ABREVIATURAS: &[&str] = &[ + "sr", "sra", "srta", "dr", "dra", "ud", "uds", "etc", "ej", "p.ej", "av", "núm", "num", + "pág", "pag", "vol", "art", "ap", "aprox", "ee.uu", "d", "dña", + ]; + !ABREVIATURAS.contains(&word.as_str()) +} + +/// Deja el texto listo para el sintetizador. +/// +/// Se limita a quitar lo que no se pronuncia; no reescribe ni resume, porque +/// lo que suena tiene que ser lo que el modelo dijo. +pub fn clean_for_speech(text: &str) -> String { + let mut out = String::with_capacity(text.len()); + let mut chars = text.chars().peekable(); + let mut at_line_start = true; + + while let Some(c) = chars.next() { + match c { + // Énfasis y código: los delimitadores se leerían en voz alta. + '*' | '_' | '`' | '~' => continue, + // Encabezados y viñetas, sólo al principio de línea: una + // almohadilla o un guion en mitad de una frase sí significan algo. + '#' if at_line_start => { + while chars.peek() == Some(&'#') { + chars.next(); + } + while chars.peek().is_some_and(|c| *c == ' ') { + chars.next(); + } + continue; + } + '-' | '•' | '–' if at_line_start && chars.peek() == Some(&' ') => { + chars.next(); + continue; + } + '>' if at_line_start => { + while chars.peek().is_some_and(|c| *c == ' ') { + chars.next(); + } + continue; + } + // Los saltos de línea se hablan como pausas. + '\n' | '\r' | '\t' => { + at_line_start = c == '\n'; + if !out.ends_with(' ') && !out.is_empty() { + out.push(' '); + } + continue; + } + _ => {} + } + at_line_start = false; + out.push(c); + } + + // Colapsa los espacios que deja la limpieza. + let mut collapsed = String::with_capacity(out.len()); + let mut space = false; + for c in out.chars() { + if c == ' ' { + space = true; + continue; + } + if space && !collapsed.is_empty() { + collapsed.push(' '); + } + space = false; + collapsed.push(c); + } + collapsed.trim().to_string() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn split_all(chunks: &[&str]) -> Vec { + let mut splitter = SentenceSplitter::new(); + let mut out: Vec = chunks.iter().flat_map(|c| splitter.push(c)).collect(); + out.extend(splitter.flush()); + out + } + + #[test] + fn la_primera_frase_sale_antes_que_las_siguientes() { + // La primera basta con que sea corta; la segunda espera a tener cuerpo. + let out = split_all(&["Claro que sí. ", "Vale. ", "Aún no.", ""]); + assert_eq!(out[0], "Claro que sí."); + assert_eq!(out[1], "Vale. Aún no."); + } + + #[test] + fn trocea_segun_llegan_los_deltas() { + let mut splitter = SentenceSplitter::new(); + assert!(splitter.push("Hola, ¿qué ").is_empty()); + assert!(splitter.push("tal est").is_empty()); + let out = splitter.push("ás hoy? "); + assert_eq!(out, vec!["Hola, ¿qué tal estás hoy?"]); + } + + #[test] + fn el_punto_decimal_no_parte_la_frase() { + let out = split_all(&["El resultado es 3.1416 exactamente y nada más."]); + assert_eq!(out.len(), 1, "un decimal no cierra frase: {out:?}"); + } + + #[test] + fn las_abreviaturas_no_parten_la_frase() { + let out = split_all(&["Avisa al Sr. Pérez cuanto antes por favor."]); + assert_eq!(out.len(), 1, "«Sr.» no cierra frase: {out:?}"); + } + + #[test] + fn una_parrafada_sin_puntuacion_se_corta_igualmente() { + // Sin este corte de emergencia, un modelo que no puntúa deja al + // asistente sin sintetizar nada hasta el final de la respuesta. + let largo = "palabra ".repeat(60); + let out = split_all(&[&largo]); + assert!(out.len() > 1, "esperaba varios cortes, hubo {}", out.len()); + assert!(out.iter().all(|s| s.len() <= 260)); + } + + #[test] + fn el_corte_de_emergencia_prefiere_una_coma() { + let mut splitter = SentenceSplitter::with_limits(4, 4, 40); + let out = splitter.push("uno dos tres, cuatro cinco seis siete ocho nueve diez once "); + assert_eq!(out[0], "uno dos tres,"); + } + + #[test] + fn se_limpia_el_markdown_que_se_leeria_en_voz_alta() { + assert_eq!(clean_for_speech("## Título"), "Título"); + assert_eq!( + clean_for_speech("Esto es **muy** importante"), + "Esto es muy importante" + ); + assert_eq!(clean_for_speech("- uno\n- dos"), "uno dos"); + assert_eq!(clean_for_speech("usa `ls -la` ahora"), "usa ls -la ahora"); + assert_eq!(clean_for_speech("> citado"), "citado"); + } + + #[test] + fn un_guion_dentro_de_la_frase_sobrevive() { + assert_eq!(clean_for_speech("teórico-práctico"), "teórico-práctico"); + } + + #[test] + fn los_fragmentos_solo_de_puntuacion_no_generan_sintesis() { + let out = split_all(&["**", "...", " "]); + assert!(out.is_empty(), "no había nada que pronunciar: {out:?}"); + } + + #[test] + fn el_flush_entrega_la_cola_sin_puntuacion_final() { + let mut splitter = SentenceSplitter::new(); + splitter.push("Una frase sin punto final"); + assert_eq!(splitter.flush().unwrap(), "Una frase sin punto final"); + assert!(splitter.flush().is_none()); + } +} diff --git a/crates/asist-core/src/tools.rs b/crates/asist-core/src/tools.rs new file mode 100644 index 0000000..0b8a646 --- /dev/null +++ b/crates/asist-core/src/tools.rs @@ -0,0 +1,563 @@ +//! Herramientas que el modelo puede invocar. +//! +//! El punto de extensión del asistente. Una herramienta es un objeto que se +//! describe a sí mismo en JSON Schema y sabe ejecutarse; el registro las +//! traduce al formato `tools` de la API de chat y despacha las llamadas que +//! vuelven. Añadir una capacidad nueva es implementar el trait y registrarla, +//! sin tocar el orquestador. + +use std::collections::BTreeMap; +use std::sync::Arc; +use std::time::Duration; + +use serde_json::{json, Value}; + +use crate::config::ToolsConfig; +use crate::error::{Error, Result}; + +/// Lo que el modelo puede pedir que se haga. +pub trait Tool: Send + Sync { + /// Identificador que usa el modelo. En minúsculas y sin espacios. + fn name(&self) -> &str; + + /// Para qué sirve. Lo lee el modelo, así que se escribe para él: concreto + /// y en la lengua de la conversación. + fn description(&self) -> &str; + + /// JSON Schema de los argumentos. + fn parameters(&self) -> Value; + + /// Ejecuta y devuelve el texto que se le entrega al modelo como resultado. + fn call(&self, args: &Value) -> Result; + + /// `true` si la herramienta cambia algo fuera del proceso. El orquestador + /// lo anuncia en voz alta antes de ejecutarla. + fn is_side_effecting(&self) -> bool { + false + } +} + +/// Una llamada tal y como la pide el modelo. +#[derive(Debug, Clone)] +pub struct ToolCall { + pub id: String, + pub name: String, + /// Argumentos en JSON, aún sin validar. + pub arguments: String, +} + +/// Resultado de ejecutarla. +#[derive(Debug, Clone)] +pub struct ToolOutcome { + pub id: String, + pub name: String, + pub ok: bool, + pub output: String, + pub took: Duration,