Servir agentes de LangGraph sin langgraph-api: autorización por hilo, drenado, retención y sondas para los cuatro casos de retail

Índice

TL;DR

El servidor HTTP de LangGraph no se despliega. langgraph-api tiene licencia Elastic 2.0, pide LANGGRAPH_CLOUD_LICENSE_KEY en producción y envía metadatos a un endpoint de LangChain salvo en una variante contractual air-gapped. Lo documentó el artículo de deepagents con fichero y línea. El motor, los checkpointers y create_agent son MIT y son lo único que hace falta.

Lo que ese servidor daba y hay que escribir: la API de hilos y ejecuciones, la autorización, la cola, el drenado, la retención, las migraciones y las sondas. Son unas cuatrocientas líneas de FastAPI por caso más un CronJob y un Job. No es poco y no es un servidor de colas: es lo justo para cuatro agentes con perfiles de carga distintos.

El checkpointer no comprueba propietarios. Cualquier thread_id que llegue en la configuración se lee y se reanuda. Verificado además que un hilo inexistente devuelve estado vacío y next vacío, sin error: un identificador mal formado no se distingue de uno ajeno. La autorización es una tabla hilos con propietario, tienda y caso, y un identificador que genera el servidor.

El drenado existe y funciona. RunControl.request_drain() desde el manejador de SIGTERM hace que la ejecución en curso termine el superstep actual, guarde el checkpoint y lance GraphDrained. Verificado: el estado queda en el límite del superstep y invoke(None, config) lo reanuda desde ahí. La herramienta en vuelo en ese superstep se repite; por eso la idempotencia del caso 2.

La retención es un CronJob que borra hilos enteros con delete_thread. Es la única operación de borrado implementada en el checkpointer de Postgres 3.1.2, y la propia clase base avisa de que una poda parcial puede vaciar canales en silencio. La fecha para decidir qué borrar está dentro del JSONB de la fila de checkpoint, en checkpoint->>'ts', no en una columna.

Las migraciones corren en un Job con permiso de DDL y sin pooler en modo transacción, porque setup() crea índices con CREATE INDEX CONCURRENTLY. Los pods del agente conectan con un rol sin DDL y con prepare_threshold=0 si hay PgBouncer.

Estás aquí: el último de la serie, y el único sin agente

Este es el artículo 5 de la serie de agentes con LangGraph para retail. Los cuatro anteriores construyen agentes y los invocan desde Python. Este los pone detrás de una URL en el cluster, con todo lo que un invoke en un cuaderno no tiene: quién puede hablar con qué hilo, qué pasa cuando el pod se reinicia a medias, cuánto tiempo se guarda una conversación, y cómo sabe Kubernetes que el servicio está vivo.

La decisión de fondo la tomó el artículo de deepagents: el motor sí, el servidor no. Aquí se paga esa decisión.

La analogía: el mostrador y la trastienda

Los cuatro agentes son cuatro dependientes que ya saben hacer su trabajo. Lo que no tienen es mostrador: quién los llama, cómo se sabe que la persona al otro lado es quien dice, qué pasa con la conversación cuando el dependiente se va a comer, y quién tira los vales de hace un mes.

El servidor de LangGraph es un mostrador de alquiler, con contrato. Lo que sigue es construir uno propio, más pequeño, con solo lo que estos cuatro necesitan.

Lo que se pierde y lo que se escribe

Lo que daba langgraph-apiLo que se hace aquí
API de asistentes, hilos y ejecucionesFastAPI con cuatro rutas por caso: crear hilo, enviar mensaje, reanudar, leer estado
Autenticación y autorización por hiloJWT de Keycloak en el gateway y tabla hilos con propietario
Cola de ejecuciones y reintentosSemáforo de ejecuciones en vuelo por pod y 503 con Retry-After; lo que necesita una cola con reintentos y visibilidad va a Temporal
CronCronJob de Kubernetes que invoca el caso 3 a las seis
StudioTrazas en Langfuse y get_state_history por hilo
Streamingastream sobre SSE
Drenado en desplieguesRunControl en el SIGTERM y reanudación desde la aplicación
RetenciónCronJob con delete_thread
MigracionesJob con setup() y rol con DDL

Lo que no se pierde: el checkpointer, el almacén, la caché, las interrupciones, los subgrafos y el middleware. Todo eso es del motor.

El servicio

Un Deployment por caso, con la misma imagen y una variable de entorno que dice qué grafo carga. Los cuatro casos tienen perfiles distintos, y separarlos es lo que permite dimensionar y limitar cada uno con su clave del gateway: el asistente de sala es interactivo y de latencia baja; el de cliente es interactivo y con hilos que duran días; el de compras es una pasada diaria de minutos; el de fichas es un lote continuo.

from contextlib import asynccontextmanager
from fastapi import FastAPI
from psycopg_pool import AsyncConnectionPool
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver

@asynccontextmanager
async def lifespan(app: FastAPI):
    app.state.pool = AsyncConnectionPool(
        os.environ["PG_AGENTES"], min_size=2, max_size=8, open=False,
        kwargs={"autocommit": True, "prepare_threshold": 0},
    )
    await app.state.pool.open()
    app.state.checkpointer = AsyncPostgresSaver(app.state.pool)
    app.state.grafo = cargar_grafo(os.environ["CASO"], app.state.checkpointer)
    app.state.en_vuelo = {}            # thread_id -> RunControl
    app.state.semaforo = asyncio.Semaphore(int(os.environ.get("MAX_EJECUCIONES", "16")))
    app.state.drenando = False
    yield
    await app.state.pool.close()

app = FastAPI(lifespan=lifespan)

AsyncPostgresSaver se construye sobre el pool, no sobre una conexión. autocommit=True es lo que el checkpointer espera, y prepare_threshold=0 desactiva las sentencias preparadas de psycopg, que PgBouncer en modo transacción no soporta. setup() no se llama aquí.

Crear un hilo

@app.post("/hilos", status_code=201)
async def crear_hilo(cuerpo: NuevoHilo, usuario: Usuario = Depends(usuario_actual)):
    thread_id = str(uuid.uuid4())
    async with app.state.pool.connection() as conn:
        await conn.execute(
            "INSERT INTO hilos (thread_id, caso, propietario, tienda, creado) VALUES (%s, %s, %s, %s, now())",
            (thread_id, os.environ["CASO"], usuario.sub, usuario.tienda),
        )
    return {"thread_id": thread_id}

El identificador lo genera el servidor. El cliente no lo elige nunca, ni siquiera con un prefijo. Y la fila de hilos se escribe antes de que exista ningún checkpoint, porque es la que autoriza todo lo demás.

usuario_actual lee el JWT que Keycloak emitió y que el gateway ya validó, y saca sub, el rol y la tienda de sus claims. Cómo se configura ese realm y qué claims lleva está en el artículo de Keycloak; aquí solo se consume.

Autorizar un hilo

async def hilo_autorizado(thread_id: str, usuario: Usuario = Depends(usuario_actual)) -> Hilo:
    async with app.state.pool.connection() as conn:
        fila = await (await conn.execute(
            "SELECT thread_id, caso, propietario, tienda FROM hilos WHERE thread_id = %s", (thread_id,)
        )).fetchone()
    if fila is None or fila[1] != os.environ["CASO"]:
        raise HTTPException(404)
    hilo = Hilo(*fila)
    if hilo.propietario != usuario.sub and not (usuario.rol == "encargado" and usuario.tienda == hilo.tienda):
        raise HTTPException(404)
    return hilo

Dos detalles. Un hilo ajeno devuelve 404, no 403, para no confirmar que existe. Y la regla del encargado es la del caso 2: un encargado de la misma tienda puede leer y reanudar los hilos de sus clientes, que es lo que la cola de aprobaciones necesita.

Esto es todo lo que el checkpointer no hace. Verificado en langgraph 1.2.11 con InMemorySaver y en el código del de Postgres: get_state con cualquier thread_id devuelve lo que haya, y con uno que no existe devuelve values={} y next=() sin lanzar nada. Sin la tabla, un cliente que adivine o construya un identificador lee la conversación de otro.

Enviar un mensaje, en streaming

from langgraph.runtime import RunControl
from langgraph.errors import GraphDrained

@app.post("/hilos/{thread_id}/mensajes")
async def enviar(thread_id: str, cuerpo: Mensaje, hilo: Hilo = Depends(hilo_autorizado), usuario: Usuario = Depends(usuario_actual)):
    if app.state.drenando:
        raise HTTPException(503, headers={"Retry-After": "5"})
    if app.state.semaforo.locked():
        raise HTTPException(503, headers={"Retry-After": "2"})
    ctx = contexto_para(usuario, hilo)
    trace_id = str(uuid.uuid4())
    config = {
        "configurable": {"thread_id": thread_id},
        "metadata": {"caso": hilo.caso, "tienda": hilo.tienda, "langfuse_trace_id": trace_id},
        "callbacks": [langfuse_handler()],
        "max_concurrency": int(os.environ.get("MAX_CONCURRENCY", "8")),
    }
    control = RunControl()

    async def flujo():
        async with app.state.semaforo:
            app.state.en_vuelo[thread_id] = control
            try:
                async for modo, chunk in app.state.grafo.astream(
                    {"messages": [{"role": "user", "content": cuerpo.texto}]},
                    config=config, context=ctx, stream_mode=["messages", "updates"],
                    durability="sync", control=control,
                ):
                    if modo == "messages" and chunk[1].get("langgraph_node") == "model" and chunk[0].content:
                        yield evento("token", chunk[0].content)
                    elif modo == "updates" and "__interrupt__" in chunk:
                        await registrar_aprobacion_pendiente(thread_id, hilo, chunk["__interrupt__"][0], trace_id)
                        yield evento("interrupcion", {"interrupt_id": chunk["__interrupt__"][0].id})
                yield evento("fin", {"trace_id": trace_id})
            except GraphDrained:
                yield evento("drenado", {"reanudable": True})
            finally:
                app.state.en_vuelo.pop(thread_id, None)

    return EventSourceResponse(flujo())

stream_mode como lista devuelve tuplas (modo, contenido); verificado que ["updates", "values"] intercala los dos modos por superstep. Con messages se filtran los trozos del nodo model para el cliente, y con updates se detecta la interrupción en cuanto se produce, en vez de esperar al final para mirar __interrupt__ en el resultado. Es lo que permite escribir la cola de aprobaciones del caso 2 dentro del flujo.

durability="sync" en todo lo interactivo. Cuesta una escritura síncrona por superstep, que en Postgres local son milisegundos, y es lo que hace que el drenado y la reanudación funcionen.

Reanudar

@app.post("/hilos/{thread_id}/reanudar")
async def reanudar(thread_id: str, cuerpo: Reanudacion, hilo: Hilo = Depends(hilo_autorizado), usuario: Usuario = Depends(usuario_actual)):
    estado = await app.state.grafo.aget_state({"configurable": {"thread_id": thread_id}})
    if not estado.next:
        raise HTTPException(409, "el hilo no está esperando nada")
    if cuerpo.decisiones is not None and usuario.rol != "encargado":
        raise HTTPException(403)
    entrada = Command(resume={"decisions": cuerpo.decisiones}) if cuerpo.decisiones is not None else None
    ctx = await contexto_desde_cola(thread_id) if cuerpo.decisiones is not None else contexto_para(usuario, hilo)
    ...  # mismo flujo que enviar, con `entrada` en vez del mensaje

Dos reanudaciones distintas pasan por aquí. Con decisiones, es la aprobación del caso 2: Command(resume=...), solo encargados, y el contexto se carga de la cola porque el runtime no lo recuerda. Sin decisiones, es la reanudación tras un drenado: invoke(None) sobre el hilo, que continúa desde el último checkpoint. La comprobación de estado.next distingue un hilo que espera de uno que terminó: reanudar un hilo terminado no falla, vuelve a ejecutar desde el final, y eso es un turno del modelo que nadie pidió.

El drenado

Un Deployment con estrategia rolling mata pods con ejecuciones en curso. Sin nada más, la ejecución muere donde esté, el último checkpoint es el del superstep anterior, y el siguiente invoke sobre el hilo repite desde ahí, herramienta en vuelo incluida.

RunControl es el mecanismo cooperativo. Se crea uno por ejecución, se pasa en control=, y request_drain() desde cualquier hilo hace que el motor termine el superstep actual, escriba el checkpoint y lance GraphDrained. Verificado con un grafo de seis pasos y un drenado a mitad: invoke lanza GraphDrained con la razón, get_state muestra tres pasos hechos y next en el cuarto, y invoke(None, config) lo termina.

import signal

def manejar_sigterm(*_):
    app.state.drenando = True
    for thread_id, control in list(app.state.en_vuelo.items()):
        control.request_drain("rolling")

signal.signal(signal.SIGTERM, manejar_sigterm)

Con terminationGracePeriodSeconds por encima de la duración del superstep más largo (una llamada al modelo con su salida, unos 30 segundos en el peor caso de los casos interactivos; un lote entero de 50 referencias en el caso 3, que puede ser un minuto), el pod termina limpio. El cliente recibe el evento drenado con reanudable: true y vuelve a llamar a reanudar sin decisiones, que cae en otro pod y continúa.

Lo que el drenado no evita es que el superstep en curso se complete: si es una herramienta lenta, hay que esperarla. Y lo que no cubre es un SIGKILL por límite de memoria, donde no hay cooperación posible: ahí vuelve la re-ejecución de la herramienta en vuelo y la idempotencia por tool_call_id del caso 2.

Para el caso 3, el drenado se combina con lotes pequeños: cuanto menos dura el superstep del abanico, antes termina el pod y menos trabajo se repite si no termina.

La retención

El checkpointer de Postgres 3.1.2 tiene una operación de borrado, delete_thread. prune, copy_thread y delete_for_runs heredan de la clase base y lanzan NotImplementedError, y el docstring de prune en esa clase base explica por qué no hay que escribirse una poda parcial a mano: los canales delta guardan solo un centinela fuera de los puntos de instantánea, y quedarse con el último checkpoint de una cadena lo deja reconstruyéndose vacío sin error.

Así que se borran hilos enteros, por antigüedad, y la antigüedad se saca del JSONB:

-- Hilos del caso 'tienda' sin actividad en 7 días
SELECT c.thread_id
FROM checkpoints c
JOIN hilos h USING (thread_id)
WHERE h.caso = 'tienda' AND c.checkpoint_ns = ''
GROUP BY c.thread_id
HAVING max((c.checkpoint->>'ts')::timestamptz) < now() - interval '7 days';

checkpoints no tiene columna de fecha. El checkpoint_id es un UUID versión 6 ordenable por tiempo, y el propio checkpoint lleva ts en ISO 8601; verificado en ambos casos. El HAVING sobre ts es más legible que ordenar por el identificador y da lo mismo.

async def retencion(caso: str, dias: int, lote: int = 200):
    async with pool.connection() as conn:
        filas = await (await conn.execute(CONSULTA, (caso, dias))).fetchall()
    for (thread_id,) in filas[:lote]:
        await checkpointer.adelete_thread(thread_id)
        async with pool.connection() as conn:
            await conn.execute("UPDATE hilos SET borrado = now() WHERE thread_id = %s", (thread_id,))

Corre como CronJob cada noche, con un lote acotado para no bloquear las tablas durante minutos. delete_thread borra checkpoints, checkpoint_blobs y checkpoint_writes para todos los espacios de nombres del hilo, subgrafos incluidos. La fila de hilos no se borra: se marca, porque las trazas de Langfuse y las filas de negocio siguen apuntando a ese thread_id.

Los plazos por caso: 7 días para el asistente de sala, 90 para el de cliente (una devolución puede reclamarse), 30 para compras (las propuestas ya están en su tabla), y 7 para fichas (el resultado ya está en fichas). Son decisiones de negocio y se leen de una tabla, no del código.

El almacén (PostgresStore) es aparte y sí tiene TTL, con un barredor en proceso; en esta serie no se usa el almacén, así que no hay barredor que dimensionar.

Las migraciones

# job_migraciones.py
from langgraph.checkpoint.postgres import PostgresSaver

with PostgresSaver.from_conn_string(os.environ["PG_AGENTES_DDL"]) as saver:
    saver.setup()

Un Job de Kubernetes, con un rol que tiene CREATE en el esquema, conectado directamente a Postgres y no al pooler. setup() aplica diez migraciones idempotentes: cuatro tablas, un ALTER, tres índices con CREATE INDEX CONCURRENTLY y una columna añadida; lo lleva en checkpoint_migrations. Los índices concurrentes no pueden ir dentro de una transacción, y PgBouncer en modo transacción envuelve cada sentencia en una: falla. Directo a Postgres, una vez, antes del primer despliegue y después de cada subida de langgraph-checkpoint-postgres.

Los pods conectan con un rol que solo tiene SELECT, INSERT, UPDATE y DELETE sobre las cuatro tablas. Si alguien despliega una versión nueva del paquete sin correr el Job, el arranque falla con un error de esquema, ruidoso, en la primera petición.

Las sondas

@app.get("/salud/vivo")
async def vivo():
    return {"ok": True}

@app.get("/salud/listo")
async def listo():
    if app.state.drenando:
        raise HTTPException(503)
    try:
        async with app.state.pool.connection() as conn:
            await conn.execute("SELECT 1 FROM checkpoint_migrations LIMIT 1")
    except Exception:
        raise HTTPException(503)
    return {"ok": True, "en_vuelo": len(app.state.en_vuelo)}

La sonda de vida no toca nada externo; un pod cuya base de datos está caída no se reinicia, porque reiniciarlo no arregla la base. La de disponibilidad comprueba que hay conexión y que las migraciones existen, y devuelve 503 en cuanto empieza el drenado para que el Service deje de enviarle tráfico nuevo mientras termina el que tiene.

El gateway de inferencia no está en la sonda. Un LiteLLM caído se ve en la primera petición como un error del modelo, que ModelRetryMiddleware reintenta y que acaba en un mensaje al usuario. Meterlo en la sonda de disponibilidad haría que los cuatro servicios se marcaran como no listos a la vez, sin que ninguno pudiera hacer nada al respecto.

Capacidad y contrapresión

Cada pod corre un proceso de uvicorn con un bucle de eventos y un semáforo de ejecuciones en vuelo. El valor del semáforo y max_concurrency se fijan por caso:

CasoMAX_EJECUCIONES por podMAX_CONCURRENCYRéplicasRazón
Tienda3242Muchas conversaciones cortas; el abanico interno es pequeño
Cliente1642Menos tráfico, hilos largos con interrupciones
Compras1241Una pasada a la vez, con el abanico limitado a la clave del gateway
Fichas821Lote continuo; el paralelismo lo da el número de fichas en vuelo

Cuando el semáforo está lleno, 503 con Retry-After. No hay cola en el servicio. Las ejecuciones que necesitan cola, reintentos con retroceso y visibilidad de qué está esperando (el caso 3 si se lanzara bajo demanda, o cualquier lote largo) van a un orquestador; el artículo sobre Temporal explica cuándo compensa y cuánto cuesta.

La suma de MAX_EJECUCIONES por MAX_CONCURRENCY por réplicas de los cuatro casos es el máximo de peticiones simultáneas que pueden llegar al gateway desde los agentes, y tiene que caber en lo que se dimensionó para la flota. Con la tabla de arriba son 256 + 128 + 24 + 16.

Observabilidad

Cada petición crea un CallbackHandler de Langfuse con el trace_id que la aplicación generó, y ese trace_id viaja en metadata para que las herramientas lo escriban en las filas de negocio. Los nodos del grafo aparecen como spans anidados; con los subgrafos desactivados en los casos 3 y 4, la traza tiene la profundidad que el grafo del padre y no más.

La alternativa OTel, con LANGSMITH_TRACING_MODE=otel y el exportador OTLP del SDK de LangSmith hacia el endpoint OTel de Langfuse, correlaciona con los spans de FastAPI y del gateway. El artículo de deepagents documenta el camino y sus versiones; aquí se usa el callback por ser menos piezas, y se acepta que la traza del agente y la de infraestructura se cruzan por trace_id en vez de anidarse.

Tres métricas propias del servicio, en Prometheus: ejecuciones en vuelo por pod, 503 por contrapresión, y drenados. Las tres deberían estar cerca de cero la mayor parte del tiempo, y la primera es la que dice cuándo subir réplicas.

Un cliente por despliegue

Los cuatro servicios son de una cadena de tiendas. Si la plataforma sirviera a varias cadenas, la tabla hilos con la tienda no bastaría: el checkpointer, el almacén y la caché son compartidos por instancia, y el artículo de deepagents documenta que el backend de ficheros fija su raíz al construirse. La regla que sale de ahí es un despliegue por cliente, con su base agentes y sus claves del gateway, y el artículo sobre aislar agentes por cliente tiene el resto.

Checklist

  • langgraph-api no está en la imagen; langgraph, los checkpointers y langchain sí.
  • Un Deployment por caso, misma imagen, CASO por entorno.
  • El thread_id lo genera el servidor y la fila de hilos se escribe antes del primer checkpoint.
  • Todo acceso a un hilo pasa por hilo_autorizado; ajeno o inexistente es 404.
  • astream con ["messages", "updates"], durability="sync" y un RunControl por ejecución.
  • SIGTERM llama a request_drain en las ejecuciones en vuelo y pone la sonda de disponibilidad en 503.
  • terminationGracePeriodSeconds cubre el superstep más largo del caso.
  • reanudar comprueba estado.next y distingue aprobación de continuación.
  • CronJob de retención con delete_thread por lotes y plazos por caso en tabla.
  • Job de migraciones con rol DDL, directo a Postgres, en cada subida del paquete.
  • Rol de los pods sin DDL; prepare_threshold=0 si hay pooler.
  • Sonda de vida sin dependencias; sonda de disponibilidad con SELECT a checkpoint_migrations.
  • Semáforo por pod y 503 con Retry-After; la suma de concurrencias cabe en el gateway.

Trampas y cosas que no son lo que parecen

Un hilo inexistente no da error. get_state devuelve values={} y next=(). Un endpoint que confíe en que el checkpointer distinga hilos válidos deja pasar cualquier identificador.

Reanudar un hilo terminado ejecuta un turno. invoke(None) sobre un hilo sin next no falla; vuelve a llamar al modelo. Comprobar estado.next antes.

GraphDrained se lanza donde se invoca. En un flujo SSE eso es dentro del generador, y hay que capturarla ahí para emitir el evento; si no, el cliente ve una conexión cortada sin saber que puede reanudar.

El drenado espera al superstep. Una herramienta que tarda dos minutos retrasa la terminación dos minutos, o se lleva un SIGKILL si el período de gracia es menor. Supersteps cortos o período de gracia largo; no hay tercera opción.

La retención no tiene columna de fecha. checkpoint->>'ts' dentro del JSONB, o el orden del UUID v6. Ninguna de las dos está indexada; con millones de filas, la consulta de retención es un recorrido y se lanza de noche.

delete_thread es lo único que hay. prune lanza NotImplementedError y una poda a mano puede vaciar canales sin error. Hilos enteros.

setup() desde el pod o a través del pooler falla. Permisos o CREATE INDEX CONCURRENTLY en transacción. Job directo con DDL.

La sonda de disponibilidad que mira al gateway tumba los cuatro servicios a la vez. El gateway se ve en la petición, no en la sonda.

Cierre

Servir los cuatro agentes sin el servidor de LangGraph es una tabla, cuatro rutas, un manejador de señal, un CronJob y un Job. La tabla es la autorización que el checkpointer no tiene. Las rutas son la API de hilos reducida a lo que estos agentes usan. El manejador de señal es el drenado, que el motor sí trae y que solo hay que conectar. El CronJob es la retención, que el motor trae a medias. El Job es lo que el motor pide y no dice dónde correr.

Lo que no hay que escribir es lo que importa: el grafo, la persistencia, las interrupciones, el abanico y la salida estructurada son del motor, con licencia MIT, y corren en un cluster propio sin pedirle permiso a nadie. Eso era la tesis del artículo que abrió la vertical, y estos seis artículos son la comprobación con cuatro problemas de tienda que un responsable entiende.

Ver también

Fuentes