#!/usr/bin/env python3 """Sprech-Loop: sprechen -> Antwort hoeren -> erneut sprechen. Nimmt vom Mikrofon auf (Push-to-Talk), schickt das Audio ueber EINE /ws/voice-Verbindung an das Gateway und spielt die Antwort ab. Das Gespraechs- gedaechtnis bleibt ueber die `session_id` erhalten. Standardmaessig folgt der Loop dem **System-Standard-Mikrofon und -Lautsprecher** (inkl. Bluetooth). Ein bestimmtes Geraet nur, wenn `--device` explizit gesetzt ist. Beispiele: python scripts/voice_loop.py # System-Standardgeraete python scripts/voice_loop.py --url ws://localhost:8003/ws/voice --session oma-anna python scripts/voice_loop.py --recorder arecord --device plughw:6,0 # bestimmtes Mikrofon python scripts/voice_loop.py --llm-provider openrouter --tts-provider openrouter python scripts/voice_loop.py --tts-provider openrouter --voice Puck # Cloud-Stimme testen python scripts/voice_loop.py --file frage.wav # ohne Mikrofon (Test) Voraussetzungen: laufendes Gateway, ein Aufnahmewerkzeug (`ffmpeg`/`parecord`/`arecord`), ein Player (`ffplay`/`paplay`/`aplay`), Python-Paket `websockets`. """ from __future__ import annotations import argparse import asyncio import io import json import os import shutil import signal import subprocess import sys import tempfile import select import termios import time import wave try: import websockets except ModuleNotFoundError: sys.exit("Fehlt: Python-Paket 'websockets' (kommt mit uvicorn[standard]).") TTS_SAMPLE_RATE = 24000 # Antwort-Audio des Gateways (s16le, mono) def _require(tool: str) -> str: path = shutil.which(tool) if not path: sys.exit(f"Fehlt: '{tool}' nicht gefunden. Bitte installieren.") return path def _first_player() -> list[str] | None: # paplay vor aplay: paplay folgt dem System-Standard-Ausgabegeraet (PulseAudio), # aplay nutzt das ALSA-Default, das auf PipeWire-Systemen oft tot ist. if shutil.which("ffplay"): return ["ffplay", "-loglevel", "quiet", "-nodisp", "-autoexit"] if shutil.which("paplay"): return ["paplay"] if shutil.which("aplay"): return ["aplay", "-q"] return None def resolve_recorder(choice: str, rate: int = 16000) -> str: """Waehlt das Aufnahmewerkzeug. 'auto' folgt dem **System-Standard-Mikrofon** und prueft per kurzem Test, dass das Werkzeug WIRKLICH Audio liefert -> es wird nie ein totes Geraet gewaehlt. Reihenfolge der Kandidaten: ffmpeg (PulseAudio-Default) -> parecord -> arecord -> pw-record. """ if choice != "auto": if not shutil.which(choice): sys.exit(f"Fehlt: Aufnahmewerkzeug '{choice}' nicht gefunden.") return choice available = [t for t in ("ffmpeg", "parecord", "arecord", "pw-record") if shutil.which(t)] if not available: sys.exit("Kein Aufnahmewerkzeug gefunden (ffmpeg / parecord / arecord / pw-record).") print(" Prüfe Standard-Aufnahmegerät …", flush=True) for tool in available: if _probe_records(tool, rate): return tool print(f" ⚠ Kein Werkzeug lieferte im Test Audio; nutze '{available[0]}'." " Ggf. --recorder/--device explizit setzen.") return available[0] def _probe_records(recorder: str, rate: int) -> bool: """Kurztest (~0,6 s), ob 'recorder' am System-Default tatsaechlich aufnimmt.""" tmp = tempfile.NamedTemporaryFile(suffix=".wav", delete=False) tmp.close() size = 0 try: proc = subprocess.Popen( _record_cmd(recorder, None, rate, tmp.name), stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) time.sleep(0.6) proc.send_signal(signal.SIGINT) try: proc.wait(timeout=3) except subprocess.TimeoutExpired: proc.kill() size = os.path.getsize(tmp.name) except Exception: # noqa: BLE001 - jedes Problem = Werkzeug taugt nicht size = 0 finally: try: os.unlink(tmp.name) except OSError: pass return size > 2000 # mehr als WAV-Header (44 B) + etwas Audio def _record_cmd(recorder: str, device: str | None, rate: int, outfile: str) -> list[str]: if recorder == "ffmpeg": # PulseAudio-Eingang folgt dem System-Standard-Mikrofon (device=None -> "default"). return ["ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error", "-y", "-f", "pulse", "-i", device or "default", "-ar", str(rate), "-ac", "1", outfile] if recorder == "arecord": cmd = ["arecord", "-q", "-f", "S16_LE", "-r", str(rate), "-c", "1"] if device: cmd += ["-D", device] return cmd + [outfile] if recorder == "pw-record": cmd = ["pw-record", "--rate", str(rate), "--channels", "1", "--format", "s16"] if device: cmd += ["--target", device] return cmd + [outfile] if recorder == "parecord": cmd = ["parecord", f"--rate={rate}", "--channels=1", "--format=s16le", "--file-format=wav"] if device: cmd += [f"--device={device}"] return cmd + [outfile] sys.exit(f"Unbekanntes Aufnahmewerkzeug: {recorder}") def record_utterance(recorder: str, device: str | None, rate: int) -> bytes: """Push-to-Talk: Enter startet, Enter stoppt die Aufnahme; gibt WAV-Bytes zurueck. Das Recorder-stderr wird abgefangen: Beim absichtlichen Stoppen per SIGINT meldet arecord ein harmloses 'Unterbrechung während des Betriebssystemaufrufs' (EINTR). Diese Meldung wird nur angezeigt, wenn die Aufnahme tatsaechlich fehlschlaegt. """ tmp = tempfile.NamedTemporaryFile(suffix=".wav", delete=False) tmp.close() cmd = _record_cmd(recorder, device, rate, tmp.name) input("\n[Enter] = Aufnahme START …") err = tempfile.TemporaryFile() proc = subprocess.Popen(cmd, stderr=err) input("[Enter] = Aufnahme STOP …") proc.send_signal(signal.SIGINT) try: proc.wait(timeout=5) except subprocess.TimeoutExpired: proc.kill() with open(tmp.name, "rb") as fh: data = fh.read() # Aufnahmedauer anzeigen (hilft, zu fruehes Stoppen sofort zu erkennen). # Hinweis: arecord schreibt beim Stoppen per Signal eine falsche Laenge in den # WAV-Header -> Dauer aus der TATSAECHLICHEN Datenmenge berechnen, nicht aus dem Header. try: with wave.open(io.BytesIO(data), "rb") as wav_in: rate = wav_in.getframerate() or 16000 channels = wav_in.getnchannels() or 1 width = wav_in.getsampwidth() or 2 raw = wav_in.readframes(10 ** 9) # liest alles Vorhandene (Header-Laenge unzuverlaessig) seconds = len(raw) / (rate * channels * width) print(f" Aufnahme: {seconds:.1f}s") except (wave.Error, EOFError): pass # Recorder-Fehlertext nur zeigen, wenn nichts/zu wenig aufgenommen wurde. if len(data) < 1000: err.seek(0) message = err.read().decode("utf-8", "replace").strip() if message: print(f" ({recorder}: {message})") err.close() return data def play_pcm(pcm: bytes) -> None: player = _first_player() if not player: print("(kein Player gefunden – Antwort-Audio wird nicht abgespielt)") return buf = io.BytesIO() with wave.open(buf, "wb") as w: w.setnchannels(1) w.setsampwidth(2) w.setframerate(TTS_SAMPLE_RATE) w.writeframes(pcm) if player[0] == "ffplay": # ffplay liest WAV von stdin via pipe:0 subprocess.run([*player, "-i", "pipe:0"], input=buf.getvalue()) else: 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) # paplay vor aplay: folgt dem System-Standard-Ausgabegeraet (PulseAudio/PipeWire). if shutil.which("ffplay"): cmd = ["ffplay", "-loglevel", "quiet", "-nodisp", "-autoexit", "-f", "s16le", "-ar", rate, "-ac", "1", "-i", "pipe:0"] elif shutil.which("paplay"): cmd = ["paplay", "--raw", f"--rate={rate}", "--format=s16le", "--channels=1"] elif shutil.which("aplay"): cmd = ["aplay", "-q", "-f", "S16_LE", "-r", rate, "-c", "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() 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) async def _stdin_watcher() -> None: """Wartet per Polling auf Enter-Tastendruck fuer Barge-in. Kein blockierender Thread: select.select mit Timeout 0 ist nicht-blockierend und funktioniert zuverlaessig auf Linux/macOS mit echtem tty-stdin. Canceln ist sicher (kein Zeichen wird konsumiert, ausser Enter wird tatsaechlich gedrueckt und diese Funktion war diejenige, die ihn gelesen hat). """ while True: ready, _, _ = select.select([sys.stdin], [], [], 0) if ready: sys.stdin.readline() # Enter-Zeilenumbruch konsumieren return await asyncio.sleep(0.05) def _start_frame(args) -> dict: frame = {"type": "start", "format": "wav"} # audio_stream ist Tri-State: True/False explizit, None -> Server-Default entscheidet. if getattr(args, "stream_audio", None) is not None: frame["audio_stream"] = args.stream_audio if getattr(args, "stream_text", False): frame["stream"] = True for key in ("stt_provider", "llm_provider", "tts_provider", "language", "voice"): value = getattr(args, key, None) if value: frame[key] = value return frame 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 `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) await ws.send(json.dumps({"type": "end"})) audio = b"" token_started = False while True: msg = await ws.recv() if isinstance(msg, (bytes, bytearray)): if queue is not None: # Nicht blockierend ablegen -> Verbindung wird durchgehend drainiert. queue.put_nowait(bytes(msg)) else: audio = bytes(msg) continue event = json.loads(msg) etype = event.get("type") if etype == "transcript": print(f" Du: {event.get('text','')!r}") elif etype == "token": # Live-Anzeige des Antworttextes, waehrend die KI ihn erzeugt. if not token_started: sys.stdout.write(" Assistent: ") token_started = True sys.stdout.write(event.get("text", "")) sys.stdout.flush() elif etype == "semantic": if token_started: sys.stdout.write("\n") # Live-Zeile abschliessen (Text steht schon) sys.stdout.flush() else: print(f" Assistent: {event.get('text','')!r}") elif etype == "emergency": print(f" ⚠ NOTFALL erkannt (Kategorie: {event.get('category')})") elif etype == "error": if token_started: sys.stdout.write("\n") print(f" Fehler {event.get('status','')}: {event.get('detail')}") return b"" elif etype == "done": break return audio async def _send_and_play(url: str, wav: bytes, start_frame: dict) -> None: """Eine kurze Verbindung pro Runde (Gedaechtnis bleibt serverseitig via session_id). Barge-in in zwei Phasen: Phase 1 (LLM-Streaming): _stdin_watcher() laeuft parallel zu one_turn(); Enter schickt {"type":"interrupt"} an den Server und bricht den Turn ab. Phase 2 (Audio-Wiedergabe): Nach dem LLM-Ende laeuft _stdin_watcher() weiter; Enter beendet den Player-Prozess sofort. Es wird immer ein durchgehender Player benutzt: der spielt sowohl viele Haeppchen (Stream-Audio) als auch ein einzelnes Komplett-Audio luekenlos ab. Wird kein Player 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"" interrupted = False # Watcher laeuft durch beide Phasen — wird erst im finally-Block beendet. watcher_task = asyncio.create_task(_stdin_watcher()) try: # 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: turn_task = asyncio.create_task(one_turn(ws, wav, start_frame, queue)) # Phase 1: LLM-Turn vs Enter done, _ = await asyncio.wait( {turn_task, watcher_task}, return_when=asyncio.FIRST_COMPLETED, ) if watcher_task in done and turn_task not in done: # Barge-in waehrend LLM-Streaming interrupted = True try: await ws.send(json.dumps({"type": "interrupt"})) except Exception: pass turn_task.cancel() try: await turn_task except asyncio.CancelledError: pass else: # Normales LLM-Ende (oder beide gleichzeitig) audio = turn_task.result() # watcher_task laeuft weiter -> Phase 2 im finally-Block finally: if feeder is not None: queue.put_nowait(None) # Feeder-Ende signalisieren if interrupted: # Phase 1 unterbrochen: Player sofort beenden if player is not None: player.terminate() if not watcher_task.done(): watcher_task.cancel() try: await feeder except Exception: pass if player is not None: try: player.wait(timeout=2) except subprocess.TimeoutExpired: player.kill() else: # Phase 2: Feeder abwarten, dann Player vs Enter racen try: await feeder except Exception: pass if player is not None: if not watcher_task.done(): # Race: Player laeuft aus vs Enter (Barge-in waehrend Audio) close_task = asyncio.create_task( asyncio.to_thread(_close_player, player) ) done2, _ = await asyncio.wait( {close_task, watcher_task}, return_when=asyncio.FIRST_COMPLETED, ) if watcher_task in done2 and close_task not in done2: # Barge-in waehrend Wiedergabe: Player sofort killen interrupted = True try: player.kill() except (ProcessLookupError, OSError): pass else: watcher_task.cancel() if not close_task.done(): try: await close_task except Exception: pass else: # Watcher lief gleichzeitig mit Turn ab — nicht als Barge-in werten await asyncio.to_thread(_close_player, player) else: # Kein Player/Feeder if not watcher_task.done(): watcher_task.cancel() if interrupted: print(" ↩ Unterbrochen — drücke [Enter] für neue Aufnahme …") elif audio: # nur wenn kein Streaming-Player verfuegbar war play_pcm(audio) # Gepufferten stdin leeren: ein Enter, der waehrend des Turns nicht als Barge-in # verarbeitet wurde, wuerde sonst den START-Prompt sofort ueberspringen. try: termios.tcflush(sys.stdin.fileno(), termios.TCIFLUSH) except Exception: pass async def run(args) -> None: url = f"{args.url}?session_id={args.session}" if args.token: url += f"&token={args.token}" start_frame = _start_frame(args) if args.file: with open(args.file, "rb") as fh: wav = fh.read() print(f"Sende Datei: {args.file}") await _send_and_play(url, wav, start_frame) return recorder = resolve_recorder(args.recorder, args.rate) print(f"Ziel: {args.url} (Session '{args.session}')") print(f"Aufnahme mit: {recorder}" + (f" (Gerät: {args.device})" if args.device else " (System-Standardgerät)")) print("Sprich nach 'START', stoppe mit Enter. Strg+C beendet den Loop.") while True: try: wav = record_utterance(recorder, args.device, args.rate) except KeyboardInterrupt: print("\nEnde.") return if len(wav) < 1000: print(" ⚠ Keine/zu kurze Aufnahme. Pruefe das System-Standard-Mikrofon" " (Ubuntu: Einstellungen → Ton → Eingabe) und ob es Pegel zeigt.") print(" Notfalls direktes ALSA-Geraet erzwingen (arecord -l zeigt die Nummer):") print(" python scripts/voice_loop.py --recorder arecord --device plughw:6,0") continue try: await _send_and_play(url, wav, start_frame) except Exception as exc: # noqa: BLE001 - Verbindung pro Runde; Fehler nicht fatal print(f" Verbindungsfehler: {exc}") def main() -> None: p = argparse.ArgumentParser(description="Sprech-Loop fuer das Voice-Assistant-Gateway") p.add_argument("--url", default="ws://127.0.0.1:8003/ws/voice") p.add_argument("--session", default="voice-loop") p.add_argument("--token", default=None, help="Bearer-Token, falls AUTH_ENABLED=true") p.add_argument("--device", default=None, help="Aufnahmegeraet explizit (ffmpeg/parecord: PulseAudio-Quelle; " "arecord: ALSA z. B. plughw:6,0; pw-record: --target). " "Ohne Angabe folgt der Loop dem System-Standardgeraet.") p.add_argument("--recorder", default="auto", choices=["auto", "ffmpeg", "pw-record", "parecord", "arecord"], help="Aufnahmewerkzeug; 'auto' folgt dem System-Standardmikrofon und " "prueft per Kurztest, dass es wirklich aufnimmt") p.add_argument("--rate", type=int, default=16000) p.add_argument("--stt-provider", dest="stt_provider", default=None) p.add_argument("--llm-provider", dest="llm_provider", default=None) p.add_argument("--tts-provider", dest="tts_provider", default=None) p.add_argument("--voice", default=None, help="TTS-Stimme (provider-spezifisch; ohne Angabe gilt der Provider-Default, " "z. B. OPENROUTER_TTS_VOICE bzw. PIPER_VOICE)") p.add_argument("--language", default=None) # Audio-Streaming: Default entscheidet der Server (AUDIO_STREAM_DEFAULT, i. d. R. an). audio_grp = p.add_mutually_exclusive_group() audio_grp.add_argument("--stream-audio", dest="stream_audio", action="store_const", const=True, default=None, help="satzweises Vorlesen erzwingen (frueher Ton)") audio_grp.add_argument("--no-stream-audio", dest="stream_audio", action="store_const", const=False, default=None, help="satzweises Vorlesen abschalten (komplett am Ende)") p.add_argument("--stream-text", dest="stream_text", action="store_true", help="Antworttext live am Monitor anzeigen, waehrend die KI ihn erzeugt") p.add_argument("--file", default=None, help="WAV statt Mikrofon senden (Test ohne Aufnahme)") args = p.parse_args() try: asyncio.run(run(args)) except KeyboardInterrupt: print("\nEnde.") if __name__ == "__main__": main()