Last active 1 month ago

Pipeline incremental de indexacion RAG: escanea, chunk, embedding y FTS5 (episodio 819)

Revision 767525517468689e7497a870d06e3084720c0808

rag_pipeline.py Raw
1"#!/usr/bin/env python3\n\"\"\"\nrag_pipeline.py — Pipeline incremental para indexar documentos markdown.\n\nFlujo completo:\n 1. Recorre /data/notas buscando archivos *.md\n 2. Calcula MD5 de cada archivo para detectar cambios\n 3. Compara con la tabla 'documentos' → nuevos, modificados, eliminados\n 4. Para cada archivo afectado: lee → extrae frontmatter → chunk → embedding → inserta\n 5. Reconstruye el índice FTS5 al final\n\"\"\"\n\nfrom __future__ import annotations\n\nimport hashlib\nimport json\nimport os\nimport sqlite3\nimport sys\nfrom typing import Any\n\nfrom rag_chunker import chunk_markdown, strip_frontmatter\nfrom rag_config import (\n BATCH_SIZE,\n DB_PATH,\n EXCLUDE_DIRS,\n EXTENSIONES,\n MAX_TOKENS,\n MODELO_EMBEDDINGS,\n OVERLAP,\n ROOT_DIR,\n)\nfrom rag_embeddings import get_embeddings_batch, pack_vector\nfrom rag_schema import crear_esquema, limpiar_chunks_por_doc_id, rebuild_fts5\n\n\n# ── Escaneo de archivos ──────────────────────────────────────────────────────\n\n\ndef walk_markdowns(root_dir: str) -> list[dict[str, Any]]:\n \"\"\"Recorre un directorio buscando archivos markdown de forma recursiva.\n\n Excluye carpetas como .git, node_modules, __pycache__, etc.\n\n Args:\n root_dir: Directorio raíz desde el que empezar la búsqueda.\n\n Returns:\n Lista de dicts con 'path' (ruta absoluta) y 'relpath' (relativa).\n \"\"\"\n archivos: list[dict[str, Any]] = []\n root_dir = os.path.abspath(root_dir)\n\n for root, dirs, files in os.walk(root_dir):\n # Excluir carpetas no deseadas (modificando dirs in-place)\n dirs[:] = [d for d in dirs if d not in EXCLUDE_DIRS]\n\n for f in files:\n if f.endswith(EXTENSIONES):\n ruta_abs = os.path.join(root, f)\n archivos.append(\n {\n \"path\": ruta_abs,\n \"relpath\": os.path.relpath(ruta_abs, root_dir),\n }\n )\n\n return archivos\n\n\n# ── Cálculo de hash ──────────────────────────────────────────────────────────\n\n\ndef compute_md5(path: str) -> str:\n \"\"\"Calcula el hash MD5 de un archivo.\n\n Lee el archivo en trozos de 64KB para no cargar ficheros grandes\n enteros en memoria.\n\n Args:\n path: Ruta absoluta al archivo.\n\n Returns:\n Hash MD5 hexadecimal en minúsculas.\n \"\"\"\n h = hashlib.md5()\n with open(path, \"rb\") as f:\n for chunk in iter(lambda: f.read(65536), b\"\"):\n h.update(chunk)\n return h.hexdigest()\n\n\n# ── Estado de documentos ─────────────────────────────────────────────────────\n\n\ndef get_doc_status(\n db: sqlite3.Connection, files: list[dict[str, Any]]\n) -> dict[str, list[str]]:\n \"\"\"Compara los archivos actuales con la base de datos para detectar cambios.\n\n Args:\n db: Conexión SQLite abierta.\n files: Lista de archivos detectados en el sistema de ficheros.\n\n Returns:\n Diccionario con tres claves:\n - 'nuevos': archivos que no están en la BD\n - 'modificados': archivos cuyo MD5 ha cambiado\n - 'eliminados': rutas que están en la BD pero ya no existen\n \"\"\"\n # Cargar estado conocido de la BD\n cur = db.execute(\"SELECT path, md5_hash FROM documentos\")\n conocidos: dict[str, str] = {row[0]: row[1] for row in cur.fetchall()}\n\n # Construir mapa de archivos actuales\n actual: dict[str, str] = {}\n for f in files:\n md5 = compute_md5(f[\"path\"])\n actual[f[\"path\"]] = md5\n\n # Clasificar\n nuevos: list[str] = []\n modificados: list[str] = []\n\n for ruta, md5 in actual.items():\n if ruta not in conocidos:\n nuevos.append(ruta)\n elif conocidos[ruta] != md5:\n modificados.append(ruta)\n\n eliminados = list(set(conocidos.keys()) - set(actual.keys()))\n\n return {\"nuevos\": nuevos, \"modificados\": modificados, \"eliminados\": eliminados}\n\n\n# ── Reindexar un archivo ─────────────────────────────────────────────────────\n\n\ndef reindex_file(db: sqlite3.Connection, file_path: str) -> int:\n \"\"\"Procesa un archivo markdown y lo inserta/actualiza en la base de datos.\n\n Pasos:\n 1. Lee el archivo\n 2. Extrae frontmatter\n 3. Chunking del contenido limpio\n 4. Inserta/actualiza en 'documentos'\n 5. Inserta chunks (los triggers sincronizan FTS5)\n 6. Genera embeddings y los guarda\n\n Args:\n db: Conexión SQLite.\n file_path: Ruta absoluta al archivo.\n\n Returns:\n Número de chunks generados para este archivo.\n \"\"\"\n md5 = compute_md5(file_path)\n relpath = os.path.relpath(file_path, ROOT_DIR)\n\n with open(file_path, encoding=\"utf-8\", errors=\"replace\") as f:\n texto = f.read()\n\n # Extraer frontmatter y contenido limpio\n metadata, contenido = strip_frontmatter(texto)\n\n if not contenido.strip():\n # Archivo vacío o solo frontmatter → lo saltamos\n return 0\n\n # Chunking\n chunks = chunk_markdown(contenido, max_tokens=MAX_TOKENS, overlap=OVERLAP)\n if not chunks:\n return 0\n\n # Extraer metadatos útiles\n titulo: str | None = metadata.get(\"title\")\n tags_raw = metadata.get(\"tags\", [])\n if isinstance(tags_raw, list):\n tags_json = json.dumps(tags_raw)\n else:\n tags_json = json.dumps([])\n\n # Verificar si el documento ya existe\n cur = db.execute(\"SELECT id FROM documentos WHERE path = ?\", (file_path,))\n row = cur.fetchone()\n\n if row:\n doc_id = row[0]\n # Eliminar chunks y embeddings viejos\n limpiar_chunks_por_doc_id(db, doc_id)\n # Actualizar metadatos\n db.execute(\n \"\"\"UPDATE documentos\n SET md5_hash = ?, last_indexed = datetime('now')\n WHERE id = ?\"\"\",\n (md5, doc_id),\n )\n else:\n cur = db.execute(\n \"INSERT INTO documentos (path, md5_hash) VALUES (?, ?)\",\n (file_path, md5),\n )\n doc_id = cur.lastrowid # type: ignore[assignment]\n\n # Insertar chunks\n for idx, chunk in enumerate(chunks):\n db.execute(\n \"\"\"INSERT INTO chunks\n (doc_id, chunk_index, content, tokens, doc_path, title, tags)\n VALUES (?, ?, ?, ?, ?, ?, ?)\"\"\",\n (\n doc_id,\n idx,\n chunk[\"content\"],\n chunk[\"tokens\"],\n relpath,\n chunk.get(\"title\") or titulo,\n tags_json,\n ),\n )\n\n # Confirmar antes de generar embeddings\n db.commit()\n\n # Generar y guardar embeddings para cada chunk\n textos = [c[\"content\"] for c in chunks]\n try:\n embeddings = get_embeddings_batch(\n textos, batch_size=BATCH_SIZE, modelo=MODELO_EMBEDDINGS\n )\n except Exception as e:\n print(f\" ⚠️ Error generando embeddings para {relpath}: {e}\")\n print(f\" Los chunks están guardados pero sin vectores.\")\n return len(chunks)\n\n # Recuperar los IDs de los chunks recién insertados\n cur = db.execute(\n \"SELECT id FROM chunks WHERE doc_id = ? ORDER BY chunk_index\",\n (doc_id,),\n )\n chunk_ids = [r[0] for r in cur.fetchall()]\n\n for chunk_id, vec in zip(chunk_ids, embeddings):\n blob = pack_vector(vec)\n db.execute(\n \"INSERT OR REPLACE INTO embeddings (chunk_id, vector, model, dimensions) VALUES (?, ?, ?, ?)\",\n (chunk_id, blob, MODELO_EMBEDDINGS, len(vec)),\n )\n\n db.commit()\n return len(chunks)\n\n\n# ── Limpiar documentos eliminados ────────────────────────────────────────────\n\n\ndef clean_deleted(db: sqlite3.Connection, active_paths: set[str]) -> int:\n \"\"\"Elimina de la BD los documentos que ya no existen en disco.\n\n Args:\n db: Conexión SQLite.\n active_paths: Conjunto de rutas activas actuales.\n\n Returns:\n Número de documentos eliminados.\n \"\"\"\n cur = db.execute(\"SELECT id, path FROM documentos\")\n eliminados = 0\n for row in cur.fetchall():\n doc_id, path = row\n if path not in active_paths:\n # ON DELETE CASCADE se encarga de chunks y embeddings\n db.execute(\"DELETE FROM documentos WHERE id = ?\", (doc_id,))\n eliminados += 1\n\n db.commit()\n return eliminados\n\n\n# ── Pipeline completo ────────────────────────────────────────────────────────\n\n\ndef run_pipeline(db_path: str = DB_PATH, root_dir: str = ROOT_DIR) -> dict[str, int]:\n \"\"\"Ejecuta el pipeline completo de indexación incremental.\n\n Args:\n db_path: Ruta al archivo de base de datos SQLite.\n root_dir: Directorio raíz con los documentos markdown.\n\n Returns:\n Dict con estadísticas del proceso.\n \"\"\"\n stats: dict[str, int] = {\n \"total_archivos\": 0,\n \"nuevos\": 0,\n \"modificados\": 0,\n \"eliminados\": 0,\n \"chunks_generados\": 0,\n \"errores\": 0,\n }\n\n print(f\"📂 Escaneando {root_dir}...\")\n files = walk_markdowns(root_dir)\n stats[\"total_archivos\"] = len(files)\n print(f\" Encontrados {len(files)} archivos markdown\\n\")\n\n print(\"🔌 Conectando a base de datos...\")\n conn = sqlite3.connect(db_path)\n try:\n # Asegurar esquema\n crear_esquema(conn)\n\n # Detectar cambios\n status = get_doc_status(conn, files)\n stats[\"nuevos\"] = len(status[\"nuevos\"])\n stats[\"modificados\"] = len(status[\"modificados\"])\n stats[\"eliminados\"] = len(status[\"eliminados\"])\n\n print(f\"\\n📊 Estado:\")\n print(f\" Nuevos: {len(status['nuevos'])}\")\n print(f\" Modificados: {len(status['modificados'])}\")\n print(f\" Eliminados: {len(status['eliminados'])}\")\n print(\n f\" Sin cambios: {len(files) - len(status['nuevos']) - len(status['modificados'])}\\n\"\n )\n\n # Procesar eliminados\n active_paths = {f[\"path\"] for f in files}\n eliminados = clean_deleted(conn, active_paths)\n\n # Procesar nuevos y modificados\n afectados = status[\"nuevos\"] + status[\"modificados\"]\n if afectados:\n print(f\"🔄 Procesando {afectados} archivos...\\n\")\n for i, ruta in enumerate(status[\"nuevos\"] + status[\"modificados\"], 1):\n rel = os.path.relpath(ruta, root_dir)\n print(f\" [{i}/{afectados}] {rel}\")\n try:\n n_chunks = reindex_file(conn, ruta)\n stats[\"chunks_generados\"] += n_chunks\n print(f\"{n_chunks} chunks\")\n except Exception as e:\n print(f\" ⚠️ ERROR: {e}\")\n stats[\"errores\"] += 1\n conn.rollback()\n\n # Reconstruir índice FTS5\n print(\"\\n🔨 Reconstruyendo índice FTS5...\")\n rebuild_fts5(conn)\n print(\" ✔ Índice FTS5 actualizado\")\n\n print(\"\\n✅ Pipeline completado.\")\n print(f\" Archivos procesados: {afectados}\")\n print(f\" Chunks generados: {stats['chunks_generados']}\")\n print(f\" Errores: {stats['errores']}\")\n print(f\" Base de datos: {db_path}\")\n\n finally:\n conn.close()\n\n return stats\n\n\n# ── CLI ──────────────────────────────────────────────────────────────────────\n\n\ndef main() -> None:\n \"\"\"Punto de entrada: ejecuta el pipeline.\"\"\"\n import argparse\n\n parser = argparse.ArgumentParser(\n description=\"Pipeline RAG: indexa documentos markdown\"\n )\n parser.add_argument(\n \"--db\",\n default=DB_PATH,\n help=f\"Ruta a la base de datos (default: {DB_PATH})\",\n )\n parser.add_argument(\n \"--root\",\n default=ROOT_DIR,\n help=f\"Directorio raíz de documentos (default: {ROOT_DIR})\",\n )\n args = parser.parse_args()\n\n run_pipeline(db_path=args.db, root_dir=args.root)\n\n\nif __name__ == \"__main__\":\n main()\n"