Multi-Agent Production System · Resiliência

Resiliência

Cinco camadas de proteção — cada uma com responsabilidade distinta. A ordem em que se compõem importa tanto quanto cada implementação individual.

Circuit BreakerBulkheadRetry Cascading TimeoutBackpressureanyio

Por que resiliência em multi-agent é diferente

Em um sistema single-agent, uma falha afeta uma requisição. Você adiciona retry, talvez um timeout, e resolve. Em multi-agent, o mesmo problema tem uma dinâmica diferente: um agente downstream que começa a falhar não afeta só as requisições que dependem dele diretamente — ele afeta todos os outros agentes que compartilham os mesmos recursos.

Imagine o agente de busca começando a responder com lentidão. Cada requisição ao agente de booking que depende do agente de busca fica segurando um slot do pool de conexões do Redis, ocupando uma thread do event loop, e consumindo memória com o estado parcial do grafo LangGraph. Com 50 requisições simultâneas travadas nesse estado, o agente de booking — que estava saudável — começa a falhar por falta de recursos. É uma falha em cascata que começa em um ponto e se espalha pelo sistema inteiro.

A solução não é um mecanismo único, mas cinco camadas com responsabilidades distintas que trabalham juntas:

Sem resiliência — falha em cascata:
search-agent começa a responder lentamente
  └─► booking-agent aguarda → segura slot do bulkhead → segura conexão Redis
  └─► 50 requests simultâneos em booking-agent → todos travados
  └─► booking-agent esgota recursos → começa a falhar também
  └─► payment-agent tenta chamar booking-agent → também falha
  └─► Sistema inteiro degradado por culpa de um agente lento

Com as cinco camadas:
search-agent começa a responder lentamente
  └─► Circuit Breaker detecta 5 falhas → OPEN → rejeita imediatamente
  └─► Slots do bulkhead são liberados — não ficam presos
  └─► Retry recebe CircuitOpenError → para imediatamente
  └─► search-agent tem chance de recuperar sem pressão adicional
  └─► booking-agent e payment-agent continuam operando normalmente

Ordem de aplicação — por que importa

As cinco camadas não são independentes — a ordem em que se encadeiam define o comportamento sob pressão. A lógica central é: falhar mais rápido e mais cedo, consumindo o mínimo de recursos possível.

O Circuit Breaker vem primeiro porque é a verificação mais barata. Se o serviço downstream está claramente degradado, não faz sentido adquirir um slot do Bulkhead, iniciar uma conexão, esperar um timeout, e só então descobrir que vai falhar. O Circuit Breaker torna isso uma rejeição em microsegundos.

O Bulkhead vem segundo porque limita quantas requisições chegam ao sistema de retry. Sem isso, um serviço lento pode ter centenas de requisições acumulando retries simultâneos — cada uma consumindo recursos enquanto espera. Com o Bulkhead, há um teto fixo de concorrência independente de quantas requisições chegam.

O Retry fica por último porque só faz sentido retentar dentro de um contexto já controlado — um slot já adquirido, um circuit já verificado. Retentar é caro: cada tentativa consome tempo de CPU, conexões de rede, e potencialmente slots de outros recursos.

Request chega
  │
  ▼
Circuit Breaker — verificação em microsegundos
  ├── OPEN → CircuitOpenError imediato (sem slot, sem conexão, sem espera)
  └── CLOSED/HALF_OPEN → continua
        │
        ▼
  Bulkhead — limite de concorrência
  ├── cheio → BulkheadFullError imediato (sem retry, sem espera)
  └── slot disponível → continua
        │
        ▼
  Cascading Timeout Budget — quanto tempo resta
        │
        ▼
  Retry + Backoff
  ├── falha retentável → espera exponencial com jitter → retenta
  ├── CircuitOpenError → para imediatamente (não gasta tentativas)
  └── non_retryable → propaga imediatamente
        │
        ▼
  @idempotent_node — garante que retry não duplica efeito externo

Circuit Breaker

O nome vem de eletricidade: um disjuntor que interrompe o circuito quando detecta sobrecarga, protegendo o sistema. A analogia é precisa — quando um serviço downstream está com problema, continuar enviando requisições não ajuda a recuperação, piora. O Circuit Breaker "abre" o circuito e rejeita requisições imediatamente, dando ao serviço downstream tempo para se recuperar sem pressão adicional.

O mecanismo funciona em três estados. Em CLOSED — o estado normal — todas as requisições passam e o circuit conta as falhas consecutivas. Quando o threshold é atingido, transita para OPEN: todas as requisições são rejeitadas imediatamente com CircuitOpenError sem nem tentar. Após um período de recuperação configurável, transita para HALF_OPEN: permite uma única requisição de teste passar. Se ela tiver sucesso, volta para CLOSED. Se falhar, volta para OPEN.

Por que HALF_OPEN existe — e por que sem ele o sistema fica preso

Sem o estado HALF_OPEN, você tem dois problemas. Se o circuit nunca fechar automaticamente, o operador precisa intervir manualmente para cada falha — inviável em produção. Se fechar automaticamente direto de OPEN para CLOSED após um timeout, você assume que o serviço recuperou sem verificar — e pode imediatamente abrir de novo se ainda estiver degradado, criando um ciclo de abrir-fechar-abrir que piora o serviço downstream.

HALF_OPEN resolve isso com uma única chamada de teste: barata para o sistema, suficiente para verificar se o serviço recuperou antes de liberar tráfego completo. O success_threshold de 2 no exemplo abaixo adiciona mais segurança — exige dois sucessos consecutivos antes de assumir que o serviço está estável.

Implementação

A implementação usa asyncio.Lock para proteger as transições de estado. O lock é adquirido apenas na verificação e transição — a execução da função acontece fora do lock para não bloquear outras corrotinas enquanto aguardam I/O. Esse detalhe é importante: um lock durante I/O transformaria o circuit breaker em um gargalo ele próprio.

circuit_breaker.py python
import asyncio, logging, time
from dataclasses import dataclass
from enum import Enum
from typing import Callable, Any

logger = logging.getLogger(__name__)


class CircuitState(Enum):
    CLOSED    = "closed"
    OPEN      = "open"
    HALF_OPEN = "half_open"


@dataclass
class CircuitBreakerConfig:
    failure_threshold: int  = 5
    success_threshold: int  = 2
    recovery_timeout: float = 30.0
    expected_exception: type = Exception


class CircuitOpenError(Exception):
    """Levantada quando o circuit está OPEN.
    Não deve ser retentada — retentativas não ajudam um serviço degradado."""


class CircuitBreaker:
    def __init__(self, name: str, cfg: CircuitBreakerConfig | None = None):
        self.name           = name
        self.cfg            = cfg or CircuitBreakerConfig()
        self._state         = CircuitState.CLOSED
        self._failure_count = 0
        self._success_count = 0
        self._last_failure:  float | None = None
        self._lock          = asyncio.Lock()
        self._total_calls        = 0
        self._total_failures     = 0
        self._total_short_circuits = 0

    async def call(self, fn: Callable, *args, **kwargs) -> Any:
        # Lock apenas na verificação de estado — não durante I/O.
        # Se o lock cobrisse a execução, uma chamada lenta bloquearia
        # todas as outras corrotinas esperando para verificar o estado.
        async with self._lock:
            if self._state == CircuitState.OPEN:
                if self._should_attempt_reset():
                    self._transition_to(CircuitState.HALF_OPEN)
                else:
                    self._total_short_circuits += 1
                    raise CircuitOpenError(
                        f"[{self.name}] Circuit OPEN. Retry em {self._seconds_until_retry():.1f}s"
                    )

        self._total_calls += 1
        try:
            result = await fn(*args, **kwargs)
            await self._on_success()
            return result
        except self.cfg.expected_exception as exc:
            await self._on_failure(exc)
            raise

    async def _on_success(self):
        async with self._lock:
            if self._state == CircuitState.HALF_OPEN:
                self._success_count += 1
                if self._success_count >= self.cfg.success_threshold:
                    self._transition_to(CircuitState.CLOSED)
            elif self._state == CircuitState.CLOSED:
                # Falhas precisam ser consecutivas — qualquer sucesso reseta o contador
                self._failure_count = 0

    async def _on_failure(self, exc: Exception):
        async with self._lock:
            self._total_failures   += 1
            self._failure_count    += 1
            self._last_failure      = time.monotonic()
            if self._state == CircuitState.HALF_OPEN:
                self._transition_to(CircuitState.OPEN)
            elif self._state == CircuitState.CLOSED:
                if self._failure_count >= self.cfg.failure_threshold:
                    self._transition_to(CircuitState.OPEN)

    def _should_attempt_reset(self) -> bool:
        if self._last_failure is None: return True
        return time.monotonic() - self._last_failure >= self.cfg.recovery_timeout

    def _seconds_until_retry(self) -> float:
        if self._last_failure is None: return 0.0
        return max(0.0, self.cfg.recovery_timeout - (time.monotonic() - self._last_failure))

    def _transition_to(self, new: CircuitState):
        logger.warning("[CB:%s] %s → %s", self.name,
                       self._state.value.upper(), new.value.upper())
        self._state = new
        self._failure_count = self._success_count = 0

    def stats(self) -> dict:
        return {"state": self._state.value, "failures": self._total_failures,
                "short_circuits": self._total_short_circuits}


class CircuitBreakerRegistry:
    def __init__(self): self._breakers: dict[str, CircuitBreaker] = {}
    def get(self, name, cfg=None) -> CircuitBreaker:
        if name not in self._breakers:
            self._breakers[name] = CircuitBreaker(name, cfg)
        return self._breakers[name]
    def stats(self) -> dict:
        return {n: cb.stats() for n, cb in self._breakers.items()}

circuit_registry = CircuitBreakerRegistry()

Configuração por serviço — por que valores diferentes

Não existe uma configuração universal. Cada serviço downstream tem uma tolerância diferente a falhas baseada em dois fatores: o impacto de uma falha naquele serviço e a frequência esperada de falhas transitórias.

O payment_api tem failure_threshold=3 e recovery_timeout=60s porque uma falha de pagamento é crítica — preferimos abrir o circuit mais rápido e ser mais conservadores na reabertura. Já o search_agent, um serviço interno, tem failure_threshold=10 porque flutuações ocasionais são esperadas e não queremos interromper o sistema por falhas transitórias normais.

config.py python
CIRCUIT_CONFIGS: dict[str, CircuitBreakerConfig] = {
    "hotel_api": CircuitBreakerConfig(
        failure_threshold=5,    # abre após 5 falhas consecutivas
        success_threshold=2,    # exige 2 sucessos no HALF_OPEN antes de fechar
        recovery_timeout=30.0, # aguarda 30s em OPEN antes de tentar HALF_OPEN
    ),
    "payment_api": CircuitBreakerConfig(
        failure_threshold=3,    # mais conservador — falha de pagamento é crítica
        success_threshold=3,    # exige 3 sucessos — mais cauteloso para reabrir
        recovery_timeout=60.0, # espera mais — pagamento precisa estar estável
    ),
    "search_agent": CircuitBreakerConfig(
        failure_threshold=10,   # mais tolerante — flutuações internas são normais
        success_threshold=1,
        recovery_timeout=15.0,
    ),
}

Bulkhead — isolamento de recursos

O nome vem de engenharia naval: os compartimentos estanques de um navio que garantem que um furo no casco não afunde o navio inteiro. A ideia é a mesma — se um agente fica lento ou travado, ele só pode consumir os recursos alocados para ele, não os recursos dos outros.

Sem bulkhead, todos os agentes compartilham o mesmo pool de resources do processo Python — o event loop, as conexões Redis, a memória. Se o booking-agent começa a acumular requisições lentas, ele pode esgotar o pool de conexões Redis que o search-agent também usa. O search-agent, que estava saudável, começa a falhar por falta de conexão — não por problema próprio.

Por que anyio.CapacityLimiter e não asyncio.Semaphore

A primeira reação é usar asyncio.Semaphore — já está na stdlib, funciona com asyncio. O problema é que Semaphore não expõe métricas: não há como saber quantos slots estão em uso, quantos estão disponíveis, ou qual a capacidade máxima. Em produção, essas métricas são essenciais para alertas e debugging.

anyio.CapacityLimiter expõe available_tokens e borrowed_tokens diretamente — você pode criar um endpoint /health/bulkheads que retorna essas métricas em tempo real e configura alarmes no CloudWatch quando o uso passar de 80%. Além disso, o move_on_after do anyio tem um modelo de cancelamento mais robusto do que asyncio.wait_for para o padrão de timeout de slot.

Implementação

O bulkhead combina dois mecanismos: CapacityLimiter para limitar concorrência e move_on_after para timeout. O timeout de aquisição de slot (2s) é curto intencionalmente — se não há slot disponível em 2 segundos, o sistema está sob pressão e é melhor rejeitar rápido do que acumular requisições esperando.

bulkhead.py python
import anyio, logging
from dataclasses import dataclass
from typing import Callable, Any

logger = logging.getLogger(__name__)


@dataclass
class BulkheadConfig:
    max_concurrent: int
    timeout: float

BULKHEAD_CONFIGS = {
    "booking_agent":  BulkheadConfig(max_concurrent=10, timeout=30.0),
    "search_agent":   BulkheadConfig(max_concurrent=20, timeout=10.0),
    "payment_agent":  BulkheadConfig(max_concurrent=5,  timeout=15.0),
    "recovery_worker": BulkheadConfig(max_concurrent=3, timeout=120.0),
}

class BulkheadFullError(Exception): pass
class BulkheadTimeoutError(Exception): pass


class Bulkhead:
    def __init__(self, name: str, cfg: BulkheadConfig):
        self.name     = name
        self.cfg      = cfg
        self._limiter = anyio.CapacityLimiter(cfg.max_concurrent)

    @property
    def available_slots(self) -> int: return int(self._limiter.available_tokens)

    @property
    def in_use(self) -> int: return int(self._limiter.borrowed_tokens)

    async def run(self, fn: Callable, *args, **kwargs) -> Any:
        # Tenta adquirir slot com timeout curto.
        # Se o sistema está com todos os slots ocupados por 2s, é sinal de
        # pressão — melhor rejeitar imediatamente do que enfileirar indefinidamente.
        async with anyio.move_on_after(2.0) as slot_scope:
            async with self._limiter:
                with anyio.move_on_after(self.cfg.timeout) as exec_scope:
                    result = await fn(*args, **kwargs)
                if exec_scope.cancelled_caught:
                    raise BulkheadTimeoutError(
                        f"[{self.name}] Timeout após {self.cfg.timeout}s")
                return result
        if slot_scope.cancelled_caught:
            raise BulkheadFullError(
                f"[{self.name}] Sem slots ({self.cfg.max_concurrent} em uso)")


class BulkheadRegistry:
    def __init__(self): self._bulkheads: dict[str, Bulkhead] = {}
    def get(self, name: str) -> Bulkhead:
        if name not in self._bulkheads:
            cfg = BULKHEAD_CONFIGS.get(name, BulkheadConfig(10, 30.0))
            self._bulkheads[name] = Bulkhead(name, cfg)
        return self._bulkheads[name]
    def stats(self) -> dict:
        return {n: {"available": bh.available_slots, "in_use": bh.in_use,
                     "max": bh.cfg.max_concurrent} for n, bh in self._bulkheads.items()}

bulkhead_registry = BulkheadRegistry()

Retry com Exponential Backoff + Full Jitter

Retry sem estratégia pode ser pior do que não retentar. O problema clássico: um serviço começa a falhar às 14h00. Todos os clientes recebem erro e tentam novamente ao mesmo tempo. O serviço recebe uma rajada de tráfego no exato momento em que está mais frágil — e piora.

O backoff exponencial resolve parcialmente isso: a cada falha, a espera dobra. Mas se todos os clientes começam com o mesmo delay base e têm o mesmo padrão de falha, eles ainda podem sincronizar — todos esperando 1s, depois 2s, depois 4s, juntos. O full jitter quebra essa sincronização: cada cliente escolhe um tempo aleatório dentro da janela exponencial. Com 100 clientes, a distribuição de retries fica uniforme ao longo do tempo em vez de em picos.

Full jitter vs equal jitter — qual a diferença real

Equal jitter: delay/2 + random(0, delay/2). Garante uma espera mínima de delay/2. Full jitter: random(0, delay). Pode resultar em espera muito curta, mas tem variância maior — o que é o que queremos para dessincronizar clientes.

A recomendação da AWS para sistemas distribuídos é full jitter. A intuição é: em um sistema com muitos clientes, quanto maior a variância individual, mais suave o agregado. Equal jitter garante uma espera mínima que parece mais "justa", mas concentra os retries num intervalo menor — o oposto do que você quer sob pressão.

Implementação

O detalhe mais importante da implementação é o tratamento de CircuitOpenError: quando o circuit breaker abre durante o processo de retry, o retry para imediatamente. Não faz sentido esperar e retentar se o circuit breaker já determinou que o serviço está degradado — seria ignorar a informação mais relevante disponível.

retry.py python
import asyncio, logging, random
from dataclasses import dataclass
from typing import Callable, Any
from circuit_breaker import CircuitBreaker, CircuitOpenError

logger = logging.getLogger(__name__)


@dataclass
class RetryConfig:
    max_attempts: int       = 3
    base_delay: float       = 1.0
    max_delay: float        = 30.0
    exponential_base: float = 2.0
    jitter: bool            = True
    retryable_exceptions: tuple     = (Exception,)
    non_retryable_exceptions: tuple = ()

RETRY_CONFIGS = {
    "hotel_api":   RetryConfig(max_attempts=3, base_delay=1.0,
                               non_retryable_exceptions=(ValueError, KeyError)),
    "payment_api": RetryConfig(max_attempts=2, base_delay=2.0,
                               non_retryable_exceptions=(ValueError,)),
    "search_agent": RetryConfig(max_attempts=4, base_delay=0.5),
}


def _calculate_delay(attempt: int, cfg: RetryConfig) -> float:
    # Exponential backoff com teto: min(base * 2^attempt, max_delay)
    # Full jitter: random(0, delay) — maximiza variância para dessincronizar clientes
    delay = min(cfg.base_delay * (cfg.exponential_base ** attempt), cfg.max_delay)
    return random.uniform(0, delay) if cfg.jitter else delay


async def with_retry(fn: Callable, *args,
                      config: RetryConfig | None = None,
                      circuit_breaker: CircuitBreaker | None = None,
                      service_name: str = "unknown", **kwargs) -> Any:
    cfg = config or RetryConfig()

    for attempt in range(cfg.max_attempts):
        try:
            if circuit_breaker:
                return await circuit_breaker.call(fn, *args, **kwargs)
            return await fn(*args, **kwargs)

        except CircuitOpenError:
            # Circuit aberto: o serviço está degradado.
            # Retentar agora apenas aumentaria a pressão — propaga imediatamente.
            logger.warning("[retry:%s] Circuit OPEN. Abortando.", service_name)
            raise

        except cfg.non_retryable_exceptions as exc:
            # Erros de negócio (hotel_id inválido, valor negativo) não melhoram com retry
            logger.warning("[retry:%s] Non-retryable: %s", service_name, exc)
            raise

        except cfg.retryable_exceptions as exc:
            if attempt == cfg.max_attempts - 1:
                logger.error("[retry:%s] Esgotou %d tentativas.", service_name, cfg.max_attempts)
                raise
            delay = _calculate_delay(attempt, cfg)
            logger.warning("[retry:%s] Tentativa %d/%d. Aguardando %.2fs.",
                           service_name, attempt + 1, cfg.max_attempts, delay)
            await asyncio.sleep(delay)

Como circuit breaker, retry e idempotência interagem

Os três mecanismos operam em camadas diferentes e cobrem cenários diferentes. Entender a distinção evita tanto lacunas quanto redundâncias na proteção.

O circuit breaker protege o serviço downstream — quando ele está degradado, para de tentar. O retry protege contra falhas transitórias — timeouts de rede, instâncias reiniciando, picos momentâneos. A idempotência protege contra efeitos colaterais duplicados quando o retry ou o recovery do LangGraph reexecutam um nó.

Quando um nó do grafo chama uma API externa, as falhas do retry alimentam o circuit breaker automaticamente — cada falha que o retry registra também é contada pelo circuit breaker. Quando o circuit abre, o retry para imediatamente. E quando o LangGraph retoma o checkpoint depois de uma falha do ECS, o @idempotent_node impede que a reserva ou cobrança seja feita duas vezes.

nodes.py — composição das três camadas python
# Ordem dos decorators: traced_node (mais externo) → idempotent_node → função real
# traced_node captura o tempo total incluindo o check de idempotência
# idempotent_node verifica cache antes de executar
# dentro do nó: circuit_breaker + retry protegem a chamada HTTP

@traced_node("book_hotel")
@idempotent_node(redis, "book_hotel", key_fields=["hotel_id", "checkin", "checkout"])
async def book_hotel_node(state: AgentState, config: dict) -> dict:
    cb = circuit_registry.get("hotel_api", CIRCUIT_CONFIGS["hotel_api"])
    try:
        booking = await with_retry(
            hotel_api.create_booking,
            hotel_id=state["hotel_id"],
            checkin=state["checkin"],
            config=RETRY_CONFIGS["hotel_api"],
            circuit_breaker=cb,
            service_name="hotel_api",
        )
        return {"booking_id": booking.id, "booking_status": "confirmed"}
    except CircuitOpenError:
        return {"booking_status": "unavailable", "error": "Hotel API indisponível."}

Cascading Timeout Budget

Timeouts independentes por serviço criam um problema sutil que só aparece sob pressão real. Imagine: o usuário tem 30 segundos de paciência. O agente de booking tem timeout de 25s. Ele chama o agente de busca com timeout de 20s. O agente de busca chama uma API externa com timeout de 15s.

Se a API externa trava por 15s, o agente de busca espera 20s, o agente de booking espera 25s — e o usuário recebe a resposta em 25s, quase no limite da paciência. Agora imagine que o booking-agent tinha feito outra chamada antes da busca que levou 10s. Agora ele chama o agente de busca com timeout de 20s, mas só tem 15s restantes. O agente de busca não sabe disso — pode demorar até 20s. O resultado é que o usuário recebe timeout depois de 30s mesmo com o agente de busca respondendo dentro do seu próprio timeout.

Por que ContextVar em vez de passar o budget como parâmetro

A alternativa óbvia é passar o budget restante como parâmetro em cada chamada de função. O problema é que isso contamina a assinatura de todas as funções no caminho — nós do LangGraph, clients A2A, chamadas de API. Qualquer ponto que esqueça de passar o budget quebra a cadeia silenciosamente.

ContextVar é thread-safe, funciona nativamente com asyncio, e cada asyncio.Task herda o contexto da Task pai automaticamente. Você inicializa o budget uma vez no início do request e qualquer ponto na cadeia pode ler o budget restante sem precisar recebê-lo como parâmetro. É a solução que não exige disciplina de todos os pontos da cadeia para funcionar.

Implementação

cascade_timeout.py python
import time
from contextvars import ContextVar

_budget: ContextVar[float] = ContextVar("timeout_budget", default=30.0)
_start:  ContextVar[float] = ContextVar("request_start",  default=0.0)

def init_budget(total_seconds: float) -> None:
    """Chama no início de cada request — antes de qualquer chamada downstream."""
    _budget.set(total_seconds)
    _start.set(time.monotonic())

def remaining_budget() -> float:
    """Quanto tempo resta — lê de qualquer ponto da cadeia sem parâmetros."""
    return max(0.0, _budget.get() - (time.monotonic() - _start.get()))


# No Agente A: propaga o budget restante ao chamar o Agente B
async def call_agent_b(endpoint: str, payload: dict, token: str) -> dict:
    budget = remaining_budget()
    if budget < 1.0:
        raise TimeoutError("Budget esgotado antes de chamar Agente B")

    headers = {
        "Authorization":   f"Bearer {token}",
        "X-Timeout-Budget": str(budget),
    }
    inject(headers)  # propaga trace context junto
    async with httpx.AsyncClient(timeout=budget) as client:
        return (await client.post(endpoint, json=payload, headers=headers)).json()


# No Agente B: inicializa o budget com o valor recebido do chamador
@app.post("/invoke")
async def invoke(request: Request, body: InvokeRequest):
    budget = float(request.headers.get("X-Timeout-Budget", "10.0"))
    init_budget(budget)  # reseta o budget local para o valor herdado do chamador
    bulkhead = bulkhead_registry.get("search_agent")
    bulkhead.cfg.timeout = min(bulkhead.cfg.timeout, budget)
    return await bulkhead.run(graph.ainvoke, body.input)

Backpressure na Recovery Queue

Sem limite de tamanho, a fila de recovery é um vetor de degradação silenciosa. Imagine um evento de spot interruption que mata 10 ECS tasks simultaneamente — 10 threads são detectados como órfãos e enfileirados. O Recovery Consumer processa um por vez, cada um levando até 120s. A fila leva 20 minutos para esvaziar. Enquanto isso, novas falhas continuam enfileirando. O Redis vai crescendo indefinidamente.

O backpressure resolve isso impondo um teto: quando a fila está cheia, novos itens vão direto para a DLQ em vez de esperar. Isso força uma decisão explícita — o operador precisa intervir, ver o que está na DLQ, entender por que o sistema está sob pressão, e agir. É melhor que a fila cresça sem limite e ninguém perceba até o Redis ficar sem memória.

Por que Lua script e não LLEN + RPUSH em pipeline

A abordagem óbvia é fazer LLEN para verificar o tamanho e RPUSH para enfileirar. Isso parece correto, mas tem uma race condition: entre o LLEN e o RPUSH, outro processo pode inserir um item. Se dois OrphanScanners rodarem em paralelo — cenário comum com múltiplos ECS tasks — ambos podem verificar que a fila está em 99 itens, ambos inserem, e a fila vai a 101.

O Lua script resolve isso porque o Redis executa Lua de forma atômica: nenhuma outra operação pode acontecer entre as instruções do script. É a única forma de garantir o teto sem locks externos. Pipelines Redis não são atômicos — apenas otimizam o round-trip enviando comandos em lote, mas cada comando ainda pode ser interleaved com comandos de outros clientes.

Implementação com Lua atômico

recovery_worker.py — QueueGuard python
QUEUE_MAX_SIZE       = int(os.getenv("RECOVERY_QUEUE_MAX_SIZE", "100"))
QUEUE_WARN_THRESHOLD = float(os.getenv("RECOVERY_QUEUE_WARN_THRESHOLD", "0.7"))


class QueueGuard:
    def __init__(self, redis, dlq):
        self.redis = redis
        self.dlq   = dlq

    async def safe_enqueue(self, thread_id: str) -> bool:
        # Lua script: verifica tamanho e enfileira atomicamente.
        # Retorna 0 se fila cheia (não enfileirou), 1 se enfileirou.
        # Nenhuma outra operação pode interromper entre LLEN e RPUSH.
        script = """
        local size = redis.call('LLEN', KEYS[1])
        if size >= tonumber(ARGV[1]) then
            return 0
        end
        redis.call('RPUSH', KEYS[1], ARGV[2])
        return 1
        """
        result = await self.redis.eval(
            script, 1,
            "recovery:queue",   # KEYS[1]
            QUEUE_MAX_SIZE,      # ARGV[1]
            thread_id,           # ARGV[2]
        )

        if result == 0:
            # Fila cheia: vai para DLQ em vez de ser descartado silenciosamente.
            # O operador vai ver isso no alarme de DLQ e pode investigar
            # por que o sistema está sob pressão suficiente para encher a fila.
            logger.error("Fila cheia. Thread %s → DLQ com reason=queue_full.", thread_id)
            await self.dlq.send_to_dlq(
                thread_id, reason="queue_full",
                error=f"recovery:queue atingiu {QUEUE_MAX_SIZE} itens",
            )
            return False

        current = await self.redis.llen("recovery:queue")
        if current / QUEUE_MAX_SIZE >= QUEUE_WARN_THRESHOLD:
            logger.warning("Fila em %.0f%% da capacidade (%d/%d).",
                           current / QUEUE_MAX_SIZE * 100, current, QUEUE_MAX_SIZE)
        return True