From 34870fed2715675c95e4d6ed67521cf2ecec99b2 Mon Sep 17 00:00:00 2001 From: mARTin-B78 Date: Fri, 29 May 2026 19:06:07 +0200 Subject: [PATCH] Pipeline TTS with LLM streaming for lower latency MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace serial LLM-wait-TTS with overlapped execution: - LLM streams via background thread → asyncio.Queue (non-blocking event loop) - _sentence_split() detects sentence boundaries in the token stream - asyncio.create_task fires TTS for each sentence immediately — TTS for sentence 1 runs while LLM is still generating sentences 2, 3, … - Audio chunks stream to frontend in order as each task completes - Time-to-first-audio drops from (LLM total + TTS total) to roughly (LLM time-to-first-sentence + TTS latency for one sentence) Frontend audio queue: - enqueueAudio() / playNextAudio() chain multi-chunk responses seamlessly - clearAudio() stops playback and cancels queue on new turn or mic click - scheduleAutoMic() waits for queue to drain before restarting mic - Error paths clear the queue to avoid stale audio playing after failure Also fix missing contextlib import (silent bug when audio temp files needed cleanup in the STT path). Co-Authored-By: Claude Sonnet 4.6 --- routes/conversation.py | 165 +++++++++++++++++++++++++++----------- static/js/conversation.js | 90 +++++++++++++-------- 2 files changed, 174 insertions(+), 81 deletions(-) diff --git a/routes/conversation.py b/routes/conversation.py index ae30330..3f68bdf 100644 --- a/routes/conversation.py +++ b/routes/conversation.py @@ -3,8 +3,11 @@ from __future__ import annotations import asyncio import base64 +import contextlib import io import json +import re +import threading import time import uuid import wave @@ -30,6 +33,20 @@ from routes.stt import _transcribe_audio, _clean_stt_backend router = APIRouter() +# ── Sentence-boundary helpers for pipelined TTS ─────────────────────────────── + +_SENT_RE = re.compile(r'(?<=[.!?])\s+') +_MIN_SENTENCE = 30 # min chars in buffer before we split + +def _sentence_split(buf: str) -> int: + """Return the index after the first sentence boundary, or -1.""" + if len(buf) < _MIN_SENTENCE: + return -1 + for m in _SENT_RE.finditer(buf): + if m.end() >= _MIN_SENTENCE: + return m.end() + return -1 + # ── LLM helpers ─────────────────────────────────────────────────────────────── def _rewrite_with_persona_sync(text: str, persona: str, llm_url: str, model: str = "") -> str: @@ -736,52 +753,100 @@ async def conversation_turn( yield sse({"type": "error", "stage": "stt", "message": "No speech detected."}) return - # 2. LLM stream + # 2. LLM stream — runs in a thread; tokens arrive via asyncio.Queue + # so the event loop is never blocked and TTS can start on sentence 1 + # while the LLM is still generating sentences 2, 3, … messages = [{"role": "system", "content": system_prompt}] messages.extend(hist[-20:]) messages.append({"role": "user", "content": transcript}) llm_payload: dict = {"messages": messages, "stream": True, "temperature": 0.7, "max_tokens": 512} if llm_model: llm_payload["model"] = llm_model + + _loop = asyncio.get_event_loop() + _token_q: asyncio.Queue[str | None] = asyncio.Queue() + + def _llm_thread() -> None: + try: + resp = requests.post( + f"{eff_llm_url}/chat/completions", json=llm_payload, + headers={"Authorization": "Bearer no-key"}, stream=True, timeout=120, + ) + resp.raise_for_status() + for raw_line in resp.iter_lines(): + if not raw_line: + continue + line = raw_line.decode("utf-8") if isinstance(raw_line, bytes) else str(raw_line) + if not line.startswith("data:"): + continue + payload_str = line[5:].strip() + if payload_str == "[DONE]": + break + try: + obj = json.loads(payload_str) + d_obj = ((obj.get("choices") or [{}])[0].get("delta") or {}) + delta = d_obj.get("content") or d_obj.get("reasoning_content") or "" + if not delta and isinstance(obj.get("message"), dict): + delta = obj["message"].get("content") or "" + if delta: + _loop.call_soon_threadsafe(_token_q.put_nowait, delta) + except Exception: + continue + except Exception as exc: + _loop.call_soon_threadsafe(_token_q.put_nowait, f"\x00ERR:{exc}") + finally: + _loop.call_soon_threadsafe(_token_q.put_nowait, None) + + threading.Thread(target=_llm_thread, daemon=True).start() + + t_llm = time.monotonic() + llm_ttft_ms: int | None = None + ttft_done = False + sent_buf = "" + tts_tasks: list[asyncio.Task] = [] + tts_first_start: float | None = None + try: - t_llm = time.monotonic() - llm_resp = await asyncio.to_thread(lambda: requests.post( - f"{eff_llm_url}/chat/completions", json=llm_payload, - headers={"Authorization": "Bearer no-key"}, stream=True, timeout=120, - )) - llm_resp.raise_for_status() - ttft_done = False - for raw_line in llm_resp.iter_lines(): - if not raw_line: - continue - line = raw_line.decode("utf-8") if isinstance(raw_line, bytes) else str(raw_line) - if not line.startswith("data:"): - continue - chunk = line[5:].strip() - if chunk == "[DONE]": + while True: + delta = await _token_q.get() + if delta is None: break - try: - obj = json.loads(chunk) - d_obj = ((obj.get("choices") or [{}])[0].get("delta") or {}) - # Primary: visible content. Fallback: reasoning_content (Qwen3 /think tokens) - delta = d_obj.get("content") or d_obj.get("reasoning_content") or "" - if not delta and isinstance(obj.get("message"), dict): - delta = obj["message"].get("content") or "" - if delta: - if not ttft_done: - llm_ttft_ms = int((time.monotonic() - t_llm) * 1000) - ttft_done = True - llm_text += delta - yield sse({"type": "token", "delta": delta}) - except Exception: - continue - llm_total_ms = int((time.monotonic() - t_llm) * 1000) - yield sse({"type": "llm_done", "text": llm_text, - "llm_ttft_ms": llm_ttft_ms, "llm_total_ms": llm_total_ms}) - except Exception as e: - yield sse({"type": "error", "stage": "llm", "message": str(e)}) + if delta.startswith("\x00ERR:"): + yield sse({"type": "error", "stage": "llm", "message": delta[5:]}) + return + if not ttft_done: + llm_ttft_ms = int((time.monotonic() - t_llm) * 1000) + ttft_done = True + llm_text += delta + sent_buf += delta + yield sse({"type": "token", "delta": delta}) + # Fire TTS on sentence boundary — runs concurrently with LLM + split = _sentence_split(sent_buf) + if split > 0: + chunk_text = sent_buf[:split].strip() + sent_buf = sent_buf[split:] + if chunk_text: + if tts_first_start is None: + tts_first_start = time.monotonic() + tts_tasks.append(asyncio.create_task( + asyncio.to_thread(_preview_request_audio, chunk_text, tts_voice, settings, "", tts_be) + )) + except Exception as exc: + yield sse({"type": "error", "stage": "llm", "message": str(exc)}) return + # Flush any remaining text as a final TTS task + if sent_buf.strip(): + if tts_first_start is None: + tts_first_start = time.monotonic() + tts_tasks.append(asyncio.create_task( + asyncio.to_thread(_preview_request_audio, sent_buf.strip(), tts_voice, settings, "", tts_be) + )) + + llm_total_ms = int((time.monotonic() - t_llm) * 1000) + yield sse({"type": "llm_done", "text": llm_text, + "llm_ttft_ms": llm_ttft_ms, "llm_total_ms": llm_total_ms}) + if not llm_text.strip(): yield sse({"type": "error", "stage": "llm", "message": "LLM returned empty response. " @@ -789,21 +854,23 @@ async def conversation_turn( "or pick a non-thinking model in the Language Model dropdown."}) return - # 3. TTS - try: - t_tts = time.monotonic() - audio_bytes, mime = await asyncio.to_thread( - _preview_request_audio, llm_text, tts_voice, settings, "", tts_be - ) - tts_ms = int((time.monotonic() - t_tts) * 1000) - total_ms = int((time.monotonic() - t0) * 1000) + # 3. Stream audio chunks in order — each chunk's TTS ran concurrently + # with LLM generation, so first audio arrives much sooner than + # waiting for the full response. + tts_start = tts_first_start or time.monotonic() + for i, task in enumerate(tts_tasks): + try: + audio_bytes, mime = await task + except Exception as exc: + yield sse({"type": "error", "stage": "tts", "message": str(exc)}) + return + if i == 0: + tts_ms = int((time.monotonic() - tts_start) * 1000) yield sse({"type": "audio", "b64": base64.b64encode(audio_bytes).decode(), "mime": mime}) - yield sse({"type": "stats", "stt_ms": stt_ms, "llm_ttft_ms": llm_ttft_ms, - "llm_total_ms": llm_total_ms, "tts_ms": tts_ms, "total_ms": total_ms}) - except Exception as e: - yield sse({"type": "error", "stage": "tts", "message": str(e)}) - return + total_ms = int((time.monotonic() - t0) * 1000) + yield sse({"type": "stats", "stt_ms": stt_ms, "llm_ttft_ms": llm_ttft_ms, + "llm_total_ms": llm_total_ms, "tts_ms": tts_ms, "total_ms": total_ms}) yield sse({"type": "done"}) return StreamingResponse(generate(), media_type="text/event-stream", diff --git a/static/js/conversation.js b/static/js/conversation.js index 276faef..ed28eec 100644 --- a/static/js/conversation.js +++ b/static/js/conversation.js @@ -186,6 +186,9 @@ $('s-import-voices-file')?.addEventListener('change', async function () { let previewTranscribing = false; let convCurrentAudio = null; let autoMicGeneration = 0; + const audioQueue = []; + let audioQueuePlaying = false; + let audioQueueDrainCb = null; // ── Populate STT backends ──────────────────────────────────────────────── async function loadConvSttBackends() { @@ -376,6 +379,45 @@ $('s-import-voices-file')?.addEventListener('change', async function () { if (micTimer) micTimer.textContent = ''; } + // ── Audio queue — plays multi-chunk TTS responses sequentially ────────────── + function clearAudio() { + if (convCurrentAudio) { try { convCurrentAudio.pause(); } catch(_){} convCurrentAudio = null; } + audioQueue.length = 0; + audioQueuePlaying = false; + audioQueueDrainCb = null; + } + + function playNextAudio() { + if (!audioQueue.length) { + audioQueuePlaying = false; + convCurrentAudio = null; + if (audioQueueDrainCb) { + const cb = audioQueueDrainCb; + audioQueueDrainCb = null; + setTimeout(cb, 150); + } else if (micStatus && !isProcessing) { + micStatus.textContent = 'Ready'; + } + return; + } + audioQueuePlaying = true; + const url = audioQueue.shift(); + const el = new Audio(url); + convCurrentAudio = el; + el.addEventListener('ended', () => { URL.revokeObjectURL(url); playNextAudio(); }); + el.play().catch(() => { URL.revokeObjectURL(url); playNextAudio(); }); + if (micStatus) micStatus.textContent = 'Speaking…'; + } + + function enqueueAudio(b64, mime) { + const binStr = atob(b64); + const arr = new Uint8Array(binStr.length); + for (let i = 0; i < binStr.length; i++) arr[i] = binStr.charCodeAt(i); + const url = URL.createObjectURL(new Blob([arr], { type: mime || 'audio/wav' })); + audioQueue.push(url); + if (!audioQueuePlaying) playNextAudio(); + } + // ── Chunked Whisper preview (fallback for browsers without SpeechRecognition) ── async function transcribeForPreview() { if (previewTranscribing || !recChunks.length) return; @@ -408,8 +450,11 @@ $('s-import-voices-file')?.addEventListener('change', async function () { if (gen !== autoMicGeneration || isProcessing) return; startRecording().catch(() => {}); } - if (convCurrentAudio && !convCurrentAudio.ended) { - convCurrentAudio.addEventListener('ended', () => setTimeout(tryStart, 200), { once: true }); + if (audioQueuePlaying || audioQueue.length > 0) { + // Audio still queued — trigger after queue drains + audioQueueDrainCb = tryStart; + } else if (convCurrentAudio && !convCurrentAudio.ended) { + convCurrentAudio.addEventListener('ended', () => setTimeout(tryStart, 150), { once: true }); } else { setTimeout(tryStart, 300); } @@ -554,7 +599,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { // ── Send turn via SSE ───────────────────────────────────────────────────── async function processBlob(blob) { autoMicGeneration++; // cancel any pending auto-mic from previous turn - convCurrentAudio = null; + clearAudio(); turnCount++; const turnN = turnCount; const t0 = Date.now(); @@ -618,17 +663,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { assistantBubble.textContent = assistantText; if (micStatus) micStatus.textContent = 'Synthesising speech…'; } else if (evt.type === 'audio') { - const mime = evt.mime || 'audio/wav'; - const binStr = atob(evt.b64); - const arr = new Uint8Array(binStr.length); - for (let i = 0; i < binStr.length; i++) arr[i] = binStr.charCodeAt(i); - const audioBlob = new Blob([arr], { type: mime }); - const url = URL.createObjectURL(audioBlob); - const audioEl = new Audio(url); - convCurrentAudio = audioEl; - audioEl.addEventListener('ended', () => { URL.revokeObjectURL(url); convCurrentAudio = null; }); - audioEl.play().catch(() => {}); - if (micStatus) micStatus.textContent = 'Speaking…'; + enqueueAudio(evt.b64, evt.mime || 'audio/wav'); } else if (evt.type === 'stats') { lastStats = evt; updateStats(evt); @@ -636,17 +671,15 @@ $('s-import-voices-file')?.addEventListener('change', async function () { conversationHistory.push({ role: 'user', content: userBubble.textContent }); conversationHistory.push({ role: 'assistant', content: assistantText }); addHistoryItem(turnN, lastStats?.total_ms ?? (Date.now() - t0), true); - scheduleAutoMic(); - if (!handsFreeToggle?.checked && micStatus) micStatus.textContent = 'Ready'; + scheduleAutoMic(); // triggers after audio queue drains } else if (evt.type === 'error') { + clearAudio(); const wrap = assistantBubble.closest('.conv-bubble-wrap'); if (wrap) wrap.remove(); const stage = evt.stage?.toUpperCase() || 'ERR'; let msg = evt.message || 'Unknown error'; - // Surface the actual server error detail, not the raw HTTP noise const detailMatch = msg.match(/HTTP \d+:\s*(.+)/s); if (detailMatch) msg = detailMatch[1].trim(); - // Truncate very long stack traces if (msg.length > 300) msg = msg.slice(0, 300) + '…'; addErrorBubble(`[${stage}] ${msg}`); addHistoryItem(turnN, Date.now() - t0, false); @@ -655,6 +688,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { } } } catch(e) { + clearAudio(); const wrap = assistantBubble.closest('.conv-bubble-wrap'); if (wrap) wrap.remove(); addErrorBubble(e.message); @@ -676,7 +710,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { text = text.trim(); if (!text || isProcessing) return; autoMicGeneration++; // cancel any pending auto-mic - convCurrentAudio = null; + clearAudio(); isProcessing = true; if (sendBtn) sendBtn.disabled = true; if (textInput) { textInput.disabled = true; textInput.value = ''; } @@ -733,17 +767,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { assistantBubble.textContent = assistantText; if (micStatus) micStatus.textContent = 'Synthesising speech…'; } else if (evt.type === 'audio') { - const mime = evt.mime || 'audio/wav'; - const binStr = atob(evt.b64); - const arr = new Uint8Array(binStr.length); - for (let i = 0; i < binStr.length; i++) arr[i] = binStr.charCodeAt(i); - const audioBlob = new Blob([arr], { type: mime }); - const url = URL.createObjectURL(audioBlob); - const audioEl = new Audio(url); - convCurrentAudio = audioEl; - audioEl.addEventListener('ended', () => { URL.revokeObjectURL(url); convCurrentAudio = null; }); - audioEl.play().catch(() => {}); - if (micStatus) micStatus.textContent = 'Speaking…'; + enqueueAudio(evt.b64, evt.mime || 'audio/wav'); } else if (evt.type === 'stats') { lastStats = evt; updateStats(evt); @@ -751,9 +775,9 @@ $('s-import-voices-file')?.addEventListener('change', async function () { conversationHistory.push({ role: 'user', content: text }); conversationHistory.push({ role: 'assistant', content: assistantText }); addHistoryItem(turnN, lastStats?.total_ms ?? (Date.now() - t0), true); - scheduleAutoMic(); - if (!handsFreeToggle?.checked && micStatus) micStatus.textContent = 'Ready'; + scheduleAutoMic(); // triggers after audio queue drains } else if (evt.type === 'error') { + clearAudio(); const wrap = assistantBubble.closest('.conv-bubble-wrap'); if (wrap) wrap.remove(); addErrorBubble(`[${evt.stage?.toUpperCase() || 'ERR'}] ${evt.message || 'Unknown error'}`); @@ -763,6 +787,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { } } } catch(e) { + clearAudio(); const wrap = assistantBubble.closest('.conv-bubble-wrap'); if (wrap) wrap.remove(); addErrorBubble(e.message); @@ -784,6 +809,7 @@ $('s-import-voices-file')?.addEventListener('change', async function () { if (mediaRecorder && mediaRecorder.state === 'recording') { stopRecording(); } else { + clearAudio(); // stop agent if still speaking startRecording(); } });