Terminal / api /memory_router.py
Baida-A
Initial clean deploy (Reverse Proxy removed)
28a08e7
Raw
History Blame
9.31 kB
"""
backend/api/memory_router.py β€” Unified Memory Router (ARCH-K2.3)
Espone /api/memory come interfaccia unica per tutti i layer di memoria:
GET /api/memory β€” lista/ricerca voci (layer, query, limit)
POST /api/memory β€” scrivi voce (layer, key, value/content)
GET /api/memory/stats β€” statistiche aggregate tutti i layer
Layer agent: persistito su Supabase (agent_memory) con fallback in-memory.
Layer episodic/semantic/reflection: delegati a MemoryManager se disponibile.
ROUTING CF PAGES: /api/memory (non /api/memory/compress o /semantic) β†’ BRAIN
Nessuna modifica a [[catchall]].ts necessaria.
NOTA: percorsi /api/memory/agent e /api/memory/decision giΓ  gestiti da _mem_router
e _decision_router. Questo router aggiunge SOLO /api/memory (radice) e /api/memory/stats.
"""
import time
import logging
from typing import Optional
from fastapi import APIRouter, Depends, Query
from fastapi.responses import JSONResponse
from pydantic import BaseModel
from .auth_guard import require_role, AuthRole
from .state import _sb, _mem_fallback
_logger = logging.getLogger("api.memory_router")
router = APIRouter(dependencies=[Depends(require_role(AuthRole.MACHINE))])
class MemoryWriteBody(BaseModel):
layer: str = "agent" # agent | episodic | semantic | reflection
key: Optional[str] = None
value: Optional[str] = None
task: Optional[str] = None # alias per episodic/semantic
content: Optional[str] = None # alias per value
category: str = "general"
success: bool = True
tags: list[str] = []
# ── GET /api/memory ───────────────────────────────────────────────────────────
@router.get("/api/memory")
async def list_memory(
layer: str = Query(default="agent", description="agent | episodic | semantic | reflection | all"),
query: Optional[str] = Query(default=None, description="testo da cercare"),
limit: int = Query(default=50, ge=1, le=500),
):
"""
Lista o cerca voci in uno o tutti i layer di memoria.
- layer=agent (default): legge da Supabase agent_memory + fallback in-memory
- layer=all + query: ricerca cross-layer tramite MemoryManager
"""
result: dict = {}
# ── AGENT layer: Supabase + in-memory fallback ────────────────────────────
if layer in ("agent", "all"):
entries: list[dict] = []
if _sb:
try:
res = (
_sb.table("agent_memory")
.select("*")
.order("updated_at", desc=True)
.limit(limit)
.execute()
)
entries = [
{
"key": r["key"],
"value": r["value"],
"category": r.get("category", "general"),
"updatedAt": r.get("updated_at", 0),
"layer": "agent",
}
for r in (res.data or [])
]
if query:
q = query.lower()
entries = [
e for e in entries
if q in e["key"].lower() or q in e["value"].lower()
]
except Exception as exc:
_logger.warning("[memory_router] Supabase agent list: %s", exc)
if not entries:
fallback_vals = list(_mem_fallback.values())[:limit]
entries = [
{**v, "layer": "agent"}
for v in fallback_vals
if not query or (
query.lower() in v.get("key", "").lower()
or query.lower() in v.get("value", "").lower()
)
]
result["agent"] = entries
# ── MemoryManager layers (episodic / semantic / reflection) ───────────────
if layer in ("semantic", "episodic", "reflection", "all"):
_mm_layer = None if layer == "all" else layer
_search_q = query or ""
if _search_q or layer != "all": # evita scan inutile su all senza query
try:
# Import lazy: MemoryManager inizializzato in _on_startup, non all'import
from memory.manager import _global_manager as _mm # type: ignore[import]
if _mm is not None:
hits = await _mm.search(_search_q, n=limit, layer=_mm_layer)
result[layer if layer != "all" else "multiLayer"] = hits
except Exception as exc:
_logger.debug("[memory_router] MemoryManager layer='%s': %s", layer, exc)
return {"layer": layer, "query": query, "results": result}
# ── POST /api/memory ──────────────────────────────────────────────────────────
@router.post("/api/memory")
async def write_memory(body: MemoryWriteBody):
"""
Scrive una voce di memoria nel layer specificato.
- layer agent: upsert su Supabase + fallback in-memory
- layer episodic/semantic/reflection: delega a MemoryManager.save_episode
"""
now = int(time.time() * 1000)
# ── AGENT layer ───────────────────────────────────────────────────────────
if body.layer == "agent":
key = body.key or f"auto_{now}"
value = body.value or body.content or ""
record = {
"key": key, "value": value,
"category": body.category,
"createdAt": now, "updatedAt": now,
}
_mem_fallback[key] = record # garanzia immediata
if _sb:
try:
_sb.table("agent_memory").upsert(
{
"key": key, "value": value,
"category": body.category,
"created_at": now, "updated_at": now,
},
on_conflict="key",
).execute()
except Exception as exc:
_logger.warning("[memory_router] Supabase write (fallback attivo): %s", exc)
return {"ok": True, "layer": "agent", "key": key}
# ── MemoryManager layers ──────────────────────────────────────────────────
if body.layer in ("episodic", "semantic", "reflection"):
try:
from memory.manager import _global_manager as _mm # type: ignore[import]
if _mm is None:
return JSONResponse(
status_code=503,
content={"ok": False, "error": "MemoryManager non inizializzato"},
)
task = body.task or body.key or f"auto_{now}"
content = body.content or body.value or ""
await _mm.save_episode(
type_=body.layer, task=task,
output=content, success=body.success,
tags=body.tags or None,
)
return {"ok": True, "layer": body.layer, "task": task}
except Exception as exc:
_logger.warning("[memory_router] MemoryManager write '%s': %s", body.layer, exc)
return JSONResponse(status_code=500, content={"ok": False, "error": str(exc)})
return JSONResponse(
status_code=400,
content={
"ok": False,
"error": (
f"layer '{body.layer}' non supportato. "
"Valori validi: agent | episodic | semantic | reflection"
),
},
)
# ── GET /api/memory/stats ─────────────────────────────────────────────────────
@router.get("/api/memory/stats")
async def memory_stats():
"""
Statistiche aggregate di tutti i layer di memoria.
Combina: conteggio Supabase agent_memory + stats MemoryManager (episodic/semantic/reflection).
"""
stats: dict = {}
# Agent layer
agent_count = len(_mem_fallback)
supabase_ok = False
if _sb:
try:
res = _sb.table("agent_memory").select("key", count="exact").execute()
if res.count is not None:
agent_count = res.count
supabase_ok = True
except Exception as exc:
_logger.debug("[memory_router] stats agent count: %s", exc)
stats["agent"] = {
"count": agent_count,
"supabase": supabase_ok,
"fallback_entries": len(_mem_fallback),
}
# MemoryManager layers
try:
from memory.manager import _global_manager as _mm # type: ignore[import]
if _mm is not None:
mgr_stats = _mm.stats()
stats.update(mgr_stats)
except Exception as exc:
_logger.debug("[memory_router] MemoryManager stats: %s", exc)
return {"stats": stats, "layers": list(stats.keys())}