#!/usr/bin/env python3 """Loop Adrian↔Elon — protocollo msg/ack su vault.db con handshake Vision. Protocollo: msg-a2e-{ts} : Adrian → Elon (pending se no ack-a2e) ack-a2e-{ts} : Elon → conferma lettura msg-e2a-{ts} : Elon → Adrian (risposta/aggiornamento) ack-e2a-{ts} : Adrian → conferma lettura State machine pre-pull (Vision-driven): CHROME_DOWN → avvia Chrome WRONG_PAGE → naviga al Project GENERATING → aspetta max 2 min READY_WARM → sessione attiva, kernel in memoria → testo diretto READY_COLD → sessione fredda, kernel da caricare → /pull AUTH_ERROR → Telegram Mauro, stop CAPTCHA → Telegram Mauro, stop ERROR → Telegram Mauro, stop Uso: python3 elon_pull.py # interattivo, aspetta risposta python3 elon_pull.py --fire-and-forget # non-blocking (manutenzione notturna) python3 elon_pull.py --retry-if-pending # cron: ritenta solo se msg-a2e senza ack >2h python3 elon_pull.py --restart-chrome # riavvia solo Chrome """ import base64 import json import os import re import signal import sqlite3 import subprocess import sys import time import urllib.request from datetime import datetime, timezone from pathlib import Path import httpx from dotenv import load_dotenv from playwright.sync_api import sync_playwright load_dotenv(Path(__file__).parent.parent / ".env") PROJECT_URL = "https://claude.ai/project/019c85a9-a4ba-737c-a803-b10f4a7769b4" PROJECT_UUID = "019c85a9-a4ba-737c-a803-b10f4a7769b4" CDP_URL = "http://localhost:9222" DB_PATH = Path("/mnt/ssd/data/adrian-ops/adrian.db") DEBUG_DIR = Path("/tmp/elon_pull_debug") DEBUG_DIR.mkdir(exist_ok=True) POLL_INTERVAL = 5 RESPONSE_TIMEOUT = 300 # secondi attesa risposta Elon (run interattivo) GENERATE_WAIT = 120 # secondi max aspetta fine generazione PENDING_AGE_H = 2 # ore prima del retry automatico TELEGRAM_TOKEN = os.getenv("ADRIAN_BOT_TOKEN", "") TELEGRAM_CHAT = os.getenv("TELEGRAM_CHAT_ID", "") OPENROUTER_URL = os.getenv("OPENROUTER_BASE_URL", "https://openrouter.ai/api/v1") OPENROUTER_KEY = os.getenv("OPENROUTER_API_KEY", "") VISION_MODEL = "anthropic/claude-sonnet-4-5" DIRECT_MAX_LEN = 2000 # msg-a2e più lunghi di questo → /pull invece di testo diretto # ── DB helpers ──────────────────────────────────────────────────────────────── def _db() -> sqlite3.Connection: conn = sqlite3.connect(DB_PATH) conn.row_factory = sqlite3.Row return conn def _ts() -> str: return datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S%fZ") def send_message(text: str) -> str: """Scrive un msg-a2e per Elon in vault.db. Ritorna il msg_id.""" ts = _ts() msg_id = f"msg-a2e-{ts}" meta = json.dumps({"tipo": "agent_memory", "from": "adrian", "to": "elon", "ts": ts}) conn = _db() conn.execute( "INSERT OR REPLACE INTO items(id, type, body, file, meta) VALUES (?,?,?,NULL,?)", (msg_id, "os", text, meta), ) conn.commit(); conn.close() return msg_id def pending_msgs_a2e() -> list[dict]: """msg-a2e senza ack-a2e — non ancora ricevuti da Elon.""" conn = _db() rows = conn.execute(""" SELECT i.id, i.body, json_extract(i.meta,'$.ts') as ts FROM items i WHERE i.id LIKE 'msg-a2e-%' AND NOT EXISTS ( SELECT 1 FROM items a WHERE a.id = 'ack-' || substr(i.id, 5) ) ORDER BY ts ASC """).fetchall() conn.close() return [dict(r) for r in rows] def stale_msgs_a2e(hours: float = PENDING_AGE_H) -> list[dict]: """msg-a2e pendenti da più di `hours` ore — candidati al retry.""" cutoff = datetime.now(timezone.utc).timestamp() - hours * 3600 pending = pending_msgs_a2e() stale = [] for m in pending: try: ts_str = (m.get("ts") or "").replace("Z", "+00:00") t = datetime.fromisoformat(ts_str).timestamp() if t < cutoff: stale.append(m) except Exception: stale.append(m) return stale def pending_responses() -> list[dict]: """msg-e2a senza ack-e2a — risposte di Elon non ancora lette.""" conn = _db() rows = conn.execute(""" SELECT i.id, i.body, json_extract(i.meta,'$.ts') as ts FROM items i WHERE i.id LIKE 'msg-e2a-%' AND NOT EXISTS ( SELECT 1 FROM items a WHERE a.id = 'ack-' || substr(i.id, 5) ) ORDER BY ts ASC """).fetchall() conn.close() return [dict(r) for r in rows] def ack_response(msg_id: str) -> None: """Adrian acka una risposta di Elon (msg-e2a → ack-e2a).""" ack_id = "ack-" + msg_id[4:] ts = _ts() meta = json.dumps({"tipo": "agent_memory", "acks": msg_id, "by": "adrian", "ts": ts}) conn = _db() conn.execute( "INSERT OR IGNORE INTO items(id, type, body, file, meta) VALUES (?,?,?,NULL,?)", (ack_id, "os", "", meta), ) conn.commit(); conn.close() # ── Telegram ───────────────────────────────────────────────────────────────── def _telegram(text: str) -> None: if not TELEGRAM_TOKEN or not TELEGRAM_CHAT: return try: httpx.post( f"https://api.telegram.org/bot{TELEGRAM_TOKEN}/sendMessage", json={"chat_id": TELEGRAM_CHAT, "text": text}, timeout=10, ) except Exception: pass # ── Chrome ──────────────────────────────────────────────────────────────────── CHROME_PROFILE = Path.home() / ".config/chrome-debug-elon" CHROME_SENTINEL = Path("/tmp/elon_pull_chrome.sentinel") # presente solo se Chrome è nostro CHROME_CMD = [ "google-chrome-stable", "--remote-debugging-port=9222", f"--user-data-dir={CHROME_PROFILE}", "--no-first-run", "--profile-directory=Default", "--start-maximized", ] def restart_chrome() -> None: result = subprocess.run(["pgrep", "-f", "chrome-debug-elon"], capture_output=True, text=True) pids = [int(p) for p in result.stdout.split() if p.strip()] for pid in pids: try: os.kill(pid, signal.SIGTERM) except ProcessLookupError: pass if pids: print(f"[chrome] Terminati {len(pids)} processi (pids: {pids})") time.sleep(3) CHROME_SENTINEL.unlink(missing_ok=True) env = {**os.environ, "DISPLAY": ":0"} subprocess.Popen(CHROME_CMD, env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) print("[chrome] Chrome riavviato — attendo CDP...") for _ in range(20): time.sleep(1) try: urllib.request.urlopen("http://localhost:9222/json", timeout=2) CHROME_SENTINEL.write_text("adrian") print("[chrome] CDP pronto.") return except Exception: pass print("[chrome] WARNING: CDP non risponde dopo 20s") def _chrome_up() -> bool: try: urllib.request.urlopen(CDP_URL + "/json", timeout=2) return True except Exception: return False def _chrome_is_ours() -> bool: """True se il Chrome su porta 9222 è stato avviato da Adrian (sentinel presente).""" return CHROME_SENTINEL.exists() and _chrome_up() # ── Vision state detection ──────────────────────────────────────────────────── VISION_PROMPT = """\ Analizza questo screenshot di Claude.ai e rispondi SOLO con un JSON: { "state": "ready_warm" | "ready_cold" | "generating" | "auth_error" | "captcha" | "error", "kernel_loaded": true | false, "details": "" } Definizioni: - ready_warm: input box visibile e libero, conversazione con messaggi visibili (sessione attiva) - ready_cold: input box visibile e libero, nessun messaggio nel thread (sessione nuova/vuota) - generating: spinner di caricamento, pulsante "Stop", o testo che si sta generando - auth_error: pagina di login, "Session expired", errore di autenticazione - captcha: captcha Cloudflare o simile - error: pagina rotta, 404, errore generico non classificabile kernel_loaded: true se vedi messaggi precedenti di una sessione strutturata (non solo UI vuota). Rispondi SOLO con il JSON, nessun testo aggiuntivo.""" def _vision_state(img_bytes: bytes) -> dict: """Manda screenshot a Claude Vision e ritorna lo stato rilevato.""" if not OPENROUTER_KEY: return {"state": "ready_cold", "kernel_loaded": False, "details": "vision non configurata"} try: b64 = base64.b64encode(img_bytes).decode() resp = httpx.post( f"{OPENROUTER_URL}/chat/completions", headers={"Authorization": f"Bearer {OPENROUTER_KEY}", "Content-Type": "application/json"}, json={ "model": VISION_MODEL, "messages": [{ "role": "user", "content": [ {"type": "image_url", "image_url": {"url": f"data:image/png;base64,{b64}"}}, {"type": "text", "text": VISION_PROMPT}, ], }], "temperature": 0, "max_tokens": 150, }, timeout=30, ) resp.raise_for_status() raw = resp.json()["choices"][0]["message"]["content"].strip() m = re.search(r'\{.*\}', raw, re.DOTALL) if m: return json.loads(m.group()) except Exception as e: print(f"[vision] Errore: {e}") return {"state": "ready_cold", "kernel_loaded": False, "details": "vision fallita, assumo cold"} def _detect_state(page) -> dict: """Screenshot + Vision → stato pagina.""" try: img = page.screenshot() fname = DEBUG_DIR / "pre_pull.png" fname.write_bytes(img) state = _vision_state(img) print(f"[state] {state['state']} | kernel={state['kernel_loaded']} | {state['details']}") return state except Exception as e: print(f"[state] Errore detection: {e}") return {"state": "error", "kernel_loaded": False, "details": str(e)} # ── Connector check ─────────────────────────────────────────────────────────── def _check_connector(page) -> bool: """Check passivo — nessun click, nessun menu aperto. Legge il testo della pagina: se Claude.ai segnala 'non disponibile' per il vault lo connettore è offline. Zero interazioni UI, zero overlay.""" try: content = page.content() if "Vault Second Brain non disponibile" in content or "vault second brain non disponibile" in content.lower(): print("[connector] ⚠️ Vault Second Brain non disponibile (rilevato passivamente)") return False return True except Exception as e: print(f"[connector] Errore check passivo: {e}") return True # ── Trigger adattivo ────────────────────────────────────────────────────────── def _send_trigger(page, state: dict, pending: list[dict]) -> None: """Manda il trigger a Elon: testo diretto (warm) o #boot (cold). Il prefisso # segnala comunicazione unattended Adrian↔Elon.""" input_box = page.locator('[contenteditable="true"]').last input_box.click() page.wait_for_timeout(300) warm = state.get("state") == "ready_warm" and state.get("kernel_loaded", False) content = "\n\n---\n\n".join(m["body"] for m in pending) if pending else "" if warm and pending and len(content) <= DIRECT_MAX_LEN: # Warm session: manda testo diretto con prefisso # (unattended) text = "# " + content + "\n\n(Acka i msg-a2e corrispondenti in vault.db come da protocollo.)" print(f"[trigger] WARM — testo diretto ({len(content)} chars, {len(pending)} msg)") else: # Cold session o msg troppo lunghi: #boot standard reason = "cold" if not warm else f"msg troppo lungo ({len(content)} chars)" print(f"[trigger] COLD ({reason}) — #boot standard") text = "#boot" input_box.fill(text) page.keyboard.press("Enter") # ── Sessione Playwright ─────────────────────────────────────────────────────── def _get_elon_page(browser): """Trova la pagina del Project Elon o None.""" for pg in browser.contexts[0].pages: if PROJECT_UUID in pg.url: return pg return None def _setup_browser(p): """Connette CDP, trova/crea la pagina Elon. Ritorna (browser, page).""" if not _chrome_up(): print("[loop] Chrome non in esecuzione — avvio...") restart_chrome() elif not _chrome_is_ours(): print("[loop] Chrome attivo ma non istanziato da Adrian — chiudo e riavvio pulito...") restart_chrome() else: print("[loop] Chrome già nostro — riuso sessione esistente") browser = p.chromium.connect_over_cdp(CDP_URL) # Maximize la finestra così Vision vede lo stato completo page = browser.contexts[0].pages[0] if browser.contexts[0].pages else None if page: page.evaluate("() => { window.moveTo(0,0); window.resizeTo(screen.width, screen.height); }") page = _get_elon_page(browser) if not page: print("[loop] Project non trovato — navigo...") page = browser.contexts[0].new_page() page.evaluate("() => { window.moveTo(0,0); window.resizeTo(screen.width, screen.height); }") page.goto(PROJECT_URL, wait_until="domcontentloaded", timeout=60000) page.wait_for_timeout(3000) return browser, page def _wait_if_generating(page, max_wait: int = GENERATE_WAIT) -> bool: """Aspetta che Elon finisca di generare. True se pronto, False se timeout.""" start = time.time() while time.time() - start < max_wait: state = _detect_state(page) if state["state"] != "generating": return True print(f"[loop] Elon genera — aspetto... ({int(time.time()-start)}s)") time.sleep(5) return False # ── Entrypoint principale ───────────────────────────────────────────────────── def _pull_core(wait_response: bool = False) -> bool: """Logica comune: state machine + trigger. Ritorna True se pull inviato.""" # Trasferisci canale legacy se presente conn = _db() legacy = conn.execute("SELECT body FROM items WHERE id='da-adrian-a-elon'").fetchone() conn.close() if legacy and (legacy["body"] or "").strip(): msg_id = send_message(legacy["body"].strip()) conn = _db() conn.execute("UPDATE items SET body='' WHERE id='da-adrian-a-elon'") conn.commit(); conn.close() print(f"[loop] Messaggio legacy → {msg_id}") pending = pending_msgs_a2e() if not pending: print("[loop] Nessun msg-a2e pendente — niente da fare.") return False print(f"[loop] {len(pending)} msg-a2e pendenti") try: with sync_playwright() as p: browser, page = _setup_browser(p) # State machine state = _detect_state(page) if state["state"] == "generating": ready = _wait_if_generating(page) if not ready: print("[loop] Timeout generazione — rinuncio") _telegram("⚠️ elon_pull: Elon stava generando da troppo tempo, /pull posticipato") browser.close() return False state = _detect_state(page) if state["state"] in ("auth_error", "captcha"): msg = f"⚠️ elon_pull: {state['state']} rilevato — intervento Mauro richiesto\n{state['details']}" print(f"[loop] {msg}") _telegram(msg) browser.close() return False if state["state"] == "error": print(f"[loop] Errore pagina: {state['details']} — ricarico Project") page.goto(PROJECT_URL, wait_until="domcontentloaded", timeout=60000) page.wait_for_timeout(3000) state = _detect_state(page) # Check connettore — non bloccante (falsi positivi da sidebar conversazioni) if not _check_connector(page): print("[loop] ⚠️ Connettore Vault potenzialmente offline — procedo comunque") _telegram("⚠️ elon_pull: connettore Vault potrebbe essere offline (Elon verificherà)") # Trigger adattivo _send_trigger(page, state, pending) print("[loop] Trigger inviato.") if wait_response: print("[loop] Attendo risposta di Elon...") elapsed = 0 responses = [] while elapsed < RESPONSE_TIMEOUT: time.sleep(POLL_INTERVAL) elapsed += POLL_INTERVAL responses = pending_responses() if responses: print(f"[loop] Elon ha risposto ({elapsed}s).") break if elapsed % 30 == 0: page.screenshot(path=str(DEBUG_DIR / f"poll_{elapsed}s.png")) print(f"[loop] Attendo... ({elapsed}s)") browser.close() if responses: for r in responses: print(f"\n--- Risposta Elon [{r['id']}] ---\n{r['body']}") ack_response(r["id"]) else: print("[loop] Timeout — nessuna risposta.") else: print("[loop] /pull inviato (fire-and-forget). Elon risponderà alla prossima sessione.") browser.close() return True except Exception as e: print(f"[loop] Errore: {e}") _telegram(f"⚠️ elon_pull fallito: {e}") return False def run() -> None: _pull_core(wait_response=True) def fire_and_forget() -> bool: return _pull_core(wait_response=False) def retry_if_pending() -> bool: """Ritenta pull solo se ci sono msg-a2e senza ack da più di PENDING_AGE_H ore.""" stale = stale_msgs_a2e() if not stale: print(f"[retry] Nessun msg-a2e stale (>{PENDING_AGE_H}h) — niente da fare.") return False print(f"[retry] {len(stale)} msg-a2e senza ack da >{PENDING_AGE_H}h — ritento pull.") return fire_and_forget() if __name__ == "__main__": if "--restart-chrome" in sys.argv: restart_chrome() elif "--fire-and-forget" in sys.argv: ok = fire_and_forget() sys.exit(0 if ok else 1) elif "--retry-if-pending" in sys.argv: ok = retry_if_pending() sys.exit(0 if ok else 1) else: run()