File size: 7,205 Bytes
e21e710
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
"""backend/api/agent_memory.py — Agent memory CRUD (S354).

GAP-MEM-FIX: aggiunta riconciliazione _mem_fallback → Supabase.
Problema confermato: quando Supabase è temporaneamente offline, le voci
finiscono solo in _mem_fallback (dict in-process). Al restart del backend
(HF Space free-tier riavvia spesso) il fallback viene perso completamente.
Fix: dopo ogni write Supabase riuscita, schedula un tentativo di sync del
fallback — se ci sono voci orfane le pubblica su Supabase e le rimuove dal
fallback locale. Nessun job periodico (troppo pesante su free-tier) — lazy
reconciliation al primo write riuscito dopo un periodo di downtime Supabase.

GAP-AUTH-MEMORY fix: POST e DELETE protetti con require_role(AuthRole.MACHINE).
GAP-MEM-RECONCILE-BREAK fix: continue invece di break per errori record-level.
GAP-AGENT-MEMORY-RECONCILE-TASK fix: create_task wrappato in try/except RuntimeError.
"""
import time, asyncio
from fastapi import APIRouter, Depends
from pydantic import BaseModel
from .state import _sb, _mem_fallback
from .auth_guard import require_role, AuthRole
from .global_state_sync import get_global_state_sync

import logging
_logger = logging.getLogger("api.agent_memory")

router = APIRouter()


class MemoryEntry(BaseModel):
    key: str
    value: str
    category: str = 'general'
    createdAt: int = 0
    updatedAt: int = 0


async def _reconcile_fallback() -> int:
    """GAP-MEM-FIX: sincronizza voci _mem_fallback → Supabase.

    Chiama dopo ogni write Supabase riuscita: se ci sono voci scritte
    solo in fallback (es. dopo un periodo di downtime Supabase), le pubblica.
    Ritorna il numero di voci sincronizzate.
    Non solleva mai eccezioni — fire-and-forget.

    GAP-MEM-RECONCILE-BREAK fix: break solo su errori network/connessione;
    continue per errori specifici al record (tipo sbagliato, valore too large, ecc.)
    per non bloccare la riconciliazione delle voci successive.
    """
    if not _sb or not _mem_fallback:
        return 0
    synced = 0
    for key, entry in list(_mem_fallback.items()):
        try:
            _sb.table('agent_memory').upsert({
                'key':        entry['key'],
                'value':      entry['value'],
                'category':   entry.get('category', 'general'),
                'created_at': entry.get('createdAt', 0),
                'updated_at': entry.get('updatedAt', 0),
            }, on_conflict='key').execute()
            synced += 1
        except Exception as _e:
            # GAP-MEM-RECONCILE-BREAK: distingui errore network (stop tutto) da errore record
            _is_network = isinstance(_e, (ConnectionError, TimeoutError, OSError))
            if _is_network:
                _logger.debug("[memory] reconcile: network error at key=%s — stopping: %s", key, _e)
                break  # Supabase non raggiungibile — interrompi, riprova al prossimo write
            # Errore specifico al record (tipo sbagliato, valore corrotto, ecc.) — salta e continua
            _logger.debug("[memory] reconcile: record-level error at key=%s — skipping: %s", key, _e)
            continue
    if synced:
        _logger.info("[memory] GAP-MEM-FIX: reconciled %d fallback entries to Supabase", synced)
    return synced


@router.get('/api/memory/agent')
async def list_agent_memory():
    # S766-GRID: Global State Sync layer (Supabase Federation)
    sync = get_global_state_sync()
    try:
        unified = await sync.get_unified_memory("all_entries")
        if unified and unified.get("data"):
            # Mappa i dati unificati nel formato atteso dal frontend
            entries = [
                {'key': r['key'], 'value': r['value'], 'category': r.get('category', 'general'),
                 'createdAt': r.get('created_at', 0), 'updatedAt': r.get('updated_at', 0)}
                for r in unified["data"]
            ]
            return {'entries': entries}
    except Exception as _exc:
        _logger.debug("[memory] grid sync list fail: %s", _exc)

    if _sb:
        try:
            data = _sb.table('agent_memory').select('*').order('updated_at', desc=True).execute()
            entries = [
                {'key': r['key'], 'value': r['value'], 'category': r.get('category', 'general'),
                 'createdAt': r.get('created_at', 0), 'updatedAt': r.get('updated_at', 0)}
                for r in (data.data or [])
            ]
            return {'entries': entries}
        except Exception as e:
            _logger.warning('[memory] Supabase list error: %s', e)
    return {'entries': list(_mem_fallback.values())}


@router.get('/api/memory/agent/{key}')
async def get_agent_memory(key: str):
    if _sb:
        try:
            data = _sb.table('agent_memory').select('*').eq('key', key).limit(1).execute()
            if data.data:
                return {'value': data.data[0]['value']}
        except Exception as e:
            _logger.warning('[memory] Supabase get error: %s', e)
    entry = _mem_fallback.get(key)
    return {'value': entry['value'] if entry else None}


@router.post('/api/memory/agent')
async def set_agent_memory(
    entry: MemoryEntry,
    _auth: AuthRole = Depends(require_role(AuthRole.MACHINE)),
):
    """GAP-AUTH-MEMORY fix: endpoint protetto con require_role(MACHINE).
    Richiede X-Internal-Token header (aggiunto dal CF Worker su tutte le route non-public).
    """
    now = int(time.time() * 1000)
    record = {
        'key': entry.key, 'value': entry.value, 'category': entry.category,
        'createdAt': entry.createdAt or now, 'updatedAt': entry.updatedAt or now,
    }
    # Sempre scrivi in fallback prima (garanzia immediata)
    _mem_fallback[entry.key] = record

    if _sb:
        try:
            _sb.table('agent_memory').upsert({
                'key': entry.key, 'value': entry.value, 'category': entry.category,
                'created_at': entry.createdAt or now, 'updated_at': entry.updatedAt or now,
            }, on_conflict='key').execute()
            # GAP-MEM-FIX: Supabase disponibile → schedula riconciliazione fallback orfano
            # GAP-AGENT-MEMORY-RECONCILE-TASK fix: wrappa in try/except per contesti senza event loop
            if len(_mem_fallback) > 1:
                try:
                    asyncio.create_task(_reconcile_fallback())
                except RuntimeError:
                    pass  # Event loop non attivo (test/startup context) — task silently dropped
        except Exception as _e:
            _logger.warning('[memory] Supabase write error (fallback attivo): %s', _e)

    return {'ok': True, 'key': entry.key}


@router.delete('/api/memory/agent/{key}')
async def delete_agent_memory(
    key: str,
    _auth: AuthRole = Depends(require_role(AuthRole.MACHINE)),
):
    """GAP-AUTH-MEMORY fix: endpoint protetto con require_role(MACHINE).
    Richiede X-Internal-Token header (aggiunto dal CF Worker su tutte le route non-public).
    """
    if _sb:
        try:
            _sb.table('agent_memory').delete().eq('key', key).execute()
        except Exception as _exc:
            _logger.debug("[agent_memory] silenced %s", type(_exc).__name__)  # noqa: BLE001
    _mem_fallback.pop(key, None)
    return {'deleted': key}