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:
- Soportar agregación por ventana de tiempo —
summary_last_minutes(n: int) -> ResumenMetricsque filtra eventos a los últimosnminutos - Agregar un método
top_costs_by_provider(n: int)que devuelve losnproviders con mayor costo en la ventana actual - Agregar un test
pytestque 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:
- ✅
MetricsCollectorcomo colaborador pasivo delUnifiedClient - ✅ 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
- Prometheus client for Python — librería oficial.
- Langfuse — Self-hosted observability — alternativa específica para LLMs.
- Helicone — proxy con observability built-in.
- Datadog APM for LLMs — si tu equipo ya tiene Datadog.
- OpenTelemetry Python — estándar emergente para tracing distribuido.