my_voice_assistant_v2/scripts/voice_loop.py
Dieter Schlüter 5296459b07 feat(voice_loop): --stream-text (Antworttext live am Monitor)
- neuer Flag --stream-text setzt stream:true im /ws/voice-Start-Frame;
  token-Events werden inline ausgegeben (Wort fuer Wort, waehrend die KI generiert)
- mit --stream-audio kombinierbar; bei semantic wird die Live-Zeile abgeschlossen
- Token-Events sind winzig -> Audio-Startzeit praktisch unveraendert (gemessen:
  Text 1. Token lokal ~1 s, Cloud ~3 s; erster Ton wie zuvor)
- Test: /ws/voice mit stream:true liefert token-Events; 66 Tests gruen
- Doku ergaenzt (Optionen + Feature-Hinweis)

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 17:47:11 +02:00

341 lines
13 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/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"}
if getattr(args, "stream_audio", False):
frame["audio_stream"] = True
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, stream_audio: bool = False) -> None:
"""Eine kurze Verbindung pro Runde (Gedaechtnis bleibt serverseitig via session_id).
Ohne Stream: komplettes Audio sammeln, NACH dem Schliessen abspielen.
Mit Stream-Audio: durchgehenden Player oeffnen, Haeppchen sofort einspeisen; der
Player wird nach dem Schliessen der Verbindung geleert (kein Keepalive-Timeout)."""
player = open_stream_player() if stream_audio else None
if stream_audio and player is None:
print(" (kein Player fuer Stream-Audio gefunden spiele am Ende komplett ab)")
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:
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, args.stream_audio)
return
recorder = resolve_recorder(args.recorder)
print(f"Ziel: {args.url} (Session '{args.session}')"
+ (" [Stream-Audio: satzweises Vorlesen]" if args.stream_audio else ""))
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, args.stream_audio)
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("--stream-audio", dest="stream_audio", action="store_true",
help="Antwort satzweise vorlesen, sobald der erste Satz fertig ist "
"(geringere Zeit bis zum ersten Ton)")
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()