//! 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); }