Módulo 8: Unified AI Client — Proyecto integrador final

Monitoring y métricas

Tu cliente ya hace fallback y routing. No sabes qué está haciendo. No sabes cuántas requests están yendo a Mistral vs OpenAI, no sabes cuánto te está costando, no sabes qué provider falla más. Sin esos datos, tu cliente es una caja negra.

En esta cápsula agregamos un MetricsCollector que registra cada request y expone métricas agregadas. Es la diferencia entre código "funciona en mi máquina" y código que operas en producción.

Al terminar vas a poder:

  • Registrar cada request: provider usado, latencia, costo, error si hubo
  • Exponer métricas agregadas: counters por provider, P50/P95 latencia, costo total
  • Exportar snapshots a JSON o Prometheus para integración con tu stack
  • Resetear métricas para tests o reportes por ventana de tiempo

Modelo mental: collector como observador pasivo

El MetricsCollector es un colaborador del UnifiedClient. No cambia la lógica de fallback/routing — solo observa y registra.

       ┌─────────────────┐
       │ UnifiedClient   │
       │                 │
       │   chat() ───────┼──> intenta adapter → OK → llama collector.record(...)
       │                 │   ↓ (si error)
       │                 │   intenta siguiente → OK → record(...)
       │                 │
       │  get_metrics()──┼──> collector.summary()
       └─────────────────┘
                │
                ▼
       ┌──────────────────┐
       │ MetricsCollector │
       │  - eventos[]     │
       │  + record(...)   │
       │  + summary()     │
       │  + reset()       │
       └──────────────────┘

El collector vive en memoria del proceso. Para tracking entre procesos (multi-worker, multi-instancia), conectas a backends externos (Prometheus, Datadog) — fuera de scope, pero el design lo permite.


Implementación

Crea unified_ai_client/metrics.py:

# unified_ai_client/metrics.py
import time
import json
from dataclasses import dataclass, field, asdict
from collections import deque


@dataclass
class EventoRequest:
    timestamp: float
    provider: str
    duration_ms: int
    tokens_input: int
    tokens_output: int
    cost_usd: float | None
    exito: bool
    error_tipo: str | None = None  # 'rate_limit', 'auth', 'timeout', 'provider', None


@dataclass
class ResumenMetrics:
    total_requests: int
    total_exitos: int
    total_errores: int
    total_costo_usd: float
    total_tokens_input: int
    total_tokens_output: int
    latencia_p50_ms: int
    latencia_p95_ms: int
    latencia_p99_ms: int
    por_provider: dict
    errores_por_tipo: dict
    ventana_segundos: float


class MetricsCollector:
    """
    Collector en memoria. Mantiene una ventana rolling de los últimos N eventos
    (default 10,000) para que la memoria no crezca sin tope.
    """

    def __init__(self, max_eventos: int = 10_000):
        self.eventos: deque[EventoRequest] = deque(maxlen=max_eventos)
        self.inicio_periodo: float = time.time()

    def record(
        self,
        provider: str,
        duration_ms: int,
        tokens_input: int = 0,
        tokens_output: int = 0,
        cost_usd: float | None = None,
        exito: bool = True,
        error_tipo: str | None = None,
    ) -> None:
        self.eventos.append(EventoRequest(
            timestamp=time.time(),
            provider=provider,
            duration_ms=duration_ms,
            tokens_input=tokens_input,
            tokens_output=tokens_output,
            cost_usd=cost_usd,
            exito=exito,
            error_tipo=error_tipo,
        ))

    def summary(self) -> ResumenMetrics:
        eventos = list(self.eventos)
        exitos = [e for e in eventos if e.exito]
        errores = [e for e in eventos if not e.exito]

        latencias = sorted(e.duration_ms for e in exitos)

        def percentil(p: int) -> int:
            if not latencias:
                return 0
            idx = int(len(latencias) * p / 100)
            return latencias[min(idx, len(latencias) - 1)]

        # Agregado por provider
        por_provider: dict[str, dict] = {}
        for e in eventos:
            if e.provider not in por_provider:
                por_provider[e.provider] = {
                    "requests": 0, "exitos": 0, "errores": 0,
                    "costo_usd": 0.0, "tokens_input": 0, "tokens_output": 0,
                    "duracion_ms_total": 0,
                }
            p = por_provider[e.provider]
            p["requests"] += 1
            if e.exito:
                p["exitos"] += 1
            else:
                p["errores"] += 1
            if e.cost_usd:
                p["costo_usd"] += e.cost_usd
            p["tokens_input"] += e.tokens_input
            p["tokens_output"] += e.tokens_output
            p["duracion_ms_total"] += e.duration_ms

        # Computar duración promedio por provider
        for p in por_provider.values():
            p["duracion_ms_promedio"] = (
                p["duracion_ms_total"] // p["requests"] if p["requests"] else 0
            )

        # Errores por tipo
        errores_por_tipo: dict[str, int] = {}
        for e in errores:
            tipo = e.error_tipo or "unknown"
            errores_por_tipo[tipo] = errores_por_tipo.get(tipo, 0) + 1

        return ResumenMetrics(
            total_requests=len(eventos),
            total_exitos=len(exitos),
            total_errores=len(errores),
            total_costo_usd=round(sum(e.cost_usd for e in eventos if e.cost_usd), 6),
            total_tokens_input=sum(e.tokens_input for e in eventos),
            total_tokens_output=sum(e.tokens_output for e in eventos),
            latencia_p50_ms=percentil(50),
            latencia_p95_ms=percentil(95),
            latencia_p99_ms=percentil(99),
            por_provider=por_provider,
            errores_por_tipo=errores_por_tipo,
            ventana_segundos=time.time() - self.inicio_periodo,
        )

    def reset(self) -> None:
        self.eventos.clear()
        self.inicio_periodo = time.time()

    def to_json(self) -> str:
        return json.dumps(asdict(self.summary()), indent=2)

Integrar en UnifiedClient

Modifica client.py para inyectar el collector y llamarlo en cada intento:

# unified_ai_client/client.py — solo cambios relevantes

from .metrics import MetricsCollector
from .exceptions import RateLimitError, AuthError, TimeoutError, ProviderError


class UnifiedClient:
    def __init__(
        self,
        config: ClientConfig,
        umbral_circuito: int = 3,
        duracion_circuito_s: int = 60,
        prioridad_default: Priority | None = None,
    ):
        # ... resto igual
        self.metrics = MetricsCollector() if config.metrics_enabled else None
        # ...

    def chat_with_messages(
        self,
        messages: list[Message],
        max_tokens: int = 256,
        temperature: float = 0.7,
        use_fallback: bool = True,
        priority: Priority | None = None,
    ) -> ChatResponse:
        prio = priority or self.prioridad_default
        adapters_ordenados = (
            ordenar_por_prioridad(self.adapters, prio) if prio else self.adapters
        )
        if not use_fallback:
            adapters_ordenados = adapters_ordenados[:1]

        errores: dict[str, Exception] = {}
        for adapter in adapters_ordenados:
            if self._circuito_abierto(adapter.name):
                errores[adapter.name] = ProviderError(adapter.name, "Circuit abierto")
                continue
            try:
                response = self._intentar_con_retry(
                    adapter, messages, max_tokens, temperature
                )
                if self.metrics:
                    self.metrics.record(
                        provider=response.provider,
                        duration_ms=response.duration_ms,
                        tokens_input=response.tokens_input,
                        tokens_output=response.tokens_output,
                        cost_usd=response.cost_usd,
                        exito=True,
                    )
                return response
            except (RateLimitError, AuthError, TimeoutError, ProviderError) as e:
                if self.metrics:
                    self.metrics.record(
                        provider=adapter.name,
                        duration_ms=0,
                        exito=False,
                        error_tipo=_tipo_error(e),
                    )
                errores[adapter.name] = e
                self._registrar_fallo(adapter.name)
                continue

        raise AllProvidersFailedError(errores)

    def get_metrics(self):
        if not self.metrics:
            return None
        return self.metrics.summary()

    def reset_metrics(self) -> None:
        if self.metrics:
            self.metrics.reset()


def _tipo_error(e: Exception) -> str:
    if isinstance(e, RateLimitError):
        return "rate_limit"
    if isinstance(e, AuthError):
        return "auth"
    if isinstance(e, TimeoutError):
        return "timeout"
    if isinstance(e, ProviderError):
        return "provider"
    return "unknown"

Verificación: dashboard simple

Crea examples/dashboard.py:

# examples/dashboard.py
import time
from unified_ai_client import UnifiedClient

client = UnifiedClient.from_yaml("examples/clients_fallback.yaml")

# Generar tráfico de muestra
prompts = [
    "Explica REST en una frase",
    "Define microservicios",
    "¿Qué es un container?",
    "Resume HTTP/2",
    "Explica DNS",
]

for i in range(20):
    prompt = prompts[i % len(prompts)]
    priority = "cost-first" if i % 2 == 0 else "quality-first"
    try:
        client.chat(prompt, max_tokens=80, priority=priority)
    except Exception as e:
        print(f"Request {i} falló: {e}")
    time.sleep(0.3)

# Reporte
summary = client.get_metrics()
print("\n=== RESUMEN ===")
print(f"Total requests:      {summary.total_requests}")
print(f"Exitosos:           {summary.total_exitos}")
print(f"Errores:            {summary.total_errores}")
print(f"Costo total:        ${summary.total_costo_usd:.6f}")
print(f"Tokens input:       {summary.total_tokens_input:,}")
print(f"Tokens output:      {summary.total_tokens_output:,}")
print(f"Latencia P50:       {summary.latencia_p50_ms}ms")
print(f"Latencia P95:       {summary.latencia_p95_ms}ms")
print(f"Latencia P99:       {summary.latencia_p99_ms}ms")
print(f"Ventana:            {summary.ventana_segundos:.1f}s")

print("\n=== POR PROVIDER ===")
for provider, stats in summary.por_provider.items():
    print(f"\n  {provider}")
    print(f"    Requests:      {stats['requests']} ({stats['exitos']} ok, {stats['errores']} err)")
    print(f"    Costo:         ${stats['costo_usd']:.6f}")
    print(f"    Duración prom: {stats['duracion_ms_promedio']}ms")

if summary.errores_por_tipo:
    print("\n=== ERRORES POR TIPO ===")
    for tipo, count in summary.errores_por_tipo.items():
        print(f"  {tipo}: {count}")

# Exportar JSON
with open("metrics_snapshot.json", "w") as f:
    f.write(client.metrics.to_json())
print("\n→ Snapshot guardado en metrics_snapshot.json")

Output esperado (depende de tu config y tráfico):

=== RESUMEN ===
Total requests:      20
Exitosos:           20
Errores:            0
Costo total:        $0.000632
Tokens input:       240
Tokens output:      1452
Latencia P50:       1842ms
Latencia P95:       2911ms
Latencia P99:       3045ms

=== POR PROVIDER ===

  ollama-mistral
    Requests:      10 (10 ok, 0 err)
    Costo:         $0.000000
    Duración prom: 4521ms

  openai-mini
    Requests:      10 (10 ok, 0 err)
    Costo:         $0.000632
    Duración prom: 1834ms

Exportar para Prometheus / Datadog

Patrón común: tu app expone un endpoint /metrics que devuelve formato Prometheus. Tu collector traduce a ese formato:

def to_prometheus(self) -> str:
    s = self.summary()
    lines = [
        f"# HELP llm_requests_total Total LLM requests",
        f"# TYPE llm_requests_total counter",
        f"llm_requests_total {s.total_requests}",
        f"# HELP llm_cost_usd_total Total cost in USD",
        f"# TYPE llm_cost_usd_total counter",
        f"llm_cost_usd_total {s.total_costo_usd}",
        f"# HELP llm_latency_p95_ms P95 latency in ms",
        f"# TYPE llm_latency_p95_ms gauge",
        f"llm_latency_p95_ms {s.latencia_p95_ms}",
    ]
    # Por provider
    for provider, stats in s.por_provider.items():
        lines.append(
            f'llm_requests_by_provider{{provider="{provider}"}} {stats["requests"]}'
        )
    return "\n".join(lines)

Si tu app es FastAPI:

@app.get("/metrics")
def metrics():
    return Response(client.metrics.to_prometheus(), media_type="text/plain")

Prometheus scrappea ese endpoint y tienes dashboards gratis.


Patrones avanzados

Pattern 1 — Alertas basadas en métricas

summary = client.get_metrics()
tasa_error = summary.total_errores / max(summary.total_requests, 1)
if tasa_error > 0.05:
    enviar_alerta(f"Error rate alto: {tasa_error:.1%}")

if summary.total_costo_usd > 100:
    enviar_alerta(f"Costo en últimos 10k requests: ${summary.total_costo_usd:.2f}")

Pattern 2 — Reset por hora

import schedule

def snapshot_y_reset():
    s = client.metrics.to_json()
    guardar_a_s3(f"metrics_{datetime.now().isoformat()}.json", s)
    client.reset_metrics()

schedule.every().hour.do(snapshot_y_reset)

Tus métricas se vuelven históricas con resolución horaria.

Pattern 3 — Métricas por tag

Si tu app tiene contexto adicional (user_tier, feature, etc.), extiende el record() para aceptar tags y agrupa por ellos.


Trampas comunes

Trampa 1 — "Memoria crece sin tope." Sin maxlen en el deque, eventos se acumulan hasta crash. Por eso uso deque(maxlen=10_000). Para casos donde necesitas todo el historial, exporta y resetea periódicamente.

Trampa 2 — "Métricas perdidas en restart." Las metrics están en memoria. Restart pierde todo. Si necesitas durabilidad, persiste cada N requests (e.g., a SQLite) o usa backend externo.

Trampa 3 — "Multi-proceso con un solo collector." Si tu app es multi-worker (gunicorn con 4 workers), cada worker tiene su propio collector. Sumas en el dashboard solo lo que vio cada worker. Para vista unificada, exporta a backend compartido.

Trampa 4 — "Cost_usd es null porque no configuré pricing." Sin price_input_per_1m y price_output_per_1m en tu ProviderConfig, cost_usd es None. Tu summary suma como cero. Configura precios para tracking real.

Trampa 5 — "Mido latencia de wall-clock pero quería de provider only." El duration_ms que mide el adapter incluye wall-clock (network + processing del provider). Si quieres separar, necesitas timing más granular (no es trivial).


Ejercicio

Extiende el MetricsCollector para:

  1. Soportar agregación por ventana de tiemposummary_last_minutes(n: int) -> ResumenMetrics que filtra eventos a los últimos n minutos
  2. Agregar un método top_costs_by_provider(n: int) que devuelve los n providers con mayor costo en la ventana actual
  3. Agregar un test pytest que mockea 100 eventos (50 exitosos, 50 con error) y verifica que el summary los cuenta correctamente
Ver solución
# En metrics.py
def summary_last_minutes(self, n: int) -> ResumenMetrics:
    """Igual que summary() pero filtra eventos de los últimos n minutos."""
    corte = time.time() - n * 60
    eventos_filtrados = [e for e in self.eventos if e.timestamp >= corte]
    # Reutiliza la lógica de summary() sobre eventos_filtrados
    # ... (refactor el método principal para aceptar lista filtrada)

def top_costs_by_provider(self, n: int = 5) -> list[tuple[str, float]]:
    summary = self.summary()
    sorted_providers = sorted(
        summary.por_provider.items(),
        key=lambda kv: kv[1]["costo_usd"],
        reverse=True,
    )
    return [(name, stats["costo_usd"]) for name, stats in sorted_providers[:n]]


# Test
def test_collector_cuenta_eventos():
    c = MetricsCollector()
    for _ in range(50):
        c.record(provider="A", duration_ms=100, tokens_input=10, tokens_output=20,
                 cost_usd=0.001, exito=True)
    for _ in range(50):
        c.record(provider="B", duration_ms=0, exito=False, error_tipo="timeout")

    s = c.summary()
    assert s.total_requests == 100
    assert s.total_exitos == 50
    assert s.total_errores == 50
    assert s.errores_por_tipo == {"timeout": 50}
    assert s.por_provider["A"]["exitos"] == 50
    assert s.por_provider["B"]["errores"] == 50
    assert abs(s.total_costo_usd - 0.05) < 1e-9

Resumen

Aprendiste:

  • MetricsCollector como colaborador pasivo del UnifiedClient
  • ✅ Tracking por request: provider, latencia, tokens, costo, exito/error
  • ✅ Summary agregado con percentiles, totales, breakdown por provider y por tipo de error
  • ✅ Export a JSON y a Prometheus para integración con stack existente
  • ✅ Patterns: alertas, reset periódico, métricas por tag
  • ✅ Trade-offs de implementación in-memory vs persistente

Checkpoint: si tu client.get_metrics() devuelve datos reales y puedes hacer un dashboard básico con esos números, estás listo.


Siguiente cápsula

07 — Testing y validation. Tu cliente tiene muchas features ahora: fallback, routing, circuit breaker, métricas. Sin tests, una refactorización rompe algo y no te enteras. Vamos a escribir suite de tests que valida cada feature con mocks y casos edge.


Recursos

  1. Prometheus client for Python — librería oficial.
  2. Langfuse — Self-hosted observability — alternativa específica para LLMs.
  3. Helicone — proxy con observability built-in.
  4. Datadog APM for LLMs — si tu equipo ya tiene Datadog.
  5. OpenTelemetry Python — estándar emergente para tracing distribuido.