Problem (Stefan): laeuft Musik, endet eine Aufnahme nie von selbst — die Energie bleibt oben, der RMS-Stille-Endpoint feuert nicht (last_voice_at immer frisch). Erst Stop druecken oder Musik leiser (dann echte Stille) beendet. Fix: waehrend lauter Phasen (rms ueber Schwelle) periodisch (~700ms, gedrosselt) mit Silero pruefen, ob im letzten endpoint_ms-Fenster ueberhaupt Sprache ist. Nur Musik/Stille → Turn beenden (finalize endpoint). Real gesprochene Turns haben Sprache im Fenster → laufen normal weiter; eine kurze Denk-Pause auch, solange innerhalb der Toleranz noch Sprache im Fenster liegt. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
826 lines
40 KiB
Python
826 lines
40 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
ARIA Voxtral-STT-3B Bridge (Transformers) — Ersatz fuer whisper.
|
|
|
|
Laeuft auf Treiber 550/CUDA 12.4 via torch cu124 (kein Treiber-Upgrade noetig).
|
|
Modell: Voxtral-Mini-3B-2507 (bf16, ~9 GB) → GPU 1 (12 GB, per Compose gepinnt).
|
|
|
|
Arbeitsweise = Zwilling der whisper-Bridge: App schickt live PCM-Chunks; wir
|
|
transkribieren alle ~STREAM_TRANSCRIBE_INTERVAL_MS auf dem Ringbuffer (Partials)
|
|
und feuern stt_endpoint, sobald der ADAPTIVE Endpointer (Rausch-Boden-VAD +
|
|
semantische Stagnation, aus M0.1) "fertig" sagt. RVS-Wire-Protokoll identisch zu
|
|
whisper → drop-in (die App merkt nur bessere Genauigkeit).
|
|
|
|
⚠️ VERIFY-ON-FIRST-RUN: Die exakte Transformers-Transkriptions-API von Voxtral
|
|
(apply_transcription_request / generate / decode) ist unten in EINER Methode
|
|
(VoxtralRunner._transcribe_blocking) gekapselt und nach dem HF-Modelcard-Muster
|
|
modelliert. Beim ersten echten Lauf gegen die Voxtral-Modelcard pruefen und dort
|
|
anpassen. Alles andere (RVS, Endpointer) ist bewaehrt.
|
|
|
|
Env:
|
|
RVS_HOST, RVS_PORT, RVS_TLS, RVS_TLS_FALLBACK, RVS_TOKEN
|
|
VOXTRAL_MODEL Default: mistralai/Voxtral-Mini-3B-2507
|
|
VOXTRAL_LANGUAGE Default: de
|
|
VOXTRAL_DEVICE Default: cuda
|
|
STREAM_TRANSCRIBE_INTERVAL_MS Default 1000 (3B ist schwerer als whisper-small)
|
|
"""
|
|
import asyncio
|
|
import base64
|
|
import json
|
|
import logging
|
|
import os
|
|
import re
|
|
import tempfile
|
|
import time
|
|
from dataclasses import dataclass, field
|
|
from typing import Optional
|
|
|
|
import numpy as np
|
|
import soundfile as sf
|
|
import websockets
|
|
|
|
import speaker_id # Speaker-ID (nur Stefans Stimme) — portiert aus der whisper-Bridge
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [%(levelname)s] %(message)s",
|
|
datefmt="%H:%M:%S",
|
|
)
|
|
logger = logging.getLogger("voxtral-bridge")
|
|
|
|
RVS_HOST = os.getenv("RVS_HOST", "").strip()
|
|
RVS_PORT = int(os.getenv("RVS_PORT", "443"))
|
|
RVS_TLS = os.getenv("RVS_TLS", "true").lower() == "true"
|
|
RVS_TLS_FALLBACK = os.getenv("RVS_TLS_FALLBACK", "true").lower() == "true"
|
|
RVS_TOKEN = os.getenv("RVS_TOKEN", "").strip()
|
|
|
|
VOXTRAL_MODEL = os.getenv("VOXTRAL_MODEL", "mistralai/Voxtral-Mini-3B-2507")
|
|
VOXTRAL_LANGUAGE = os.getenv("VOXTRAL_LANGUAGE", "de")
|
|
VOXTRAL_DEVICE = os.getenv("VOXTRAL_DEVICE", "cuda")
|
|
|
|
STREAM_TRANSCRIBE_INTERVAL_MS = int(os.getenv("STREAM_TRANSCRIBE_INTERVAL_MS", "1000"))
|
|
STREAM_DEFAULT_ENDPOINT_MS = 2400
|
|
STREAM_DEFAULT_HARD_CAP_MS = 300000
|
|
STREAM_MIN_AUDIO_MS = 600
|
|
STREAM_SPEAKER_CHECK_MS = 1500 # ab so viel Audio einmalig Speaker-ID pruefen
|
|
STREAM_SESSION_TTL_S = 120
|
|
STREAM_ENERGY_WINDOW_MS = 300
|
|
STREAM_SEMANTIC_BACKUP_FACTOR = 2.0
|
|
# Adaptiver Voice-Schwellwert (M0.1): relativ zum gemessenen Rausch-Boden.
|
|
STREAM_VOICE_FACTOR = 2.5
|
|
STREAM_VOICE_RMS_MIN = 0.005
|
|
STREAM_VOICE_RMS_MAX = 0.020
|
|
# Mindest-Stimme (in ~200ms-Endpointer-Frames), ab der eine Aufnahme ueberhaupt
|
|
# als Sprache gilt. Darunter = Stille / kurzer Geraeusch-Blip → KEIN Transkript
|
|
# (Voxtral halluziniert aus Fast-Nichts sonst einen Fuellsatz). 2 ≈ 400ms.
|
|
STREAM_MIN_VOICED_FRAMES = int(os.getenv("STREAM_MIN_VOICED_FRAMES", "2"))
|
|
|
|
# Halluzinations-Filter (2. Netz NACH der Transkription). Der voiced_frames-Guard
|
|
# oben faengt die reine Stille; hier kommt das "borderline"-Band dazu: wenn wenig
|
|
# echte Stimme da war UND das Transkript ein bekanntes Voxtral-Silence-Artefakt
|
|
# ist (Untertitel-Credits, Staedte-/Geo-Fakten "Flaeche von X km2"), ist es fast
|
|
# sicher ein Phantom aus Fast-Nichts → verwerfen. Gegated auf wenig voiced_frames,
|
|
# damit eine ECHTE Geografie-Frage (die hat normale Stimm-Energie) durchgeht.
|
|
STREAM_HALLUC_GUARD_FRAMES = int(os.getenv("STREAM_HALLUC_GUARD_FRAMES",
|
|
str(STREAM_MIN_VOICED_FRAMES * 4))) # ~1.6s
|
|
_HALLUCINATION_RE = re.compile(
|
|
r"untertitel"
|
|
r"|amara\.org"
|
|
r"|vielen\s+dank\s+f[uü]r'?s?\s+(zuschauen|zusehen|zuh[oö]ren)"
|
|
r"|bis\s+zum\s+n[aä]chsten\s+mal"
|
|
r"|abonnier"
|
|
r"|fl[aä]che\s+von\s+[\d.,]+\s*(km|quadratkilometer)"
|
|
r"|[\d.,]+\s*(km²|quadratkilometern?|einwohnern?)\b",
|
|
re.IGNORECASE,
|
|
)
|
|
|
|
# Kollabiert unmittelbar wiederholte Phrasen (Voxtral-Repetition-Loop) auf EINE
|
|
# Kopie. Zweites Netz hinter no_repeat_ngram in der Generation. Phrase 5-80 Zeichen,
|
|
# 3+ mal hintereinander → eine. Kurze legitime Doppelungen ('ja ja', 'sehr sehr')
|
|
# bleiben (Unit < 5 Zeichen bzw. < 3 Wiederholungen).
|
|
_REPEAT_RE = re.compile(r"(.{5,80}?)(?:\s*\1){2,}", re.IGNORECASE | re.DOTALL)
|
|
|
|
|
|
def _collapse_repetitions(text: str) -> str:
|
|
if not text:
|
|
return text
|
|
out = text
|
|
for _ in range(3): # mehrfach fuer verschachtelte/ungleiche Loops
|
|
new = _REPEAT_RE.sub(r"\1", out)
|
|
if new == out:
|
|
break
|
|
out = new
|
|
return out.strip()
|
|
|
|
|
|
# ── Silero VAD: echte Sprach-Erkennung VOR dem Transkribieren ──────────────
|
|
# Der Muster-Filter oben kennt nur spezifische Artefakte. Generische Phantome
|
|
# ("Ich bin ein guter Mann" aus Fast-Stille) kann ein Text-Regex nicht fangen —
|
|
# aber ein VAD schon, weil es am AUDIO entscheidet, nicht am Text. Silero trennt
|
|
# Sprache zuverlaessig von Stille / Rauschen / MUSIK. Kein Speech-Segment →
|
|
# no-speech → nicht transkribieren → kein Phantom (und Musik/Instrumental fliegt
|
|
# gleich mit raus).
|
|
# FAIL-OPEN: klappt das VAD nicht (Import/Load/Inferenz), wird trotzdem normal
|
|
# transkribiert. Die STT darf NIE komplett sterben (Speaker-ID-Lektion).
|
|
SILERO_VAD_ENABLED = os.getenv("SILERO_VAD_ENABLED", "true").lower() in ("1", "true", "yes")
|
|
SILERO_VAD_THRESHOLD = float(os.getenv("SILERO_VAD_THRESHOLD", "0.5"))
|
|
SILERO_MIN_SPEECH_MS = int(os.getenv("SILERO_MIN_SPEECH_MS", "150"))
|
|
SILERO_PAD_MS = int(os.getenv("SILERO_PAD_MS", "200"))
|
|
# Speech-Endpoint gegen laute Umgebungsmusik: der RMS-Stille-Endpoint feuert bei
|
|
# durchgehender Musik NIE (Energie bleibt oben). Deshalb waehrend lauter Phasen
|
|
# periodisch (alle X ms) mit Silero pruefen, ob im letzten endpoint_ms-Fenster
|
|
# ueberhaupt noch Sprache ist — wenn nicht (nur Musik/Stille), Turn beenden.
|
|
STREAM_SPEECH_ENDPOINT_CHECK_MS = int(os.getenv("STREAM_SPEECH_ENDPOINT_CHECK_MS", "700"))
|
|
|
|
_vad_state = {"model": None, "get_ts": None, "failed": False}
|
|
|
|
|
|
def _speech_segments(audio_f32):
|
|
"""Silero-VAD-Sprachsegmente (Liste von {start,end} Sample-Indizes) im
|
|
16kHz-float32-Audio. Rueckgabe:
|
|
[] → kein Speech (Stille/Rauschen/Musik) → Phantom-Verdacht, verwerfen.
|
|
[...] → Speech vorhanden.
|
|
None → VAD nicht verfuegbar → fail-open (Aufrufer transkribiert normal)."""
|
|
if not SILERO_VAD_ENABLED or _vad_state["failed"]:
|
|
return None
|
|
if _vad_state["model"] is None:
|
|
try:
|
|
from silero_vad import load_silero_vad, get_speech_timestamps
|
|
_vad_state["model"] = load_silero_vad()
|
|
_vad_state["get_ts"] = get_speech_timestamps
|
|
logger.info("Silero VAD geladen (threshold=%.2f, min_speech=%dms)",
|
|
SILERO_VAD_THRESHOLD, SILERO_MIN_SPEECH_MS)
|
|
except Exception:
|
|
logger.exception("Silero VAD Laden fehlgeschlagen — dauerhaft aus (fail-open)")
|
|
_vad_state["failed"] = True
|
|
return None
|
|
try:
|
|
import torch as _torch
|
|
segs = _vad_state["get_ts"](
|
|
_torch.from_numpy(audio_f32), _vad_state["model"],
|
|
sampling_rate=16000, threshold=SILERO_VAD_THRESHOLD,
|
|
min_speech_duration_ms=SILERO_MIN_SPEECH_MS,
|
|
)
|
|
return segs or []
|
|
except Exception:
|
|
logger.exception("Silero VAD Inferenz fehlgeschlagen — dieser Turn fail-open")
|
|
return None
|
|
|
|
# Speaker-ID Gating global an/aus. DEFAULT AUS (fail-open) — die "nur meine Stimme"-
|
|
# Pruefung ist ein BEWUSSTER Schalter, kein Automatismus: ein einziger schlechter
|
|
# Enroll darf nie die ganze STT lahmlegen (genau das ist passiert). Wird per config-
|
|
# Broadcast (voiceIdEnabled, aus dem Diagnostic) zur Laufzeit gesetzt. Kann per ENV
|
|
# vorbelegt werden.
|
|
SPEAKER_ID_ENABLED = os.getenv("VOICE_ID_ENABLED", "false").lower() in ("1", "true", "yes")
|
|
|
|
|
|
def _set_speaker_id_enabled(val: bool) -> None:
|
|
global SPEAKER_ID_ENABLED
|
|
SPEAKER_ID_ENABLED = bool(val)
|
|
|
|
|
|
def pcm_s16le_to_float32(data: bytes) -> np.ndarray:
|
|
if not data:
|
|
return np.zeros(0, dtype=np.float32)
|
|
return np.frombuffer(data, dtype=np.int16).astype(np.float32) / 32768.0
|
|
|
|
|
|
async def _send(ws, mtype: str, payload: dict) -> None:
|
|
try:
|
|
await ws.send(json.dumps({
|
|
"type": mtype, "payload": payload, "timestamp": int(time.time() * 1000),
|
|
}))
|
|
except Exception as e:
|
|
logger.warning("RVS-Send fehlgeschlagen (%s): %s", mtype, e)
|
|
|
|
|
|
class VoxtralRunner:
|
|
"""Haelt das Voxtral-Modell (Transformers). transcribe() blockiert → aus dem
|
|
Event-Loop via run_in_executor aufrufen. Ein Lock serialisiert GPU-Zugriffe."""
|
|
|
|
def __init__(self) -> None:
|
|
self.model = None
|
|
self.processor = None
|
|
self._lock = asyncio.Lock()
|
|
|
|
def load(self) -> None:
|
|
import torch
|
|
from transformers import AutoProcessor, VoxtralForConditionalGeneration
|
|
t0 = time.time()
|
|
logger.info("Lade Voxtral '%s' (device=%s, bf16)…", VOXTRAL_MODEL, VOXTRAL_DEVICE)
|
|
self.processor = AutoProcessor.from_pretrained(VOXTRAL_MODEL)
|
|
self.model = VoxtralForConditionalGeneration.from_pretrained(
|
|
VOXTRAL_MODEL, torch_dtype=torch.bfloat16, device_map=VOXTRAL_DEVICE,
|
|
)
|
|
logger.info("Voxtral geladen in %.1fs", time.time() - t0)
|
|
|
|
def _transcribe_blocking(self, audio_f32: np.ndarray, language: str) -> str:
|
|
import torch
|
|
proc, model = self.processor, self.model
|
|
if proc is None or model is None or audio_f32.size == 0:
|
|
return ""
|
|
# VoxtralProcessor verlangt bei rohen Arrays ein 'format'. Robuster:
|
|
# in ein temp-WAV schreiben und den PFAD uebergeben — der Processor liest
|
|
# Format + Samplerate selbst, kein 'format'-Argument noetig.
|
|
wav_path = None
|
|
try:
|
|
with tempfile.NamedTemporaryFile(suffix=".wav", delete=False) as tf:
|
|
wav_path = tf.name
|
|
sf.write(wav_path, audio_f32, 16000, subtype="PCM_16")
|
|
inputs = proc.apply_transcription_request(
|
|
language=language, audio=wav_path, model_id=VOXTRAL_MODEL,
|
|
)
|
|
inputs = inputs.to(VOXTRAL_DEVICE, dtype=torch.bfloat16)
|
|
with torch.no_grad():
|
|
# hoch genug fuer lange Diktate (stoppt eh am EOS); 512 hat
|
|
# mehrminutige Aufnahmen abgeschnitten.
|
|
# Repetition-Bremse: Voxtral kippt bei Stille/Rauschen am Ende
|
|
# gern in eine Schleife und wiederholt einen Satz zig-mal
|
|
# ("Vergiss das, das ist nur... Vergiss das, das ist nur..."
|
|
# x15). no_repeat_ngram_size=4 laesst die ERSTE echte Nennung
|
|
# durch, verbietet aber die exakte 4-Gramm-Wiederholung → Loop
|
|
# bricht ab; repetition_penalty daempft zusaetzlich. Beides mild,
|
|
# damit normale Sprache (auch mal ein doppeltes Wort) unberuehrt
|
|
# bleibt.
|
|
outputs = model.generate(
|
|
**inputs,
|
|
max_new_tokens=4096,
|
|
no_repeat_ngram_size=4,
|
|
repetition_penalty=1.15,
|
|
)
|
|
trimmed = outputs[:, inputs.input_ids.shape[1]:]
|
|
text = proc.batch_decode(trimmed, skip_special_tokens=True)
|
|
return (text[0] if text else "").strip()
|
|
finally:
|
|
if wav_path:
|
|
try:
|
|
os.unlink(wav_path)
|
|
except Exception:
|
|
pass
|
|
|
|
async def transcribe(self, audio_f32: np.ndarray, language: str) -> str:
|
|
loop = asyncio.get_running_loop()
|
|
async with self._lock:
|
|
return await loop.run_in_executor(None, self._transcribe_blocking, audio_f32, language)
|
|
|
|
|
|
@dataclass
|
|
class StreamSession:
|
|
request_id: str
|
|
audio_request_id: str
|
|
language: str
|
|
endpoint_ms: int
|
|
hard_cap_ms: int
|
|
voice: str = ""
|
|
speed: float = 1.0
|
|
interrupted: bool = False
|
|
location: Optional[dict] = None
|
|
sample_rate: int = 16000
|
|
voice_factor: float = STREAM_VOICE_FACTOR
|
|
voice_rms_min: float = STREAM_VOICE_RMS_MIN
|
|
voice_rms_max: float = STREAM_VOICE_RMS_MAX
|
|
pcm_buffer: bytearray = field(default_factory=bytearray)
|
|
started_at: float = field(default_factory=time.time)
|
|
last_chunk_at: float = field(default_factory=time.time)
|
|
last_partial: str = ""
|
|
last_growth_at: float = 0.0
|
|
last_transcribe_at: float = 0.0
|
|
last_voice_at: float = 0.0
|
|
last_speech_check_at: float = 0.0 # Drossel fuer den Silero-Speech-Endpoint
|
|
noise_floor: float = 0.0
|
|
closed: bool = False
|
|
endpoint_sent: bool = False
|
|
# Einmaliges "Sprache erkannt"-Signal an die App gesendet? Voxtral schickt
|
|
# keine Live-Partials, aber der App-No-Speech-Watchdog wartet auf ein
|
|
# stt_partial, um "der User redet" zu erkennen — sonst cancelt er mitten im
|
|
# Satz. Wir feuern EIN leeres stt_partial beim ersten Voice-Frame.
|
|
speech_signaled: bool = False
|
|
# Anzahl Endpointer-Frames (~200ms) mit echter Stimme. Gate gegen Halluzination
|
|
# aus Stille/Blips: unter STREAM_MIN_VOICED_FRAMES wird nicht transkribiert.
|
|
voiced_frames: int = 0
|
|
# Speaker-ID Gating (einmalig auf die ersten ~1.5s der Aufnahme)
|
|
speaker_checked: bool = False
|
|
speaker_match: Optional[bool] = None
|
|
speaker_similarity: float = 0.0
|
|
|
|
|
|
class SessionManager:
|
|
def __init__(self, runner: VoxtralRunner) -> None:
|
|
self.runner = runner
|
|
self._sessions: dict[str, StreamSession] = {}
|
|
self._ws = None
|
|
|
|
def attach_ws(self, ws) -> None:
|
|
self._ws = ws
|
|
|
|
def start_session(self, payload: dict) -> None:
|
|
rid = (payload.get("requestId") or "").strip()
|
|
if not rid:
|
|
return
|
|
try:
|
|
endpoint_ms = int(payload.get("endpointMs") or STREAM_DEFAULT_ENDPOINT_MS)
|
|
except (TypeError, ValueError):
|
|
endpoint_ms = STREAM_DEFAULT_ENDPOINT_MS
|
|
try:
|
|
hard_cap_ms = int(payload.get("hardCapMs") or STREAM_DEFAULT_HARD_CAP_MS)
|
|
except (TypeError, ValueError):
|
|
hard_cap_ms = STREAM_DEFAULT_HARD_CAP_MS
|
|
try:
|
|
voice_factor = float(payload.get("voiceFactor") or STREAM_VOICE_FACTOR)
|
|
except (TypeError, ValueError):
|
|
voice_factor = STREAM_VOICE_FACTOR
|
|
self._sessions[rid] = StreamSession(
|
|
request_id=rid,
|
|
audio_request_id=payload.get("audioRequestId", "") or "",
|
|
language=payload.get("language") or VOXTRAL_LANGUAGE,
|
|
endpoint_ms=endpoint_ms,
|
|
hard_cap_ms=hard_cap_ms,
|
|
voice=payload.get("voice", "") or "",
|
|
speed=float(payload.get("speed") or 1.0),
|
|
voice_factor=voice_factor,
|
|
interrupted=bool(payload.get("interrupted", False)),
|
|
location=payload.get("location") or None,
|
|
sample_rate=int(payload.get("sampleRate") or 16000),
|
|
)
|
|
logger.info("Voxtral-Session offen: id=%s lang=%s endpointMs=%d",
|
|
rid[:8], self._sessions[rid].language, endpoint_ms)
|
|
|
|
def feed_chunk(self, payload: dict) -> bool:
|
|
sess = self._sessions.get(payload.get("requestId", ""))
|
|
if sess is None or sess.closed:
|
|
return False
|
|
pcm_b64 = payload.get("pcm", "")
|
|
if pcm_b64:
|
|
try:
|
|
sess.pcm_buffer.extend(base64.b64decode(pcm_b64))
|
|
except Exception:
|
|
pass
|
|
sess.last_chunk_at = time.time()
|
|
return True
|
|
|
|
def end_session(self, request_id: str) -> None:
|
|
sess = self._sessions.get(request_id)
|
|
if sess is not None:
|
|
sess.closed = True
|
|
|
|
def drop(self, request_id: str) -> None:
|
|
self._sessions.pop(request_id, None)
|
|
|
|
# ── Endpointer (adaptiv, M0.1) ──
|
|
def _buffer_ms(self, sess: StreamSession) -> float:
|
|
samples = len(sess.pcm_buffer) // 2
|
|
return (samples / sess.sample_rate) * 1000.0 if samples else 0.0
|
|
|
|
def _tail_rms(self, sess: StreamSession) -> float:
|
|
win = int(sess.sample_rate * STREAM_ENERGY_WINDOW_MS / 1000) * 2
|
|
if win <= 0:
|
|
return 0.0
|
|
tail = sess.pcm_buffer[-win:]
|
|
if len(tail) < 2:
|
|
return 0.0
|
|
arr = pcm_s16le_to_float32(bytes(tail))
|
|
return float(np.sqrt(np.mean(arr * arr))) if arr.size else 0.0
|
|
|
|
def _voice_threshold(self, sess: StreamSession) -> float:
|
|
nf = sess.noise_floor
|
|
if nf <= 0.0:
|
|
return sess.voice_rms_min
|
|
return min(max(nf * sess.voice_factor, sess.voice_rms_min), sess.voice_rms_max)
|
|
|
|
def _update_noise_floor(self, sess: StreamSession, rms: float) -> None:
|
|
nf = sess.noise_floor
|
|
if nf <= 0.0:
|
|
sess.noise_floor = rms
|
|
elif rms < nf:
|
|
sess.noise_floor = 0.90 * nf + 0.10 * rms
|
|
else:
|
|
sess.noise_floor = 0.98 * nf + 0.02 * rms
|
|
|
|
async def _check_speaker(self, sess: StreamSession) -> None:
|
|
"""Einmalig: erste ~1.5s → Embedding → Vergleich mit Fingerprint.
|
|
Ohne Fingerprint fail-open (match=True). Bei Mismatch: Session beenden."""
|
|
sess.speaker_checked = True
|
|
# Schalter aus (Default) → gar keine Pruefung, alles durchlassen.
|
|
if not SPEAKER_ID_ENABLED:
|
|
sess.speaker_match = True
|
|
return
|
|
head = bytes(sess.pcm_buffer[: STREAM_SPEAKER_CHECK_MS * 32])
|
|
if len(head) < speaker_id.MIN_SAMPLE_BYTES:
|
|
sess.speaker_match = True
|
|
return
|
|
try:
|
|
loop = asyncio.get_running_loop()
|
|
is_match, sim = await loop.run_in_executor(None, speaker_id.verify, head)
|
|
except Exception as exc:
|
|
logger.warning("Stream %s: speaker-check crashed (%s) — fail-open",
|
|
sess.request_id[:8], exc)
|
|
sess.speaker_match = True
|
|
return
|
|
sess.speaker_match = is_match
|
|
sess.speaker_similarity = sim
|
|
logger.info("Stream %s: speaker-check sim=%.2f → %s (thr=%.2f)",
|
|
sess.request_id[:8], sim, "MATCH" if is_match else "REJECT",
|
|
speaker_id.DEFAULT_THRESHOLD)
|
|
if not is_match:
|
|
await self._finalize_speaker_mismatch(sess, sim)
|
|
|
|
async def _finalize_speaker_mismatch(self, sess: StreamSession, similarity: float) -> None:
|
|
"""Fremde Stimme: synthetisches leeres stt_endpoint (reason=speaker_mismatch),
|
|
Session droppen — kein Voxtral-Transcribe, kein Brain-Call."""
|
|
if sess.endpoint_sent:
|
|
return
|
|
sess.endpoint_sent = True
|
|
duration_s = self._buffer_ms(sess) / 1000.0
|
|
logger.info("Stream %s: speaker-mismatch (sim=%.2f) — DROP nach %.1fs",
|
|
sess.request_id[:8], similarity, duration_s)
|
|
if self._ws is not None:
|
|
payload = {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": "speaker_mismatch",
|
|
"durationS": duration_s, "sttMs": 0,
|
|
"voice": sess.voice, "speed": sess.speed,
|
|
"interrupted": sess.interrupted,
|
|
"speakerSimilarity": float(similarity),
|
|
}
|
|
if sess.location:
|
|
payload["location"] = sess.location
|
|
await _send(self._ws, "stt_endpoint", payload)
|
|
await _send(self._ws, "stt_stream_done", {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": "speaker_mismatch",
|
|
})
|
|
self.drop(sess.request_id)
|
|
|
|
async def run_endpointer(self) -> None:
|
|
logger.info("Voxtral-Endpointer gestartet (adaptiver VAD, interval=%dms)",
|
|
STREAM_TRANSCRIBE_INTERVAL_MS)
|
|
while True:
|
|
await asyncio.sleep(0.2)
|
|
now = time.time()
|
|
for sid, sess in list(self._sessions.items()):
|
|
try:
|
|
await self._tick(sess, now)
|
|
except Exception:
|
|
logger.exception("Tick crashed (session=%s)", sid[:8])
|
|
for sid, sess in list(self._sessions.items()):
|
|
if now - sess.last_chunk_at > STREAM_SESSION_TTL_S:
|
|
logger.info("Stream %s: TTL — drop", sid[:8])
|
|
self.drop(sid)
|
|
|
|
async def _tick(self, sess: StreamSession, now: float) -> None:
|
|
if sess.endpoint_sent:
|
|
return
|
|
if (now - sess.started_at) * 1000.0 > sess.hard_cap_ms and not sess.closed:
|
|
await self._finalize(sess, "hardcap")
|
|
return
|
|
if sess.closed:
|
|
await self._finalize(sess, "stream_end")
|
|
return
|
|
if self._buffer_ms(sess) < STREAM_MIN_AUDIO_MS:
|
|
return
|
|
# Speaker-ID einmalig: ist es Stefans Stimme? Fremde → Session verwerfen
|
|
# (kein Transcribe, kein Brain-Call). Ohne Enrollment fail-open.
|
|
if not sess.speaker_checked and self._buffer_ms(sess) >= STREAM_SPEAKER_CHECK_MS:
|
|
await self._check_speaker(sess)
|
|
if sess.speaker_match is False:
|
|
return
|
|
# Adaptive akustische Sprach-Aktivitaet (M0.1). KEINE Live-Partials mehr:
|
|
# Voxtral-3B transkribiert den ganzen WACHSENDEN Buffer und braucht dafuer
|
|
# bei langen Aufnahmen 5-6 s — zu langsam fuer Live-Text, UND diese Latenz
|
|
# hat den semantischen Endpoint faelschlich ausgeloest (Partial-Latenz >
|
|
# Timeout → willkuerliche Abbrueche nach 20-40 s). Deshalb: Turn-Ende rein
|
|
# AKUSTISCH, transkribiert wird nur EINMAL im _finalize.
|
|
rms = self._tail_rms(sess)
|
|
if rms >= self._voice_threshold(sess):
|
|
sess.last_voice_at = now
|
|
sess.voiced_frames += 1
|
|
# Einmalig der App melden, dass Sprache begonnen hat — aber ERST ab genug
|
|
# echter Stimme (>= STREAM_MIN_VOICED_FRAMES). Ein einzelner Geraeusch-
|
|
# Blip darf den No-Speech-Watchdog NICHT loeschen, sonst transkribiert
|
|
# Voxtral das Fast-Nichts und HALLUZINIERT einen Phantom-Satz. Ohne Live-
|
|
# Partials wuerde der Watchdog die Aufnahme sonst am Konversationsfenster
|
|
# canceln, obwohl der User redet ("beendet nach ~4s"-Repro). Leeres
|
|
# stt_partial: App setzt streamGotPartial=true + loescht den Watchdog.
|
|
# Nach der Speaker-ID-Pruefung (oben) → fremde Stimmen signalisieren NICHT.
|
|
if (not sess.speech_signaled and self._ws is not None
|
|
and sess.voiced_frames >= STREAM_MIN_VOICED_FRAMES):
|
|
sess.speech_signaled = True
|
|
await _send(self._ws, "stt_partial", {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "",
|
|
})
|
|
# Speech-Endpoint gegen laute Umgebungsmusik: es ist gerade laut (rms
|
|
# ueber Schwelle) — aber ist es Sprache oder Musik? Der RMS-Endpoint
|
|
# unten wuerde bei Musik NIE feuern (last_voice_at bleibt frisch).
|
|
# Deshalb gedrosselt mit Silero das letzte endpoint_ms-Fenster pruefen:
|
|
# KEINE Sprache drin (nur Musik) → Turn ist zu Ende. Real gesprochene
|
|
# Turns haben Sprache im Fenster → laufen weiter.
|
|
if (self._buffer_ms(sess) >= sess.endpoint_ms
|
|
and (now - sess.last_speech_check_at) * 1000.0 >= STREAM_SPEECH_ENDPOINT_CHECK_MS):
|
|
sess.last_speech_check_at = now
|
|
try:
|
|
tail_bytes = int(sess.endpoint_ms / 1000.0 * sess.sample_rate) * 2
|
|
tail = pcm_s16le_to_float32(bytes(sess.pcm_buffer[-tail_bytes:]))
|
|
segs = _speech_segments(tail)
|
|
if segs is not None and len(segs) == 0:
|
|
logger.info("Stream %s: Speech-Endpoint (keine Sprache im letzten %dms — Musik/Stille) → finalize",
|
|
sess.request_id[:8], sess.endpoint_ms)
|
|
await self._finalize(sess, "endpoint")
|
|
return
|
|
except Exception:
|
|
logger.exception("Speech-Endpoint-Check fehlgeschlagen — ignoriert")
|
|
else:
|
|
self._update_noise_floor(sess, rms)
|
|
# No-Speech-Timeout: wurde die GANZE Zeit KEINE Stimme erkannt
|
|
# (last_voice_at==0), feuert der normale Endpoint unten NIE — der braucht
|
|
# last_voice_at>0. Ohne das bleibt ein reines Stille-Fenster offen bis
|
|
# Hardcap/manuellem Stop → genau Stefans Repro: "die Stille-Ende wird nie
|
|
# erreicht, stop ich selbst ist es weg". Nach endpoint_ms Stille ab Start
|
|
# schliessen wir das Fenster selbst als no-speech (leer, lautlos, zurueck
|
|
# aufs Wake-Word). voiced_frames==0 → _finalize verwirft ohne Transkript,
|
|
# also KEIN Phantom.
|
|
if sess.last_voice_at == 0 and (now - sess.started_at) * 1000.0 >= sess.endpoint_ms:
|
|
await self._finalize(sess, "no_speech")
|
|
return
|
|
# Endpoint: hat der User schon gesprochen UND ist es seit endpoint_ms still?
|
|
if sess.last_voice_at > 0 and (now - sess.last_voice_at) * 1000.0 >= sess.endpoint_ms:
|
|
await self._finalize(sess, "endpoint")
|
|
|
|
async def _emit_no_speech(self, sess: "StreamSession", reason_label: str) -> None:
|
|
"""Leeres no-speech-Endpoint senden + Session droppen (kein Transkript).
|
|
App re-armt still, zurueck aufs Wake-Word."""
|
|
if self._ws is not None:
|
|
payload = {"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": reason_label,
|
|
"durationS": 0.0, "sttMs": 0}
|
|
await _send(self._ws, "stt_endpoint", payload)
|
|
await _send(self._ws, "stt_stream_done", {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": reason_label})
|
|
self.drop(sess.request_id)
|
|
|
|
async def _finalize(self, sess: StreamSession, reason: str) -> None:
|
|
if sess.endpoint_sent:
|
|
return
|
|
sess.endpoint_sent = True
|
|
# Halluzinations-Guard: zu wenig echte Stimme (Stille / kurzer Blip im
|
|
# Passiv-/Wake-Fenster) → NICHT transkribieren. Voxtral (wie Whisper) baut
|
|
# aus Fast-Nichts gern einen Fuellsatz ("Die Stadt hat eine Flaeche von
|
|
# 1,5 km2"), der dann als PHANTOM-Nachricht ans Brain geht und das Gespraech
|
|
# entgleisen laesst (Stefans Repro: "kam Nachricht von mir, obwohl ich
|
|
# nichts sagte"). Leeres Endpoint = no-speech → App re-armt still.
|
|
#
|
|
# WICHTIG (aus dem ai-box-Log gelernt): die Phantome kommen mit
|
|
# reason=stream_end — Passiv-/Wake-Fenster enden AUCH per stream_end, wenn
|
|
# sie auf Stille zumachen. stream_end ist also NICHT gleich "manueller Stop".
|
|
# Deshalb greift der Guard jetzt auch bei stream_end, aber mit niedrigerer
|
|
# Schwelle (voiced==0 = gar keine Stimme), damit ein kurzes bewusstes Wort
|
|
# ('ja', 'stopp') am Aufnahme-Button noch durchgeht, echte Stille aber nicht.
|
|
_min_voiced = STREAM_MIN_VOICED_FRAMES if reason != "stream_end" else 1
|
|
if sess.voiced_frames < _min_voiced:
|
|
logger.info("Stream %s: no-speech (voiced_frames=%d<%d, reason=%s) — leeres Endpoint",
|
|
sess.request_id[:8], sess.voiced_frames, _min_voiced, reason)
|
|
if self._ws is not None:
|
|
nospeech = {"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": f"no_speech:{reason}",
|
|
"durationS": 0.0, "sttMs": 0}
|
|
await _send(self._ws, "stt_endpoint", nospeech)
|
|
await _send(self._ws, "stt_stream_done", {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": f"no_speech:{reason}"})
|
|
self.drop(sess.request_id)
|
|
return
|
|
audio = pcm_s16le_to_float32(bytes(sess.pcm_buffer))
|
|
|
|
# Silero VAD: ist ueberhaupt echte Sprache im Audio? Das entscheidet am
|
|
# AUDIO, nicht am Text — faengt also generische Phantome ("Ich bin ein
|
|
# guter Mann") UND Musik/Rauschen, die der Muster-Filter nicht kennt.
|
|
# Kein Speech-Segment → no-speech, gar nicht erst transkribieren.
|
|
# fail-open: segs=None (VAD nicht verfuegbar) → normal weiter.
|
|
segs = _speech_segments(audio)
|
|
if segs is not None and len(segs) == 0:
|
|
logger.info("Stream %s: Silero VAD — keine Sprache (%.1fs, reason=%s) → no-speech",
|
|
sess.request_id[:8], audio.size / 16000.0, reason)
|
|
await self._emit_no_speech(sess, f"vad_no_speech:{reason}")
|
|
return
|
|
if segs:
|
|
# Auf die Sprach-Spanne trimmen (Stille-Raender weg → Voxtral
|
|
# halluziniert an den Enden weniger). Kleiner Pad gegen abgeschnittene
|
|
# leise Wort-Anfaenge/-Enden.
|
|
pad = int(SILERO_PAD_MS / 1000.0 * 16000)
|
|
s0 = max(0, segs[0]["start"] - pad)
|
|
s1 = min(int(audio.size), segs[-1]["end"] + pad)
|
|
if s1 > s0 and (s1 - s0) < audio.size:
|
|
audio = audio[s0:s1]
|
|
|
|
t0 = time.time()
|
|
try:
|
|
final_text = (await self.runner.transcribe(audio, sess.language)).strip()
|
|
except Exception:
|
|
logger.exception("Stream %s: Final-Transcribe crashed", sess.request_id[:8])
|
|
final_text = sess.last_partial
|
|
stt_ms = int((time.time() - t0) * 1000)
|
|
duration_s = audio.size / 16000.0
|
|
# Repetition-Loop einkassieren, falls trotz no_repeat_ngram was durchkam.
|
|
_collapsed = _collapse_repetitions(final_text)
|
|
if _collapsed != final_text:
|
|
logger.info("Stream %s: Repetition-Loop kollabiert (%d→%d Zeichen)",
|
|
sess.request_id[:8], len(final_text), len(_collapsed))
|
|
final_text = _collapsed
|
|
logger.info("Stream %s: FINAL (reason=%s, %.1fs, %dms): %r",
|
|
sess.request_id[:8], reason, duration_s, stt_ms, final_text[:120])
|
|
|
|
# Halluzinations-Filter (2. Netz): leeres/Artefakt-Transkript im borderline-
|
|
# Band → als no-speech verwerfen statt ein Phantom ("Die Stadt hat eine
|
|
# Flaeche von 1,5 km2") ans Brain zu schicken. Gilt fuer ALLE reasons inkl.
|
|
# stream_end (dort kamen die realen Phantome!) — aber das borderline-Band
|
|
# (wenig voiced_frames) schuetzt echte, klar gesprochene Eingaben: eine echte
|
|
# Geografie-FRAGE hat normale Stimm-Energie (voiced_frames >> Schwelle) und
|
|
# geht durch; das Phantom aus Stille hat ~0 und wird verworfen. Ein leeres
|
|
# Transkript wird immer verworfen (nichts gesagt = nichts senden).
|
|
_clean = final_text.strip(" .,!?…-\t\n\r")
|
|
_borderline = sess.voiced_frames < STREAM_HALLUC_GUARD_FRAMES
|
|
_is_phantom = (not _clean) or (_borderline and bool(_HALLUCINATION_RE.search(final_text)))
|
|
if _is_phantom:
|
|
logger.info("Stream %s: Halluzination verworfen (voiced_frames=%d<%d, %.1fs, text=%r)",
|
|
sess.request_id[:8], sess.voiced_frames, STREAM_HALLUC_GUARD_FRAMES,
|
|
duration_s, final_text[:80])
|
|
if self._ws is not None:
|
|
nospeech = {"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": f"hallucination:{reason}",
|
|
"durationS": 0.0, "sttMs": stt_ms}
|
|
await _send(self._ws, "stt_endpoint", nospeech)
|
|
await _send(self._ws, "stt_stream_done", {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": "", "reason": f"hallucination:{reason}"})
|
|
self.drop(sess.request_id)
|
|
return
|
|
|
|
if self._ws is not None:
|
|
payload = {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": final_text,
|
|
"reason": reason,
|
|
"durationS": duration_s,
|
|
"sttMs": stt_ms,
|
|
"voice": sess.voice,
|
|
"speed": sess.speed,
|
|
"interrupted": sess.interrupted,
|
|
}
|
|
if sess.location:
|
|
payload["location"] = sess.location
|
|
await _send(self._ws, "stt_endpoint", payload)
|
|
await _send(self._ws, "stt_stream_done", {
|
|
"requestId": sess.request_id,
|
|
"audioRequestId": sess.audio_request_id,
|
|
"text": final_text,
|
|
"reason": reason,
|
|
})
|
|
self.drop(sess.request_id)
|
|
|
|
|
|
async def _broadcast_status(ws, state: str, **extra) -> None:
|
|
payload = {"service": "voxtral", "state": state}
|
|
payload.update(extra)
|
|
await _send(ws, "service_status", payload)
|
|
|
|
|
|
async def run_loop(sessions: SessionManager) -> None:
|
|
use_tls = RVS_TLS
|
|
retry_s = 2
|
|
tls_fallback_tried = False
|
|
while True:
|
|
scheme = "wss" if use_tls else "ws"
|
|
url = f"{scheme}://{RVS_HOST}:{RVS_PORT}/ws?token={RVS_TOKEN}"
|
|
masked = url.replace(RVS_TOKEN, "***") if RVS_TOKEN else url
|
|
try:
|
|
logger.info("Verbinde zu RVS: %s", masked)
|
|
async with websockets.connect(url, ping_interval=20, ping_timeout=10,
|
|
max_size=50 * 1024 * 1024) as ws:
|
|
logger.info("RVS verbunden")
|
|
retry_s = 2
|
|
tls_fallback_tried = False
|
|
sessions.attach_ws(ws)
|
|
await _broadcast_status(ws, "ready", model=VOXTRAL_MODEL)
|
|
await _send(ws, "config_request", {"service": "voxtral"})
|
|
async for raw in ws:
|
|
try:
|
|
msg = json.loads(raw)
|
|
except Exception:
|
|
continue
|
|
mtype = msg.get("type", "")
|
|
payload = msg.get("payload", {}) or {}
|
|
if mtype == "stt_stream_start":
|
|
sessions.start_session(payload)
|
|
elif mtype == "stt_audio_chunk":
|
|
sessions.feed_chunk(payload)
|
|
elif mtype == "stt_stream_end":
|
|
sessions.end_session(payload.get("requestId", ""))
|
|
elif mtype == "stt_transcribe_blob":
|
|
# One-Shot-Transkription eines PCM-Schnipsels (kein Live-
|
|
# Stream) — fuer die Wake-Wort-Bestaetigung: die App schickt
|
|
# den Vor-Trigger-Audio, wir sagen was gesagt wurde, die App
|
|
# prueft ob "Computer" drin ist. Silero vorgeschaltet:
|
|
# Musik/Rauschen → leerer Text (nicht bestaetigt).
|
|
req_id = payload.get("requestId", "")
|
|
try:
|
|
pcm = base64.b64decode(payload.get("pcm", ""))
|
|
audio = pcm_s16le_to_float32(pcm)
|
|
segs = _speech_segments(audio)
|
|
if segs is not None and len(segs) == 0:
|
|
text = ""
|
|
else:
|
|
text = (await sessions.runner.transcribe(
|
|
audio, payload.get("language", "de"))).strip()
|
|
logger.info("stt_transcribe_blob (%.1fs) → %r",
|
|
audio.size / 16000.0, text[:60])
|
|
await _send(ws, "stt_transcribe_result",
|
|
{"requestId": req_id, "text": text})
|
|
except Exception as exc:
|
|
logger.warning("stt_transcribe_blob fehlgeschlagen: %s", exc)
|
|
await _send(ws, "stt_transcribe_result",
|
|
{"requestId": req_id, "text": "", "error": str(exc)[:200]})
|
|
elif mtype == "voice_id_status_request":
|
|
req_id = payload.get("requestId", "")
|
|
try:
|
|
status = speaker_id.status()
|
|
await _send(ws, "voice_id_status_response",
|
|
{"requestId": req_id, "ok": True, **status})
|
|
except Exception as exc:
|
|
await _send(ws, "voice_id_status_response",
|
|
{"requestId": req_id, "ok": False, "error": str(exc)[:200]})
|
|
elif mtype == "voice_id_enroll_request":
|
|
req_id = payload.get("requestId", "")
|
|
samples = payload.get("samples") or []
|
|
logger.info("voice_id_enroll_request: %d Samples (id=%s)", len(samples), req_id[:8])
|
|
try:
|
|
result = await asyncio.get_running_loop().run_in_executor(
|
|
None, speaker_id.enroll_from_samples, samples)
|
|
await _send(ws, "voice_id_enroll_response", {
|
|
"requestId": req_id, "ok": True,
|
|
"sample_count": result.get("sample_count", 0),
|
|
"rejected": result.get("rejected", []),
|
|
"updated_at": result.get("updated_at"),
|
|
"embedding_dim": result.get("embedding_dim"),
|
|
})
|
|
except Exception as exc:
|
|
logger.warning("voice_id_enroll failed: %s", exc)
|
|
await _send(ws, "voice_id_enroll_response",
|
|
{"requestId": req_id, "ok": False, "error": str(exc)[:300]})
|
|
elif mtype == "voice_id_delete_request":
|
|
req_id = payload.get("requestId", "")
|
|
removed = speaker_id.delete_fingerprint()
|
|
await _send(ws, "voice_id_delete_response",
|
|
{"requestId": req_id, "ok": True, "removed": removed})
|
|
elif mtype == "config":
|
|
if "voiceIdThreshold" in payload:
|
|
try:
|
|
t = float(payload.get("voiceIdThreshold", 0.5))
|
|
if 0.0 <= t <= 1.0:
|
|
speaker_id.DEFAULT_THRESHOLD = t
|
|
logger.info("[speaker-id] threshold gesetzt: %.2f", t)
|
|
except (TypeError, ValueError):
|
|
pass
|
|
if "voiceIdEnabled" in payload:
|
|
_set_speaker_id_enabled(payload.get("voiceIdEnabled"))
|
|
logger.info("[speaker-id] Gating %s (voiceIdEnabled)",
|
|
"AN" if SPEAKER_ID_ENABLED else "AUS")
|
|
except Exception as e:
|
|
logger.warning("RVS-Verbindung verloren: %s — retry in %ds", e, retry_s)
|
|
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
|
use_tls = False
|
|
tls_fallback_tried = True
|
|
continue
|
|
await asyncio.sleep(retry_s)
|
|
retry_s = min(retry_s * 2, 30)
|
|
use_tls = RVS_TLS
|
|
|
|
|
|
async def main() -> None:
|
|
if not RVS_HOST or not RVS_TOKEN:
|
|
logger.error("RVS_HOST/RVS_TOKEN fehlen — .env pruefen. Abbruch.")
|
|
return
|
|
runner = VoxtralRunner()
|
|
loop = asyncio.get_running_loop()
|
|
await loop.run_in_executor(None, runner.load) # Modell laden (blockierend)
|
|
sessions = SessionManager(runner)
|
|
logger.info("Voxtral-Bridge startet — Modell=%s", VOXTRAL_MODEL)
|
|
await asyncio.gather(run_loop(sessions), sessions.run_endpointer())
|
|
|
|
|
|
if __name__ == "__main__":
|
|
try:
|
|
asyncio.run(main())
|
|
except KeyboardInterrupt:
|
|
pass
|