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:
- Sanitiza inputs: Normalización Unicode, eliminación de caracteres peligrosos, length limits, HTML stripping
- Valida outputs: Schema enforcement con Pydantic, JSON repair, retry strategies, fallback chains
- Filtra contenido: Toxicidad (local), PII detection básica, off-topic detection, policy enforcement
- Aplica guardrails: Topic boundaries, language consistency, confidence calibration, custom business rules
- 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 -
InputSanitizercon normalización Unicode, zero-width removal, HTML strip, length limits -
OutputValidatorcon JSON extraction, repair, Pydantic validation, fallback -
ContentFiltercon toxicity, PII, off-topic detection -
GuardrailChaincon 3+ guardrails, fail-fast, timing -
SanitizationPipelinecon 5 etapas integradas, error handling, audit log - Todos los tests pasan (
pytest tests/ -v) - FastAPI app funcional con endpoint
/chaty/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ía | Puntos | Criterios clave |
|---|---|---|
| InputSanitizer | 15 | Unicode normalization (3), zero-width removal (3), HTML strip (3), length limits (3), configurable (3) |
| OutputValidator | 15 | JSON extraction (3), repair (3), Pydantic validation (3), fallback (3), retry concept (3) |
| ContentFilter | 15 | Toxicity patterns (4), PII detection (4), off-topic detection (4), configurable (3) |
| GuardrailChain | 15 | 3+ guardrails (5), chain execution (3), fail-fast (3), timing (2), configurable (2) |
| Pipeline Integration | 20 | 5 stages integrated (5), error handling per stage (5), audit log (5), fallback strategies (5) |
| Tests | 10 | Sanitizer tests (2), validator tests (2), filter tests (2), guardrail tests (2), pipeline tests (2) |
| Code Quality | 10 | Clean structure (3), typing (2), config centralized (2), no hardcoded values (3) |
Distribución de notas
| Rango | Calificación |
|---|---|
| 90-100 | Excelente — Pipeline production-ready |
| 80-89 | Muy bien — Pipeline sólido con mejoras menores |
| 70-79 | Bien — Cubre lo básico pero necesita más robustez |
| 60-69 | Aceptable — Faltan componentes o profundidad |
| < 60 | Necesita 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ódulo | Cómo se conecta |
|---|---|
| Módulo 5: Secrets Management | Las API keys que usa tu pipeline (OpenAI) se protegen con Vault/KMS en lugar de env vars |
| Módulo 6: PII Protection | Tu PII detection básica (regex) se reemplaza con Presidio para detección enterprise-grade |
| Módulo 7: Security Testing | Testeas tu pipeline con adversarial inputs, pen testing, y red team exercises |
| Módulo 8: Integración | Tu 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
- Pydantic V2 Documentation — Referencia completa para schemas de validación y configuración
- FastAPI Documentation — Framework web para el servidor del pipeline
- OpenAI API Reference — API del LLM que el pipeline wrappea
- Pytest Documentation — Framework de testing para los tests del pipeline
- OWASP LLM05: Improper Output Handling — La vulnerabilidad que el pipeline mitiga
- Python Logging — Referencia para audit logging
Creado: Marzo 2026 Versión: 1.0