atareao / rag_consulta.py

0 likes
0 forks
1 files
Last active 1 month ago
Cliente de consulta FTS5 para la base de conocimiento RAG con output coloreado (episodio 819)
1 "#!/usr/bin/env python3\n\"\"\"\nrag_consulta.py — Consulta FTS5 sobre la base de conocimiento RAG.\n\nBusca en el índice FTS5 y devuelve los fragmentos más relevantes\ncon resaltado de términos, información del documento y score BM25.\n\"\"\"\n\nfrom __future__ import annotations\n\nimport argparse\nimport json\nimport sqlite3\nimport sys\nfrom typing import Any\n\nfrom rag_config import DB_PATH\n\n\n# ── Colores ANSI ─────────────────────────────────────────────────────────────\n\n\nclass ANSI:\n \"\"\"Códigos de color ANSI para la salida por terminal.\"\"\"\n\n RESET = \"\\033[0m\"\n BOLD = \"\\033[1m\"\n DIM = \"\\033[2m\"\n CYAN = \"\\033[36m\"\n GREEN = \"\\033[32m\"\n YELLOW = \"\\033[33m\"\n RED = \"\\033[31m\"\n MAGENTA = \"\\033[35m\"\n BLUE = \"\\033[34m\"\n\n\n# ── Consulta FTS5 ────────────────────────────────────────────────────────────\n\n\ndef consulta_fts5(\n db: sqlite3.Connection,\n query: str,\n limit: int = 5,\n) -> list[dict[str, Any]]:\n \"\"\"Ejecuta una consulta FTS5 y devuelve los resultados ordenados por relevancia.\n\n Args:\n db: Conexión SQLite.\n query: Término(s) de búsqueda en sintaxis FTS5.\n limit: Número máximo de resultados.\n\n Returns:\n Lista de dicts con: id, doc_id, doc_path, title, snippet, rank, tokens.\n \"\"\"\n sql = \"\"\"\n SELECT\n c.id,\n c.doc_id,\n c.doc_path,\n c.title,\n snippet(chunks_fts, 0, '<<<', '>>>', '...', 50) AS fragmento,\n rank,\n c.tokens,\n c.tags,\n c.chunk_index\n FROM chunks_fts\n JOIN chunks c ON chunks_fts.rowid = c.id\n WHERE chunks_fts MATCH ?\n ORDER BY rank\n LIMIT ?\n \"\"\"\n cur = db.execute(sql, (query, limit))\n resultados: list[dict[str, Any]] = []\n for row in cur.fetchall():\n resultados.append(\n {\n \"id\": row[0],\n \"doc_id\": row[1],\n \"doc_path\": row[2],\n \"title\": row[3] or \"(sin título)\",\n \"snippet\": row[4] or \"(sin fragmento disponible)\",\n \"rank\": row[5],\n \"tokens\": row[6],\n \"tags\": row[7],\n \"chunk_index\": row[8],\n }\n )\n return resultados\n\n\n# ── Salida formateada ────────────────────────────────────────────────────────\n\n\ndef mostrar_resultados(resultados: list[dict[str, Any]], query: str) -> None:\n \"\"\"Muestra los resultados de la consulta con colores y formato bonito.\n\n Args:\n resultados: Lista de resultados de consulta_fts5.\n query: La consulta original (para mostrarla en el encabezado).\n \"\"\"\n if not resultados:\n print(\n f\"\\n{ANSI.YELLOW}😕 No se encontraron resultados para: {ANSI.BOLD}{query}{ANSI.RESET}\"\n )\n print(f\" Prueba con términos más generales o revisa la ortografía.\")\n print(\n f' Sintaxis FTS5: \"frase exacta\", termino AND otro, termino OR otro, term*'\n )\n return\n\n print(f\"\\n{ANSI.CYAN}{ANSI.BOLD}🔍 Resultados para: {query}{ANSI.RESET}\")\n print(f\"{ANSI.DIM} {len(resultados)} resultado(s) encontrado(s){ANSI.RESET}\\n\")\n\n for i, r in enumerate(resultados, 1):\n # Encabezado del resultado\n print(f\"{ANSI.CYAN}{ANSI.BOLD}─── Resultado #{i} ───{ANSI.RESET}\")\n\n # Título\n print(f\" {ANSI.GREEN}{ANSI.BOLD}📄 {r['title']}{ANSI.RESET}\")\n\n # Ruta\n print(f\" {ANSI.DIM}📍 {r['doc_path']} (chunk {r['chunk_index']}){ANSI.RESET}\")\n\n # Tags si existen\n if r[\"tags\"] and r[\"tags\"] != \"[]\" and r[\"tags\"] != \"null\":\n try:\n tags = json.loads(r[\"tags\"])\n if tags:\n print(f\" {ANSI.MAGENTA}🏷️ {', '.join(tags)}{ANSI.RESET}\")\n except (json.JSONDecodeError, TypeError):\n pass\n\n # Snippet con resaltado\n snippet = r[\"snippet\"]\n # Reemplazar marcadores por colores ANSI\n snippet_coloreado = snippet.replace(\"<<<\", f\"{ANSI.YELLOW}{ANSI.BOLD}\")\n snippet_coloreado = snippet_coloreado.replace(\">>>\", f\"{ANSI.RESET}\")\n print(f\" {ANSI.DIM}Fragmento:{ANSI.RESET}\")\n print(f\" {snippet_coloreado}\")\n\n # Score y tokens\n print(\n f\" {ANSI.BLUE}Score: {ANSI.BOLD}{r['rank']:.4f}{ANSI.RESET}\"\n f\" {ANSI.DIM}|{ANSI.RESET} \"\n f\"{ANSI.BLUE}Tokens: {r['tokens']}{ANSI.RESET}\"\n )\n print()\n\n\n# ── Estadísticas de la base de datos ─────────────────────────────────────────\n\n\ndef mostrar_estadisticas(db: sqlite3.Connection) -> None:\n \"\"\"Muestra estadísticas básicas de la base de datos.\"\"\"\n try:\n docs = db.execute(\"SELECT COUNT(*) FROM documentos\").fetchone()[0]\n chunks = db.execute(\"SELECT COUNT(*) FROM chunks\").fetchone()[0]\n embeds = db.execute(\"SELECT COUNT(*) FROM embeddings\").fetchone()[0]\n\n # Última indexación\n ultima = db.execute(\"SELECT MAX(last_indexed) FROM documentos\").fetchone()[0]\n\n print(\n f\"\\n{ANSI.CYAN}{ANSI.BOLD}📊 Estadísticas de la base de datos{ANSI.RESET}\"\n )\n print(f\" Documentos: {docs}\")\n print(f\" Chunks: {chunks}\")\n print(f\" Embeddings: {embeds}\")\n if ultima:\n print(f\" Última indexación: {ultima}\")\n print()\n except Exception:\n pass\n\n\n# ── CLI ──────────────────────────────────────────────────────────────────────\n\n\ndef main() -> None:\n \"\"\"Punto de entrada: procesa argumentos y ejecuta la consulta.\"\"\"\n parser = argparse.ArgumentParser(\n description=\"Consulta la base de conocimiento RAG con FTS5\",\n formatter_class=argparse.RawDescriptionHelpFormatter,\n epilog=\"\"\"\nEjemplos:\n %(prog)s \"cómo configurar docker\"\n %(prog)s \"systemd AND timers\"\n %(prog)s \"docker compose volumes\"\n %(prog)s \"inteligencia artificial\" -l 3\n %(prog)s --stats\n %(prog)s \"nginx AND proxy AND ssl\" -l 10\n \"\"\",\n )\n parser.add_argument(\n \"query\",\n nargs=\"?\",\n help=\"Término(s) de búsqueda en formato FTS5\",\n )\n parser.add_argument(\n \"-l\",\n \"--limit\",\n type=int,\n default=5,\n help=\"Número de resultados (default: 5)\",\n )\n parser.add_argument(\n \"--db\",\n help=f\"Ruta a la base de datos (default: {DB_PATH})\",\n )\n parser.add_argument(\n \"--stats\",\n action=\"store_true\",\n help=\"Muestra estadísticas de la base de datos\",\n )\n\n args = parser.parse_args()\n\n db_path = args.db or DB_PATH\n\n # Conectar a la base de datos\n try:\n conn = sqlite3.connect(db_path)\n except sqlite3.Error as e:\n print(f\"{ANSI.RED}Error conectando a la base de datos: {e}{ANSI.RESET}\")\n sys.exit(1)\n\n try:\n # Mostrar estadísticas si se pide\n if args.stats:\n mostrar_estadisticas(conn)\n return\n\n # Consulta\n if not args.query:\n parser.print_help()\n return\n\n resultados = consulta_fts5(conn, query=args.query, limit=args.limit)\n mostrar_resultados(resultados, query=args.query)\n\n finally:\n conn.close()\n\n\nif __name__ == \"__main__\":\n main()\n"

atareao / rag_config.py

0 likes
0 forks
1 files
Last active 1 month ago
Configuracion central para la base de conocimiento RAG con SQLite y Ollama (episodio 819)
1 "#!/usr/bin/env python3\n\"\"\"\nrag_config.py — Configuración central para la base de conocimiento RAG.\n\nTodas las constantes del pipeline en un solo sitio para no andar\ncambiando valores en seis ficheros distintos.\n\"\"\"\n\nfrom __future__ import annotations\nimport os\nfrom typing import Final\n\n# ── Rutas ────────────────────────────────────────────────────────────────────\n\n# Ruta donde se guarda la base de datos SQLite\nDB_PATH: Final[str] = os.path.join(\n os.path.dirname(os.path.abspath(__file__)),\n \"..\",\n \"rag_conocimiento.db\",\n)\n\n# Directorio raíz donde están los documentos markdown\nROOT_DIR: Final[str] = os.path.abspath(\n os.path.join(os.path.dirname(os.path.abspath(__file__)), \"..\", \"..\", \"..\", \"..\")\n)\n\n# ── Embeddings ───────────────────────────────────────────────────────────────\n\n# Modelo de embeddings en Ollama (bge-m3, BAAI, 568M params, 1024 dims)\nMODELO_EMBEDDINGS: Final[str] = \"bge-m3\"\n\n# Dimensión del vector de embedding (bge-m3 → 1024)\nDIMENSIONES_EMBEDDING: Final[int] = 1024\n\n# ── Chunking ─────────────────────────────────────────────────────────────────\n\n# Máximo de tokens por chunk\nMAX_TOKENS: Final[int] = 512\n\n# Solapamiento entre chunks consecutivos (25% de MAX_TOKENS)\nOVERLAP: Final[int] = 128\n\n# ── Pipeline ─────────────────────────────────────────────────────────────────\n\n# Tamaño del lote para enviar embeddings a Ollama\nBATCH_SIZE: Final[int] = 16\n\n# Carpetas que se excluyen al recorrer el árbol de documentos\nEXCLUDE_DIRS: Final[set[str]] = {\n \".git\",\n \"node_modules\",\n \"__pycache__\",\n \"venv\",\n \".venv\",\n \".opencode\",\n}\n\n# Extensiones de archivo a procesar\nEXTENSIONES: Final[tuple[str, ...]] = (\".md\",)\n"

atareao / rag_schema.py

0 likes
0 forks
1 files
Last active 1 month ago
Esquema SQLite completo: tablas, FTS5, triggers y embeddings para RAG (episodio 819)
1 "#!/usr/bin/env python3\n\"\"\"\nrag_schema.py — Esquema SQLite para la base de conocimiento RAG.\n\nCrea y gestiona las tablas:\n - documentos → metadatos de cada archivo markdown\n - chunks → fragmentos de texto indexados por documento\n - chunks_fts → índice FTS5 (contenido externo sobre chunks)\n - embeddings → vectores numéricos por chunk\n\"\"\"\n\nfrom __future__ import annotations\n\nimport sqlite3\nfrom typing import Final\n\n# ── SQL DDL ──────────────────────────────────────────────────────────────────\n\nSQL_DOCUMENTOS: Final[str] = \"\"\"\nCREATE TABLE IF NOT EXISTS documentos (\n id INTEGER PRIMARY KEY AUTOINCREMENT,\n path TEXT UNIQUE NOT NULL,\n md5_hash TEXT NOT NULL,\n last_indexed TEXT DEFAULT (datetime('now'))\n);\n\"\"\"\n\nSQL_CHUNKS: Final[str] = \"\"\"\nCREATE TABLE IF NOT EXISTS chunks (\n id INTEGER PRIMARY KEY AUTOINCREMENT,\n doc_id INTEGER NOT NULL,\n chunk_index INTEGER NOT NULL,\n content TEXT NOT NULL,\n tokens INTEGER,\n doc_path TEXT NOT NULL,\n title TEXT,\n tags TEXT,\n created_at TEXT DEFAULT (datetime('now')),\n FOREIGN KEY (doc_id) REFERENCES documentos(id) ON DELETE CASCADE\n);\n\"\"\"\n\nSQL_CHUNKS_FTS: Final[str] = \"\"\"\nCREATE VIRTUAL TABLE IF NOT EXISTS chunks_fts USING fts5(\n content,\n tokenize='porter unicode61 remove_diacritics 1',\n content='chunks',\n content_rowid='id'\n);\n\"\"\"\n\nSQL_EMBEDDINGS: Final[str] = \"\"\"\nCREATE TABLE IF NOT EXISTS embeddings (\n chunk_id INTEGER PRIMARY KEY,\n vector BLOB NOT NULL,\n model TEXT DEFAULT 'bge-m3',\n dimensions INTEGER DEFAULT 1024,\n FOREIGN KEY (chunk_id) REFERENCES chunks(id) ON DELETE CASCADE\n);\n\"\"\"\n\nSQL_INDEX_DOC_PATH: Final[str] = (\n \"CREATE INDEX IF NOT EXISTS idx_chunks_doc_path ON chunks(doc_path);\"\n)\nSQL_INDEX_DOC_ID: Final[str] = (\n \"CREATE INDEX IF NOT EXISTS idx_chunks_doc_id ON chunks(doc_id);\"\n)\n\nSQL_TRIGGER_FTS_INSERT: Final[str] = \"\"\"\nCREATE TRIGGER IF NOT EXISTS chunks_ai AFTER INSERT ON chunks\nBEGIN\n INSERT INTO chunks_fts(rowid, content) VALUES (new.id, new.content);\nEND;\n\"\"\"\n\nSQL_TRIGGER_FTS_DELETE: Final[str] = \"\"\"\nCREATE TRIGGER IF NOT EXISTS chunks_ad AFTER DELETE ON chunks\nBEGIN\n INSERT INTO chunks_fts(chunks_fts, rowid, content) VALUES('delete', old.id, old.content);\nEND;\n\"\"\"\n\nSQL_TRIGGER_FTS_UPDATE: Final[str] = \"\"\"\nCREATE TRIGGER IF NOT EXISTS chunks_au AFTER UPDATE ON chunks\nBEGIN\n INSERT INTO chunks_fts(chunks_fts, rowid, content) VALUES('delete', old.id, old.content);\n INSERT INTO chunks_fts(rowid, content) VALUES (new.id, new.content);\nEND;\n\"\"\"\n\nSQL_PRAGMAS: Final[str] = \"\"\"\nPRAGMA journal_mode = WAL;\nPRAGMA foreign_keys = ON;\nPRAGMA synchronous = NORMAL;\n\"\"\"\n\n\n# ── Funciones públicas ───────────────────────────────────────────────────────\n\n\ndef crear_esquema(conn: sqlite3.Connection) -> None:\n \"\"\"Crea todas las tablas, índices y triggers si no existen.\"\"\"\n conn.executescript(SQL_PRAGMAS)\n conn.execute(SQL_DOCUMENTOS)\n conn.execute(SQL_CHUNKS)\n conn.execute(SQL_CHUNKS_FTS)\n conn.execute(SQL_EMBEDDINGS)\n conn.execute(SQL_INDEX_DOC_PATH)\n conn.execute(SQL_INDEX_DOC_ID)\n conn.execute(SQL_TRIGGER_FTS_INSERT)\n conn.execute(SQL_TRIGGER_FTS_DELETE)\n conn.execute(SQL_TRIGGER_FTS_UPDATE)\n conn.commit()\n\n\ndef rebuild_fts5(conn: sqlite3.Connection) -> None:\n \"\"\"Reconstruye el índice FTS5 desde los datos actuales de la tabla chunks.\n\n Necesario después de inserciones/actualizaciones masivas porque\n la tabla es de contenido externo.\n \"\"\"\n conn.executescript(\"INSERT INTO chunks_fts(chunks_fts) VALUES('rebuild');\")\n\n\ndef limpiar_chunks_por_doc_id(conn: sqlite3.Connection, doc_id: int) -> None:\n \"\"\"Elimina todos los chunks y embeddings de un documento.\n\n Los triggers FTS5 se encargan de limpiar el índice textual.\n El CASCADE de embeddings se activa por la FK.\n \"\"\"\n conn.execute(\n \"DELETE FROM embeddings WHERE chunk_id IN (SELECT id FROM chunks WHERE doc_id = ?)\",\n (doc_id,),\n )\n conn.execute(\"DELETE FROM chunks WHERE doc_id = ?\", (doc_id,))\n conn.commit()\n\n\ndef main() -> None:\n \"\"\"Demo: crea el esquema en la base de datos configurada.\"\"\"\n from rag_config import DB_PATH\n\n print(f\"Creando esquema en: {DB_PATH}\")\n conn = sqlite3.connect(DB_PATH)\n try:\n crear_esquema(conn)\n print(\"✔ Esquema creado correctamente.\")\n\n # Listar las tablas para verificar\n cur = conn.execute(\n \"SELECT name FROM sqlite_master WHERE type='table' ORDER BY name;\"\n )\n tablas = [r[0] for r in cur.fetchall()]\n print(f\" Tablas: {', '.join(tablas)}\")\n\n cur = conn.execute(\n \"SELECT name FROM sqlite_master WHERE type='trigger' ORDER BY name;\"\n )\n triggers = [r[0] for r in cur.fetchall()]\n print(f\" Triggers: {', '.join(triggers)}\")\n finally:\n conn.close()\n\n\nif __name__ == \"__main__\":\n main()\n"

atareao / rag_embeddings.py

0 likes
0 forks
1 files
Last active 1 month ago
Generacion de embeddings con Ollama + bge-m3 para RAG local (episodio 819)
1 "#!/usr/bin/env python3\n\"\"\"\nrag_embeddings.py — Generación de embeddings con Ollama + bge-m3.\n\nEnvía texto al modelo de embeddings de Ollama y empaqueta los vectores\nresultantes como BLOBs binarios para almacenar en SQLite.\n\"\"\"\n\nfrom __future__ import annotations\n\nimport struct\nfrom typing import Any\n\nimport ollama\n\nfrom rag_config import MODELO_EMBEDDINGS, DIMENSIONES_EMBEDDING\n\n\n# ── Funciones de embedding ───────────────────────────────────────────────────\n\n\ndef get_embedding(texto: str, modelo: str | None = None) -> list[float]:\n \"\"\"Genera el embedding de un texto usando el modelo de Ollama.\n\n Args:\n texto: Texto a convertir en vector.\n modelo: Nombre del modelo en Ollama (por defecto bge-m3).\n\n Returns:\n Lista de floats con el vector de embedding (1024 dimensiones).\n \"\"\"\n resp: Any = ollama.embed(model=modelo or MODELO_EMBEDDINGS, input=texto)\n # La API devuelve una lista de embeddings; para un solo input\n # el primer (y único) elemento es el que nos interesa.\n return resp[\"embeddings\"][0]\n\n\ndef get_embeddings_batch(\n textos: list[str],\n batch_size: int = 16,\n modelo: str | None = None,\n) -> list[list[float]]:\n \"\"\"Genera embeddings en lotes, que es más eficiente que llamadas individuales.\n\n Args:\n textos: Lista de textos a embedder.\n batch_size: Tamaño del lote para cada llamada a Ollama.\n modelo: Nombre del modelo en Ollama.\n\n Returns:\n Lista de embeddings (cada uno es una lista de 1024 floats).\n \"\"\"\n resultados: list[list[float]] = []\n modelo = modelo or MODELO_EMBEDDINGS\n\n for i in range(0, len(textos), batch_size):\n batch = textos[i : i + batch_size]\n resp: Any = ollama.embed(model=modelo, input=batch)\n resultados.extend(resp[\"embeddings\"])\n\n return resultados\n\n\n# ── Empaquetado binario ──────────────────────────────────────────────────────\n\n\ndef pack_vector(vec: list[float]) -> bytes:\n \"\"\"Convierte una lista de floats a un BLOB binario (4 bytes por float).\n\n Args:\n vec: Lista de floats (1024 dimensiones para bge-m3).\n\n Returns:\n Bytes empaquetados con struct.pack (little-endian, float32).\n \"\"\"\n return struct.pack(f\"{len(vec)}f\", *vec)\n\n\ndef unpack_vector(data: bytes) -> list[float]:\n \"\"\"Convierte un BLOB binario de vuelta a una lista de floats.\n\n Args:\n data: Bytes empaquetados (de pack_vector).\n\n Returns:\n Lista de floats.\n \"\"\"\n # El formato 'f' usa float32 (4 bytes)\n n = len(data) // 4\n return list(struct.unpack(f\"{n}f\", data))\n\n\n# ── CLI de prueba ────────────────────────────────────────────────────────────\n\n\ndef main() -> None:\n \"\"\"Prueba rápida: genera un embedding y muestra estadísticas.\"\"\"\n texto = \"¿Cómo configuro systemd timers en Linux?\"\n print(f'Generando embedding para: \"{texto}\"\\n')\n\n embedding = get_embedding(texto)\n print(f\"Dimensiones: {len(embedding)}\")\n print(f\"Primeros 5 valores: {embedding[:5]}\")\n print(f\"Norma L2: {sum(x * x for x in embedding) ** 0.5:.6f}\")\n\n # Probar empaquetado\n blob = pack_vector(embedding)\n print(f\"\\nTamaño del BLOB: {len(blob)} bytes (esperado: {len(embedding) * 4})\")\n\n recovered = unpack_vector(blob)\n print(f\"Recuperados: {len(recovered)} floats\")\n print(f\"Coinciden: {abs(embedding[0] - recovered[0]) < 1e-6}\")\n\n\nif __name__ == \"__main__\":\n main()\n"

atareao / rag_pipeline.py

0 likes
0 forks
1 files
Last active 1 month ago
Pipeline incremental de indexacion RAG: escanea, chunk, embedding y FTS5 (episodio 819)
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"