#!/usr/bin/env python3 """vault_guardian.py — Census orfani su vault.db + invio anomalie a Elon. Sola segnalazione. Nessuna correzione — il giudizio «collegare o lasciare» resta a Elon/Mauro. Gira via cron (04:30 ogni notte): - Produce census (nodi 0-out, 0-in, relazioni dangling) - Se ci sono anomalie che richiedono giudizio → msg-a2e a Elon - Notifica Telegram a Mauro se anomalie nuove rispetto al run precedente """ import json import os import sqlite3 import subprocess import sys from datetime import datetime, timezone from pathlib import Path import httpx from dotenv import load_dotenv load_dotenv(Path(__file__).parent.parent / ".env") VAULT_DB = Path("/mnt/ssd/data/vault-secondbrain/vault.db") TG_TOKEN = os.getenv("ADRIAN_BOT_TOKEN", "") TG_CHAT_ID = os.getenv("TELEGRAM_CHAT_ID", "") REPORT_ID = "vault-guardian-report" def tg_notify(text: str): if not TG_TOKEN or not TG_CHAT_ID: return try: httpx.post( f"https://api.telegram.org/bot{TG_TOKEN}/sendMessage", json={"chat_id": TG_CHAT_ID, "text": text}, timeout=10, ) except Exception: pass def run_census(conn: sqlite3.Connection) -> dict: zero_out = conn.execute(""" SELECT i.id, json_extract(i.meta,'$.tipo') as tipo, json_extract(i.meta,'$.nome') as nome FROM items i WHERE i.type='nodo' AND json_extract(i.meta,'$.tipo') != 'limbo' AND NOT EXISTS (SELECT 1 FROM edges e WHERE e.src=i.id) """).fetchall() zero_in = conn.execute(""" SELECT i.id, json_extract(i.meta,'$.tipo') as tipo, json_extract(i.meta,'$.nome') as nome FROM items i WHERE i.type='nodo' AND json_extract(i.meta,'$.tipo') != 'limbo' AND NOT EXISTS (SELECT 1 FROM edges e WHERE e.dst=i.id) """).fetchall() nodi = conn.execute("SELECT id, meta FROM items WHERE type='nodo'").fetchall() dangling = [] for nid, meta_str in nodi: meta = json.loads(meta_str or "{}") for rel in meta.get("relazioni", []): dst = rel.get("id") if dst: exists = conn.execute( "SELECT 1 FROM items WHERE id=? AND type='nodo'", (dst,) ).fetchone() if not exists: dangling.append({"src": nid, "dst": dst}) return { "zero_out": [{"id": r[0], "tipo": r[1], "nome": r[2]} for r in zero_out], "zero_in": [{"id": r[0], "tipo": r[1], "nome": r[2]} for r in zero_in], "dangling": dangling, } def format_report(census: dict, ts: str) -> str: lines = [f"vault-guardian — {ts}"] lines.append( f"nodi 0-out: {len(census['zero_out'])} · " f"nodi 0-in: {len(census['zero_in'])} · " f"dangling: {len(census['dangling'])}" ) if census["zero_out"]: lines.append("\n0 archi uscenti:") for n in census["zero_out"]: lines.append(f" {n['id']} [{n['tipo']}]") if census["zero_in"]: lines.append("\n0 archi entranti:") for n in census["zero_in"]: lines.append(f" {n['id']} [{n['tipo']}]") if census["dangling"]: lines.append("\nRelazioni irrisolte (dst non esiste):") for d in census["dangling"]: lines.append(f" {d['src']} → {d['dst']}") return "\n".join(lines) def load_previous(conn: sqlite3.Connection) -> dict | None: row = conn.execute("SELECT meta FROM items WHERE id=?", (REPORT_ID,)).fetchone() if not row: return None try: return json.loads(row[0] or "{}") except Exception: return None def has_new_anomalies(prev: dict | None, curr: dict) -> bool: if prev is None: return bool(curr["zero_out"] or curr["zero_in"] or curr["dangling"]) prev_ids = { "zero_out": {n["id"] for n in prev.get("zero_out", [])}, "zero_in": {n["id"] for n in prev.get("zero_in", [])}, "dangling": {(d["src"], d["dst"]) for d in prev.get("dangling", [])}, } curr_ids = { "zero_out": {n["id"] for n in curr["zero_out"]}, "zero_in": {n["id"] for n in curr["zero_in"]}, "dangling": {(d["src"], d["dst"]) for d in curr["dangling"]}, } return bool( curr_ids["zero_out"] - prev_ids["zero_out"] or curr_ids["zero_in"] - prev_ids["zero_in"] or curr_ids["dangling"] - prev_ids["dangling"] ) def needs_elon(census: dict) -> bool: """Gatekeeper: solo anomalie che richiedono giudizio semantico di Elon. Nodi 0-out/0-in potrebbero essere foglie/radici intenzionali — conta solo se ci sono relazioni dangling (link nel vuoto = errore certo) o se il numero di orfani è aumentato rispetto al run precedente.""" return bool(census["dangling"]) def send_elon_maintenance(conn: sqlite3.Connection, census: dict, ts: str) -> str: """Crea msg-a2e con le anomalie per Elon. Ritorna il msg_id.""" ts_micro = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S%fZ") msg_id = f"msg-a2e-{ts_micro}" lines = [ f"Manutenzione notturna — {ts}", "", "vault_guardian ha trovato anomalie che richiedono il tuo giudizio.", "Sola segnalazione — nessuna azione automatica è stata presa.", "", ] if census["dangling"]: lines.append("RELAZIONI IRRISOLTE (dst non esiste in vault.db):") for d in census["dangling"]: lines.append(f" {d['src']} → {d['dst']}") lines.append(" → collegare a nodo esistente, creare stub, o rimuovere la relazione\n") if census["zero_out"]: lines.append("NODI SENZA ARCHI USCENTI:") for n in census["zero_out"]: lines.append(f" {n['id']} [{n['tipo']}]") lines.append(" → collegare ad altri nodi o valutare se sono foglie intenzionali\n") if census["zero_in"]: lines.append("NODI SENZA ARCHI ENTRANTI:") for n in census["zero_in"]: lines.append(f" {n['id']} [{n['tipo']}]") lines.append(" → collegare da altri nodi o valutare se sono radici intenzionali\n") lines.append("Manda ack-e2a quando hai gestito. — Adrian") body = "\n".join(lines) meta = json.dumps({ "tipo": "agent_memory", "from": "adrian", "to": "elon", "ts": ts_micro, "subject": "manutenzione-notturna", "guardian_ts": ts, }) conn.execute( "INSERT OR REPLACE INTO items(id, type, body, file, meta) VALUES (?,?,?,NULL,?)", (msg_id, "os", body, meta), ) return msg_id def main(): ts = datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") conn = sqlite3.connect(VAULT_DB) conn.row_factory = sqlite3.Row census = run_census(conn) report_body = format_report(census, ts) prev = load_previous(conn) meta = json.dumps({**census, "last_run": ts}) conn.execute( "INSERT OR REPLACE INTO items(id, type, body, file, meta) VALUES (?,?,?,NULL,?)", (REPORT_ID, "os", report_body, meta), ) elon_msg_id = None if needs_elon(census): elon_msg_id = send_elon_maintenance(conn, census, ts) conn.commit() conn.close() print(report_body) if elon_msg_id: print(f"→ msg-a2e scritto per Elon: {elon_msg_id}") new_anomalies = has_new_anomalies(prev, census) if new_anomalies: total = len(census["zero_out"]) + len(census["zero_in"]) + len(census["dangling"]) tg_notify( f"⚠️ vault-guardian: {total} anomalie nuove in vault.db\n" f"{report_body[:300]}" + ("\n→ Elon notificato via #boot" if elon_msg_id else "") ) # Triggera Elon via #boot solo se c'è qualcosa per lui (gatekeeper) if elon_msg_id: elon_pull = Path(__file__).parent.parent / "vault-secondbrain" / "elon_pull.py" env = {**os.environ, "DISPLAY": ":0"} result = subprocess.run( [sys.executable, str(elon_pull), "--fire-and-forget"], env=env, capture_output=True, text=True, timeout=120, ) print(result.stdout) if result.returncode != 0: print(f"[guardian] elon_pull fallito (exit {result.returncode}): {result.stderr[:200]}") else: print("→ Nessuna anomalia che richiede Elon — /pull non inviato.") if __name__ == "__main__": main()