perf(stream): TTS-Synthese vom Token-Streaming entkoppeln (Bubble läuft voraus)

Im Streaming-Pfad lief die satzweise TTS-Synthese bisher inline in der
LLM-Leseschleife: pro Satzgrenze wartete `chat_stream` per `await` auf die
Synthese, wodurch das Token-Streaming (Text in der Antwort-Bubble) satzweise
pausierte.

Jetzt Producer/Consumer: Die LLM-Leseschleife legt fertige Sätze in eine
asyncio.Queue und liest sofort weiter; ein einzelner Consumer-Task
synthetisiert sequenziell und ruft `on_audio` in Reihenfolge auf. Ein
einzelner Consumer garantiert die Audio-Reihenfolge (gapless-Invariante im
Client). Sentinel + await am Ende, Exceptions aus dem Consumer werden
propagiert, verwaiste Tasks werden im finally abgeräumt.

Deterministisch gemessen (Fake-LLM + Fake-TTS 3x800ms Synthese):
Token-Spanne 1845ms -> 242ms, größtes Token-Loch 821ms -> 20ms.
Audio-Chunks bleiben in Reihenfolge (seq 0..n). text_only- und
Nicht-Streaming-Pfad unverändert.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Dieter Schlüter 2026-06-26 11:33:04 +02:00
commit 1b1e1bb1d5

View file

@ -1,3 +1,5 @@
import asyncio
from app.schemas import AudioChunk, PipelineTrace
from app.pipeline.sentence_chunker import SentenceChunker
from app.metrics import timer, metrics
@ -197,49 +199,92 @@ class Orchestrator:
parts: list[str] = []
stream_fn = getattr(self.llm, "stream", None)
if stream_fn is not None:
async for delta in stream_fn(trace.cleaned_transcript or "", history=history, language=effective_language):
parts.append(delta)
if on_token:
await on_token(delta)
if chunker:
for sentence in chunker.feed(delta):
await _emit_sentence(sentence)
else:
result = await self.llm.complete(trace.cleaned_transcript or "", history=history, language=effective_language)
parts.append(result)
if on_token:
await on_token(result)
if chunker:
for sentence in chunker.feed(result):
# Producer/Consumer: Die satzweise TTS-Synthese läuft in einem eigenen
# Consumer-Task, damit das Token-Streaming (Text in der Bubble) nicht
# mehr satzweise auf die Synthese wartet. Ein *einzelner* Consumer
# garantiert die Reihenfolge der Audio-Chunks — harte Invariante fürs
# gapless Playback im Client; out-of-order würde das Audio zerstören.
queue: asyncio.Queue | None = asyncio.Queue() if chunker else None
consumer_task: asyncio.Task | None = None
async def _consume() -> None:
while True:
sentence = await queue.get()
try:
if sentence is None: # Sentinel: keine Sätze mehr
return
await _emit_sentence(sentence)
finally:
queue.task_done()
trace.semantic_response = "".join(parts)
if not trace.semantic_response:
raise RuntimeError("LLM returned an empty response")
async def _dispatch(sentence: str) -> None:
# Satz in die Queue legen und sofort weiterlesen (kein await auf die
# Synthese). Ist der Consumer zuvor an einer Synthese-Exception
# gestorben, diese sofort hochreichen, statt weiter zu puffern.
if consumer_task is not None and consumer_task.done():
await consumer_task
await queue.put(sentence)
trace.spoken_response = await self.spoken_adapter.run(
trace.semantic_response,
language=effective_language,
)
trace.tts_ready_text = await self.tts_normalizer.run(
trace.spoken_response,
language=effective_language,
level=self.normalize_level,
)
if queue is not None:
consumer_task = asyncio.create_task(_consume())
if text_only:
return trace, b""
try:
if stream_fn is not None:
async for delta in stream_fn(trace.cleaned_transcript or "", history=history, language=effective_language):
parts.append(delta)
if on_token:
await on_token(delta)
if chunker:
for sentence in chunker.feed(delta):
await _dispatch(sentence)
else:
result = await self.llm.complete(trace.cleaned_transcript or "", history=history, language=effective_language)
parts.append(result)
if on_token:
await on_token(result)
if chunker:
for sentence in chunker.feed(result):
await _dispatch(sentence)
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, language=effective_language
trace.semantic_response = "".join(parts)
if not trace.semantic_response:
raise RuntimeError("LLM returned an empty response")
trace.spoken_response = await self.spoken_adapter.run(
trace.semantic_response,
language=effective_language,
)
trace.tts_ready_text = await self.tts_normalizer.run(
trace.spoken_response,
language=effective_language,
level=self.normalize_level,
)
if text_only:
return trace, b""
if chunker:
tail = chunker.flush()
if tail:
await _dispatch(tail)
await queue.put(None) # Sentinel: Consumer beenden
await consumer_task # auf restliche Synthese warten + Exceptions propagieren
consumer_task = None
audio = b"".join(audio_parts)
else:
audio = await self.tts.synthesize(
trace.tts_ready_text, voice=voice, language=effective_language
)
finally:
# Bei Fehler/Abbruch den noch laufenden Consumer-Task abräumen,
# damit kein verwaister Task zurückbleibt.
if consumer_task is not None and not consumer_task.done():
consumer_task.cancel()
try:
await consumer_task
except (asyncio.CancelledError, Exception):
pass
await self._emit_to_output(audio, output)
return trace, audio