Pular para conteúdo

Projeção assíncrona Inbox → CRM

Fonte de verdade e componentes

O DynamoDB guarda o atendimento autoritativo. O CRM recebe uma projeção eventual dos dados e do estado, produzida por CrmSyncProcessor em src/core/application/crm_sync.py. src/crm_sync/app.py é o handler agendado; src/core/integrations/dynamodb/crm_sync.py contém CrmSyncRepository, desired_projection() e target_stage().

DynamoSessionRepository grava sessão e intenção CRM na mesma transação, incluindo fechamento e mudança de assunto. Um turno ultrapassado ou conflito de posse não deve gerar intenção externa. A alocação inicial NEW sem deal é ignorada; só um resultado confirmado libera a primeira projeção. Hash SHA-256 do conteúdo desejado evita enfileirar revisões semanticamente iguais.

Este é um worker sobre DynamoDB, não uma fila SQS específica de CRM. O transporte remoto continua sendo o adapter HubSpot compatível.

Item persistido e operação

Cada sessão tem PK=CRM_SYNC#<session_id>, SK=META. Campos relevantes:

Campo Finalidade
desired stage, terminal, payload contact e payload deal
revision / confirmed_revision Hash da projeção desejada e última revisão confirmada
contact_id / deal_id IDs externos duráveis, independentes do ponteiro de sessão ativa
status PENDING, SYNCED, BLOCKED ou UNCERTAIN
etag Token de compare-and-swap para alteração concorrente
attempts, last_error, next_attempt_at Diagnóstico e agendamento de retry
pending_since, updated_at Idade operacional e última mudança
lease_token, lease_until Posse temporária da execução por até cinco minutos
creation_started, creation_token Marcador durável antes do POST de negócio
GSI1PK, GSI1SK CRM_SYNC#<status> e <next_attempt_at>#<session_id>

get() usa leitura consistente. prepare() prepara item e etag esperado para a transação da sessão. change() relê e tenta CAS até oito vezes; esgotamento lança crm_sync_concurrent_update. list_status() pagina a GSI; problems() agrega os três estados não confirmados. Novas revisões preservam IDs, lease e marca de criação em andamento; atualizar uma sessão antiga não deve apagar o vínculo remoto.

Processamento e garantias

flowchart TD
    A[PENDING e prazo atingido] --> B[Claim com lease de 5 minutos]
    B --> C{Criação iniciada sem deal_id?}
    C -->|Sim| U[UNCERTAIN]
    C -->|Não| D{Terminal sem deal?}
    D -->|Sim| S[Confirmar sem criar]
    D -->|Não| E[Validar stage e upsert contato]
    E --> F{Deal conhecido?}
    F -->|Sim| G[PATCH]
    F -->|Não| H[Persistir marca antes do POST]
    H --> I[POST uma vez]
    I --> J[Salvar deal_id]
    G --> K[Confirmar revisão processada]
    J --> K
    K --> L{Revisão ainda é atual?}
    L -->|Sim| M[SYNCED]
    L -->|Não| N[PENDING para convergir]

Um resultado antigo confirma apenas a revisão que processou. Se a sessão mudar durante o PATCH/POST, a próxima execução converge para a nova intenção. Persistir a marca antes de criar evita POST duplicado entre workers ou após crash; a consequência conservadora é UNCERTAIN mesmo se o crash ocorreu antes do envio.

Não existe suposição de suporte remoto a Idempotency-Key. O código não compara o conteúdo real do CRM a cada execução: SYNCED significa que uma revisão foi aceita, não que ninguém editou o CRM depois. Encerramento sem deal e sem criação em andamento confirma sem criar um card terminal.

Estágios, timeout e falhas

target_stage() usa novo_lead em NEW, em_triagem nos estados ativos de agente, atendimento_humano em HANDOFF_PENDING/HUMAN, e crm_terminal_stage ou finalizado nos terminais. O encerramento por abandono grava perdido.

O handler lê os cinco HUBSPOT_STAGE_* diretamente, sem fallbacks das propriedades de Settings. Stage ausente ou terminal igual a qualquer stage ativo resulta em BLOCKED com stage_configuration_invalid. Configure valores explícitos.

Falha Estado e recuperação automática
Temporária de busca/PATCH ou rejeição explícita 429 na criação PENDING; 1, 2, 4, 8, 16, 32, depois 60 minutos
Permanente, incluindo PATCH 404 BLOCKED; não cria deal substituto
Resultado incerto de POST UNCERTAIN; não repete POST
creation_started=true sem ID após lease UNCERTAIN com creation_result_unknown
Exceção ao trabalhar uma sessão crm_sync_worker_failed; marca/lease persistem

CrmSyncFunction em template.yaml usa schedule de 1 minuto, timeout Lambda 120 s e concorrência reservada 1. run(max_seconds=70) para de iniciar novos itens ao alcançar o orçamento; uma operação já iniciada ainda pode levar tempo adicional. O HTTP deste worker limita HUBSPOT_TIMEOUT_SECONDS a 1–10 s e usa retry_max=0, independentemente de HUBSPOT_RETRY_MAX.

HUBSPOT_ENABLED=false pausa o handler e a geração de novas intenções, sem apagar pendências. A construção separada lê apenas CRM e DynamoDB: não chama load_settings() nem precisa resolver segredo OpenAI/WhatsApp para funcionar.

Observação e intervenção

/atendimento/crm mostra pendências, bloqueios e incertezas a operadores autenticados; não é uma tela de reparo. scripts/diagnose_crm_sync.py consulta a projeção em modo de leitura. Com snapshot externo e mapa de estágios, também compara o último estado observado do CRM. Nenhum diagnóstico sem consulta/snapshot externo comprova convergência remota atual.

python scripts/diagnose_crm_sync.py --table <TABELA_ALVO>

Métricas GoGenetic/CRM: Pending, Blocked, Uncertain, OldestPendingSeconds, por FunctionName. A idade é calculada sobre todos os itens retornados por problems(), inclusive bloqueados/incertos. Alarmes cobrem Blocked/Uncertain ≥ 1, idade > 900 s e erro de Lambda; notificações dependem de AlertsSnsTopicArn.

Para BLOCKED, confira configuração, código HTTP e existência do ID. Para UNCERTAIN, confronte horário e dados no CRM antes de liberar criação. A rotina não oferece comando de apply: uma recuperação exige manutenção controlada do item com etag recém-lido, preservando revisão/intenção/IDs e ajustando status e índice juntos. A especificação operacional detalhada está em docs/crm-sync-operacao.md. Esta documentação não autoriza saneamento histórico ou deploy.

Testes e evolução

tests/unit/application/test_crm_sync.py cobre ID reutilizado, revisão invariável, concorrência entre workers, fechamento durante POST/PATCH, transação com turno ultrapassado, 404 sem recriação, resultado incerto, stage terminal inválido e projeções de sessões independentes. tests/unit/lambdas/test_crm_sync.py cobre handler, métricas e flag desabilitado.

Ao alterar a projeção, preserve hashing determinístico, transação da sessão, etag, token de criação e confirmação por revisão. Ao trocar provider, implemente a classificação transient/permanent/uncertain e os contratos do CRM; preservar apenas o formato HTTP não é suficiente para evitar duplicações.