From b5913b0a4441caefd9a7f79c1bd42534c7742c86 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dieter=20Schl=C3=BCter?= Date: Wed, 17 Jun 2026 04:43:50 +0200 Subject: [PATCH] =?UTF-8?q?feat:=20Audio-Streaming=20(chunked=20TTS)=20?= =?UTF-8?q?=C3=BCber=20WebSocket=20(#4=20Ausbau)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - SentenceChunker (pipeline/sentence_chunker.py): inkrementelle Satzsegmentierung - Orchestrator.chat_stream(on_audio): satzweise TTS, Audio-Chunk pro fertigem Satz; Gesamtaudio zusaetzlich an den Output-Endpunkt - WS /ws/chat {"audio_stream":true}: audio-Events (json seq + binaerer Frame) live, kein finales Vollaudio; mit stream kombinierbar - Tests: 45 gruen (+2: Sentence-Chunker, Audio-Streaming) - Doku aktualisiert (README, BEDIENUNGSANLEITUNG, Architektur) Co-Authored-By: Claude Opus 4.8 --- BEDIENUNGSANLEITUNG.md | 4 ++- Docs/voice-assistant-architecture.md | 7 ++-- README.md | 16 +++++---- app/api/ws.py | 27 ++++++++++++--- app/core/orchestrator.py | 36 +++++++++++++++++--- app/pipeline/sentence_chunker.py | 36 ++++++++++++++++++++ tests/test_streaming.py | 50 ++++++++++++++++++++++++++++ 7 files changed, 157 insertions(+), 19 deletions(-) create mode 100644 app/pipeline/sentence_chunker.py diff --git a/BEDIENUNGSANLEITUNG.md b/BEDIENUNGSANLEITUNG.md index d49d4ff..7436afe 100644 --- a/BEDIENUNGSANLEITUNG.md +++ b/BEDIENUNGSANLEITUNG.md @@ -274,7 +274,9 @@ Diese Erinnerungen gibt der Assistent bei jedem Chat als Kontext mit — auch oh `{"text": "..."}`; Antwort kommt als Event-Folge (`ack`, `semantic`, Audio, `done`). Token per Query (`?token=…`), Gedächtnis per `?session_id=…`. Mit `{"text": "...", "stream": true}` kommt die Antwort schon während der Generierung -als `token`-Events (geringere wahrgenommene Latenz). +als `token`-Events (geringere wahrgenommene Latenz). Mit `"audio_stream": true` +kommt zusätzlich das Audio satzweise (`audio`-Event + binärer Frame), sobald ein +Satz fertig ist. > **Für lokale Entwicklung** ist in der mitgelieferten `.env` `AUTH_ENABLED=false` > gesetzt — dann ist kein Token nötig (anonymer Nutzer). diff --git a/Docs/voice-assistant-architecture.md b/Docs/voice-assistant-architecture.md index da60379..4ca8d87 100644 --- a/Docs/voice-assistant-architecture.md +++ b/Docs/voice-assistant-architecture.md @@ -187,7 +187,8 @@ Device Router (strikt, Singleton); Output-Lifecycle; **Authentifizierung + dauerhafte Nutzer-Präferenzen**; **Gesprächsgedächtnis pro Session (Verlauf im Store, fließt ins LLM)**; **Langzeit-Erinnerungen pro Nutzer (als LLM-Kontext)**; **WebSocket-Streaming-Chat (`/ws/chat`) inkl. Token-Level-LLM-Streaming (SSE, -opt-in via `stream:true`)**; automatisierte Tests. +`stream:true`) und satzweisem Audio-Streaming (chunked TTS, `audio_stream:true`)**; +automatisierte Tests. **Platzhalter (Gerüst):** Audio-Endpunkte (`local-default`, `bluetooth`, `mobile-ws`, `mobile-webrtc`) liefern leere Chunks — nur Auswahl/Lifecycle sind @@ -202,7 +203,7 @@ Reihenfolge der Weiterentwicklung: 1. **(erledigt)** Konfig- & Routing-Fundament: Profile, Device Router, Registry, Pro-Request-Override. 2. **(erledigt)** Cloud-Fundament: Bearer-Token-Auth, Mehrbenutzer, persistenter SQLite-Store, Mandanten-Trennung, dauerhafte Nutzer-Präferenzen. Offen: Skalierung auf gemeinsamen Store (Postgres/Redis) für mehrere Instanzen. 3. **(erledigt)** Konversationsgedächtnis: Kurzzeit-Gesprächsverlauf pro Session + Langzeit-Erinnerungen pro Nutzer (manuell gepflegt, als LLM-Kontext). Offen: **automatische** Extraktion/Zusammenfassung von Erinnerungen aus Gesprächen. -4. **(teilweise erledigt)** Echtzeit: WebSocket-Streaming-Chat (`/ws/chat`) mit Event-Folge (ack/semantic/audio/done) **und Token-Level-LLM-Streaming (SSE, opt-in `stream:true`)** sind umgesetzt. Offen: **Audio-Streaming (chunked TTS)**, **Audio-Eingang/Streaming-STT**, **Barge-in/Turn-Manager**, **WebRTC**. +4. **(teilweise erledigt)** Echtzeit: WebSocket-Streaming-Chat (`/ws/chat`), **Token-Level-LLM-Streaming (SSE, `stream:true`)** und **Audio-Streaming (chunked TTS satzweise, `audio_stream:true`)** sind umgesetzt. Offen: **Audio-Eingang/Streaming-STT**, **Barge-in/Turn-Manager**, **WebRTC**. 5. **Resilienz:** Fallback-Policy (remote KI fällt aus → lokaler/alternativer Provider), Metriken/Tracing. 6. **Betrieb:** Kosten-/Quota-Kontrolle pro Nutzer; Notfall-/Eskalationskonzept (Senioren-Kontext). 7. **TransportRouter** als eigene lokal/remote-Achse aktivieren. @@ -226,7 +227,7 @@ voice-assistant-scaffold/ │ ├── api/ # health, chat, speak, transcribe, devices, sessions, config, admin, me, ws │ ├── core/ # orchestrator │ ├── audio/ # router, transport_router, endpoints/input|output/* -│ ├── pipeline/ # input_cleaner, spoken_response_adapter, tts_normalizer +│ ├── pipeline/ # input_cleaner, spoken_response_adapter, tts_normalizer, sentence_chunker │ └── providers/ # stt/ llm/ tts/ (openrouter + lokale Stubs) ├── config/ # voice-assistant.example.toml (+ lokale .toml, gitignored) ├── data/ # SQLite-DB (gitignored) diff --git a/README.md b/README.md index 74be773..9ac9288 100644 --- a/README.md +++ b/README.md @@ -133,13 +133,17 @@ der Server streamt strukturierte Events zurück: `ack` → `semantic` → Audio Erinnerungen gelten wie bei `POST /api/chat`. **Token-Streaming:** Mit `{"text": "...", "stream": true}` schickt der Server die -LLM-Antwort schon während der Generierung als `token`-Events -(`ack` → `token*` → `semantic` → Audio → `done`) — spürbar geringere wahrgenommene -Latenz. OpenRouter und der lokale OpenAI-kompatible Provider streamen via SSE; -Provider ohne Streaming liefern die komplette Antwort als ein `token`-Event. +LLM-Antwort schon während der Generierung als `token`-Events — spürbar geringere +wahrgenommene Latenz. OpenRouter und der lokale OpenAI-kompatible Provider streamen +via SSE; Provider ohne Streaming liefern die komplette Antwort als ein `token`-Event. -> Audio-Streaming (chunked TTS), Audio-Eingang/Streaming-STT, Barge-in und WebRTC -> sind als nächste Increments vorgesehen (siehe Architektur-Dokument). +**Audio-Streaming:** Mit `{"text": "...", "audio_stream": true}` wird das Audio +**satzweise** erzeugt (chunked TTS) und pro fertigem Satz als `audio`-Event (JSON +mit `seq` + binärer Frame) gesendet — die Ausgabe beginnt, bevor die Antwort fertig +ist. `stream` und `audio_stream` lassen sich kombinieren. + +> Audio-Eingang/Streaming-STT, Barge-in/Turn-Manager und WebRTC sind als nächste +> Increments vorgesehen (siehe Architektur-Dokument). ## Authentifizierung diff --git a/app/api/ws.py b/app/api/ws.py index 6570a7f..630b564 100644 --- a/app/api/ws.py +++ b/app/api/ws.py @@ -96,11 +96,25 @@ async def ws_chat( voice = msg.get("voice") or settings.openrouter_tts_voice stream = bool(msg.get("stream")) - try: - if stream: - async def on_token(delta): - await websocket.send_json({"type": "token", "text": delta}) + audio_stream = bool(msg.get("audio_stream")) + on_token = None + if stream: + async def on_token(delta): + await websocket.send_json({"type": "token", "text": delta}) + + on_audio = None + if audio_stream: + audio_seq = 0 + + async def on_audio(chunk): + nonlocal audio_seq + await websocket.send_json({"type": "audio", "seq": audio_seq}) + audio_seq += 1 + await websocket.send_bytes(chunk) + + try: + if stream or audio_stream: trace, audio = await orchestrator.chat_stream( text, language=route.language, @@ -108,6 +122,7 @@ async def ws_chat( output=output, history=llm_context, on_token=on_token, + on_audio=on_audio, ) else: trace, audio = await orchestrator.chat_text( @@ -132,7 +147,9 @@ async def ws_chat( "spoken": trace.spoken_response, } ) - await websocket.send_bytes(audio) + # Bei audio_stream wurden die Audio-Chunks bereits live gesendet. + if not audio_stream: + await websocket.send_bytes(audio) await websocket.send_json( {"type": "done", "audio_format": "pcm", "sample_rate": 24000} ) diff --git a/app/core/orchestrator.py b/app/core/orchestrator.py index de4cfb6..15ef63b 100644 --- a/app/core/orchestrator.py +++ b/app/core/orchestrator.py @@ -1,4 +1,5 @@ from app.schemas import AudioChunk, PipelineTrace +from app.pipeline.sentence_chunker import SentenceChunker # Festes Ausgabeformat der TTS-Stufe (s16le PCM, 24 kHz, mono). TTS_AUDIO_FORMAT = "pcm" @@ -114,16 +115,30 @@ class Orchestrator: output=None, history: list[dict] | None = None, on_token=None, + on_audio=None, ): - """Wie chat_text, aber die LLM-Antwort wird tokenweise gestreamt. + """Wie chat_text, aber gestreamt. - `on_token(delta)` (async) wird pro Token-Delta aufgerufen. Audio/Output - werden erst nach der vollstaendigen Antwort erzeugt (TTS ist nicht streamend). + `on_token(delta)` (async) wird pro LLM-Token-Delta aufgerufen. + Ist `on_audio(chunk)` gesetzt, wird das Audio satzweise erzeugt (chunked TTS) + und pro fertigem Satz ausgeliefert, statt erst am Ende komplett. """ trace = PipelineTrace() trace.raw_transcript = text trace.cleaned_transcript = await self.input_cleaner.run(text or "") + chunker = SentenceChunker() if on_audio else None + audio_parts: list[bytes] = [] + + async def _emit_sentence(sentence: str) -> None: + spoken = await self.spoken_adapter.run(sentence, language=language) + ready = await self.tts_normalizer.run(spoken, language=language) + if not ready.strip(): + return + chunk = await self.tts.synthesize(ready, voice=voice) + audio_parts.append(chunk) + await on_audio(chunk) + parts: list[str] = [] stream_fn = getattr(self.llm, "stream", None) if stream_fn is not None: @@ -131,12 +146,18 @@ class Orchestrator: parts.append(delta) if on_token: await on_token(delta) + if chunker: + for sentence in chunker.feed(delta): + await _emit_sentence(sentence) else: # Provider ohne Streaming -> komplette Antwort als ein Token. result = await self.llm.complete(trace.cleaned_transcript or "", history=history) parts.append(result) if on_token: await on_token(result) + if chunker: + for sentence in chunker.feed(result): + await _emit_sentence(sentence) trace.semantic_response = "".join(parts) if not trace.semantic_response: @@ -151,6 +172,13 @@ class Orchestrator: language=language, ) - audio = await self.tts.synthesize(trace.tts_ready_text, voice=voice) + if chunker: + tail = chunker.flush() + if tail: + await _emit_sentence(tail) + audio = b"".join(audio_parts) + else: + audio = await self.tts.synthesize(trace.tts_ready_text, voice=voice) + await self._emit_to_output(audio, output) return trace, audio diff --git a/app/pipeline/sentence_chunker.py b/app/pipeline/sentence_chunker.py new file mode 100644 index 0000000..39e0afb --- /dev/null +++ b/app/pipeline/sentence_chunker.py @@ -0,0 +1,36 @@ +import re + +# Satzende: . ! ? … gefolgt von Whitespace (oder Stringende beim flush). +_SENTENCE_END = re.compile(r"[.!?…]+(?=\s)") + + +class SentenceChunker: + """Inkrementelle Satzsegmentierung fuer gestreamte LLM-Token. + + `feed(delta)` liefert die seit dem letzten Aufruf fertig gewordenen Saetze, + `flush()` den verbleibenden Rest (z. B. der letzte Satz ohne abschliessendes + Leerzeichen). Damit kann pro Satz schon TTS erzeugt werden, waehrend das LLM + noch weiterschreibt. + """ + + def __init__(self): + self._buffer = "" + + def feed(self, text: str) -> list[str]: + self._buffer += text + sentences: list[str] = [] + while True: + match = _SENTENCE_END.search(self._buffer) + if not match: + break + end = match.end() + sentence = self._buffer[:end].strip() + self._buffer = self._buffer[end:] + if sentence: + sentences.append(sentence) + return sentences + + def flush(self) -> str: + rest = self._buffer.strip() + self._buffer = "" + return rest diff --git a/tests/test_streaming.py b/tests/test_streaming.py index 6f5a6d8..e5db1fc 100644 --- a/tests/test_streaming.py +++ b/tests/test_streaming.py @@ -5,10 +5,21 @@ from fastapi.testclient import TestClient import app.dependencies as deps from app.main import app from app.providers.llm.base import sse_delta, LLMProvider +from app.pipeline.sentence_chunker import SentenceChunker client = TestClient(app) +def test_sentence_chunker_incremental(): + ch = SentenceChunker() + emitted = [] + for tok in ["Hallo", " Anna", ". ", "Wie", " geht", " es", "? ", "Tschuess"]: + emitted += ch.feed(tok) + assert emitted == ["Hallo Anna.", "Wie geht es?"] + assert ch.flush() == "Tschuess" + assert SentenceChunker().feed("Eins. Zwei! Drei? Vier") == ["Eins.", "Zwei!", "Drei?"] + + def test_sse_delta_parsing(): assert sse_delta('data: {"choices":[{"delta":{"content":"Hal"}}]}') == "Hal" assert sse_delta("data: [DONE]") is None @@ -92,3 +103,42 @@ def test_ws_stream_fallback_for_nonstreaming_llm(monkeypatch): token = ws.receive_json() assert token["type"] == "token" and token["text"] == "komplett" assert ws.receive_json()["type"] == "semantic" + + +def test_ws_audio_stream_sends_chunks_per_sentence(monkeypatch): + tts_calls = [] + + class StreamLLM: + async def complete(self, text, history=None, session_id=None): + return "Satz eins. Satz zwei." + + async def stream(self, text, history=None, session_id=None): + for tok in ["Satz ", "eins. ", "Satz ", "zwei."]: + yield tok + + class CountTTS: + async def synthesize(self, text, voice=None, audio_format="pcm"): + tts_calls.append(text) + return b"A" * len(tts_calls) + + monkeypatch.setitem(deps.LLM_REGISTRY, "stream", lambda s: StreamLLM()) + monkeypatch.setitem(deps.TTS_REGISTRY, "cnt", lambda s: CountTTS()) + base = {"llm_provider": "stream", "tts_provider": "cnt", "output_endpoint": "loopback"} + + with client.websocket_connect("/ws/chat") as ws: + ws.send_json({"text": "x", "audio_stream": True, **base}) + assert ws.receive_json()["type"] == "ack" + + audio_events = 0 + event = ws.receive_json() + while event["type"] != "semantic": + assert event["type"] == "audio" + assert ws.receive_bytes() # binärer Audio-Chunk folgt + audio_events += 1 + event = ws.receive_json() + + assert audio_events == 2 # zwei Sätze -> zwei Chunks + # Kein finales Vollaudio mehr -> direkt done. + assert ws.receive_json()["type"] == "done" + + assert len(tts_calls) == 2 # TTS pro Satz