feat: Barge-in/Turn-Manager und VAD-Aeusserungserkennung (#4 Ausbau)
- Barge-in: Antwort-Turn als abbrechbarer asyncio.Task; {"type":"interrupt"} oder
neue Eingabe bricht laufende Antwort ab -> interrupted-Event (/ws/chat + /ws/voice)
- VAD (app/audio/vad.py): energie-basierte Stille-Erkennung (reines Python, int16-PCM)
- /ws/voice opt-in {"type":"start","vad":true}: automatisches Aeusserungsende ohne end
- Tests: 52 gruen (+5: VAD-Unit, Barge-in, VAD-Auto-Segmentierung)
- Doku aktualisiert (README, BEDIENUNGSANLEITUNG, Architektur)
Offen (schwere Deps/Dienste): echte partielle Live-Transkripte (Streaming-STT),
WebRTC (aiortc). STT laeuft heute pro Aeusserung.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
9340d3f998
commit
093da817d8
6 changed files with 312 additions and 48 deletions
144
app/api/ws.py
144
app/api/ws.py
|
|
@ -3,14 +3,18 @@
|
|||
- /ws/chat : Text rein (JSON pro Turn), Antwort als Event-Folge zurueck.
|
||||
- /ws/voice: Audio rein (binaere Chunks + Control), Transkription -> selbe Pipeline.
|
||||
|
||||
Event-Folge der Antwort: ack -> [token*] -> [audio*] -> semantic -> done.
|
||||
Antwort-Events: ack -> [token*] -> [audio*] -> semantic -> done.
|
||||
Mit {"stream":true} kommen LLM-Token live, mit {"audio_stream":true} das Audio
|
||||
satzweise (chunked TTS). /ws/voice sendet zuvor ein transcript-Event.
|
||||
|
||||
Spaeter (eigene Increments): partielle Live-Transkripte (Streaming-STT mit VAD),
|
||||
Barge-in/Turn-Manager und WebRTC.
|
||||
Barge-in: Ein {"type":"interrupt"}-Frame oder eine neue Eingabe bricht eine laufende
|
||||
Antwort ab (-> interrupted-Event). Der Antwort-Turn laeuft als abbrechbarer Task.
|
||||
|
||||
Spaeter (eigene Increments): echte partielle Live-Transkripte (Streaming-STT-Dienst),
|
||||
WebRTC.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
|
||||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
||||
|
|
@ -24,6 +28,7 @@ from app.dependencies import (
|
|||
resolve_output_endpoint,
|
||||
)
|
||||
from app.store import SessionOwnershipError
|
||||
from app.audio.vad import EnergyVAD
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
|
@ -126,6 +131,51 @@ async def _run_turn(websocket, store, user, session_id, route, orchestrator, out
|
|||
await websocket.send_json({"type": "done", "audio_format": "pcm", "sample_rate": 24000})
|
||||
|
||||
|
||||
async def _cancel_active(task, websocket) -> None:
|
||||
"""Bricht einen laufenden Antwort-Turn ab (Barge-in) und meldet 'interrupted'."""
|
||||
if task is None or task.done():
|
||||
return
|
||||
task.cancel()
|
||||
try:
|
||||
await task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
await websocket.send_json({"type": "interrupted"})
|
||||
|
||||
|
||||
async def _chat_turn(websocket, store, user, session_id, text, options):
|
||||
try:
|
||||
route, orchestrator, output = await _resolve(user, session_id, options)
|
||||
except SessionOwnershipError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 403, "detail": str(exc)})
|
||||
return
|
||||
except RoutingError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 422, "detail": str(exc)})
|
||||
return
|
||||
await _run_turn(websocket, store, user, session_id, route, orchestrator, output, text, options)
|
||||
|
||||
|
||||
async def _voice_turn(websocket, store, user, session_id, audio, fmt, options):
|
||||
try:
|
||||
route, orchestrator, output = await _resolve(user, session_id, options)
|
||||
except SessionOwnershipError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 403, "detail": str(exc)})
|
||||
return
|
||||
except RoutingError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 422, "detail": str(exc)})
|
||||
return
|
||||
try:
|
||||
transcript = await orchestrator.stt.transcribe(audio, fmt=fmt, language=route.language)
|
||||
except Exception as exc:
|
||||
await websocket.send_json({"type": "error", "status": 502, "detail": str(exc)})
|
||||
return
|
||||
await websocket.send_json({"type": "transcript", "text": transcript})
|
||||
if not transcript or not transcript.strip():
|
||||
await websocket.send_json({"type": "error", "detail": "empty transcript"})
|
||||
return
|
||||
await _run_turn(websocket, store, user, session_id, route, orchestrator, output, transcript, options)
|
||||
|
||||
|
||||
@router.websocket("/ws/chat")
|
||||
async def ws_chat(websocket: WebSocket, session_id: str | None = None, token: str | None = None):
|
||||
user = _authenticate(token)
|
||||
|
|
@ -134,24 +184,26 @@ async def ws_chat(websocket: WebSocket, session_id: str | None = None, token: st
|
|||
return
|
||||
await websocket.accept()
|
||||
store = get_store()
|
||||
active = None
|
||||
|
||||
try:
|
||||
while True:
|
||||
msg = await websocket.receive_json()
|
||||
if msg.get("type") == "interrupt":
|
||||
await _cancel_active(active, websocket)
|
||||
active = None
|
||||
continue
|
||||
text = (msg.get("text") or "").strip()
|
||||
if not text:
|
||||
await websocket.send_json({"type": "error", "detail": "empty text"})
|
||||
continue
|
||||
try:
|
||||
route, orchestrator, output = await _resolve(user, session_id, msg)
|
||||
except SessionOwnershipError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 403, "detail": str(exc)})
|
||||
continue
|
||||
except RoutingError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 422, "detail": str(exc)})
|
||||
continue
|
||||
await _run_turn(websocket, store, user, session_id, route, orchestrator, output, text, msg)
|
||||
await _cancel_active(active, websocket) # Barge-in bei neuer Eingabe
|
||||
active = asyncio.create_task(
|
||||
_chat_turn(websocket, store, user, session_id, text, msg)
|
||||
)
|
||||
except WebSocketDisconnect:
|
||||
if active and not active.done():
|
||||
active.cancel()
|
||||
return
|
||||
|
||||
|
||||
|
|
@ -166,15 +218,33 @@ async def ws_voice(websocket: WebSocket, session_id: str | None = None, token: s
|
|||
|
||||
audio_buffer = bytearray()
|
||||
fmt = "wav"
|
||||
active = None
|
||||
vad = None
|
||||
vad_options: dict = {}
|
||||
|
||||
async def _start_voice(audio: bytes, options: dict):
|
||||
nonlocal active
|
||||
await _cancel_active(active, websocket) # Barge-in bei neuer Aeusserung
|
||||
active = asyncio.create_task(
|
||||
_voice_turn(websocket, store, user, session_id, audio, fmt, options)
|
||||
)
|
||||
|
||||
try:
|
||||
while True:
|
||||
message = await websocket.receive()
|
||||
if message["type"] == "websocket.disconnect":
|
||||
if active and not active.done():
|
||||
active.cancel()
|
||||
return
|
||||
|
||||
if message.get("bytes") is not None:
|
||||
audio_buffer.extend(message["bytes"])
|
||||
# VAD: Aeusserungsende automatisch erkennen (opt-in via start-Frame).
|
||||
if vad is not None and vad.feed(message["bytes"]):
|
||||
audio = bytes(audio_buffer)
|
||||
audio_buffer.clear()
|
||||
vad.reset()
|
||||
await _start_voice(audio, vad_options)
|
||||
continue
|
||||
|
||||
raw = message.get("text")
|
||||
|
|
@ -187,9 +257,22 @@ async def ws_voice(websocket: WebSocket, session_id: str | None = None, token: s
|
|||
continue
|
||||
|
||||
ctype = control.get("type")
|
||||
if ctype == "interrupt":
|
||||
await _cancel_active(active, websocket)
|
||||
active = None
|
||||
continue
|
||||
if ctype == "start":
|
||||
audio_buffer.clear()
|
||||
fmt = control.get("format", "wav")
|
||||
if control.get("vad"):
|
||||
vad = EnergyVAD(
|
||||
sample_rate=control.get("sample_rate", 16000),
|
||||
threshold=control.get("vad_threshold", 500.0),
|
||||
silence_ms=control.get("vad_silence_ms", 700.0),
|
||||
)
|
||||
vad_options = control
|
||||
else:
|
||||
vad = None
|
||||
continue
|
||||
if ctype != "end":
|
||||
continue
|
||||
|
|
@ -198,35 +281,12 @@ async def ws_voice(websocket: WebSocket, session_id: str | None = None, token: s
|
|||
await websocket.send_json({"type": "error", "detail": "no audio received"})
|
||||
continue
|
||||
|
||||
try:
|
||||
route, orchestrator, output = await _resolve(user, session_id, control)
|
||||
except SessionOwnershipError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 403, "detail": str(exc)})
|
||||
audio_buffer.clear()
|
||||
continue
|
||||
except RoutingError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 422, "detail": str(exc)})
|
||||
audio_buffer.clear()
|
||||
continue
|
||||
|
||||
try:
|
||||
transcript = await orchestrator.stt.transcribe(
|
||||
bytes(audio_buffer), fmt=fmt, language=route.language
|
||||
)
|
||||
except Exception as exc:
|
||||
await websocket.send_json({"type": "error", "status": 502, "detail": str(exc)})
|
||||
audio_buffer.clear()
|
||||
continue
|
||||
finally:
|
||||
audio_buffer.clear()
|
||||
|
||||
await websocket.send_json({"type": "transcript", "text": transcript})
|
||||
if not transcript or not transcript.strip():
|
||||
await websocket.send_json({"type": "error", "detail": "empty transcript"})
|
||||
continue
|
||||
|
||||
await _run_turn(
|
||||
websocket, store, user, session_id, route, orchestrator, output, transcript, control
|
||||
)
|
||||
audio = bytes(audio_buffer)
|
||||
audio_buffer.clear()
|
||||
if vad is not None:
|
||||
vad.reset()
|
||||
await _start_voice(audio, control)
|
||||
except WebSocketDisconnect:
|
||||
if active and not active.done():
|
||||
active.cancel()
|
||||
return
|
||||
|
|
|
|||
59
app/audio/vad.py
Normal file
59
app/audio/vad.py
Normal file
|
|
@ -0,0 +1,59 @@
|
|||
"""Einfache energie-basierte Sprachaktivitaetserkennung (VAD).
|
||||
|
||||
Reines Python (stdlib `array`), arbeitet auf s16le-PCM (mono). Erkennt das Ende
|
||||
einer Aeusserung anhand andauernder Stille nach erkannter Sprache. Damit kann der
|
||||
Server in /ws/voice Aeusserungen automatisch segmentieren, ohne dass der Client
|
||||
ein explizites Ende-Signal schickt.
|
||||
|
||||
Hinweis: Das ersetzt keinen echten Streaming-STT-Dienst (keine wortweisen
|
||||
Teil-Transkripte) - es bestimmt nur die Aeusserungsgrenzen.
|
||||
"""
|
||||
|
||||
import array
|
||||
import math
|
||||
|
||||
|
||||
def rms(pcm: bytes) -> float:
|
||||
"""Lautstaerke (RMS) eines s16le-PCM-Puffers; 0.0 bei leerem Puffer."""
|
||||
usable = len(pcm) - (len(pcm) % 2)
|
||||
if usable <= 0:
|
||||
return 0.0
|
||||
samples = array.array("h")
|
||||
samples.frombytes(pcm[:usable])
|
||||
if not samples:
|
||||
return 0.0
|
||||
return math.sqrt(sum(s * s for s in samples) / len(samples))
|
||||
|
||||
|
||||
class EnergyVAD:
|
||||
def __init__(self, sample_rate: int = 16000, threshold: float = 500.0, silence_ms: float = 700.0):
|
||||
self.sample_rate = sample_rate
|
||||
self.threshold = threshold
|
||||
self.silence_ms = silence_ms
|
||||
self._speech_started = False
|
||||
self._silence_ms = 0.0
|
||||
|
||||
def feed(self, pcm: bytes) -> bool:
|
||||
"""Verarbeitet einen Audio-Chunk.
|
||||
|
||||
Liefert True, sobald nach erkannter Sprache genug Stille (silence_ms)
|
||||
vergangen ist - die Aeusserung gilt dann als beendet.
|
||||
"""
|
||||
level = rms(pcm)
|
||||
n_samples = len(pcm) // 2
|
||||
chunk_ms = (n_samples / self.sample_rate) * 1000.0 if self.sample_rate else 0.0
|
||||
|
||||
if level >= self.threshold:
|
||||
self._speech_started = True
|
||||
self._silence_ms = 0.0
|
||||
return False
|
||||
|
||||
if self._speech_started:
|
||||
self._silence_ms += chunk_ms
|
||||
if self._silence_ms >= self.silence_ms:
|
||||
return True
|
||||
return False
|
||||
|
||||
def reset(self) -> None:
|
||||
self._speech_started = False
|
||||
self._silence_ms = 0.0
|
||||
Loading…
Add table
Add a link
Reference in a new issue