feat: Resilienz (Fallback-Ketten) und Metriken (#5)
- Fallback-Provider (app/providers/fallback.py) fuer STT/LLM/TTS: Provider-Kette der Reihe nach; Config *_FALLBACK; build_orchestrator baut Ketten (dedupliziert) - LLM-Stream-Fallback nur solange kein Token gesendet wurde - Metriken (app/metrics.py): In-Memory Counter/Timer, keine externe Dependency - HTTP-Middleware (Requests/Latenz/Status je Pfad); Pipeline-Stufen-Timing stt/llm/tts; Fallback-/Fehlerzaehler; GET /api/metrics (JSON + Prometheus) - Tests: 58 gruen (+6); Doku aktualisiert (README, BEDIENUNGSANLEITUNG, Architektur, .env.example) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
093da817d8
commit
6422444017
13 changed files with 452 additions and 23 deletions
13
app/api/metrics.py
Normal file
13
app/api/metrics.py
Normal file
|
|
@ -0,0 +1,13 @@
|
|||
from fastapi import APIRouter, Query
|
||||
from fastapi.responses import PlainTextResponse
|
||||
|
||||
from app.metrics import metrics
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
@router.get("/metrics")
|
||||
async def get_metrics(format: str = Query(default="json", description="json | prometheus")):
|
||||
if format == "prometheus":
|
||||
return PlainTextResponse(metrics.prometheus(), media_type="text/plain; version=0.0.4")
|
||||
return metrics.snapshot()
|
||||
|
|
@ -124,6 +124,9 @@ class Settings(BaseSettings):
|
|||
admin_api_key: str = ""
|
||||
auth_enabled: bool = True
|
||||
history_max_messages: int = 10
|
||||
stt_fallback: str = "" # kommaseparierte Provider-Namen (Fallback-Kette)
|
||||
llm_fallback: str = ""
|
||||
tts_fallback: str = ""
|
||||
model_config = SettingsConfigDict(
|
||||
env_file=ENV_FILE, case_sensitive=False, extra="ignore"
|
||||
)
|
||||
|
|
|
|||
|
|
@ -1,5 +1,10 @@
|
|||
from app.schemas import AudioChunk, PipelineTrace
|
||||
from app.pipeline.sentence_chunker import SentenceChunker
|
||||
from app.metrics import timer, metrics
|
||||
|
||||
|
||||
def _stage(name: str):
|
||||
return timer("stage_duration_seconds", {"stage": name})
|
||||
|
||||
# Festes Ausgabeformat der TTS-Stufe (s16le PCM, 24 kHz, mono).
|
||||
TTS_AUDIO_FORMAT = "pcm"
|
||||
|
|
@ -48,11 +53,12 @@ class Orchestrator:
|
|||
# input dient hier nur der Validierung/Metadaten; das Audio kommt per Upload.
|
||||
if input is not None:
|
||||
await input.capabilities()
|
||||
trace.raw_transcript = await self.stt.transcribe(
|
||||
audio_bytes,
|
||||
fmt=fmt,
|
||||
language=language,
|
||||
)
|
||||
with _stage("stt"):
|
||||
trace.raw_transcript = await self.stt.transcribe(
|
||||
audio_bytes,
|
||||
fmt=fmt,
|
||||
language=language,
|
||||
)
|
||||
trace.cleaned_transcript = await self.input_cleaner.run(
|
||||
trace.raw_transcript or ""
|
||||
)
|
||||
|
|
@ -67,7 +73,8 @@ class Orchestrator:
|
|||
):
|
||||
spoken = await self.spoken_adapter.run(text, language=language)
|
||||
normalized = await self.tts_normalizer.run(spoken, language=language)
|
||||
audio = await self.tts.synthesize(normalized, voice=voice)
|
||||
with _stage("tts"):
|
||||
audio = await self.tts.synthesize(normalized, voice=voice)
|
||||
await self._emit_to_output(audio, output)
|
||||
return audio
|
||||
|
||||
|
|
@ -84,10 +91,11 @@ class Orchestrator:
|
|||
trace.raw_transcript = text
|
||||
trace.cleaned_transcript = await self.input_cleaner.run(text or "")
|
||||
|
||||
trace.semantic_response = await self.llm.complete(
|
||||
trace.cleaned_transcript or "",
|
||||
history=history,
|
||||
)
|
||||
with _stage("llm"):
|
||||
trace.semantic_response = await self.llm.complete(
|
||||
trace.cleaned_transcript or "",
|
||||
history=history,
|
||||
)
|
||||
if not trace.semantic_response:
|
||||
raise RuntimeError("LLM returned an empty response")
|
||||
|
||||
|
|
@ -100,10 +108,11 @@ class Orchestrator:
|
|||
language=language,
|
||||
)
|
||||
|
||||
audio = await self.tts.synthesize(
|
||||
trace.tts_ready_text,
|
||||
voice=voice,
|
||||
)
|
||||
with _stage("tts"):
|
||||
audio = await self.tts.synthesize(
|
||||
trace.tts_ready_text,
|
||||
voice=voice,
|
||||
)
|
||||
await self._emit_to_output(audio, output)
|
||||
return trace, audio
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,11 @@ from app.providers.llm.openrouter import OpenRouterLLMProvider
|
|||
from app.providers.tts.openrouter import OpenRouterTTSProvider
|
||||
from app.providers.tts.chatterbox import ChatterboxTTSProvider
|
||||
from app.providers.tts.piper import PiperTTSProvider
|
||||
from app.providers.fallback import (
|
||||
FallbackSTTProvider,
|
||||
FallbackLLMProvider,
|
||||
FallbackTTSProvider,
|
||||
)
|
||||
from app.pipeline.input_cleaner import InputCleaner
|
||||
from app.pipeline.spoken_response_adapter import SpokenResponseAdapter
|
||||
from app.pipeline.tts_normalizer import TTSNormalizer
|
||||
|
|
@ -198,11 +203,32 @@ def resolve_route(
|
|||
return ResolvedRoute(**resolved)
|
||||
|
||||
|
||||
_FALLBACK_CLASS = {
|
||||
"stt": FallbackSTTProvider,
|
||||
"llm": FallbackLLMProvider,
|
||||
"tts": FallbackTTSProvider,
|
||||
}
|
||||
|
||||
|
||||
def _provider_chain(registry, primary: str, fallback_csv: str, module: str, cfg: Settings):
|
||||
"""Baut primaeren Provider + optionale Fallback-Kette (dedupliziert, Reihenfolge erhalten)."""
|
||||
names = [primary] + [n.strip() for n in (fallback_csv or "").split(",") if n.strip()]
|
||||
seen, ordered = set(), []
|
||||
for name in names:
|
||||
if name not in seen:
|
||||
seen.add(name)
|
||||
ordered.append(name)
|
||||
entries = [(name, _from_registry(registry, name, module.upper(), cfg)) for name in ordered]
|
||||
if len(entries) == 1:
|
||||
return entries[0][1]
|
||||
return _FALLBACK_CLASS[module](module, entries)
|
||||
|
||||
|
||||
def build_orchestrator(route: ResolvedRoute, cfg: Settings = settings) -> Orchestrator:
|
||||
return Orchestrator(
|
||||
stt=get_stt_provider(route.stt_provider, cfg),
|
||||
llm=get_llm_provider(route.llm_provider, cfg),
|
||||
tts=get_tts_provider(route.tts_provider, cfg),
|
||||
stt=_provider_chain(STT_REGISTRY, route.stt_provider, cfg.stt_fallback, "stt", cfg),
|
||||
llm=_provider_chain(LLM_REGISTRY, route.llm_provider, cfg.llm_fallback, "llm", cfg),
|
||||
tts=_provider_chain(TTS_REGISTRY, route.tts_provider, cfg.tts_fallback, "tts", cfg),
|
||||
input_cleaner=InputCleaner(),
|
||||
spoken_adapter=SpokenResponseAdapter(),
|
||||
tts_normalizer=TTSNormalizer(),
|
||||
|
|
|
|||
25
app/main.py
25
app/main.py
|
|
@ -1,4 +1,8 @@
|
|||
from fastapi import FastAPI
|
||||
import time
|
||||
|
||||
from fastapi import FastAPI, Request
|
||||
|
||||
from app.metrics import metrics
|
||||
from app.api.health import router as health_router
|
||||
from app.api.chat import router as chat_router
|
||||
from app.api.transcribe import router as transcribe_router
|
||||
|
|
@ -8,9 +12,27 @@ from app.api.sessions import router as sessions_router
|
|||
from app.api.config import router as config_router
|
||||
from app.api.admin import router as admin_router
|
||||
from app.api.me import router as me_router
|
||||
from app.api.metrics import router as metrics_router
|
||||
from app.api.ws import router as ws_router
|
||||
|
||||
app = FastAPI(title="Voice Assistant Gateway")
|
||||
|
||||
|
||||
@app.middleware("http")
|
||||
async def record_metrics(request: Request, call_next):
|
||||
start = time.perf_counter()
|
||||
response = await call_next(request)
|
||||
duration = time.perf_counter() - start
|
||||
# Route-Template (z. B. /api/sessions/{session_id}/route) statt konkreter URL,
|
||||
# um die Label-Kardinalitaet niedrig zu halten.
|
||||
route = request.scope.get("route")
|
||||
path = getattr(route, "path", request.url.path)
|
||||
labels = {"method": request.method, "path": path}
|
||||
metrics.inc("http_requests_total", {**labels, "status": response.status_code})
|
||||
metrics.observe("http_request_duration_seconds", duration, labels)
|
||||
return response
|
||||
|
||||
|
||||
app.include_router(health_router)
|
||||
app.include_router(chat_router, prefix="/api")
|
||||
app.include_router(transcribe_router, prefix="/api")
|
||||
|
|
@ -20,4 +42,5 @@ app.include_router(sessions_router, prefix="/api")
|
|||
app.include_router(config_router, prefix="/api")
|
||||
app.include_router(admin_router, prefix="/api")
|
||||
app.include_router(me_router, prefix="/api")
|
||||
app.include_router(metrics_router, prefix="/api")
|
||||
app.include_router(ws_router)
|
||||
|
|
|
|||
82
app/metrics.py
Normal file
82
app/metrics.py
Normal file
|
|
@ -0,0 +1,82 @@
|
|||
"""Schlanke In-Memory-Metriken (Counter + Timer) fuer einen Prozess.
|
||||
|
||||
Bewusst ohne externe Dependency. Fuer mehrere Instanzen/Prozesse spaeter durch
|
||||
einen gemeinsamen Backend (z. B. Prometheus-Exporter) ersetzbar.
|
||||
"""
|
||||
|
||||
import threading
|
||||
import time
|
||||
from collections import defaultdict
|
||||
|
||||
|
||||
class Metrics:
|
||||
def __init__(self):
|
||||
self._lock = threading.Lock()
|
||||
self._counters: dict[str, float] = defaultdict(float)
|
||||
self._timers: dict[str, list] = defaultdict(lambda: [0.0, 0]) # [sum, count]
|
||||
|
||||
@staticmethod
|
||||
def _key(name: str, labels: dict | None) -> str:
|
||||
if not labels:
|
||||
return name
|
||||
rendered = ",".join(f'{k}="{v}"' for k, v in sorted(labels.items()))
|
||||
return f"{name}{{{rendered}}}"
|
||||
|
||||
def inc(self, name: str, labels: dict | None = None, value: float = 1.0) -> None:
|
||||
with self._lock:
|
||||
self._counters[self._key(name, labels)] += value
|
||||
|
||||
def observe(self, name: str, seconds: float, labels: dict | None = None) -> None:
|
||||
with self._lock:
|
||||
agg = self._timers[self._key(name, labels)]
|
||||
agg[0] += seconds
|
||||
agg[1] += 1
|
||||
|
||||
def snapshot(self) -> dict:
|
||||
with self._lock:
|
||||
counters = dict(self._counters)
|
||||
timers = {
|
||||
key: {
|
||||
"sum": agg[0],
|
||||
"count": agg[1],
|
||||
"avg": (agg[0] / agg[1] if agg[1] else 0.0),
|
||||
}
|
||||
for key, agg in self._timers.items()
|
||||
}
|
||||
return {"counters": counters, "timers": timers}
|
||||
|
||||
def prometheus(self) -> str:
|
||||
snap = self.snapshot()
|
||||
lines = []
|
||||
for key, value in sorted(snap["counters"].items()):
|
||||
lines.append(f"{key} {value}")
|
||||
for key, agg in sorted(snap["timers"].items()):
|
||||
base, _, labels = key.partition("{")
|
||||
suffix = ("{" + labels) if labels else ""
|
||||
lines.append(f"{base}_sum{suffix} {agg['sum']}")
|
||||
lines.append(f"{base}_count{suffix} {agg['count']}")
|
||||
return "\n".join(lines) + "\n"
|
||||
|
||||
def reset(self) -> None:
|
||||
with self._lock:
|
||||
self._counters.clear()
|
||||
self._timers.clear()
|
||||
|
||||
|
||||
metrics = Metrics()
|
||||
|
||||
|
||||
class timer:
|
||||
"""Context-Manager: misst die Dauer und schreibt sie als Timer-Beobachtung."""
|
||||
|
||||
def __init__(self, name: str, labels: dict | None = None):
|
||||
self.name = name
|
||||
self.labels = labels
|
||||
|
||||
def __enter__(self):
|
||||
self._start = time.perf_counter()
|
||||
return self
|
||||
|
||||
def __exit__(self, *exc):
|
||||
metrics.observe(self.name, time.perf_counter() - self._start, self.labels)
|
||||
return False
|
||||
85
app/providers/fallback.py
Normal file
85
app/providers/fallback.py
Normal file
|
|
@ -0,0 +1,85 @@
|
|||
"""Fallback-Ketten: versuchen mehrere Provider der Reihe nach.
|
||||
|
||||
Faellt der primaere Provider aus (Timeout/Fehler), wird transparent der naechste
|
||||
versucht. Erfolgreicher Fallback und Provider-Fehler werden als Metrik erfasst.
|
||||
"""
|
||||
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from app.metrics import metrics
|
||||
|
||||
|
||||
class _Chain:
|
||||
def __init__(self, module: str, entries: list[tuple[str, object]]):
|
||||
self.module = module
|
||||
self.entries = entries # [(provider_name, provider), ...]
|
||||
|
||||
def _on_error(self, name: str) -> None:
|
||||
metrics.inc("provider_error_total", {"module": self.module, "provider": name})
|
||||
|
||||
def _on_fallback(self) -> None:
|
||||
metrics.inc("provider_fallback_total", {"module": self.module})
|
||||
|
||||
|
||||
class FallbackSTTProvider(_Chain):
|
||||
async def transcribe(self, audio_bytes, fmt, language=None) -> str:
|
||||
last_exc = None
|
||||
for index, (name, provider) in enumerate(self.entries):
|
||||
try:
|
||||
result = await provider.transcribe(audio_bytes, fmt, language=language)
|
||||
if index > 0:
|
||||
self._on_fallback()
|
||||
return result
|
||||
except Exception as exc: # noqa: BLE001 - bewusst breit fuer Resilienz
|
||||
last_exc = exc
|
||||
self._on_error(name)
|
||||
raise last_exc
|
||||
|
||||
|
||||
class FallbackLLMProvider(_Chain):
|
||||
async def complete(self, text, history=None, session_id=None) -> str:
|
||||
last_exc = None
|
||||
for index, (name, provider) in enumerate(self.entries):
|
||||
try:
|
||||
result = await provider.complete(text, history=history, session_id=session_id)
|
||||
if index > 0:
|
||||
self._on_fallback()
|
||||
return result
|
||||
except Exception as exc: # noqa: BLE001
|
||||
last_exc = exc
|
||||
self._on_error(name)
|
||||
raise last_exc
|
||||
|
||||
async def stream(self, text, history=None, session_id=None) -> AsyncIterator[str]:
|
||||
last_exc = None
|
||||
for index, (name, provider) in enumerate(self.entries):
|
||||
produced = False
|
||||
try:
|
||||
async for delta in provider.stream(text, history=history, session_id=session_id):
|
||||
produced = True
|
||||
yield delta
|
||||
if index > 0:
|
||||
self._on_fallback()
|
||||
return
|
||||
except Exception as exc: # noqa: BLE001
|
||||
last_exc = exc
|
||||
self._on_error(name)
|
||||
if produced:
|
||||
# Schon Token gesendet -> kein Fallback mehr moeglich.
|
||||
raise
|
||||
raise last_exc
|
||||
|
||||
|
||||
class FallbackTTSProvider(_Chain):
|
||||
async def synthesize(self, text, voice=None, audio_format="pcm") -> bytes:
|
||||
last_exc = None
|
||||
for index, (name, provider) in enumerate(self.entries):
|
||||
try:
|
||||
result = await provider.synthesize(text, voice=voice, audio_format=audio_format)
|
||||
if index > 0:
|
||||
self._on_fallback()
|
||||
return result
|
||||
except Exception as exc: # noqa: BLE001
|
||||
last_exc = exc
|
||||
self._on_error(name)
|
||||
raise last_exc
|
||||
Loading…
Add table
Add a link
Reference in a new issue