diff --git a/app/core/orchestrator.py b/app/core/orchestrator.py index cff898a..af2e6ac 100644 --- a/app/core/orchestrator.py +++ b/app/core/orchestrator.py @@ -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