""" Pipeline di ingestione documenti per vault-secondbrain. Estrae testo, classifica con LLM, chunk + embedding → LanceDB, archivia in docs/. """ import hashlib import json import os import re import shutil import sqlite3 import textwrap from datetime import datetime from pathlib import Path import httpx import lancedb import numpy as np from dotenv import load_dotenv from sentence_transformers import SentenceTransformer load_dotenv(Path(__file__).parent.parent / ".env") # --- Paths --- VAULT_DB = Path("/mnt/ssd/data/vault-secondbrain/vault.db") LANCEDB = Path("/mnt/ssd/data/adrian-ops/lancedb") INBOX = Path("/mnt/ssd/data/adrian-ops/inbox") DOCS = Path("/mnt/ssd/data/adrian-ops/docs") # --- Config --- CHUNK_SIZE = 1800 CHUNK_OVERLAP = 250 EMBED_MODEL = "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2" LLM_MODEL = "anthropic/claude-haiku-4-5" PREVIEW_CHARS = 3000 # chars da inviare all'LLM per classificazione OPENROUTER_URL = os.getenv("OPENROUTER_BASE_URL", "https://openrouter.ai/api/v1") OPENROUTER_KEY = os.getenv("OPENROUTER_API_KEY", "") TELEGRAM_TOKEN = os.getenv("ADRIAN_BOT_TOKEN", "") TELEGRAM_CHAT = os.getenv("TELEGRAM_CHAT_ID", "") _embedder = None def get_embedder() -> SentenceTransformer: global _embedder if _embedder is None: _embedder = SentenceTransformer(EMBED_MODEL) return _embedder # --- Estrazione testo --- def extract_text(path: Path) -> str: suffix = path.suffix.lower() if suffix == ".pdf": return _extract_pdf(path) elif suffix in (".docx", ".doc"): return _extract_docx(path) elif suffix in (".txt", ".md"): return path.read_text(errors="ignore") else: raise ValueError(f"Tipo file non supportato: {suffix}") def _extract_pdf(path: Path) -> str: from pypdf import PdfReader reader = PdfReader(str(path)) pages = [] for page in reader.pages: t = page.extract_text() or "" pages.append(t) return "\n".join(pages) def _vision_pdf(path: Path) -> tuple[str, dict | None]: """Converte PDF scansionato in immagini e le manda a Claude Vision. Ritorna (testo_estratto, classificazione) — classificazione può essere None se fallisce.""" import base64 from pdf2image import convert_from_path from io import BytesIO print(f"[ingest] Vision in corso: {path.name}") images = convert_from_path(str(path), dpi=200, fmt="jpeg") # Codifica immagini in base64 (max 5 pagine per non superare i limiti token) content = [] for img in images[:5]: buf = BytesIO() img.save(buf, format="JPEG", quality=85) b64 = base64.b64encode(buf.getvalue()).decode() content.append({ "type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{b64}"}, }) content.append({ "type": "text", "text": ( "Analizza questo documento di Mauro Gagliardi (famiglia: moglie Nadia, figli Samuele e Francesco).\n" "Restituisci un JSON con questi campi:\n" "{\n" ' "ambito": "lavoro" | "personale" | "famiglia",\n' ' "categoria": "",\n' ' "soggetto": "",\n' ' "data_doc": "",\n' ' "confident": true | false,\n' ' "motivo": "",\n' ' "testo_estratto": ""\n' "}\n" "Regole:\n" "- ambito lavoro: documenti aziendali Mediaset, buste paga, contratti di lavoro\n" "- ambito famiglia: documenti di Samuele, Francesco, Nadia\n" "- ambito personale: documenti intestati a Mauro ma non lavorativi\n" "- categoria = NATURA del documento (cosa è), NON la persona. Es: fattura dentista di Samuele → categoria='medico/odontoiatria', soggetto='samuele'\n" "Rispondi SOLO con il JSON." ), }) try: resp = httpx.post( f"{OPENROUTER_URL}/chat/completions", headers={ "Authorization": f"Bearer {OPENROUTER_KEY}", "Content-Type": "application/json", }, json={ "model": "anthropic/claude-sonnet-4-5", "messages": [{"role": "user", "content": content}], "temperature": 0, "max_tokens": 2000, }, timeout=60, ) resp.raise_for_status() raw = resp.json()["choices"][0]["message"]["content"].strip() match = re.search(r'\{.*\}', raw, re.DOTALL) if not match: raise ValueError(f"Nessun JSON nella risposta Vision: {raw[:200]}") data = json.loads(match.group()) testo = data.pop("testo_estratto", "") or "" print(f"[ingest] Vision completato: {len(testo)} chars estratti") return testo, data except Exception as e: print(f"[ingest] Vision fallito: {e}") return "", None def _extract_docx(path: Path) -> str: from docx import Document doc = Document(str(path)) return "\n".join(p.text for p in doc.paragraphs if p.text.strip()) # --- Classificazione LLM --- CLASSIFY_PROMPT = """\ Sei un assistente che classifica documenti personali di Mauro Gagliardi (famiglia: moglie Nadia, figli Samuele e Francesco). Analizza il testo fornito e restituisci un JSON con questi campi: {{ "ambito": "lavoro" | "personale" | "famiglia", "categoria": "", "soggetto": "", "data_doc": "", "confident": true | false, "motivo": "" }} Regole: - ambito "lavoro": documenti aziendali, contratti di lavoro, buste paga, comunicazioni ufficio Mediaset - ambito "famiglia": documenti che riguardano i figli (Samuele, Francesco) o la moglie (Nadia) - ambito "personale": documenti intestati a Mauro ma non lavorativi (casa, fiscale, medico, auto, ecc.) - categoria = NATURA del documento (cosa è il documento), NON la persona — es. una fattura dentista è "medico/odontoiatria" anche se intestata a Samuele - soggetto = persona fisica coinvolta (null se il documento riguarda Mauro stesso o è generico) - confident=false se il testo è troppo breve, illeggibile, o ambiguo - Rispondi SOLO con il JSON, nessun testo aggiuntivo. Testo documento (prime {n} parole): --- {testo} --- """ def classify_document(text: str) -> dict: preview = text[:PREVIEW_CHARS].strip() prompt = CLASSIFY_PROMPT.format(testo=preview, n=len(preview.split())) try: resp = httpx.post( f"{OPENROUTER_URL}/chat/completions", headers={ "Authorization": f"Bearer {OPENROUTER_KEY}", "Content-Type": "application/json", }, json={ "model": LLM_MODEL, "messages": [{"role": "user", "content": prompt}], "temperature": 0, "max_tokens": 200, }, timeout=30, ) resp.raise_for_status() content = resp.json()["choices"][0]["message"]["content"].strip() # Estrai il blocco JSON robusto: cerca il primo { ... } nell'intera risposta match = re.search(r'\{.*\}', content, re.DOTALL) if not match: raise ValueError(f"Nessun JSON trovato nella risposta LLM: {content[:200]}") return json.loads(match.group()) except Exception as e: return { "ambito": "da-organizzare", "categoria": "inbox", "data_doc": None, "confident": False, "motivo": f"errore classificazione: {e}", } # --- Chunking --- def chunk_text(text: str) -> list[str]: text = " ".join(text.split()) # normalizza whitespace if not text: return [] chunks = [] start = 0 while start < len(text): end = start + CHUNK_SIZE chunks.append(text[start:end]) start += CHUNK_SIZE - CHUNK_OVERLAP return chunks # --- Embedding --- def embed_chunks(chunks: list[str]) -> list[list[float]]: if not chunks: return [] model = get_embedder() vecs = model.encode(chunks, show_progress_bar=False, normalize_embeddings=True) return [v.tolist() for v in vecs] # --- Vault DB --- def compute_doc_id(path: Path) -> str: h = hashlib.sha256(path.read_bytes()).hexdigest() return h[:16] def is_already_ingested(doc_id: str) -> bool: conn = sqlite3.connect(VAULT_DB) row = conn.execute( "SELECT id FROM items WHERE type='documento' AND json_extract(meta,'$.doc_id')=?", (doc_id,) ).fetchone() conn.close() return row is not None def save_documento(doc_id: str, nome_file: str, classificazione: dict, n_chunks: int, path_archivio: str) -> None: item_id = f"doc-{doc_id}" meta = { "doc_id": doc_id, "nome_file": nome_file, "fonte": "inbox", "ambito": classificazione["ambito"], "categoria": classificazione["categoria"], "soggetto": classificazione.get("soggetto"), "data_doc": classificazione.get("data_doc"), "tipo": Path(nome_file).suffix.lower().lstrip("."), "stato": "ingerito", "path_archivio": path_archivio, "n_chunks": n_chunks, "classificato_auto": classificazione.get("confident", False), "ingested_at": datetime.utcnow().isoformat(), } conn = sqlite3.connect(VAULT_DB) conn.execute( "INSERT OR REPLACE INTO items(id, type, body, meta) VALUES(?,?,?,?)", (item_id, "documento", classificazione.get("motivo", ""), json.dumps(meta)) ) conn.commit() conn.close() # --- LanceDB --- def save_chunks_lancedb(doc_id: str, nome_file: str, chunks: list[str], vettori: list[list[float]], classificazione: dict) -> None: db = lancedb.connect(str(LANCEDB)) try: tbl = db.open_table("documenti") except Exception: # Crea la tabella al primo documento import pyarrow as pa schema = pa.schema([ pa.field("doc_id", pa.utf8()), pa.field("chunk_id", pa.utf8()), pa.field("testo", pa.utf8()), pa.field("fonte", pa.utf8()), pa.field("ambito", pa.utf8()), pa.field("categoria", pa.utf8()), pa.field("soggetto", pa.utf8()), pa.field("data_doc", pa.utf8()), pa.field("tipo", pa.utf8()), pa.field("vettore", pa.list_(pa.float32(), 384)), ]) tbl = db.create_table("documenti", schema=schema) anno = (classificazione.get("data_doc") or "")[:4] or str(datetime.utcnow().year) records = [] for i, (chunk, vec) in enumerate(zip(chunks, vettori)): records.append({ "doc_id": doc_id, "chunk_id": f"{doc_id}-{i:04d}", "testo": chunk, "fonte": "inbox", "ambito": classificazione["ambito"], "categoria": classificazione["categoria"], "soggetto": classificazione.get("soggetto") or "", "data_doc": classificazione.get("data_doc") or "", "tipo": Path(nome_file).suffix.lower().lstrip("."), "vettore": vec, }) if records: tbl.add(records) # --- Archiviazione --- def archive_file(path: Path, classificazione: dict) -> Path: anno = (classificazione.get("data_doc") or "")[:4] or str(datetime.utcnow().year) dest_dir = DOCS / classificazione["ambito"] / classificazione["categoria"] / anno dest_dir.mkdir(parents=True, exist_ok=True) dest = dest_dir / path.name # Evita collisioni di nome if dest.exists(): stem = path.stem suffix = path.suffix dest = dest_dir / f"{stem}_{datetime.utcnow().strftime('%H%M%S')}{suffix}" shutil.move(str(path), str(dest)) return dest # --- Telegram --- def telegram_notify(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, "parse_mode": "HTML"}, timeout=10, ) except Exception: pass # --- Limbo node --- def _create_limbo_node(doc_id: str, nome_file: str, classificazione: dict, n_chunks: int) -> None: item_id = f"limbo-{doc_id}" meta = json.dumps({ "tipo": "limbo", "doc_id": doc_id, "nome_file": nome_file, "ambito": classificazione["ambito"], "categoria": classificazione["categoria"], "data_doc": classificazione.get("data_doc"), "n_chunks": n_chunks, "limbo_stato": "pending", }) body = ( f"{nome_file} — {classificazione['ambito']}/{classificazione['categoria']}" + (f" ({classificazione['data_doc']})" if classificazione.get("data_doc") else "") + f" — {n_chunks} chunk in LanceDB. Da collegare ai nodi esistenti." ) conn = sqlite3.connect(VAULT_DB) conn.execute( "INSERT OR IGNORE INTO items(id, type, body, meta) VALUES(?,?,?,?)", (item_id, "nodo", body, meta) ) conn.commit() conn.close() # --- Entry point --- def ingest_file(path: Path) -> bool: """ Processa un file da inbox/. Ritorna True se ingerito, False se saltato. """ print(f"[ingest] {path.name}") doc_id = compute_doc_id(path) if is_already_ingested(doc_id): print(f"[ingest] già ingerito, skip: {path.name}") return False # Estrazione testo try: text = extract_text(path) except Exception as e: print(f"[ingest] errore estrazione {path.name}: {e}") telegram_notify(f"⚠️ Ingestione fallita: {path.name}\n{e}") return False vision_classificazione = None if not text.strip() and path.suffix.lower() == ".pdf": # PDF scansionato: Vision fa classificazione + estrazione testo in un colpo text, vision_classificazione = _vision_pdf(path) if not text.strip(): msg = f"⚠️ Ingestione: {path.name} — Vision non ha estratto testo (file corrotto?)\nFile rimasto in inbox/." print(f"[ingest] Vision: testo vuoto, skip: {path.name}") telegram_notify(msg) return False # Classificazione (salta se Vision l'ha già fatta) if vision_classificazione: classificazione = vision_classificazione else: classificazione = classify_document(text) print(f"[ingest] classificato: {classificazione['ambito']}/{classificazione['categoria']} " f"(confident={classificazione['confident']})") # Notifica Telegram se non sicuro if not classificazione["confident"]: telegram_notify( f"🗂 Nuovo documento: {path.name}\n" f"Proposta: {classificazione['ambito']}/{classificazione['categoria']}\n" f"Motivo incertezza: {classificazione.get('motivo', '—')}\n" f"(ingerito comunque — correggi i tag in vault.db se serve)" ) # Chunking + embedding chunks = chunk_text(text) vettori = embed_chunks(chunks) print(f"[ingest] {len(chunks)} chunk, embedding done") # Archiviazione file dest = archive_file(path, classificazione) print(f"[ingest] archiviato in: {dest}") # Pulizia directory vuote in inbox/ (non la root) parent = path.parent if parent != INBOX and parent.is_dir() and not any(parent.iterdir()): parent.rmdir() print(f"[ingest] rimossa directory vuota: {parent.name}") # Salvataggio vault.db + LanceDB save_documento(doc_id, path.name, classificazione, len(chunks), str(dest)) save_chunks_lancedb(doc_id, path.name, chunks, vettori, classificazione) _create_limbo_node(doc_id, path.name, classificazione, len(chunks)) telegram_notify( f"✅ {path.name} ingerito\n" f"{classificazione['ambito']}/{classificazione['categoria']} · {len(chunks)} chunk" ) print(f"[ingest] completato: {path.name}") return True def ingest_file_with_classify(path: Path, doc_id: str, classificazione: dict) -> bool: """Completa l'ingestione con classificazione già fornita da Adrian.""" print(f"[ingest] {path.name} (classificazione da Adrian)") try: text = extract_text(path) except Exception as e: telegram_notify(f"⚠️ Ingestione fallita: {path.name}\n{e}") return False if not text.strip() and path.suffix.lower() == ".pdf": # PDF scansionato: usa testo_estratto dalla classificazione Adrian se disponibile text = classificazione.pop("testo_estratto", "") or "" if not text.strip(): telegram_notify(f"⚠️ {path.name} — testo vuoto anche dopo classificazione Adrian.") return False print(f"[ingest] classificato: {classificazione['ambito']}/{classificazione['categoria']} " f"(confident={classificazione['confident']})") if not classificazione["confident"]: telegram_notify( f"🗂 Nuovo documento: {path.name}\n" f"Proposta: {classificazione['ambito']}/{classificazione['categoria']}\n" f"Motivo: {classificazione.get('motivo', '—')}\n" f"(ingerito comunque)" ) chunks = chunk_text(text) vettori = embed_chunks(chunks) dest = archive_file(path, classificazione) parent = path.parent if parent != INBOX and parent.is_dir() and not any(parent.iterdir()): parent.rmdir() save_documento(doc_id, path.name, classificazione, len(chunks), str(dest)) save_chunks_lancedb(doc_id, path.name, chunks, vettori, classificazione) _create_limbo_node(doc_id, path.name, classificazione, len(chunks)) telegram_notify( f"✅ {path.name} ingerito\n" f"{classificazione['ambito']}/{classificazione['categoria']} · {len(chunks)} chunk" ) print(f"[ingest] completato: {path.name}") return True if __name__ == "__main__": import sys if len(sys.argv) > 1: ingest_file(Path(sys.argv[1])) else: print("Uso: python ingest.py ")