#!/usr/bin/env python3 """ sync_missing_docs.py Gleicht Dokumente zwischen der Job-Matching-DB (NAS, job_matching.profile_documents) und der apply4jobs.de-RAG-DB (DO-Droplet, anythingllm_vectors) ab. Matching-Kriterium: filename (job_matching) <-> metadata->>'title' (anythingllm_vectors) Fehlende Dokumente werden: 1. aus profile_documents vollständig zusammengesetzt (Chunks in chunk_index-Reihenfolge) 2. neu gechunkt (gleiche Logik wie ingest.py: token-basiert, CHUNK_SIZE/CHUNK_OVERLAP) 3. mit dem LOKALEN Modell multilingual-e5-small neu embedded (die OpenAI-1536-Dim-Embeddings aus job_matching sind NICHT kompatibel mit der 384-Dim anythingllm_vectors-Spalte -> Text wird übernommen, Vektor neu berechnet) 4. via SSH-Tunnel in anythingllm_vectors eingefügt (gleiches Format wie ingest.py) Nutzung: python sync_missing_docs.py diff # nur anzeigen, was fehlt python sync_missing_docs.py sync # fehlende Dokumente übertragen python sync_missing_docs.py sync --dry-run # simulieren, nichts schreiben python sync_missing_docs.py sync --workspace-id 1 # nur bestimmten Workspace abgleichen python sync_missing_docs.py sync --force # ALLE Dokumente neu übertragen # (bestehende Titel werden vorher # in anythingllm_vectors gelöscht # und mit dem aktuellen Stand aus # job_matching neu eingefügt) Voraussetzung: SSH_HOST in .env.ingest muss auf den aktuellen DO-Droplet zeigen (aktuell noch die alte Hetzner-IP -> bitte vor dem Lauf aktualisieren). """ import os import uuid import json import argparse from datetime import datetime from pathlib import Path from dotenv import load_dotenv import psycopg2 from psycopg2.extras import execute_values, RealDictCursor from sshtunnel import SSHTunnelForwarder import tiktoken # --- .env laden --- load_dotenv(".env.ingest") # Ziel-DB (anythingllm_vectors) via SSH-Tunnel load_dotenv(".env.job_matching") # Quelle (job_matching, NAS, direkt) # Ziel (apply4jobs.de RAG-DB) SSH_HOST = os.getenv("SSH_HOST") SSH_USER = os.getenv("SSH_USER", "root") SSH_KEY_PATH = os.path.expanduser(os.getenv("SSH_KEY_PATH", "~/.ssh/id_ed25519")) SSH_KEY_PASSPHRASE = os.getenv("SSH_KEY_PASSPHRASE", None) DB_HOST = os.getenv("DB_HOST", "127.0.0.1") DB_PORT = int(os.getenv("DB_PORT", 5432)) DB_NAME = os.getenv("DB_NAME", "anythingllm") DB_USER = os.getenv("DB_USER", "anythingllm") DB_PASSWORD = os.getenv("DB_PASSWORD") NAMESPACE = os.getenv("NAMESPACE", "mein-workspace") CHUNK_SIZE = int(os.getenv("CHUNK_SIZE", 500)) CHUNK_OVERLAP = int(os.getenv("CHUNK_OVERLAP", 50)) EMBED_MODEL = "intfloat/multilingual-e5-small" # Quelle (job_matching, NAS) JM_DB_HOST = os.getenv("JM_DB_HOST", "192.168.178.128") JM_DB_PORT = int(os.getenv("JM_DB_PORT", 5433)) JM_DB_NAME = os.getenv("JM_DB_NAME", "job_matching") JM_DB_USER = os.getenv("JM_DB_USER") JM_DB_PASSWORD = os.getenv("JM_DB_PASSWORD") embed_model = None # lazy load, wie in ingest.py # --- Quelle: job_matching (NAS, direkt) --- def get_jm_connection(): return psycopg2.connect( host=JM_DB_HOST, port=JM_DB_PORT, dbname=JM_DB_NAME, user=JM_DB_USER, password=JM_DB_PASSWORD ) def list_jm_filenames(workspace_id: int | None) -> set[str]: """Alle distinct filenames aus profile_documents, optional nach workspace_id gefiltert.""" conn = get_jm_connection() cur = conn.cursor() if workspace_id is not None: cur.execute( "SELECT DISTINCT filename FROM profile_documents WHERE workspace_id = %s", (workspace_id,) ) else: cur.execute("SELECT DISTINCT filename FROM profile_documents") rows = {r[0] for r in cur.fetchall()} cur.close() conn.close() return rows def fetch_jm_document_text(filename: str, workspace_id: int | None) -> str: """Setzt ein Dokument aus seinen Chunks (chunk_index-Reihenfolge) wieder zusammen.""" conn = get_jm_connection() cur = conn.cursor() if workspace_id is not None: cur.execute( """SELECT content FROM profile_documents WHERE filename = %s AND workspace_id = %s ORDER BY chunk_index ASC""", (filename, workspace_id) ) else: cur.execute( """SELECT content FROM profile_documents WHERE filename = %s ORDER BY chunk_index ASC""", (filename,) ) rows = cur.fetchall() cur.close() conn.close() return "\n\n".join(r[0] for r in rows) # --- Ziel: anythingllm_vectors (DO-Droplet, via SSH-Tunnel) --- def get_target_connection(tunnel_port: int): return psycopg2.connect( host="127.0.0.1", port=tunnel_port, dbname=DB_NAME, user=DB_USER, password=DB_PASSWORD ) def list_target_titles(cur) -> set[str]: cur.execute( "SELECT DISTINCT metadata->>'title' FROM anythingllm_vectors WHERE namespace = %s", (NAMESPACE,) ) return {r[0] for r in cur.fetchall()} def delete_target_title(cur, title: str): """Löscht alle Chunks eines Titels im Ziel-Namespace (für --force Re-Sync).""" cur.execute( "DELETE FROM anythingllm_vectors WHERE metadata->>'title' = %s AND namespace = %s", (title, NAMESPACE) ) # --- Chunking (identisch zu ingest.py) --- def chunk_text(text: str, chunk_size: int = CHUNK_SIZE, overlap: int = CHUNK_OVERLAP) -> list[str]: enc = tiktoken.get_encoding("cl100k_base") tokens = enc.encode(text) chunks = [] start = 0 while start < len(tokens): end = min(start + chunk_size, len(tokens)) chunks.append(enc.decode(tokens[start:end])) start += chunk_size - overlap return chunks # --- Embeddings (lokal, wie ingest.py) --- def get_embeddings(texts: list[str]) -> list[list[float]]: global embed_model if embed_model is None: from sentence_transformers import SentenceTransformer print(" ⏳ Lade Embedding-Modell (einmalig)...") embed_model = SentenceTransformer(EMBED_MODEL) prefixed = [f"passage: {t}" for t in texts] return embed_model.encode(prefixed).tolist() def insert_chunks(cur, chunks: list[str], embeddings: list[list[float]], filename: str, source_note: str): now = datetime.now().isoformat() records = [] for chunk, embedding in zip(chunks, embeddings): metadata = { "id": str(uuid.uuid4()), "url": f"jobmatch-sync://{filename}", "text": chunk, "title": filename, "docSource": source_note, "published": now, "wordCount": len(chunk.split()), "chunkSource": f"jobmatch-sync://{filename}", } records.append((str(uuid.uuid4()), json.dumps(metadata), NAMESPACE, embedding)) execute_values( cur, "INSERT INTO anythingllm_vectors (id, metadata, namespace, embedding) VALUES %s", records, template="(%s::uuid, %s::jsonb, %s, %s::vector)" ) # --- Diff --- def compute_missing(workspace_id: int | None) -> list[str]: jm_files, _target_titles, missing = _compute_state(workspace_id) return missing def _compute_state(workspace_id: int | None) -> tuple[set[str], set[str], list[str]]: """Liefert (jm_files, target_titles, missing) in einem Rutsch, damit sync/--force nicht doppelt gegen beide DBs fragen muss.""" print("⏳ Lade Dateinamen aus job_matching (NAS)...") jm_files = list_jm_filenames(workspace_id) print(f" ✅ {len(jm_files)} Dokumente in job_matching") print("⏳ Verbinde mit apply4jobs.de-DB via SSH-Tunnel...") with SSHTunnelForwarder( (SSH_HOST, 22), ssh_username=SSH_USER, ssh_pkey=SSH_KEY_PATH, ssh_private_key_password=SSH_KEY_PASSPHRASE, remote_bind_address=("127.0.0.1", DB_PORT) ) as tunnel: conn = get_target_connection(tunnel.local_bind_port) cur = conn.cursor() target_titles = list_target_titles(cur) cur.close() conn.close() print(f" ✅ {len(target_titles)} Dokumente in anythingllm_vectors (Namespace '{NAMESPACE}')") missing = sorted(jm_files - target_titles) return jm_files, target_titles, missing # --- Sync --- def sync_missing(workspace_id, dry_run, force=False): jm_files, target_titles, missing = _compute_state(workspace_id) if force: to_process = sorted(jm_files) updates = sorted(jm_files & target_titles) new_docs = sorted(jm_files - target_titles) else: to_process = missing updates = [] new_docs = missing if not to_process: print("\n✅ Keine fehlenden Dokumente. Beide DBs sind synchron.") return if force: print(f"\n📋 {len(to_process)} Dokumente werden verarbeitet ({len(new_docs)} neu, {len(updates)} Update/Re-Sync):") else: print(f"\n📋 {len(to_process)} fehlende Dokumente:") for f in to_process: tag = " (Update)" if force and f in updates else "" print(f" - {f}{tag}") if dry_run: print("\n🔎 Dry-Run: nichts geschrieben.") return print(f"\n⏳ Verbinde mit apply4jobs.de-DB via SSH-Tunnel...") with SSHTunnelForwarder( (SSH_HOST, 22), ssh_username=SSH_USER, ssh_pkey=SSH_KEY_PATH, ssh_private_key_password=SSH_KEY_PASSPHRASE, remote_bind_address=("127.0.0.1", DB_PORT) ) as tunnel: conn = get_target_connection(tunnel.local_bind_port) cur = conn.cursor() success = 0 for filename in to_process: print(f"\n📄 {filename}") text = fetch_jm_document_text(filename, workspace_id) if not text.strip(): print(" ⚠️ Leerer Inhalt, übersprungen.") continue chunks = chunk_text(text) print(f" ✅ {len(chunks)} Chunks") embeddings = get_embeddings(chunks) print(f" ✅ {len(embeddings)} Embeddings (lokal, multilingual-e5-small)") if force and filename in target_titles: delete_target_title(cur, filename) print(" 🗑️ Alte Chunks gelöscht (Re-Sync)") insert_chunks(cur, chunks, embeddings, filename, source_note="jobmatch-sync-script") conn.commit() print(f" ✅ in anythingllm_vectors gespeichert (Namespace '{NAMESPACE}')") success += 1 cur.close() conn.close() print(f"\n✅ Fertig: {success}/{len(to_process)} Dokumente übertragen.") # --- CLI --- if __name__ == "__main__": parser = argparse.ArgumentParser(description="Dokumentenabgleich job_matching <-> apply4jobs.de") sub = parser.add_subparsers(dest="command") p_diff = sub.add_parser("diff", help="Nur anzeigen, welche Dokumente fehlen") p_diff.add_argument("--workspace-id", type=int, default=None) p_sync = sub.add_parser("sync", help="Fehlende Dokumente übertragen") p_sync.add_argument("--workspace-id", type=int, default=None) p_sync.add_argument("--dry-run", action="store_true") p_sync.add_argument("--force", action="store_true", help="Auch bereits vorhandene Titel neu einlesen (löscht alte Chunks und ersetzt sie)") args = parser.parse_args() if args.command == "diff": missing = compute_missing(args.workspace_id) if missing: print(f"\n📋 {len(missing)} fehlende Dokumente:") for f in missing: print(f" - {f}") else: print("\n✅ Keine fehlenden Dokumente.") elif args.command == "sync": sync_missing(args.workspace_id, args.dry_run, args.force) else: parser.print_help()