Consistência de Estado
Três problemas distintos que surgem quando múltiplas instâncias competem pelo mesmo thread, quando o código evolui mas os checkpoints não, e quando efeitos externos precisam ser protegidos contra reexecução.
Três problemas com naturezas diferentes
Consistência de estado em LangGraph não é um problema único — são três problemas distintos que surgem em momentos diferentes do ciclo de vida do sistema.
O primeiro é a race condition de concorrência: dois requests chegam com o mesmo thread_id, ambos leem o checkpoint, ambos modificam o estado, o último a escrever sobrescreve o trabalho do outro. Isso pode acontecer por retry automático do cliente, por dois agentes A2A chamando o mesmo thread, ou simplesmente por um double-click do usuário.
O segundo é o efeito colateral duplicado: o LangGraph retoma um thread após falha do ECS e reexecuta um nó que já tinha feito uma reserva ou cobrança. O checkpoint garante onde o grafo estava, mas não sabe o que já aconteceu no mundo externo. Uma reserva duplicada ou cobrança dupla é dano real ao usuário.
O terceiro é a incompatibilidade de schema: você evolui o AgentState adicionando um campo, mas há checkpoints antigos no Redis que não têm esse campo. Quando o LangGraph tenta carregar um checkpoint v1 com o código v2, o comportamento é indefinido — pode ser KeyError, pode ser None silencioso, pode ser comportamento incorreto sem erro visível.
Cada problema exige um mecanismo diferente:
| Problema | Quando ocorre | Solução |
|---|---|---|
| Race condition de concorrência | Múltiplos requests com mesmo thread_id | Distributed Lock (Redis SET NX) |
| Efeito colateral duplicado | Recovery ou retry reexecuta nó | @idempotent_node com cache Redis |
| Schema incompatível | Deploy com AgentState evoluído | schema_version + migration chain |
Distributed Lock
O RedisSaver não tem locking nativo. Ele é um checkpointer — salva e lê estado, sem se preocupar com concorrência. Se dois requests chegam com o mesmo thread_id e ambos tentam executar ao mesmo tempo, os dois vão ler o mesmo checkpoint, os dois vão executar nós em paralelo partindo do mesmo estado, e o resultado final vai depender de qual termina por último e escreve o checkpoint. O trabalho do primeiro é silenciosamente descartado.
O sintoma em produção é sutil: o grafo produz resultados incorretos sem erro explícito. Você vai ver no log que ambas as execuções terminaram com sucesso, mas o estado final não reflete os dois — apenas o último.
asyncio.Lock existe apenas na memória do processo. Em produção com 3 ECS tasks do agente, cada task tem seu próprio asyncio.Lock — são locks completamente independentes que não se comunicam. Request A vai para Task 1, adquire o lock. Request B vai para Task 2, também adquire o lock — Task 2 não sabe que Task 1 tem o lock. Ambos executam em paralelo e corrompem o estado.
O erro é silencioso porque ambas as tasks acreditam que têm exclusividade. Não há exceção, não há log de conflito — apenas um resultado incorreto difícil de reproduzir porque depende de timing. Redis SET NX é a única forma de ter um lock genuinamente distribuído entre múltiplos processos que compartilham o mesmo Redis.
SET NX e Lua script para check-and-delete atômico
O lock usa dois mecanismos Redis. SET NX EX para aquisição: cria a chave apenas se não existir, com TTL para evitar locks eternos em caso de crash. Lua script para liberação: verifica se o valor ainda é o nosso token antes de deletar.
O Lua script na liberação resolve um problema sutil: sem ele, um processo poderia deletar o lock de outro. Imagine que o TTL expira enquanto o processo A ainda está executando. O processo B adquire o lock. O processo A termina e chama DEL sem verificar — deleta o lock do processo B. Processo C adquire o lock. Agora B e C estão executando em paralelo. O check-and-delete atômico via Lua garante que só quem adquiriu o lock pode liberá-lo.
import asyncio, uuid
import redis.asyncio as aioredis
from contextlib import asynccontextmanager
class LockTimeoutError(Exception): pass
class RedisThreadLock:
def __init__(self, redis: aioredis.Redis, ttl: int = 120):
self.redis = redis
self.ttl = ttl
@asynccontextmanager
async def acquire(self, thread_id: str, timeout: float = 5.0):
lock_key = f"lock:thread:{thread_id}"
# UUID único por aquisição: se o TTL expirar e outro processo adquirir,
# este processo não vai deletar o lock alheio na liberação.
lock_value = str(uuid.uuid4())
acquired = False
deadline = asyncio.get_event_loop().time() + timeout
while asyncio.get_event_loop().time() < deadline:
# SET NX EX: atômico — só cria se não existir, com TTL.
# Retorna None se a chave já existe (outro processo tem o lock).
ok = await self.redis.set(lock_key, lock_value, nx=True, ex=self.ttl)
if ok:
acquired = True
break
await asyncio.sleep(0.1)
if not acquired:
raise LockTimeoutError(
f"Thread {thread_id} está sendo processado por outra instância"
)
try:
yield
finally:
# Lua script: check-and-delete atômico.
# Verifica que o valor ainda é nosso antes de deletar.
# Sem isso: se o TTL expirou e outro processo adquiriu,
# deletaríamos o lock dele — corrompendo a exclusividade.
script = """
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
"""
await self.redis.eval(script, 1, lock_key, lock_value)
thread_lock = RedisThreadLock(redis, ttl=120)
@app.post("/agents/invoke")
async def invoke(body: InvokeRequest):
try:
async with thread_lock.acquire(body.thread_id, timeout=5.0):
return await run_agent(body.thread_id, body.input)
except LockTimeoutError:
# 409: thread já está sendo processado por outra instância.
# O cliente pode retentar — quando o lock for liberado, a próxima
# tentativa vai adquirir e processar normalmente.
raise HTTPException(
status_code=409,
detail="Este thread já está sendo processado. Tente em instantes."
)
Renovação automática para execuções longas
O TTL de 120s funciona bem para a maioria dos agentes. Mas se um agente pode levar mais tempo — processamento de documentos longos, por exemplo — o TTL expira antes da execução terminar. Outro processo adquire o lock enquanto o primeiro ainda está executando. O resultado é o mesmo que não ter lock.
A solução é renovar o TTL periodicamente enquanto a execução está em andamento. A renovação usa metade do TTL como intervalo — se o TTL é 120s, renova a cada 60s. Isso garante uma margem de 60s para detectar uma falha de renovação antes que o lock expire.
@asynccontextmanager
async def acquire_with_renewal(self, thread_id: str, timeout: float = 5.0):
lock_key = f"lock:thread:{thread_id}"
lock_value = str(uuid.uuid4())
# ... aquisição igual ao método anterior ...
async def _renew():
while True:
await asyncio.sleep(self.ttl * 0.5)
# EXPIRE não verifica o valor — mas se o TTL expirou e outro
# processo adquiriu, estamos renovando o lock alheio.
# Em cenários de alta disponibilidade, use Redlock para isso.
await self.redis.expire(lock_key, self.ttl)
renew_task = asyncio.create_task(_renew())
try:
yield
finally:
renew_task.cancel()
await self._safe_release(lock_key, lock_value)
Idempotência — o que o checkpoint não garante
Existe uma confusão comum: achar que o checkpoint do LangGraph garante que efeitos externos não serão duplicados. Não garante. O checkpoint salva o estado do grafo — as variáveis em memória — mas não sabe o que aconteceu no mundo externo.
Considere: o nó book_hotel executa, a reserva é criada na API, a resposta chega, o nó está prestes a salvar o resultado no estado — e o ECS task morre. O checkpoint salvo é o do nó anterior, antes de book_hotel começar. Quando o Recovery Consumer retoma, o LangGraph não sabe que a reserva já foi criada. Executa book_hotel de novo. Reserva duplicada.
O @idempotent_node resolve isso na camada do nó: antes de executar qualquer ação externa, verifica se já existe um resultado cacheado no Redis para aquele input específico. Se existir, retorna o resultado cacheado sem executar nada. A reserva já existe no sistema externo, e o grafo continua com o booking_id correto — sem duplicar.
A chave de idempotência é um hash do input do nó. A questão é: hash de quê exatamente? O estado LangGraph completo inclui o histórico de mensagens, que cresce a cada turno. Se você usa o estado inteiro, a mesma ação de reserva com os mesmos parâmetros em turnos diferentes vai gerar chaves diferentes — a proteção não funciona.
key_fields define exatamente quais campos determinam a unicidade da ação — geralmente os parâmetros da chamada externa. Para book_hotel, o que torna uma reserva única é hotel_id, checkin e checkout — não o histórico completo da conversa. Usar apenas esses campos garante que a mesma reserva em qualquer contexto de conversa seja reconhecida como idempotente.
Geração da chave
import hashlib, json, logging
from functools import wraps
import redis.asyncio as aioredis
logger = logging.getLogger(__name__)
IDEMPOTENCY_TTL = 60 * 60 * 24 # 24h — janela de proteção
def _make_key(thread_id: str, node_name: str, state: dict) -> str:
# sort_keys=True: garante que a mesma dict sempre gera o mesmo JSON,
# independente da ordem de inserção das chaves em Python 3.7+.
# default=str: converte datetime, UUID e outros tipos não-serializáveis.
state_bytes = json.dumps(state, sort_keys=True, default=str).encode()
state_hash = hashlib.sha256(state_bytes).hexdigest()[:16]
return f"idempotency:{thread_id}:{node_name}:{state_hash}"
Decorator completo
def idempotent_node(
redis: aioredis.Redis,
node_name: str,
key_fields: list[str] | None = None,
ttl: int = IDEMPOTENCY_TTL,
):
def decorator(func):
@wraps(func)
async def wrapper(state: dict, *args, **kwargs):
config = kwargs.get("config", {})
thread_id = config.get("configurable", {}).get("thread_id", "unknown")
relevant = (
{k: state.get(k) for k in key_fields if k in state}
if key_fields else state
)
key = _make_key(thread_id, node_name, relevant)
cached = await redis.get(key)
if cached:
logger.info("[idempotent] Cache hit: nó '%s' thread '%s'.",
node_name, thread_id)
return json.loads(cached)
result = await func(state, *args, **kwargs)
await redis.setex(key, ttl, json.dumps(result, default=str))
return result
return wrapper
return decorator
Duas camadas de proteção
O decorator protege o LangGraph — evita reexecutar a ação. Mas existe uma janela de race condition que o decorator não cobre: o request chega à API externa, ela processa, mas a resposta se perde antes de voltar. O decorator não salvou o resultado (nunca recebeu), mas a ação já aconteceu no sistema externo. Na próxima tentativa, o decorator tenta de novo — e a API processa de novo.
A solução completa usa duas camadas: o decorator no LangGraph e um idempotency_key na chamada à API externa. Stripe, Adyen, Braintree e a maioria das APIs de pagamento modernas suportam esse padrão. Com as duas camadas, mesmo que a resposta se perca, a API externa reconhece o idempotency_key na próxima tentativa e retorna o resultado já processado sem duplicar.
@traced_node("charge_card")
@idempotent_node(redis, "charge_card",
key_fields=["booking_id", "amount", "currency"])
async def charge_card_node(state: AgentState, config: dict) -> dict:
# Camada 1: @idempotent_node verifica o cache Redis.
# Se o resultado já foi salvo (execução anterior bem-sucedida),
# retorna sem chamar a API — sem chance de duplicar.
# Camada 2: idempotency_key na API externa.
# Se a resposta da API se perdeu antes de salvar no Redis,
# a próxima tentativa retorna o resultado anterior sem cobrar de novo.
charge = await payment_api.charge(
booking_id=state["booking_id"],
amount=state["amount"],
currency=state["currency"],
idempotency_key=f"{state['booking_id']}:{state['amount']}",
)
return {"charge_id": charge.id, "payment_status": "paid"}
Quais nós precisam de idempotência
| Tipo de nó | Precisa? | Motivo |
|---|---|---|
| Leitura / busca / consulta | Não | Sem efeito colateral — repetir é seguro e barato |
| Decisão / roteamento / LLM call | Opcional | LLM calls: economia de tokens; resultado pode variar entre execuções |
| POST em API externa | Sim | Efeito colateral irreversível sem rollback automático |
| Escrita em banco próprio | Sim | Duplicata causa inconsistência; depende do banco ter upsert ou não |
| Cobrança / pagamento | Sim — crítico | Cobrança dupla é dano financeiro direto e visível ao usuário |
| Reserva / booking | Sim — crítico | Reserva dupla gera custo operacional, pode causar overbooking |
| Email / notificação push | Sim | Mensagem duplicada é visível ao usuário e prejudica confiança |
Versionamento de Schema
O LangGraph não tem migração nativa de checkpoints. Quando você evolui o AgentState — adiciona um campo, muda um tipo, renomeia uma chave — os checkpoints salvos no Redis continuam com o schema antigo. Quando o código novo tenta carregar um checkpoint antigo, o comportamento depende de como você acessa o campo ausente.
O pior cenário não é uma exceção visível — é um None silencioso. Se você adiciona o campo budget: float e um nó faz state.get("budget", 0.0), o código vai funcionar com valor default em checkpoints antigos — mas o comportamento do agente vai ser diferente do esperado sem nenhum log de aviso. O bug só aparece quando você percebe que os agentes estão ignorando o orçamento do usuário.
Um campo ausente com .get() retorna None ou o default. O agente continua funcionando — só que com o comportamento errado. Não há exception, não há log de erro, não há alerta. Você só descobre quando o usuário reclama que o agente ignorou uma instrução ou o teste de negócio falha. Em produção com milhares de checkpoints antigos, isso pode afetar uma fração significativa das sessões ativas após um deploy.
Estratégias de migração
Existem três abordagens, cada uma com trade-offs diferentes:
Versionamento no thread_id (v2::thread-123): threads antigos ficam no namespace v1 e não são afetados pelo novo código. Simples de implementar, mas abandona sessões ativas no momento do deploy. Aceitável quando o negócio permite reiniciar conversas em andamento.
Campos opcionais (budget: Optional[float] = None): funciona para adições simples sem lógica de migração. O problema é que o código novo precisa lidar com None em todos os lugares que usam o campo — e é fácil esquecer um lugar. Não funciona para renomeações ou mudanças de tipo.
Migration chain com nó __coerce__: é a abordagem mais robusta. Adiciona um campo schema_version ao estado e um nó de entrada que migra o estado progressivamente até a versão atual. Um checkpoint v1 passa por v1→v2→v3→v4 antes de qualquer nó de negócio executar.
Migration chain — implementação
O nó __coerce__ é configurado como entry point do grafo. Cada vez que o grafo é invocado — novo request ou recovery — ele roda primeiro. Se o estado já está na versão atual, retorna sem modificação. Se está em uma versão anterior, aplica as migrations em sequência.
from typing import Callable
# Cada entrada: versão atual → função que produz a próxima versão.
# Funções puras: recebem o estado antigo, retornam o novo com schema_version atualizado.
MIGRATIONS: dict[int, Callable[[dict], dict]] = {
1: lambda s: {**s, "budget": 0.0, "schema_version": 2},
2: lambda s: {**s, "preferred_currency": "USD", "schema_version": 3},
3: lambda s: {**s, "loyalty_tier": "standard", "schema_version": 4},
}
CURRENT_VERSION = 4
def coerce_state(state: dict, config: dict) -> dict:
"""Nó __coerce__: migra o estado para a versão atual antes de qualquer nó de negócio."""
version = state.get("schema_version", 1) # default 1: estados sem versão são v1
while version < CURRENT_VERSION:
migration = MIGRATIONS.get(version)
if not migration:
raise ValueError(f"Sem migration para versão {version} → {version + 1}")
state = migration(state)
version = state["schema_version"]
return state
from langgraph.graph import StateGraph, END
from migrations import coerce_state
def build_graph(checkpointer):
graph = StateGraph(AgentState)
# __coerce__ é o primeiro nó — todo request e todo recovery passa por aqui.
# Garante que nenhum nó de negócio recebe estado em versão desatualizada.
graph.add_node("__coerce__", coerce_state)
graph.add_node("router", router_node)
graph.add_node("search", search_hotels_node)
graph.add_node("book", book_hotel_node)
graph.add_node("charge", charge_card_node)
graph.set_entry_point("__coerce__")
graph.add_edge("__coerce__", "router")
return graph.compile(checkpointer=checkpointer)
Remover uma entrada de MIGRATIONS significa que qualquer checkpoint com aquela versão vai lançar ValueError. Antes de remover, você precisa de garantia de que não existem mais checkpoints naquela versão no Redis. Na prática: configure o TTL dos checkpoints no RedisSaver para um período razoável (7-30 dias), aguarde esse período após o deploy da versão nova, e então remova as migrations obsoletas em um deploy subsequente. Se você usa o RedisSaver sem TTL configurado, faça um scan manual antes de remover.
Testes de consistência
import pytest, asyncio
import fakeredis.aioredis as fakeredis
from idempotency import idempotent_node
from migrations import coerce_state, CURRENT_VERSION
from distributed_lock import RedisThreadLock, LockTimeoutError
@pytest.mark.asyncio
async def test_idempotent_node_nao_duplica_em_retry():
redis = fakeredis.FakeRedis(decode_responses=True)
call_count = 0
@idempotent_node(redis, "charge_card", key_fields=["booking_id", "amount"])
async def charge(state, config):
nonlocal call_count
call_count += 1
return {"charge_id": "ch-abc"}
state = {"booking_id": "b-123", "amount": 500.0}
config = {"configurable": {"thread_id": "t-001"}}
r1 = await charge(state, config=config)
r2 = await charge(state, config=config) # simula retry do recovery
assert call_count == 1 # ação executou só uma vez
assert r1 == r2 # resultado idêntico nas duas chamadas
def test_migration_chain_v1_para_versao_atual():
state_v1 = {"schema_version": 1, "messages": [], "destination": "Lisboa"}
migrated = coerce_state(state_v1, config={})
assert migrated["schema_version"] == CURRENT_VERSION
assert migrated["budget"] == 0.0
assert migrated["preferred_currency"] == "USD"
assert migrated["loyalty_tier"] == "standard"
assert migrated["destination"] == "Lisboa" # campos originais preservados
@pytest.mark.asyncio
async def test_lock_executa_em_sequencia_nunca_paralelo():
redis = fakeredis.FakeRedis(decode_responses=True)
lock = RedisThreadLock(redis, ttl=10)
executed = []
async def task(name):
async with lock.acquire("thread-x", timeout=3.0):
executed.append(("start", name))
await asyncio.sleep(0.05)
executed.append(("end", name))
await asyncio.gather(task("A"), task("B"))
# A deve terminar completamente antes de B começar (ou vice-versa)
start_a = executed.index(("start", "A"))
end_a = executed.index(("end", "A"))
start_b = executed.index(("start", "B"))
assert end_a < start_b or executed.index(("end", "B")) < start_a
O lock distribuído via Redis SET NX, as chaves de idempotência cacheadas no Redis e os checkpoints versionados que sustentam o __coerce__ dependem diretamente de como o RedisSaver está configurado em produção — para os detalhes de clustering, TTLs e tabelas DynamoDB que armazenam esse estado, veja Memória Redis & DynamoDB.