fix(voice_loop): 1011-Keepalive-Timeout bei langen Antworten beheben
Empfang und Echtzeit-Wiedergabe entkoppeln: ein _player_feeder spielt PCM- Haeppchen aus einer asyncio.Queue ab, der Empfangs-Loop legt nur per put_nowait ab und drainiert die WebSocket-Verbindung durchgehend. Vorher blockierte das Schreiben in den Player (Echtzeit, voller Puffer) den Empfangs-Loop -> Frames stauten sich -> websockets pausierte per Backpressure den Reader -> Pings/Pongs unbearbeitet -> 1011 keepalive ping timeout (trat bei langen Antworten auf). Zusaetzlich defensiv: ping_timeout=60, max_queue=None. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
ece64711a5
commit
84d6ac6e7b
1 changed files with 33 additions and 9 deletions
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue