Módulo 4: Input & Output Sanitization

6. Pipeline Input → Output Completo

Descripción

En las cápsulas 02-05 construiste cuatro componentes individuales: InputSanitizer, OutputValidator, ContentFilter, y GuardrailChain. Cada uno funciona de forma independiente. Pero en producción, no los ejecutas por separado — los integras en un pipeline unificado que procesa cada request de principio a fin: desde que llega el input del usuario hasta que sale la respuesta sanitizada.

Esta cápsula conecta todas las piezas en un middleware FastAPI que implementa el flujo completo: Input → Sanitize → Validate → LLM → Validate Output → Filter → Sanitize Output → Respond. No es código nuevo — es la integración de todo lo que ya construiste, con error handling en cada etapa, fallback strategies, y logging de performance.

El pipeline que construyes aquí es la versión simplificada del proyecto de la cápsula 08. La diferencia: aquí entiendes la arquitectura y los patrones de integración; allá implementas la versión completa production-ready con tests y rúbrica.


El flujo completo

Request (usuario)
│
├── STAGE 1: Input Sanitization ──────────────── InputSanitizer (Cap 02)
│   ├── Unicode normalization (NFKC)
│   ├── Zero-width character removal
│   ├── HTML/markdown stripping
│   ├── Whitespace normalization
│   └── Length enforcement
│   │
│   └── Si falla → 400 Bad Request con mensaje genérico
│
├── STAGE 2: Input Validation ────────────────── Pydantic (request model)
│   ├── Type checking
│   └── Field constraints
│   │
│   └── Si falla → 422 Validation Error
│
├── STAGE 3: Injection Detection ─────────────── Módulo 3 Pipeline
│   ├── Pattern matching
│   └── Classifier
│   │
│   └── Si falla → 400 con mensaje genérico (no revelar detección)
│
├── STAGE 4: LLM Call ────────────────────────── OpenAI API
│   ├── System prompt + sanitized input
│   └── Temperature, max_tokens config
│   │
│   └── Si falla → 503 Service Unavailable + retry
│
├── STAGE 5: Output Validation ───────────────── OutputValidator (Cap 03)
│   ├── JSON extraction
│   ├── Schema validation (Pydantic)
│   └── Type coercion
│   │
│   └── Si falla → retry LLM (max 2) → fallback response
│
├── STAGE 6: Content Filtering ───────────────── ContentFilter (Cap 04)
│   ├── PII detection
│   ├── Toxicity check
│   ├── Off-topic detection
│   └── Policy enforcement
│   │
│   └── Si falla → safe response (no revelar razón)
│
├── STAGE 7: Guardrails ──────────────────────── GuardrailChain (Cap 05)
│   ├── Language consistency
│   ├── Confidence calibration
│   ├── Topic boundaries
│   └── Custom business rules
│   │
│   └── Si falla → safe response
│
├── STAGE 8: Output Sanitization ─────────────── OutputSanitizer
│   ├── HTML escaping
│   └── Encoding normalization
│
└── STAGE 9: Audit Log ──────────────────────── Logger
    ├── Request metadata
    ├── Pipeline stages passed/failed
    ├── Timing per stage
    └── Flags triggered

Response (usuario)

Implementación del Pipeline como Middleware FastAPI

import time
import uuid
import json
import logging
from dataclasses import dataclass, field
from typing import Optional, Callable, Any
from enum import Enum

from fastapi import FastAPI, Request, Response
from fastapi.responses import JSONResponse
from pydantic import BaseModel, Field

logger = logging.getLogger("sanitization_pipeline")


class PipelineStage(Enum):
    INPUT_SANITIZE = "input_sanitize"
    INPUT_VALIDATE = "input_validate"
    INJECTION_DETECT = "injection_detect"
    LLM_CALL = "llm_call"
    OUTPUT_VALIDATE = "output_validate"
    CONTENT_FILTER = "content_filter"
    GUARDRAILS = "guardrails"
    OUTPUT_SANITIZE = "output_sanitize"


@dataclass
class PipelineContext:
    request_id: str = field(default_factory=lambda: str(uuid.uuid4())[:8])
    user_id: str = "anonymous"
    endpoint: str = ""
    original_input: str = ""
    sanitized_input: Optional[str] = None
    raw_llm_output: Optional[str] = None
    validated_output: Optional[dict] = None
    filtered_output: Optional[dict] = None
    final_output: Optional[dict] = None
    stages_passed: list[str] = field(default_factory=list)
    stages_failed: list[str] = field(default_factory=list)
    flags: list[dict] = field(default_factory=list)
    timings: dict[str, float] = field(default_factory=dict)
    total_time_ms: float = 0.0

    def add_flag(self, stage: str, flag_type: str, detail: str):
        self.flags.append({
            "stage": stage,
            "type": flag_type,
            "detail": detail,
            "timestamp": time.time(),
        })

    def record_timing(self, stage: str, duration_ms: float):
        self.timings[stage] = duration_ms


class PipelineError(Exception):
    def __init__(self, stage: str, status_code: int, message: str, user_message: str):
        self.stage = stage
        self.status_code = status_code
        self.message = message
        self.user_message = user_message


class SanitizationPipeline:
    """Pipeline completo de sanitización input → output."""

    def __init__(
        self,
        input_sanitizer,
        output_validator,
        content_filter,
        guardrail_chain,
        llm_caller: Optional[Callable] = None,
        fallback_response: Optional[dict] = None,
        max_output_retries: int = 2,
    ):
        self.input_sanitizer = input_sanitizer
        self.output_validator = output_validator
        self.content_filter = content_filter
        self.guardrail_chain = guardrail_chain
        self.llm_caller = llm_caller
        self.fallback_response = fallback_response or {
            "answer": "No pude procesar tu solicitud. Intenta reformular tu pregunta.",
            "confidence": 0.0,
        }
        self.max_output_retries = max_output_retries

    async def process(
        self,
        user_input: str,
        context: PipelineContext,
        system_prompt: str = "",
    ) -> dict:
        """Ejecuta el pipeline completo."""

        context.original_input = user_input
        pipeline_start = time.perf_counter()

        # STAGE 1: Input Sanitization
        stage_start = time.perf_counter()
        sanitize_result = self.input_sanitizer.sanitize(user_input)
        context.record_timing(
            PipelineStage.INPUT_SANITIZE.value,
            (time.perf_counter() - stage_start) * 1000,
        )

        if not sanitize_result.passed:
            context.stages_failed.append(PipelineStage.INPUT_SANITIZE.value)
            context.add_flag(
                PipelineStage.INPUT_SANITIZE.value,
                "rejected",
                str(sanitize_result.issues),
            )
            raise PipelineError(
                stage=PipelineStage.INPUT_SANITIZE.value,
                status_code=400,
                message=f"Input rejected: {sanitize_result.issues}",
                user_message="Tu mensaje no pudo ser procesado. Intenta con un texto más corto y sin caracteres especiales.",
            )

        context.sanitized_input = sanitize_result.sanitized
        context.stages_passed.append(PipelineStage.INPUT_SANITIZE.value)

        if sanitize_result.issues:
            context.add_flag(
                PipelineStage.INPUT_SANITIZE.value,
                "cleaned",
                str(sanitize_result.issues),
            )

        # STAGE 4: LLM Call (stages 2-3 handled by FastAPI + M3)
        stage_start = time.perf_counter()
        try:
            raw_output = await self._call_llm(
                context.sanitized_input, system_prompt
            )
            context.raw_llm_output = raw_output
            context.stages_passed.append(PipelineStage.LLM_CALL.value)
        except Exception as e:
            context.stages_failed.append(PipelineStage.LLM_CALL.value)
            logger.error(f"[{context.request_id}] LLM call failed: {e}")
            return self.fallback_response
        finally:
            context.record_timing(
                PipelineStage.LLM_CALL.value,
                (time.perf_counter() - stage_start) * 1000,
            )

        # STAGE 5: Output Validation (with retry)
        stage_start = time.perf_counter()
        validation_result = self.output_validator.validate(raw_output)

        if not validation_result.success:
            for retry in range(self.max_output_retries):
                raw_output = await self._call_llm(
                    context.sanitized_input,
                    system_prompt + "\nIMPORTANTE: Responde SOLO con JSON válido.",
                )
                validation_result = self.output_validator.validate(raw_output)
                if validation_result.success:
                    context.add_flag(
                        PipelineStage.OUTPUT_VALIDATE.value,
                        "retry_succeeded",
                        f"Retry {retry + 1} succeeded",
                    )
                    break

        context.record_timing(
            PipelineStage.OUTPUT_VALIDATE.value,
            (time.perf_counter() - stage_start) * 1000,
        )

        if validation_result.success:
            context.validated_output = validation_result.data
            context.stages_passed.append(PipelineStage.OUTPUT_VALIDATE.value)
        else:
            context.stages_failed.append(PipelineStage.OUTPUT_VALIDATE.value)
            context.add_flag(
                PipelineStage.OUTPUT_VALIDATE.value,
                "validation_failed",
                "Using fallback after retries exhausted",
            )
            return self.fallback_response

        # STAGE 6: Content Filtering
        stage_start = time.perf_counter()
        output_text = json.dumps(context.validated_output)
        filter_result = self.content_filter.filter(output_text)
        context.record_timing(
            PipelineStage.CONTENT_FILTER.value,
            (time.perf_counter() - stage_start) * 1000,
        )

        if filter_result.verdict.value == "blocked":
            context.stages_failed.append(PipelineStage.CONTENT_FILTER.value)
            context.add_flag(
                PipelineStage.CONTENT_FILTER.value,
                "blocked",
                str(filter_result.flags),
            )
            return self.fallback_response
        context.stages_passed.append(PipelineStage.CONTENT_FILTER.value)

        # STAGE 7: Guardrails
        stage_start = time.perf_counter()
        guardrail_result = self.guardrail_chain.run(
            output_text,
            context={
                "user_id": context.user_id,
                "input_language": "es",
                "confidence": context.validated_output.get("confidence", 0.5),
            },
        )
        context.record_timing(
            PipelineStage.GUARDRAILS.value,
            (time.perf_counter() - stage_start) * 1000,
        )

        if not guardrail_result.passed:
            context.stages_failed.append(PipelineStage.GUARDRAILS.value)
            context.add_flag(
                PipelineStage.GUARDRAILS.value,
                "blocked",
                f"Blocked by {guardrail_result.blocked_by}",
            )
            return self.fallback_response
        context.stages_passed.append(PipelineStage.GUARDRAILS.value)

        # STAGE 8: Output Sanitization
        context.final_output = context.validated_output
        context.stages_passed.append(PipelineStage.OUTPUT_SANITIZE.value)

        # STAGE 9: Audit Log
        context.total_time_ms = (time.perf_counter() - pipeline_start) * 1000
        self._audit_log(context)

        return context.final_output

    async def _call_llm(self, user_input: str, system_prompt: str) -> str:
        if self.llm_caller:
            return await self.llm_caller(user_input, system_prompt)

        from openai import AsyncOpenAI
        client = AsyncOpenAI()
        response = await client.chat.completions.create(
            model="gpt-4o-mini",
            messages=[
                {"role": "system", "content": system_prompt},
                {"role": "user", "content": user_input},
            ],
            temperature=0.3,
        )
        return response.choices[0].message.content

    def _audit_log(self, context: PipelineContext):
        log_entry = {
            "request_id": context.request_id,
            "user_id": context.user_id,
            "endpoint": context.endpoint,
            "stages_passed": context.stages_passed,
            "stages_failed": context.stages_failed,
            "flags_count": len(context.flags),
            "timings_ms": context.timings,
            "total_time_ms": round(context.total_time_ms, 2),
        }

        if context.flags:
            log_entry["flags"] = context.flags

        logger.info(json.dumps(log_entry))

FastAPI Integration

Montando el pipeline como middleware en una aplicación FastAPI:

from fastapi import FastAPI, HTTPException
from pydantic import BaseModel, Field

app = FastAPI(title="Sanitized AI API")


class ChatRequest(BaseModel):
    message: str = Field(min_length=1, max_length=10000)
    user_id: str = Field(default="anonymous")


class ChatResponse(BaseModel):
    answer: str
    confidence: float = Field(ge=0.0, le=1.0, default=0.5)


# Inicializar componentes del pipeline
# (usando las clases de las cápsulas anteriores)

def create_pipeline() -> SanitizationPipeline:
    """Factory que crea el pipeline con configuración default."""

    # Componentes de cápsulas anteriores (simplificados para ejemplo)
    from dataclasses import dataclass

    @dataclass
    class SimpleSanitizer:
        max_length: int = 4000

        def sanitize(self, text):
            import unicodedata
            cleaned = unicodedata.normalize("NFKC", text.strip())
            issues = []
            if cleaned != text:
                issues.append("normalized")
            if len(cleaned) > self.max_length:
                cleaned = cleaned[:self.max_length]
                issues.append("truncated")

            class Result:
                pass
            r = Result()
            r.passed = bool(cleaned)
            r.sanitized = cleaned
            r.issues = issues
            return r

    @dataclass
    class SimpleValidator:
        def validate(self, raw):
            import json, re
            try:
                match = re.search(r"\{[\s\S]*\}", raw)
                if match:
                    data = json.loads(match.group())
                else:
                    data = json.loads(raw)

                class Result:
                    pass
                r = Result()
                r.success = True
                r.data = data
                return r
            except Exception:
                class Result:
                    pass
                r = Result()
                r.success = False
                r.data = None
                return r

    @dataclass
    class SimpleFilter:
        def filter(self, text):
            class Result:
                pass
            r = Result()
            r.verdict = type("V", (), {"value": "clean"})()
            r.flags = []
            return r

    @dataclass
    class SimpleChain:
        def run(self, text, context=None):
            class Result:
                pass
            r = Result()
            r.passed = True
            r.blocked_by = None
            r.results = []
            r.total_time_ms = 0
            return r

    return SanitizationPipeline(
        input_sanitizer=SimpleSanitizer(),
        output_validator=SimpleValidator(),
        content_filter=SimpleFilter(),
        guardrail_chain=SimpleChain(),
        fallback_response={"answer": "No pude procesar tu solicitud.", "confidence": 0.0},
    )


pipeline = create_pipeline()

SYSTEM_PROMPT = """Eres un asistente de servicio al cliente de TechStore.
Responde en JSON con campos: answer (string), confidence (float 0-1).
Solo responde preguntas sobre productos electrónicos."""


@app.post("/chat", response_model=ChatResponse)
async def chat(request: ChatRequest):
    context = PipelineContext(
        user_id=request.user_id,
        endpoint="/chat",
    )

    try:
        result = await pipeline.process(
            user_input=request.message,
            context=context,
            system_prompt=SYSTEM_PROMPT,
        )
        return ChatResponse(**result)

    except PipelineError as e:
        raise HTTPException(
            status_code=e.status_code,
            detail=e.user_message,
        )
    except Exception as e:
        logger.error(f"Unexpected error: {e}")
        raise HTTPException(
            status_code=500,
            detail="Error interno. Intenta de nuevo.",
        )


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

Error Handling por etapa

Cada etapa del pipeline tiene una estrategia de error diferente:

ERROR_STRATEGIES = {
    PipelineStage.INPUT_SANITIZE: {
        "on_failure": "reject",
        "http_status": 400,
        "user_message": "Tu mensaje no pudo ser procesado.",
        "reveal_reason": False,
        "log_level": "warning",
    },
    PipelineStage.INJECTION_DETECT: {
        "on_failure": "reject",
        "http_status": 400,
        "user_message": "Tu mensaje no pudo ser procesado.",
        "reveal_reason": False,
        "log_level": "warning",
    },
    PipelineStage.LLM_CALL: {
        "on_failure": "fallback",
        "http_status": 200,
        "user_message": None,
        "reveal_reason": False,
        "log_level": "error",
    },
    PipelineStage.OUTPUT_VALIDATE: {
        "on_failure": "retry_then_fallback",
        "max_retries": 2,
        "http_status": 200,
        "user_message": None,
        "reveal_reason": False,
        "log_level": "warning",
    },
    PipelineStage.CONTENT_FILTER: {
        "on_failure": "fallback",
        "http_status": 200,
        "user_message": None,
        "reveal_reason": False,
        "log_level": "warning",
    },
    PipelineStage.GUARDRAILS: {
        "on_failure": "fallback",
        "http_status": 200,
        "user_message": None,
        "reveal_reason": False,
        "log_level": "info",
    },
}

print("Error Handling Strategies:")
print(f"{'Stage':<22} {'On Failure':<22} {'HTTP Status':<14} {'Reveal?':<10}")
print("-" * 68)
for stage, strategy in ERROR_STRATEGIES.items():
    print(
        f"{stage.value:<22} "
        f"{strategy['on_failure']:<22} "
        f"{strategy['http_status']:<14} "
        f"{'Yes' if strategy['reveal_reason'] else 'No':<10}"
    )

# Output esperado:
# Error Handling Strategies:
# Stage                  On Failure             HTTP Status    Reveal?
# --------------------------------------------------------------------
# input_sanitize         reject                 400            No
# injection_detect       reject                 400            No
# llm_call               fallback               200            No
# output_validate        retry_then_fallback    200            No
# content_filter         fallback               200            No
# guardrails             fallback               200            No

Principio clave: nunca reveles información de seguridad

Observa que reveal_reason es False para todas las etapas. Esto es intencional:

  • ❌ "Tu mensaje fue bloqueado por el detector de injection" → Le dice al atacante que tienes un detector
  • ❌ "Output bloqueado por contenido tóxico" → Le dice al atacante cómo evadir
  • ❌ "PII detectada en output" → Confirma que tienes PII detection
  • ✅ "No pude procesar tu solicitud. Intenta reformular tu pregunta." → Genérico, no da pistas

Fallback Strategies

from dataclasses import dataclass
from typing import Optional


@dataclass
class FallbackConfig:
    strategy: str  # "static", "cached", "degraded", "error"
    static_response: Optional[dict] = None
    cache_ttl_seconds: int = 300
    degraded_system_prompt: Optional[str] = None


FALLBACK_STRATEGIES = {
    "static": FallbackConfig(
        strategy="static",
        static_response={
            "answer": "No pude procesar tu solicitud. Intenta reformular tu pregunta.",
            "confidence": 0.0,
        },
    ),
    "cached": FallbackConfig(
        strategy="cached",
        cache_ttl_seconds=300,
    ),
    "degraded": FallbackConfig(
        strategy="degraded",
        degraded_system_prompt="Responde de forma muy breve y conservadora.",
    ),
    "error": FallbackConfig(
        strategy="error",
    ),
}


class FallbackHandler:
    def __init__(self, config: FallbackConfig):
        self.config = config
        self.cache: dict[str, dict] = {}

    def get_fallback(self, original_input: str, stage_failed: str) -> dict:
        if self.config.strategy == "static":
            return self.config.static_response

        if self.config.strategy == "cached":
            cache_key = hash(original_input[:100])
            if cache_key in self.cache:
                return self.cache[cache_key]
            return self.config.static_response or {
                "answer": "Servicio temporalmente no disponible.",
                "confidence": 0.0,
            }

        if self.config.strategy == "error":
            return {
                "error": True,
                "answer": "Error procesando la solicitud.",
                "confidence": 0.0,
            }

        return {"answer": "Fallback genérico.", "confidence": 0.0}


handler = FallbackHandler(FALLBACK_STRATEGIES["static"])
fallback = handler.get_fallback("¿Cuánto cuesta el iPhone?", "content_filter")
print(f"Fallback: {fallback}")

# Output esperado:
# Fallback: {'answer': 'No pude procesar tu solicitud. Intenta reformular tu pregunta.', 'confidence': 0.0}

Performance Optimization

import time

PERFORMANCE_BUDGET = {
    "target_total_ms": 2000,
    "budget_by_stage": {
        PipelineStage.INPUT_SANITIZE: {"target_ms": 5, "max_ms": 20},
        PipelineStage.INJECTION_DETECT: {"target_ms": 10, "max_ms": 50},
        PipelineStage.LLM_CALL: {"target_ms": 1500, "max_ms": 5000},
        PipelineStage.OUTPUT_VALIDATE: {"target_ms": 5, "max_ms": 20},
        PipelineStage.CONTENT_FILTER: {"target_ms": 10, "max_ms": 300},
        PipelineStage.GUARDRAILS: {"target_ms": 10, "max_ms": 50},
        PipelineStage.OUTPUT_SANITIZE: {"target_ms": 2, "max_ms": 10},
    },
}


def analyze_pipeline_performance(context: PipelineContext) -> dict:
    """Analiza el performance del pipeline contra el budget."""
    analysis = {
        "total_ms": context.total_time_ms,
        "within_budget": context.total_time_ms <= PERFORMANCE_BUDGET["target_total_ms"],
        "stages": {},
    }

    for stage_name, timing in context.timings.items():
        try:
            stage = PipelineStage(stage_name)
        except ValueError:
            continue
        budget = PERFORMANCE_BUDGET["budget_by_stage"].get(stage, {})
        target = budget.get("target_ms", 0)
        max_ms = budget.get("max_ms", 0)

        analysis["stages"][stage_name] = {
            "actual_ms": round(timing, 2),
            "target_ms": target,
            "max_ms": max_ms,
            "status": (
                "ok" if timing <= target
                else "warning" if timing <= max_ms
                else "critical"
            ),
        }

    return analysis


# Ejemplo con timings simulados
mock_context = PipelineContext()
mock_context.total_time_ms = 1850
mock_context.timings = {
    "input_sanitize": 3.2,
    "llm_call": 1600,
    "output_validate": 4.1,
    "content_filter": 220,
    "guardrails": 8.5,
}

analysis = analyze_pipeline_performance(mock_context)
print(f"Total: {analysis['total_ms']:.0f}ms (within budget: {analysis['within_budget']})")
for stage, info in analysis["stages"].items():
    print(f"  {stage}: {info['actual_ms']}ms [{info['status']}] (target: {info['target_ms']}ms)")

# Output esperado:
# Total: 1850ms (within budget: True)
#   input_sanitize: 3.2ms [ok] (target: 5ms)
#   llm_call: 1600ms [warning] (target: 1500ms)
#   output_validate: 4.1ms [ok] (target: 5ms)
#   content_filter: 220ms [warning] (target: 10ms)
#   guardrails: 8.5ms [ok] (target: 10ms)

Request/Response Lifecycle completo

from datetime import datetime, timezone


def trace_request_lifecycle(context: PipelineContext) -> str:
    """Genera un trace visual del lifecycle de una request."""
    lines = [
        f"=== Request Lifecycle: {context.request_id} ===",
        f"User: {context.user_id}",
        f"Endpoint: {context.endpoint}",
        f"Time: {datetime.now(timezone.utc).isoformat()}",
        "",
    ]

    all_stages = [
        PipelineStage.INPUT_SANITIZE,
        PipelineStage.INPUT_VALIDATE,
        PipelineStage.INJECTION_DETECT,
        PipelineStage.LLM_CALL,
        PipelineStage.OUTPUT_VALIDATE,
        PipelineStage.CONTENT_FILTER,
        PipelineStage.GUARDRAILS,
        PipelineStage.OUTPUT_SANITIZE,
    ]

    for stage in all_stages:
        name = stage.value
        timing = context.timings.get(name, 0)

        if name in context.stages_passed:
            status = "✅ PASS"
        elif name in context.stages_failed:
            status = "❌ FAIL"
        else:
            status = "⏭️ SKIP"

        lines.append(f"  {status} {name:<20} {timing:.1f}ms")

    lines.append("")
    lines.append(f"  Total: {context.total_time_ms:.1f}ms")
    lines.append(f"  Flags: {len(context.flags)}")

    if context.flags:
        lines.append("  Flag details:")
        for flag in context.flags:
            lines.append(f"    - [{flag['stage']}] {flag['type']}: {flag['detail'][:60]}")

    return "\n".join(lines)


# Simulación
demo_context = PipelineContext(
    request_id="abc12345",
    user_id="user-42",
    endpoint="/chat",
)
demo_context.stages_passed = [
    "input_sanitize", "llm_call", "output_validate",
    "content_filter", "guardrails", "output_sanitize",
]
demo_context.timings = {
    "input_sanitize": 2.5,
    "llm_call": 1200,
    "output_validate": 3.8,
    "content_filter": 15.2,
    "guardrails": 5.1,
}
demo_context.total_time_ms = 1226.6
demo_context.flags = [
    {"stage": "input_sanitize", "type": "cleaned", "detail": "['Unicode normalized (NFKC)']", "timestamp": 0},
]

print(trace_request_lifecycle(demo_context))

# Output esperado:
# === Request Lifecycle: abc12345 ===
# User: user-42
# Endpoint: /chat
# ...
#   ✅ PASS input_sanitize       2.5ms
#   ⏭️ SKIP input_validate       0.0ms
#   ⏭️ SKIP injection_detect     0.0ms
#   ✅ PASS llm_call              1200.0ms
#   ✅ PASS output_validate       3.8ms
#   ✅ PASS content_filter        15.2ms
#   ✅ PASS guardrails            5.1ms
#   ✅ PASS output_sanitize       0.0ms
#
#   Total: 1226.6ms
#   Flags: 1

Troubleshooting

Problema 1: "El pipeline agrega demasiada latencia total"

Cuando todas las etapas se ejecutan secuencialmente, la latencia se acumula.

Solución: Identifica las etapas más lentas (usualmente LLM call y Moderation API). Las etapas de sanitización local (<10ms) no son el problema. Para la Moderation API, considera ejecutarla en paralelo con los guardrails si son independientes. El LLM call es inherentemente serial — optimiza con streaming (cápsula 07).

Problema 2: "Los retries de output validation duplican el costo"

Cada retry es una llamada adicional al LLM con el costo de tokens correspondiente.

Solución: Limita retries a 2. Usa temperatura decreciente (0.7 → 0.3 → 0.0). Monitorea tu retry rate — si es > 10%, el problema es tu prompt, no la validación. Mejora el system prompt antes de agregar retries.

Problema 3: "El fallback response es aburrido y genérico"

Siempre retornar "No pude procesar tu solicitud" es una mala experiencia de usuario.

Solución: Personaliza los fallbacks según la etapa que falló y el contexto:

CONTEXTUAL_FALLBACKS = {
    "output_validate": "Déjame intentar responder de otra forma: {simplified_answer}",
    "content_filter": "Tengo información sobre ese tema pero necesito reformularla. ¿Podrías especificar más tu pregunta?",
    "guardrails": "Esa pregunta está fuera de mi especialidad. ¿Puedo ayudarte con productos o soporte técnico?",
}

Problema 4: "Los logs del audit trail son demasiado verbosos"

Con miles de requests, los logs completos del pipeline llenan el disco rápido.

Solución: Loggea siempre: request_id, stages_passed/failed, total_time_ms, flag_count. Loggea condicionalmente: flags details (solo si hay flags), timings (solo si total > budget), raw output (solo si validation failed).


Ejercicios

Ejercicio 1: Pipeline con circuit breaker

Implementa un circuit breaker que deshabilite temporalmente el LLM call si falla más de N veces consecutivas.

Ver solución
import time


class CircuitBreaker:
    def __init__(self, failure_threshold: int = 5, reset_timeout: int = 60):
        self.failure_threshold = failure_threshold
        self.reset_timeout = reset_timeout
        self.failure_count = 0
        self.last_failure_time = 0.0
        self.state = "closed"

    def can_execute(self) -> bool:
        if self.state == "closed":
            return True
        if self.state == "open":
            if time.time() - self.last_failure_time > self.reset_timeout:
                self.state = "half-open"
                return True
            return False
        return True  # half-open

    def record_success(self):
        self.failure_count = 0
        self.state = "closed"

    def record_failure(self):
        self.failure_count += 1
        self.last_failure_time = time.time()
        if self.failure_count >= self.failure_threshold:
            self.state = "open"


cb = CircuitBreaker(failure_threshold=3, reset_timeout=30)

for i in range(5):
    if cb.can_execute():
        cb.record_failure()
        print(f"Attempt {i+1}: failed (state: {cb.state}, failures: {cb.failure_count})")
    else:
        print(f"Attempt {i+1}: circuit open, using fallback")

# Output esperado:
# Attempt 1: failed (state: closed, failures: 1)
# Attempt 2: failed (state: closed, failures: 2)
# Attempt 3: failed (state: open, failures: 3)
# Attempt 4: circuit open, using fallback
# Attempt 5: circuit open, using fallback

Explicación: El circuit breaker evita que fallos cascading del LLM saturen tu sistema. Si el LLM falla 5 veces seguidas, el circuit breaker se abre y retorna fallback responses sin llamar al LLM, reduciendo latencia y costos.

Ejercicio 2: Pipeline metrics dashboard

Crea un endpoint /metrics que reporte estadísticas del pipeline.

Ver solución
from collections import Counter, defaultdict

class PipelineMetrics:
    def __init__(self):
        self.total_requests = 0
        self.stage_pass_counts = Counter()
        self.stage_fail_counts = Counter()
        self.flag_counts = Counter()
        self.total_times: list[float] = []

    def record(self, context: PipelineContext):
        self.total_requests += 1
        for stage in context.stages_passed:
            self.stage_pass_counts[stage] += 1
        for stage in context.stages_failed:
            self.stage_fail_counts[stage] += 1
        for flag in context.flags:
            self.flag_counts[flag["type"]] += 1
        self.total_times.append(context.total_time_ms)

    def summary(self) -> dict:
        avg_time = sum(self.total_times) / max(len(self.total_times), 1)
        return {
            "total_requests": self.total_requests,
            "avg_time_ms": round(avg_time, 1),
            "stage_pass_rates": {
                stage: count / max(self.total_requests, 1)
                for stage, count in self.stage_pass_counts.items()
            },
            "top_flags": dict(self.flag_counts.most_common(5)),
        }


metrics = PipelineMetrics()
# In production: metrics.record(context) after each request
# @app.get("/metrics")
# async def get_metrics():
#     return metrics.summary()

print("Metrics endpoint ready at /metrics")

Explicación: Las métricas te dicen: cuántos requests pasan cada etapa, cuál es la latencia promedio, y cuáles flags se activan más. Esto guía las decisiones de calibración.

Ejercicio 3: A/B testing de configuraciones de pipeline

Implementa un sistema que permita ejecutar dos configuraciones de pipeline en paralelo y comparar resultados.

Ver solución
import random

class ABTestPipeline:
    def __init__(self, pipeline_a, pipeline_b, split_ratio: float = 0.5):
        self.pipeline_a = pipeline_a
        self.pipeline_b = pipeline_b
        self.split_ratio = split_ratio
        self.results_a: list[dict] = []
        self.results_b: list[dict] = []

    async def process(self, user_input: str, context: PipelineContext) -> dict:
        variant = "A" if random.random() < self.split_ratio else "B"
        pipeline = self.pipeline_a if variant == "A" else self.pipeline_b

        result = await pipeline.process(user_input, context, system_prompt="")

        record = {
            "variant": variant,
            "passed": len(context.stages_failed) == 0,
            "flags": len(context.flags),
            "time_ms": context.total_time_ms,
        }

        if variant == "A":
            self.results_a.append(record)
        else:
            self.results_b.append(record)

        return result

    def compare(self) -> dict:
        def stats(results):
            if not results:
                return {"n": 0}
            pass_rate = sum(1 for r in results if r["passed"]) / len(results)
            avg_time = sum(r["time_ms"] for r in results) / len(results)
            avg_flags = sum(r["flags"] for r in results) / len(results)
            return {"n": len(results), "pass_rate": pass_rate, "avg_time": avg_time, "avg_flags": avg_flags}

        return {"A": stats(self.results_a), "B": stats(self.results_b)}


print("A/B test framework ready. Use ABTestPipeline(strict_pipeline, permissive_pipeline)")

Explicación: El A/B testing te permite comparar una configuración estricta vs permisiva en producción. Si la variante estricta bloquea 20% de requests legítimas, sabes que necesitas ajustar.

Ejercicio 4: Pipeline con timeout por etapa

Implementa timeouts individuales para cada etapa del pipeline.

Ver solución
import asyncio

STAGE_TIMEOUTS = {
    "input_sanitize": 1.0,
    "llm_call": 10.0,
    "output_validate": 2.0,
    "content_filter": 3.0,
    "guardrails": 2.0,
}


async def execute_with_timeout(
    stage_name: str, coro, fallback_value=None
) -> tuple[Any, bool]:
    timeout = STAGE_TIMEOUTS.get(stage_name, 5.0)
    try:
        result = await asyncio.wait_for(coro, timeout=timeout)
        return result, True
    except asyncio.TimeoutError:
        logger.warning(f"Stage {stage_name} timed out after {timeout}s")
        return fallback_value, False


print("Stage timeouts configured:")
for stage, timeout in STAGE_TIMEOUTS.items():
    print(f"  {stage}: {timeout}s")

# Output esperado:
# Stage timeouts configured:
#   input_sanitize: 1.0s
#   llm_call: 10.0s
#   output_validate: 2.0s
#   content_filter: 3.0s
#   guardrails: 2.0s

Explicación: Los timeouts por etapa previenen que una etapa lenta bloquee todo el pipeline. Si la Moderation API tarda más de 3 segundos, el timeout la salta y el pipeline continúa con los guardrails locales.


Resumen

  • 🔑 El pipeline completo integra InputSanitizer → OutputValidator → ContentFilter → GuardrailChain en un flujo secuencial con error handling por etapa
  • 🔑 Cada etapa tiene una estrategia de error diferente: las etapas de input rechazan (400), las etapas de output usan fallback (200 con respuesta segura)
  • 🔑 Nunca reveles al usuario por qué se bloqueó un input o output — los mensajes genéricos previenen que los atacantes hagan fingerprinting de tus defensas
  • 🔑 Las fallback strategies deben ser contextuales: respuesta estática, respuesta cacheada, modo degradado, o error — según la etapa que falló y el contexto
  • 🔑 El performance budget distribuye los ~2 segundos totales: ~1.5s para LLM, ~200ms para Moderation API, <50ms para sanitización/validación/guardrails locales
  • 🔑 El audit log registra cada request con: stages passed/failed, flags, timings — es la base para monitoring y calibración
  • 🔑 El circuit breaker protege contra fallos cascading del LLM: después de N fallos consecutivos, retorna fallback sin llamar al LLM
  • 🔑 Este pipeline es la versión simplificada del proyecto de la cápsula 08 — allá implementas la versión completa con tests y rúbrica

Recursos adicionales

  1. FastAPI Middleware — Documentación oficial de middleware en FastAPI para integrar el pipeline
  2. Circuit Breaker Pattern — Patrón de resiliencia para manejar fallos en servicios externos
  3. OWASP LLM05: Improper Output Handling — La vulnerabilidad que el pipeline completo mitiga
  4. Structured Logging with Python — Guía de logging estructurado para audit trails
  5. OpenTelemetry for Python — Framework de observabilidad para tracing distribuido del pipeline
  6. A/B Testing Best Practices — Principios de A/B testing aplicables a configuraciones de pipeline
  7. Resilience Patterns — Patrones de resiliencia: retry, circuit breaker, bulkhead, timeout

Creado: Marzo 2026 Versión: 1.0