Módulo 7: Flujos Avanzados

Patterns de Producción

Descripción de la cápsula

Los patterns de las cápsulas 02-06 hacen que tu sistema sea robusto: retry con backoff, branching paralelo, subgraphs reutilizables, map-reduce para colecciones, y error handling con fallback. Pero hay una capa más que separa un prototipo funcional de un sistema de producción real: timeouts por nodo, rate limiting interno, logging estructurado, y métricas de rendimiento.

Son los "últimos 10%" que hacen que tu sistema sea operable — que puedas entender qué pasa cuando algo falla a las 3am, que puedas responder "¿por qué esta request tardó 12 segundos?" con datos concretos, y que puedas escalar de 1 usuario en desarrollo a 100 usuarios concurrentes sin que todo se caiga.


El problema

Tu agente funciona perfecto en desarrollo. Lo pruebas con un tema, ejecuta en 3 segundos, genera un reporte limpio. Listo para producción, ¿verdad?

No.

En producción con 100 usuarios concurrentes:

  • Un nodo de búsqueda tarda 45 segundos porque la API externa está lenta. Los otros 99 usuarios esperan.
  • Tu agente hace 200 llamadas por minuto a una API que permite 60. Te bloquean el acceso.
  • Un usuario reporta "mi investigación falló". No tienes idea de dónde ni por qué — solo sabes que falló.
  • El equipo pregunta "¿cuántos tokens gastamos ayer?" y no puedes responder.

Estos son problemas de operabilidad, no de funcionalidad. Tu agente hace lo correcto — pero no puedes operar lo, monitorear lo, ni escalar lo. Esta cápsula cierra esa brecha.


Timeouts por nodo

El problema que resuelve

Un nodo de tu grafo llama a una API externa. Normalmente responde en 2 segundos. Pero un día la API está lenta y tarda 60 segundos. Sin timeout, ese nodo bloquea todo el pipeline. Con 100 requests concurrentes, tienes 100 threads bloqueados esperando una API que no va a mejorar.

Implementación con asyncio

from dotenv import load_dotenv
load_dotenv()

import asyncio
import time
from langgraph.func import entrypoint, task


@task
async def fast_search(query: str) -> dict:
    """Búsqueda que responde rápido."""
    await asyncio.sleep(0.5)
    return {"source": "fast_api", "content": f"Resultados rápidos para '{query}'"}


@task
async def slow_search(query: str) -> dict:
    """Búsqueda que simula una API lenta."""
    await asyncio.sleep(10)
    return {"source": "slow_api", "content": f"Resultados lentos para '{query}'"}


@task
async def search_with_timeout(query: str, timeout_seconds: float = 3.0) -> dict:
    """Ejecuta búsqueda con timeout. Si excede el tiempo, retorna error."""
    try:
        result = await asyncio.wait_for(
            slow_search.acall(query),
            timeout=timeout_seconds,
        )
        return {"status": "success", "result": result}
    except asyncio.TimeoutError:
        return {
            "status": "timeout",
            "error": f"Búsqueda excedió {timeout_seconds}s",
            "query": query,
        }


@entrypoint()
async def research_with_timeouts(topic: str) -> dict:
    start = time.time()

    fast_future = fast_search(topic)
    slow_future = search_with_timeout(topic, timeout_seconds=3.0)

    fast_result = await fast_future
    slow_result = await slow_future

    elapsed = time.time() - start

    return {
        "topic": topic,
        "fast_result": fast_result,
        "slow_result": slow_result,
        "total_time_seconds": round(elapsed, 2),
    }


result = asyncio.run(
    research_with_timeouts.ainvoke("machine learning en medicina")
)
print(f"Tiempo total: {result['total_time_seconds']}s")
print(f"Fast: {result['fast_result']['source']}")
print(f"Slow: {result['slow_result']['status']}")
# Output esperado:
# Tiempo total: ~3.0s (no 10s)
# Fast: fast_api
# Slow: timeout

Sin el timeout, el pipeline tardaría 10 segundos esperando a slow_search. Con timeout de 3 segundos, falla rápido y el pipeline continúa con los resultados disponibles.

Timeout con threading (para código síncrono)

Si tus tasks son síncronas, puedes usar concurrent.futures para timeouts:

from dotenv import load_dotenv
load_dotenv()

import time
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langgraph.func import entrypoint, task


def _slow_api_call(query: str) -> str:
    """Simula una API que tarda mucho."""
    time.sleep(15)
    return f"Resultado para {query}"


@task
def search_with_thread_timeout(query: str, timeout_seconds: float = 3.0) -> dict:
    """Ejecuta una búsqueda síncrona con timeout usando ThreadPoolExecutor."""
    with ThreadPoolExecutor(max_workers=1) as executor:
        future = executor.submit(_slow_api_call, query)
        try:
            result = future.result(timeout=timeout_seconds)
            return {"status": "success", "content": result}
        except FuturesTimeoutError:
            return {
                "status": "timeout",
                "error": f"API excedió {timeout_seconds}s",
            }


@entrypoint()
def pipeline_with_sync_timeout(topic: str) -> dict:
    result = search_with_thread_timeout(topic, timeout_seconds=2.0).result()
    return result


result = pipeline_with_sync_timeout.invoke("quantum computing")
print(f"Status: {result['status']}")
print(f"Error: {result.get('error', 'ninguno')}")
# Output esperado:
# Status: timeout
# Error: API excedió 2.0s

Cuándo usar cada approach

EscenarioApproachPor qué
Tasks async (APIs con aiohttp)asyncio.wait_forNativo, eficiente, no crea threads extra
Tasks síncronas (requests, SDK)ThreadPoolExecutorFunciona con cualquier código blocking
Timeout global del pipelineTimeout en el .invoke() callerNo modifica el grafo

Rate limiting interno

El problema que resuelve

Tu agente busca en 3 fuentes para cada sub-query, con 4 sub-queries. Son 12 llamadas a APIs externas en paralelo. Si la API permite 10 requests por minuto, las últimas 2 llamadas fallan con 429 Too Many Requests. Peor: si tienes 10 usuarios concurrentes, son 120 llamadas simultáneas.

Implementación con InMemoryRateLimiter

LangChain incluye InMemoryRateLimiter para controlar la tasa de llamadas:

from dotenv import load_dotenv
load_dotenv()

import time
from langchain_core.rate_limiters import InMemoryRateLimiter
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task

rate_limiter = InMemoryRateLimiter(
    requests_per_second=2,
    check_every_n_seconds=0.1,
    max_bucket_size=5,
)

model = init_chat_model(
    "openai:gpt-4.1-mini",
    rate_limiter=rate_limiter,
)


@task
def analyze_topic(topic: str, index: int) -> dict:
    """Analiza un sub-tema con rate limiting automático."""
    start = time.time()
    response = model.invoke(
        f"Resume en una oración el tema: {topic}"
    )
    elapsed = time.time() - start
    return {
        "index": index,
        "topic": topic,
        "summary": response.content,
        "elapsed_seconds": round(elapsed, 2),
    }


@entrypoint()
def rate_limited_pipeline(topics: list) -> dict:
    start = time.time()

    futures = [analyze_topic(t, i) for i, t in enumerate(topics)]
    results = [f.result() for f in futures]

    total_elapsed = time.time() - start

    return {
        "results": results,
        "total_time_seconds": round(total_elapsed, 2),
        "requests_made": len(results),
    }


topics = [
    "inteligencia artificial",
    "computación cuántica",
    "energías renovables",
    "biotecnología",
    "exploración espacial",
    "ciberseguridad",
]

result = rate_limited_pipeline.invoke(topics)
print(f"Total: {result['total_time_seconds']}s para {result['requests_made']} requests")
for r in result["results"]:
    print(f"  [{r['index']}] {r['topic']}: {r['elapsed_seconds']}s")
# Output esperado:
# Total: ~3-4s para 6 requests (rate limited a 2/s)
#   [0] inteligencia artificial: 0.8s
#   [1] computación cuántica: 1.2s
#   [2] energías renovables: 1.5s
#   ...

Rate limiter por servicio

En producción, diferentes APIs tienen diferentes límites. Crea un rate limiter por servicio:

from langchain_core.rate_limiters import InMemoryRateLimiter

rate_limiters = {
    "openai": InMemoryRateLimiter(
        requests_per_second=5,
        check_every_n_seconds=0.1,
        max_bucket_size=10,
    ),
    "search_api": InMemoryRateLimiter(
        requests_per_second=1,
        check_every_n_seconds=0.1,
        max_bucket_size=3,
    ),
    "news_api": InMemoryRateLimiter(
        requests_per_second=2,
        check_every_n_seconds=0.1,
        max_bucket_size=5,
    ),
}

Cada modelo o cliente usa su propio rate limiter. Así, las llamadas a OpenAI no bloquean las llamadas a la API de búsqueda y viceversa.


Logging estructurado

El problema que resuelve

Tu agente falla. El log dice:

ERROR: Something went wrong

Inútil. No sabes qué nodo falló, cuánto tardó antes de fallar, qué input recibió, ni qué request del usuario causó el error. Con logging estructurado, el log dice:

{"node": "search_academic", "duration_ms": 4500, "status": "error", "error": "TimeoutError", "request_id": "abc123", "query": "quantum computing", "timestamp": "2026-03-08T14:30:00"}

Ahora puedes filtrar por request_id, ver que search_academic es el nodo problemático, y saber que tardó 4.5 segundos antes de fallar.

Implementación completa

from dotenv import load_dotenv
load_dotenv()

import json
import time
import uuid
import logging
from langgraph.func import entrypoint, task

logging.basicConfig(
    level=logging.INFO,
    format="%(message)s",
)
logger = logging.getLogger("research_agent")


class StructuredLogger:
    """Logger que emite JSON estructurado con contexto de request."""

    def __init__(self, logger_instance: logging.Logger):
        self._logger = logger_instance
        self._request_id = None

    def set_request_id(self, request_id: str):
        self._request_id = request_id

    def _log(self, level: str, node: str, **kwargs):
        entry = {
            "timestamp": time.strftime("%Y-%m-%dT%H:%M:%S"),
            "level": level,
            "request_id": self._request_id,
            "node": node,
            **kwargs,
        }
        self._logger.info(json.dumps(entry, ensure_ascii=False))

    def node_start(self, node: str, **kwargs):
        self._log("INFO", node, event="node_start", **kwargs)

    def node_end(self, node: str, duration_ms: float, **kwargs):
        self._log("INFO", node, event="node_end", duration_ms=round(duration_ms, 1), **kwargs)

    def node_error(self, node: str, error: str, duration_ms: float, **kwargs):
        self._log("ERROR", node, event="node_error", error=error, duration_ms=round(duration_ms, 1), **kwargs)


slog = StructuredLogger(logger)


@task
def search_with_logging(query: str, source: str) -> dict:
    """Búsqueda con logging estructurado automático."""
    slog.node_start("search", source=source, query=query)
    start = time.time()

    try:
        time.sleep(0.3)

        if source == "academic" and "quantum" in query.lower():
            raise ConnectionError("Academic API unavailable")

        result = {
            "source": source,
            "content": f"Resultados de {source} para '{query}'",
        }
        duration_ms = (time.time() - start) * 1000
        slog.node_end("search", duration_ms=duration_ms, source=source, status="success")
        return result

    except Exception as e:
        duration_ms = (time.time() - start) * 1000
        slog.node_error("search", error=str(e), duration_ms=duration_ms, source=source)
        raise


@task
def synthesize_with_logging(results: list) -> str:
    """Síntesis con logging estructurado."""
    slog.node_start("synthesize", num_sources=len(results))
    start = time.time()

    summary = f"Síntesis de {len(results)} fuentes completada."

    duration_ms = (time.time() - start) * 1000
    slog.node_end("synthesize", duration_ms=duration_ms, num_sources=len(results))
    return summary


@entrypoint()
def logged_pipeline(topic: str) -> dict:
    request_id = uuid.uuid4().hex[:8]
    slog.set_request_id(request_id)

    slog.node_start("pipeline", topic=topic)
    pipeline_start = time.time()

    sources = ["web", "academic", "news"]
    futures = [search_with_logging(topic, s) for s in sources]

    results = []
    errors = []
    for i, future in enumerate(futures):
        try:
            results.append(future.result())
        except Exception as e:
            errors.append({"source": sources[i], "error": str(e)})

    summary = synthesize_with_logging(results).result()

    pipeline_duration = (time.time() - pipeline_start) * 1000
    slog.node_end("pipeline", duration_ms=pipeline_duration, sources_ok=len(results), sources_failed=len(errors))

    return {
        "request_id": request_id,
        "summary": summary,
        "sources_ok": len(results),
        "errors": errors,
        "pipeline_duration_ms": round(pipeline_duration, 1),
    }


result = logged_pipeline.invoke("quantum computing")
print(f"\nRequest {result['request_id']}: {result['sources_ok']} fuentes OK, {len(result['errors'])} errores")
print(f"Duración total: {result['pipeline_duration_ms']}ms")
# Output esperado (logs en stderr, resultado en stdout):
# {"timestamp": "2026-03-08T14:30:00", "level": "INFO", "request_id": "a1b2c3d4", "node": "pipeline", "event": "node_start", "topic": "quantum computing"}
# {"timestamp": "2026-03-08T14:30:00", "level": "INFO", "request_id": "a1b2c3d4", "node": "search", "event": "node_start", "source": "web", "query": "quantum computing"}
# {"timestamp": "2026-03-08T14:30:00", "level": "INFO", "request_id": "a1b2c3d4", "node": "search", "event": "node_end", "duration_ms": 302.1, "source": "web", "status": "success"}
# {"timestamp": "2026-03-08T14:30:01", "level": "ERROR", "request_id": "a1b2c3d4", "node": "search", "event": "node_error", "error": "Academic API unavailable", "duration_ms": 300.5, "source": "academic"}
# {"timestamp": "2026-03-08T14:30:01", "level": "INFO", "request_id": "a1b2c3d4", "node": "search", "event": "node_end", "duration_ms": 301.3, "source": "news", "status": "success"}
# {"timestamp": "2026-03-08T14:30:01", "level": "INFO", "request_id": "a1b2c3d4", "node": "synthesize", "event": "node_start", "num_sources": 2}
# {"timestamp": "2026-03-08T14:30:01", "level": "INFO", "request_id": "a1b2c3d4", "node": "pipeline", "event": "node_end", "duration_ms": 920.5, "sources_ok": 2, "sources_failed": 1}
#
# Request a1b2c3d4: 2 fuentes OK, 1 errores
# Duración total: 920.5ms

Correlation IDs: trazando una request completa

El request_id es la pieza clave. Cuando un usuario reporta "mi investigación falló", te da el request_id y puedes filtrar todos los logs de esa request:

# Filtrar todos los logs de una request específica
cat logs.jsonl | jq 'select(.request_id == "a1b2c3d4")'

Verás cada nodo que se ejecutó, cuánto tardó, cuáles fallaron, y en qué orden. Sin correlation IDs, los logs de 100 requests concurrentes se mezclan y son imposibles de analizar.


Métricas de rendimiento

El problema que resuelve

"¿Cuánto tarda nuestro pipeline en promedio?" "¿Qué nodo es el cuello de botella?" "¿Cuántos tokens gastamos ayer?" Sin métricas, respondes con "no sé" o con una anécdota de la última vez que probaste. Con métricas, respondes con números.

Implementación: colector de métricas in-memory

from dotenv import load_dotenv
load_dotenv()

import time
import statistics
from collections import defaultdict
from langgraph.func import entrypoint, task
from langchain.chat_models import init_chat_model


class MetricsCollector:
    """Colector de métricas en memoria para pipelines de LangGraph."""

    def __init__(self):
        self._latencies: dict[str, list[float]] = defaultdict(list)
        self._success_count: dict[str, int] = defaultdict(int)
        self._error_count: dict[str, int] = defaultdict(int)
        self._token_usage: dict[str, int] = defaultdict(int)

    def record_latency(self, node: str, duration_ms: float):
        self._latencies[node].append(duration_ms)

    def record_success(self, node: str):
        self._success_count[node] += 1

    def record_error(self, node: str):
        self._error_count[node] += 1

    def record_tokens(self, node: str, tokens: int):
        self._token_usage[node] += tokens

    def get_summary(self) -> dict:
        summary = {}
        for node in set(
            list(self._latencies.keys())
            + list(self._success_count.keys())
            + list(self._error_count.keys())
        ):
            latencies = self._latencies.get(node, [])
            successes = self._success_count.get(node, 0)
            errors = self._error_count.get(node, 0)
            total = successes + errors

            summary[node] = {
                "total_calls": total,
                "successes": successes,
                "errors": errors,
                "success_rate": round(successes / total, 2) if total > 0 else 0,
                "avg_latency_ms": round(statistics.mean(latencies), 1) if latencies else 0,
                "p50_latency_ms": round(statistics.median(latencies), 1) if latencies else 0,
                "p95_latency_ms": round(
                    sorted(latencies)[int(len(latencies) * 0.95)] if latencies else 0, 1
                ),
                "max_latency_ms": round(max(latencies), 1) if latencies else 0,
                "total_tokens": self._token_usage.get(node, 0),
            }
        return summary

    def reset(self):
        self._latencies.clear()
        self._success_count.clear()
        self._error_count.clear()
        self._token_usage.clear()


metrics = MetricsCollector()
model = init_chat_model("openai:gpt-4.1-mini")


@task
def search_node(query: str, source: str) -> dict:
    start = time.time()
    try:
        time.sleep(0.2)
        result = {"source": source, "content": f"Resultado de {source}"}
        duration_ms = (time.time() - start) * 1000
        metrics.record_latency(f"search_{source}", duration_ms)
        metrics.record_success(f"search_{source}")
        return result
    except Exception:
        duration_ms = (time.time() - start) * 1000
        metrics.record_latency(f"search_{source}", duration_ms)
        metrics.record_error(f"search_{source}")
        raise


@task
def synthesize_node(topic: str, results: list) -> str:
    start = time.time()
    context = "\n".join(r["content"] for r in results)
    response = model.invoke(
        f"Resume en una oración estos resultados sobre '{topic}': {context}"
    )

    duration_ms = (time.time() - start) * 1000
    metrics.record_latency("synthesize", duration_ms)
    metrics.record_success("synthesize")

    token_count = response.usage_metadata.get("total_tokens", 0) if response.usage_metadata else 0
    metrics.record_tokens("synthesize", token_count)

    return response.content


@entrypoint()
def pipeline_with_metrics(topic: str) -> dict:
    start = time.time()

    sources = ["web", "academic", "news"]
    futures = [search_node(topic, s) for s in sources]
    results = [f.result() for f in futures]

    summary = synthesize_node(topic, results).result()

    pipeline_ms = (time.time() - start) * 1000
    metrics.record_latency("pipeline_total", pipeline_ms)
    metrics.record_success("pipeline_total")

    return {"summary": summary}


for topic in ["AI en educación", "quantum computing", "energías renovables"]:
    pipeline_with_metrics.invoke(topic)

summary = metrics.get_summary()
print("\n📊 MÉTRICAS DE RENDIMIENTO")
print("=" * 60)
for node, data in sorted(summary.items()):
    print(f"\n  {node}:")
    print(f"    Calls: {data['total_calls']} ({data['success_rate']:.0%} éxito)")
    print(f"    Latencia: avg={data['avg_latency_ms']}ms, p50={data['p50_latency_ms']}ms, p95={data['p95_latency_ms']}ms")
    if data['total_tokens'] > 0:
        print(f"    Tokens: {data['total_tokens']}")
# Output esperado:
# 📊 MÉTRICAS DE RENDIMIENTO
# ============================================================
#
#   pipeline_total:
#     Calls: 3 (100% éxito)
#     Latencia: avg=1200.5ms, p50=1180.3ms, p95=1250.1ms
#
#   search_academic:
#     Calls: 3 (100% éxito)
#     Latencia: avg=201.2ms, p50=200.8ms, p95=202.1ms
#
#   search_news:
#     Calls: 3 (100% éxito)
#     Latencia: avg=200.5ms, p50=200.3ms, p95=201.0ms
#
#   search_web:
#     Calls: 3 (100% éxito)
#     Latencia: avg=200.8ms, p50=200.5ms, p95=201.5ms
#
#   synthesize:
#     Calls: 3 (100% éxito)
#     Latencia: avg=980.3ms, p50=950.1ms, p95=1050.2ms
#     Tokens: 450

Las métricas revelan que synthesize es el cuello de botella (980ms vs 200ms de búsqueda). Con esta información, puedes decidir si optimizar el prompt, usar un modelo más rápido para síntesis, o cachear resultados similares.


Combinando todo: grafo production-ready

Ahora combina los 4 patterns en un solo pipeline. Este es el patrón completo que usarías en producción:

from dotenv import load_dotenv
load_dotenv()

import json
import time
import uuid
import logging
import statistics
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langchain_core.rate_limiters import InMemoryRateLimiter
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task


# === INFRAESTRUCTURA ===

logging.basicConfig(level=logging.INFO, format="%(message)s")
log = logging.getLogger("production_agent")


class ProdLogger:
    def __init__(self):
        self.request_id = None

    def set_request(self, rid: str):
        self.request_id = rid

    def info(self, node: str, event: str, **kw):
        entry = {"ts": time.strftime("%H:%M:%S"), "rid": self.request_id, "node": node, "event": event, **kw}
        log.info(json.dumps(entry, ensure_ascii=False))

    def error(self, node: str, event: str, **kw):
        entry = {"ts": time.strftime("%H:%M:%S"), "rid": self.request_id, "node": node, "event": event, "level": "ERROR", **kw}
        log.info(json.dumps(entry, ensure_ascii=False))


class ProdMetrics:
    def __init__(self):
        self.latencies = defaultdict(list)
        self.counts = defaultdict(lambda: {"ok": 0, "err": 0})
        self.tokens = defaultdict(int)

    def record(self, node: str, duration_ms: float, success: bool, tokens: int = 0):
        self.latencies[node].append(duration_ms)
        self.counts[node]["ok" if success else "err"] += 1
        if tokens:
            self.tokens[node] += tokens

    def summary(self) -> dict:
        out = {}
        for node in self.latencies:
            lats = self.latencies[node]
            c = self.counts[node]
            total = c["ok"] + c["err"]
            out[node] = {
                "calls": total,
                "success_rate": f"{c['ok']/total:.0%}" if total else "N/A",
                "avg_ms": round(statistics.mean(lats), 1),
                "p95_ms": round(sorted(lats)[int(len(lats) * 0.95)], 1) if lats else 0,
                "tokens": self.tokens.get(node, 0),
            }
        return out


plog = ProdLogger()
pmetrics = ProdMetrics()

rate_limiter = InMemoryRateLimiter(
    requests_per_second=3,
    check_every_n_seconds=0.1,
    max_bucket_size=5,
)

model = init_chat_model("openai:gpt-4.1-mini", rate_limiter=rate_limiter)

SEARCH_TIMEOUT_SECONDS = 5.0


# === TASKS ===

def _mock_search(query: str, source: str) -> str:
    """Simula búsqueda externa con latencia variable."""
    import random
    delay = random.uniform(0.2, 0.8)
    time.sleep(delay)
    return f"[{source}] Resultados para '{query}': información relevante encontrada."


@task
def production_search(query: str, source: str) -> dict:
    """Búsqueda con timeout + logging + métricas."""
    plog.info("search", "start", source=source, query=query[:50])
    start = time.time()

    with ThreadPoolExecutor(max_workers=1) as executor:
        future = executor.submit(_mock_search, query, source)
        try:
            content = future.result(timeout=SEARCH_TIMEOUT_SECONDS)
            duration_ms = (time.time() - start) * 1000

            plog.info("search", "end", source=source, duration_ms=round(duration_ms))
            pmetrics.record("search", duration_ms, success=True)

            return {"source": source, "content": content, "status": "ok"}

        except FuturesTimeoutError:
            duration_ms = (time.time() - start) * 1000
            plog.error("search", "timeout", source=source, duration_ms=round(duration_ms))
            pmetrics.record("search", duration_ms, success=False)
            return {"source": source, "content": "", "status": "timeout"}


@task
def production_synthesize(topic: str, results: list) -> dict:
    """Síntesis con rate limiting + logging + métricas + tokens."""
    plog.info("synthesize", "start", num_sources=len(results))
    start = time.time()

    context = "\n".join(r["content"] for r in results if r["content"])
    response = model.invoke(
        f"Resume en 2-3 oraciones estos resultados sobre '{topic}':\n{context}"
    )

    duration_ms = (time.time() - start) * 1000
    tokens = response.usage_metadata.get("total_tokens", 0) if response.usage_metadata else 0

    plog.info("synthesize", "end", duration_ms=round(duration_ms), tokens=tokens)
    pmetrics.record("synthesize", duration_ms, success=True, tokens=tokens)

    return {"summary": response.content, "tokens_used": tokens}


# === PIPELINE ===

@entrypoint()
def production_pipeline(topic: str) -> dict:
    request_id = uuid.uuid4().hex[:8]
    plog.set_request(request_id)
    plog.info("pipeline", "start", topic=topic)

    pipeline_start = time.time()

    sources = ["web", "academic", "news"]
    search_futures = [production_search(topic, s) for s in sources]
    search_results = [f.result() for f in search_futures]

    successful = [r for r in search_results if r["status"] == "ok"]
    failed = [r for r in search_results if r["status"] != "ok"]

    if not successful:
        plog.error("pipeline", "all_sources_failed")
        pmetrics.record("pipeline", (time.time() - pipeline_start) * 1000, success=False)
        return {"error": "Todas las fuentes fallaron", "request_id": request_id}

    synthesis = production_synthesize(topic, successful).result()

    pipeline_ms = (time.time() - pipeline_start) * 1000
    plog.info("pipeline", "end", duration_ms=round(pipeline_ms), sources_ok=len(successful), sources_failed=len(failed))
    pmetrics.record("pipeline", pipeline_ms, success=True)

    return {
        "request_id": request_id,
        "summary": synthesis["summary"],
        "sources_ok": len(successful),
        "sources_failed": len(failed),
        "tokens_used": synthesis["tokens_used"],
        "duration_ms": round(pipeline_ms),
    }


result = production_pipeline.invoke("impacto de LLMs en producción de software")
print(f"\n{'=' * 50}")
print(f"Request: {result['request_id']}")
print(f"Fuentes: {result['sources_ok']} OK, {result['sources_failed']} fallidas")
print(f"Tokens: {result['tokens_used']}")
print(f"Duración: {result['duration_ms']}ms")
print(f"Resumen: {result['summary'][:120]}...")

print(f"\n{'=' * 50}")
print("MÉTRICAS ACUMULADAS:")
for node, data in pmetrics.summary().items():
    print(f"  {node}: {data['calls']} calls, {data['success_rate']} éxito, avg {data['avg_ms']}ms")
# Output esperado:
# {"ts": "14:30:00", "rid": "a1b2c3d4", "node": "pipeline", "event": "start", "topic": "impacto de LLMs en producción de software"}
# {"ts": "14:30:00", "rid": "a1b2c3d4", "node": "search", "event": "start", "source": "web", ...}
# {"ts": "14:30:00", "rid": "a1b2c3d4", "node": "search", "event": "end", "source": "web", "duration_ms": 450}
# ... (más logs)
#
# ==================================================
# Request: a1b2c3d4
# Fuentes: 3 OK, 0 fallidas
# Tokens: 150
# Duración: 1850ms
# Resumen: Los LLMs están transformando la producción de software...
#
# ==================================================
# MÉTRICAS ACUMULADAS:
#   search: 3 calls, 100% éxito, avg 420.5ms
#   synthesize: 1 calls, 100% éxito, avg 980.3ms
#   pipeline: 1 calls, 100% éxito, avg 1850.2ms

Este pipeline tiene todo lo que necesitas para producción:

  • Timeouts por búsqueda (5s máximo, no bloquea el pipeline)
  • Rate limiting en las llamadas al modelo (3 req/s, respeta límites de OpenAI)
  • Logging estructurado con correlation ID (filtrable por request)
  • Métricas de latencia, success rate, y tokens (responde preguntas de negocio)
  • Graceful degradation (si todas las fuentes fallan, no crashea)

Qué necesitarías sin LangGraph

Sin LangGraph, implementar este pipeline requiere:

ComponenteCon LangGraphSin LangGraph
ParalelismoFutures de @taskconcurrent.futures.ThreadPoolExecutor manual, gestión de threads
CheckpointingAutomático con MemorySaverBase de datos + lógica de checkpoint + serialización custom
RetryPython try/except dentro de @task con checkpointRetry library + state management manual para saber qué reintentar
Composición@entrypoint anida @task naturalmenteFunciones anidadas con contexto compartido via globals o inyección
Timeoutasyncio.wait_for o ThreadPool en el @taskIgual, pero sin el beneficio de checkpoint automático si timeout
State recoveryResume desde el último checkpointReimplementar todo el state management desde cero

El boilerplate sin LangGraph es 3-5x más código para la misma funcionalidad, y no incluye checkpointing (que es la parte más difícil de implementar correctamente).


Troubleshooting

Problema 1: "El rate limiter bloquea todo y el pipeline es muy lento"

Síntoma: Con requests_per_second=1 y 10 tasks paralelas, el pipeline tarda 10 segundos.

Causa: El rate limiter serializa las requests. Si tienes 10 tasks y solo 1 req/s permitida, cada task espera su turno.

Solución: Ajusta requests_per_second al límite real de tu API. Si la API permite 60 req/min, usa requests_per_second=1 (correcto). Si permite 600 req/min, usa requests_per_second=10. El max_bucket_size permite ráfagas cortas:

rate_limiter = InMemoryRateLimiter(
    requests_per_second=10,
    check_every_n_seconds=0.05,
    max_bucket_size=20,
)

Problema 2: "Los logs de diferentes requests se mezclan y no puedo rastrear una"

Síntoma: Con 10 requests concurrentes, los logs JSON están intercalados y no sabes cuáles pertenecen a cuál request.

Causa: El request_id se pierde entre nodos, o estás usando un logger global sin correlation.

Solución: El StructuredLogger de esta cápsula usa set_request_id() al inicio del pipeline. Para concurrencia real, usa contextvars de Python para que cada thread tenga su propio request_id:

import contextvars

_request_id_var = contextvars.ContextVar("request_id", default="unknown")

class ThreadSafeLogger:
    def set_request_id(self, rid: str):
        _request_id_var.set(rid)

    def _log(self, node: str, **kwargs):
        entry = {"request_id": _request_id_var.get(), "node": node, **kwargs}
        print(json.dumps(entry))

Problema 3: "Las métricas de p95 no son confiables"

Síntoma: El p95 reportado no coincide con la realidad — a veces es más bajo que el promedio.

Causa: Con pocas muestras (menos de 20), los percentiles estadísticos no son representativos. Con 3 muestras, el p95 es básicamente el valor más alto.

Solución: Solo reporta percentiles cuando tienes suficientes muestras:

def get_p95(self, node: str) -> str:
    lats = self.latencies.get(node, [])
    if len(lats) < 20:
        return f"~{max(lats):.1f}ms (solo {len(lats)} muestras)"
    return f"{sorted(lats)[int(len(lats) * 0.95)]:.1f}ms"

Problema 4: "El timeout no cancela la tarea, solo ignora el resultado"

Síntoma: Pones timeout de 3 segundos, pero la función sigue ejecutándose en background consumiendo recursos.

Causa: ThreadPoolExecutor.submit().result(timeout=3) lanza TimeoutError en el caller, pero el thread sigue corriendo. Python no puede matar threads de forma limpia.

Solución: Para timeouts reales que liberen recursos, usa asyncio con tasks cancelables:

async def search_with_real_cancel(query: str, timeout: float):
    task = asyncio.create_task(async_search(query))
    try:
        return await asyncio.wait_for(task, timeout=timeout)
    except asyncio.TimeoutError:
        task.cancel()
        return {"status": "timeout"}

Con asyncio, task.cancel() realmente cancela la coroutine. Con threads, no hay equivalente limpio.


Ejercicios

Ejercicio 1: Timeout configurable por fuente (Fácil)

Modifica el pipeline de timeouts para que cada fuente tenga su propio timeout: web = 2s, academic = 5s, news = 3s. Simula que academic tarda 4s (dentro del timeout) y web tarda 3s (fuera de su timeout).

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langgraph.func import entrypoint, task

SOURCE_TIMEOUTS = {
    "web": 2.0,
    "academic": 5.0,
    "news": 3.0,
}

SOURCE_DELAYS = {
    "web": 3.0,
    "academic": 4.0,
    "news": 1.0,
}


def _mock_search(query: str, source: str) -> str:
    time.sleep(SOURCE_DELAYS[source])
    return f"[{source}] Resultados para '{query}'"


@task
def search_with_custom_timeout(query: str, source: str) -> dict:
    timeout = SOURCE_TIMEOUTS[source]
    with ThreadPoolExecutor(max_workers=1) as executor:
        future = executor.submit(_mock_search, query, source)
        try:
            content = future.result(timeout=timeout)
            return {"source": source, "status": "ok", "content": content}
        except FuturesTimeoutError:
            return {"source": source, "status": "timeout", "timeout_seconds": timeout}


@entrypoint()
def pipeline(topic: str) -> dict:
    sources = ["web", "academic", "news"]
    futures = [search_with_custom_timeout(topic, s) for s in sources]
    results = [f.result() for f in futures]

    ok = [r for r in results if r["status"] == "ok"]
    failed = [r for r in results if r["status"] != "ok"]

    return {
        "ok_sources": [r["source"] for r in ok],
        "failed_sources": [f"{r['source']} (timeout {r.get('timeout_seconds', '?')}s)" for r in failed],
    }


result = pipeline.invoke("deep learning")
print(f"OK: {result['ok_sources']}")
print(f"Failed: {result['failed_sources']}")
# Output esperado:
# OK: ['academic', 'news']  (academic tarda 4s pero timeout es 5s; news tarda 1s)
# Failed: ['web (timeout 2.0s)']  (web tarda 3s pero timeout es 2s)

Ejercicio 2: Rate limiter con tracking visual (Fácil)

Crea un pipeline que haga 8 llamadas a un modelo con rate limit de 2/segundo. Muestra el timestamp de cada llamada para verificar visualmente que se respeta el rate limit.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
from langchain_core.rate_limiters import InMemoryRateLimiter
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task

rate_limiter = InMemoryRateLimiter(
    requests_per_second=2,
    check_every_n_seconds=0.1,
    max_bucket_size=2,
)

model = init_chat_model("openai:gpt-4.1-mini", rate_limiter=rate_limiter)
start_time = time.time()


@task
def timed_call(index: int) -> dict:
    elapsed = time.time() - start_time
    response = model.invoke(f"Di 'hola {index}' y nada más.")
    after = time.time() - start_time
    return {
        "index": index,
        "request_at": round(elapsed, 2),
        "response_at": round(after, 2),
        "content": response.content.strip(),
    }


@entrypoint()
def rate_limit_demo(count: int) -> list:
    futures = [timed_call(i) for i in range(count)]
    return [f.result() for f in futures]


results = rate_limit_demo.invoke(8)
for r in sorted(results, key=lambda x: x["request_at"]):
    print(f"  Call {r['index']}: request @{r['request_at']}s → response @{r['response_at']}s | {r['content']}")
# Output esperado (tiempos aproximados — rate limited a 2/s):
#   Call 0: request @0.01s → response @0.85s | hola 0
#   Call 1: request @0.01s → response @0.90s | hola 1
#   Call 2: request @0.51s → response @1.35s | hola 2
#   Call 3: request @0.51s → response @1.40s | hola 3
#   Call 4: request @1.01s → response @1.85s | hola 4
#   ... (cada par de llamadas separado por ~0.5s)

Las llamadas salen en pares de 2 (el rate limit), con ~0.5 segundos entre cada par.

Ejercicio 3: Logger con niveles y filtrado (Medio)

Extiende el StructuredLogger para soportar niveles (DEBUG, INFO, WARN, ERROR) y un nivel mínimo configurable. Si el nivel mínimo es WARN, los logs de INFO y DEBUG no se emiten. Prueba con un pipeline que genere logs de todos los niveles.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import json
import time
import logging
from langgraph.func import entrypoint, task

logging.basicConfig(level=logging.INFO, format="%(message)s")


class LeveledLogger:
    LEVELS = {"DEBUG": 0, "INFO": 1, "WARN": 2, "ERROR": 3}

    def __init__(self, min_level: str = "INFO"):
        self._min = self.LEVELS.get(min_level, 1)
        self._logger = logging.getLogger("leveled")
        self._rid = None

    def set_request_id(self, rid: str):
        self._rid = rid

    def _emit(self, level: str, node: str, message: str, **kw):
        if self.LEVELS.get(level, 0) < self._min:
            return
        entry = {
            "ts": time.strftime("%H:%M:%S"),
            "level": level,
            "rid": self._rid,
            "node": node,
            "msg": message,
            **kw,
        }
        self._logger.info(json.dumps(entry, ensure_ascii=False))

    def debug(self, node, msg, **kw): self._emit("DEBUG", node, msg, **kw)
    def info(self, node, msg, **kw): self._emit("INFO", node, msg, **kw)
    def warn(self, node, msg, **kw): self._emit("WARN", node, msg, **kw)
    def error(self, node, msg, **kw): self._emit("ERROR", node, msg, **kw)


log_info = LeveledLogger(min_level="INFO")
log_warn = LeveledLogger(min_level="WARN")


@task
def demo_task(value: int, logger: LeveledLogger) -> dict:
    logger.debug("demo", f"Procesando valor {value}", value=value)
    logger.info("demo", f"Valor recibido: {value}")

    if value > 5:
        logger.warn("demo", f"Valor alto: {value}", threshold=5)
    if value > 8:
        logger.error("demo", f"Valor crítico: {value}", threshold=8)

    return {"value": value, "processed": True}


@entrypoint()
def logging_demo(config: dict) -> dict:
    logger = log_info if config.get("verbose") else log_warn
    logger.set_request_id("demo-001")

    values = config["values"]
    futures = [demo_task(v, logger) for v in values]
    results = [f.result() for f in futures]

    return {"processed": len(results)}


print("=== Con min_level=INFO (verbose) ===")
logging_demo.invoke({"values": [2, 6, 9], "verbose": True})

print("\n=== Con min_level=WARN (solo warnings+) ===")
logging_demo.invoke({"values": [2, 6, 9], "verbose": False})
# Output esperado:
# === Con min_level=INFO (verbose) ===
# {"ts": "14:30:00", "level": "INFO", "rid": "demo-001", "node": "demo", "msg": "Valor recibido: 2"}
# {"ts": "14:30:00", "level": "INFO", "rid": "demo-001", "node": "demo", "msg": "Valor recibido: 6"}
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 6", "threshold": 5}
# {"ts": "14:30:00", "level": "INFO", "rid": "demo-001", "node": "demo", "msg": "Valor recibido: 9"}
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 9", "threshold": 5}
# {"ts": "14:30:00", "level": "ERROR", "rid": "demo-001", "node": "demo", "msg": "Valor crítico: 9", "threshold": 8}
#
# === Con min_level=WARN (solo warnings+) ===
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 6", "threshold": 5}
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 9", "threshold": 5}
# {"ts": "14:30:00", "level": "ERROR", "rid": "demo-001", "node": "demo", "msg": "Valor crítico: 9", "threshold": 8}

Con min_level="WARN", los logs de DEBUG e INFO se filtran. Solo ves warnings y errores — exactamente lo que quieres en producción cuando no estás debugging.

Ejercicio 4: Dashboard de métricas con resumen (Medio)

Extiende el MetricsCollector para generar un dashboard de texto que incluya: top 3 nodos más lentos, success rate total, total de tokens consumidos, y alertas para nodos con success rate < 90%.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
import statistics
from collections import defaultdict
from langgraph.func import entrypoint, task


class DashboardMetrics:
    def __init__(self):
        self.latencies = defaultdict(list)
        self.counts = defaultdict(lambda: {"ok": 0, "err": 0})
        self.tokens = defaultdict(int)

    def record(self, node: str, ms: float, ok: bool, tokens: int = 0):
        self.latencies[node].append(ms)
        self.counts[node]["ok" if ok else "err"] += 1
        if tokens:
            self.tokens[node] += tokens

    def dashboard(self) -> str:
        lines = []
        lines.append("╔══════════════════════════════════════════════════╗")
        lines.append("║          📊 PRODUCTION DASHBOARD                ║")
        lines.append("╚══════════════════════════════════════════════════╝")

        total_ok = sum(c["ok"] for c in self.counts.values())
        total_err = sum(c["err"] for c in self.counts.values())
        total_all = total_ok + total_err
        total_tokens = sum(self.tokens.values())

        lines.append(f"\n  🎯 Success rate global: {total_ok}/{total_all} ({total_ok/total_all:.0%})" if total_all else "")
        lines.append(f"  🪙 Tokens totales: {total_tokens}")
        lines.append(f"  📦 Nodos tracked: {len(self.latencies)}")

        avg_by_node = {
            n: statistics.mean(lats) for n, lats in self.latencies.items()
        }
        slowest = sorted(avg_by_node.items(), key=lambda x: x[1], reverse=True)[:3]

        lines.append(f"\n  🐌 TOP 3 NODOS MÁS LENTOS:")
        for i, (node, avg) in enumerate(slowest, 1):
            lines.append(f"     {i}. {node}: {avg:.1f}ms avg")

        alerts = []
        for node, c in self.counts.items():
            total = c["ok"] + c["err"]
            if total > 0 and c["ok"] / total < 0.9:
                rate = c["ok"] / total
                alerts.append(f"     ⚠️ {node}: {rate:.0%} success rate ({c['err']} errores)")

        if alerts:
            lines.append(f"\n  🚨 ALERTAS:")
            lines.extend(alerts)
        else:
            lines.append(f"\n  ✅ Sin alertas — todos los nodos > 90% success rate")

        return "\n".join(lines)


dm = DashboardMetrics()


@task
def reliable_node(name: str, delay: float) -> str:
    time.sleep(delay)
    dm.record(name, delay * 1000, ok=True, tokens=50)
    return f"{name} OK"


@task
def flaky_node(name: str, delay: float, fail_rate: float) -> str:
    import random
    time.sleep(delay)
    if random.random() < fail_rate:
        dm.record(name, delay * 1000, ok=False)
        raise RuntimeError(f"{name} falló")
    dm.record(name, delay * 1000, ok=True, tokens=30)
    return f"{name} OK"


@entrypoint()
def demo_dashboard(runs: int) -> str:
    for _ in range(runs):
        reliable_node("search_web", 0.2).result()
        reliable_node("synthesize", 0.8).result()
        try:
            flaky_node("search_academic", 0.3, 0.3).result()
        except Exception:
            pass

    return dm.dashboard()


print(demo_dashboard.invoke(10))
# Output esperado:
# ╔══════════════════════════════════════════════════╗
# ║          📊 PRODUCTION DASHBOARD                ║
# ╚══════════════════════════════════════════════════╝
#
#   🎯 Success rate global: 27/30 (90%)
#   🪙 Tokens totales: 1210
#   📦 Nodos tracked: 3
#
#   🐌 TOP 3 NODOS MÁS LENTOS:
#      1. synthesize: 800.0ms avg
#      2. search_academic: 300.0ms avg
#      3. search_web: 200.0ms avg
#
#   🚨 ALERTAS:
#      ⚠️ search_academic: 70% success rate (3 errores)

Ejercicio 5: Pipeline con circuit breaker (Avanzado)

Implementa un circuit breaker: si un nodo falla más de 3 veces consecutivas, las siguientes llamadas a ese nodo retornan un error inmediatamente (sin intentar) durante 30 segundos. Después de 30 segundos, permite una llamada de prueba ("half-open"). Si esa llamada tiene éxito, el circuit se cierra (normal). Si falla, el circuit se abre de nuevo.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
from langgraph.func import entrypoint, task


class CircuitBreaker:
    """Circuit breaker: CLOSED → OPEN → HALF_OPEN → CLOSED/OPEN."""

    def __init__(self, failure_threshold: int = 3, recovery_timeout: float = 30.0):
        self.failure_threshold = failure_threshold
        self.recovery_timeout = recovery_timeout
        self._states: dict[str, dict] = {}

    def _get_state(self, name: str) -> dict:
        if name not in self._states:
            self._states[name] = {
                "status": "CLOSED",
                "consecutive_failures": 0,
                "opened_at": 0,
            }
        return self._states[name]

    def can_execute(self, name: str) -> tuple[bool, str]:
        state = self._get_state(name)

        if state["status"] == "CLOSED":
            return True, "CLOSED"

        if state["status"] == "OPEN":
            elapsed = time.time() - state["opened_at"]
            if elapsed >= self.recovery_timeout:
                state["status"] = "HALF_OPEN"
                return True, "HALF_OPEN"
            return False, f"OPEN ({self.recovery_timeout - elapsed:.0f}s remaining)"

        if state["status"] == "HALF_OPEN":
            return True, "HALF_OPEN"

        return False, state["status"]

    def record_success(self, name: str):
        state = self._get_state(name)
        state["consecutive_failures"] = 0
        state["status"] = "CLOSED"

    def record_failure(self, name: str):
        state = self._get_state(name)
        state["consecutive_failures"] += 1

        if state["consecutive_failures"] >= self.failure_threshold:
            state["status"] = "OPEN"
            state["opened_at"] = time.time()


cb = CircuitBreaker(failure_threshold=3, recovery_timeout=5.0)

call_count = 0


@task
def protected_search(query: str, source: str) -> dict:
    global call_count
    can_exec, status = cb.can_execute(source)

    if not can_exec:
        return {"source": source, "status": "circuit_open", "circuit": status}

    try:
        call_count += 1
        if source == "flaky_api" and call_count <= 4:
            raise ConnectionError(f"flaky_api fallo #{call_count}")

        cb.record_success(source)
        return {"source": source, "status": "ok", "content": f"Resultado de {source}"}

    except Exception as e:
        cb.record_failure(source)
        return {"source": source, "status": "error", "error": str(e)}


@entrypoint()
def circuit_demo(config: dict) -> list:
    results = []

    for i in range(config["num_calls"]):
        r = protected_search(config["query"], config["source"]).result()
        results.append({"call": i + 1, **r})

        if config.get("sleep_between"):
            time.sleep(config["sleep_between"])

    return results


results = circuit_demo.invoke({
    "query": "test",
    "source": "flaky_api",
    "num_calls": 8,
    "sleep_between": 1.0,
})

for r in results:
    print(f"  Call {r['call']}: {r['status']} | {r.get('circuit', r.get('error', r.get('content', '')))}")
# Output esperado:
#   Call 1: error | flaky_api fallo #1
#   Call 2: error | flaky_api fallo #2
#   Call 3: error | flaky_api fallo #3         ← circuit OPENS here
#   Call 4: circuit_open | OPEN (2s remaining)  ← calls blocked
#   Call 5: circuit_open | OPEN (1s remaining)  ← calls blocked
#   Call 6: error | flaky_api fallo #4          ← HALF_OPEN test call fails → re-opens
#   Call 7: circuit_open | OPEN (4s remaining)  ← blocked again
#   Call 8: circuit_open | OPEN (3s remaining)  ← blocked again

El circuit breaker protege al sistema de seguir llamando a un servicio que está caído. En vez de fallar 100 veces en 100 requests, falla 3 veces, bloquea las siguientes, y eventualmente prueba si el servicio se recuperó.

Ejercicio 6: Pipeline completo production-ready (Avanzado)

Combina timeouts, rate limiting, logging, métricas y circuit breaker en un pipeline que busque en 3 fuentes. Una fuente debe ser "flaky" (falla el 50% del tiempo). El pipeline debe: respetar timeouts, loguear cada operación, registrar métricas, y usar circuit breaker para la fuente flaky.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import json
import time
import random
import logging
import statistics
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langgraph.func import entrypoint, task

logging.basicConfig(level=logging.INFO, format="%(message)s")


class MiniCircuitBreaker:
    def __init__(self, threshold=3, timeout=10.0):
        self.threshold = threshold
        self.timeout = timeout
        self._fails = defaultdict(int)
        self._opened = {}

    def allow(self, name):
        if name in self._opened:
            if time.time() - self._opened[name] > self.timeout:
                del self._opened[name]
                return True
            return False
        return True

    def success(self, name):
        self._fails[name] = 0
        self._opened.pop(name, None)

    def failure(self, name):
        self._fails[name] += 1
        if self._fails[name] >= self.threshold:
            self._opened[name] = time.time()


class MiniMetrics:
    def __init__(self):
        self.data = defaultdict(lambda: {"lats": [], "ok": 0, "err": 0})

    def record(self, node, ms, ok):
        self.data[node]["lats"].append(ms)
        self.data[node]["ok" if ok else "err"] += 1

    def report(self):
        for node, d in sorted(self.data.items()):
            total = d["ok"] + d["err"]
            rate = d["ok"] / total if total else 0
            avg = statistics.mean(d["lats"]) if d["lats"] else 0
            print(f"  {node}: {total} calls, {rate:.0%} ok, avg {avg:.0f}ms")


cb = MiniCircuitBreaker(threshold=2, timeout=5.0)
met = MiniMetrics()


def _search_impl(query, source):
    delays = {"reliable_a": 0.3, "flaky_b": 0.4, "reliable_c": 0.2}
    time.sleep(delays.get(source, 0.3))
    if source == "flaky_b" and random.random() < 0.5:
        raise ConnectionError(f"{source} failed")
    return f"[{source}] Results for '{query}'"


@task
def guarded_search(query: str, source: str) -> dict:
    if not cb.allow(source):
        met.record(source, 0, False)
        return {"source": source, "status": "circuit_open"}

    start = time.time()
    with ThreadPoolExecutor(max_workers=1) as ex:
        fut = ex.submit(_search_impl, query, source)
        try:
            content = fut.result(timeout=2.0)
            ms = (time.time() - start) * 1000
            cb.success(source)
            met.record(source, ms, True)
            return {"source": source, "status": "ok", "content": content}
        except FuturesTimeoutError:
            ms = (time.time() - start) * 1000
            cb.failure(source)
            met.record(source, ms, False)
            return {"source": source, "status": "timeout"}
        except Exception as e:
            ms = (time.time() - start) * 1000
            cb.failure(source)
            met.record(source, ms, False)
            return {"source": source, "status": "error", "error": str(e)}


@entrypoint()
def full_production_pipeline(config: dict) -> dict:
    topic = config["topic"]
    sources = config["sources"]
    all_run_results = []

    for run in range(config.get("runs", 1)):
        futures = [guarded_search(topic, s) for s in sources]
        results = [f.result() for f in futures]
        ok = [r for r in results if r["status"] == "ok"]
        all_run_results.append({"run": run + 1, "ok": len(ok), "total": len(results)})

    return {"runs": all_run_results}


result = full_production_pipeline.invoke({
    "topic": "AI production patterns",
    "sources": ["reliable_a", "flaky_b", "reliable_c"],
    "runs": 6,
})

print("\n📊 Resultados por run:")
for r in result["runs"]:
    print(f"  Run {r['run']}: {r['ok']}/{r['total']} fuentes OK")

print("\n📈 Métricas acumuladas:")
met.report()
# Output esperado (varía por el random de flaky_b):
# 📊 Resultados por run:
#   Run 1: 3/3 fuentes OK      (o 2/3 si flaky_b falló)
#   Run 2: 2/3 fuentes OK
#   Run 3: 2/3 fuentes OK      (circuit opens for flaky_b)
#   Run 4: 2/3 fuentes OK      (circuit_open — no intenta)
#   Run 5: 2/3 fuentes OK      (circuit_open)
#   Run 6: 2/3 fuentes OK      (circuit_open or half_open)
#
# 📈 Métricas acumuladas:
#   flaky_b: 6 calls, 33% ok, avg 180ms
#   reliable_a: 6 calls, 100% ok, avg 302ms
#   reliable_c: 6 calls, 100% ok, avg 201ms

Este es un pipeline production-ready completo: timeout por búsqueda, circuit breaker para fuentes inestables, y métricas para monitorear la salud del sistema.


Resumen

En esta cápsula aprendiste los 4 patterns que separan un prototipo funcional de un sistema de producción:

  • Timeouts por nodo — Previenen que una API lenta bloquee todo el pipeline. Usa asyncio.wait_for para código async o ThreadPoolExecutor para código síncrono. Sin timeout, un nodo lento congela a todos los usuarios
  • Rate limiting internoInMemoryRateLimiter de LangChain controla cuántas requests por segundo envías a cada servicio. Sin rate limiting, excedes los límites de la API y te bloquean
  • Logging estructurado — JSON con request_id, nombre del nodo, duración, y status. Correlation IDs para rastrear una request a través de todos los nodos. Sin logging structured, "algo falló" es tu mejor diagnóstico
  • Métricas de rendimiento — Latencia por nodo (avg, p50, p95), success rate, tokens consumidos. Identifican cuellos de botella y responden preguntas de negocio. Sin métricas, optimizas a ciegas

Combinados, estos patterns hacen que tu agente sea operable: puedes diagnosticar problemas, responder preguntas de negocio, y escalar con confianza. Son el equivalente a instrumentar un servidor web con logs, métricas, y health checks — estándar en la industria, pero frecuentemente olvidado en AI engineering.

Próxima cápsula: Proyecto Evolutivo — vas a aplicar retry logic, branching paralelo, y estos patterns de producción al Research Agent v1 para convertirlo en v2.


Recursos adicionales

  1. LangChain Rate Limiters — Documentación de InMemoryRateLimiter y su integración con modelos
  2. Python asyncio.wait_for — Timeouts con asyncio para código asíncrono
  3. Python logging JSON — Guía oficial de logging con formatters custom
  4. Circuit Breaker Pattern — Martin Fowler sobre circuit breakers (el artículo original)
  5. LangGraph Functional API — Referencia de @entrypoint y @task
  6. Python concurrent.futures — ThreadPoolExecutor para timeouts síncronos

Módulo 7 — LangChain & LangGraph: From Chains to Agents