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/pipeline.rs | 634 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 634 insertions(+) create mode 100644 crates/asist-app/src/pipeline.rs (limited to 'crates/asist-app/src/pipeline.rs') 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); +} -- cgit v1.2.3