Spaces:
Running
Running
| """ | |
| 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 βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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 ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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 βββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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())} | |