feat: Tageskontingent und Notfall-Eskalation (#6)
- Quota (app/quota.py): Anfragen pro Nutzer/Tag (usage-Tabelle), DAILY_REQUEST_LIMIT, pro Nutzer via prefs uebersteuerbar; 429 (REST) bzw. error-Event (WS); Metrik - Notfall (app/safety/emergency.py): heuristische Erkennung (de/en); Log im Store (emergency_events) + optionaler Webhook (best-effort) + X-Emergency/emergency-Event; Metrik emergency_total; Notfaelle umgehen das Kontingent - verdrahtet in chat/speak/transcribe + WS-Turns - Tests: 64 gruen (+6); Doku aktualisiert (README, BEDIENUNGSANLEITUNG, Architektur, .env.example) Hinweis: Notfall-Erkennung ist eine Heuristik (kein Lebensretter); erkannte Texte sind sensibel -> DSGVO beachten. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
6422444017
commit
1eb79c1f09
14 changed files with 378 additions and 2 deletions
|
|
@ -13,6 +13,8 @@ from app.dependencies import (
|
|||
resolve_output_endpoint,
|
||||
get_store,
|
||||
)
|
||||
from app.quota import enforce_quota, record_usage, QuotaExceededError
|
||||
from app.safety.emergency import handle_emergency
|
||||
from app.schemas import ChatRequest
|
||||
|
||||
router = APIRouter()
|
||||
|
|
@ -77,6 +79,15 @@ async def chat(
|
|||
)
|
||||
llm_context = [{"role": "system", "content": memory_text}] + llm_context
|
||||
|
||||
# Notfall-Erkennung zuerst (immer eskalieren, auch bei Quota-Limit).
|
||||
emergency = handle_emergency(user, payload.text, store)
|
||||
|
||||
if emergency is None:
|
||||
try:
|
||||
enforce_quota(user, store)
|
||||
except QuotaExceededError as exc:
|
||||
raise HTTPException(status_code=429, detail=str(exc))
|
||||
|
||||
try:
|
||||
trace, audio = await orchestrator.chat_text(
|
||||
payload.text,
|
||||
|
|
@ -88,6 +99,8 @@ async def chat(
|
|||
except Exception as exc:
|
||||
raise HTTPException(status_code=502, detail=str(exc))
|
||||
|
||||
record_usage(user, store, len(payload.text) + len(trace.semantic_response or ""))
|
||||
|
||||
# Turn persistieren (User-Eingabe + semantische Antwort) fuer das Gedaechtnis.
|
||||
if session_id:
|
||||
store.append_message(session_id, user.id, "user", payload.text)
|
||||
|
|
@ -101,6 +114,7 @@ async def chat(
|
|||
"route": route.as_dict(),
|
||||
"history_len": len(conversation),
|
||||
"memories_len": len(memories),
|
||||
"emergency": emergency,
|
||||
"trace": {
|
||||
"raw_transcript": trace.raw_transcript,
|
||||
"cleaned_transcript": trace.cleaned_transcript,
|
||||
|
|
@ -119,4 +133,6 @@ async def chat(
|
|||
"X-Audio-Sample-Width": "16",
|
||||
**_route_headers(route),
|
||||
}
|
||||
if emergency:
|
||||
headers["X-Emergency"] = emergency["category"]
|
||||
return StreamingResponse(BytesIO(audio), media_type="audio/pcm", headers=headers)
|
||||
|
|
|
|||
|
|
@ -11,7 +11,9 @@ from app.dependencies import (
|
|||
resolve_route,
|
||||
build_orchestrator,
|
||||
resolve_output_endpoint,
|
||||
get_store,
|
||||
)
|
||||
from app.quota import enforce_quota, record_usage, QuotaExceededError
|
||||
from app.schemas import SpeakRequest
|
||||
|
||||
router = APIRouter()
|
||||
|
|
@ -42,6 +44,12 @@ async def speak(
|
|||
except RoutingError as exc:
|
||||
raise HTTPException(status_code=422, detail=str(exc))
|
||||
|
||||
store = get_store()
|
||||
try:
|
||||
enforce_quota(user, store)
|
||||
except QuotaExceededError as exc:
|
||||
raise HTTPException(status_code=429, detail=str(exc))
|
||||
|
||||
try:
|
||||
audio = await orchestrator.speak_only(
|
||||
payload.text,
|
||||
|
|
@ -49,6 +57,7 @@ async def speak(
|
|||
language=route.language,
|
||||
output=output,
|
||||
)
|
||||
record_usage(user, store, len(payload.text))
|
||||
|
||||
headers = {
|
||||
"Content-Language": route.language,
|
||||
|
|
|
|||
|
|
@ -7,7 +7,9 @@ from app.dependencies import (
|
|||
resolve_route,
|
||||
build_orchestrator,
|
||||
resolve_input_endpoint,
|
||||
get_store,
|
||||
)
|
||||
from app.quota import enforce_quota, record_usage, QuotaExceededError
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
|
@ -39,6 +41,12 @@ async def transcribe(
|
|||
except RoutingError as exc:
|
||||
raise HTTPException(status_code=422, detail=str(exc))
|
||||
|
||||
store = get_store()
|
||||
try:
|
||||
enforce_quota(user, store)
|
||||
except QuotaExceededError as exc:
|
||||
raise HTTPException(status_code=429, detail=str(exc))
|
||||
|
||||
content = await file.read()
|
||||
suffix = (file.filename or "audio.wav").rsplit(".", 1)[-1].lower()
|
||||
|
||||
|
|
@ -52,4 +60,6 @@ async def transcribe(
|
|||
except Exception as exc:
|
||||
raise HTTPException(status_code=502, detail=str(exc))
|
||||
|
||||
record_usage(user, store, len(trace.raw_transcript or ""))
|
||||
|
||||
return {"route": route.as_dict(), "trace": trace.model_dump()}
|
||||
|
|
|
|||
|
|
@ -29,6 +29,8 @@ from app.dependencies import (
|
|||
)
|
||||
from app.store import SessionOwnershipError
|
||||
from app.audio.vad import EnergyVAD
|
||||
from app.quota import enforce_quota, record_usage, QuotaExceededError
|
||||
from app.safety.emergency import handle_emergency
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
|
@ -75,6 +77,17 @@ async def _run_turn(websocket, store, user, session_id, route, orchestrator, out
|
|||
)
|
||||
llm_context = [{"role": "system", "content": memory_text}] + llm_context
|
||||
|
||||
# Notfall-Erkennung zuerst (immer eskalieren, auch bei Quota-Limit).
|
||||
emergency = handle_emergency(user, text, store)
|
||||
if emergency:
|
||||
await websocket.send_json({"type": "emergency", "category": emergency["category"]})
|
||||
else:
|
||||
try:
|
||||
enforce_quota(user, store)
|
||||
except QuotaExceededError as exc:
|
||||
await websocket.send_json({"type": "error", "status": 429, "detail": str(exc)})
|
||||
return
|
||||
|
||||
await websocket.send_json({"type": "ack", "route": route.as_dict()})
|
||||
|
||||
voice = options.get("voice") or settings.openrouter_tts_voice
|
||||
|
|
@ -119,6 +132,8 @@ async def _run_turn(websocket, store, user, session_id, route, orchestrator, out
|
|||
await websocket.send_json({"type": "error", "status": 502, "detail": str(exc)})
|
||||
return
|
||||
|
||||
record_usage(user, store, len(text) + len(trace.semantic_response or ""))
|
||||
|
||||
if session_id:
|
||||
store.append_message(session_id, user.id, "user", text)
|
||||
store.append_message(session_id, user.id, "assistant", trace.semantic_response)
|
||||
|
|
|
|||
|
|
@ -127,6 +127,8 @@ class Settings(BaseSettings):
|
|||
stt_fallback: str = "" # kommaseparierte Provider-Namen (Fallback-Kette)
|
||||
llm_fallback: str = ""
|
||||
tts_fallback: str = ""
|
||||
daily_request_limit: int = 0 # 0 = unbegrenzt; Anfragen pro Nutzer pro Tag
|
||||
emergency_webhook_url: str = "" # optionaler Eskalations-Webhook
|
||||
model_config = SettingsConfigDict(
|
||||
env_file=ENV_FILE, case_sensitive=False, extra="ignore"
|
||||
)
|
||||
|
|
|
|||
40
app/quota.py
Normal file
40
app/quota.py
Normal file
|
|
@ -0,0 +1,40 @@
|
|||
"""Pro-Nutzer-Tageskontingent (Kostenkontrolle).
|
||||
|
||||
Limit aus Settings (`daily_request_limit`), pro Nutzer ueber `prefs.daily_request_limit`
|
||||
ueberschreibbar. 0 bedeutet unbegrenzt.
|
||||
"""
|
||||
|
||||
from app.config import settings, Settings
|
||||
from app.metrics import metrics
|
||||
|
||||
|
||||
class QuotaExceededError(Exception):
|
||||
def __init__(self, limit: int, count: int):
|
||||
self.limit = limit
|
||||
self.count = count
|
||||
super().__init__(f"Daily request limit reached ({count}/{limit})")
|
||||
|
||||
|
||||
def effective_limit(user, cfg: Settings = settings) -> int:
|
||||
pref = user.prefs.get("daily_request_limit") if user and user.prefs else None
|
||||
if pref is not None:
|
||||
try:
|
||||
return int(pref)
|
||||
except (TypeError, ValueError):
|
||||
pass
|
||||
return cfg.daily_request_limit
|
||||
|
||||
|
||||
def enforce_quota(user, store, cfg: Settings = settings) -> None:
|
||||
"""Wirft QuotaExceededError, wenn das Tageslimit erreicht ist."""
|
||||
limit = effective_limit(user, cfg)
|
||||
if limit and limit > 0:
|
||||
count = store.get_request_count(user.id)
|
||||
if count >= limit:
|
||||
metrics.inc("quota_exceeded_total")
|
||||
raise QuotaExceededError(limit, count)
|
||||
|
||||
|
||||
def record_usage(user, store, units: int = 0) -> None:
|
||||
store.add_usage(user.id, units)
|
||||
metrics.inc("turns_total")
|
||||
0
app/safety/__init__.py
Normal file
0
app/safety/__init__.py
Normal file
84
app/safety/emergency.py
Normal file
84
app/safety/emergency.py
Normal file
|
|
@ -0,0 +1,84 @@
|
|||
"""Heuristische Notfall-Erkennung und Eskalation (Senioren-Kontext).
|
||||
|
||||
WICHTIG: Schluesselwort-Heuristik, KEIN Ersatz fuer eine echte Klassifikation.
|
||||
Sie kann Notlagen verpassen oder Fehlalarme ausloesen. Die erkannten Textauszuege
|
||||
sind hochsensibel und werden bewusst protokolliert (DSGVO beachten: Einwilligung,
|
||||
Aufbewahrung, Zugriff).
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
|
||||
import httpx
|
||||
|
||||
from app.config import settings, Settings
|
||||
from app.metrics import metrics
|
||||
|
||||
# Phrasen je Kategorie (de/en), bewusst eher spezifisch gegen Fehlalarme.
|
||||
_PATTERNS: dict[str, list[str]] = {
|
||||
"medical": [
|
||||
"brustschmerz", "schmerzen in der brust", "kann nicht atmen", "keine luft",
|
||||
"atemnot", "herzinfarkt", "schlaganfall", "bewusstlos", "gestuerzt", "gestürzt",
|
||||
"gefallen und komme nicht hoch", "starke blutung",
|
||||
"chest pain", "can't breathe", "cannot breathe", "heart attack", "stroke",
|
||||
"i fell and can't", "bleeding badly",
|
||||
],
|
||||
"self_harm": [
|
||||
"nicht mehr leben", "mich umbringen", "selbstmord", "suizid", "will sterben",
|
||||
"kill myself", "end my life", "suicide", "want to die",
|
||||
],
|
||||
"help": [
|
||||
"notruf", "notarzt", "krankenwagen", "ruf einen arzt", "es brennt",
|
||||
"call an ambulance", "call 911", "call 112",
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
def detect(text: str):
|
||||
"""Liefert (category, matched_phrase) oder None."""
|
||||
if not text:
|
||||
return None
|
||||
low = text.lower()
|
||||
for category, phrases in _PATTERNS.items():
|
||||
for phrase in phrases:
|
||||
if phrase in low:
|
||||
return category, phrase
|
||||
return None
|
||||
|
||||
|
||||
async def _fire_webhook(url: str, user, category: str, snippet: str) -> None:
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=5) as client:
|
||||
await client.post(
|
||||
url,
|
||||
json={
|
||||
"user_id": user.id,
|
||||
"display_name": user.display_name,
|
||||
"category": category,
|
||||
"text": snippet,
|
||||
},
|
||||
)
|
||||
except Exception: # noqa: BLE001 - best effort, darf den Chat nicht brechen
|
||||
metrics.inc("emergency_webhook_error_total")
|
||||
|
||||
|
||||
def handle_emergency(user, text: str, store, cfg: Settings = settings):
|
||||
"""Erkennt, protokolliert und eskaliert ein Notfall-Signal.
|
||||
|
||||
Gibt {"category", "matched"} zurueck, wenn etwas erkannt wurde, sonst None.
|
||||
Der Webhook (falls konfiguriert) wird nicht-blockierend ausgeloest.
|
||||
"""
|
||||
match = detect(text)
|
||||
if not match:
|
||||
return None
|
||||
category, phrase = match
|
||||
snippet = text[:500]
|
||||
store.log_emergency(user.id, category, snippet)
|
||||
metrics.inc("emergency_total", {"category": category})
|
||||
if cfg.emergency_webhook_url:
|
||||
try:
|
||||
asyncio.get_running_loop().create_task(
|
||||
_fire_webhook(cfg.emergency_webhook_url, user, category, snippet)
|
||||
)
|
||||
except RuntimeError:
|
||||
pass # kein laufender Event-Loop (z. B. im Test) -> Webhook ueberspringen
|
||||
return {"category": category, "matched": phrase}
|
||||
64
app/store.py
64
app/store.py
|
|
@ -97,6 +97,18 @@ class Store(ABC):
|
|||
def delete_memory(self, user_id: str, memory_id: int) -> bool:
|
||||
"""Loescht eine Erinnerung des Nutzers. True, wenn etwas geloescht wurde."""
|
||||
|
||||
@abstractmethod
|
||||
def get_request_count(self, user_id: str, day: str | None = None) -> int:
|
||||
"""Anzahl der Anfragen des Nutzers am angegebenen Tag (Default: heute, UTC)."""
|
||||
|
||||
@abstractmethod
|
||||
def add_usage(self, user_id: str, units: int = 0, day: str | None = None) -> int:
|
||||
"""Zaehlt eine Anfrage (+units) und liefert die neue Tages-Anfragezahl."""
|
||||
|
||||
@abstractmethod
|
||||
def log_emergency(self, user_id: str, category: str, snippet: str) -> None:
|
||||
"""Protokolliert ein erkanntes Notfall-Signal (sensibel!)."""
|
||||
|
||||
|
||||
class SQLiteStore(Store):
|
||||
def __init__(self, db_path: str):
|
||||
|
|
@ -146,6 +158,20 @@ class SQLiteStore(Store):
|
|||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_memories_user
|
||||
ON memories(user_id, id);
|
||||
CREATE TABLE IF NOT EXISTS usage (
|
||||
user_id TEXT NOT NULL,
|
||||
day TEXT NOT NULL,
|
||||
requests INTEGER NOT NULL DEFAULT 0,
|
||||
units INTEGER NOT NULL DEFAULT 0,
|
||||
PRIMARY KEY (user_id, day)
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS emergency_events (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
user_id TEXT NOT NULL,
|
||||
category TEXT NOT NULL,
|
||||
snippet TEXT NOT NULL,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
"""
|
||||
)
|
||||
|
||||
|
|
@ -301,3 +327,41 @@ class SQLiteStore(Store):
|
|||
(memory_id, user_id),
|
||||
)
|
||||
return cur.rowcount > 0
|
||||
|
||||
# ----- Nutzung / Quota --------------------------------------------------
|
||||
@staticmethod
|
||||
def _today() -> str:
|
||||
return datetime.now(timezone.utc).date().isoformat()
|
||||
|
||||
def get_request_count(self, user_id: str, day: str | None = None) -> int:
|
||||
day = day or self._today()
|
||||
with self._connect() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT requests FROM usage WHERE user_id = ? AND day = ?",
|
||||
(user_id, day),
|
||||
).fetchone()
|
||||
return int(row["requests"]) if row else 0
|
||||
|
||||
def add_usage(self, user_id: str, units: int = 0, day: str | None = None) -> int:
|
||||
day = day or self._today()
|
||||
with self._connect() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO usage (user_id, day, requests, units) VALUES (?, ?, 1, ?)"
|
||||
" ON CONFLICT(user_id, day) DO UPDATE SET"
|
||||
" requests = requests + 1, units = units + excluded.units",
|
||||
(user_id, day, units),
|
||||
)
|
||||
row = conn.execute(
|
||||
"SELECT requests FROM usage WHERE user_id = ? AND day = ?",
|
||||
(user_id, day),
|
||||
).fetchone()
|
||||
return int(row["requests"])
|
||||
|
||||
# ----- Notfall-Protokoll ------------------------------------------------
|
||||
def log_emergency(self, user_id: str, category: str, snippet: str) -> None:
|
||||
with self._connect() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO emergency_events (user_id, category, snippet, created_at)"
|
||||
" VALUES (?, ?, ?, ?)",
|
||||
(user_id, category, snippet, _now()),
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue