From ced63bac4ce876d93f4ee1dc0b3b0a055a7f518c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dieter=20Schl=C3=BCter?= Date: Wed, 17 Jun 2026 17:16:05 +0200 Subject: [PATCH] feat(voice_loop): --stream-audio (satzweises Vorlesen waehrend der Generierung) - neuer Flag --stream-audio setzt audio_stream:true im /ws/voice-Start-Frame - durchgehender Roh-PCM-Player (ffplay/aplay/paplay), Haeppchen werden sofort eingespeist -> Antwort beginnt nach dem ersten Satz, nicht erst nach der ganzen - Schreiben via asyncio.to_thread (Event-Loop bleibt frei -> kein Keepalive-Timeout); Player wird nach Verbindungsschluss geleert - ohne Flag unveraendert (komplettes Audio nach Verbindungsschluss) - Doku-Hinweis ergaenzt; live verifiziert (STT lokal + LLM lokal + TTS remote, satzweise) Co-Authored-By: Claude Opus 4.8 --- BEDIENUNGSANLEITUNG.md | 1 + scripts/voice_loop.py | 88 +++++++++++++++++++++++++++++++++++------- 2 files changed, 75 insertions(+), 14 deletions(-) diff --git a/BEDIENUNGSANLEITUNG.md b/BEDIENUNGSANLEITUNG.md index 6f82b1d..7bd2c35 100644 --- a/BEDIENUNGSANLEITUNG.md +++ b/BEDIENUNGSANLEITUNG.md @@ -109,6 +109,7 @@ Ablauf je Runde: Nützliche Optionen: ```bash +python scripts/voice_loop.py --stream-audio # satzweises Vorlesen: Antwort beginnt frueher python scripts/voice_loop.py --recorder pw-record # PipeWire-Aufnahme (Standard bei 'auto') python scripts/voice_loop.py --recorder arecord --device hw:1,0 # ALSA, bestimmtes Mikrofon python scripts/voice_loop.py --llm-provider openrouter --tts-provider openrouter diff --git a/scripts/voice_loop.py b/scripts/voice_loop.py index 610f0c6..924b6b7 100644 --- a/scripts/voice_loop.py +++ b/scripts/voice_loop.py @@ -150,8 +150,48 @@ def play_pcm(pcm: bytes) -> None: subprocess.run(player + ["/dev/stdin"], input=buf.getvalue()) +def open_stream_player(): + """Startet einen durchgehenden Player, der rohes s16le-PCM von stdin abspielt. + + So koennen Audio-Haeppchen (satzweises TTS) sofort und luekenlos abgespielt werden, + waehrend die KI noch weiter generiert. None, wenn kein Player gefunden wird. + """ + rate = str(TTS_SAMPLE_RATE) + if shutil.which("ffplay"): + cmd = ["ffplay", "-loglevel", "quiet", "-nodisp", "-autoexit", + "-f", "s16le", "-ar", rate, "-ac", "1", "-i", "pipe:0"] + elif shutil.which("aplay"): + cmd = ["aplay", "-q", "-f", "S16_LE", "-r", rate, "-c", "1"] + elif shutil.which("paplay"): + cmd = ["paplay", "--raw", f"--rate={rate}", "--format=s16le", "--channels=1"] + else: + return None + return subprocess.Popen(cmd, stdin=subprocess.PIPE) + + +def _feed_player(player, chunk: bytes) -> None: + try: + player.stdin.write(chunk) + player.stdin.flush() + except (BrokenPipeError, ValueError): + pass + + +def _close_player(player) -> None: + try: + player.stdin.close() + except (BrokenPipeError, ValueError): + pass + try: + player.wait(timeout=60) + except subprocess.TimeoutExpired: + player.kill() + + def _start_frame(args) -> dict: frame = {"type": "start", "format": "wav"} + if getattr(args, "stream_audio", False): + frame["audio_stream"] = True for key in ("stt_provider", "llm_provider", "tts_provider", "language"): value = getattr(args, key, None) if value: @@ -159,9 +199,13 @@ def _start_frame(args) -> dict: return frame -async def one_turn(ws, wav: bytes, start_frame: dict) -> bytes: - """Sendet eine Aeusserung, gibt das Antwort-Audio (PCM) zurueck. Spielt NICHT ab, - damit die Wiedergabe ausserhalb der offenen Verbindung passieren kann.""" +async def one_turn(ws, wav: bytes, start_frame: dict, player=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). + """ await ws.send(json.dumps(start_frame)) await ws.send(wav) await ws.send(json.dumps({"type": "end"})) @@ -170,7 +214,11 @@ async def one_turn(ws, wav: bytes, start_frame: dict) -> bytes: while True: msg = await ws.recv() if isinstance(msg, (bytes, bytearray)): - audio = bytes(msg) + if player is not None: + # Schreiben im Thread, damit der Event-Loop frei bleibt (Keepalive). + await asyncio.to_thread(_feed_player, player, bytes(msg)) + else: + audio = bytes(msg) continue event = json.loads(msg) etype = event.get("type") @@ -188,13 +236,21 @@ async def one_turn(ws, wav: bytes, start_frame: dict) -> bytes: return audio -async def _send_and_play(url: str, wav: bytes, start_frame: dict) -> None: - """Oeffnet eine kurze Verbindung pro Runde (Gedaechtnis bleibt serverseitig via - session_id erhalten), empfaengt die Antwort und spielt sie NACH dem Schliessen ab. - So blockiert die (lange) Aufnahme/Wiedergabe nie eine offene WebSocket-Verbindung - (kein Keepalive-Timeout).""" - async with websockets.connect(url, max_size=None) as ws: - audio = await one_turn(ws, wav, start_frame) +async def _send_and_play(url: str, wav: bytes, start_frame: dict, stream_audio: bool = False) -> None: + """Eine kurze Verbindung pro Runde (Gedaechtnis bleibt serverseitig via session_id). + + Ohne Stream: komplettes Audio sammeln, NACH dem Schliessen abspielen. + Mit Stream-Audio: durchgehenden Player oeffnen, Haeppchen sofort einspeisen; der + Player wird nach dem Schliessen der Verbindung geleert (kein Keepalive-Timeout).""" + player = open_stream_player() if stream_audio else None + if stream_audio and player is None: + print(" (kein Player fuer Stream-Audio gefunden – spiele am Ende komplett ab)") + try: + async with websockets.connect(url, max_size=None) as ws: + audio = await one_turn(ws, wav, start_frame, player) + finally: + if player is not None: + await asyncio.to_thread(_close_player, player) if audio: play_pcm(audio) @@ -209,11 +265,12 @@ async def run(args) -> None: with open(args.file, "rb") as fh: wav = fh.read() print(f"Sende Datei: {args.file}") - await _send_and_play(url, wav, start_frame) + await _send_and_play(url, wav, start_frame, args.stream_audio) return recorder = resolve_recorder(args.recorder) - print(f"Ziel: {args.url} (Session '{args.session}')") + print(f"Ziel: {args.url} (Session '{args.session}')" + + (" [Stream-Audio: satzweises Vorlesen]" if args.stream_audio else "")) print(f"Aufnahme mit: {recorder}" + (f" (Gerät: {args.device})" if args.device else " (Standardgerät)")) print("Sprich nach 'START', stoppe mit Enter. Strg+C beendet den Loop.") @@ -231,7 +288,7 @@ async def run(args) -> None: print(" python scripts/voice_loop.py --recorder arecord --device plughw:2,0") continue try: - await _send_and_play(url, wav, start_frame) + await _send_and_play(url, wav, start_frame, args.stream_audio) except Exception as exc: # noqa: BLE001 - Verbindung pro Runde; Fehler nicht fatal print(f" Verbindungsfehler: {exc}") @@ -251,6 +308,9 @@ def main() -> None: p.add_argument("--llm-provider", dest="llm_provider", default=None) p.add_argument("--tts-provider", dest="tts_provider", default=None) p.add_argument("--language", default=None) + p.add_argument("--stream-audio", dest="stream_audio", action="store_true", + help="Antwort satzweise vorlesen, sobald der erste Satz fertig ist " + "(geringere Zeit bis zum ersten Ton)") p.add_argument("--file", default=None, help="WAV statt Mikrofon senden (Test ohne Aufnahme)") args = p.parse_args() try: