diff --git a/.env.example b/.env.example index 11df940..c1af288 100644 --- a/.env.example +++ b/.env.example @@ -77,6 +77,14 @@ TTS_SAMPLE_RATE=24000 # Ziel-Sample-Rate (ffmpeg # LLM_FALLBACK=local-openai-compatible # TTS_FALLBACK=piper +# --- Automatische Erinnerungs-Extraktion ----------------------------------- +# Das LLM destilliert nach je N Turns dauerhafte Fakten/Vorlieben aus dem Gespraech +# und legt sie als Nutzer-Erinnerungen ab (best-effort, nicht-blockierend). +# MEMORY_EXTRACTION_ENABLED=true +# MEMORY_EXTRACTION_EVERY_N_TURNS=3 # wie oft extrahiert wird +# MEMORY_EXTRACTION_MAX=50 # Obergrenze gespeicherter Erinnerungen +# MEMORY_EXTRACTION_PROVIDER= # leer = Default-LLM; sonst Registry-Name + # --- Betrieb: Kontingent & Notfall ----------------------------------------- DAILY_REQUEST_LIMIT=0 # Anfragen pro Nutzer/Tag (0 = unbegrenzt) # EMERGENCY_WEBHOOK_URL=https://example.org/alert # optionale Eskalation diff --git a/Docs/voice-assistant-architecture.md b/Docs/voice-assistant-architecture.md index b987e4f..b76ccce 100644 --- a/Docs/voice-assistant-architecture.md +++ b/Docs/voice-assistant-architecture.md @@ -214,7 +214,7 @@ Reihenfolge der Weiterentwicklung: 1. **(erledigt)** Konfig- & Routing-Fundament: Profile, Device Router, Registry, Pro-Request-Override. 2. **(erledigt)** Cloud-Fundament: Bearer-Token-Auth, Mehrbenutzer, persistenter SQLite-Store, Mandanten-Trennung, dauerhafte Nutzer-Präferenzen. Offen: Skalierung auf gemeinsamen Store (Postgres/Redis) für mehrere Instanzen. -3. **(erledigt)** Konversationsgedächtnis: Kurzzeit-Gesprächsverlauf pro Session + Langzeit-Erinnerungen pro Nutzer (manuell gepflegt, als LLM-Kontext). Offen: **automatische** Extraktion/Zusammenfassung von Erinnerungen aus Gesprächen. +3. **(erledigt)** Konversationsgedächtnis: Kurzzeit-Gesprächsverlauf pro Session + Langzeit-Erinnerungen pro Nutzer (manuell **und automatisch** gepflegt, als LLM-Kontext). **Automatische Extraktion** (`app/core/memory_extractor.py`): nach je N Turns destilliert ein LLM dauerhafte Fakten/Vorlieben aus dem Verlauf und legt sie dedupliziert als Erinnerungen ab — best-effort, nicht-blockierend (Hintergrund-Task), konfigurierbar (`MEMORY_EXTRACTION_*`). Offen: periodische Verdichtung/Zusammenfassung wachsender Erinnerungslisten. 4. **(weitgehend erledigt)** Echtzeit: WebSocket-Streaming-Chat (`/ws/chat`), **Token-Level-LLM-Streaming (SSE, `stream:true`)**, **Audio-Streaming (chunked TTS satzweise, `audio_stream:true`)**, **Audio-Eingang (`/ws/voice`)**, **Barge-in/Turn-Manager (`interrupt` bricht laufende Antwort ab)** und **VAD-Aeusserungserkennung (energie-basiert, opt-in)** sind umgesetzt. Offen: **echte partielle Live-Transkripte (Streaming-STT-Dienst, wortweise)** und **WebRTC (aiortc)** — beide brauchen schwere Abhaengigkeiten/Dienste. Heute laeuft STT pro Aeusserung. 5. **(weitgehend erledigt)** Resilienz: Fallback-Ketten je Modul (`*_FALLBACK`, Provider faellt aus → naechster) und In-Memory-Metriken (`/api/metrics`: Request/Latenz, Pipeline-Stufen, Fallback/Fehler; JSON + Prometheus). Offen: verteiltes Tracing, Alerting. 6. **(weitgehend erledigt)** Betrieb: Tageskontingent pro Nutzer (`DAILY_REQUEST_LIMIT`, 429) und heuristische Notfall-Eskalation (Erkennung -> Log + optionaler Webhook + Flag/Event). Offen: echte Klassifikation statt Schluesselwort-Heuristik, Telefon-/Angehoerigen-Integration, Abrechnung. diff --git a/README.md b/README.md index 0f4f025..277bfde 100644 --- a/README.md +++ b/README.md @@ -226,6 +226,13 @@ curl -X POST http://localhost:8080/api/me/memories \ -H 'Content-Type: application/json' -d '{"content":"Mag morgens Kamillentee."}' ``` +**Automatische Erinnerungen:** Zusätzlich zur manuellen Pflege destilliert das LLM +nach je N Turns (Default 3) dauerhafte Fakten/Vorlieben aus dem Gespräch und legt sie +dedupliziert als Erinnerungen ab — best-effort und **nicht-blockierend** (Hintergrund-Task, +erhöht die Antwortlatenz nicht). Steuerung über `MEMORY_EXTRACTION_*` (siehe `.env.example`); +`MEMORY_EXTRACTION_ENABLED=false` schaltet es ab. Sinnvoll mit einem lokalen LLM, da pro +Turn ein zusätzlicher (kostenloser) Modellaufruf anfällt. + ## Echtzeit-Chat (WebSocket) `/ws/chat` bietet einen dauerhaften, bidirektionalen Kanal. Der Client sendet pro diff --git a/app/api/chat.py b/app/api/chat.py index df8ffbd..6deb509 100644 --- a/app/api/chat.py +++ b/app/api/chat.py @@ -13,6 +13,7 @@ from app.dependencies import ( resolve_output_endpoint, get_store, ) +from app.core.memory_extractor import maybe_schedule_extraction from app.quota import enforce_quota, record_usage, QuotaExceededError from app.safety.emergency import handle_emergency from app.schemas import ChatRequest @@ -106,6 +107,7 @@ async def chat( if session_id: store.append_message(session_id, user.id, "user", payload.text) store.append_message(session_id, user.id, "assistant", trace.semantic_response) + maybe_schedule_extraction(store, user.id, session_id) if debug: return JSONResponse( diff --git a/app/api/ws.py b/app/api/ws.py index 86c6b22..01ce254 100644 --- a/app/api/ws.py +++ b/app/api/ws.py @@ -29,6 +29,7 @@ from app.dependencies import ( ) from app.store import SessionOwnershipError from app.audio.vad import EnergyVAD +from app.core.memory_extractor import maybe_schedule_extraction from app.quota import enforce_quota, record_usage, QuotaExceededError from app.safety.emergency import handle_emergency @@ -141,6 +142,7 @@ async def _run_turn(websocket, store, user, session_id, route, orchestrator, out if session_id: store.append_message(session_id, user.id, "user", text) store.append_message(session_id, user.id, "assistant", trace.semantic_response) + maybe_schedule_extraction(store, user.id, session_id) await websocket.send_json( {"type": "semantic", "text": trace.semantic_response, "spoken": trace.spoken_response} diff --git a/app/config.py b/app/config.py index d995fe2..bed84ac 100644 --- a/app/config.py +++ b/app/config.py @@ -142,6 +142,13 @@ class Settings(BaseSettings): admin_api_key: str = "" auth_enabled: bool = True history_max_messages: int = 10 + # Automatische Erinnerungs-Extraktion: das LLM destilliert dauerhafte Fakten + # aus dem Gespraech und legt sie als Nutzer-Erinnerungen ab (best-effort, + # nicht-blockierend). Leerer Provider = Default-LLM-Provider. + memory_extraction_enabled: bool = True + memory_extraction_every_n_turns: int = 3 + memory_extraction_max: int = 50 + memory_extraction_provider: str = "" audio_stream_default: bool = True # satzweises TTS als Default (Admin kann abschalten) # TTS-Text-Normalisierung: auto|full|light|off. "auto" = piper -> full, Cloud -> light. tts_normalize_level: str = "auto" diff --git a/app/core/memory_extractor.py b/app/core/memory_extractor.py new file mode 100644 index 0000000..c84101c --- /dev/null +++ b/app/core/memory_extractor.py @@ -0,0 +1,164 @@ +"""Automatische Erinnerungs-Extraktion. + +Nach einigen Gespraechsturns destilliert ein LLM dauerhafte Fakten/Vorlieben +ueber den Nutzer aus dem Verlauf und legt sie als Nutzer-Erinnerungen ab. + +Bewusst **best-effort und nicht-blockierend**: Die Extraktion laeuft als +Hintergrund-Task und darf die Antwortlatenz nie erhoehen. Schlaegt sie fehl +(LLM-Fehler, kaputtes JSON), ist die Folge nur "kein neuer Fakt" - niemals ein +Fehler im Antwort-Turn. +""" + +import asyncio +import json +import logging +import re + +from app.config import Settings, settings + +logger = logging.getLogger(__name__) + +# Turn-Zaehler pro Session (in-memory, bewusst kein DB-Schema-Eingriff). +_turn_counts: dict[str, int] = {} +# Referenzen auf laufende Tasks halten, damit sie nicht vorzeitig vom GC kassiert werden. +_pending: set[asyncio.Task] = set() + +_EXTRACTION_SYSTEM_PROMPT = ( + "Du extrahierst dauerhafte, langfristig relevante Fakten und Vorlieben ueber den " + "Nutzer aus einem Gespraech (z. B. Name, Wohnort, Familie, Gesundheit, Hobbys, " + "Vorlieben, Abneigungen, feste Routinen). Gib AUSSCHLIESSLICH ein JSON-Array " + "kurzer deutscher Strings zurueck, ohne Erklaerung und ohne Markdown. Nimm nur " + "NEUE Fakten auf, die nicht bereits bekannt sind. Ignoriere fluechtige oder rein " + "situative Aussagen. Gibt es nichts Neues, antworte mit []." +) + + +def _build_extractor_llm(cfg: Settings): + """Baut eine eigene LLM-Instanz fuer die Extraktion (nicht der Sprach-Provider). + + Fuer den lokalen Provider wird der Extraktions-System-Prompt direkt gesetzt + (der Chat-Provider ist auf kurze, vorlesbare Saetze getrimmt und taugt nicht + fuer JSON). Fuer andere Provider wird der generische Provider verwendet; die + Anweisung steckt dann zusaetzlich in der Nachricht selbst. + """ + provider = cfg.memory_extraction_provider or cfg.default_llm_provider + if provider == "local-openai-compatible": + from app.providers.llm.local_openai_compatible import LocalOpenAICompatibleLLM + + return LocalOpenAICompatibleLLM( + cfg.local_llm_base_url, + cfg.local_llm_api_key, + cfg.local_llm_model, + system_prompt=_EXTRACTION_SYSTEM_PROMPT, + disable_reasoning=True, + max_tokens=512, + temperature=0.1, + ) + + from app.dependencies import get_llm_provider + + return get_llm_provider(provider, cfg) + + +def _format_conversation(messages: list[dict]) -> str: + lines = [] + for msg in messages: + role = "Nutzer" if msg.get("role") == "user" else "Assistent" + content = (msg.get("content") or "").strip() + if content: + lines.append(f"{role}: {content}") + return "\n".join(lines) + + +def parse_facts(raw: str) -> list[str]: + """Liest ein JSON-Array von Fakt-Strings aus der (evtl. verrauschten) LLM-Antwort.""" + if not raw: + return [] + match = re.search(r"\[.*\]", raw, re.DOTALL) + if not match: + return [] + try: + data = json.loads(match.group(0)) + except ValueError: + return [] + if not isinstance(data, list): + return [] + facts = [] + for item in data: + if isinstance(item, str): + fact = item.strip() + if fact: + facts.append(fact) + return facts + + +def _norm(text: str) -> str: + return " ".join(text.lower().split()) + + +async def extract_and_store(store, user_id: str, session_id: str, cfg: Settings = settings) -> int: + """Extrahiert neue Fakten und speichert sie. Liefert die Anzahl neu gespeicherter.""" + messages = store.get_recent_messages(session_id, cfg.history_max_messages) + conversation = _format_conversation(messages) + if not conversation: + return 0 + + existing = store.get_memories(user_id) + if len(existing) >= cfg.memory_extraction_max: + return 0 + + known = [m.content for m in existing] + known_text = "\n".join(f"- {k}" for k in known) if known else "(noch nichts bekannt)" + user_prompt = ( + f"Bereits bekannt:\n{known_text}\n\n" + f"Gespraech:\n{conversation}\n\n" + "Neue Fakten als JSON-Array:" + ) + + llm = _build_extractor_llm(cfg) + raw = await llm.complete(user_prompt) + facts = parse_facts(raw) + if not facts: + return 0 + + seen = {_norm(k) for k in known} + added = 0 + for fact in facts: + if len(existing) + added >= cfg.memory_extraction_max: + break + key = _norm(fact) + if key in seen: + continue + seen.add(key) + store.add_memory(user_id, fact) + added += 1 + if added: + logger.info("memory-extraction: %d neue Erinnerung(en) fuer %s", added, user_id) + return added + + +async def _run_safe(store, user_id: str, session_id: str, cfg: Settings) -> None: + try: + await extract_and_store(store, user_id, session_id, cfg) + except Exception: # best-effort: niemals den Turn beeintraechtigen + logger.exception("memory-extraction fehlgeschlagen (ignoriert)") + + +def maybe_schedule_extraction(store, user_id: str, session_id: str | None, + cfg: Settings = settings) -> asyncio.Task | None: + """Plant die Extraktion als Hintergrund-Task, sofern aktiviert und N Turns erreicht. + + Gibt den geplanten Task zurueck (oder None) - blockiert nie. + """ + if not cfg.memory_extraction_enabled or not session_id: + return None + every = max(1, cfg.memory_extraction_every_n_turns) + count = _turn_counts.get(session_id, 0) + 1 + _turn_counts[session_id] = count + if count % every != 0: + return None + + task = asyncio.create_task(_run_safe(store, user_id, session_id, cfg)) + _pending.add(task) + task.add_done_callback(_pending.discard) + return task diff --git a/tests/test_memory_extraction.py b/tests/test_memory_extraction.py new file mode 100644 index 0000000..2c5cae0 --- /dev/null +++ b/tests/test_memory_extraction.py @@ -0,0 +1,135 @@ +import asyncio + +import pytest + +import app.dependencies as deps +from app.config import settings +from app.core import memory_extractor as me + + +class StubLLM: + def __init__(self, raw): + self.raw = raw + self.calls = 0 + + async def complete(self, text, history=None, session_id=None): + self.calls += 1 + return self.raw + + +@pytest.fixture(autouse=True) +def _clear_turn_counts(): + me._turn_counts.clear() + yield + me._turn_counts.clear() + + +def _seed_conversation(store, session_id="conv", user_id="anonymous"): + store.append_message(session_id, user_id, "user", "Ich heisse Anna und wohne in Kiel.") + store.append_message(session_id, user_id, "assistant", "Schoen, Anna!") + + +def test_parse_facts_variants(): + assert me.parse_facts('["heisst Anna", "wohnt in Kiel"]') == ["heisst Anna", "wohnt in Kiel"] + # umschlossen von Geschwafel/Markdown + assert me.parse_facts('Hier:\n```json\n["x"]\n```') == ["x"] + assert me.parse_facts("[]") == [] + assert me.parse_facts("kein json") == [] + assert me.parse_facts("") == [] + # Nicht-Strings werden ignoriert + assert me.parse_facts('["ok", 5, null, " "]') == ["ok"] + + +def test_extracts_and_stores_new_facts(monkeypatch): + store = deps.get_store() + _seed_conversation(store) + monkeypatch.setattr(me, "_build_extractor_llm", + lambda cfg: StubLLM('["heisst Anna", "wohnt in Kiel"]')) + + added = asyncio.run(me.extract_and_store(store, "anonymous", "conv", settings)) + + assert added == 2 + contents = [m.content for m in store.get_memories("anonymous")] + assert contents == ["heisst Anna", "wohnt in Kiel"] + + +def test_dedup_skips_known(monkeypatch): + store = deps.get_store() + _seed_conversation(store) + store.add_memory("anonymous", "heisst Anna") + # LLM liefert einen bekannten (anders gross-/kleingeschrieben) + einen neuen Fakt. + monkeypatch.setattr(me, "_build_extractor_llm", + lambda cfg: StubLLM('["Heisst Anna", "wohnt in Kiel"]')) + + added = asyncio.run(me.extract_and_store(store, "anonymous", "conv", settings)) + + assert added == 1 + contents = [m.content for m in store.get_memories("anonymous")] + assert contents == ["heisst Anna", "wohnt in Kiel"] + + +def test_malformed_output_no_crash(monkeypatch): + store = deps.get_store() + _seed_conversation(store) + monkeypatch.setattr(me, "_build_extractor_llm", + lambda cfg: StubLLM("Tut mir leid, kein JSON hier.")) + + added = asyncio.run(me.extract_and_store(store, "anonymous", "conv", settings)) + + assert added == 0 + assert store.get_memories("anonymous") == [] + + +def test_cap_respected(monkeypatch): + store = deps.get_store() + _seed_conversation(store) + monkeypatch.setattr(settings, "memory_extraction_max", 1) + monkeypatch.setattr(me, "_build_extractor_llm", + lambda cfg: StubLLM('["fakt a", "fakt b", "fakt c"]')) + + added = asyncio.run(me.extract_and_store(store, "anonymous", "conv", settings)) + + assert added == 1 + assert len(store.get_memories("anonymous")) == 1 + + +def test_empty_conversation_skips_llm(monkeypatch): + store = deps.get_store() + called = StubLLM("[]") + monkeypatch.setattr(me, "_build_extractor_llm", lambda cfg: called) + + added = asyncio.run(me.extract_and_store(store, "anonymous", "leer", settings)) + + assert added == 0 + assert called.calls == 0 # ohne Gespraech kein LLM-Aufruf + + +def test_schedule_only_every_n_turns(monkeypatch): + store = deps.get_store() + monkeypatch.setattr(settings, "memory_extraction_enabled", True) + monkeypatch.setattr(settings, "memory_extraction_every_n_turns", 3) + + async def run(): + # Session "s1" hat keine Nachrichten -> der geplante Task endet sofort (kein LLM). + results = [me.maybe_schedule_extraction(store, "anonymous", "s1") for _ in range(3)] + for task in results: + if task is not None: + await task + return results + + tasks = asyncio.run(run()) + # nur der 3. Aufruf plant einen Task + assert tasks[0] is None and tasks[1] is None + assert tasks[2] is not None + + +def test_schedule_disabled_is_noop(monkeypatch): + store = deps.get_store() + monkeypatch.setattr(settings, "memory_extraction_enabled", False) + assert me.maybe_schedule_extraction(store, "anonymous", "s1") is None + + +def test_schedule_without_session_is_noop(monkeypatch): + store = deps.get_store() + monkeypatch.setattr(settings, "memory_extraction_enabled", True) + assert me.maybe_schedule_extraction(store, "anonymous", None) is None