Módulo 4: Input & Output Sanitization

8. Proyecto: Sanitization Pipeline

Descripción del proyecto

Este proyecto cierra el Módulo 4 con la implementación de un Sanitization Pipeline completo — un middleware FastAPI reutilizable que integra todas las capas de sanitización, validación, filtrado y guardrails que construiste en las cápsulas 02-07. No es un ejercicio parcial — es el artefacto production-ready que demuestra que tu sistema AI no solo resiste ataques (Módulo 3) sino que mantiene la higiene completa del flujo de datos.

En las cápsulas anteriores construiste las piezas individuales: InputSanitizer (02), OutputValidator (03), ContentFilter (04), GuardrailChain (05), el pipeline integrado (06), y los patrones de producción (07). Ahora todo se consolida en un solo sistema que:

  1. Sanitiza inputs: Normalización Unicode, eliminación de caracteres peligrosos, length limits, HTML stripping
  2. Valida outputs: Schema enforcement con Pydantic, JSON repair, retry strategies, fallback chains
  3. Filtra contenido: Toxicidad (local), PII detection básica, off-topic detection, policy enforcement
  4. Aplica guardrails: Topic boundaries, language consistency, confidence calibration, custom business rules
  5. Registra auditoría: Logging de cada etapa con timings, flags, y métricas

Al terminar vas a tener un directorio con código Python ejecutable, tests unitarios, y un servidor FastAPI que demuestra el pipeline en acción. Este es el cuarto artefacto de la guía y se integra con el Injection Defense Pipeline (M3) en el Módulo 8 (Secured AI System).


Objetivo del proyecto

Construir un middleware FastAPI de sanitización input/output que procese cada request a través de 5 capas (input sanitizer → output validator → content filter → guardrails → audit logger), con error handling por etapa, fallback strategies, y métricas de performance — todo en un paquete reutilizable.


Especificaciones técnicas

Stack

Python >= 3.10
pydantic >= 2.0
fastapi >= 0.100
uvicorn >= 0.20
openai >= 1.0
httpx >= 0.24
bleach >= 6.0

Estructura del entregable

sanitization-pipeline/
├── pipeline/
│   ├── __init__.py
│   ├── sanitizer.py          # InputSanitizer
│   ├── validator.py           # OutputValidator
│   ├── content_filter.py      # ContentFilter
│   ├── guardrails.py          # Guardrails + GuardrailChain
│   ├── pipeline.py            # SanitizationPipeline (integración)
│   └── config.py              # Configuración centralizada
├── app.py                     # FastAPI application
├── tests/
│   ├── test_sanitizer.py
│   ├── test_validator.py
│   ├── test_content_filter.py
│   ├── test_guardrails.py
│   └── test_pipeline.py
├── requirements.txt
└── README.md

Funcionalidades obligatorias

1. InputSanitizer

  • ✅ Normalización Unicode (NFKC)
  • ✅ Eliminación de caracteres de ancho cero (zero-width)
  • ✅ Eliminación de caracteres de control
  • ✅ HTML tag stripping
  • ✅ Normalización de whitespace
  • ✅ Length limits configurables por endpoint
  • ✅ Resultado con action (pass/cleaned/truncated/rejected) e issues

2. OutputValidator

  • ✅ Extracción de JSON de texto libre (con/sin markdown)
  • ✅ Validación contra Pydantic schema
  • ✅ Reparación de JSON parcial (truncado por max_tokens)
  • ✅ Fallback response cuando la validación falla
  • ✅ Resultado con action (valid/coerced/repaired/fallback/failed)

3. ContentFilter

  • ✅ Toxicity detection local (patterns)
  • ✅ PII detection básica (email, phone, SSN, credit card)
  • ✅ Off-topic detection configurable
  • ✅ Content policy engine con reglas custom
  • ✅ Resultado con verdict (clean/flagged/blocked/redacted)

4. GuardrailChain

  • ✅ Mínimo 3 guardrails implementados
  • ✅ Ejecución en cadena con fail-fast
  • ✅ Timing por guardrail
  • ✅ Resultado con passed/blocked_by/warnings

5. SanitizationPipeline

  • ✅ Integración de los 4 componentes en flujo secuencial
  • ✅ Error handling por etapa (reject vs fallback)
  • ✅ Audit logging con request_id, stages, timings, flags
  • ✅ Fallback response configurable
  • ✅ PipelineContext con tracking completo

Código de implementación

Paso 1: Setup del proyecto

mkdir sanitization-pipeline && cd sanitization-pipeline
mkdir pipeline tests

Crea requirements.txt:

pydantic>=2.0
fastapi>=0.100
uvicorn>=0.20
openai>=1.0
httpx>=0.24
bleach>=6.0
pytest>=7.0
pip install -r requirements.txt

Paso 2: Configuración centralizada

Crea pipeline/config.py:

from pydantic import BaseModel, Field
from typing import Optional


class SanitizerConfig(BaseModel):
    max_length: int = 4000
    max_lines: int = 50
    normalize_unicode: bool = True
    strip_html: bool = True
    strip_markdown: bool = False
    remove_zero_width: bool = True
    normalize_whitespace: bool = True
    on_overlength: str = "truncate"


class ValidatorConfig(BaseModel):
    allow_repair: bool = True
    allow_coercion: bool = True
    fallback_response: dict = Field(
        default_factory=lambda: {
            "answer": "No pude procesar tu solicitud. Intenta reformular tu pregunta.",
            "confidence": 0.0,
        }
    )


class ContentFilterConfig(BaseModel):
    check_toxicity: bool = True
    check_pii: bool = True
    check_off_topic: bool = True
    forbidden_topics: list[str] = Field(default_factory=lambda: [
        r"\b(política|elecciones|gobierno)\b",
        r"\b(religión|iglesia)\b",
    ])
    policy_rules: list[dict] = Field(default_factory=list)


class GuardrailConfig(BaseModel):
    max_output_length: int = 2000
    blocked_topics: list[str] = Field(default_factory=list)
    check_language_consistency: bool = True
    low_confidence_threshold: float = 0.5
    disclaimer: str = (
        "\n\n⚠️ Esta respuesta puede no ser completamente precisa. "
        "Verifica con fuentes oficiales."
    )


class PipelineConfig(BaseModel):
    sanitizer: SanitizerConfig = Field(default_factory=SanitizerConfig)
    validator: ValidatorConfig = Field(default_factory=ValidatorConfig)
    content_filter: ContentFilterConfig = Field(default_factory=ContentFilterConfig)
    guardrails: GuardrailConfig = Field(default_factory=GuardrailConfig)
    max_output_retries: int = 2
    enable_audit_log: bool = True
    system_prompt: str = (
        "Eres un asistente de servicio al cliente de TechStore. "
        "Responde en JSON con campos: answer (string), confidence (float 0-1). "
        "Solo respondes preguntas sobre productos electrónicos."
    )

Paso 3: InputSanitizer

Crea pipeline/sanitizer.py:

import re
import unicodedata
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional

try:
    import bleach
    HAS_BLEACH = True
except ImportError:
    HAS_BLEACH = False

from .config import SanitizerConfig


class SanitizationAction(Enum):
    PASS = "pass"
    CLEANED = "cleaned"
    TRUNCATED = "truncated"
    REJECTED = "rejected"


@dataclass
class SanitizationResult:
    original: str
    sanitized: Optional[str]
    action: SanitizationAction
    issues: list[str] = field(default_factory=list)
    metrics: dict = field(default_factory=dict)

    @property
    def passed(self) -> bool:
        return self.action != SanitizationAction.REJECTED


class InputSanitizer:
    ZERO_WIDTH_CHARS = set(
        "\u200b\u200c\u200d\u200e\u200f"
        "\u202a\u202b\u202c\u202d\u202e"
        "\u2060\u2061\u2062\u2063\u2064"
        "\ufeff"
    )
    CONTROL_CHAR_RE = re.compile(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f-\x9f]")
    HTML_TAG_RE = re.compile(r"<[^>]+>")

    def __init__(self, config: Optional[SanitizerConfig] = None):
        self.config = config or SanitizerConfig()

    def sanitize(self, text: str) -> SanitizationResult:
        if not text or not text.strip():
            return SanitizationResult(
                original=text, sanitized=None,
                action=SanitizationAction.REJECTED,
                issues=["Empty or whitespace-only input"],
            )

        issues: list[str] = []
        cleaned = text
        metrics = {"original_length": len(text)}

        if self.config.normalize_unicode:
            normalized = unicodedata.normalize("NFKC", cleaned)
            if normalized != cleaned:
                issues.append("Unicode normalized (NFKC)")
                cleaned = normalized

        if self.config.remove_zero_width:
            zw_count = sum(1 for c in cleaned if c in self.ZERO_WIDTH_CHARS)
            if zw_count > 0:
                cleaned = "".join(c for c in cleaned if c not in self.ZERO_WIDTH_CHARS)
                issues.append(f"Removed {zw_count} zero-width characters")

        control_matches = self.CONTROL_CHAR_RE.findall(cleaned)
        if control_matches:
            cleaned = self.CONTROL_CHAR_RE.sub("", cleaned)
            issues.append(f"Removed {len(control_matches)} control characters")

        if self.config.strip_html:
            html_tags = self.HTML_TAG_RE.findall(cleaned)
            if html_tags:
                if HAS_BLEACH:
                    cleaned = bleach.clean(cleaned, tags=[], strip=True)
                else:
                    cleaned = self.HTML_TAG_RE.sub("", cleaned)
                issues.append(f"Stripped {len(html_tags)} HTML tags")

        if self.config.normalize_whitespace:
            before_len = len(cleaned)
            cleaned = re.sub(r" {2,}", " ", cleaned)
            cleaned = re.sub(r"\n{3,}", "\n\n", cleaned)
            cleaned = cleaned.strip()
            diff = before_len - len(cleaned)
            if diff > 0:
                issues.append(f"Normalized whitespace (saved {diff} chars)")

        metrics["cleaned_length"] = len(cleaned)
        metrics["estimated_tokens"] = len(cleaned) // 4

        if len(cleaned) > self.config.max_length:
            if self.config.on_overlength == "reject":
                return SanitizationResult(
                    original=text, sanitized=None,
                    action=SanitizationAction.REJECTED,
                    issues=[f"Exceeds max length: {len(cleaned)} > {self.config.max_length}"],
                    metrics=metrics,
                )
            cleaned = cleaned[:self.config.max_length]
            issues.append(f"Truncated to {self.config.max_length} chars")
            metrics["truncated"] = True

        action = SanitizationAction.PASS
        if metrics.get("truncated"):
            action = SanitizationAction.TRUNCATED
        elif issues:
            action = SanitizationAction.CLEANED

        return SanitizationResult(
            original=text, sanitized=cleaned,
            action=action, issues=issues, metrics=metrics,
        )

Paso 4: OutputValidator

Crea pipeline/validator.py:

import json
import re
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional, Any

from pydantic import BaseModel, ValidationError

from .config import ValidatorConfig


class ValidationAction(Enum):
    VALID = "valid"
    REPAIRED = "repaired"
    FALLBACK = "fallback"
    FAILED = "failed"


@dataclass
class ValidationResult:
    action: ValidationAction
    data: Optional[dict]
    raw_output: str
    issues: list[str] = field(default_factory=list)

    @property
    def success(self) -> bool:
        return self.action in (ValidationAction.VALID, ValidationAction.REPAIRED)


class OutputValidator:
    JSON_BLOCK_RE = re.compile(r"```(?:json)?\s*([\s\S]*?)```")
    JSON_OBJECT_RE = re.compile(r"\{[\s\S]*\}")

    def __init__(
        self,
        schema: type[BaseModel],
        config: Optional[ValidatorConfig] = None,
    ):
        self.schema = schema
        self.config = config or ValidatorConfig()

    def validate(self, raw_output: str) -> ValidationResult:
        issues: list[str] = []

        extracted = self._extract_json(raw_output)

        if extracted is None:
            repaired = self._repair_json(raw_output)
            if repaired is not None:
                extracted = repaired
                issues.append("JSON repaired from partial output")

        if extracted is None:
            return ValidationResult(
                action=ValidationAction.FALLBACK,
                data=self.config.fallback_response,
                raw_output=raw_output,
                issues=["No JSON found in output, using fallback"],
            )

        try:
            validated = self.schema(**extracted)
            action = ValidationAction.REPAIRED if issues else ValidationAction.VALID
            return ValidationResult(
                action=action,
                data=validated.model_dump(),
                raw_output=raw_output,
                issues=issues,
            )
        except ValidationError as e:
            return ValidationResult(
                action=ValidationAction.FALLBACK,
                data=self.config.fallback_response,
                raw_output=raw_output,
                issues=issues + [f"Validation error: {e.error_count()} errors"],
            )

    def _extract_json(self, text: str) -> Optional[dict]:
        match = self.JSON_BLOCK_RE.search(text)
        if match:
            try:
                return json.loads(match.group(1).strip())
            except json.JSONDecodeError:
                pass

        try:
            return json.loads(text)
        except json.JSONDecodeError:
            pass

        match = self.JSON_OBJECT_RE.search(text)
        if match:
            try:
                return json.loads(match.group())
            except json.JSONDecodeError:
                pass

        return None

    def _repair_json(self, text: str) -> Optional[dict]:
        cleaned = text.strip()
        open_braces = cleaned.count("{")
        close_braces = cleaned.count("}")

        if open_braces > close_braces:
            last_comma = cleaned.rfind(",")
            last_colon = cleaned.rfind(":")
            if last_comma > last_colon:
                cleaned = cleaned[:last_comma]
            cleaned += "}" * (open_braces - close_braces)
            try:
                return json.loads(cleaned)
            except json.JSONDecodeError:
                pass

        return None

Paso 5: ContentFilter

Crea pipeline/content_filter.py:

import re
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional

from .config import ContentFilterConfig


class FilterVerdict(Enum):
    CLEAN = "clean"
    FLAGGED = "flagged"
    BLOCKED = "blocked"


@dataclass
class FilterResult:
    verdict: FilterVerdict
    text: Optional[str]
    flags: list[dict] = field(default_factory=list)


class ContentFilter:
    TOXICITY_PATTERNS = [
        (re.compile(r"\b(idiota|estúpido|imbécil|inútil)\b", re.I), "insult", 0.7),
        (re.compile(r"\b(matar|destruir|explotar|asesinar)\b", re.I), "violence", 0.9),
        (re.compile(r"\b(odio|repugnante|asco)\b", re.I), "hate", 0.6),
    ]

    PII_PATTERNS = [
        (re.compile(r"\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Z|a-z]{2,}\b"), "email"),
        (re.compile(r"\b(?:\+\d{1,3}\s?)?\(?\d{3}\)?[-.\s]?\d{3}[-.\s]?\d{4}\b"), "phone"),
        (re.compile(r"\b\d{3}-\d{2}-\d{4}\b"), "ssn"),
        (re.compile(r"\b(?:\d{4}[-\s]?){3}\d{4}\b"), "credit_card"),
    ]

    def __init__(self, config: Optional[ContentFilterConfig] = None):
        self.config = config or ContentFilterConfig()
        self.forbidden_topic_patterns = [
            re.compile(p, re.IGNORECASE)
            for p in self.config.forbidden_topics
        ]

    def filter(self, text: str) -> FilterResult:
        flags = []
        verdict = FilterVerdict.CLEAN

        if self.config.check_toxicity:
            for pattern, category, severity in self.TOXICITY_PATTERNS:
                matches = pattern.findall(text)
                if matches:
                    flags.append({
                        "layer": "toxicity",
                        "category": category,
                        "severity": severity,
                        "matches": matches,
                    })
                    if severity >= 0.8:
                        verdict = FilterVerdict.BLOCKED

        if self.config.check_pii:
            for pattern, pii_type in self.PII_PATTERNS:
                matches = pattern.findall(text)
                if matches:
                    flags.append({
                        "layer": "pii",
                        "type": pii_type,
                        "count": len(matches),
                    })
                    if verdict != FilterVerdict.BLOCKED:
                        verdict = FilterVerdict.FLAGGED

        if self.config.check_off_topic:
            for pattern in self.forbidden_topic_patterns:
                matches = pattern.findall(text)
                if matches:
                    flags.append({
                        "layer": "off_topic",
                        "matches": matches,
                    })
                    if verdict != FilterVerdict.BLOCKED:
                        verdict = FilterVerdict.FLAGGED

        return FilterResult(
            verdict=verdict,
            text=text if verdict != FilterVerdict.BLOCKED else None,
            flags=flags,
        )

Paso 6: Guardrails

Crea pipeline/guardrails.py:

import re
import time
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional

from .config import GuardrailConfig


class GuardrailAction(Enum):
    PASS = "pass"
    WARN = "warn"
    BLOCK = "block"
    MODIFY = "modify"


@dataclass
class GuardrailResult:
    name: str
    action: GuardrailAction
    message: str = ""
    modified_text: Optional[str] = None
    execution_time_ms: float = 0.0


class Guardrail(ABC):
    def __init__(self, name: str, enabled: bool = True):
        self.name = name
        self.enabled = enabled

    @abstractmethod
    def check(self, text: str, context: Optional[dict] = None) -> GuardrailResult:
        pass

    def execute(self, text: str, context: Optional[dict] = None) -> GuardrailResult:
        if not self.enabled:
            return GuardrailResult(name=self.name, action=GuardrailAction.PASS)
        start = time.perf_counter()
        result = self.check(text, context)
        result.execution_time_ms = (time.perf_counter() - start) * 1000
        return result


class LengthGuardrail(Guardrail):
    def __init__(self, max_length: int = 2000):
        super().__init__("length_check")
        self.max_length = max_length

    def check(self, text, context=None):
        if len(text) > self.max_length:
            return GuardrailResult(
                name=self.name, action=GuardrailAction.MODIFY,
                message=f"Truncated from {len(text)} to {self.max_length}",
                modified_text=text[:self.max_length] + "...",
            )
        return GuardrailResult(name=self.name, action=GuardrailAction.PASS)


class TopicGuardrail(Guardrail):
    def __init__(self, blocked_topics: list[str]):
        super().__init__("topic_boundary")
        self.patterns = [re.compile(t, re.IGNORECASE) for t in blocked_topics]

    def check(self, text, context=None):
        for pattern in self.patterns:
            if pattern.search(text):
                return GuardrailResult(
                    name=self.name, action=GuardrailAction.BLOCK,
                    message=f"Blocked topic: {pattern.pattern}",
                )
        return GuardrailResult(name=self.name, action=GuardrailAction.PASS)


class ConfidenceGuardrail(Guardrail):
    def __init__(self, threshold: float = 0.5, disclaimer: str = ""):
        super().__init__("confidence_calibration")
        self.threshold = threshold
        self.disclaimer = disclaimer or (
            "\n\n⚠️ Esta respuesta puede no ser completamente precisa."
        )

    def check(self, text, context=None):
        confidence = (context or {}).get("confidence", 1.0)
        if confidence < self.threshold:
            return GuardrailResult(
                name=self.name, action=GuardrailAction.MODIFY,
                message=f"Low confidence ({confidence}), adding disclaimer",
                modified_text=text + self.disclaimer,
            )
        return GuardrailResult(name=self.name, action=GuardrailAction.PASS)


@dataclass
class ChainResult:
    passed: bool
    final_text: Optional[str]
    results: list[GuardrailResult] = field(default_factory=list)
    blocked_by: Optional[str] = None
    total_time_ms: float = 0.0


class GuardrailChain:
    def __init__(self, guardrails: list[Guardrail]):
        self.guardrails = guardrails

    def run(self, text: str, context: Optional[dict] = None) -> ChainResult:
        results = []
        current = text
        start = time.perf_counter()

        for guardrail in self.guardrails:
            result = guardrail.execute(current, context)
            results.append(result)

            if result.action == GuardrailAction.BLOCK:
                return ChainResult(
                    passed=False, final_text=None, results=results,
                    blocked_by=guardrail.name,
                    total_time_ms=(time.perf_counter() - start) * 1000,
                )
            if result.action == GuardrailAction.MODIFY and result.modified_text:
                current = result.modified_text

        return ChainResult(
            passed=True, final_text=current, results=results,
            total_time_ms=(time.perf_counter() - start) * 1000,
        )


def create_guardrail_chain(config: GuardrailConfig) -> GuardrailChain:
    guardrails = [
        LengthGuardrail(max_length=config.max_output_length),
        ConfidenceGuardrail(
            threshold=config.low_confidence_threshold,
            disclaimer=config.disclaimer,
        ),
    ]
    if config.blocked_topics:
        guardrails.insert(1, TopicGuardrail(config.blocked_topics))
    return GuardrailChain(guardrails)

Paso 7: Pipeline integrado

Crea pipeline/pipeline.py:

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

from pydantic import BaseModel

from .sanitizer import InputSanitizer, SanitizationAction
from .validator import OutputValidator, ValidationAction
from .content_filter import ContentFilter, FilterVerdict
from .guardrails import GuardrailChain, GuardrailAction
from .config import PipelineConfig

logger = logging.getLogger("sanitization_pipeline")


@dataclass
class PipelineContext:
    request_id: str = field(default_factory=lambda: str(uuid.uuid4())[:8])
    user_id: str = "anonymous"
    endpoint: str = ""
    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


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


class SanitizationPipeline:
    def __init__(
        self,
        input_sanitizer: InputSanitizer,
        output_validator: OutputValidator,
        content_filter: ContentFilter,
        guardrail_chain: GuardrailChain,
        config: PipelineConfig,
        llm_caller: Optional[Callable] = None,
    ):
        self.input_sanitizer = input_sanitizer
        self.output_validator = output_validator
        self.content_filter = content_filter
        self.guardrail_chain = guardrail_chain
        self.config = config
        self.llm_caller = llm_caller

    async def process(
        self,
        user_input: str,
        context: PipelineContext,
    ) -> dict:
        start = time.perf_counter()

        # Stage 1: Input Sanitization
        t = time.perf_counter()
        san_result = self.input_sanitizer.sanitize(user_input)
        context.timings["input_sanitize"] = (time.perf_counter() - t) * 1000

        if not san_result.passed:
            context.stages_failed.append("input_sanitize")
            raise PipelineError(
                "input_sanitize", 400,
                "Tu mensaje no pudo ser procesado. Intenta con un texto más corto.",
            )
        context.stages_passed.append("input_sanitize")
        if san_result.issues:
            context.flags.append({"stage": "input_sanitize", "issues": san_result.issues})

        # Stage 2: LLM Call
        t = time.perf_counter()
        try:
            raw_output = await self._call_llm(san_result.sanitized)
            context.stages_passed.append("llm_call")
        except Exception as e:
            context.stages_failed.append("llm_call")
            logger.error(f"[{context.request_id}] LLM error: {e}")
            context.total_time_ms = (time.perf_counter() - start) * 1000
            return self.config.validator.fallback_response
        finally:
            context.timings["llm_call"] = (time.perf_counter() - t) * 1000

        # Stage 3: Output Validation (with retry)
        t = time.perf_counter()
        val_result = self.output_validator.validate(raw_output)

        if not val_result.success:
            for retry in range(self.config.max_output_retries):
                raw_output = await self._call_llm(
                    san_result.sanitized + "\n\nResponde SOLO con JSON válido."
                )
                val_result = self.output_validator.validate(raw_output)
                if val_result.success:
                    context.flags.append({
                        "stage": "output_validate",
                        "retry": retry + 1,
                    })
                    break

        context.timings["output_validate"] = (time.perf_counter() - t) * 1000

        if val_result.success:
            context.stages_passed.append("output_validate")
        else:
            context.stages_failed.append("output_validate")
            context.total_time_ms = (time.perf_counter() - start) * 1000
            return self.config.validator.fallback_response

        # Stage 4: Content Filter
        t = time.perf_counter()
        output_text = json.dumps(val_result.data, ensure_ascii=False)
        filter_result = self.content_filter.filter(output_text)
        context.timings["content_filter"] = (time.perf_counter() - t) * 1000

        if filter_result.verdict == FilterVerdict.BLOCKED:
            context.stages_failed.append("content_filter")
            context.flags.append({"stage": "content_filter", "flags": filter_result.flags})
            context.total_time_ms = (time.perf_counter() - start) * 1000
            return self.config.validator.fallback_response
        context.stages_passed.append("content_filter")

        # Stage 5: Guardrails
        t = time.perf_counter()
        gr_result = self.guardrail_chain.run(
            output_text,
            context={"confidence": val_result.data.get("confidence", 0.5)},
        )
        context.timings["guardrails"] = (time.perf_counter() - t) * 1000

        if not gr_result.passed:
            context.stages_failed.append("guardrails")
            context.total_time_ms = (time.perf_counter() - start) * 1000
            return self.config.validator.fallback_response
        context.stages_passed.append("guardrails")

        # Audit log
        context.total_time_ms = (time.perf_counter() - start) * 1000
        if self.config.enable_audit_log:
            self._audit_log(context)

        return val_result.data

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

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

    def _audit_log(self, context: PipelineContext):
        entry = {
            "request_id": context.request_id,
            "stages_passed": context.stages_passed,
            "stages_failed": context.stages_failed,
            "flags": len(context.flags),
            "timings": {k: round(v, 2) for k, v in context.timings.items()},
            "total_ms": round(context.total_time_ms, 2),
        }
        logger.info(json.dumps(entry))

Paso 8: pipeline/__init__.py

from .sanitizer import InputSanitizer, SanitizationResult, SanitizationAction
from .validator import OutputValidator, ValidationResult, ValidationAction
from .content_filter import ContentFilter, FilterResult, FilterVerdict
from .guardrails import (
    Guardrail, GuardrailChain, GuardrailAction, GuardrailResult, ChainResult,
    LengthGuardrail, TopicGuardrail, ConfidenceGuardrail,
    create_guardrail_chain,
)
from .pipeline import SanitizationPipeline, PipelineContext, PipelineError
from .config import PipelineConfig

Paso 9: FastAPI Application

Crea app.py:

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

from pipeline import (
    InputSanitizer, OutputValidator, ContentFilter,
    SanitizationPipeline, PipelineContext, PipelineError, PipelineConfig,
    create_guardrail_chain,
)

logging.basicConfig(level=logging.INFO)

app = FastAPI(title="Sanitized AI API", version="1.0")


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


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


config = PipelineConfig()

pipeline = SanitizationPipeline(
    input_sanitizer=InputSanitizer(config.sanitizer),
    output_validator=OutputValidator(ChatResponse, config.validator),
    content_filter=ContentFilter(config.content_filter),
    guardrail_chain=create_guardrail_chain(config.guardrails),
    config=config,
)


@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(request.message, context)
        return ChatResponse(**result)
    except PipelineError as e:
        raise HTTPException(status_code=e.status_code, detail=e.user_message)
    except Exception:
        raise HTTPException(status_code=500, detail="Error interno.")


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

Paso 10: Tests

Crea tests/test_sanitizer.py:

from pipeline.sanitizer import InputSanitizer, SanitizationAction
from pipeline.config import SanitizerConfig


def test_clean_input_passes():
    sanitizer = InputSanitizer()
    result = sanitizer.sanitize("¿Cuánto cuesta el iPhone 15?")
    assert result.passed
    assert result.action == SanitizationAction.PASS


def test_empty_input_rejected():
    sanitizer = InputSanitizer()
    result = sanitizer.sanitize("   ")
    assert not result.passed
    assert result.action == SanitizationAction.REJECTED


def test_zero_width_removed():
    sanitizer = InputSanitizer()
    result = sanitizer.sanitize("Hola\u200b mundo\u200d")
    assert result.passed
    assert result.sanitized == "Hola mundo"
    assert result.action == SanitizationAction.CLEANED


def test_html_stripped():
    sanitizer = InputSanitizer()
    result = sanitizer.sanitize("<b>Hello</b> <script>alert(1)</script>")
    assert result.passed
    assert "<" not in result.sanitized
    assert result.action == SanitizationAction.CLEANED


def test_unicode_normalized():
    sanitizer = InputSanitizer()
    result = sanitizer.sanitize("Hello")
    assert result.passed
    assert result.sanitized == "Hello"


def test_length_truncation():
    sanitizer = InputSanitizer(SanitizerConfig(max_length=10))
    result = sanitizer.sanitize("a" * 100)
    assert result.passed
    assert len(result.sanitized) == 10
    assert result.action == SanitizationAction.TRUNCATED


def test_length_rejection():
    sanitizer = InputSanitizer(SanitizerConfig(max_length=10, on_overlength="reject"))
    result = sanitizer.sanitize("a" * 100)
    assert not result.passed
    assert result.action == SanitizationAction.REJECTED

Crea tests/test_validator.py:

from pydantic import BaseModel, Field
from pipeline.validator import OutputValidator, ValidationAction


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


def test_valid_json():
    validator = OutputValidator(TestSchema)
    result = validator.validate('{"answer": "Paris", "confidence": 0.9}')
    assert result.success
    assert result.action == ValidationAction.VALID


def test_json_in_markdown():
    validator = OutputValidator(TestSchema)
    result = validator.validate('```json\n{"answer": "Paris"}\n```')
    assert result.success


def test_partial_json_repaired():
    validator = OutputValidator(TestSchema)
    result = validator.validate('{"answer": "Paris", "confidence": 0.8')
    assert result.success
    assert result.action == ValidationAction.REPAIRED


def test_no_json_uses_fallback():
    validator = OutputValidator(TestSchema)
    result = validator.validate("No tengo esa información.")
    assert not result.success
    assert result.action == ValidationAction.FALLBACK
    assert result.data is not None


def test_invalid_schema_uses_fallback():
    validator = OutputValidator(TestSchema)
    result = validator.validate('{"answer": "", "confidence": 2.0}')
    assert not result.success
    assert result.action == ValidationAction.FALLBACK

Crea tests/test_content_filter.py:

from pipeline.content_filter import ContentFilter, FilterVerdict


def test_clean_text():
    f = ContentFilter()
    result = f.filter("El iPhone 15 cuesta $799.")
    assert result.verdict == FilterVerdict.CLEAN


def test_toxic_blocked():
    f = ContentFilter()
    result = f.filter("Voy a destruir todo y matar a todos.")
    assert result.verdict == FilterVerdict.BLOCKED


def test_pii_flagged():
    f = ContentFilter()
    result = f.filter("Contacta a juan@email.com.")
    assert result.verdict == FilterVerdict.FLAGGED


def test_off_topic_flagged():
    f = ContentFilter()
    result = f.filter("Las próximas elecciones serán en noviembre.")
    assert result.verdict == FilterVerdict.FLAGGED

Crea tests/test_guardrails.py:

from pipeline.guardrails import (
    LengthGuardrail, TopicGuardrail, ConfidenceGuardrail,
    GuardrailChain, GuardrailAction,
)


def test_length_pass():
    gr = LengthGuardrail(max_length=100)
    result = gr.execute("Short text")
    assert result.action == GuardrailAction.PASS


def test_length_modify():
    gr = LengthGuardrail(max_length=10)
    result = gr.execute("A very long text that exceeds the limit")
    assert result.action == GuardrailAction.MODIFY
    assert len(result.modified_text) < 20


def test_topic_block():
    gr = TopicGuardrail(blocked_topics=[r"\bpolítica\b"])
    result = gr.execute("La política del gobierno...")
    assert result.action == GuardrailAction.BLOCK


def test_confidence_modify():
    gr = ConfidenceGuardrail(threshold=0.5)
    result = gr.execute("Respuesta", context={"confidence": 0.3})
    assert result.action == GuardrailAction.MODIFY
    assert "⚠️" in result.modified_text


def test_chain_pass():
    chain = GuardrailChain([
        LengthGuardrail(100),
        ConfidenceGuardrail(0.5),
    ])
    result = chain.run("Short text", context={"confidence": 0.9})
    assert result.passed


def test_chain_block():
    chain = GuardrailChain([
        LengthGuardrail(100),
        TopicGuardrail([r"\bpolítica\b"]),
    ])
    result = chain.run("La política del gobierno")
    assert not result.passed
    assert result.blocked_by == "topic_boundary"

Crea tests/test_pipeline.py:

import pytest
from pipeline import (
    InputSanitizer, OutputValidator, ContentFilter,
    SanitizationPipeline, PipelineContext, PipelineError, PipelineConfig,
    create_guardrail_chain,
)
from pydantic import BaseModel, Field


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


@pytest.fixture
def pipeline():
    config = PipelineConfig()

    async def mock_llm(user_input: str) -> str:
        return '{"answer": "Test response", "confidence": 0.85}'

    return SanitizationPipeline(
        input_sanitizer=InputSanitizer(config.sanitizer),
        output_validator=OutputValidator(TestResponse, config.validator),
        content_filter=ContentFilter(config.content_filter),
        guardrail_chain=create_guardrail_chain(config.guardrails),
        config=config,
        llm_caller=mock_llm,
    )


@pytest.mark.asyncio
async def test_pipeline_happy_path(pipeline):
    context = PipelineContext()
    result = await pipeline.process("¿Cuánto cuesta el iPhone?", context)
    assert result["answer"] == "Test response"
    assert "input_sanitize" in context.stages_passed
    assert "output_validate" in context.stages_passed


@pytest.mark.asyncio
async def test_pipeline_empty_input_rejected(pipeline):
    context = PipelineContext()
    with pytest.raises(PipelineError) as exc_info:
        await pipeline.process("   ", context)
    assert exc_info.value.status_code == 400


@pytest.mark.asyncio
async def test_pipeline_llm_failure_uses_fallback(pipeline):
    async def failing_llm(user_input):
        raise Exception("API unavailable")

    pipeline.llm_caller = failing_llm
    context = PipelineContext()
    result = await pipeline.process("Test", context)
    assert result["confidence"] == 0.0
    assert "llm_call" in context.stages_failed

Ejecución

Correr los tests

cd sanitization-pipeline
pytest tests/ -v

# Output esperado:
# tests/test_sanitizer.py::test_clean_input_passes PASSED
# tests/test_sanitizer.py::test_empty_input_rejected PASSED
# tests/test_sanitizer.py::test_zero_width_removed PASSED
# tests/test_sanitizer.py::test_html_stripped PASSED
# tests/test_sanitizer.py::test_unicode_normalized PASSED
# tests/test_sanitizer.py::test_length_truncation PASSED
# tests/test_sanitizer.py::test_length_rejection PASSED
# tests/test_validator.py::test_valid_json PASSED
# tests/test_validator.py::test_json_in_markdown PASSED
# tests/test_validator.py::test_partial_json_repaired PASSED
# tests/test_validator.py::test_no_json_uses_fallback PASSED
# tests/test_validator.py::test_invalid_schema_uses_fallback PASSED
# tests/test_content_filter.py::test_clean_text PASSED
# tests/test_content_filter.py::test_toxic_blocked PASSED
# tests/test_content_filter.py::test_pii_flagged PASSED
# tests/test_content_filter.py::test_off_topic_flagged PASSED
# tests/test_guardrails.py::test_length_pass PASSED
# ...
# All tests passed!

Correr el servidor

uvicorn app:app --reload --port 8000

Probar con curl

# Happy path
curl -X POST http://localhost:8000/chat \
  -H "Content-Type: application/json" \
  -d '{"message": "¿Cuánto cuesta el iPhone 15?", "user_id": "test-user"}'

# Input con HTML (se limpia)
curl -X POST http://localhost:8000/chat \
  -H "Content-Type: application/json" \
  -d '{"message": "<script>alert(1)</script> ¿Precio del laptop?", "user_id": "test-user"}'

# Input vacío (rechazado)
curl -X POST http://localhost:8000/chat \
  -H "Content-Type: application/json" \
  -d '{"message": "   ", "user_id": "test-user"}'

# Health check
curl http://localhost:8000/health

Criterios de éxito

Tu proyecto está completo cuando puedas verificar estos puntos:

  • Estructura de directorios con pipeline/, tests/, app.py
  • InputSanitizer con normalización Unicode, zero-width removal, HTML strip, length limits
  • OutputValidator con JSON extraction, repair, Pydantic validation, fallback
  • ContentFilter con toxicity, PII, off-topic detection
  • GuardrailChain con 3+ guardrails, fail-fast, timing
  • SanitizationPipeline con 5 etapas integradas, error handling, audit log
  • Todos los tests pasan (pytest tests/ -v)
  • FastAPI app funcional con endpoint /chat y /health
  • Configuración centralizada en PipelineConfig
  • Mensajes de error genéricos (no revelan detalles de seguridad)

Rúbrica de evaluación

Total: 100 puntos

CategoríaPuntosCriterios clave
InputSanitizer15Unicode normalization (3), zero-width removal (3), HTML strip (3), length limits (3), configurable (3)
OutputValidator15JSON extraction (3), repair (3), Pydantic validation (3), fallback (3), retry concept (3)
ContentFilter15Toxicity patterns (4), PII detection (4), off-topic detection (4), configurable (3)
GuardrailChain153+ guardrails (5), chain execution (3), fail-fast (3), timing (2), configurable (2)
Pipeline Integration205 stages integrated (5), error handling per stage (5), audit log (5), fallback strategies (5)
Tests10Sanitizer tests (2), validator tests (2), filter tests (2), guardrail tests (2), pipeline tests (2)
Code Quality10Clean structure (3), typing (2), config centralized (2), no hardcoded values (3)

Distribución de notas

RangoCalificación
90-100Excelente — Pipeline production-ready
80-89Muy bien — Pipeline sólido con mejoras menores
70-79Bien — Cubre lo básico pero necesita más robustez
60-69Aceptable — Faltan componentes o profundidad
< 60Necesita revisión — Brechas en el pipeline

Errores comunes

1. Hardcodear valores en lugar de usar configuración

❌ if len(text) > 4000:
✅ if len(text) > self.config.max_length:

Cada límite, umbral, y pattern debe venir de la configuración. Esto permite ajustar el pipeline sin cambiar código.

2. Revelar razones de bloqueo al usuario

❌ HTTPException(detail="Input blocked: injection pattern detected")
✅ HTTPException(detail="Tu mensaje no pudo ser procesado.")

Los mensajes detallados le dan información al atacante. Los mensajes genéricos son más seguros.

3. No testear edge cases

Tests solo para happy paths no son suficientes. Testea: inputs vacíos, inputs con solo whitespace, inputs con solo Unicode invisible, outputs sin JSON, outputs con JSON parcial, outputs tóxicos.

4. Ignorar el audit log

Sin audit log, no puedes saber qué pasa en producción. El log debe incluir al menos: request_id, stages passed/failed, flags count, total time.

5. No manejar fallos del LLM

Si la API de OpenAI falla, tu pipeline debe retornar un fallback, no un error 500. El fallback es una respuesta de baja calidad, pero es mejor que un error.

6. Guardrails que no se pueden deshabilitar

Cada guardrail debe tener un flag enabled. En debugging, necesitas deshabilitar guardrails individuales para aislar problemas.

7. Content filter sin configuración por negocio

Los patterns de toxicidad y PII son genéricos. Las policy rules (competidores, precios internos, disclaimers) son específicas de tu negocio. Asegúrate de que las policy rules sean configurables.

8. Pipeline sin fallback en ninguna etapa

Si cada etapa lanza excepciones sin fallback, un solo fallo rompe todo. Las etapas de output (validation, filter, guardrails) deben tener fallback responses.


Conexión con los módulos siguientes

Tu Sanitization Pipeline es el cuarto artefacto. Cuando avances:

MóduloCómo se conecta
Módulo 5: Secrets ManagementLas API keys que usa tu pipeline (OpenAI) se protegen con Vault/KMS en lugar de env vars
Módulo 6: PII ProtectionTu PII detection básica (regex) se reemplaza con Presidio para detección enterprise-grade
Módulo 7: Security TestingTesteas tu pipeline con adversarial inputs, pen testing, y red team exercises
Módulo 8: IntegraciónTu pipeline se integra con Injection Defense (M3) + Secrets (M5) + PII (M6) en el sistema completo

Resumen

  • El Sanitization Pipeline es el artefacto central del Módulo 4 — integra input sanitization, output validation, content filtering, y guardrails en un middleware FastAPI reutilizable
  • 5 componentes trabajan en secuencia: InputSanitizer → OutputValidator → ContentFilter → GuardrailChain → AuditLogger
  • Configuración centralizada con Pydantic permite ajustar el pipeline sin cambiar código
  • Error handling por etapa distingue entre rechazar inputs (400) y usar fallbacks para outputs (200 con respuesta segura)
  • Tests unitarios verifican cada componente de forma aislada y el pipeline integrado
  • El pipeline se integra con el Injection Defense Pipeline del Módulo 3 y se consolida en el Módulo 8

Recursos para el proyecto

  1. Pydantic V2 Documentation — Referencia completa para schemas de validación y configuración
  2. FastAPI Documentation — Framework web para el servidor del pipeline
  3. OpenAI API Reference — API del LLM que el pipeline wrappea
  4. Pytest Documentation — Framework de testing para los tests del pipeline
  5. OWASP LLM05: Improper Output Handling — La vulnerabilidad que el pipeline mitiga
  6. Python Logging — Referencia para audit logging

Creado: Marzo 2026 Versión: 1.0