RAG

RAG - Retrieval-Augmented Generation

Diese Technologie ist derzeit sehr populär. Doch ehrlich gesagt bleiben Fragen offen, da die Knowledge Base vektoriell und letztlich statistisch ist und sich daher nicht für den Flugzeugbau oder die Durchführung medizinischer Operationen eignet. Ich bevorzuge neuro-symbolische KI, bei der die Wissensbasis beispielsweise aus den exakten S-Ausdrücken von Lisp besteht, die wir hier noch genauer betrachten werden.

1. ALL-IN-ONE Lösung

Wir stellen ein Schweizer Taschenmesser für alle Ihre Bedürfnisse her, sodass Sie nur für LLM-API bezahlen müssen (falls Sie kein eigenes besitzen).

Welche Klingen hat unser Schweizer Taschenmesser?

  • Knowledge Bases - PostgreSQL(pgvector)
  • Knowledge Base Pumping - ".txt", ".md", ".json", ".log", ".py", ".csv", ".pdf", ".docx", ".html", ".htm"
  • Textvektorisierung mit Sentence-Transformers
  • Generieren von Suchoptionen durch vLLM.
  • NLP-SpacyChunker
  • SEARCH ENGINE (BM25 + DENSE + RRF FUSION)
  • MCP-Generator (FastMCP)

Wir benötigen: 3 Server (vorzugsweise über ein lokales Netzwerk verbunden), 3 Subdomains mit A-Records.

  • Server 1 / Subdomain 1: llm.domain.de (10.0.0.5)
  • Server 2 / Subdomain 2: vpg.domain.de (10.0.0.7)
  • Server 3 / Subdomain 3: rag.domain.de (10.0.0.9)

2. Anweisungen zum Bereitstellen der Umgebung Server 1

Qwen2.5-0.5B-Instruct - Das Modell wird von Alibaba unter der Apache-2.0-Open-Source-Lizenz vertrieben.

  • Laden Sie es herunter und führen Sie es auf Ihrem Server oder Computer ohne Lizenzgebühren aus.
  • Nutzen Sie es für private, kommerzielle oder Forschungsprojekte.
  • Ändern und in Ihre Anwendungen integrieren

Referenz:

  • Das Qwen2.5-0.5B-Instruct-Modell ist ein extrem schlankes, universelles Sprachmodell mit nur 0,5 Milliarden Parametern (0,5 B). Der Zusatz „Instruct“ im Namen des Modells deutet darauf hin, dass es speziell für das Verständnis von Textanweisungen und Dialogen optimiert ist und nicht einfach nur Text fortsetzt.
  • Mikroaufgaben und Hilfslogik in RAG: Aufgrund seiner geringen Größe eignet es sich ideal als schnelles lokales "Gehirn" in Skripten - zum Beispiel zum Umformulieren von Suchanfragen (QueryRewriter), zum Extrahieren von Entitäten oder zum Filtern von Datenmüll, wenn große Modelle unpraktisch sind.
  • Funktioniert auch auf leistungsschwachen Servern und CPUs: Es benötigt nur minimalen Arbeitsspeicher (buchstäblich einige hundert Megabyte) und ist in der Lage, selbst auf einem Standardprozessor (CPU) ohne den Einsatz teurer Grafikkarten hohe Geschwindigkeiten zu erreichen.
  • Lokale Automatisierung: Geeignet für einfache Aufgaben wie Textklassifizierung, Zusammenfassung kurzer Passagen, Generierung von JSON-Strukturen aus einer Vorlage und einfache Chatbot-Konversationen, die keine tiefgreifende Expertenlogik erfordern.

Wenn es über das Internet funktioniert


python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip
pip install vllm
python -m vllm.entrypoints.openai.api_server --help

sudo -s

tee /etc/systemd/system/vllm.service > /dev/null << 'EOF'
[Unit]
Description=vLLM OpenAI-Compatible API Server (CPU Mode)
After=network.target

[Service]
Type=simple
User=root
WorkingDirectory=/path/to/your/project
Environment="VLLM_TARGET_DEVICE=cpu"
ExecStart=/path/to/your/.venv/bin/python -m vllm.entrypoints.openai.api_server --model Qwen/Qwen2.5-0.5B-Instruct --host 127.0.0.1 --port 8000 --device cpu
Restart=always
RestartSec=10

[Install]
WantedBy=multi-user.target
EOF

systemctl daemon-reload
systemctl start vllm
systemctl enable vllm
systemctl status vllm

apt install nginx

tee /etc/nginx/sites-available/vllm > /dev/null << 'EOF'
server {
    listen 80;
    server_name llm.domain.de;

    location / {
        proxy_pass http://127.0.0.1:8000;
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
        proxy_set_header X-Forwarded-Proto $scheme;

        proxy_read_timeout 600s;
        proxy_send_timeout 600s;
    }
}
EOF

ln -s /etc/nginx/sites-available/vllm /etc/nginx/sites-enabled/
nginx -t
/etc/init.d/nginx restart

apt install certbot python3-certbot-nginx -y
certbot --nginx

crontab -e
0 3 * * * /usr/bin/certbot renew --quiet --deploy-hook "systemctl reload nginx"

certbot --nginx --redirect
nginx -t
/etc/init.d/nginx restart
exit

Wenn ein lokales Netzwerk vorhanden ist


python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip
pip install vllm
python -m vllm.entrypoints.openai.api_server --help

sudo -s

tee /etc/systemd/system/vllm.service > /dev/null << 'EOF'
[Unit]
Description=vLLM OpenAI-Compatible API Server (CPU Mode)
After=network.target

[Service]
Type=simple
User=root
WorkingDirectory=/path/to/your/project
Environment="VLLM_TARGET_DEVICE=cpu"
ExecStart=/path/to/your/.venv/bin/python -m vllm.entrypoints.openai.api_server --model Qwen/Qwen2.5-0.5B-Instruct --host 127.0.0.1 --port 8000 --device cpu
Restart=always
RestartSec=10

[Install]
WantedBy=multi-user.target
EOF

systemctl daemon-reload
systemctl start vllm
systemctl enable vllm
systemctl status vllm

apt install nginx

tee /etc/nginx/sites-available/vllm > /dev/null << 'EOF'
server {
    listen 10.0.0.5:80;

    location / {
        proxy_pass http://127.0.0.1:8000;
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
        proxy_set_header X-Forwarded-Proto $scheme;

        proxy_read_timeout 600s;
        proxy_send_timeout 600s;
    }
}
EOF

ln -s /etc/nginx/sites-available/vllm /etc/nginx/sites-enabled/
nginx -t
/etc/init.d/nginx restart
exit

3. Anweisungen zum Bereitstellen der Umgebung Server 2

Eine Architekturlösung zur Aufgabe von Multi-Datenbank-Konfigurationen (Qdrant/Chroma) zugunsten eines einzigen PostgreSQL-Speichers:

  • 1. BM25 Indexproblem (OOM-Engpass): Bei der Verwendung von In-Memory-Lösungen (z. B. rank_bm25 über lokales SQLite/ChromaDB/Qdrant) muss der gesamte Lemma-Korpus in den Arbeitsspeicher geladen werden, um BM25Okapi zu generieren bei jedem Aufruf von reload_index(). Bei Millionen von Datenblöcken führt dies zu Speichermangel (OOM) und anhaltenden CPU-Sperren. Die Verwendung von PostgreSQL ermöglicht den Wechsel zur nativen Volltextsuche (FTS) oder zu indiziertem JSONB, wodurch das Laden des Korpus in den Arbeitsspeicher entfällt.
  • 2. Optimierung der Vektorsuche in PostgreSQL mit HNSW: Standardmäßig führt die Kosinusdistanzsuche (`vector <=> query_vector`) einen sequenziellen Scan durch (eine vollständige Suche aller Tabellenzeilen). Die Verwendung eines HNSW-Index (Hierarchical Navigable Small World) erstellt einen Nächste-Nachbarn-Graphen direkt im DBMS und reduziert so die Suchkomplexität von O(N) auf O(log N).

sudo -s
apt install -y wget gnupg
sh -c 'echo "deb http://apt.postgresql.org/pub/repos/apt $(lsb_release -cs)-pgdg main" > /etc/apt/sources.list.d/pgdg.list'
wget --quiet -O - https://www.postgresql.org/media/keys/ACCC4CF8.asc | sudo apt-key add -
apt update
apt install -y postgresql-16 postgresql-client-16 postgresql-16-pgvector

su - postgres -c psql
CREATE DATABASE rag_db;
\c rag_db;
CREATE EXTENSION vector;
\dx
\q

vi /etc/postgresql/16/main/postgresql.conf
    listen_addresses = '*'

vi /etc/postgresql/16/main/pg_hba.conf
    # dostup tolko po localnoj seti
    host    rag_db       all             10.0.0.0/24          scram-sha-256
    # dospup cherez iternet, no s konkretnogo IP
    hostssl rag_db       all             xxx.x.xxx.xx/32         scram-sha-256
    # polnih dostup
    hostssl rag_db       all             0.0.0.0/0         scram-sha-256

su - postgres -c psql
CREATE USER rag_user WITH ENCRYPTED PASSWORD 'XXXXXXXXXXXXXXXXX';
GRANT ALL PRIVILEGES ON DATABASE rag_db TO rag_user;
\c rag_db;
GRANT ALL ON SCHEMA public TO rag_user;
\q
systemctl restart postgresql

netstat -na | grep 5432
# Wenn UFW vorhanden ist und der Port geschlossen ist
# für lokale Netzwerke
ufw allow from 10.0.0.0/24 to any port 5432 proto tcp
# Für externen Zugriff:
ufw allow 5432/tcp

# wenn hostssl
apt install certbot -y
certbot certonly --standalone -d vpg.domain.de
mkdir -p /etc/postgresql/16/main/certs
cp /etc/letsencrypt/live/vpg.domain.de/fullchain.pem /etc/postgresql/16/main/certs/server.crt
cp /etc/letsencrypt/live/vpg.domain.de/privkey.pem /etc/postgresql/16/main/certs/server.key
chown -R postgres:postgres /etc/postgresql/16/main/certs
chmod 600 /etc/postgresql/16/main/certs/server.key
chmod 644 /etc/postgresql/16/main/certs/server.crt
vi /etc/postgresql/16/main/postgresql.conf
    ssl = on
    ssl_cert_file = 'certs/server.crt'
    ssl_key_file = 'certs/server.key'

systemctl restart postgresql

tee /etc/letsencrypt/renewal-hooks/deploy/postgresql-ssl.sh > /dev/null << 'EOF'
#!/bin/bash
DOMAIN="vpg.domain.de"
CERTS_DIR="/etc/postgresql/16/main/certs"
if [ "$RENEWED_LINEAGE" = "/etc/letsencrypt/live/$DOMAIN" ]; then
    cp /etc/letsencrypt/live/$DOMAIN/fullchain.pem $CERTS_DIR/server.crt
    cp /etc/letsencrypt/live/$DOMAIN/privkey.pem $CERTS_DIR/server.key
    chown -R postgres:postgres $CERTS_DIR
    chmod 600 $CERTS_DIR/server.key
    chmod 644 $CERTS_DIR/server.crt
    systemctl reload postgresql
fi
EOF
chmod +x /etc/letsencrypt/renewal-hooks/deploy/postgresql-ssl.sh
exit

4. Anweisungen zum Bereitstellen der Umgebung Server 3

Unsere API befindet sich hier.


sudo -s
apt install nginx
cd /var/www
mkdir rag-api
cd rag-api
python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip setuptools wheel
                

Erstellen Sie eine Datei requirements.txt


spacy>=3.7.0
sentence-transformers>=2.3.0
torch>=2.0.0
aiohttp>=3.8.0
psycopg2-binary>=2.9.0
rank-bm25>=0.2.2
spacy>=3.0.0
sentence-transformers>=2.2.0
pypdf>=3.0.0
python-docx>=0.8.11
beautifulsoup4>=4.9.0
Flask>=2.0.0
openai>=1.0.0
gunicorn>=20.0.0
                

Wir installieren alles Notwendige


pip install -r requirements.txt
                

Unsere internen Neuronen netze installieren


python3 -m spacy download de_core_news_lg
# "de" - kann ersetzt werden durch "ru", "en", "fr", "es", ...
# "lg" - kann ersetzt werden durch "sm", "md"
python -c "from sentence_transformers import SentenceTransformer; SentenceTransformer('sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2')"
# Ein wichtiger Punkt: Wir müssen die Frage klären: GPU oder CPU? Falls Sie sich für die GPU entscheiden, können wir fortfahren. Falls Sie die CPU wählen, überprüfen wir, ob alles ordnungsgemäß funktioniert.
python -c "from sentence_transformers import SentenceTransformer; m = SentenceTransformer('sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2', device='cpu'); print(m.device)"
                

Daher ist in den meisten Fällen keine Konfiguration erforderlich – das Modell wird automatisch auf der CPU ausgeführt. Wenn Sie die Ausführung jedoch fest auf der CPU festlegen möchten (beispielsweise auf einem Server ohne GPU, um unnötige Geräteprüfungen zu vermeiden), müssen Sie das Gerät bei der Initialisierung des Modells explizit angeben (device="cpu").

OpenMP/MKL Thread-Begrenzung


import torch
# Geben Sie die Anzahl der Threads an. Wenn Sie 8 Threads haben, geben Sie 4 Threads an. Sie müssen die Hälfte der Threads freilassen, da wir "ThreadPoolExecutor" verwenden.
torch.set_num_threads(4)
                

Wenn Geschwindigkeit bei schwachen CPUs entscheidend ist, kann das Modell in das int8-Format konvertiert werden, aber für ein leichtes Modell wie das MiniLM-L12-v2 ist dies auf Standardprozessoren normalerweise nicht notwendig – es läuft bereits sehr schnell (zehn bis hunderte von Texten pro Sekunde).

Erstellen Sie eine Datei rag.py


import argparse
import asyncio
import json
import math
import os
from concurrent.futures import ThreadPoolExecutor
from dataclasses import asdict, dataclass, field
from pathlib import Path
from typing import Any, Dict, List, Optional, Set, Union, cast
import aiohttp
from rank_bm25 import BM25Okapi

try:
    import psycopg2
    import psycopg2.extras
    POSTGRES_AVAILABLE = True
except ImportError:
    POSTGRES_AVAILABLE = False

# --- INTERFACE DATA ---
@dataclass
class TextChunk:
    id: str
    doc_id: str
    source_file: str
    text: str
    lemmas: List[str]
    entities: Dict[str, List[str]]
    vector: Optional[List[float]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def to_dict(self) -> Dict[str, Any]:
        d = dict(self.__dict__)
        if d.get("vector"):
            d["vector_dim"] = len(d["vector"])
            del d["vector"]
        return d

# --- INTERFACE SEARCH ---
@dataclass
class SearchResult:
    chunk_id: str
    doc_id: str
    source_file: str
    text: str
    rrf_score: float
    bm25_score: float
    dense_score: float
    metadata: Dict[str, Any]

# --- DOCUMENT READER ENGINE ---
class DocumentReader:
    """Extrahiert Text und Metadaten aus lokalen Dateien verschiedener Formate."""
    @staticmethod
    def read_file(file_path: Union[str, Path]) -> Dict[str, Any]:
        path = Path(file_path)
        if not path.exists():
            raise FileNotFoundError(f"Файл не найден: {path}")
        ext = path.suffix.lower().replace(".", "")
        content = ""
        meta = {
            "file_name": path.name,
            "file_path": str(path.absolute()),
            "file_size": path.stat().st_size,
            "extension": ext,
        }
        # Textdateien, Markdown und Quellcode
        if ext in ["txt", "md", "json", "log", "py", "csv"]:
            with open(path, "r", encoding="utf-8", errors="ignore") as f:
                content = f.read()
        # PDF-Dokumente
        elif ext == "pdf":
            try:
                import pypdf
                reader = pypdf.PdfReader(path)
                pages = [
                    page.extract_text()
                    for page in reader.pages
                    if page.extract_text()
                ]
                content = "\n\n".join(pages)
                meta["num_pages"] = len(reader.pages)
            except ImportError:
                raise ImportError("Zum Lesen von PDFs installieren: pip install pypdf")

        # Microsoft Word
        elif ext == "docx":
            try:
                import docx
                doc = docx.Document(path)
                content = "\n".join(
                    [p.text for p in doc.paragraphs if p.text.strip()]
                )
            except ImportError:
                raise ImportError(
                    "Zum Lesen von DOCX-Dateien installieren: pip install python-docx"
                )
        # HTML
        elif ext in ["html", "htm"]:
            try:
                from bs4 import BeautifulSoup
                with open(path, "r", encoding="utf-8", errors="ignore") as f:
                    soup = BeautifulSoup(f.read(), "html.parser")
                    for script in soup(["script", "style"]):
                        script.decompose()
                    content = soup.get_text(separator="\n")
            except ImportError:
                raise ImportError(
                    "Um HTML lesen zu können, installieren: pip install beautifulsoup4"
                )
        else:
            with open(path, "r", encoding="utf-8", errors="ignore") as f:
                content = f.read()
        return {"content": content, "metadata": meta}

# --- SPACY NLP CORE ENGINE (OPTIMIZED) ---
class SpacyNLPProcessor:
    """Ein eigenständiger NLP-Prozessor auf Basis von spaCy mit optimierter Single-Pass-Verarbeitung."""

    def __init__(self, model_name: str = "de_core_news_lg"):
        try:
            import spacy
            self.nlp = spacy.load(model_name)
        except OSError:
            import spacy.cli
            spacy.cli.download(model_name)
            self.nlp = spacy.load(model_name)

    def analyze_doc(self, text: str):
        """Führt einen einzelnen spaCy-Durchlauf über den Text durch."""
        return self.nlp(text)

    def extract_features_from_doc(self, doc) -> Dict[str, Any]:
        """Extrahiert Lemmata und Entitäten aus einem einzelnen spaCy-Dokument ohne erneutes Parsen."""
        lemmas = [
            token.lemma_
            for token in doc
            if not token.is_stop
            and not token.is_punct
            and not token.is_space
            and len(token.lemma_) > 1
        ]
        entities: Dict[str, List[str]] = {}
        for ent in doc.ents:
            label = ent.label_
            if label not in entities:
                entities[label] = []
            if ent.text not in entities[label]:
                entities[label].append(ent.text)
        return {"lemmas": lemmas, "entities": entities}

    def lemmatize_query(self, text: str) -> List[str]:
        doc = self.nlp(text.lower())
        return [
            token.lemma_
            for token in doc
            if not token.is_stop
            and not token.is_punct
            and not token.is_space
            and len(token.lemma_) > 1
        ]

# --- SEMANTIC SPACY CHUNKER ---
class SemanticSpacyChunker:
    """Ein Chunker, der Satzgrenzen respektiert und optimiertes Single-Pass-spaCy verwendet."""

    def __init__(
        self,
        nlp_processor: SpacyNLPProcessor,
        max_chunk_size: int = 512,
        overlap_sents: int = 1,
    ):
        self.nlp = nlp_processor
        self.max_chunk_size = max_chunk_size
        self.overlap_sents = overlap_sents

    def chunk_document(
        self,
        text: str,
        doc_id: str,
        source_file: str,
        extra_meta: Optional[Dict[str, Any]] = None,
    ) -> List[TextChunk]:
        full_doc = self.nlp.analyze_doc(text)
        sentences = [
            sent for sent in full_doc.sents if len(sent.text.strip()) > 0
        ]
        chunks: List[TextChunk] = []
        current_sent_objs = []
        current_len = 0
        chunk_idx = 0
        for sent in sentences:
            sent_text = sent.text.strip()
            sent_len = len(sent_text)
            if current_len + sent_len > self.max_chunk_size and current_sent_objs:
                chunk_text = " ".join([s.text.strip() for s in current_sent_objs])
                chunk_doc = self.nlp.analyze_doc(chunk_text)
                features = self.nlp.extract_features_from_doc(chunk_doc)
                meta = {"chunk_index": chunk_idx}
                if extra_meta:
                    meta.update(extra_meta)
                chunks.append(
                    TextChunk(
                        id=f"{doc_id}_c{chunk_idx}",
                        doc_id=doc_id,
                        source_file=source_file,
                        text=chunk_text,
                        lemmas=features["lemmas"],
                        entities=features["entities"],
                        metadata=meta,
                    )
                )
                chunk_idx += 1
                current_sent_objs = (
                    current_sent_objs[-self.overlap_sents :]
                    if self.overlap_sents > 0
                    else []
                )
                current_len = sum(len(s.text.strip()) for s in current_sent_objs)
            current_sent_objs.append(sent)
            current_len += sent_len
        if current_sent_objs:
            chunk_text = " ".join([s.text.strip() for s in current_sent_objs])
            chunk_doc = self.nlp.analyze_doc(chunk_text)
            features = self.nlp.extract_features_from_doc(chunk_doc)
            meta = {"chunk_index": chunk_idx}
            if extra_meta:
                meta.update(extra_meta)
            chunks.append(
                TextChunk(
                    id=f"{doc_id}_c{chunk_idx}",
                    doc_id=doc_id,
                    source_file=source_file,
                    text=chunk_text,
                    lemmas=features["lemmas"],
                    entities=features["entities"],
                    metadata=meta,
                )
            )
        return chunks

# --- MICRO-LLM QUERY REWRITER (STRICT VLLM ENGINE) ---
class QueryRewriter:
    """Generierung von Suchoptionen über vLLM (OpenAI API-kompatibel)."""

    def __init__(
        self,
        endpoint_url: Optional[str] = None,
        model_name: str = "Qwen/Qwen2.5-0.5B-Instruct",
    ):
        self.endpoint_url = endpoint_url
        self.model_name = model_name

    async def rewrite(self, original_query: str) -> List[str]:
        if not self.endpoint_url:
            return [original_query]
        payload = {
            "model": self.model_name,
            "messages": [
                {
                    "role": "system",
                    "content": (
                        "Du bist das RAG SEO-Modul. "
                        "Wandeln Sie Ihre Suchanfrage in zwei alternative Schlüsselbegriffe um. "
                        'Bitte antworten Sie NUR mit einem JSON-Objekt im folgenden Format: {"queries": ["formulation 1", "formulation 2"]}'
                    ),
                },
                {"role": "user", "content": original_query},
            ],
            "response_format": {"type": "json_object"},
            "temperature": 0.3,
            "max_tokens": 150,
        }
        try:
            async with aiohttp.ClientSession() as session:
                async with session.post(
                    self.endpoint_url, json=payload, timeout=3.0
                ) as resp:
                    if resp.status == 200:
                        data = await resp.json()
                        raw_content = data["choices"][0]["message"]["content"]
                        parsed = json.loads(raw_content)
                        queries = (
                            parsed.get("queries", [])
                            if isinstance(parsed, dict)
                            else parsed
                        )
                        if isinstance(queries, list) and len(queries) > 0:
                            return list(
                                dict.fromkeys(
                                    [original_query] + [str(q) for q in queries]
                                )
                            )
        except Exception:
            pass
        return [original_query]

# --- DENSE EMBEDDER ENGINE ---
class LocalDenseEmbedder:
    """Textvektorisierung mittels Satztransformatoren (Fallback: Hash-Vektorisierer)."""

    def __init__(
        self,
        model_name: str = "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2",
    ):
        try:
            from sentence_transformers import SentenceTransformer
            self.model = SentenceTransformer(model_name)
            self.dim = self.model.get_embedding_dimension()
            self.use_st = True
        except (ImportError, OSError, ValueError):
            self.use_st = False
            self.dim = 384

    def embed(self, text: str) -> List[float]:
        if self.use_st:
            return self.model.encode(text, normalize_embeddings=True).tolist()
        import hashlib
        words = text.lower().split()
        vec = [0.0] * self.dim
        for w in words:
            idx = int(hashlib.md5(w.encode()).hexdigest(), 16) % self.dim
            vec[idx] += 1.0
        norm = math.sqrt(sum(v * v for v in vec)) or 1.0
        return [v / norm for v in vec]

    def embed_batch(self, texts: List[str]) -> List[List[float]]:
        if not texts:
            return []
        if self.use_st:
            embeddings = self.model.encode(
                texts, normalize_embeddings=True, show_progress_bar=False
            )
            return [e.tolist() for e in embeddings]
        return [self.embed(t) for t in texts]

# --- POSTGRESQL STORAGE BACKEND ---
class StorageBackend:

    def __init__(self, connection_str: str, vector_dim: int = 384):
        if not POSTGRES_AVAILABLE:
            raise ImportError(
                "psycopg2-binary is missing. Install via: pip install psycopg2-binary"
            )
        self.connection_str = connection_str
        self.vector_dim = vector_dim
        self._init_backend()

    def _init_backend(self):
        with psycopg2.connect(self.connection_str) as conn:
            with conn.cursor() as cur:
                cur.execute("CREATE EXTENSION IF NOT EXISTS vector;")
                cur.execute(
                    f"""
                    CREATE TABLE IF NOT EXISTS rag_chunks (
                        id VARCHAR(255) PRIMARY KEY,
                        doc_id VARCHAR(255),
                        source_file VARCHAR(512),
                        text TEXT,
                        lemmas JSONB,
                        entities JSONB,
                        metadata JSONB,
                        vector vector({self.vector_dim})
                    );
                """
                )
                # Optimierung der Vektorsuche mittels Kosinusdistanz unter Verwendung des HNSW-Index
                cur.execute(
                    """
                    CREATE INDEX IF NOT EXISTS idx_rag_chunks_vector
                    ON rag_chunks USING hnsw (vector vector_cosine_ops);
                """
                )
            conn.commit()

    def save_chunk(self, chunk: TextChunk):
        self.save_chunks_batch([chunk])

    def save_chunks_batch(self, chunks: List[TextChunk]):
        if not chunks:
            return

        with psycopg2.connect(self.connection_str) as conn:
            with conn.cursor() as cur:
                psycopg2.extras.execute_values(
                    cur,
                    """
                    INSERT INTO rag_chunks (id, doc_id, source_file, text, lemmas, entities, metadata, vector)
                    VALUES %s
                    ON CONFLICT (id) DO UPDATE SET
                        text = EXCLUDED.text, lemmas = EXCLUDED.lemmas,
                        entities = EXCLUDED.entities, metadata = EXCLUDED.metadata, vector = EXCLUDED.vector;
                    """,
                    [
                        (
                            c.id,
                            c.doc_id,
                            c.source_file,
                            c.text,
                            json.dumps(c.lemmas, ensure_ascii=False),
                            json.dumps(c.entities, ensure_ascii=False),
                            json.dumps(c.metadata, ensure_ascii=False),
                            c.vector,
                        )
                        for c in chunks
                    ],
                )
            conn.commit()

    def fetch_all(self) -> List[TextChunk]:
        chunks = []
        with psycopg2.connect(self.connection_str) as conn:
            with conn.cursor(cursor_factory=psycopg2.extras.DictCursor) as cur:
                cur.execute(
                    "SELECT id, doc_id, source_file, text, lemmas, entities, metadata, vector FROM rag_chunks;"
                )
                for r in cur.fetchall():
                    vec = (
                        [float(x) for x in r["vector"].strip("[]").split(",")]
                        if r["vector"]
                        else None
                    )
                    chunks.append(
                        TextChunk(
                            id=r["id"],
                            doc_id=r["doc_id"],
                            source_file=r["source_file"],
                            text=r["text"],
                            lemmas=(
                                r["lemmas"]
                                if isinstance(r["lemmas"], list)
                                else json.loads(r["lemmas"])
                            ),
                            entities=(
                                r["entities"]
                                if isinstance(r["entities"], dict)
                                else json.loads(r["entities"])
                            ),
                            metadata=(
                                r["metadata"]
                                if isinstance(r["metadata"], dict)
                                else json.loads(r["metadata"])
                            ),
                            vector=vec,
                        )
                    )
        return chunks

    def search_dense(
        self, query_vec: List[float], top_k: int = 5
    ) -> List[tuple[str, float]]:
        vec_str = f"[{','.join(map(str, query_vec))}]"
        with psycopg2.connect(self.connection_str) as conn:
            with conn.cursor() as cur:
                cur.execute(
                    """
                    SELECT id, 1 - (vector <=> %s::vector) AS similarity
                    FROM rag_chunks ORDER BY vector <=> %s::vector ASC LIMIT %s;
                """,
                    (vec_str, vec_str, top_k),
                )
                return [(r[0], float(r[1])) for r in cur.fetchall()]

# --- SEARCH ENGINE (BM25 + DENSE + RRF FUSION) ---
class HybridSearchEngine:
    def __init__(
        self,
        nlp_processor: SpacyNLPProcessor,
        embedder: LocalDenseEmbedder,
        storage: StorageBackend,
    ):
        self.nlp = nlp_processor
        self.embedder = embedder
        self.storage = storage
        self.chunks_cache: Dict[str, TextChunk] = {}
        self.bm25: Optional[BM25Okapi] = None
        self.chunk_ids_order: List[str] = []
        self.reload_index()

    def reload_index(self):
        chunks = self.storage.fetch_all()
        self.chunks_cache = {c.id: c for c in chunks}
        self.chunk_ids_order = list(self.chunks_cache.keys())
        corpus_lemmas = [c.lemmas for c in chunks]
        if corpus_lemmas:
            self.bm25 = BM25Okapi(corpus_lemmas)

    def search_sparse_only(
        self, query: str, top_k: int = 5
    ) -> List[Dict[str, Any]]:
        if not self.bm25:
            return []
        q_lemmas = self.nlp.lemmatize_query(query)
        scores = self.bm25.get_scores(q_lemmas)
        ranked = sorted(
            zip(self.chunk_ids_order, scores), key=lambda x: x[1], reverse=True
        )[:top_k]
        return [
            {
                "chunk_id": cid,
                "bm25_score": round(float(s), 4),
                "text": self.chunks_cache[cid].text,
            }
            for cid, s in ranked
        ]

    def search_dense_only(
        self, query: str, top_k: int = 5
    ) -> List[Dict[str, Any]]:
        q_vec = self.embedder.embed(query)
        results = self.storage.search_dense(q_vec, top_k=top_k)
        return [
            {
                "chunk_id": cid,
                "dense_score": round(float(s), 4),
                "text": self.chunks_cache[cid].text,
            }
            for cid, s in results
        ]

    def search_hybrid_rrf(
        self, queries: List[str], top_k: int = 5, rrf_k: int = 60
    ) -> List[SearchResult]:
        if not self.chunks_cache or not self.bm25:
            return []
        rrf_scores: Dict[str, float] = {cid: 0.0 for cid in self.chunks_cache}
        bm25_best: Dict[str, float] = {cid: 0.0 for cid in self.chunks_cache}
        dense_best: Dict[str, float] = {cid: 0.0 for cid in self.chunks_cache}
        candidate_depth = max(top_k * 5, 50)
        for q in queries:
            q_lemmas = self.nlp.lemmatize_query(q)
            q_vec = self.embedder.embed(q)
            # 1. BM25 Sparse Search
            bm25_raw = (
                self.bm25.get_scores(q_lemmas)
                if q_lemmas
                else [0.0] * len(self.chunk_ids_order)
            )
            bm25_ranked = sorted(
                zip(self.chunk_ids_order, bm25_raw),
                key=lambda x: x[1],
                reverse=True,
            )[:candidate_depth]
            for rank, (cid, score) in enumerate(bm25_ranked):
                rrf_scores[cid] += 1.0 / (rrf_k + rank + 1)
                bm25_best[cid] = max(bm25_best[cid], float(score))
            # 2. PostgreSQL Vector Search
            dense_ranked = self.storage.search_dense(
                q_vec, top_k=candidate_depth
            )
            for rank, (cid, score) in enumerate(dense_ranked):
                rrf_scores[cid] += 1.0 / (rrf_k + rank + 1)
                dense_best[cid] = max(dense_best[cid], float(score))
        sorted_rrf = sorted(
            rrf_scores.items(), key=lambda x: x[1], reverse=True
        )[:top_k]
        output = []
        for cid, score in sorted_rrf:
            if score == 0.0:
                continue
            c = self.chunks_cache[cid]
            output.append(
                SearchResult(
                    chunk_id=c.id,
                    doc_id=c.doc_id,
                    source_file=c.source_file,
                    text=c.text,
                    rrf_score=round(score, 5),
                    bm25_score=round(bm25_best[cid], 4),
                    dense_score=round(dense_best[cid], 4),
                    metadata={**c.metadata, "entities": c.entities},
                )
            )
        return output

# --- FACADE SYSTEM MANAGER & DIRECTORY SCANNER ---
class RAGEngine:

    def __init__(
        self,
        db_conn: str,
        spacy_model: str = "de_core_news_lg",
        vllm_url: Optional[str] = None,
        max_workers: int = 4,
    ):
        self.nlp = SpacyNLPProcessor(model_name=spacy_model)
        self.chunker = SemanticSpacyChunker(self.nlp)
        self.embedder = LocalDenseEmbedder()
        self.rewriter = QueryRewriter(endpoint_url=vllm_url)
        self.storage = StorageBackend(
            connection_str=db_conn, vector_dim=self.embedder.dim
        )
        self.search_engine = HybridSearchEngine(
            self.nlp, self.embedder, self.storage
        )
        self.max_workers = max_workers
        self.pool = ThreadPoolExecutor(max_workers=max_workers)

    def ingest_file(
        self, file_path: str, max_chunk_size: int = 512
    ) -> Dict[str, Any]:
        """Dokumentenanalyse mittels DocumentReader und Vektorisierung in PostgreSQL."""
        doc_data = DocumentReader.read_file(file_path)
        content = doc_data["content"]
        file_meta = doc_data["metadata"]
        if not content.strip():
            return {
                "file": file_path,
                "chunks_created": 0,
                "status": "skipped_empty",
            }
        self.chunker.max_chunk_size = max_chunk_size
        doc_id = file_meta["file_name"]
        chunks = self.chunker.chunk_document(
            content, doc_id=doc_id, source_file=file_path, extra_meta=file_meta
        )
        texts = [c.text for c in chunks]
        vectors = self.embedder.embed_batch(texts)
        for chunk, vec in zip(chunks, vectors):
            chunk.vector = vec
        self.storage.save_chunks_batch(chunks)
        return {
            "file": file_path,
            "chunks_created": len(chunks),
            "storage_backend": "POSTGRESQL",
            "metadata": file_meta,
        }

    async def ingest_file_async(
        self, file_path: str, max_chunk_size: int = 512
    ) -> Dict[str, Any]:
        loop = asyncio.get_running_loop()
        return await loop.run_in_executor(
            self.pool, self.ingest_file, file_path, max_chunk_size
        )

    def ingest_directory(
        self,
        dir_path: str,
        extensions: Optional[Set[str]] = None,
        max_chunk_size: int = 512,
    ) -> Dict[str, Any]:
        if extensions is None:
            extensions = {
                ".txt",
                ".md",
                ".json",
                ".log",
                ".py",
                ".csv",
                ".pdf",
                ".docx",
                ".html",
                ".htm",
            }
        root = Path(dir_path)
        if not root.exists() or not root.is_dir():
            raise NotADirectoryError(
                f"Das Verzeichnis {dir_path} existiert nicht oder ist kein Ordner."
            )
        files_to_process = [
            str(path)
            for path in root.rglob("*")
            if path.is_file() and path.suffix.lower() in extensions
        ]
        file_summary = []
        processed_files = 0
        total_chunks = 0
        futures = {
            self.pool.submit(self.ingest_file, fp, max_chunk_size): fp
            for fp in files_to_process
        }
        for future in futures:
            try:
                res = future.result()
                processed_files += 1
                total_chunks += res.get("chunks_created", 0)
                file_summary.append(res)
            except Exception as e:
                fp = futures[future]
                file_summary.append(
                    {"file": fp, "error": str(e), "chunks_created": 0}
                )
        self.search_engine.reload_index()
        return {
            "directory": str(root.absolute()),
            "processed_files_count": processed_files,
            "total_chunks_created": total_chunks,
            "total_chunks_in_db": len(self.search_engine.chunks_cache),
            "storage_backend": "POSTGRESQL",
            "files": file_summary,
        }

    async def ingest_directory_async(
        self,
        dir_path: str,
        extensions: Optional[Set[str]] = None,
        max_chunk_size: int = 512,
    ) -> Dict[str, Any]:
        if extensions is None:
            extensions = {
                ".txt",
                ".md",
                ".json",
                ".log",
                ".py",
                ".csv",
                ".pdf",
                ".docx",
                ".html",
                ".htm",
            }
        root = Path(dir_path)
        if not root.exists() or not root.is_dir():
            raise NotADirectoryError(
                f"Das Verzeichnis {dir_path} existiert nicht oder ist kein Ordner."
            )
        files_to_process = [
            str(path)
            for path in root.rglob("*")
            if path.is_file() and path.suffix.lower() in extensions
        ]
        loop = asyncio.get_running_loop()
        tasks = [
            loop.run_in_executor(
                self.pool, self.ingest_file, fp, max_chunk_size
            )
            for fp in files_to_process
        ]
        results = await asyncio.gather(*tasks, return_exceptions=True)
        file_summary = []
        processed_files = 0
        total_chunks = 0
        for fp, res in zip(files_to_process, results):
            if isinstance(res, BaseException):
                file_summary.append(
                    {"file": fp, "error": str(res), "chunks_created": 0}
                )
            else:
                res_dict = cast(Dict[str, Any], res)
                processed_files += 1
                total_chunks += res_dict.get("chunks_created", 0)
                file_summary.append(res_dict)
        self.search_engine.reload_index()
        return {
            "directory": str(root.absolute()),
            "processed_files_count": processed_files,
            "total_chunks_created": total_chunks,
            "total_chunks_in_db": len(self.search_engine.chunks_cache),
            "storage_backend": "POSTGRESQL",
            "files": file_summary,
        }

# --- CLI INTERFACE ---
def main():
    parser = argparse.ArgumentParser(
        description="RAGEngine PostgreSQL Architecture"
    )
    parser.add_argument(
        "--db-conn",
        type=str,
        default=os.getenv(
            "POSTGRES_CONN_STR",
            "postgresql://postgres:postgres@localhost:5432/rag_db",
        ),
        help="PostgreSQL Connection String",
    )
    parser.add_argument(
        "--spacy-model", type=str, default="de_core_news_lg"
    )
    parser.add_argument(
        "--vllm-url",
        type=str,
        default=None,
        help="vLLM Chat Completions URL (e.g. http://localhost:8000/v1/chat/completions)",
    )
    parser.add_argument(
        "--max-workers",
        type=int,
        default=4,
        help="ThreadPoolExecutor worker threads count",
    )
    subparsers = parser.add_subparsers(dest="command", required=True)
    # NLP Commands
    chunk_p = subparsers.add_parser(
        "chunk", help="[Chunker] Test semantic chunking"
    )
    chunk_p.add_argument("--text", type=str, required=True)
    chunk_p.add_argument("--max-size", type=int, default=200)
    rw_p = subparsers.add_parser(
        "rewrite", help="[vLLM] Test query expansion"
    )
    rw_p.add_argument("--query", type=str, required=True)
    # Search Commands
    ss_p = subparsers.add_parser(
        "search-sparse", help="[Search] BM25-only search"
    )
    ss_p.add_argument("--query", type=str, required=True)
    ss_p.add_argument("--top-k", type=int, default=3)
    sd_p = subparsers.add_parser(
        "search-dense", help="[Search] Dense vector search in PostgreSQL"
    )
    sd_p.add_argument("--query", type=str, required=True)
    sd_p.add_argument("--top-k", type=int, default=3)
    sh_p = subparsers.add_parser(
        "search-hybrid",
        help="[Facade Search] Full RRF Hybrid Search (BM25 + PG Vector)",
    )
    sh_p.add_argument("--query", type=str, required=True)
    sh_p.add_argument("--top-k", type=int, default=3)
    # Ingestion Commands
    ing_f = subparsers.add_parser(
        "ingest-file", help="[Storage] Index single file into PostgreSQL"
    )
    ing_f.add_argument("--file", type=str, required=True)
    ing_d = subparsers.add_parser(
        "ingest-dir",
        help="[Storage] Recursively scan and index directory into PostgreSQL",
    )
    ing_d.add_argument("--dir", type=str, required=True)
    ing_d.add_argument(
        "--exts",
        type=str,
        default=".txt,.md,.json,.log,.py,.csv,.pdf,.docx,.html,.htm",
        help="Comma-separated file extensions",
    )
    subparsers.add_parser("mcp", help="[MCP] Launch FastMCP server")
    args = parser.parse_args()
    engine = RAGEngine(
        db_conn=args.db_conn,
        spacy_model=args.spacy_model,
        vllm_url=args.vllm_url,
        max_workers=args.max_workers,
    )
    if args.command == "chunk":
        engine.chunker.max_chunk_size = args.max_size
        chunks = engine.chunker.chunk_document(
            args.text, doc_id="cli_test", source_file="stdin"
        )
        print(
            json.dumps(
                [c.to_dict() for c in chunks], ensure_ascii=False, indent=2
            )
        )
    elif args.command == "rewrite":
        rewritten = asyncio.run(engine.rewriter.rewrite(args.query))
        print(
            json.dumps(
                {"expanded_queries": rewritten}, ensure_ascii=False, indent=2
            )
        )
    elif args.command == "search-sparse":
        res = engine.search_engine.search_sparse_only(
            args.query, top_k=args.top_k
        )
        print(
            json.dumps(
                {"type": "BM25-Sparse", "results": res},
                ensure_ascii=False,
                indent=2,
            )
        )
    elif args.command == "search-dense":
        res = engine.search_engine.search_dense_only(
            args.query, top_k=args.top_k
        )
        print(
            json.dumps(
                {"type": "Dense-Vector (POSTGRESQL)", "results": res},
                ensure_ascii=False,
                indent=2,
            )
        )
    elif args.command == "search-hybrid":
        queries = asyncio.run(engine.rewriter.rewrite(args.query))
        results = engine.search_engine.search_hybrid_rrf(
            queries, top_k=args.top_k
        )
        print(
            json.dumps(
                {
                    "expanded_queries": queries,
                    "backend": "POSTGRESQL",
                    "results": [asdict(cast(Any, r)) for r in results],
                },
                ensure_ascii=False,
                indent=2,
            )
        )
    elif args.command == "ingest-file":
        info = asyncio.run(engine.ingest_file_async(args.file))
        engine.search_engine.reload_index()
        print(
            json.dumps(
                {"status": "success", "info": info},
                ensure_ascii=False,
                indent=2,
            )
        )
    elif args.command == "ingest-dir":
        x_set = {f".{ext.strip().lstrip('.')}" for ext in args.exts.split(",")}
        info = asyncio.run(
            engine.ingest_directory_async(args.dir, extensions=x_set)
        )
        print(
            json.dumps(
                {"status": "success", "info": info},
                ensure_ascii=False,
                indent=2,
            )
        )
    elif args.command == "mcp":
        try:
            from mcp.server.fastmcp import FastMCP
        except ImportError:
            try:
                from fastmcp import FastMCP
            except ImportError:
                print("FastMCP is missing. Install via: pip install fastmcp")
                return
        mcp = FastMCP("RAGEngine PostgreSQL Service")
        @mcp.tool()

        async def query_rag(query: str) -> str:
            q_list = await engine.rewriter.rewrite(query)
            res_ = engine.search_engine.search_hybrid_rrf(q_list, top_k=3)
            if not res_:
                return "Совпадений не найдено."
            return "\n\n".join(
                [
                    f"[{r.source_file} | RRF: {r.rrf_score}]\n{r.text}"
                    for r in res_
                ]
            )

        mcp.run()

if __name__ == "__main__":
    main()
                

Erstellen Sie eine Datei wrap.py


import logging
import os
from typing import Any, Dict, List, Optional

# Importieren Sie die Original-Engine aus rag.py
import rag as core_rag

logger = logging.getLogger(__name__)

# Zwischenspeicherung von RAGEngine-Instanzen anhand der Verbindungszeichenfolge, um unnötige Initialisierung zu vermeiden
_engines_cache: Dict[str, core_rag.RAGEngine] = {}

def _get_engine(db_conn: str) -> core_rag.RAGEngine:
    if db_conn not in _engines_cache:
        _engines_cache[db_conn] = core_rag.RAGEngine(db_conn=db_conn)
    return _engines_cache[db_conn]

def search(
    query: str,
    db_type: str,
    db_conn: str,
    top_k: int = 5,
    filters: Optional[Dict[str, Any]] = None,
) -> List[Dict[str, Any]]:
    """Dokumentensuche mit Hilfe der hybriden RAGEngine-Engine (BM25 + Dense RRF in PostgreSQL)."""
    try:
        engine = _get_engine(db_conn)
        results = engine.search_engine.search_hybrid_rrf([query], top_k=top_k)
        formatted_docs = []
        for r in results:
            # Anwenden optionaler Metadatenfilter
            if filters:
                skip = False
                for k, v in filters.items():
                    if r.metadata.get(k) != v:
                        skip = True
                        break
                if skip:
                    continue
            formatted_docs.append(
                {
                    "id": r.chunk_id,
                    "doc_id": r.doc_id,
                    "text": r.text,
                    "score": r.rrf_score,
                    "source": r.source_file,
                    "bm25_score": r.bm25_score,
                    "dense_score": r.dense_score,
                    "metadata": r.metadata,
                }
            )
        return formatted_docs
    except Exception as e:
        logger.error(f"Error in wrap.search: {e}", exc_info=True)
        return []

def rerank(
    query: str, documents: List[Dict[str, Any]], threshold: float = 0.7
) -> List[Dict[str, Any]]:
    """Gefundene Dokumente werden anhand eines Relevanzschwellenwerts (Score) gefiltert."""
    return [doc for doc in documents if doc.get("score", 0.0) >= threshold]

def chunk_and_trim(
    documents: List[Dict[str, Any]], max_tokens: int = 2000
) -> List[Dict[str, Any]]:
    """Die Liste der Chunks wird so gekürzt, dass sie dem Kontext-Token-Limit entspricht."""
    trimmed = []
    current_tokens = 0
    for doc in documents:
        text = doc.get("text", "")
        estimated_tokens = len(text.split()) * 1.3
        if current_tokens + estimated_tokens <= max_tokens:
            trimmed.append(doc)
            current_tokens += estimated_tokens
        else:
            break
    return trimmed

def format_context(documents: List[Dict[str, Any]]) -> str:
    """Formatieren einer Liste von Dokumenten in eine einzelne Kontexttextzeichenfolge."""
    if not documents:
        return "Kontext nicht gefunden."
    blocks = []
    for doc in documents:
        source = doc.get("source", "Unknown")
        score = doc.get("score", "N/A")
        text = doc.get("text", "")
        blocks.append(f"--- Quelle: {source} (Score: {score}) ---\n{text}")
    return "\n\n".join(blocks)

def build_system_prompt(
    context: str, custom_instructions: Optional[str] = None
) -> str:
    """Erstellung einer Systemabfrage mit eingebettetem Wissensdatenbankkontext."""
    base_instructions = custom_instructions or (
        "Sie sind eine erfahrene Unternehmensassistentin. Beantworten Sie Benutzerfragen, "
        "ausschließlich basierend auf dem bereitgestellten Wissensbasiskontext. "
        "Wenn der Kontext keine Antwort liefert, machen Sie dies deutlich."
    )
    return f"{base_instructions}\n\n=== WISSENSBASISKONTEXT ===\n{context}"

def extract_sources(documents: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
    """Extrahieren von Metadaten über zitierte Quellen."""
    return [
        {
            "id": doc.get("id"),
            "source": doc.get("source"),
            "score": doc.get("score"),
        }
        for doc in documents
    ]

def get_mcp_context(query: str, options: Dict[str, Any]) -> str:
    """Extraktion dynamischer Kontextinformationen aus dem Model Context Protocol (MCP)."""
    return f"MCP Dynamic Context for query '{query}': PostgreSQL RAGEngine storage is active and operational."

def get_mcp_tools(options: Dict[str, Any]) -> List[Dict[str, Any]]:
    """Abrufen von Beschreibungen der MCP-Funktionen/Tools im OpenAI Tools-Format."""
    return [
        {
            "type": "function",
            "function": {
                "name": "query_system_metrics",
                "description": "Erhalten Sie aktuelle Kennzahlen zum Zustand des RAG-Systems.",
                "parameters": {
                    "type": "object",
                    "properties": {
                        "metric_name": {
                            "type": "string",
                            "description": "Metrikname (Beispiel, total_chunks, db_status)",
                        }
                    },
                    "required": ["metric_name"],
                },
            },
        }
    ]

def execute_mcp_tool(
    name: str, arguments: Dict[str, Any], options: Dict[str, Any]
) -> Dict[str, Any]:
    """Führe einen MCP-Funktionsaufruf aus und gib das Ergebnis zurück."""
    if name == "query_system_metrics":
        metric = arguments.get("metric_name", "db_status")
        db_conn = options.get("db_conn", "postgresql://postgres:postgres@localhost:5432/rag_db")
        try:
            engine = _get_engine(db_conn)
            total_chunks = len(engine.search_engine.chunks_cache)
        except Exception:
            total_chunks = "unknown"
        return {
            "metric": metric,
            "total_chunks_in_cache": total_chunks,
            "status": "operational",
        }
    return {"error": f"Tool '{name}' not found in RAGEngine MCP registry."}
                

Erstellen Sie eine Datei app.py


import json
import logging
import os
import requests
from flask import Flask, Response, jsonify, request, stream_with_context
from openai import OpenAI
import wrap as rag

app = Flask(__name__)
logging.basicConfig(level=logging.INFO)

OPENAI_API_KEY = os.getenv("OPENAI_API_KEY")
if not OPENAI_API_KEY:
    raise ValueError("OPENAI_API_KEY environment variable is required")

client = OpenAI(api_key=OPENAI_API_KEY)

DEFAULT_DB_TYPE = os.getenv("DB_TYPE")
DEFAULT_DB_CONN = os.getenv("DB_CONN")
DEFAULT_LLM_URL = os.getenv("LLM_URL")
MAX_TOOL_ITERATIONS = 5

def extract_user_query(messages: list[dict]) -> str:
    for msg in reversed(messages):
        if msg.get("role") == "user":
            content = msg.get("content", "")
            if isinstance(content, str):
                return content
            elif isinstance(content, list):
                return " ".join([
                    part.get("text", "") for part in content
                    if isinstance(part, dict) and part.get("type") == "text"
                ])
    return ""

@app.route("/health", methods=["GET"])
def health_check():
    return jsonify({"status": "ok"}), 200

@app.route("/v1/chat/completions", methods=["POST"])
def chat_completions():
    payload = request.get_json() or {}
    messages = payload.get("messages", [])
    model = payload.get("model", "gpt-4o-mini")
    temperature = payload.get("temperature", 0.2)
    stream = payload.get("stream", False)
    rag_opts = payload.get("rag_options", {})
    db_type = rag_opts.get("db_type", DEFAULT_DB_TYPE)
    db_conn = rag_opts.get("db_conn", DEFAULT_DB_CONN)
    llm_url = rag_opts.get("llm_url", DEFAULT_LLM_URL)
    top_k = rag_opts.get("top_k", 5)
    score_threshold = rag_opts.get("score_threshold", 0.7)
    max_context_tokens = rag_opts.get("max_context_tokens", 2000)
    filters = rag_opts.get("filters", {})
    system_instructions = rag_opts.get("system_instructions", None)
    enable_query_rewrite = rag_opts.get("rewrite_query", False)
    rewrite_model = rag_opts.get("rewrite_model", "qwen")
    use_mcp = rag_opts.get("use_mcp", False)
    raw_user_query = extract_user_query(messages)
    if not raw_user_query:
        return jsonify({
            "error": {
                "message": "User query is missing in 'messages'",
                "type": "invalid_request_error",
                "code": 400
            }
        }), 400
    try:
        # 1. Query Rewrite
        if enable_query_rewrite:
            endpoint = f"{llm_url.rstrip('/')}/chat/completions"
            rewrite_payload = {
                "model": rewrite_model,
                "temperature": 0.0,
                "messages": [
                    {
                        "role": "system",
                        "content": "Rephrase the input query for optimal vector search. Return ONLY the rewritten query text."
                    },
                    {"role": "user", "content": raw_user_query}
                ]
            }
            resp = requests.post(endpoint, json=rewrite_payload, timeout=15)
            resp.raise_for_status()
            search_query = resp.json()["choices"][0]["message"]["content"].strip()
        else:
            search_query = raw_user_query
        # 2. RAG Pipeline
        raw_docs = rag.search(
            query=search_query,
            db_type=db_type,
            db_conn=db_conn,
            top_k=top_k,
            filters=filters
        )
        reranked_docs = rag.rerank(
            query=search_query,
            documents=raw_docs,
            threshold=score_threshold
        )
        trimmed_chunks = rag.chunk_and_trim(
            documents=reranked_docs,
            max_tokens=max_context_tokens
        )
        context_text = rag.format_context(trimmed_chunks)
        system_prompt = rag.build_system_prompt(
            context=context_text,
            custom_instructions=system_instructions
        )
        sources_metadata = rag.extract_sources(trimmed_chunks)
        # 3. MCP Integration
        openai_tools = None
        mcp_executed_tools = []
        if use_mcp:
            mcp_context = rag.get_mcp_context(query=search_query, options=rag_opts)
            if mcp_context:
                system_prompt = f"{system_prompt}\n\n--- Additional MCP Context ---\n{mcp_context}"
            openai_tools = rag.get_mcp_tools(options=rag_opts)
        # 4. OpenAI Completion Loop
        conversation_history = [m for m in messages if m.get("role") != "system"]
        final_messages = [{"role": "system", "content": system_prompt}] + conversation_history
        completion_kwargs = {
            "model": model,
            "messages": final_messages,
            "temperature": temperature
        }
        if openai_tools:
            completion_kwargs["tools"] = openai_tools
        iterations = 0
        completion = None
        while iterations < MAX_TOOL_ITERATIONS:
            if stream and not (openai_tools and iterations == 0):
                break
            completion = client.chat.completions.create(**completion_kwargs)
            response_message = completion.choices[0].message
            if response_message.tool_calls:
                final_messages.append(response_message.model_dump(exclude_none=True))
                for tool_call in response_message.tool_calls:
                    tool_name = tool_call.function.name
                    tool_args = json.loads(tool_call.function.arguments or "{}")
                    tool_output = rag.execute_mcp_tool(
                        name=tool_name,
                        arguments=tool_args,
                        options=rag_opts
                    )
                    mcp_executed_tools.append({
                        "name": tool_name,
                        "args": tool_args,
                        "output": tool_output
                    })
                    final_messages.append({
                        "role": "tool",
                        "tool_call_id": tool_call.id,
                        "content": json.dumps(tool_output, ensure_ascii=False)
                    })
                completion_kwargs["messages"] = final_messages
                iterations += 1
            else:
                break
        # 5. Output Response
        if stream:
            completion_kwargs["stream"] = True
            response_stream = client.chat.completions.create(**completion_kwargs)
            def generate_sse():
                for chunk in response_stream:
                    yield f"data: {chunk.model_dump_json()}\n\n"
                yield "data: [DONE]\n\n"
            return Response(
                stream_with_context(generate_sse()),
                mimetype="text/event-stream"
            )
        if not completion:
            raise RuntimeError("Failed to receive completion from OpenAI API.")
        response_dict = completion.model_dump()
        response_dict["rag_metadata"] = {
            "db_type": db_type,
            "query_rewritten": enable_query_rewrite,
            "llm_url_used": llm_url if enable_query_rewrite else None,
            "original_query": raw_user_query,
            "search_query_used": search_query,
            "retrieved_chunks_count": len(trimmed_chunks),
            "sources": sources_metadata,
            "mcp": {
                "enabled": use_mcp,
                "executed_tools": mcp_executed_tools
            }
        }
        return jsonify(response_dict), 200
    except Exception as e:
        app.logger.error(f"Execution Error: {str(e)}", exc_info=True)
        return jsonify({
            "error": {
                "message": f"Internal RAG Execution Error: {str(e)}",
                "type": "api_error",
                "code": 500
            }
        }), 500
                

Erstellen Sie eine Datei /etc/systemd/system/rag-api.service


[Unit]
Description=Flask RAG & MCP Proxy API Service
After=network.target

[Service]
User=www-data
Group=www-data
WorkingDirectory=/var/www/rag-api
RuntimeDirectory=rag-api
Environment="PATH=/var/www/rag-api/.venv/bin"
Environment="OPENAI_API_KEY=sk-proj-your-actual-api-key-here"
Environment="DB_TYPE=postgres"
Environment="DB_CONN=postgresql://postgres:pass@10.0.0.7:5432/rag"
Environment="LLM_URL=http://10.0.0.5/v1"
ExecStart=/var/www/rag-api/venv/bin/gunicorn \
    --workers 4 \
    --threads 2 \
    --worker-class gthread \
    --bind unix:/run/rag-api/rag-api.sock \
    --access-logfile /var/log/rag-api/access.log \
    --error-logfile /var/log/rag-api/error.log \
    --umask 007 \
    --timeout 120 \
    app:app

Restart=always
RestartSec=5
StandardOutput=journal
StandardError=journal

[Install]
WantedBy=multi-user.target
                

Erstellen Sie eine Datei /etc/nginx/sites-available/rag-api


server {
    listen 80;
    server_name rag.domain.de;
    client_max_body_size 20M;
    location / {
        proxy_pass http://unix:/run/rag-api/rag-api.sock;
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
        proxy_set_header X-Forwarded-Proto $scheme;
        # Unterstützung SSE (Server-Sent Events / Streaming)
        proxy_http_version 1.1;
        proxy_set_header Connection "";
        proxy_buffering off;
        proxy_cache off;
        chunked_transfer_encoding on;
        # Timeouts für langfristige LLM- und RAG-Generationen
        proxy_read_timeout 180s;
        proxy_connect_timeout 60s;
        proxy_send_timeout 180s;
    }
}
                

Erstellen Sie Protokolldateien, starten Sie den systemd-Dienst neu und überprüfen Sie, ob die Socket-Datei erstellt wurde:


mkdir -p /var/log/rag-api
chown -R www-data:www-data /var/log/rag-api
systemctl daemon-reload
systemctl restart rag-api
systemctl status rag-api
ls -la /run/rag-api/rag-api.sock
                

Wir erheben den vhost


ln -s /etc/nginx/sites-available/rag-api /etc/nginx/sites-enabled/
nginx -t
systemctl reload nginx
                

Verbindung schützen


apt install certbot python3-certbot-nginx -y
certbot --nginx
crontab -e
0 3 * * * /usr/bin/certbot renew --quiet --deploy-hook "systemctl reload nginx"
certbot --nginx --redirect
nginx -t
/etc/init.d/nginx restart
                

Erstellen Sie eine Datei: /etc/logrotate.d/rag-api


/var/log/rag-api/*.log {
    daily
    rotate 14
    missingok
    notifempty
    compress
    delaycompress
    create 0640 www-data www-data
    sharedscripts
    postrotate
        # Senden eines USR1-Signals an Gunicorn, um Protokolldateien wieder zu öffnen, ohne Prozesse neu zu starten
        /bin/systemctl kill -s USR1 rag-api.service > /dev/null 2>&1 || true
    endscript
}
                

Wir lassen die root zurück


exit
                

4. Vervollständigung der Wissensdatenbank CLI

Grundlegende Indizierung einer Textdatei


python rag.py ingest-file --file /home/bogatyrev/perelman_poincare.txt
                

Indizieren einer Datei mithilfe einer benutzerdefinierten PostgreSQL-Verbindungszeichenfolge


python rag.py --db-conn "postgresql://rag_user:xxxxxxxxxxxxx@10.0.0.7:5432/rag_db" ingest-file --file /home/bogatyrev/postgres_tuning.md
                

Indizierung eines Markdown-Dokuments mit dem alternativen Sprachmodell spaCy (NLP)


python rag.py --spacy-model de_core_news_lg ingest-file --file /home/bogatyrev/kubectl_k3s.md
                

Indizierung des Quellcodes (Python-Skripte) in die Wissensbasis


python rag.py ingest-file --file /home/bogatyrev/anthologie.py
                

Indizieren der JSON-Konfigurationsdatei


python rag.py --db-conn "postgresql://rag_user:xxxxxxxxxxxxx@10.0.0.7:5432/rag_db" ingest-file --file config/cluster_topology.json
                

Indizieren eines PDF-Dokuments


python rag.py ingest-file --file /home/bogatyrev/reports/security_audit_2026.pdf
                

Indizierung eines unternehmensweiten MS Word-Dokuments (.docx)


python rag.py ingest-file --file /home/bogatyrev/sla_agreement.docx
                

Rekursives Scannen und Indizieren des gesamten Verzeichnisses


python rag.py ingest-dir --dir /home/bogatyrev/knowledge_base/
                

Rekursive Indizierung mit Filterung nach bestimmten Dateierweiterungen


python rag.py ingest-dir --dir /home/bogatyrev/projects/docs --exts ".md,.txt"
                

Rekursive Indizierung mit Filterung nach bestimmten Dateierweiterungen


python rag.py --max-workers 8 --vllm-url "https://llm.domain.de/v1/chat/completions" ingest-dir --dir ~/enterprise_docs/ --exts ".pdf,.docx,.md"
                

Liebe Leserinnen und Leser, ich freue mich sehr über Ihr Interesse an meinem Blog.

Lieber René aus Belgien (Ihrer IP-Adresse nach zu urteilen), ich beantworte Ihre Frage (psql: error: FATAL: password authentication failed for user "rene"): Ich gehe davon aus, dass Sie PostgreSQL nicht aus dem PostgreSQL-Repository, sondern aus dem Debian11-Repository installiert haben, wo aktuell Version 13 verfügbar ist. Diese Version verwendet einen anderen Passwor­tverschlüsselungs­algorithmus (scram-sha-256):

So lösen Sie das Problem


# Stellen wir den Benutzer und die Authentifizierungsmethode auf die altbewährte Methode md5 um.
sudo vi /etc/postgresql/13/main/pg_hba.conf
    <ESC>i
    host    rag_db      rene        127.0.0.1/32            md5
    host    rag_db      rene        ::1/128                 md5
    <ESC>:wq<ENTER>
sudo -u postgres psql
    SET password_encryption = 'md5';
    ALTER USER rene WITH ENCRYPTED PASSWORD 'xxxxxxxxxxxxxx';
    \q
sudo systemctl restart postgresql
psql -h 127.0.0.1 -U rene -d rag_db

# Was könnte sonst noch passieren, wenn Sie scram-sha-256 weiterhin beibehalten möchten?
sudo vi /etc/postgresql/13/main/postgresql.conf
    <ESC>i
    password_encryption = scram-sha-256
    # Falls die Einstellung abweicht oder auskommentiert ist, ändern Sie sie.
    <ESC>:wq<ENTER>
    # Falls dies der Fall ist, wechseln Sie zu MD5.
    <ESC>:q!<ENTER>
sudo -u postgres psql
    ALTER USER rene WITH ENCRYPTED PASSWORD 'xxxxxxxxxxxxxx';
    \q
sudo systemctl restart postgresql
psql -h 127.0.0.1 -U rene -d rag_db
                    

J'espère vous avoir aidé.

Diese Website verwendet den Browser-Cache für den Offline-Modus und verarbeitet Ihre Daten aus dem Kontaktformular gemäß unserer Datenschutzerklärung.