""" Daemon che monitora inbox/ e ingerisce nuovi documenti. Polling ogni 30 secondi. Avviato come systemd service. Classificazione: delega ad Adrian (via classify-request in vault.db). Adrian legge il file, classifica e scrive classify-result entro 5 minuti. Fallback API se Adrian non risponde in tempo. """ import json import sqlite3 import time from datetime import datetime, timezone from pathlib import Path from loguru import logger from ingest import INBOX, VAULT_DB, compute_doc_id LOG_PATH = Path(__file__).parent.parent / "logs" / "inbox_watch.log" POLL_INTERVAL = 30 # secondi ADRIAN_TIMEOUT = 300 # 5 minuti — poi fallback API SUPPORTED_EXT = {".pdf", ".docx", ".doc", ".txt", ".md"} logger.add(str(LOG_PATH), rotation="10 MB", retention="30 days", level="INFO") def _db(): conn = sqlite3.connect(VAULT_DB) conn.row_factory = sqlite3.Row return conn def _ts(): return datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S%fZ") def request_classify(doc_id: str, path: Path) -> None: """Scrive classify-request-{doc_id} per Adrian.""" req_id = f"classify-request-{doc_id}" meta = json.dumps({"doc_id": doc_id, "path": str(path), "ts": _ts()}) conn = _db() conn.execute( "INSERT OR REPLACE INTO items(id, type, body, file, meta) VALUES (?,?,?,NULL,?)", (req_id, "os", str(path), meta), ) conn.commit() conn.close() def poll_classify_result(doc_id: str, timeout: int = ADRIAN_TIMEOUT) -> dict | None: """Aspetta classify-result-{doc_id} da Adrian. Ritorna il dict o None se timeout.""" result_id = f"classify-result-{doc_id}" deadline = time.time() + timeout while time.time() < deadline: conn = _db() row = conn.execute( "SELECT body FROM items WHERE id=?", (result_id,) ).fetchone() conn.close() if row and row["body"]: try: return json.loads(row["body"]) except Exception: return None time.sleep(5) return None def cleanup_classify(doc_id: str) -> None: """Rimuove request e result dopo l'uso.""" conn = _db() conn.execute("DELETE FROM items WHERE id IN (?,?)", (f"classify-request-{doc_id}", f"classify-result-{doc_id}")) conn.commit() conn.close() def scan_inbox() -> None: from ingest import extract_text, is_already_ingested, ingest_file_with_classify for path in sorted(INBOX.rglob("*")): if not path.is_file() or path.name.startswith("."): continue if path.suffix.lower() not in SUPPORTED_EXT: logger.warning(f"Tipo non supportato, skip: {path.name}") continue doc_id = compute_doc_id(path) if is_already_ingested(doc_id): continue # Controlla se c'è già una request pendente per questo file conn = _db() req = conn.execute( "SELECT id FROM items WHERE id=?", (f"classify-request-{doc_id}",) ).fetchone() conn.close() if req: # Request già inviata — controlla se Adrian ha risposto result = poll_classify_result(doc_id, timeout=5) # check rapido if result: logger.info(f"Classificazione ricevuta da Adrian: {path.name}") cleanup_classify(doc_id) ingest_file_with_classify(path, doc_id, result) # altrimenti aspetta il prossimo ciclo continue # Nuovo file — chiedi classificazione ad Adrian logger.info(f"Nuovo file, richiesta classificazione ad Adrian: {path.name}") request_classify(doc_id, path) def main() -> None: logger.info("inbox_watch avviato — polling ogni {}s su {}", POLL_INTERVAL, INBOX) INBOX.mkdir(parents=True, exist_ok=True) while True: try: scan_inbox() except Exception as e: logger.error(f"Errore scan: {e}") time.sleep(POLL_INTERVAL) if __name__ == "__main__": main()