Files
avatar/plugins/documenti_rag.py

260 lines
12 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Documenti interrogabili: indicizza le cartelle scelte (PDF, Word, testo, pagine…) in pezzi da ~800 caratteri,
con ricerca ibrida per parole (SQLite FTS5) e per significato (embeddings multilingual-e5-small). Tutto locale.
L'indice si aggiorna da solo ogni ora, solo per i file nuovi o modificati."""
from __future__ import annotations
import os
import re
import sqlite3
import sys
import threading
import time
from pathlib import Path
import numpy as np
sys.path.insert(0, str(Path(__file__).resolve().parent))
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from avatar.settings import DATA_DIR, Settings # noqa: E402
DB = DATA_DIR / "documenti.sqlite"
EMB_MODEL = "intfloat/multilingual-e5-small"
EXT = {".pdf", ".docx", ".doc", ".txt", ".md", ".rtf", ".odt", ".pages", ".html", ".csv"}
SKIP_DIRS = {"node_modules", ".git", "Library", ".venv", "venv", "__pycache__", "build", "dist", ".Trash"}
MAX_MB = 30
_emb: dict = {}
_stato = {"in_corso": False, "ultimo": "", "errore": ""}
def _cartelle() -> list[Path]:
raw = str(Settings().get("documenti_cartelle") or "~/Documents, ~/Desktop")
return [Path(c.strip()).expanduser() for c in raw.split(",") if c.strip()]
def _con() -> sqlite3.Connection:
DB.parent.mkdir(parents=True, exist_ok=True)
con = sqlite3.connect(DB, timeout=30)
con.executescript("""
CREATE TABLE IF NOT EXISTS files (path TEXT PRIMARY KEY, mtime REAL, chunks INTEGER);
CREATE TABLE IF NOT EXISTS chunks (id INTEGER PRIMARY KEY, path TEXT, n INTEGER, text TEXT, emb BLOB);
CREATE INDEX IF NOT EXISTS chunks_path ON chunks(path);
CREATE VIRTUAL TABLE IF NOT EXISTS fts USING fts5(text, content='chunks', content_rowid='id', tokenize='unicode61 remove_diacritics 2');
CREATE TRIGGER IF NOT EXISTS chunks_ai AFTER INSERT ON chunks BEGIN INSERT INTO fts(rowid, text) VALUES (new.id, new.text); END;
CREATE TRIGGER IF NOT EXISTS chunks_ad AFTER DELETE ON chunks BEGIN INSERT INTO fts(fts, rowid, text) VALUES ('delete', old.id, old.text); END;
""")
return con
def _embed(texts: list[str], query: bool = False) -> np.ndarray:
"""Vettori normalizzati (384 dimensioni); prefissi 'query:'/'passage:' richiesti dal modello e5."""
import torch
from transformers import AutoModel, AutoTokenizer
if not _emb:
_emb["tok"] = AutoTokenizer.from_pretrained(EMB_MODEL)
dev = "mps" if torch.backends.mps.is_available() else "cpu" # GPU del Mac: ~10× più veloce della CPU
_emb["model"] = AutoModel.from_pretrained(EMB_MODEL).eval().to(dev)
_emb["dev"] = dev
tok, model = _emb["tok"], _emb["model"]
out = []
for i in range(0, len(texts), 32):
batch = [("query: " if query else "passage: ") + t for t in texts[i:i + 32]]
enc = tok(batch, padding=True, truncation=True, max_length=512, return_tensors="pt").to(_emb["dev"])
with torch.no_grad():
h = model(**enc).last_hidden_state
m = enc["attention_mask"].unsqueeze(-1).float()
v = (h * m).sum(1) / m.sum(1)
out.append(torch.nn.functional.normalize(v, dim=1).float().cpu().numpy().astype(np.float32))
return np.concatenate(out) if out else np.zeros((0, 384), np.float32)
def _testo(path: Path) -> str:
import file as fileplugin # stesso estrattore del plugin file (pdf, docx, testo, textutil per gli altri)
return fileplugin._read(path)
def _pezzi(testo: str, size: int = 800, overlap: int = 150) -> list[str]:
testo = re.sub(r"[ \t]+", " ", re.sub(r"\n{3,}", "\n\n", testo)).strip()
out, i = [], 0
while i < len(testo):
out.append(testo[i:i + size]); i += size - overlap
return [p for p in out if len(p.strip()) > 40]
def _escludi() -> set[str]:
raw = str(Settings().get("documenti_escludi") or "Progetti2026")
return {x.strip() for x in raw.split(",") if x.strip()}
def _files():
skip = SKIP_DIRS | _escludi()
for root in _cartelle():
if not root.exists():
continue
for dirpath, dirnames, filenames in os.walk(root):
dirnames[:] = [d for d in dirnames if d not in skip and not d.startswith(".") and not d.endswith(".app")]
for f in filenames:
p = Path(dirpath) / f
if p.suffix.lower() in EXT and not f.startswith("~$"):
try:
if p.stat().st_size <= MAX_MB * 1e6:
yield p
except OSError:
pass
def aggiorna(log=None) -> str:
"""Indicizzazione incrementale: aggiunge i file nuovi/modificati, toglie quelli spariti."""
import fcntl
DB.parent.mkdir(parents=True, exist_ok=True)
lockf = open(DATA_DIR / "documenti.lock", "w")
try: # blocco tra processi: app e indicizzazioni manuali non si sovrappongono
fcntl.flock(lockf, fcntl.LOCK_EX | fcntl.LOCK_NB)
except OSError:
lockf.close()
return "Indicizzazione già in corso."
_stato["in_corso"] = True
try:
con = _con()
noti = dict(con.execute("SELECT path, mtime FROM files").fetchall())
visti, nuovi, pezzi_tot = set(), 0, 0
for p in _files():
sp = str(p); visti.add(sp)
mt = p.stat().st_mtime
if noti.get(sp) == mt:
continue
try:
pezzi = _pezzi(_testo(p))[:400]
except Exception:
pezzi = []
vec = _embed(pezzi) if pezzi else np.zeros((0, 384), np.float32)
con.execute("DELETE FROM chunks WHERE path=?", (sp,))
con.executemany("INSERT INTO chunks(path, n, text, emb) VALUES (?,?,?,?)", [(sp, i, t, vec[i].tobytes()) for i, t in enumerate(pezzi)])
con.execute("INSERT OR REPLACE INTO files(path, mtime, chunks) VALUES (?,?,?)", (sp, mt, len(pezzi)))
con.commit()
nuovi += 1; pezzi_tot += len(pezzi)
if log and nuovi % 25 == 0:
log(f"[documenti] indicizzati {nuovi} file…")
spariti = [p for p in noti if p not in visti]
for sp in spariti:
con.execute("DELETE FROM chunks WHERE path=?", (sp,)); con.execute("DELETE FROM files WHERE path=?", (sp,))
con.commit()
n_files, n_chunks = con.execute("SELECT COUNT(*), COALESCE(SUM(chunks),0) FROM files").fetchone()
con.close()
_stato["ultimo"] = time.strftime("%d/%m %H:%M")
return f"Indice aggiornato: {nuovi} file nuovi o modificati, {len(spariti)} rimossi. Totale {n_files} documenti, {n_chunks} passaggi."
finally:
_stato["in_corso"] = False
fcntl.flock(lockf, fcntl.LOCK_UN); lockf.close()
def cerca_testo(domanda: str, k: int = 6) -> list[dict]:
con = _con()
rows = con.execute("SELECT id, path, text, emb FROM chunks").fetchall()
if not rows:
con.close(); return []
ids = np.array([r[0] for r in rows]); M = np.frombuffer(b"".join(r[3] for r in rows), dtype=np.float32).reshape(len(rows), -1)
q = _embed([domanda], query=True)[0]
sem = M @ q # coseno (vettori normalizzati)
score = {int(i): float(s) for i, s in zip(ids, sem)}
words = [w for w in re.findall(r"\w{3,}", domanda.lower())]
if words:
try:
for (rid,) in con.execute("SELECT rowid FROM fts WHERE fts MATCH ? ORDER BY rank LIMIT 30", (" OR ".join(w + "*" for w in words),)):
score[rid] = score.get(rid, 0) + 0.15 # bonus per le parole esatte
except sqlite3.OperationalError:
pass
best = sorted(score.items(), key=lambda kv: -kv[1])[:k * 3]
by_id = {r[0]: r for r in rows}
out, per_file = [], {}
for rid, sc in best:
_i, path, text, _e = by_id[rid]
if per_file.get(path, 0) >= 2:
continue
per_file[path] = per_file.get(path, 0) + 1
out.append({"path": path, "testo": text, "score": round(sc, 3)})
if len(out) >= k:
break
con.close()
return out
# ── strumenti ────────────────────────────────────────────────────────────────
def t_cerca(params: dict, ctx: dict) -> str:
domanda = str(params.get("domanda") or "").strip()
if not domanda:
return "Errore: serve la domanda."
if not DB.exists():
return "Non ho ancora indicizzato i documenti: usa documenti_aggiorna (ci vuole qualche minuto la prima volta)."
res = cerca_testo(domanda, max(3, min(int(params.get("quanti") or 6), 12)))
if not res:
return "Nessun passaggio pertinente nei documenti indicizzati."
home = str(Path.home())
righe = [f"[{i + 1}] {r['path'].replace(home, '~')} (pertinenza {r['score']}):\n{r['testo']}" for i, r in enumerate(res)]
return ("Passaggi più pertinenti dai documenti dell'utente. Rispondi basandoti su questi, cita il nome del file, "
"e di' chiaramente se l'informazione non c'è:\n\n" + "\n\n".join(righe))
def t_aggiorna(params: dict, ctx: dict) -> str:
p = ctx.get("player")
if p is not None and hasattr(p, "set_state"):
p.set_state("PROCESSING · indicizzo i documenti")
return aggiorna(ctx.get("log"))
def t_stato(params: dict, ctx: dict) -> str:
if not DB.exists():
return "Indice non ancora creato. Cartelle: " + ", ".join(str(c) for c in _cartelle())
con = _con(); n_files, n_chunks = con.execute("SELECT COUNT(*), COALESCE(SUM(chunks),0) FROM files").fetchone(); con.close()
return (f"Documenti indicizzati: {n_files} file, {n_chunks} passaggi. Cartelle: {', '.join(str(c) for c in _cartelle())}. "
+ ("Indicizzazione in corso." if _stato["in_corso"] else f"Ultimo aggiornamento: {_stato['ultimo'] or 'all\'avvio'}."))
_proc: dict = {"p": None}
def pausa(attiva: bool) -> None:
"""Sospende/riprende l'indicizzazione mentre LuZa risponde: niente contesa di GPU e CPU con il modello."""
import signal
p = _proc["p"]
if p is not None and p.poll() is None:
try:
p.send_signal(signal.SIGSTOP if attiva else signal.SIGCONT)
except Exception:
pass
def _servizio() -> None:
"""Ogni ora indicizza in un processo separato a priorità bassa (niente blocchi del GIL nell'app)."""
import subprocess
time.sleep(120)
while True:
if Settings().get("documenti_enabled", True):
code = ("import sys; sys.argv=['memory_mcp.py']; sys.path.insert(0, %r); import documenti_rag as d; print('[documenti]', d.aggiorna(), flush=True)"
% str(Path(__file__).resolve().parent))
_proc["p"] = subprocess.Popen(["nice", "-n", "15", sys.executable, "-c", code], cwd=str(Path(__file__).resolve().parent.parent))
_proc["p"].wait()
_proc["p"] = None
time.sleep(3600)
if Path(sys.argv[0]).name != "memory_mcp.py": # solo nell'app
threading.Thread(target=_servizio, daemon=True, name="documenti").start()
TOOLS = [
{"name": "documenti_cerca", "description": "Cerca nel CONTENUTO dei documenti dell'utente (PDF, Word, testi nelle cartelle indicizzate: Documenti, Scrivania) per rispondere a domande come 'cosa dice il contratto sulle penali?', 'quando scade l'assicurazione?', 'trova la ricevuta del dentista'. Restituisce i passaggi più pertinenti con il file di provenienza.",
"parameters": {"type": "object", "properties": {"domanda": {"type": "string"}, "quanti": {"type": "integer"}}, "required": ["domanda"]}, "run": t_cerca},
{"name": "documenti_aggiorna", "description": "Aggiorna subito l'indice dei documenti (file nuovi o modificati). Normalmente avviene da solo ogni ora.",
"parameters": {"type": "object", "properties": {}}, "run": t_aggiorna},
{"name": "documenti_stato", "description": "Quanti documenti sono indicizzati, quali cartelle e quando è stato l'ultimo aggiornamento.",
"parameters": {"type": "object", "properties": {}}, "run": t_stato},
]
if __name__ == "__main__":
assert _pezzi("a" * 2000)[1].startswith("a") and len(_pezzi("a" * 2000)) == 4
v = _embed(["Il contratto prevede una penale del 5% per ritardo."]); q = _embed(["quanto è la penale?"], query=True)[0]
w = _embed(["La ricetta della carbonara usa guanciale."])
assert float(v[0] @ q) > float(w[0] @ q), "il significato deve battere il testo non pertinente"
print("ok: pertinente", round(float(v[0] @ q), 3), "> non pertinente", round(float(w[0] @ q), 3))