Pulka commited on
Commit
2a17653
·
verified ·
1 Parent(s): c4132a4

sync: 2 changed, 0 deleted — from push d9b29fe0 (2026-06-18T20:42 UTC)

Browse files
agents/skill_tracker.py CHANGED
@@ -1,10 +1,8 @@
1
  """
2
  skill_tracker.py — GAP-SKILL-SYNC: Session-scoped adaptive tool success/failure tracker.
 
3
 
4
- Chiude il ciclo di feedback mancante: il backend impara quali tool funzionano meglio
5
- per una data sessione e li preferisce nelle iterazioni successive via Wilson score.
6
-
7
- Design:
8
  - In-memory Dict[session_id, Dict[tool_name, SkillStats]] — zero DB dep, zero latency
9
  - SkillStats: success_count, fail_count, last_used, total_latency_ms
10
  - record() sincrono — GIL-safe in CPython, asyncio single-threaded
@@ -12,24 +10,52 @@ Design:
12
  - get_stats() — JSON-serializable per /api/agent/skill-stats endpoint
13
  - clear_session() — cleanup opzionale fine task (evita memory leak su run lunghissimi)
14
 
 
 
 
 
 
 
 
 
 
15
  Integrazione:
16
  - unified_loop.py: record() dopo ogni executor.run_tool() — registra successo/fallimento
 
17
  - api/agent.py: GET /api/agent/skill-stats/{session_id} per merge Dexie frontend
18
 
19
  Singleton: get_skill_tracker() restituisce sempre lo stesso SkillTracker globale.
20
  """
21
  from __future__ import annotations
22
 
 
23
  import math
 
24
  import time
25
  import logging
26
  from collections import defaultdict
27
  from dataclasses import dataclass
28
  from typing import Any
29
 
 
 
30
  _logger = logging.getLogger("agente_ai.skill_tracker")
31
 
32
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
33
  # ─── SkillStats ───────────────────────────────────────────────────────────────
34
 
35
  @dataclass
@@ -72,16 +98,115 @@ class SkillStats:
72
  den = 1 + z * z / n
73
  return num / den
74
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
75
 
76
  # ─── SkillTracker ─────────────────────────────────────────────────────────────
77
 
78
  class SkillTracker:
79
- """Singleton session-scoped tracker: impara quali tool funzionano per ogni sessione."""
 
 
 
80
 
81
  def __init__(self) -> None:
82
  self._sessions: dict[str, dict[str, SkillStats]] = defaultdict(
83
  lambda: defaultdict(SkillStats)
84
  )
 
 
85
 
86
  # ── Write ─────────────────────────────────────────────────────────────────
87
 
@@ -95,7 +220,7 @@ class SkillTracker:
95
  """Registra il risultato di una chiamata tool (sincrono, GIL-safe).
96
 
97
  Chiamato da unified_loop.py dopo ogni executor.run_tool().
98
- Zero overhead: solo incremento contatori in-memory.
99
  """
100
  s = self._sessions[session_id][tool_name]
101
  if success:
@@ -105,6 +230,64 @@ class SkillTracker:
105
  s.last_used = time.monotonic()
106
  s.total_latency_ms += latency_ms
107
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
108
  # ── Read / routing ────────────────────────────────────────────────────────
109
 
110
  def get_sorted_fallbacks(
@@ -158,6 +341,7 @@ class SkillTracker:
158
  def clear_session(self, session_id: str) -> None:
159
  """Libera memoria per una sessione terminata."""
160
  removed = self._sessions.pop(session_id, None)
 
161
  if removed is not None:
162
  _logger.debug(
163
  "[skill_tracker] cleared session %s (%d tools tracked)",
 
1
  """
2
  skill_tracker.py — GAP-SKILL-SYNC: Session-scoped adaptive tool success/failure tracker.
3
+ GAP-SUP1: Supabase persistence — skill stats sopravvivono ai riavvii del backend.
4
 
5
+ Design originale:
 
 
 
6
  - In-memory Dict[session_id, Dict[tool_name, SkillStats]] — zero DB dep, zero latency
7
  - SkillStats: success_count, fail_count, last_used, total_latency_ms
8
  - record() sincrono — GIL-safe in CPython, asyncio single-threaded
 
10
  - get_stats() — JSON-serializable per /api/agent/skill-stats endpoint
11
  - clear_session() — cleanup opzionale fine task (evita memory leak su run lunghissimi)
12
 
13
+ GAP-SUP1 (Supabase persistence):
14
+ - _supabase_upsert(): httpx POST → PostgREST /rest/v1/skill_stats (upsert conflict)
15
+ - _supabase_load_session(): httpx GET → ripristina sessione precedente al boot
16
+ - record() fire-and-forget ogni _SYNC_EVERY_N chiamate per tool/sessione
17
+ - load_session_from_cloud(): chiamato da unified_loop al boot sessione
18
+ - Fallback silente se SUPABASE_URL/SUPABASE_ANON_KEY assenti → comportamento invariato
19
+ - Timeout conservativo 8s — mai blocca il loop principale
20
+ - Schema SQL: backend/migrations/gap1_skill_stats.sql
21
+
22
  Integrazione:
23
  - unified_loop.py: record() dopo ogni executor.run_tool() — registra successo/fallimento
24
+ - unified_loop.py: load_session_from_cloud() al boot sessione (se Supabase abilitato)
25
  - api/agent.py: GET /api/agent/skill-stats/{session_id} per merge Dexie frontend
26
 
27
  Singleton: get_skill_tracker() restituisce sempre lo stesso SkillTracker globale.
28
  """
29
  from __future__ import annotations
30
 
31
+ import asyncio
32
  import math
33
+ import os
34
  import time
35
  import logging
36
  from collections import defaultdict
37
  from dataclasses import dataclass
38
  from typing import Any
39
 
40
+ import httpx
41
+
42
  _logger = logging.getLogger("agente_ai.skill_tracker")
43
 
44
 
45
+ # ─── Supabase config (GAP-SUP1) ───────────────────────────────────────────────
46
+
47
+ _SUPA_URL = os.getenv("SUPABASE_URL", "").rstrip("/")
48
+ _SUPA_KEY = os.getenv("SUPABASE_ANON_KEY", "")
49
+ _SUPA_ENABLED = bool(_SUPA_URL and _SUPA_KEY)
50
+ _SUPA_TABLE = "skill_stats" # tabella PostgREST — vedi gap1_skill_stats.sql
51
+ _SYNC_EVERY_N = 5 # upsert ogni N record() per tool/sessione (throttle)
52
+
53
+ if _SUPA_ENABLED:
54
+ _logger.info("[skill_tracker] Supabase persistence ABILITATA → %s/rest/v1/%s", _SUPA_URL, _SUPA_TABLE)
55
+ else:
56
+ _logger.debug("[skill_tracker] Supabase non configurato — solo in-memory")
57
+
58
+
59
  # ─── SkillStats ───────────────────────────────────────────────────────────────
60
 
61
  @dataclass
 
98
  den = 1 + z * z / n
99
  return num / den
100
 
101
+ def to_supabase_row(self, session_id: str, tool_name: str) -> dict:
102
+ """Serializza per upsert PostgREST."""
103
+ return {
104
+ "session_id": session_id,
105
+ "tool_name": tool_name,
106
+ "success_count": self.success_count,
107
+ "fail_count": self.fail_count,
108
+ "last_used": self.last_used,
109
+ "total_latency_ms": self.total_latency_ms,
110
+ }
111
+
112
+ @classmethod
113
+ def from_supabase_row(cls, row: dict) -> "SkillStats":
114
+ """Deserializza da riga PostgREST."""
115
+ return cls(
116
+ success_count = int(row.get("success_count", 0)),
117
+ fail_count = int(row.get("fail_count", 0)),
118
+ last_used = float(row.get("last_used", 0.0)),
119
+ total_latency_ms = float(row.get("total_latency_ms", 0.0)),
120
+ )
121
+
122
+
123
+ # ─── Supabase helpers (GAP-SUP1) ──────────────────────────────────────────────
124
+
125
+ async def _supabase_upsert(session_id: str, tool_name: str, stats: SkillStats) -> None:
126
+ """Fire-and-forget: upsert riga skill_stats su Supabase (PostgREST).
127
+
128
+ Fallback silente su qualsiasi errore — mai blocca il loop principale.
129
+ Timeout 8s conservativo.
130
+ """
131
+ if not _SUPA_ENABLED:
132
+ return
133
+ row = stats.to_supabase_row(session_id, tool_name)
134
+ try:
135
+ async with httpx.AsyncClient(timeout=8.0) as client:
136
+ resp = await client.post(
137
+ f"{_SUPA_URL}/rest/v1/{_SUPA_TABLE}",
138
+ json=row,
139
+ headers={
140
+ "apikey": _SUPA_KEY,
141
+ "Authorization": f"Bearer {_SUPA_KEY}",
142
+ "Content-Type": "application/json",
143
+ "Prefer": "resolution=merge-duplicates,return=minimal",
144
+ },
145
+ )
146
+ if resp.status_code not in (200, 201, 204):
147
+ _logger.debug(
148
+ "[skill_tracker] supabase upsert %s: HTTP %d %s",
149
+ tool_name[:20], resp.status_code, resp.text[:120],
150
+ )
151
+ except Exception as exc: # noqa: BLE001
152
+ _logger.debug("[skill_tracker] supabase upsert silenced: %s", type(exc).__name__)
153
+
154
+
155
+ async def _supabase_load_session(session_id: str) -> dict[str, SkillStats]:
156
+ """Carica tutti i tool stats di una sessione da Supabase.
157
+
158
+ Ritorna dict vuoto su qualsiasi errore (fallback silente).
159
+ Chiamato da load_session_from_cloud() al boot sessione.
160
+ """
161
+ if not _SUPA_ENABLED:
162
+ return {}
163
+ try:
164
+ async with httpx.AsyncClient(timeout=8.0) as client:
165
+ resp = await client.get(
166
+ f"{_SUPA_URL}/rest/v1/{_SUPA_TABLE}",
167
+ params={"session_id": f"eq.{session_id}", "select": "*"},
168
+ headers={
169
+ "apikey": _SUPA_KEY,
170
+ "Authorization": f"Bearer {_SUPA_KEY}",
171
+ },
172
+ )
173
+ if resp.status_code != 200:
174
+ _logger.debug(
175
+ "[skill_tracker] supabase load %s: HTTP %d",
176
+ session_id[:12], resp.status_code,
177
+ )
178
+ return {}
179
+ rows: list[dict] = resp.json()
180
+ loaded = {
181
+ row["tool_name"]: SkillStats.from_supabase_row(row)
182
+ for row in rows
183
+ if "tool_name" in row
184
+ }
185
+ if loaded:
186
+ _logger.info(
187
+ "[skill_tracker] GAP-SUP1: ripristinati %d tool stats per sessione %s",
188
+ len(loaded), session_id[:12],
189
+ )
190
+ return loaded
191
+ except Exception as exc: # noqa: BLE001
192
+ _logger.debug("[skill_tracker] supabase load silenced: %s", type(exc).__name__)
193
+ return {}
194
+
195
 
196
  # ─── SkillTracker ─────────────────────────────────────────────────────────────
197
 
198
  class SkillTracker:
199
+ """Singleton session-scoped tracker: impara quali tool funzionano per ogni sessione.
200
+
201
+ GAP-SUP1: i dati persistono su Supabase e vengono ripristinati al boot sessione.
202
+ """
203
 
204
  def __init__(self) -> None:
205
  self._sessions: dict[str, dict[str, SkillStats]] = defaultdict(
206
  lambda: defaultdict(SkillStats)
207
  )
208
+ # Contatori per throttle upsert (session_id → tool_name → count_since_last_sync)
209
+ self._sync_counters: dict[str, dict[str, int]] = defaultdict(lambda: defaultdict(int))
210
 
211
  # ── Write ─────────────────────────────────────────────────────────────────
212
 
 
220
  """Registra il risultato di una chiamata tool (sincrono, GIL-safe).
221
 
222
  Chiamato da unified_loop.py dopo ogni executor.run_tool().
223
+ GAP-SUP1: fire-and-forget upsert Supabase ogni _SYNC_EVERY_N chiamate.
224
  """
225
  s = self._sessions[session_id][tool_name]
226
  if success:
 
230
  s.last_used = time.monotonic()
231
  s.total_latency_ms += latency_ms
232
 
233
+ # GAP-SUP1: sync throttled — ogni _SYNC_EVERY_N record per questo tool/sessione
234
+ if _SUPA_ENABLED:
235
+ self._sync_counters[session_id][tool_name] += 1
236
+ if self._sync_counters[session_id][tool_name] >= _SYNC_EVERY_N:
237
+ self._sync_counters[session_id][tool_name] = 0
238
+ try:
239
+ loop = asyncio.get_running_loop()
240
+ loop.create_task(
241
+ _supabase_upsert(session_id, tool_name, s),
242
+ name=f"skill_sync_{tool_name[:20]}",
243
+ )
244
+ except RuntimeError:
245
+ pass # no running loop (test context) — silente
246
+
247
+ # ── Cloud bootstrap (GAP-SUP1) ────────────────────────────────────────────
248
+
249
+ async def load_session_from_cloud(self, session_id: str) -> int:
250
+ """Ripristina stats precedenti da Supabase per la sessione (chiamare al boot task).
251
+
252
+ Merge con in-memory: somma i contatori (in-memory è vuoto al boot, ma sicuro).
253
+ Ritorna numero di tool ripristinati (0 se Supabase non configurato).
254
+ Idempotente: chiamate multiple sommano i dati — chiamare una sola volta per sessione.
255
+ """
256
+ loaded = await _supabase_load_session(session_id)
257
+ if not loaded:
258
+ return 0
259
+ session = self._sessions[session_id]
260
+ for tool_name, cloud_stats in loaded.items():
261
+ mem = session[tool_name]
262
+ # Merge additivo — in-memory è tipicamente vuoto al boot
263
+ mem.success_count += cloud_stats.success_count
264
+ mem.fail_count += cloud_stats.fail_count
265
+ mem.total_latency_ms += cloud_stats.total_latency_ms
266
+ # last_used: prendi il più recente
267
+ if cloud_stats.last_used > mem.last_used:
268
+ mem.last_used = cloud_stats.last_used
269
+ return len(loaded)
270
+
271
+ async def flush_session_to_cloud(self, session_id: str) -> int:
272
+ """Forza upsert di tutti i tool di una sessione su Supabase (chiamare a fine task).
273
+
274
+ Ritorna numero di tool sincronizzati. Fallback silente su errori.
275
+ """
276
+ if not _SUPA_ENABLED:
277
+ return 0
278
+ session = self._sessions.get(session_id, {})
279
+ tasks = [
280
+ _supabase_upsert(session_id, tool_name, stats)
281
+ for tool_name, stats in session.items()
282
+ ]
283
+ if tasks:
284
+ await asyncio.gather(*tasks, return_exceptions=True)
285
+ _logger.info(
286
+ "[skill_tracker] GAP-SUP1: flush %d tool stats per sessione %s",
287
+ len(tasks), session_id[:12],
288
+ )
289
+ return len(tasks)
290
+
291
  # ── Read / routing ────────────────────────────────────────────────────────
292
 
293
  def get_sorted_fallbacks(
 
341
  def clear_session(self, session_id: str) -> None:
342
  """Libera memoria per una sessione terminata."""
343
  removed = self._sessions.pop(session_id, None)
344
+ self._sync_counters.pop(session_id, None)
345
  if removed is not None:
346
  _logger.debug(
347
  "[skill_tracker] cleared session %s (%d tools tracked)",
migrations/gap1_skill_stats.sql ADDED
@@ -0,0 +1,85 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ -- GAP-SUP1: Supabase Skill Tracker persistence
2
+ -- Eseguire nel SQL Editor di Supabase:
3
+ -- https://app.supabase.com/project/_/sql/new
4
+ --
5
+ -- Tabella skill_stats: statistiche tool per sessione (backend skill_tracker.py)
6
+ -- Tabella skill_patterns: pattern tool sequenze (frontend skillPersistence.ts)
7
+
8
+ -- ─── skill_stats ──────────────────────────────────────────────────────────────
9
+
10
+ CREATE TABLE IF NOT EXISTS skill_stats (
11
+ session_id text NOT NULL,
12
+ tool_name text NOT NULL,
13
+ success_count integer NOT NULL DEFAULT 0,
14
+ fail_count integer NOT NULL DEFAULT 0,
15
+ last_used double precision NOT NULL DEFAULT 0,
16
+ total_latency_ms double precision NOT NULL DEFAULT 0,
17
+ updated_at timestamptz NOT NULL DEFAULT now(),
18
+ PRIMARY KEY (session_id, tool_name)
19
+ );
20
+
21
+ -- Index per query per sessione (GET ?session_id=eq.xxx)
22
+ CREATE INDEX IF NOT EXISTS idx_skill_stats_session ON skill_stats (session_id);
23
+
24
+ -- Row Level Security: anon può leggere e scrivere (chiave anon pubblica)
25
+ ALTER TABLE skill_stats ENABLE ROW LEVEL SECURITY;
26
+
27
+ DO $$ BEGIN
28
+ IF NOT EXISTS (
29
+ SELECT 1 FROM pg_policies
30
+ WHERE tablename = 'skill_stats' AND policyname = 'anon full access'
31
+ ) THEN
32
+ CREATE POLICY "anon full access" ON skill_stats
33
+ FOR ALL TO anon
34
+ USING (true)
35
+ WITH CHECK (true);
36
+ END IF;
37
+ END $$;
38
+
39
+ -- Trigger: aggiorna updated_at ad ogni upsert
40
+ CREATE OR REPLACE FUNCTION _touch_updated_at()
41
+ RETURNS TRIGGER LANGUAGE plpgsql AS $$
42
+ BEGIN NEW.updated_at = now(); RETURN NEW; END;
43
+ $$;
44
+
45
+ DROP TRIGGER IF EXISTS trg_skill_stats_updated_at ON skill_stats;
46
+ CREATE TRIGGER trg_skill_stats_updated_at
47
+ BEFORE UPDATE ON skill_stats
48
+ FOR EACH ROW EXECUTE FUNCTION _touch_updated_at();
49
+
50
+
51
+ -- ─── skill_patterns ────────────────────────────────────────────────────────────
52
+ -- Persistenza cross-device per skillPersistence.ts (frontend Dexie → Supabase cloud)
53
+
54
+ CREATE TABLE IF NOT EXISTS skill_patterns (
55
+ id text PRIMARY KEY, -- djb2(taskSignature)
56
+ task_signature text NOT NULL,
57
+ tool_sequence text[] NOT NULL,
58
+ success_count integer NOT NULL DEFAULT 0,
59
+ total_count integer NOT NULL DEFAULT 0,
60
+ last_used bigint NOT NULL DEFAULT 0,
61
+ confidence double precision NOT NULL DEFAULT 0,
62
+ updated_at timestamptz NOT NULL DEFAULT now()
63
+ );
64
+
65
+ CREATE INDEX IF NOT EXISTS idx_skill_patterns_confidence ON skill_patterns (confidence DESC);
66
+ CREATE INDEX IF NOT EXISTS idx_skill_patterns_last_used ON skill_patterns (last_used DESC);
67
+
68
+ ALTER TABLE skill_patterns ENABLE ROW LEVEL SECURITY;
69
+
70
+ DO $$ BEGIN
71
+ IF NOT EXISTS (
72
+ SELECT 1 FROM pg_policies
73
+ WHERE tablename = 'skill_patterns' AND policyname = 'anon full access'
74
+ ) THEN
75
+ CREATE POLICY "anon full access" ON skill_patterns
76
+ FOR ALL TO anon
77
+ USING (true)
78
+ WITH CHECK (true);
79
+ END IF;
80
+ END $$;
81
+
82
+ DROP TRIGGER IF EXISTS trg_skill_patterns_updated_at ON skill_patterns;
83
+ CREATE TRIGGER trg_skill_patterns_updated_at
84
+ BEFORE UPDATE ON skill_patterns
85
+ FOR EACH ROW EXECUTE FUNCTION _touch_updated_at();