feat(voice_loop): --stream-audio (satzweises Vorlesen waehrend der Generierung)
- neuer Flag --stream-audio setzt audio_stream:true im /ws/voice-Start-Frame - durchgehender Roh-PCM-Player (ffplay/aplay/paplay), Haeppchen werden sofort eingespeist -> Antwort beginnt nach dem ersten Satz, nicht erst nach der ganzen - Schreiben via asyncio.to_thread (Event-Loop bleibt frei -> kein Keepalive-Timeout); Player wird nach Verbindungsschluss geleert - ohne Flag unveraendert (komplettes Audio nach Verbindungsschluss) - Doku-Hinweis ergaenzt; live verifiziert (STT lokal + LLM lokal + TTS remote, satzweise) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
70c7e2ec0c
commit
ced63bac4c
2 changed files with 75 additions and 14 deletions
|
|
@ -109,6 +109,7 @@ Ablauf je Runde:
|
|||
|
||||
Nützliche Optionen:
|
||||
```bash
|
||||
python scripts/voice_loop.py --stream-audio # satzweises Vorlesen: Antwort beginnt frueher
|
||||
python scripts/voice_loop.py --recorder pw-record # PipeWire-Aufnahme (Standard bei 'auto')
|
||||
python scripts/voice_loop.py --recorder arecord --device hw:1,0 # ALSA, bestimmtes Mikrofon
|
||||
python scripts/voice_loop.py --llm-provider openrouter --tts-provider openrouter
|
||||
|
|
|
|||
|
|
@ -150,8 +150,48 @@ def play_pcm(pcm: bytes) -> None:
|
|||
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
|
||||
for key in ("stt_provider", "llm_provider", "tts_provider", "language"):
|
||||
value = getattr(args, key, None)
|
||||
if value:
|
||||
|
|
@ -159,9 +199,13 @@ def _start_frame(args) -> dict:
|
|||
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."""
|
||||
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"}))
|
||||
|
|
@ -170,6 +214,10 @@ async def one_turn(ws, wav: bytes, start_frame: dict) -> bytes:
|
|||
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)
|
||||
|
|
@ -188,13 +236,21 @@ async def one_turn(ws, wav: bytes, start_frame: dict) -> bytes:
|
|||
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 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)
|
||||
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)
|
||||
|
||||
|
|
@ -209,11 +265,12 @@ async def run(args) -> None:
|
|||
with open(args.file, "rb") as fh:
|
||||
wav = fh.read()
|
||||
print(f"Sende Datei: {args.file}")
|
||||
await _send_and_play(url, wav, start_frame)
|
||||
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}')")
|
||||
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.")
|
||||
|
|
@ -231,7 +288,7 @@ async def run(args) -> None:
|
|||
print(" python scripts/voice_loop.py --recorder arecord --device plughw:2,0")
|
||||
continue
|
||||
try:
|
||||
await _send_and_play(url, wav, start_frame)
|
||||
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}")
|
||||
|
||||
|
|
@ -251,6 +308,9 @@ def main() -> 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("--file", default=None, help="WAV statt Mikrofon senden (Test ohne Aufnahme)")
|
||||
args = p.parse_args()
|
||||
try:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue