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_querywrapper obligatorio conTenantContext - ✅ 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:
- Correr el benchmark completo sobre tu eval set.
- Validar tests de aislamiento al 100%.
- Review del README con compliance/security.
- Deploy gradual con feature flag.
Módulo 7 (Production con Pinecone):
Lo que construiste con ChromaDB se traduce directo a Pinecone:
workspace_idfilter → 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
- Pydantic Documentation — Validación
- ChromaDB — Metadata Filtering — Reference
- SOC2 Compliance Framework — Estándar
- OWASP Multi-Tenancy Security
- GDPR Article 32 — Security context
- Anthropic — Contextual Retrieval
Tiempo estimado: 4-6 horas + 1 hora de análisis Siguiente módulo: Módulo 7 — Production con Pinecone