feat: Audio-Streaming (chunked TTS) über WebSocket (#4 Ausbau)
- 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 <noreply@anthropic.com>
This commit is contained in:
parent
379e002460
commit
b5913b0a44
7 changed files with 157 additions and 19 deletions
|
|
@ -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`).
|
`{"text": "..."}`; Antwort kommt als Event-Folge (`ack`, `semantic`, Audio, `done`).
|
||||||
Token per Query (`?token=…`), Gedächtnis per `?session_id=…`. Mit
|
Token per Query (`?token=…`), Gedächtnis per `?session_id=…`. Mit
|
||||||
`{"text": "...", "stream": true}` kommt die Antwort schon während der Generierung
|
`{"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`
|
> **Für lokale Entwicklung** ist in der mitgelieferten `.env` `AUTH_ENABLED=false`
|
||||||
> gesetzt — dann ist kein Token nötig (anonymer Nutzer).
|
> gesetzt — dann ist kein Token nötig (anonymer Nutzer).
|
||||||
|
|
|
||||||
|
|
@ -187,7 +187,8 @@ Device Router (strikt, Singleton); Output-Lifecycle; **Authentifizierung
|
||||||
+ dauerhafte Nutzer-Präferenzen**; **Gesprächsgedächtnis pro Session (Verlauf im
|
+ dauerhafte Nutzer-Präferenzen**; **Gesprächsgedächtnis pro Session (Verlauf im
|
||||||
Store, fließt ins LLM)**; **Langzeit-Erinnerungen pro Nutzer (als LLM-Kontext)**;
|
Store, fließt ins LLM)**; **Langzeit-Erinnerungen pro Nutzer (als LLM-Kontext)**;
|
||||||
**WebSocket-Streaming-Chat (`/ws/chat`) inkl. Token-Level-LLM-Streaming (SSE,
|
**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`,
|
**Platzhalter (Gerüst):** Audio-Endpunkte (`local-default`, `bluetooth`,
|
||||||
`mobile-ws`, `mobile-webrtc`) liefern leere Chunks — nur Auswahl/Lifecycle sind
|
`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.
|
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.
|
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.
|
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.
|
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).
|
6. **Betrieb:** Kosten-/Quota-Kontrolle pro Nutzer; Notfall-/Eskalationskonzept (Senioren-Kontext).
|
||||||
7. **TransportRouter** als eigene lokal/remote-Achse aktivieren.
|
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
|
│ ├── api/ # health, chat, speak, transcribe, devices, sessions, config, admin, me, ws
|
||||||
│ ├── core/ # orchestrator
|
│ ├── core/ # orchestrator
|
||||||
│ ├── audio/ # router, transport_router, endpoints/input|output/*
|
│ ├── 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)
|
│ └── providers/ # stt/ llm/ tts/ (openrouter + lokale Stubs)
|
||||||
├── config/ # voice-assistant.example.toml (+ lokale .toml, gitignored)
|
├── config/ # voice-assistant.example.toml (+ lokale .toml, gitignored)
|
||||||
├── data/ # SQLite-DB (gitignored)
|
├── data/ # SQLite-DB (gitignored)
|
||||||
|
|
|
||||||
16
README.md
16
README.md
|
|
@ -133,13 +133,17 @@ der Server streamt strukturierte Events zurück: `ack` → `semantic` → Audio
|
||||||
Erinnerungen gelten wie bei `POST /api/chat`.
|
Erinnerungen gelten wie bei `POST /api/chat`.
|
||||||
|
|
||||||
**Token-Streaming:** Mit `{"text": "...", "stream": true}` schickt der Server die
|
**Token-Streaming:** Mit `{"text": "...", "stream": true}` schickt der Server die
|
||||||
LLM-Antwort schon während der Generierung als `token`-Events
|
LLM-Antwort schon während der Generierung als `token`-Events — spürbar geringere
|
||||||
(`ack` → `token*` → `semantic` → Audio → `done`) — spürbar geringere wahrgenommene
|
wahrgenommene Latenz. OpenRouter und der lokale OpenAI-kompatible Provider streamen
|
||||||
Latenz. OpenRouter und der lokale OpenAI-kompatible Provider streamen via SSE;
|
via SSE; Provider ohne Streaming liefern die komplette Antwort als ein `token`-Event.
|
||||||
Provider ohne Streaming liefern die komplette Antwort als ein `token`-Event.
|
|
||||||
|
|
||||||
> Audio-Streaming (chunked TTS), Audio-Eingang/Streaming-STT, Barge-in und WebRTC
|
**Audio-Streaming:** Mit `{"text": "...", "audio_stream": true}` wird das Audio
|
||||||
> sind als nächste Increments vorgesehen (siehe Architektur-Dokument).
|
**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
|
## Authentifizierung
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -96,11 +96,25 @@ async def ws_chat(
|
||||||
|
|
||||||
voice = msg.get("voice") or settings.openrouter_tts_voice
|
voice = msg.get("voice") or settings.openrouter_tts_voice
|
||||||
stream = bool(msg.get("stream"))
|
stream = bool(msg.get("stream"))
|
||||||
try:
|
audio_stream = bool(msg.get("audio_stream"))
|
||||||
if stream:
|
|
||||||
async def on_token(delta):
|
|
||||||
await websocket.send_json({"type": "token", "text": delta})
|
|
||||||
|
|
||||||
|
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(
|
trace, audio = await orchestrator.chat_stream(
|
||||||
text,
|
text,
|
||||||
language=route.language,
|
language=route.language,
|
||||||
|
|
@ -108,6 +122,7 @@ async def ws_chat(
|
||||||
output=output,
|
output=output,
|
||||||
history=llm_context,
|
history=llm_context,
|
||||||
on_token=on_token,
|
on_token=on_token,
|
||||||
|
on_audio=on_audio,
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
trace, audio = await orchestrator.chat_text(
|
trace, audio = await orchestrator.chat_text(
|
||||||
|
|
@ -132,7 +147,9 @@ async def ws_chat(
|
||||||
"spoken": trace.spoken_response,
|
"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(
|
await websocket.send_json(
|
||||||
{"type": "done", "audio_format": "pcm", "sample_rate": 24000}
|
{"type": "done", "audio_format": "pcm", "sample_rate": 24000}
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
from app.schemas import AudioChunk, PipelineTrace
|
from app.schemas import AudioChunk, PipelineTrace
|
||||||
|
from app.pipeline.sentence_chunker import SentenceChunker
|
||||||
|
|
||||||
# Festes Ausgabeformat der TTS-Stufe (s16le PCM, 24 kHz, mono).
|
# Festes Ausgabeformat der TTS-Stufe (s16le PCM, 24 kHz, mono).
|
||||||
TTS_AUDIO_FORMAT = "pcm"
|
TTS_AUDIO_FORMAT = "pcm"
|
||||||
|
|
@ -114,16 +115,30 @@ class Orchestrator:
|
||||||
output=None,
|
output=None,
|
||||||
history: list[dict] | None = None,
|
history: list[dict] | None = None,
|
||||||
on_token=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
|
`on_token(delta)` (async) wird pro LLM-Token-Delta aufgerufen.
|
||||||
werden erst nach der vollstaendigen Antwort erzeugt (TTS ist nicht streamend).
|
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 = PipelineTrace()
|
||||||
trace.raw_transcript = text
|
trace.raw_transcript = text
|
||||||
trace.cleaned_transcript = await self.input_cleaner.run(text or "")
|
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] = []
|
parts: list[str] = []
|
||||||
stream_fn = getattr(self.llm, "stream", None)
|
stream_fn = getattr(self.llm, "stream", None)
|
||||||
if stream_fn is not None:
|
if stream_fn is not None:
|
||||||
|
|
@ -131,12 +146,18 @@ class Orchestrator:
|
||||||
parts.append(delta)
|
parts.append(delta)
|
||||||
if on_token:
|
if on_token:
|
||||||
await on_token(delta)
|
await on_token(delta)
|
||||||
|
if chunker:
|
||||||
|
for sentence in chunker.feed(delta):
|
||||||
|
await _emit_sentence(sentence)
|
||||||
else:
|
else:
|
||||||
# Provider ohne Streaming -> komplette Antwort als ein Token.
|
# Provider ohne Streaming -> komplette Antwort als ein Token.
|
||||||
result = await self.llm.complete(trace.cleaned_transcript or "", history=history)
|
result = await self.llm.complete(trace.cleaned_transcript or "", history=history)
|
||||||
parts.append(result)
|
parts.append(result)
|
||||||
if on_token:
|
if on_token:
|
||||||
await on_token(result)
|
await on_token(result)
|
||||||
|
if chunker:
|
||||||
|
for sentence in chunker.feed(result):
|
||||||
|
await _emit_sentence(sentence)
|
||||||
|
|
||||||
trace.semantic_response = "".join(parts)
|
trace.semantic_response = "".join(parts)
|
||||||
if not trace.semantic_response:
|
if not trace.semantic_response:
|
||||||
|
|
@ -151,6 +172,13 @@ class Orchestrator:
|
||||||
language=language,
|
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)
|
await self._emit_to_output(audio, output)
|
||||||
return trace, audio
|
return trace, audio
|
||||||
|
|
|
||||||
36
app/pipeline/sentence_chunker.py
Normal file
36
app/pipeline/sentence_chunker.py
Normal file
|
|
@ -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
|
||||||
|
|
@ -5,10 +5,21 @@ from fastapi.testclient import TestClient
|
||||||
import app.dependencies as deps
|
import app.dependencies as deps
|
||||||
from app.main import app
|
from app.main import app
|
||||||
from app.providers.llm.base import sse_delta, LLMProvider
|
from app.providers.llm.base import sse_delta, LLMProvider
|
||||||
|
from app.pipeline.sentence_chunker import SentenceChunker
|
||||||
|
|
||||||
client = TestClient(app)
|
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():
|
def test_sse_delta_parsing():
|
||||||
assert sse_delta('data: {"choices":[{"delta":{"content":"Hal"}}]}') == "Hal"
|
assert sse_delta('data: {"choices":[{"delta":{"content":"Hal"}}]}') == "Hal"
|
||||||
assert sse_delta("data: [DONE]") is None
|
assert sse_delta("data: [DONE]") is None
|
||||||
|
|
@ -92,3 +103,42 @@ def test_ws_stream_fallback_for_nonstreaming_llm(monkeypatch):
|
||||||
token = ws.receive_json()
|
token = ws.receive_json()
|
||||||
assert token["type"] == "token" and token["text"] == "komplett"
|
assert token["type"] == "token" and token["text"] == "komplett"
|
||||||
assert ws.receive_json()["type"] == "semantic"
|
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
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue