Módulo 8: Proyecto Integrador RAG con ChromaDB

Cápsula 04: Retrieval, Generation y API

Descripción de la cápsula

Vas a tomar el pipeline RAG end-to-end que construiste en M4/11 (función ask() mínima en ~150 líneas) y exponerlo como API REST profesional con FastAPI: endpoints documentados, validación de inputs, manejo de errores, citas estructuradas y trace_id para debugging.

Prerequisito clave (M4/11): ya tienes la lógica core del RAG — query → retrieve top-k → construir prompt con anti-alucinación → generar con GPT → respuesta con fuentes. Esta cápsula no reinventa esa lógica; la encapsula detrás de una API que terceros (frontend, mobile, otros servicios) pueden consumir.

Lo nuevo de esta cápsula es la capa de servicio: contratos de API estables, validación con Pydantic, observabilidad (latency por endpoint, error tracking, request tracing), y los detalles que separan un script Python de un servicio que un equipo puede operar.


Endpoints de la API

MétodoRutaDescripción
POST/ingestDispara ingestion de un directorio (async o sync según diseño)
GET/searchBúsqueda por similitud (query, top_k, filtros opcionales)
POST/askPregunta RAG: retrieval + generation, respuesta con fuentes
GET/collectionsLista colecciones en ChromaDB
DELETE/collections/{name}Borra una colección
GET/healthHealth check para Docker/K8s

Contrato del endpoint /ask

{
  "answer": "Texto generado basado en el contexto recuperado.",
  "sources": [
    {
      "doc_id": "doc_001",
      "source": "/path/to/doc.txt",
      "title": "Documento Ejemplo",
      "score": 0.87
    }
  ],
  "confidence": 0.82,
  "trace_id": "550e8400-e29b-41d4-a716-446655440000",
  "fallback_reason": null
}

Si no hay evidencia suficiente: fallback_reason = "insufficient retrieval", confidence = 0, sources = [].


Estructura del proyecto

project/
├── api/
│   ├── __init__.py
│   ├── main.py          # FastAPI app
│   ├── routes/
│   │   ├── ask.py
│   │   ├── search.py
│   │   └── ingest.py
│   └── services/
│       ├── retrieval.py
│       └── generation.py
├── ingestion/           # (de cápsula 03)
├── config.py
└── run_ingestion.py

Código completo

1. Configuración

# config.py
import os
from dotenv import load_dotenv

load_dotenv()

CHROMA_PATH = os.getenv("CHROMA_PATH", "./chroma_data")
COLLECTION_NAME = os.getenv("COLLECTION_NAME", "rag_docs")
EMBEDDING_MODEL = os.getenv("EMBEDDING_MODEL", "text-embedding-3-small")
LLM_MODEL = os.getenv("LLM_MODEL", "gpt-3.5-turbo")
TOP_K = int(os.getenv("TOP_K", "5"))
SCORE_THRESHOLD = float(os.getenv("SCORE_THRESHOLD", "0.5"))
MAX_QUESTION_LENGTH = int(os.getenv("MAX_QUESTION_LENGTH", "2000"))

2. Servicio de Retrieval

# api/services/retrieval.py
import chromadb
from chromadb.config import Settings
from openai import OpenAI
from typing import List, Optional
import os

client_openai = OpenAI()
chroma_client = chromadb.PersistentClient(path=os.getenv("CHROMA_PATH", "./chroma_data"))


def embed_text(text: str, model: str = "text-embedding-3-small") -> List[float]:
    """Genera embedding para un único texto (query)."""
    response = client_openai.embeddings.create(model=model, input=[text])
    return response.data[0].embedding


def retrieve(
    query: str,
    collection_name: str = "rag_docs",
    top_k: int = 5,
    score_threshold: float = 0.5,
    where: Optional[dict] = None,
) -> dict:
    """
    Retrieval: embed query, similarity search, filtrar por score.
    Returns: documents, metadatas, ids, distances (lower = more similar).
    """
    collection = chroma_client.get_collection(collection_name)
    query_embedding = embed_text(query)

    kwargs = {
        "query_embeddings": [query_embedding],
        "n_results": top_k,
        "include": ["documents", "metadatas", "distances"],
    }
    if where:
        kwargs["where"] = where

    result = collection.query(**kwargs)

    docs = result["documents"][0] if result["documents"] else []
    metas = result["metadatas"][0] if result["metadatas"] else []
    ids = result["ids"][0] if result["ids"] else []
    distances = result["distances"][0] if result.get("distances") else []

    # ChromaDB usa L2 por defecto; convertir a similarity score (1 / (1 + distance))
    # O si usas cosine, distances puede ser cosine distance (1 - similarity)
    scores = [1 / (1 + d) if d is not None else 0 for d in distances]

    # Filtrar por umbral
    filtered = [
        (doc, meta, id_, score)
        for doc, meta, id_, score in zip(docs, metas, ids, scores)
        if score >= score_threshold
    ]
    if not filtered:
        return {"documents": [], "metadatas": [], "ids": [], "scores": []}

    docs_f, metas_f, ids_f, scores_f = zip(*filtered)
    return {
        "documents": list(docs_f),
        "metadatas": list(metas_f),
        "ids": list(ids_f),
        "scores": list(scores_f),
    }

3. Servicio de Generation

# api/services/generation.py
from openai import OpenAI
from typing import List, Optional

client = OpenAI()


def build_context(documents: List[str], metadatas: List[dict]) -> str:
    """Ensambla el contexto para el prompt a partir de los chunks recuperados."""
    parts = []
    for i, (doc, meta) in enumerate(zip(documents, metadatas), 1):
        source = meta.get("source", "unknown")
        title = meta.get("title", meta.get("doc_id", "Doc"))
        parts.append(f"[Fuente {i}: {title}]\n{doc}")
    return "\n\n---\n\n".join(parts)


def generate_answer(
    question: str,
    documents: List[str],
    metadatas: List[dict],
    model: str = "gpt-3.5-turbo",
    temperature: float = 0.2,
) -> str:
    """
    Genera respuesta usando el contexto recuperado.
    Si no hay documentos, devuelve mensaje de fallback sin llamar al LLM.
    """
    if not documents or not metadatas:
        return "No tengo suficiente información en mi base de conocimiento para responder esta pregunta."

    context = build_context(documents, metadatas)
    prompt = f"""Eres un asistente que responde preguntas basándote ÚNICAMENTE en el contexto proporcionado.
No inventes información. Si el contexto no contiene la respuesta, di que no tienes esa información.

Contexto:
{context}

Pregunta: {question}

Responde de forma clara y concisa. Cita las fuentes cuando sea relevante."""

    response = client.chat.completions.create(
        model=model,
        messages=[
            {"role": "system", "content": "Respondes basándote solo en el contexto dado. No inventes datos."},
            {"role": "user", "content": prompt},
        ],
        temperature=temperature,
    )
    return response.choices[0].message.content


def build_sources_response(metadatas: List[dict], scores: List[float]) -> List[dict]:
    """Formatea las fuentes para el contrato de respuesta."""
    return [
        {
            "doc_id": m.get("doc_id", ""),
            "source": m.get("source", ""),
            "title": m.get("title", m.get("doc_id", "Doc")),
            "score": round(s, 2),
        }
        for m, s in zip(metadatas, scores)
    ]

4. Rutas FastAPI

# api/routes/ask.py
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel, Field
from uuid import uuid4
import os

from api.services.retrieval import retrieve
from api.services.generation import generate_answer, build_sources_response

router = APIRouter(prefix="/ask", tags=["RAG"])
MAX_QUESTION_LENGTH = int(os.getenv("MAX_QUESTION_LENGTH", "2000"))


class AskRequest(BaseModel):
    question: str = Field(..., min_length=1, max_length=MAX_QUESTION_LENGTH)


class SourceItem(BaseModel):
    doc_id: str
    source: str
    title: str
    score: float


class AskResponse(BaseModel):
    answer: str
    sources: list[SourceItem]
    confidence: float
    trace_id: str
    fallback_reason: str | None = None


@router.post("", response_model=AskResponse)
def ask(req: AskRequest):
    trace_id = str(uuid4())
    question = req.question.strip()

    # Retrieval
    result = retrieve(
        query=question,
        top_k=int(os.getenv("TOP_K", "5")),
        score_threshold=float(os.getenv("SCORE_THRESHOLD", "0.5")),
    )
    docs = result["documents"]
    metas = result["metadatas"]
    scores = result["scores"]

    # Fallback si no hay evidencia suficiente
    if not docs or not metas:
        return AskResponse(
            answer="No encontré información relevante para responder tu pregunta.",
            sources=[],
            confidence=0.0,
            trace_id=trace_id,
            fallback_reason="insufficient retrieval",
        )

    # Calcular confidence (promedio de scores normalizado)
    confidence = round(sum(scores) / len(scores), 2) if scores else 0.0

    # Generation
    answer = generate_answer(
        question=question,
        documents=docs,
        metadatas=metas,
        model=os.getenv("LLM_MODEL", "gpt-3.5-turbo"),
    )
    sources = build_sources_response(metas, scores)

    return AskResponse(
        answer=answer,
        sources=[SourceItem(**s) for s in sources],
        confidence=confidence,
        trace_id=trace_id,
        fallback_reason=None,
    )

# api/routes/search.py
from fastapi import APIRouter, Query
from pydantic import BaseModel
import os

from api.services.retrieval import retrieve

router = APIRouter(prefix="/search", tags=["Search"])


class SearchResponse(BaseModel):
    documents: list[str]
    metadatas: list[dict]
    ids: list[str]
    scores: list[float]


@router.get("", response_model=SearchResponse)
def search(
    q: str = Query(..., min_length=1),
    top_k: int = Query(5, ge=1, le=20),
    score_threshold: float = Query(0.0, ge=0, le=1),
):
    result = retrieve(
        query=q,
        top_k=top_k,
        score_threshold=score_threshold,
    )
    return SearchResponse(
        documents=result["documents"],
        metadatas=result["metadatas"],
        ids=result["ids"],
        scores=result["scores"],
    )

# api/routes/ingest.py
from fastapi import APIRouter, HTTPException, BackgroundTasks
from pydantic import BaseModel

router = APIRouter(prefix="/ingest", tags=["Ingestion"])


class IngestRequest(BaseModel):
    dir_path: str


@router.post("")
def trigger_ingest(req: IngestRequest, background_tasks: BackgroundTasks):
    """Dispara ingestion asíncrona de un directorio."""
    from pathlib import Path
    if not Path(req.dir_path).exists():
        raise HTTPException(status_code=400, detail="Directorio no encontrado")
    from ingestion.pipeline import ingest_directory
    from ingestion.config import IngestionConfig
    config = IngestionConfig()
    def run():
        ingest_directory(req.dir_path, config=config)
    background_tasks.add_task(run)
    return {"status": "ingestion started", "dir_path": req.dir_path}

# api/routes/collections.py
from fastapi import APIRouter
import chromadb
import os

router = APIRouter(prefix="/collections", tags=["Collections"])
chroma_client = chromadb.PersistentClient(path=os.getenv("CHROMA_PATH", "./chroma_data"))


@router.get("")
def list_collections():
    cols = chroma_client.list_collections()
    return [{"name": c.name, "count": c.count()} for c in cols]


@router.delete("/{name}")
def delete_collection(name: str):
    try:
        chroma_client.delete_collection(name)
        return {"status": "deleted", "collection": name}
    except Exception as e:
        from fastapi import HTTPException
        raise HTTPException(status_code=404, detail=str(e))

# api/main.py
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware

from api.routes.ask import router as ask_router
from api.routes.search import router as search_router
from api.routes.ingest import router as ingest_router
from api.routes.collections import router as collections_router

app = FastAPI(
    title="RAG API",
    description="Sistema RAG con ChromaDB y OpenAI",
    version="1.0.0",
    docs_url="/docs",   # Swagger
    redoc_url="/redoc", # ReDoc
)

app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

app.include_router(ask_router)
app.include_router(search_router)
app.include_router(collections_router)


@app.get("/health")
def health():
    return {"status": "ok"}


@app.get("/")
def root():
    return {
        "message": "RAG API - Usa /docs para Swagger, /redoc para ReDoc",
        "endpoints": ["/ask", "/search", "/ingest", "/collections", "/health"],
    }

5. Ejecutar la API

# Instalar dependencias
pip install fastapi uvicorn chromadb openai python-dotenv pydantic

# Ejecutar
uvicorn api.main:app --reload --host 0.0.0.0 --port 8000

Acceso:


Reglas de calidad de respuesta

ReglaImplementación
No inventar sin evidenciaFallback explícito cuando retrieval vacío o scores bajos
Incluir fuentessources con doc_id, source, title, score
Evitar respuestas sin contextoNo llamar al LLM si no hay chunks; devolver mensaje estándar
Trazabilidadtrace_id en cada respuesta para correlación con logs
Validación de inputmax_length en pregunta para evitar abuse

Ejemplo de uso

POST /ask

curl -X POST "http://localhost:8000/ask" \
  -H "Content-Type: application/json" \
  -d '{"question": "¿Qué es un vector database?"}'

Respuesta esperada:

{
  "answer": "Un vector database es un sistema de almacenamiento optimizado para...",
  "sources": [
    {
      "doc_id": "intro_vectors",
      "source": "/docs/intro.txt",
      "title": "intro",
      "score": 0.87
    }
  ],
  "confidence": 0.85,
  "trace_id": "550e8400-e29b-41d4-a716-446655440000",
  "fallback_reason": null
}

GET /search

curl "http://localhost:8000/search?q=vector%20database&top_k=3"

Ejercicios prácticos

Ejercicio 1: Extender /ask con metadata filtering

Añade un parámetro opcional doc_id o source a la request para filtrar la búsqueda solo en ciertos documentos.

Solución: En AskRequest añade doc_id: str | None = None. En la llamada a retrieve, si doc_id está presente, usa where={"doc_id": doc_id}.


Ejercicio 2: Respuesta "no tengo evidencia" cuando confidence < 0.5

Si el promedio de scores es menor a 0.5, no llames al LLM. Devuelve el mensaje de fallback y fallback_reason="low_confidence".

Solución: Antes de generate_answer, calcula confidence. Si confidence < 0.5, return AskResponse(..., fallback_reason="low_confidence") sin llamar al LLM.


Ejercicio 3: Incluir trace_id en logs

Configura un middleware o dependency que capture el trace_id y lo añada a cada log del request. Usa structlog o logging con contexto.

Solución: Middleware que genera trace_id al inicio del request y lo guarda en request.state.trace_id. En cada endpoint, pasas ese valor a la respuesta. Para logs, usa logging.LoggerAdapter con extra={"trace_id": trace_id}.


Ejercicio 4: Timeout y retry para OpenAI

Añade timeout de 30s a las llamadas de OpenAI y retry (1 intento) si falla por timeout.

Solución: client.chat.completions.create(..., timeout=30). Envuelve en try/except; en Timeout o APIConnectionError, retry una vez con time.sleep(2).


Ejercicio 5: Endpoint POST /ingest

Crea un endpoint que reciba {"dir_path": "/path/to/docs"} y ejecute el pipeline de ingestion de la cápsula 03. Para no bloquear, considera BackgroundTasks de FastAPI.

Solución:

from fastapi import BackgroundTasks
from ingestion.pipeline import ingest_directory

@router.post("/ingest")
def trigger_ingest(payload: dict, bg: BackgroundTasks):
    dir_path = payload.get("dir_path")
    if not dir_path:
        raise HTTPException(400, "dir_path required")
    def run():
        ingest_directory(dir_path)
    bg.add_task(run)
    return {"status": "ingestion started", "dir": dir_path}

Ejercicio 6: Documentación de errores

Define un esquema de error estándar: {"error": str, "trace_id": str, "detail": str}. Usa un exception handler global para devolverlo en 4xx/5xx.

Solución:

from fastapi import Request
from fastapi.responses import JSONResponse

@app.exception_handler(Exception)
def global_exception_handler(request: Request, exc: Exception):
    trace_id = getattr(request.state, "trace_id", "unknown")
    return JSONResponse(
        status_code=500,
        content={"error": "internal_error", "trace_id": trace_id, "detail": str(exc)},
    )

Troubleshooting

"Respuestas correctas pero sin trazabilidad"

Causa: No se incluyen fuentes ni trace_id.

Solución: Asegúrate de que AskResponse siempre incluya sources (aunque vacío) y trace_id. Genera trace_id al inicio del request y pásalo por todo el flujo.


"Respuestas inventadas cuando no hay contexto"

Causa: Se llama al LLM incluso con retrieval vacío o scores muy bajos.

Solución: Añade validación: si len(docs) == 0 o max(scores) < 0.5, no llames a generate_answer. Devuelve fallback explícito con fallback_reason.


"API inestable bajo carga"

Causa: Sin límites, timeouts ni control de concurrencia.

Solución: Limita longitud de pregunta (max_length=2000). Añade timeout a OpenAI. Considera rate limiting con slowapi o similar. Revisa connection pooling de ChromaDB.


"Swagger/ReDoc no muestran schemas correctos"

Causa: Modelos Pydantic sin response_model o con tipos complejos no documentados.

Solución: Usa response_model=AskResponse en el decorator. Define todos los modelos con BaseModel y campos tipados. FastAPI genera los schemas automáticamente.


"Errores 500 sin detalle útil"

Causa: Excepciones no capturadas, logs insuficientes.

Solución: Exception handler global que loguee traceback y devuelva mensaje genérico al cliente (no revelar internals). En desarrollo, debug=True puede ayudar; en producción, logs estructurados con trace_id.


Resumen

  • La API expone /ask, /search, /collections, /health con FastAPI.
  • El endpoint /ask integra retrieval (ChromaDB + OpenAI embeddings) y generation (GPT-3.5-turbo).
  • Contrato de respuesta: answer, sources, confidence, trace_id, fallback_reason.
  • Fallback explícito cuando no hay evidencia suficiente; no se inventan respuestas.
  • Swagger (/docs) y ReDoc (/redoc) documentan la API automáticamente.
  • Separar retrieval y generation en funciones independientes facilita tests y evolución.
  • Validación de input, timeout y manejo de errores mejoran robustez.

Recursos adicionales


Tiempo estimado: 35-40 minutos
Siguiente: 05-testing-evaluacion.md