feat: Langzeit-Erinnerungen (#3b) und WebSocket-Streaming-Chat (#4, erster Increment)

#3b Langzeit-Erinnerungen:
- Store: memories-Tabelle + add/get/delete_memory (pro Nutzer)
- API: GET/POST/DELETE /api/me/memories
- chat.py injiziert Nutzer-Erinnerungen als System-Kontext ins LLM (sessionunabhaengig)
- ?debug zeigt memories_len

#4 Echtzeit (erster Increment):
- WS /ws/chat: dauerhafter Kanal, Event-Folge ack -> semantic -> audio (binaer) -> done
- Auth (Token-Query), Session-Gedaechtnis und Erinnerungen wie bei POST /api/chat
- Fehler als error-Event (422/403/502)

- Tests: 38 gruen (Erinnerungs-CRUD/Injektion, WebSocket-Streaming/Auth)
- Doku aktualisiert (README, BEDIENUNGSANLEITUNG, Architektur)

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Dieter Schlüter 2026-06-17 04:26:55 +02:00
commit 531b57e08d
12 changed files with 389 additions and 9 deletions

View file

@ -58,23 +58,32 @@ async def chat(
orchestrator = build_orchestrator(route)
output = await resolve_output_endpoint(route)
# Gespraechsverlauf laden (nur bei gesetzter session_id -> sonst zustandslos).
history = (
conversation = (
store.get_recent_messages(session_id, settings.history_max_messages)
if session_id
else []
)
# Langzeit-Erinnerungen sind nutzerbezogen und gelten auch ohne Session.
memories = store.get_memories(user.id)
except SessionOwnershipError as exc:
raise HTTPException(status_code=403, detail=str(exc))
except RoutingError as exc:
raise HTTPException(status_code=422, detail=str(exc))
llm_context = list(conversation)
if memories:
memory_text = "Was du ueber den Nutzer weisst:\n" + "\n".join(
f"- {m.content}" for m in memories
)
llm_context = [{"role": "system", "content": memory_text}] + llm_context
try:
trace, audio = await orchestrator.chat_text(
payload.text,
language=route.language,
voice=voice,
output=output,
history=history,
history=llm_context,
)
except Exception as exc:
raise HTTPException(status_code=502, detail=str(exc))
@ -90,7 +99,8 @@ async def chat(
"ok": True,
"voice": voice,
"route": route.as_dict(),
"history_len": len(history),
"history_len": len(conversation),
"memories_len": len(memories),
"trace": {
"raw_transcript": trace.raw_transcript,
"cleaned_transcript": trace.cleaned_transcript,

View file

@ -1,8 +1,8 @@
from fastapi import APIRouter, Depends
from fastapi import APIRouter, Depends, HTTPException
from app.auth import require_user
from app.dependencies import get_store
from app.schemas import UserPrefs
from app.schemas import UserPrefs, MemoryCreate, MemoryOut
from app.store import User
router = APIRouter()
@ -19,3 +19,21 @@ async def set_my_prefs(payload: UserPrefs, user: User = Depends(require_user)):
prefs = {k: v for k, v in payload.model_dump().items() if v is not None}
updated = get_store().set_user_prefs(user.id, prefs)
return {"user_id": updated.id, "prefs": updated.prefs}
@router.get("/me/memories", response_model=list[MemoryOut])
async def list_memories(user: User = Depends(require_user)):
return get_store().get_memories(user.id)
@router.post("/me/memories", response_model=MemoryOut)
async def add_memory(payload: MemoryCreate, user: User = Depends(require_user)):
"""Speichert einen dauerhaften Fakt/eine Vorliebe ueber den Nutzer."""
return get_store().add_memory(user.id, payload.content.strip())
@router.delete("/me/memories/{memory_id}")
async def delete_memory(memory_id: int, user: User = Depends(require_user)):
if not get_store().delete_memory(user.id, memory_id):
raise HTTPException(status_code=404, detail="Memory not found")
return {"ok": True, "deleted": memory_id}

126
app/api/ws.py Normal file
View file

@ -0,0 +1,126 @@
"""WebSocket-Streaming-Chat (Echtzeit-Transport, erster Increment).
Etabliert einen dauerhaften, bidirektionalen Kanal: der Client schickt pro Turn
eine JSON-Nachricht, der Server streamt strukturierte Events zurueck
(ack -> semantic -> audio (binaer) -> done). Auth, Session-Gedaechtnis und
Langzeit-Erinnerungen gelten wie bei POST /api/chat.
Bewusst spaeter (eigene Increments): Token-Level-LLM-Streaming, Audio-Eingang/
Streaming-STT, Barge-in/Interrupt und WebRTC.
"""
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from app.config import settings
from app.errors import RoutingError
from app.dependencies import (
get_store,
resolve_route,
build_orchestrator,
resolve_output_endpoint,
)
from app.store import SessionOwnershipError
router = APIRouter()
def _authenticate(token: str | None):
store = get_store()
if not settings.auth_enabled:
return store.ensure_anonymous_user()
if not token:
return None
return store.get_user_by_token(token)
@router.websocket("/ws/chat")
async def ws_chat(
websocket: WebSocket,
session_id: str | None = None,
token: str | None = None,
):
user = _authenticate(token)
if user is None:
# Vor accept() schliessen -> Handshake wird mit 403 abgelehnt.
await websocket.close(code=1008)
return
await websocket.accept()
store = get_store()
try:
while True:
msg = await websocket.receive_json()
text = (msg.get("text") or "").strip()
if not text:
await websocket.send_json({"type": "error", "detail": "empty text"})
continue
overrides = {
key: msg.get(key)
for key in (
"input_endpoint",
"output_endpoint",
"language",
"stt_provider",
"llm_provider",
"tts_provider",
)
}
try:
route = resolve_route(user, session_id, overrides)
orchestrator = build_orchestrator(route)
output = await resolve_output_endpoint(route)
conversation = (
store.get_recent_messages(session_id, settings.history_max_messages)
if session_id
else []
)
memories = store.get_memories(user.id)
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
llm_context = list(conversation)
if memories:
memory_text = "Was du ueber den Nutzer weisst:\n" + "\n".join(
f"- {m.content}" for m in memories
)
llm_context = [{"role": "system", "content": memory_text}] + llm_context
await websocket.send_json({"type": "ack", "route": route.as_dict()})
voice = msg.get("voice") or settings.openrouter_tts_voice
try:
trace, audio = await orchestrator.chat_text(
text,
language=route.language,
voice=voice,
output=output,
history=llm_context,
)
except Exception as exc:
await websocket.send_json({"type": "error", "status": 502, "detail": str(exc)})
continue
if session_id:
store.append_message(session_id, user.id, "user", text)
store.append_message(session_id, user.id, "assistant", trace.semantic_response)
await websocket.send_json(
{
"type": "semantic",
"text": trace.semantic_response,
"spoken": trace.spoken_response,
}
)
await websocket.send_bytes(audio)
await websocket.send_json(
{"type": "done", "audio_format": "pcm", "sample_rate": 24000}
)
except WebSocketDisconnect:
return

View file

@ -8,6 +8,7 @@ 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.ws import router as ws_router
app = FastAPI(title="Voice Assistant Gateway")
app.include_router(health_router)
@ -19,3 +20,4 @@ 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(ws_router)

View file

@ -88,3 +88,13 @@ class UserPrefs(BaseModel):
tts_provider: str | None = None
language: str | None = None
class MemoryCreate(BaseModel):
content: str = Field(min_length=1)
class MemoryOut(BaseModel):
id: int
content: str
created_at: str

View file

@ -42,6 +42,13 @@ class Session:
data: dict = field(default_factory=dict)
@dataclass
class Memory:
id: int
content: str
created_at: str = ""
class SessionOwnershipError(Exception):
"""Eine Session gehoert einem anderen Nutzer (-> HTTP 403)."""
@ -78,6 +85,18 @@ class Store(ABC):
def get_recent_messages(self, session_id: str, limit: int) -> list[dict]:
"""Liefert die letzten `limit` Nachrichten chronologisch ([{'role','content'}, ...])."""
@abstractmethod
def add_memory(self, user_id: str, content: str) -> Memory:
"""Speichert eine dauerhafte Erinnerung (Fakt/Vorliebe) zum Nutzer."""
@abstractmethod
def get_memories(self, user_id: str) -> list[Memory]:
"""Liefert alle Erinnerungen des Nutzers (chronologisch)."""
@abstractmethod
def delete_memory(self, user_id: str, memory_id: int) -> bool:
"""Loescht eine Erinnerung des Nutzers. True, wenn etwas geloescht wurde."""
class SQLiteStore(Store):
def __init__(self, db_path: str):
@ -119,6 +138,14 @@ class SQLiteStore(Store):
);
CREATE INDEX IF NOT EXISTS idx_messages_session
ON messages(session_id, id);
CREATE TABLE IF NOT EXISTS memories (
id INTEGER PRIMARY KEY AUTOINCREMENT,
user_id TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_memories_user
ON memories(user_id, id);
"""
)
@ -244,3 +271,33 @@ class SQLiteStore(Store):
(session_id, limit),
).fetchall()
return [{"role": row["role"], "content": row["content"]} for row in reversed(rows)]
# ----- Langzeit-Erinnerungen --------------------------------------------
def add_memory(self, user_id: str, content: str) -> Memory:
now = _now()
with self._connect() as conn:
cur = conn.execute(
"INSERT INTO memories (user_id, content, created_at) VALUES (?, ?, ?)",
(user_id, content, now),
)
memory_id = cur.lastrowid
return Memory(id=memory_id, content=content, created_at=now)
def get_memories(self, user_id: str) -> list[Memory]:
with self._connect() as conn:
rows = conn.execute(
"SELECT id, content, created_at FROM memories WHERE user_id = ? ORDER BY id",
(user_id,),
).fetchall()
return [
Memory(id=row["id"], content=row["content"], created_at=row["created_at"])
for row in rows
]
def delete_memory(self, user_id: str, memory_id: int) -> bool:
with self._connect() as conn:
cur = conn.execute(
"DELETE FROM memories WHERE id = ? AND user_id = ?",
(memory_id, user_id),
)
return cur.rowcount > 0