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.
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.
| Componente | O que faz | Loop | Onde roda |
|---|---|---|---|
| AgentHeartbeat | Renova chave com TTL — anuncia que está vivo | a cada 10s | Agent Task |
| OrphanScanner | Detecta threads sem heartbeat, enfileira | a cada 20s | Recovery Worker |
| RecoveryConsumer | Consome fila, retoma grafos do checkpoint | BLPOP contínuo | Recovery Worker |
| DLQManager | Conta tentativas, move para DLQ quando esgota | sob demanda | Recovery Worker |
| QueueMonitor | Emite métricas estruturadas para CloudWatch | a cada 30s | Recovery Worker |
As chaves Redis que sustentam o sistema:
| Chave | Tipo | TTL | Significado |
|---|---|---|---|
threads:in_progress | SET | — | Thread_ids iniciados e não finalizados |
heartbeat:{thread_id} | STRING | 30s | Task viva. Expira automaticamente se o processo morre. |
recovery:queue | LIST | — | Fila de threads para retomar. Máx 100 itens. |
recovery:attempts:{thread_id} | STRING | 24h | Contador de tentativas. Reseta automaticamente. |
recovery:dlq | LIST | — | 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.
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.
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
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.
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.
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.
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
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.
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.
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
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"}}
)