#!/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 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) 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). 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"" 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: 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 play_pcm(audio) 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()