#!/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. Beispiele: python scripts/voice_loop.py python scripts/voice_loop.py --url ws://localhost:8003/ws/voice --session oma-anna python scripts/voice_loop.py --device hw:1,0 # bestimmtes Mikrofon (arecord -L) python scripts/voice_loop.py --llm-provider openrouter --tts-provider openrouter python scripts/voice_loop.py --file frage.wav # ohne Mikrofon (Test) Voraussetzungen: laufendes Gateway, `arecord` (Aufnahme), ein Player (`ffplay`/`aplay`/`paplay`), Python-Paket `websockets`. """ from __future__ import annotations import argparse import asyncio import io import json import shutil import signal import subprocess import sys import tempfile 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: if shutil.which("ffplay"): return ["ffplay", "-loglevel", "quiet", "-nodisp", "-autoexit"] if shutil.which("aplay"): return ["aplay", "-q"] if shutil.which("paplay"): return ["paplay"] return None def resolve_recorder(choice: str) -> str: """Waehlt das Aufnahmewerkzeug. 'auto' bevorzugt PipeWire (pw-record).""" if choice != "auto": if not shutil.which(choice): sys.exit(f"Fehlt: Aufnahmewerkzeug '{choice}' nicht gefunden.") return choice for tool in ("pw-record", "parecord", "arecord"): if shutil.which(tool): return tool sys.exit("Kein Aufnahmewerkzeug gefunden (pw-record / parecord / arecord).") def _record_cmd(recorder: str, device: str | None, rate: int, outfile: str) -> list[str]: 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 _start_frame(args) -> dict: frame = {"type": "start", "format": "wav"} for key in ("stt_provider", "llm_provider", "tts_provider", "language"): value = getattr(args, key, None) if value: frame[key] = value 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.""" await ws.send(json.dumps(start_frame)) await ws.send(wav) await ws.send(json.dumps({"type": "end"})) audio = b"" while True: msg = await ws.recv() if isinstance(msg, (bytes, bytearray)): audio = bytes(msg) continue event = json.loads(msg) etype = event.get("type") if etype == "transcript": print(f" Du: {event.get('text','')!r}") elif etype == "semantic": print(f" Assistent: {event.get('text','')!r}") elif etype == "emergency": print(f" ⚠ NOTFALL erkannt (Kategorie: {event.get('category')})") elif etype == "error": 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: """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) if audio: 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) print(f"Ziel: {args.url} (Session '{args.session}')") 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.") 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. Moegliche Ursache: Aufnahmewerkzeug" " oder Geraet nicht nutzbar.") print(" Verfuegbare Mikrofone: arecord -l") print(" Direktes ALSA-Geraet verwenden, z. B.:") print(" python scripts/voice_loop.py --recorder arecord --device plughw:2,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 (arecord: -L; pw-record: --target; parecord: --device)") p.add_argument("--recorder", default="auto", choices=["auto", "pw-record", "parecord", "arecord"], help="Aufnahmewerkzeug; 'auto' bevorzugt PipeWire (pw-record)") 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("--language", default=None) 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()