Multi-Agent Production System · Lifecycle

Lifecycle de Execução

Do início de um thread LangGraph até sua conclusão — incluindo o que acontece quando o ECS task morre no meio do caminho e como o sistema detecta, retoma e trata falhas repetidas.

LangGraphRedisSaverHeartbeat OrphanScannerRecoveryConsumerDLQ

O problema do thread zumbi

O RedisSaver do LangGraph resolve a durabilidade do estado — quando um nó termina, o checkpoint é salvo no Redis. Mas durabilidade de estado não é o mesmo que durabilidade de execução. Se o ECS task morre no meio de um nó, o checkpoint do nó anterior está salvo, mas ninguém sabe que aquele thread precisa continuar.

O resultado é um thread zumbi: o checkpoint existe no Redis, o grafo está em um estado intermediário válido, mas nenhum processo está executando nem vai executar. Do ponto de vista do usuário, a requisição simplesmente não termina. Não há erro, não há timeout, não há log de falha — o thread simplesmente some.

Sem lifecycle:
t=0   booking-agent Task A executa thread-123, node_b em progresso
t=15  ECS mata Task A por spot interruption
      └─► checkpoint do node_a está no Redis ✓
      └─► node_b em execução: sem checkpoint salvo ✗
      └─► nenhum processo sabe que thread-123 existe
      └─► thread-123 fica suspenso indefinidamente — thread zumbi

Com lifecycle:
t=15  Task A morre → heartbeat para de ser renovado
t=45  TTL do heartbeat expira → chave some do Redis
t=60  OrphanScanner detecta: in_progress mas sem heartbeat → enfileira
t=61  RecoveryConsumer: ainvoke(None, thread-123) → retoma do node_a
      └─► LangGraph lê checkpoint → executa node_b → grafo completa

A solução é um sistema de presença ativa: o agente anuncia continuamente que está vivo via heartbeat com TTL curto. Quando morre, para de anunciar. O OrphanScanner detecta a ausência e aciona o recovery. O estado continua no Redis — o que falta é apenas o processo que vai continuar.

Componentes e Redis como intermediário

Um equívoco comum ao pensar nesse sistema: imaginar que o Recovery Worker precisa se comunicar com o Agent Task. Não precisa — e não pode, porque o Task morreu. O Redis é o único intermediário. Todos os componentes leem e escrevem no Redis, sem comunicação direta entre processos.

ComponenteO que fazLoopOnde roda
AgentHeartbeatRenova chave com TTL — anuncia que está vivoa cada 10sAgent Task
OrphanScannerDetecta threads sem heartbeat, enfileiraa cada 20sRecovery Worker
RecoveryConsumerConsome fila, retoma grafos do checkpointBLPOP contínuoRecovery Worker
DLQManagerConta tentativas, move para DLQ quando esgotasob demandaRecovery Worker
QueueMonitorEmite métricas estruturadas para CloudWatcha cada 30sRecovery Worker

As chaves Redis que sustentam o sistema:

ChaveTipoTTLSignificado
threads:in_progressSET—Thread_ids iniciados e não finalizados
heartbeat:{thread_id}STRING30sTask viva. Expira automaticamente se o processo morre.
recovery:queueLIST—Fila de threads para retomar. Máx 100 itens.
recovery:attempts:{thread_id}STRING24hContador de tentativas. Reseta automaticamente.
recovery:dlqLIST—Threads que esgotaram tentativas — intervenção manual.
recovery:dlq:meta:{thread_id}HASH—reason, error, attempts, failed_at para diagnóstico.

Heartbeat — TTL como sinal de vida

A elegância do heartbeat está em usar o mecanismo de expiração do Redis como detector de falha. Em vez de precisar de um sistema de monitoramento externo que detecta quando um processo morre, você simplesmente não precisa fazer nada — se o processo para de renovar a chave, o Redis a expira automaticamente. A ausência da chave é o sinal.

A janela de TTL de 30s com renovação a cada 10s não é arbitrária. A razão de 3x entre renovação e TTL garante que flutuações normais — GC pauses, pico de CPU, latência de rede momentânea — não causem falsos positivos. Se o processo está vivo mas com alguma sobrecarga, ainda vai renovar dentro de 30s. Se morreu, a chave some em no máximo 30s — que é a janela máxima de invisibilidade de um thread órfão antes de ser detectado.

Por que o TTL curto não causa falsos positivos no scanner

O OrphanScanner roda a cada 20s. O heartbeat tem TTL de 30s renovado a cada 10s. Se o scanner rodar e o heartbeat tiver sido renovado há 25s, a chave ainda existe — não é falso positivo. Para um falso positivo acontecer, o scanner precisaria rodar exatamente na janela entre o processo parar de renovar e o TTL expirar. Mesmo que isso aconteça, o RecoveryConsumer verifica state.next antes de retomar — se o grafo está em estado terminal (terminou normalmente), ignora e limpa o registro sem reexecutar.

Implementação

O heartbeat usa asynccontextmanager para garantir que a chave seja deletada quando o agente termina normalmente. A deleção explícita na saída do contexto é importante: sem ela, o OrphanScanner esperaria até o TTL expirar (30s) antes de considerar o thread como encerrado — o que criaria um falso positivo para threads que terminaram com sucesso.

agent_runner.py python
import asyncio
import redis.asyncio as aioredis
from contextlib import asynccontextmanager

HEARTBEAT_TTL      = 30   # segundos — TTL da chave no Redis
HEARTBEAT_INTERVAL = 10   # segundos — frequência de renovação (razão 3x com TTL)


class AgentHeartbeat:
    def __init__(self, redis: aioredis.Redis, thread_id: str, task_id: str):
        self.redis     = redis
        self.thread_id = thread_id
        self.task_id   = task_id  # ECS task ID — útil para correlacionar logs
        self._key      = f"heartbeat:{thread_id}"

    async def _beat(self):
        # SETEX: SET + EXPIRE atômico.
        # Valor = task_id — permite saber qual ECS task estava rodando o thread.
        await self.redis.setex(self._key, HEARTBEAT_TTL, self.task_id)

    @asynccontextmanager
    async def running(self):
        async def _loop():
            while True:
                await self._beat()
                await asyncio.sleep(HEARTBEAT_INTERVAL)

        task = asyncio.create_task(_loop())
        try:
            await self._beat()  # bate imediatamente ao entrar — não espera 10s
            yield
        finally:
            task.cancel()
            # Deleção explícita: OrphanScanner não vai detectar este thread como órfão.
            # Sem isso, o TTL de 30s causaria um falso positivo para threads normais.
            await self.redis.delete(self._key)


async def run_agent(thread_id: str, input: dict, ecs_task_id: str):
    await redis.sadd("threads:in_progress", thread_id)

    heartbeat = AgentHeartbeat(redis, thread_id, ecs_task_id)
    async with heartbeat.running():
        result = await graph.ainvoke(
            input, config={"configurable": {"thread_id": thread_id}},
        )

    await redis.srem("threads:in_progress", thread_id)
    return result

OrphanScanner — detecção de threads órfãos

A lógica do scanner é direta: um thread é órfão quando está em threads:in_progress (foi iniciado e não terminou) mas não tem heartbeat ativo (o processo que o executava morreu). A combinação dos dois critérios é necessária — um thread pode estar em in_progress sem heartbeat legitimamente por alguns segundos durante a deleção explícita do contexto do heartbeat.

O scanner usa SMEMBERS para ler todos os thread_ids em andamento. Para conjuntos grandes (acima de alguns milhares), considere SSCAN com cursor — mas na prática o número de threads simultâneos raramente justifica isso. O custo de SMEMBERS é O(n) no tamanho do set, e com o Bulkhead limitando concorrência a 10-20 threads por agente, o set tende a ser pequeno.

Implementação

recovery_worker.py — OrphanScanner python
class OrphanScanner:
    def __init__(self, redis):
        self.redis = redis

    async def scan_once(self) -> list[str]:
        thread_ids: set[str] = await self.redis.smembers("threads:in_progress")
        if not thread_ids:
            return []

        orphans = []
        for thread_id in thread_ids:
            alive = await self.redis.exists(f"heartbeat:{thread_id}")
            if not alive:
                orphans.append(thread_id)
        return orphans

    async def enqueue_orphans(self, orphans: list[str], guard: QueueGuard) -> None:
        for thread_id in orphans:
            # Verifica se já está na fila — o scanner pode rodar enquanto
            # um ciclo anterior ainda está sendo processado.
            already = await self.redis.lpos("recovery:queue", thread_id)
            if already is not None:
                continue

            await self.redis.srem("threads:in_progress", thread_id)
            enqueued = await guard.safe_enqueue(thread_id)
            if enqueued:
                logger.warning("Órfão detectado e enfileirado: %s", thread_id)

    async def run(self, guard: QueueGuard) -> None:
        logger.info("OrphanScanner iniciado (intervalo=20s)")
        while True:
            try:
                orphans = await self.scan_once()
                if orphans:
                    await self.enqueue_orphans(orphans, guard)
            except Exception:
                logger.exception("Erro no OrphanScanner")
            await asyncio.sleep(20)

Recovery — como ainvoke(None) retoma o grafo

Este é o mecanismo central do recovery e merece atenção especial porque parece mágico até você entender o que está acontecendo. Quando você chama graph.ainvoke(None, config={"configurable": {"thread_id": "thread-123"}}), o LangGraph não começa um novo grafo do zero. Ele chama o checkpointer para buscar o último estado salvo daquele thread_id e continua a partir dos nós pendentes em state.next.

O None como input significa exatamente isso: não há novo input do usuário — use o estado que está no checkpoint. Se você passasse um dict como input, o LangGraph mesclaria com o estado existente, o que é útil para retomadas com novo contexto do usuário, mas não para recovery transparente onde você quer continuar exatamente de onde parou.

O que state.next significa e por que é a peça central

Cada checkpoint do LangGraph salva não só o estado atual (state.values), mas também os próximos nós a executar (state.next). Quando um nó termina, o grafo atualiza state.next com os nós subsequentes antes de salvar o checkpoint. Quando o processo morre no meio de um nó — antes de o nó terminar — o último checkpoint salvo é o do nó anterior, e state.next ainda aponta para o nó que morreu no meio.

O recovery retoma exatamente do nó que não terminou. Não do começo do grafo, não do começo do nó anterior — do ponto exato onde o checkpoint foi salvo pela última vez. É por isso que a idempotência dos nós (doc 03) é crítica: o nó pode ser reexecutado do início, então precisa ser seguro reexecutar.

graph.ainvoke(None, config={"configurable": {"thread_id": "thread-123"}})
  │
  ▼
RedisSaver.get(thread_id="thread-123")
  └─► lê último checkpoint do Redis
  └─► deserializa: state.values + state.next + state.metadata
  │
  ▼
state.next = ("booking_agent",)  ← nós pendentes
  │
  ▼
LangGraph executa apenas os nós em state.next
  └─► booking_agent_node(state, config)
  └─► salva novo checkpoint após cada nó
  └─► atualiza state.next com os próximos
  │
  ▼
state.next = ()  ← estado terminal → ainvoke retorna

RecoveryConsumer

O consumer usa BLPOP em vez de polling com sleep. BLPOP bloqueia até haver um item na fila — sem gastar CPU verificando uma fila vazia repetidamente. O timeout=5 evita que o processo fique bloqueado indefinidamente em caso de problema de conexão com o Redis — depois de 5 segundos sem item, retorna None e o loop continua.

recovery_worker.py — RecoveryConsumer python
class RecoveryConsumer:
    def __init__(self, redis, graph, dlq: "DLQManager"):
        self.redis = redis
        self.graph = graph
        self.dlq   = dlq

    async def consume_one(self, thread_id: str) -> None:
        config = {"configurable": {"thread_id": thread_id}}
        try:
            state = await self.graph.aget_state(config)

            if not state:
                # Checkpoint não existe — pode ter expirado ou nunca ter sido salvo.
                # Reenfileirar não vai resolver: vai para DLQ diretamente.
                await self.dlq.send_to_dlq(
                    thread_id, reason="no_checkpoint",
                    error="Checkpoint ausente ou expirado",
                )
                return

            if not state.next:
                # Estado terminal — grafo já terminou.
                # Provavelmente falso positivo do scanner. Limpa sem reexecutar.
                logger.info("Thread %s em estado terminal. Falso positivo.", thread_id)
                await self.redis.delete(f"recovery:attempts:{thread_id}")
                return

            logger.info("Retomando %s nos nós: %s", thread_id, list(state.next))
            await asyncio.wait_for(
                self.graph.ainvoke(None, config=config),
                timeout=120,
            )
            await self.redis.delete(f"recovery:attempts:{thread_id}")
            logger.info("Thread %s retomado com sucesso.", thread_id)

        except asyncio.TimeoutError:
            await self.dlq.handle_failure(
                thread_id, reason="timeout", error="Execução excedeu 120s")
        except Exception as exc:
            await self.dlq.handle_failure(
                thread_id, reason="exception", error=str(exc))

    async def run(self) -> None:
        while True:
            try:
                # BLPOP bloqueia até item disponível — sem polling busy.
                # timeout=5: retorna None após 5s sem item — evita bloqueio eterno.
                result = await self.redis.blpop("recovery:queue", timeout=5)
                if result:
                    _, thread_id = result
                    await self.consume_one(thread_id)
            except Exception:
                logger.exception("Erro no RecoveryConsumer")
                await asyncio.sleep(1)

Dead-Letter Queue — o problema do retry infinito

Sem DLQ, um thread que falha consistentemente entra em loop eterno: falha → reenfileira → falha → reenfileira. O sistema desperdiça recursos tentando recuperar algo que nunca vai funcionar — um checkpoint corrompido, um nó com bug, uma API permanentemente indisponível. E o pior: o operador não tem visibilidade disso acontecendo.

O DLQManager resolve com uma lógica simples: cada thread tem um contador de tentativas com TTL de 24h. Após atingir o máximo configurado, vai para a DLQ com metadados completos. O TTL de 24h no contador é intencional — se o problema for corrigido (deploy de um fix, API voltando), o contador reseta automaticamente e o thread pode ser reenfileirado sem intervenção manual para limpar o contador.

Dois caminhos distintos para a DLQ — e por que isso importa

handle_failure() é para erros recuperáveis — timeout, exception genérica. Conta tentativas e decide entre reenfileirar ou DLQ. A suposição é que o problema pode ser transitório.

send_to_dlq() direto é para erros irrecuperáveis — checkpoint ausente é o exemplo principal. Se o checkpoint não existe, não importa quantas vezes você tentar: o grafo não tem estado para retomar. Reenfileirar apenas gastaria as tentativas antes de ir para a DLQ de qualquer forma. Ir direto economiza tempo e deixa o motivo mais claro nos metadados.

DLQManager

recovery_worker.py — DLQManager python
from datetime import datetime, timezone

MAX_ATTEMPTS = int(os.getenv("RECOVERY_MAX_ATTEMPTS", "3"))
ATTEMPTS_TTL = 60 * 60 * 24  # 24h — reseta automaticamente


class DLQManager:
    def __init__(self, redis):
        self.redis = redis

    async def get_attempts(self, thread_id: str) -> int:
        val = await self.redis.get(f"recovery:attempts:{thread_id}")
        return int(val) if val else 0

    async def increment_attempts(self, thread_id: str) -> int:
        key      = f"recovery:attempts:{thread_id}"
        attempts = await self.redis.incr(key)
        if attempts == 1:
            # Define TTL apenas na primeira falha — INCR não aceita EX.
            await self.redis.expire(key, ATTEMPTS_TTL)
        return attempts

    async def send_to_dlq(self, thread_id: str, reason: str, error: str = "") -> None:
        attempts = await self.get_attempts(thread_id)
        meta = {
            "thread_id": thread_id,
            "reason":    reason,
            "error":     error[:500],
            "attempts":  str(attempts),
            "failed_at": datetime.now(timezone.utc).isoformat(),
        }
        pipe = self.redis.pipeline()
        pipe.rpush("recovery:dlq", thread_id)
        pipe.hset(f"recovery:dlq:meta:{thread_id}", mapping=meta)
        pipe.delete(f"recovery:attempts:{thread_id}")
        await pipe.execute()
        logger.error("Thread %s → DLQ após %d tentativas. reason=%s",
                     thread_id, attempts, reason)

    async def handle_failure(self, thread_id: str, reason: str, error: str = "") -> None:
        attempts = await self.increment_attempts(thread_id)
        if attempts >= MAX_ATTEMPTS:
            await self.send_to_dlq(thread_id, reason=reason, error=error)
        else:
            logger.warning("Thread %s falhou (%d/%d). Reenfileirando.",
                           thread_id, attempts, MAX_ATTEMPTS)
            await self.redis.rpush("recovery:queue", thread_id)

    async def list_dlq(self) -> list[dict]:
        ids    = await self.redis.lrange("recovery:dlq", 0, -1)
        result = []
        for tid in ids:
            meta = await self.redis.hgetall(f"recovery:dlq:meta:{tid}")
            result.append(meta or {"thread_id": tid})
        return result

Requeue manual

Após identificar e corrigir o problema que causou as falhas — um bug no nó corrigido em deploy, uma API externa voltando, um checkpoint corrompido removido manualmente — o operador pode mover threads da DLQ de volta para a fila principal. O requeue reseta o contador de tentativas para dar uma nova chance limpa.

recovery_worker.py — DLQManager.requeue_from_dlq python
async def requeue_from_dlq(self, thread_id: str) -> bool:
    pos = await self.redis.lpos("recovery:dlq", thread_id)
    if pos is None:
        logger.warning("Thread %s não encontrado na DLQ.", thread_id)
        return False

    # Pipeline: remove da DLQ, limpa meta, reseta contador, reenfileira.
    # Tudo atomicamente — não deixa estado parcial.
    pipe = self.redis.pipeline()
    pipe.lrem("recovery:dlq", 1, thread_id)
    pipe.delete(f"recovery:dlq:meta:{thread_id}")
    pipe.delete(f"recovery:attempts:{thread_id}")
    pipe.rpush("recovery:queue", thread_id)
    await pipe.execute()
    logger.info("Thread %s removido da DLQ e reenfileirado.", thread_id)
    return True

Por que o Recovery Worker precisa ser um ECS Service separado

A primeira ideia costuma ser rodar o OrphanScanner no mesmo container do agente — como uma corrotina em background com asyncio.gather. O problema é estrutural: o Recovery Worker existe exatamente para recuperar o agente quando ele morre. Se os dois estão no mesmo container, quando o container morre, o worker morre junto e não pode fazer nada.

A separação em ECS Services distintos resolve isso. O Recovery Worker tem desired_count=1 e o ECS garante que sempre há uma instância rodando, independente do que aconteça com os containers dos agentes. Mesma imagem Docker, entrypoint diferente via command no task definition — sem manter dois Dockerfiles separados.

Graceful shutdown no SIGTERM

Quando o ECS drena um container para um deploy, ele envia SIGTERM. Sem tratamento, o processo termina imediatamente — interrompendo qualquer agente em execução no meio de um nó. Com graceful shutdown, o container para de aceitar novos requests e aguarda os em andamento terminarem antes de sair.

app.py — graceful shutdown python
from contextlib import asynccontextmanager

_active_tasks: set[asyncio.Task] = set()
_shutting_down = False


@asynccontextmanager
async def lifespan(app):
    yield
    global _shutting_down
    _shutting_down = True
    if _active_tasks:
        logger.info("Aguardando %d task(s) ativas...", len(_active_tasks))
        # timeout=25 — stop_timeout do ECS task definition deve ser 120s,
        # deixando margem para overhead de shutdown após as tasks terminarem.
        await asyncio.wait(_active_tasks, timeout=25.0)

app = FastAPI(lifespan=lifespan)


@app.post("/agents/invoke")
async def invoke(body: InvokeRequest):
    if _shutting_down:
        # 503 durante shutdown — ALB roteia para outro target saudável.
        raise HTTPException(status_code=503, detail="Serviço em manutenção.")
    task = asyncio.current_task()
    _active_tasks.add(task)
    try:
        return await run_agent(body.thread_id, body.input)
    finally:
        _active_tasks.discard(task)

Testes do lifecycle

tests/integration/test_recovery.py python
import pytest
from unittest.mock import AsyncMock, MagicMock
import fakeredis.aioredis as fakeredis
from recovery_worker import OrphanScanner, RecoveryConsumer, DLQManager, QueueGuard

@pytest.fixture
def redis(): return fakeredis.FakeRedis(decode_responses=True)

@pytest.fixture
def dlq(redis): return DLQManager(redis)


@pytest.mark.asyncio
async def test_scanner_detecta_orfao(redis):
    # Thread em andamento sem heartbeat → deve ser detectado como órfão
    await redis.sadd("threads:in_progress", "t-001")
    orphans = await OrphanScanner(redis).scan_once()
    assert "t-001" in orphans


@pytest.mark.asyncio
async def test_scanner_ignora_thread_vivo(redis):
    # Thread com heartbeat ativo → não é órfão
    await redis.sadd("threads:in_progress", "t-002")
    await redis.setex("heartbeat:t-002", 30, "task-xyz")
    orphans = await OrphanScanner(redis).scan_once()
    assert "t-002" not in orphans


@pytest.mark.asyncio
async def test_dlq_apos_max_attempts(redis, dlq):
    for i in range(3):  # MAX_ATTEMPTS = 3
        await dlq.handle_failure("t-003", reason="exception", error=f"err {i}")

    items = await redis.lrange("recovery:dlq", 0, -1)
    assert "t-003" in items
    meta  = await redis.hgetall("recovery:dlq:meta:t-003")
    assert meta["reason"] == "exception"
    assert meta["attempts"] == "3"


@pytest.mark.asyncio
async def test_consumer_limpa_attempts_em_sucesso(redis, dlq):
    # Simula uma falha anterior — tentativa 1
    await dlq.increment_attempts("t-004")

    mock_graph = MagicMock()
    mock_graph.aget_state = AsyncMock(return_value=MagicMock(next=("node_b",)))
    mock_graph.ainvoke    = AsyncMock(return_value={"ok": True})

    await RecoveryConsumer(redis, mock_graph, dlq).consume_one("t-004")

    # Após sucesso: contador deve ser zerado, ainvoke chamado com None
    assert await dlq.get_attempts("t-004") == 0
    mock_graph.ainvoke.assert_called_once_with(
        None, config={"configurable": {"thread_id": "t-004"}}
    )