my_voice_assistant_v2/scripts/voice_loop.py

250 lines
9.1 KiB
Python
Raw Normal View History

#!/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()
# 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()