Módulo 6: Metadata Filtering — el componente que casi nadie implementa primero pero todos terminan necesitando

Cápsula 08: Proyecto integrador — Metadata-Filtered RAG end-to-end

Descripción del proyecto

Este es el cierre del módulo 6 — y la culminación de la arquitectura "estado del arte" de RAG production. Vas a construir un sistema que integra todo: metadata filtering + hybrid search + re-ranking + tests de aislamiento + auditoría. Es el proyecto que va al portfolio o al codebase real.

El sistema implementa pipeline completo (filter → hybrid → rerank), tests automatizados de seguridad (cero data leak entre tenants), benchmark vs baseline para justificar la inversión, y documentación que un Tech Lead pueda revisar para aprobar deploy a producción.

Al finalizar este proyecto vas a tener:

  • ✅ Pipeline RAG completo con metadata filtering integrado
  • ✅ Schema de metadata validado con Pydantic
  • secure_query wrapper obligatorio con TenantContext
  • ✅ Tests automatizados de aislamiento que corren en CI
  • ✅ Eval set propio + benchmark comparativo
  • ✅ Audit logging para compliance
  • ✅ README que justifica la arquitectura para deploy aprobación

Tiempo estimado: 4-6 horas + 1 hora de análisis.


Arquitectura del proyecto

metadata_rag_project/
├── src/
│   ├── schema/
│   │   ├── __init__.py
│   │   ├── document_metadata.py     # Pydantic schemas
│   │   └── tenant_context.py
│   ├── retrievers/
│   │   ├── base.py
│   │   ├── semantic.py
│   │   ├── bm25.py
│   │   └── hybrid_filtered.py        # Pipeline completo
│   ├── security/
│   │   ├── secure_query.py           # Wrapper obligatorio
│   │   └── audit_log.py
│   ├── ingestion/
│   │   └── ingest_pipeline.py        # Validation + indexing
│   ├── evaluation/
│   │   ├── eval_set.py
│   │   ├── metrics.py
│   │   └── ab_test.py
│   └── pipeline.py                   # Top-level pipeline
├── tests/
│   ├── test_isolation.py             # Cross-tenant
│   ├── test_security.py              # Permission escalation
│   └── test_pipeline.py
├── data/
│   ├── corpus/
│   └── golden_set.json
├── benchmarks/
│   ├── run_benchmark.py
│   └── reports/
├── scripts/
│   └── pre_commit_check.py           # Anti-bypass linter
├── .env
├── requirements.txt
└── README.md

Paso 1: Schema de metadata con Pydantic

# src/schema/document_metadata.py
from pydantic import BaseModel, Field
from typing import Literal


class DocumentMetadata(BaseModel):
    """Schema obligatorio para todos los documentos."""

    # Universales
    workspace_id: str = Field(..., description="Aislamiento multi-tenant.")
    doc_id: str
    chunk_index: int = Field(..., ge=0)
    source: str
    created_at: int  # Unix timestamp

    # Categorización
    type: Literal["tutorial", "reference", "faq", "changelog", "alert"]
    visibility: Literal["public", "team", "private", "admin_only"] = "public"

    # Metadata específica del dominio (ajustar según tu caso)
    language: Literal["en", "es", "pt"] = "en"
    tags: list[str] = Field(default_factory=list)

    class Config:
        extra = "forbid"

Paso 2: TenantContext y secure_query

# src/schema/tenant_context.py
from dataclasses import dataclass


class TenantIsolationError(Exception):
    pass


@dataclass(frozen=True)
class TenantContext:
    workspace_id: str
    user_id: str
    user_role: str = "member"
    project_ids: list[str] = None

    def __post_init__(self):
        if not self.workspace_id or not self.user_id:
            raise TenantIsolationError("workspace_id y user_id obligatorios")


# src/security/secure_query.py
def secure_query(
    collection,
    query_text: str,
    tenant: TenantContext,
    additional_filters: dict = None,
    n_results: int = 5,
):
    if not tenant or not tenant.workspace_id:
        raise TenantIsolationError("TenantContext obligatorio")

    # Validar additional_filters no sobrescribe workspace_id
    if additional_filters and "workspace_id" in additional_filters:
        raise TenantIsolationError("No sobrescribir workspace_id")

    where = {"workspace_id": tenant.workspace_id}
    if additional_filters:
        where = {"$and": [where, additional_filters]}

    audit_log_query(tenant, query_text, n_results)

    return collection.query(
        query_texts=[query_text],
        where=where,
        n_results=n_results,
    )

Paso 3: Hybrid filtered retriever

# src/retrievers/hybrid_filtered.py
from concurrent.futures import ThreadPoolExecutor
from sentence_transformers import CrossEncoder
from rank_bm25 import BM25Okapi


class HybridFilteredRetriever:
    """Pipeline completo: filter → hybrid → rerank."""

    def __init__(
        self,
        collection,
        bm25_index: BM25Okapi,
        all_doc_ids: list[str],
        all_metadatas: list[dict],
    ):
        self.collection = collection
        self.bm25_index = bm25_index
        self.all_doc_ids = all_doc_ids
        self.all_metadatas = all_metadatas
        self.reranker = CrossEncoder("cross-encoder/ms-marco-MiniLM-L-12-v2")

    def search(
        self,
        query: str,
        tenant: TenantContext,
        additional_filters: dict = None,
        n_candidates: int = 30,
        n_final: int = 5,
    ) -> list[dict]:

        # 1. Build secure filter
        where = self._build_secure_filter(tenant, additional_filters)

        # 2. Hybrid retrieval en paralelo
        with ThreadPoolExecutor(max_workers=2) as executor:
            sem_future = executor.submit(self._semantic_search, query, where, n_candidates)
            bm25_future = executor.submit(self._bm25_search, query, where, n_candidates)
            sem_ids = sem_future.result()
            bm25_ids = bm25_future.result()

        # 3. RRF fusion
        fused_ids = self._reciprocal_rank_fusion([sem_ids, bm25_ids])

        # 4. Recuperar contenido + rerank
        candidates = self.collection.get(
            ids=fused_ids[:n_candidates],
            include=["documents", "metadatas"],
        )

        # 5. Cross-encoder rerank
        return self._rerank(query, candidates, top_k=n_final)

    def _build_secure_filter(self, tenant, additional):
        if not tenant or not tenant.workspace_id:
            raise TenantIsolationError("workspace_id obligatorio")

        base = {"workspace_id": tenant.workspace_id}
        if additional and "workspace_id" in additional:
            raise TenantIsolationError("No sobrescribir workspace_id")

        return {"$and": [base, additional]} if additional else base

    def _semantic_search(self, query, where, n):
        results = self.collection.query(
            query_texts=[query], where=where, n_results=n
        )
        return results["ids"][0]

    def _bm25_search(self, query, where, n):
        # Obtener IDs allowed por filter
        allowed = set(self.collection.get(where=where, include=[])["ids"])

        query_tokens = query.lower().split()
        scores = self.bm25_index.get_scores(query_tokens)

        scored = [
            (self.all_doc_ids[i], scores[i])
            for i in range(len(scores))
            if self.all_doc_ids[i] in allowed and scores[i] > 0
        ]
        scored.sort(key=lambda x: -x[1])
        return [doc_id for doc_id, _ in scored[:n]]

    def _reciprocal_rank_fusion(self, rankings, k=60):
        from collections import defaultdict
        scores = defaultdict(float)
        for ranking in rankings:
            for rank, doc_id in enumerate(ranking, 1):
                scores[doc_id] += 1.0 / (k + rank)
        return [doc_id for doc_id, _ in sorted(scores.items(), key=lambda x: -x[1])]

    def _rerank(self, query, candidates, top_k):
        if not candidates["documents"]:
            return []

        pairs = [(query, doc) for doc in candidates["documents"]]
        scores = self.reranker.predict(pairs, batch_size=32, show_progress_bar=False)

        ranked = sorted(
            range(len(scores)),
            key=lambda i: -scores[i],
        )[:top_k]

        return [
            {
                "doc_id": candidates["ids"][i],
                "document": candidates["documents"][i],
                "metadata": candidates["metadatas"][i],
                "score": float(scores[i]),
            }
            for i in ranked
        ]

Paso 4: Tests de aislamiento

# tests/test_isolation.py
import pytest
from src.schema.tenant_context import TenantContext, TenantIsolationError
from src.retrievers.hybrid_filtered import HybridFilteredRetriever


@pytest.fixture
def retriever_with_two_tenants():
    """Setup con docs de tenant_a y tenant_b."""
    # ... setup ChromaDB + BM25 con 100 docs por tenant ...
    return retriever


def test_tenant_a_cannot_see_tenant_b_docs(retriever_with_two_tenants):
    ctx = TenantContext(workspace_id="tenant_a", user_id="user_1")
    results = retriever_with_two_tenants.search("test", ctx, n_final=20)

    for r in results:
        assert r["metadata"]["workspace_id"] == "tenant_a", (
            f"LEAK: tenant_a vio {r['metadata']['workspace_id']}"
        )


def test_workspace_id_is_required():
    with pytest.raises(TenantIsolationError):
        TenantContext(workspace_id="", user_id="user_1")


def test_additional_filters_cannot_override_tenant():
    ctx = TenantContext(workspace_id="tenant_a", user_id="user_1")

    with pytest.raises(TenantIsolationError):
        retriever.search(
            "test",
            ctx,
            additional_filters={"workspace_id": "tenant_b"},  # bypass attempt
        )


def test_concurrent_tenants_no_leakage(retriever_with_two_tenants):
    """Queries concurrentes de tenants distintos no se cruzan."""
    import threading

    leaks = []

    def query_as(tenant_name):
        ctx = TenantContext(workspace_id=tenant_name, user_id=f"u_{tenant_name}")
        results = retriever_with_two_tenants.search("test", ctx, n_final=20)
        for r in results:
            if r["metadata"]["workspace_id"] != tenant_name:
                leaks.append((tenant_name, r["metadata"]["workspace_id"]))

    threads = [
        threading.Thread(target=query_as, args=("tenant_a",))
        for _ in range(10)
    ] + [
        threading.Thread(target=query_as, args=("tenant_b",))
        for _ in range(10)
    ]

    for t in threads:
        t.start()
    for t in threads:
        t.join()

    assert not leaks, f"Concurrent leakage: {leaks}"

Paso 5: Eval set + benchmark

# data/golden_set.json
[
  {
    "query": "OAuth2PasswordBearer scopes",
    "tenant_id": "tenant_a",
    "expected_doc_ids": ["a_doc_42"],
    "category": "exact_match"
  },
  {
    "query": "how do I authenticate users",
    "tenant_id": "tenant_a",
    "expected_doc_ids": ["a_doc_103"],
    "category": "conceptual"
  }
]


# benchmarks/run_benchmark.py
def benchmark_pipeline_variants(eval_set):
    """Compara 4 arquitecturas: semantic, +hybrid, +rerank, +filter."""

    metrics = {}

    for arch_name, retriever_fn in [
        ("A_semantic_baseline", baseline_semantic_only),
        ("B_hybrid", hybrid_no_filter),
        ("C_hybrid_rerank", hybrid_with_rerank_no_filter),
        ("D_full_pipeline", pipeline_with_everything),
    ]:
        results = []
        for item in eval_set:
            tenant = TenantContext(workspace_id=item["tenant_id"], user_id="bench")
            res = retriever_fn(item["query"], tenant)
            retrieved_ids = [r["doc_id"] for r in res]

            relevant = set(item["expected_doc_ids"])
            in_top_5 = retrieved_ids[:5]
            precision = sum(1 for d in in_top_5 if d in relevant) / 5
            recall = sum(1 for d in in_top_5 if d in relevant) / max(len(relevant), 1)
            results.append({"precision": precision, "recall": recall})

        metrics[arch_name] = {
            "precision_at_5": sum(r["precision"] for r in results) / len(results),
            "recall_at_5": sum(r["recall"] for r in results) / len(results),
        }

    return metrics

Paso 6: README de la solución

# Metadata-Filtered RAG System

## Arquitectura

Pipeline production-ready: `metadata filter → hybrid retrieval → rerank → LLM`.

Cada componente ataca un problema específico:
- Metadata filter: aislamiento multi-tenant + recencia + categorización
- Hybrid (semantic + BM25 + RRF): cobertura conceptual + exact match
- Cross-encoder rerank: refinar top-K con mejor scoring
- LLM con citas: generation con anti-alucinación

## Resultados del benchmark (eval set de 80 queries)

| Arquitectura | Precision@5 | Recall@5 | Latency p95 | Data Leak Risk |
|--------------|-------------|----------|--------------|----------------|
| A: Semantic baseline | 72% | 65% | 220ms | ALTO |
| B: + Hybrid | 84% | 78% | 320ms | ALTO |
| C: + Rerank | 91% | 84% | 470ms | ALTO |
| D: + Metadata filter (full) | **92%** | **87%** | **310ms** | **CERO** |

**Destacado:** la arquitectura D (con filter) es **más rápida** que C porque busca en
menos vectores, además de eliminar data leak risk.

## Tests de aislamiento

Suite automatizada en `tests/test_isolation.py`:
- ✅ Cross-tenant leakage: 0 detectados
- ✅ Bypass via additional_filters: bloqueado
- ✅ Concurrent queries de tenants distintos: aislados
- ✅ workspace_id obligatorio: enforced

## Compliance

- ✅ Audit logging de cada query (workspace_id, user_id, timestamp, query_hash)
- ✅ Logs estructurados con retention 7 años
- ✅ Tests de aislamiento en CI antes de cada deploy
- ✅ Pre-commit hook bloquea collection.query() directo

## Setup

\`\`\`bash
pip install -r requirements.txt
cp .env.example .env  # llenar OPENAI_API_KEY
python scripts/build_indexes.py
python benchmarks/run_benchmark.py
\`\`\`

## Próximos pasos

- Migración a Pinecone para escalar (M07)
- RAG evaluation continuo (M08)

Checklist de entrega del proyecto

  • Schema con Pydantic validado en ingest pipeline
  • TenantContext + secure_query obligatorios en todo el código
  • HybridFilteredRetriever end-to-end funcional
  • Eval set propio de 50+ queries con tenant_id y expected_doc_ids
  • Tests de aislamiento que corren en CI (5 tests mínimos)
  • Pre-commit hook que bloquea collection.query() directo
  • Audit logging estructurado
  • Benchmark report comparando las 4 arquitecturas
  • README con arquitectura, métricas, compliance
  • Documentación de cómo deployar y monitorear

Extensiones opcionales

Extensión 1: dashboard de métricas

Construir Grafana o similar mostrando:

  • Latencia p50/p95/p99 por tenant
  • Recall rate por categoría de query
  • Audit log search interface
  • Alertas si data leak detectado

Extensión 2: A/B testing en producción

Feature flag para rollout gradual:

if user.workspace_id in ROLLOUT_GROUP:
    pipeline = HybridFilteredPipeline(...)
else:
    pipeline = LegacyPipeline(...)

Comparar métricas entre los dos grupos durante 2 semanas.

Extensión 3: tenant onboarding automatizado

Script que onboard un nuevo tenant:

  • Valida que el schema sigue compliance
  • Indexa los docs iniciales
  • Crea TenantContext template
  • Configura monitoring específico

Extensión 4: pipeline observability

Cada paso del pipeline emite traces (OpenTelemetry):

  • Filter time, hybrid time, rerank time
  • Hit rate de cada componente
  • Análisis de bottlenecks

Trampas y errores comunes en el proyecto

Trampa 1: tests de aislamiento que solo cubren happy path

Asegurarse de testear:

  • Concurrencia (threads de tenants distintos)
  • Bypass attempts (additional_filters con workspace_id)
  • Edge cases (tenant_id vacío, None, valores inválidos)

Trampa 2: ignorar audit log retention

Compliance requiere retention típicamente 7 años. No es opcional.

Trampa 3: BM25 sin filter

Si BM25 busca en todo el corpus aunque vector search esté filtrado, RRF fusiona con docs cross-tenant.

Trampa 4: deploy sin smoke test

Después de deployar, primer paso obligatorio: query real desde la app verificando que devuelve docs del tenant correcto.

Trampa 5: documentation desactualizada

README es para Tech Lead que va a aprobar el deploy. Si documenta arquitectura vieja, falla auditoría.


Resumen del proyecto

Construiste:

  • ✅ Pipeline completo (filter + hybrid + rerank) con security by design
  • ✅ Schema validado con Pydantic
  • ✅ Tests de aislamiento automatizados
  • ✅ Pre-commit hooks anti-bypass
  • ✅ Audit logging para compliance
  • ✅ Benchmark comparativo
  • ✅ README para approval

Mejora típica esperada con full pipeline:

Métrica            Baseline         Full pipeline    Mejora
Precision@5        72%              92%              +20 pts
Recall@5           65%              87%              +22 pts
Latency p95        220ms            310ms            +90ms (mejor que sin filter)
Data leak risk     ALTO             CERO             eliminado

Próximos pasos

Inmediato:

  1. Correr el benchmark completo sobre tu eval set.
  2. Validar tests de aislamiento al 100%.
  3. Review del README con compliance/security.
  4. Deploy gradual con feature flag.

Módulo 7 (Production con Pinecone):

Lo que construiste con ChromaDB se traduce directo a Pinecone:

  • workspace_id filter → namespaces nativos
  • BM25 → Pinecone soporta sparse vectors (BM25 nativo desde 2024)
  • Audit logging → Pinecone tiene logs nativos en planes enterprise
  • Tests de aislamiento → ya los tienes, solo cambia el backend

M07 cubre la migración + cómo escalar a 10M+ vectores manteniendo aislamiento.


Recursos

  1. Pydantic Documentation — Validación
  2. ChromaDB — Metadata Filtering — Reference
  3. SOC2 Compliance Framework — Estándar
  4. OWASP Multi-Tenancy Security
  5. GDPR Article 32 — Security context
  6. Anthropic — Contextual Retrieval

Tiempo estimado: 4-6 horas + 1 hora de análisis Siguiente módulo: Módulo 7 — Production con Pinecone