Multi-Agent Production System · Consistência

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.

Distributed LockIdempotência Schema MigrationRedis SET NXLua Script

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:

ProblemaQuando ocorreSolução
Race condition de concorrênciaMúltiplos requests com mesmo thread_idDistributed Lock (Redis SET NX)
Efeito colateral duplicadoRecovery ou retry reexecuta nó@idempotent_node com cache Redis
Schema incompatívelDeploy com AgentState evoluídoschema_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.

Por que asyncio.Lock não resolve — e o erro seria completamente silencioso

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.

distributed_lock.py python
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.py — uso do lock por endpoint python
@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.

distributed_lock.py — renovação python
@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.

Por que key_fields e não o estado inteiro no hash

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

idempotency.py python
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

idempotency.py — decorator python
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.

nodes.py — duas camadas de idempotência python
@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 / consultaNãoSem efeito colateral — repetir é seguro e barato
Decisão / roteamento / LLM callOpcionalLLM calls: economia de tokens; resultado pode variar entre execuções
POST em API externaSimEfeito colateral irreversível sem rollback automático
Escrita em banco próprioSimDuplicata causa inconsistência; depende do banco ter upsert ou não
Cobrança / pagamentoSim — críticoCobrança dupla é dano financeiro direto e visível ao usuário
Reserva / bookingSim — críticoReserva dupla gera custo operacional, pode causar overbooking
Email / notificação pushSimMensagem 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.

O sintoma mais insidioso: comportamento silenciosamente incorreto

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.

migrations.py python
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
graph.py — __coerce__ como entry point python
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)
Quando é seguro remover código de migration antigo

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

tests/unit/test_consistencia.py python
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
Ver também

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.