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.