diff --git a/scripts/voice_loop.py b/scripts/voice_loop.py index 0798968..f03e08a 100644 --- a/scripts/voice_loop.py +++ b/scripts/voice_loop.py @@ -188,6 +188,22 @@ def _close_player(player) -> None: player.kill() +async def _player_feeder(player, queue: "asyncio.Queue") -> None: + """Spielt PCM-Haeppchen aus der Queue ab — entkoppelt vom Empfangs-Loop. + + Wichtig gegen 1011-Keepalive-Timeouts: Die Echtzeit-Wiedergabe (Player-stdin + blockiert, wenn der Puffer voll ist) darf NICHT den WebSocket-Empfang bremsen. + Sonst stauen sich eingehende Frames, websockets pausiert per Backpressure den + Transport-Reader, Pings/Pongs werden nicht mehr verarbeitet -> Verbindung bricht. + Der Empfangs-Loop legt Haeppchen nur in die Queue (put_nowait) und drainiert die + Verbindung durchgehend; hier wird in Ruhe (im Thread) abgespielt. None = Ende.""" + while True: + chunk = await queue.get() + if chunk is None: + return + await asyncio.to_thread(_feed_player, player, chunk) + + def _start_frame(args) -> dict: frame = {"type": "start", "format": "wav"} # audio_stream ist Tri-State: True/False explizit, None -> Server-Default entscheidet. @@ -202,12 +218,13 @@ def _start_frame(args) -> dict: return frame -async def one_turn(ws, wav: bytes, start_frame: dict, player=None) -> bytes: +async def one_turn(ws, wav: bytes, start_frame: dict, queue: "asyncio.Queue | None" = None) -> bytes: """Sendet eine Aeusserung und verarbeitet die Antwort-Events. - Ist `player` gesetzt (Stream-Audio): jedes PCM-Haeppchen wird sofort in den Player - geschrieben (gespielt waehrend die KI weiter generiert) und b"" zurueckgegeben. - Sonst wird das komplette Audio gesammelt und zurueckgegeben (Wiedergabe spaeter). + Ist `queue` gesetzt (Stream-Audio): jedes PCM-Haeppchen wird sofort (nicht blockierend) + in die Queue gelegt und von `_player_feeder` abgespielt; gibt b"" zurueck. Der Empfang + bleibt dadurch reaktiv (kein Keepalive-Timeout bei langen Antworten). Ohne Queue wird + das komplette Audio gesammelt und zurueckgegeben (Wiedergabe spaeter). """ await ws.send(json.dumps(start_frame)) await ws.send(wav) @@ -218,9 +235,9 @@ async def one_turn(ws, wav: bytes, start_frame: dict, player=None) -> bytes: while True: msg = await ws.recv() if isinstance(msg, (bytes, bytearray)): - if player is not None: - # Schreiben im Thread, damit der Event-Loop frei bleibt (Keepalive). - await asyncio.to_thread(_feed_player, player, bytes(msg)) + if queue is not None: + # Nicht blockierend ablegen -> Verbindung wird durchgehend drainiert. + queue.put_nowait(bytes(msg)) else: audio = bytes(msg) continue @@ -261,10 +278,17 @@ async def _send_and_play(url: str, wav: bytes, start_frame: dict) -> None: gefunden, wird das Audio gesammelt und als WAV abgespielt. Die Wiedergabe blockiert nie die offene Verbindung (Schreiben im Thread; Player wird nach Schliessen geleert).""" player = open_stream_player() + queue: asyncio.Queue | None = asyncio.Queue() if player is not None else None + feeder = asyncio.create_task(_player_feeder(player, queue)) if queue is not None else None + audio = b"" try: - async with websockets.connect(url, max_size=None) as ws: - audio = await one_turn(ws, wav, start_frame, player) + # ping_timeout/max_queue defensiv: lange Antworten + Echtzeit-Wiedergabe ueberleben. + async with websockets.connect(url, max_size=None, max_queue=None, ping_timeout=60) as ws: + audio = await one_turn(ws, wav, start_frame, queue) finally: + if feeder is not None: + queue.put_nowait(None) # Feeder beenden, restliche Haeppchen noch abspielen + await feeder if player is not None: await asyncio.to_thread(_close_player, player) if audio: # nur wenn kein Streaming-Player verfuegbar war