- Config audio_stream_default=True; _run_turn nutzt es, wenn die Anfrage audio_stream nicht explizit setzt (explizite Anfrage gewinnt) - voice_loop: immer durchgehender Player (spielt 1 oder N Haeppchen luekenlos); --stream-audio / --no-stream-audio als Tri-State (sonst entscheidet der Server-Default) - conftest: audio_stream_default=False fuer deterministische Tests; neuer Test fuer Default-an-Streaming ueber /ws/voice; 67 Tests gruen - Doku + .env.example aktualisiert (AUDIO_STREAM_DEFAULT, --no-stream-audio) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
343 lines
13 KiB
Python
343 lines
13 KiB
Python
#!/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 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"}
|
||
# 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"):
|
||
value = getattr(args, key, None)
|
||
if value:
|
||
frame[key] = value
|
||
return frame
|
||
|
||
|
||
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"}))
|
||
|
||
audio = b""
|
||
token_started = False
|
||
while True:
|
||
msg = await ws.recv()
|
||
if isinstance(msg, (bytes, bytearray)):
|
||
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")
|
||
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()
|
||
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: # 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)
|
||
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)
|
||
# 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()
|