Pular para conteúdo

SQS FIFO e entrega

Por que a fila existe

DebounceQueue desacopla o POST do webhook das chamadas lentas de IA e permite repetir processamento/entrega quando a Lambda falha. A fila atual é FIFO, agrupada por telefone. O nome “Debounce” permanece por compatibilidade, mas mensagens novas são despachadas imediatamente.

Fontes: template.yaml, src/webhook_handler/app.py::_send_sqs_dispatch, src/orchestrator/app.py, src/core/integrations/dynamodb/debounce.py e src/core/application/debounce.py.

Configuração realmente versionada

Configuração Valor no SAM Efeito
Nome principal gogenetic-agent-debounce-${Environment}.fifo Isolamento por ambiente
FifoQueue true Ordem por grupo
ContentBasedDeduplication false Produtor fornece ID explícito
MessageGroupId turn.phone_number no código Mesmo telefone compartilha ordenação
MessageDeduplicationId SHA-256 de turn.message_id Evita enqueue duplicado na janela FIFO
DelaySeconds 0 Sem atraso deliberado
VisibilityTimeout 120 segundos Oculta mensagem durante processamento
Timeout do Orchestrator 90 segundos Limite total da Lambda
Retenção principal 345600 segundos, 4 dias Janela de permanência na fila
BatchSize 1 Um registro por invocação do event source
Resposta de falha ReportBatchItemFailures Lambda identifica registros que devem retornar
maxReceiveCount 3 Redrive para DLQ após repetidas recepções
Nome DLQ gogenetic-agent-debounce-dlq-${Environment}.fifo Fila de diagnóstico
Retenção DLQ 1209600 segundos, 14 dias Tempo para investigar

Esses são valores do template, não leitura da AWS. A soma de chamadas LLM, retries, resumos e integrações pode aproximar o timeout da Lambda; o teto de 10 iterações do caso de uso não garante concluir em 90 segundos.

Deduplicação durável e ordem real

O corpo contém apenas phone_number e message_id. Antes do enqueue, register_inbound() grava uma transação com receipt exclusivo (attribute_not_exists(PK)) e marcador LATEST_INBOUND. O receipt tem TTL de sete dias, além da janela de deduplicação de cinco minutos da FIFO mencionada pelo próprio código. Uma repetição do webhook não torna uma mensagem antiga “a mais recente”.

O marcador mantém pending_message_ids em ordem atômica de registro. O consumidor percorre essa lista, não ordena pelo corpo SQS nem pelo timestamp do provider. pending_inbounds() ajusta timestamps empatados/regressivos em um microssegundo para manter o histórico ordenado. Só receipts roteados para o agente ou despachados por implantação compatível podem ser drenados.

Se o envio SQS falha após registrar o receipt, o marcador dispatched_at permanece ausente. Uma repetição do webhook pode recuperar o envio. Se uma execução antiga já está drenando o telefone, agent_routed_at permite que ela encontre o receipt com segurança.

Latest-message-wins

FIFO organiza consumidores, mas não sabe que B chegou enquanto a IA calculava A. Essa decisão usa DynamoDB:

  1. Webhook registra B e atualiza latest imediatamente.
  2. A checa o guard após a IA e antes da ação durável.
  3. A obsoleta é concluída sem entregar seu texto; B recebe o histórico do cliente.
  4. Após a persistência existe outra checagem antes de enviar/reproduzir resposta.

Pedidos de contexto e encadeamento de agentes também verificam latest. Handoff e confirmação possuem cuidados adicionais: depois de uma mutação externa já iniciada, não se desfaz a sequência de cadastro. Consulte o ponto de integração.

Retry e replay

Cenário Comportamento
Corpo sem telefone ou JSON malformado bad_sqs_body; não reencaminha o registro inválido
Nenhum receipt/buffer/resposta pendente No-op idempotente; pode registrar buffer_already_flushed
Exceção de execução batchItemFailures=[{"itemIdentifier":"<ID_SQS>"}]
WhatsApp retorna falha Conserva resposta persistida e solicita retry
Retry com resposta válida Reenvia texto existente, sem chamar IA novamente
Retry após nova mensagem/takeover Marca resposta antiga como superseded e não a envia
Repetidas falhas de infraestrutura Redrive configurado para DLQ
{"batchItemFailures":[{"itemIdentifier":"sqs-exemplo-001"}]}

Um handoff técnico tratado com sucesso é resultado de aplicação, não uma falha SQS. Por isso uma resposta de fallback pode concluir o registro sem DLQ. Já uma IllegalStateTransitionError não capturada ou indisponibilidade DynamoDB pode provocar retry.

Compatibilidade com debounce anterior

DebounceBuffer usa PK=DEBOUNCE#{phone}, SK=BUFFER, expires_at e TTL de 30 segundos desde a primeira mensagem. flush_buffer() devolve texto concatenado e IDs; _process_record() ainda aceita esse formato quando não há receipts novos. PENDING_REPLY mantém uma resposta legada que precisa sobreviver à substituição do buffer. delete_if_first_message_at() só remove o buffer esperado, evitando apagar um posterior.

Divergência documental

Descrições antigas de fila Standard com espera de debounce não representam o POST atual. DEBOUNCE_WINDOW_SECONDS=0 no SAM e DelaySeconds=0 confirmam o despacho imediato. O TTL de 30 segundos pertence ao buffer legado; receipts usam sete dias.

Operação e testes

Investigue a causa antes de redrive: compare receipt, última mensagem, estado ativo e existência de resposta já enviada. Redrive é ação operacional com potencial de envio real; não é etapa automática de build da documentação. Use observabilidade e troubleshooting.

tests/unit/lambdas/test_orchestrator.py cobre falha de envio, replay, buffer legado substituído, ordenação por registro, mensagem durante IA, corrida de takeover e tombstone de debug. tests/unit/lambdas/test_webhook_handler.py cobre registro e enqueue; tests/unit/application/test_debounce.py e tests/unit/integrations/test_dynamodb_repos.py cobrem buffers/receipts. Alterações de fila devem preservar a chave por telefone, o ID do provider e a política de conclusão durável.