aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-app/src/pipeline.rs
diff options
context:
space:
mode:
authorelvis <elvis@claros.ar>2026-09-06 19:21:16 -0300
committerelvis <elvis@claros.ar>2026-09-06 19:23:37 -0300
commitf8f98e83481e7376a235fb305a095552433f79ab (patch)
tree6d0ecddb26bef8430a086431de40f8ce9b1039e6 /crates/asist-app/src/pipeline.rs
downloadasist-p-f8f98e83481e7376a235fb305a095552433f79ab.tar.gz
asist-p-f8f98e83481e7376a235fb305a095552433f79ab.zip
Local voice assistant on top of Canary, llama.cpp and qwentts
Pipeline en Rust de hilos y canales que une los tres motores: el micrófono alimenta un segmentador con VAD, las intervenciones cerradas van al reconocedor, la transcripción al modelo y cada frase que este cierra sale hacia el sintetizador sin esperar al resto de la respuesta. Dos reglas sostienen el diseño: ninguna etapa bloquea a la anterior —quien va sobrado descarta trabajo en lugar de acumular retraso— y todo lo que viaja por los canales lleva el turno al que pertenece, así que interrumpir es subir el contador y levantar dos banderas de cancelación. Seis crates: core (configuración, eventos, HTTP, telemetría, herramientas), audio (cpal, VAD, anillo de reproducción), asr, llm, tts y app (supervisor de procesos y orquestador). Los motores van como submódulos fijados a un commit, con los cambios locales en vendor/patches. Midiendo el pipeline aparecieron tres cuellos de botella de configuración que valieron más que cualquier cambio de código, todos documentados en docs/RENDIMIENTO.md: - tts-server decodificaba el audio en bloques de 24 s, de modo que el modo «streaming» llegaba de una pieza: 4948 ms -> 585 ms hasta el primer audio. - La plantilla de chat del modelo abre <think> y no lo cierra nunca, sin variable que lo apague: 8630 ms -> 413 ms hasta el primer token, con una copia de la plantilla que deja el bloque cerrado de entrada. - Cualquier indicación de estilo junto a la guía de herramientas hace que este modelo de 2B deje de llamarlas y se invente el dato (8/8 aciertos con la guía sola, 0/8 con la persona de asistente de voz). El turno alterna ahora entre dos instrucciones de sistema. La ejecución de órdenes del sistema queda implementada y apagada, tras cuatro barreras: lista blanca sobre el ejecutable, rutas rechazadas, sin shell que interprete metacaracteres y plazo máximo. 68 pruebas unitarias sin modelos, más seis de integración que se saltan solas si no hay servidores y se turnan la GPU: en paralelo, los dos servidores no caben en 4 GB y miden contención en vez de latencia. Claude-Session: https://claude.ai/code/session_01FNxz5cSdQSscJH9H7b8uGU
Diffstat (limited to 'crates/asist-app/src/pipeline.rs')
-rw-r--r--crates/asist-app/src/pipeline.rs634
1 files changed, 634 insertions, 0 deletions
diff --git a/crates/asist-app/src/pipeline.rs b/crates/asist-app/src/pipeline.rs
new file mode 100644
index 0000000..9e9037e
--- /dev/null
+++ b/crates/asist-app/src/pipeline.rs
@@ -0,0 +1,634 @@
+//! El orquestador: los hilos del asistente y los canales que los unen.
+//!
+//! ```text
+//! micrófono ──muestras──> segmentador ──intervención──> ASR ──texto──┐
+//! (cpal, tiempo real) (VAD, turnos) (Canary) │
+//! v
+//! altavoz <──muestras── síntesis <──frases── conversación <──────────┘
+//! (cpal, anillo) (qwentts) (llama.cpp + herramientas)
+//! ```
+//!
+//! Cada caja es un hilo y cada flecha un canal. Dos consecuencias que valen
+//! por todo el diseño:
+//!
+//! * **Nada bloquea al que va delante.** El micrófono nunca espera al ASR, y
+//! el modelo nunca espera al sintetizador; quien va sobrado descarta trabajo
+//! en lugar de acumular retraso.
+//! * **Se responde por frases, no por respuestas.** La conversación entrega
+//! cada frase al sintetizador en cuanto está cerrada, así que el asistente
+//! empieza a hablar mientras el modelo sigue escribiendo. Es lo que separa
+//! una latencia de medio segundo de una de cinco.
+
+use std::sync::Arc;
+use std::thread::{self, JoinHandle};
+use std::time::{Duration, Instant};
+
+use crossbeam_channel::{bounded, unbounded, Receiver, Select, Sender};
+
+use asist_asr::{AsrJob, AsrResult, Recognizer};
+use asist_audio::{
+ CaptureBlock, Gate, InputFormat, PlaybackHandle, Segmenter, VoiceEvent, ASR_SAMPLE_RATE,
+};
+use asist_core::config::Config;
+use asist_core::error::Result;
+use asist_core::event::{Event, InterruptReason, TurnId};
+use asist_core::telemetry::{Metrics, Stage, TurnTimer};
+use asist_core::text::SentenceSplitter;
+use asist_core::tools::ToolRegistry;
+use asist_llm::chat::Conversation;
+use asist_llm::{Delta, LlmClient, Message};
+use asist_tts::TtsClient;
+
+use crate::session::Session;
+
+/// Una frase pendiente de sintetizar.
+#[derive(Debug)]
+pub struct SpeakJob {
+ pub turn: TurnId,
+ pub index: usize,
+ pub text: String,
+ /// Origen desde el que se mide la latencia percibida del turno.
+ pub turn_started: Instant,
+}
+
+/// Los hilos en marcha, para poder esperarlos al cerrar.
+pub struct Pipeline {
+ handles: Vec<JoinHandle<()>>,
+ pub events: Receiver<Event>,
+ pub metrics: Arc<Metrics>,
+ session: Session,
+ capture_tx: Option<Sender<CaptureBlock>>,
+}
+
+impl Pipeline {
+ /// Monta el pipeline entero. `capture_rx` viene del micrófono.
+ #[allow(clippy::too_many_arguments)]
+ pub fn spawn(
+ config: &Config,
+ session: Session,
+ recognizer: Recognizer,
+ llm: LlmClient,
+ tts: TtsClient,
+ tools: ToolRegistry,
+ playback: PlaybackHandle,
+ capture_rx: Receiver<CaptureBlock>,
+ capture_tx: Sender<CaptureBlock>,
+ input_format: InputFormat,
+ ) -> Result<Self> {
+ let (event_tx, event_rx) = unbounded::<Event>();
+ let (asr_tx, asr_rx) = unbounded::<AsrJob>();
+ let (result_tx, result_rx) = unbounded::<AsrResult>();
+ // Acotado a propósito: si la síntesis se atasca, el hilo de
+ // conversación debe notarlo y frenar en lugar de acumular frases que
+ // llegarán tarde a un turno que quizá ya se ha interrumpido.
+ let (speak_tx, speak_rx) = bounded::<SpeakJob>(8);
+
+ let metrics = Arc::new(Metrics::default());
+ let mut handles = Vec::new();
+
+ handles.push(spawn_named("segmentador", {
+ let config = config.clone();
+ let session = session.clone();
+ let events = event_tx.clone();
+ let playback = playback.clone();
+ move || {
+ run_segmenter(
+ &config,
+ session,
+ capture_rx,
+ asr_tx,
+ events,
+ playback,
+ input_format,
+ )
+ }
+ }));
+
+ handles.push(spawn_named("asr", move || {
+ recognizer.run(asr_rx, result_tx)
+ }));
+
+ handles.push(spawn_named("conversación", {
+ let config = config.clone();
+ let session = session.clone();
+ let events = event_tx.clone();
+ let metrics = Arc::clone(&metrics);
+ move || {
+ run_brain(
+ &config, session, llm, tools, result_rx, speak_tx, events, metrics,
+ )
+ }
+ }));
+
+ handles.push(spawn_named("síntesis", {
+ let session = session.clone();
+ let events = event_tx;
+ move || run_speech(session, tts, speak_rx, playback, events)
+ }));
+
+ Ok(Self {
+ handles,
+ events: event_rx,
+ metrics,
+ session,
+ capture_tx: Some(capture_tx),
+ })
+ }
+
+ /// Cierra en orden: se corta lo que suena, se suelta el emisor del
+ /// micrófono y con eso la cadena de hilos se desmonta sola.
+ pub fn shutdown(mut self) {
+ self.session.request_stop();
+ self.capture_tx.take();
+ for handle in self.handles.drain(..) {
+ let _ = handle.join();
+ }
+ }
+}
+
+fn spawn_named(name: &str, body: impl FnOnce() + Send + 'static) -> JoinHandle<()> {
+ thread::Builder::new()
+ .name(name.to_string())
+ .spawn(body)
+ .expect("no se pudo crear un hilo del pipeline")
+}
+
+// ---------------------------------------------------------------------------
+// Segmentador: micrófono -> intervenciones
+// ---------------------------------------------------------------------------
+
+fn run_segmenter(
+ config: &Config,
+ session: Session,
+ capture: Receiver<CaptureBlock>,
+ asr: Sender<AsrJob>,
+ events: Sender<Event>,
+ playback: PlaybackHandle,
+ input_format: InputFormat,
+) {
+ let mut segmenter = Segmenter::new(&config.vad);
+ let mut turn = session.current();
+
+ while let Ok(block) = capture.recv() {
+ if session.is_stopping() {
+ break;
+ }
+ // El micrófono se cierra mientras suena el altavoz. Con barge-in
+ // activo no se cierra, sólo sube el listón de volumen.
+ segmenter.set_gate(if session.is_speaking() || playback.is_active() {
+ Gate::Speaking
+ } else {
+ Gate::Open
+ });
+
+ // El dispositivo se abre a 16 kHz mono siempre que puede, y entonces
+ // esto no hace nada; cuando no puede, es aquí donde se convierte.
+ let mono = asist_audio::to_mono_at(&block.samples, input_format, ASR_SAMPLE_RATE);
+
+ for event in segmenter.push(&mono) {
+ match event {
+ VoiceEvent::BargeIn => {
+ // Se corta antes de abrir el turno nuevo: así lo que
+ // quedaba en el anillo no se cuela por encima.
+ playback.stop();
+ session.interrupt();
+ let _ = events.send(Event::Interrupted {
+ turn: session.current(),
+ reason: InterruptReason::UserSpoke,
+ });
+ }
+ VoiceEvent::Started => {
+ turn = session.begin_turn();
+ let _ = events.send(Event::SpeechStarted { turn, at: block.at });
+ }
+ VoiceEvent::Audio(samples) => {
+ if config.asr.partials
+ && asr
+ .send(AsrJob::Window {
+ turn,
+ samples,
+ at: block.at,
+ })
+ .is_err()
+ {
+ return;
+ }
+ }
+ VoiceEvent::Ended(utterance) => {
+ let at = Instant::now();
+ if asr
+ .send(AsrJob::Utterance {
+ turn,
+ samples: utterance.samples,
+ at,
+ })
+ .is_err()
+ {
+ return;
+ }
+ }
+ VoiceEvent::Discarded => {
+ if asr.send(AsrJob::Reset { turn }).is_err() {
+ return;
+ }
+ }
+ }
+ }
+ }
+ // Lo que quedara a medias al cerrar todavía merece transcribirse.
+ if let Some(utterance) = segmenter.flush() {
+ let _ = asr.send(AsrJob::Utterance {
+ turn,
+ samples: utterance.samples,
+ at: Instant::now(),
+ });
+ }
+}
+
+// ---------------------------------------------------------------------------
+// Conversación: transcripción -> frases
+// ---------------------------------------------------------------------------
+
+#[allow(clippy::too_many_arguments)]
+fn run_brain(
+ config: &Config,
+ session: Session,
+ llm: LlmClient,
+ tools: ToolRegistry,
+ results: Receiver<AsrResult>,
+ speak: Sender<SpeakJob>,
+ events: Sender<Event>,
+ metrics: Arc<Metrics>,
+) {
+ let prompts = Prompts::new(config, &tools);
+ let mut chat = Conversation::new(prompts.deciding.clone(), config.general.history_turns);
+
+ while let Ok(result) = results.recv() {
+ if session.is_stopping() {
+ break;
+ }
+ match result {
+ AsrResult::Partial {
+ turn,
+ committed,
+ volatile,
+ ..
+ } => {
+ let _ = events.send(Event::Partial {
+ turn,
+ committed,
+ volatile,
+ });
+ }
+ AsrResult::Empty { turn } => {
+ let _ = events.send(Event::Discarded { turn });
+ }
+ AsrResult::Error { turn, message } => {
+ let _ = events.send(Event::Warning { turn, message });
+ }
+ AsrResult::Final {
+ turn,
+ text,
+ audio_secs,
+ decode,
+ spoken_at,
+ } => {
+ // Una transcripción de un turno superado llega tarde: el
+ // usuario ya ha dicho otra cosa.
+ if !session.is_current(turn) {
+ tracing::debug!(target: "brain", turno = turn.0, "transcripción obsoleta");
+ continue;
+ }
+ let mut timer = TurnTimer::new(turn);
+ timer.input_secs = audio_secs;
+ timer.record(Stage::Asr, decode);
+
+ let _ = events.send(Event::Transcript {
+ turn,
+ text: text.clone(),
+ audio_secs,
+ decode,
+ });
+
+ chat.push(Message::user(text));
+ answer(
+ config, &session, &llm, &tools, &prompts, &mut chat, &speak, &events,
+ &mut timer, spoken_at,
+ );
+
+ if config.general.report_latency {
+ tracing::info!(target: "latencia", "\n{}", timer.report());
+ }
+ metrics.record_turn(&timer);
+ }
+ }
+ }
+}
+
+/// Las dos instrucciones de sistema entre las que alterna un turno.
+///
+/// Que sean dos no es un capricho de diseño sino un resultado medido: con
+/// Qwen3.5-2B, cualquier indicación de estilo junto a la guía de herramientas
+/// hace que el modelo deje de llamarlas y se invente el dato (8/8 aciertos con
+/// la guía sola, 0/8 con la persona de asistente de voz añadida). Así que la
+/// pasada en que el modelo *decide* si actuar lleva la guía a solas, y la
+/// instrucción de voz se reserva para redactar lo que se va a pronunciar.
+struct Prompts {
+ /// Para la pasada en que el modelo decide si llamar a una herramienta.
+ deciding: String,
+ /// Para redactar la respuesta hablada, ya con los resultados en la mano.
+ speaking: String,
+ /// `false` cuando no hay herramientas: entonces ambas son la misma.
+ two_phase: bool,
+}
+
+impl Prompts {
+ fn new(config: &Config, tools: &ToolRegistry) -> Self {
+ let speaking = config.general.system_prompt.trim().to_string();
+ let two_phase = config.tools.enabled && config.tools.dedicated_prompt && !tools.is_empty();
+ Self {
+ deciding: if two_phase {
+ config.general.tools_prompt.trim().to_string()
+ } else {
+ speaking.clone()
+ },
+ speaking,
+ two_phase,
+ }
+ }
+}
+
+/// Genera la respuesta de un turno, resolviendo herramientas si el modelo las
+/// pide, y va soltando frases al sintetizador según se cierran.
+#[allow(clippy::too_many_arguments)]
+fn answer(
+ config: &Config,
+ session: &Session,
+ llm: &LlmClient,
+ tools: &ToolRegistry,
+ prompts: &Prompts,
+ chat: &mut Conversation,
+ speak: &Sender<SpeakJob>,
+ events: &Sender<Event>,
+ timer: &mut TurnTimer,
+ turn_started: Instant,
+) {
+ let turn = timer.turn;
+ let tools = (!tools.is_empty() && config.tools.enabled).then_some(tools);
+ let mut splitter = SentenceSplitter::new();
+ let mut spoken = String::new();
+
+ for round in 0..=config.llm.max_tool_rounds {
+ let llm_started = Instant::now();
+ let mut first_token = None;
+
+ let outcome = llm.stream(chat, tools, session.llm_cancel(), |delta| {
+ if !session.is_current(turn) {
+ return false;
+ }
+ match delta {
+ Delta::Text(text) => {
+ if first_token.is_none() {
+ first_token = Some(llm_started.elapsed());
+ let _ = events.send(Event::ReplyStarted {
+ turn,
+ ttft: llm_started.elapsed(),
+ });
+ }
+ let _ = events.send(Event::ReplyDelta {
+ turn,
+ text: text.clone(),
+ });
+ spoken.push_str(text);
+ // Aquí está el solape: cada frase cerrada sale hacia el
+ // sintetizador sin esperar al resto de la respuesta.
+ for sentence in splitter.push(text) {
+ let index = splitter.emitted() - 1;
+ let _ = events.send(Event::Sentence {
+ turn,
+ index,
+ text: sentence.clone(),
+ });
+ if speak
+ .send(SpeakJob {
+ turn,
+ index,
+ text: sentence,
+ turn_started,
+ })
+ .is_err()
+ {
+ return false;
+ }
+ }
+ true
+ }
+ Delta::ToolCalls(_) => true,
+ }
+ });
+
+ if let Some(ttft) = first_token {
+ timer.record(Stage::LlmFirstToken, ttft);
+ timer.record(Stage::LlmRest, llm_started.elapsed().saturating_sub(ttft));
+ }
+
+ let outcome = match outcome {
+ Ok(outcome) => outcome,
+ Err(err) => {
+ let _ = events.send(Event::Failed {
+ turn,
+ message: err.to_string(),
+ });
+ return;
+ }
+ };
+
+ if !session.is_current(turn) {
+ return;
+ }
+
+ if outcome.tool_calls.is_empty() {
+ if let Some(rest) = splitter.flush() {
+ let index = splitter.emitted() - 1;
+ let _ = events.send(Event::Sentence {
+ turn,
+ index,
+ text: rest.clone(),
+ });
+ let _ = speak.send(SpeakJob {
+ turn,
+ index,
+ text: rest,
+ turn_started,
+ });
+ }
+ timer.sentences = splitter.emitted();
+ chat.push(Message::assistant(outcome.text.clone()));
+ // El turno siguiente vuelve a empezar decidiendo.
+ if prompts.two_phase {
+ chat.set_system(&prompts.deciding);
+ }
+ let _ = events.send(Event::ReplyDone {
+ turn,
+ text: outcome.text,
+ });
+ return;
+ }
+
+ // El modelo ha pedido herramientas: se ejecutan, se le devuelve el
+ // resultado y se le vuelve a preguntar.
+ if round == config.llm.max_tool_rounds {
+ let _ = events.send(Event::Warning {
+ turn,
+ message: format!(
+ "el modelo siguió pidiendo herramientas tras {} vueltas; se corta",
+ config.llm.max_tool_rounds
+ ),
+ });
+ return;
+ }
+
+ let tools = tools.expect("no puede haber llamadas sin herramientas registradas");
+ chat.push(Message::tool_request(
+ outcome.text.clone(),
+ outcome.tool_calls.clone(),
+ ));
+ let tools_started = Instant::now();
+ for call in &outcome.tool_calls {
+ let _ = events.send(Event::ToolRequested {
+ turn,
+ name: call.name.clone(),
+ arguments: call.arguments.clone(),
+ });
+ let result = tools.dispatch(call);
+ let _ = events.send(Event::ToolFinished {
+ turn,
+ name: result.name.clone(),
+ ok: result.ok,
+ output: result.output.clone(),
+ took: result.took,
+ });
+ chat.push(Message::tool_result(&result));
+ timer.tool_calls += 1;
+ }
+ timer.record(Stage::Tools, tools_started.elapsed());
+
+ // Con el resultado ya en la conversación, la siguiente vuelta sólo
+ // tiene que redactar: es el momento de recuperar la instrucción de
+ // voz, que en la pasada anterior habría impedido la llamada.
+ if prompts.two_phase {
+ chat.set_system(&prompts.speaking);
+ }
+ }
+}
+
+// ---------------------------------------------------------------------------
+// Síntesis: frases -> altavoz
+// ---------------------------------------------------------------------------
+
+fn run_speech(
+ session: Session,
+ tts: TtsClient,
+ jobs: Receiver<SpeakJob>,
+ playback: PlaybackHandle,
+ events: Sender<Event>,
+) {
+ // 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<TurnId> = 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);
+}