Contatos
Sumário
- Visão Geral do Módulo de Contatos
- Estrutura de Dados — Tabelas Principais
- Component
contact-maintenance - Importação Legada (fluxo antigo)
- Nova Importação Assíncrona (fluxo novo)
- Rotas HTTP
- Workers e Filas RabbitMQ
- Deduplicação e Estratégias de Conflito
- Integração com Jornadas
- Integração com Disparos Massivos
- Migrações de Banco de Dados
- Variáveis de Ambiente Relevantes
- Diagrama de Fluxo Comparativo
1. Visão Geral do Módulo de Contatos
O módulo de contatos é responsável por gerenciar o ciclo de vida completo dos contatos no Suite: criação, atualização, busca, importação em massa e associação com canais de comunicação, campos personalizados e jornadas.
Entidades centrais
| Entidade | Descrição |
|---|---|
contact | Registro principal do contato com nome, CPF/CNPJ e endereço |
contact_details | Dados de canal do contato (telefone, e-mail, etc.) por canal |
contact_custom_fields | Valores de campos personalizados configurados pelo cliente |
mailing_file | Arquivo CSV enviado para importação |
mailing_import_intent | Intenção de importação — controla o ciclo de vida do processamento |
mailing_import_pending_merge | Contatos em conflito aguardando resolução de merge (novo) |
Principais componentes e responsabilidades
| Componente | Responsabilidade |
|---|---|
contact-maintenance | CRUD de contatos, importação nova (processamento de chunks) |
mailing-intent-import | Importação legada — leitura de CSV e inserção contato a contato |
mailing-intent-maintenance | Criação de mailing file, intent de importação e gestão de metadados |
firing-intent-maintenance | Criação de intenção de disparo massivo (HSM, SMS, Email) |
journey-maintenance | Integração com jornadas — classificações e contatos de jornada |
2. Estrutura de Dados — Tabelas Principais
contact
Tabela central. Cada registro representa um contato único dentro do escopo de um cliente.
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
tenant_id | integer | ID do tenant (sempre 1000) |
client_id | integer | ID da empresa |
name | varchar | Nome do contato |
cpf_cnpj | varchar | CPF ou CNPJ (usado como chave de deduplicação) |
street_address | varchar | Endereço |
deleted_at | timestamp | Soft delete |
Índice único novo (refact-massive):
CREATE UNIQUE INDEX idx_contact_client_cpf_cnpj_unique
ON contact (client_id, cpf_cnpj)
WHERE cpf_cnpj IS NOT NULL AND deleted_at IS NULL;
Garante que não existam dois contatos com o mesmo CPF/CNPJ para o mesmo cliente.
contact_details
Armazena os dados de cada canal de comunicação de um contato (WhatsApp, SMS, e-mail, jornada, etc.).
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
contact_id | integer | FK para contact |
client_id | integer | ID da empresa |
channel_id | integer | ID do canal (ver enum channel-id) |
value | varchar | Valor do canal (telefone, e-mail, etc.) |
country_code | integer | DDI do número |
country | varchar | País |
deleted_at | timestamp | Soft delete |
contact_custom_fields
Valores dos campos personalizados configurados pelo cliente para cada contato.
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
contact_id | integer | FK para contact |
client_id | integer | ID da empresa |
custom_field_id | integer | FK para custom_fields |
value | text | Valor do campo |
mailing_file
Representa o arquivo CSV enviado para importação.
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
client_id | integer | ID da empresa |
file_url | varchar | URL do arquivo no S3 |
import_layout_id | integer | FK para layout de importação (opcional) |
has_sensitive_data | boolean | Dados sensíveis (LGPD) |
has_underage_data | boolean | Dados de menores (LGPD) |
data_processing_purpose | text | Finalidade do tratamento (LGPD) |
lgpd_processing_hypothesis | integer | Hipótese legal (LGPD) |
mailing_import_intent
Controla o ciclo de vida de uma importação de contatos. Cada importação tem uma intent associada.
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
client_id | integer | ID da empresa |
mailing_file_id | integer | FK para mailing_file |
conflict_strategy | integer | Estratégia de conflito (1=IGNORE/DISCARD, 2=MERGE) |
status | varchar | Status atual (ver abaixo) |
journey_id | integer | FK para jornada (opcional) |
auto_start_journey_mailing | boolean | Iniciar jornada automaticamente após importação |
process_started_at | timestamp | Início do processamento |
process_ended_at | timestamp | Fim do processamento |
fail_reason | text | Motivo de falha (se status = FAILED) |
total_chunks | integer | total de chunks publicados na fila |
processed_chunks | integer | chunks já processados (incrementado atomicamente) |
inserted_contacts_amount | integer | contatos inseridos |
updated_contacts_amount | integer | contatos atualizados |
invalid_contacts_amount | integer | contatos inválidos no CSV |
error_contacts_amount | integer | linhas que causaram erro de parsing |
ignored_contacts_amount | integer | contatos ignorados por conflito com DISCARD |
Status possíveis:
| Status | Descrição |
|---|---|
PENDING | Aguardando início do processamento |
IMPORTING | Em processamento ativo |
DONE | Concluído com sucesso |
FAILED | Falhou durante o processamento |
mailing_import_pending_merge (nova)
Tabela temporária criada pela branch refact-massive. Armazena contatos que precisam ser mesclados mas cujo lock Redis não pôde ser adquirido durante o processamento paralelo de chunks.
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
mailing_import_intent_id | integer | FK para mailing_import_intent (CASCADE DELETE) |
client_id | integer | ID da empresa |
cpf_cnpj | varchar(20) | CPF/CNPJ do contato pendente |
mapped_contact | jsonb | Payload completo do contato mapeado do CSV |
created_at | timestamp | Data de criação |
Os registros desta tabela são consumidos e removidos atomicamente (DELETE ... RETURNING) ao final da importação, quando todos os chunks já foram processados.
contact_import_layout_fields
Define o mapeamento de colunas do CSV para atributos, canais e campos personalizados de um layout de importação.
| Coluna | Tipo | Descrição |
|---|---|---|
id | bigserial | PK |
import_layout_id | integer | FK para o layout |
index | integer | Posição da coluna no CSV |
index_name | varchar | Nome da coluna no CSV |
contact_attribute_id | integer | FK para atributo do contato (nome, CPF, endereço) |
channel_id | integer | FK para canal (WhatsApp, SMS, etc.) |
custom_field_id | integer | FK para campo personalizado |
3. Component contact-maintenance
Localização: component/contact-maintenance/contact-maintenance.js
É o component central de gestão de contatos. Além de operações CRUD (busca por ID, por CPF/CNPJ, por origem), passou a concentrar toda a lógica da nova importação assíncrona.
Dependências injetadas no construtor
new ContactMaintenance(
integrationRepository, // integrações de ticket (Movidesk, etc.)
integrationMapper, // mapper de integrações
contactRepository, // acesso direto à tabela contact
clientRepository, // dados do cliente
vtexRepository, // integração VTEX (cliente 2109)
logger,
cacheRepository, // Redis — usado para dedup de importação
intentRepository, // acesso a mailing_import_intent e S3
processingQueueRepository, // publicação em filas RabbitMQ
journeyMaintenance // integração com jornadas
)
Repositórios filhos
| Arquivo | Responsabilidade |
|---|---|
data/contact-repository.js | Queries em contact, contact_details, contact_custom_fields |
data/intent-repository.js | Queries em mailing_import_intent, leitura de CSV do S3 |
data/cache-repository.js | Operações Redis: SET NX, SADD, SMEMBERS, DEL, EXPIRE |
data/processing-queue-repository.js | Publicação de mensagens no RabbitMQ |
4. Importação Legada (fluxo antigo)
Ainda em uso para importações iniciadas pelo fluxo antigo de duas etapas (upload do arquivo + criação da intent separados).
Visão geral
O fluxo legado é síncrono por contato: o worker lê o arquivo CSV, processa cada linha individualmente e insere ou atualiza o contato diretamente no banco, sem paralelismo entre chunks.
Componentes envolvidos
component/mailing-intent-import/mailing-intent-import.js— orquestrador principalcomponent/mailing-intent-import/data/contact-repository.js— inserção/atualização de contatoscomponent/mailing-intent-import/data/intent-repository.js— leitura da intent e metadadosservice/import-contact/app.js— worker consumidor da fila legada
Fila utilizada
| Fila | Produtor | Consumidor |
|---|---|---|
contact-import/process (via RABBITMQ_QUEUE) | mailing-import-api ao finalizar criação da intent | service/import-contact |
Fluxo passo a passo
1. Frontend envia POST para criar mailing file (upload CSV)
2. Frontend envia PUT com conflict_strategy, LGPD e demais campos
3. Sistema atualiza status da intent para PENDING e publica na fila contact-import/process
4. Worker import-contact consome a mensagem
5. BulkImportComponent.bulkImport() é chamado
6. Lê o CSV do S3 inteiro em memória
7. Para cada linha do CSV:
a. Mapeia colunas para atributos do contato
b. Verifica se contato já existe (por chave primária do cliente)
c. Se não existe: INSERT em contact, contact_details, contact_custom_fields
d. Se existe: UPDATE dos dados, mesclando campos conforme conflict_strategy
8. Ao final: atualiza status da intent para DONE ou FAILED
Limitações do fluxo legado
- Síncrono: todo o CSV é processado sequencialmente por um único worker
- Sem paralelismo: arquivos grandes (centenas de milhares de linhas) podem demorar horas
- Sem progresso em tempo real: não há contadores parciais durante o processamento
- Sem deduplicação concorrente: se dois arquivos do mesmo cliente são processados ao mesmo tempo, podem ocorrer conflitos de inserção
- Tags de contato acopladas ao banco: antes da refatoração, as tags eram armazenadas em tabelas relacionais (
contact_tag,contact_tag_contact) e precisavam ser informadas obrigatoriamente na intent antes do processamento. Essa obrigatoriedade foi removida na branchrefact-massive
O que mudou no fluxo legado com a refatoração
Mesmo o fluxo legado recebeu ajustes na refact-massive:
- Remoção da obrigatoriedade de tags: a validação que exigia
contact_tagsna intent foi removida. O processamento continua mesmo sem tags associadas - Remoção das filas CRC alternativas: antes existiam filas
contact-import/process-crc-1/2/3para distribuição de carga por cliente. Foram removidas em favor do novo sistema de chunks - Remoção das tabelas de tags relacionais:
contact_tag,contact_tag_contact,contact_tag_sms_intent,hsm_intent_contact_tag,mailing_import_intent_contact_tagforam dropadas. Tags passaram a ser armazenadas como JSON inline nas intenções
5. Nova Importação Assíncrona (fluxo novo)
Introduzida pela branch
refact-massive. Substitui o modelo síncrono por um pipeline assíncrono em duas fases com processamento paralelo de chunks.
Visão geral
O novo fluxo divide o processamento em duas fases independentes orquestradas por filas RabbitMQ:
- Fase 1 — Leitura e particionamento: lê o CSV do S3 linha a linha (stream), valida e mapeia os contatos, e publica chunks na fila de inserção
- Fase 2 — Inserção paralela: múltiplos workers processam os chunks em paralelo, com deduplicação em três camadas (Map, Redis, PostgreSQL)
Componentes envolvidos
| Componente | Arquivo | Responsabilidade |
|---|---|---|
ContactMaintenance | contact-maintenance.js | Orquestração completa das duas fases |
IntentRepository | contact-maintenance/data/intent-repository.js | Queries em mailing_import_intent, leitura streaming do S3 |
ContactRepository | contact-maintenance/data/contact-repository.js | Bulk insert/upsert de contatos, detalhes e custom fields |
CacheRepository | contact-maintenance/data/cache-repository.js | Locks de deduplicação no Redis |
ProcessingQueueRepository | contact-maintenance/data/processing-queue-repository.js | Publicação nos workers |
mailing-import service | service/mailing-import/mailing-import.js | Fase 1 — consome contact-import/process-csv |
import-contact service | service/import-contact/app.js | Fase 2 — consome contact-import/insert |
Fluxo passo a passo detalhado
Método processContactImportFile — detalhes internos
Arquivo: component/contact-maintenance/contact-maintenance.js
Responsável pela Fase 1. Pontos importantes:
- Usa
readlineinterface sobre o stream S3 — não carrega o CSV inteiro em memória - O tamanho do chunk é configurável via
process.env.CONTACT_IMPORT_CHUNK_SIZE(min: 100, max: 5000, padrão: 500) - Cada chunk publicado na fila contém:
intentId,clientId,chunkIndex,mappedContacts[],journeyClassificationContactId - Caso nenhum contato válido seja encontrado no CSV, a importação é finalizada imediatamente como DONE
Método processContactImportChunk — detalhes internos
Arquivo: component/contact-maintenance/contact-maintenance.js
Responsável pela Fase 2. Executado por cada chunk recebido da fila contact-import/insert.
6. Rotas HTTP
Nova rota unificada — importação + disparo HSM
API: hsm-api
Endpoint: POST /hsm-api/hsm/intent/import
Tipo: multipart/form-data
Autenticação: token JWT + permissão HSM.CREATE
Cria um mailing import E uma intenção de disparo em uma única chamada. O disparo só é enviado para a fila após a importação ser concluída (skipQueue: true).
Campos obrigatórios:
| Campo | Tipo | Descrição |
|---|---|---|
file | File (CSV) | Arquivo de contatos |
has_sensitive_data | boolean | Dados sensíveis (LGPD) |
has_underage_data | boolean | Dados de menores (LGPD) |
data_processing_purpose | string | Finalidade do tratamento |
lgpd_processing_hypothesis | integer | Hipótese legal |
conflict_strategy | MERGE | IGNORE | Estratégia de conflito |
provider_channel_id | integer | Canal de disparo |
template_id | string | ID do template |
fieldsTemplate | JSON array | Variáveis do template |
Campos opcionais: import_layout_id, has_internationalization, tags, scheduling_at (ISO 8601), cadence, cadence_interval, conditionalHsm, conditionalStatusBlocked, conditionalInterval, hsmTag, flowActionId, check_contacts_in_blacklist, file_url, headerContent
Validação: schema Zod post-hsm-intent-import.schema.js aplicado antes do handler. Erros retornam HTTP 400 com array [{ field, message }].
Nova rota de importação pura
API: mailing-import-api
Endpoint: POST /mailing-import-api/intent/import
Tipo: multipart/form-data
Autenticação: token JWT + permissão CONTACT_IMPORT.CREATE
Importação de contatos sem disparo vinculado.
Campos obrigatórios: file, has_sensitive_data, has_underage_data, data_processing_purpose, lgpd_processing_hypothesis, conflict_strategy
Campos opcionais: import_layout_id, has_internationalization, journey_id, auto_start_journey_mailing (obrigatório se journey_id for informado)
Validação: schema Zod post-intent-import.schema.js. Inclui .refine() que valida a dependência entre journey_id e auto_start_journey_mailing.
Middleware de validação (validate-schema.js)
Localização: api/middleware/commonjs/validate-schema.js
Factory que vincula o applyResult de cada API ao middleware de validação Zod, garantindo que erros de validação usem o padrão ResultValidation do projeto.
// Em cada arquivo de rotas:
const validateSchema = createValidateSchema(applyResult);
const validateSchemaMiddleware = {
postHsmIntentImport: (req, res, next) => validateSchema(
PostHsmIntentImportSchema,
(req) => ({ file: req.file, user_id: req.user.id, client_id: req.user.client_id })
)(req, res, next)
};
O getExtra injeta campos que não estão no req.body (arquivo do multer e dados do token JWT). Esses campos têm precedência sobre o body em caso de chave duplicada, garantindo que client_id e user_id venham sempre do token.
7. Workers e Filas RabbitMQ
service/mailing-import
Responsabilidade: Fase 1 da nova importação + importações legadas + outros imports (blacklist, tabulation, user, break, tag)
| Fila consumida | Prefetch | Handler |
|---|---|---|
RABBITMQ_QUEUE (legado) | 1 | bulkImportComponent.bulkImport() |
contact-import/process-csv | 1 | contactMaintenance.processContactImportFile() |
RABBITMQ_BLACKLIST_IMPORT_QUEUE | 1 | blackList.processBlacklistCsv() |
RABBITMQ_TABULATION_IMPORT_QUEUE | 1 | tabulation.processTabulationCsv() |
RABBITMQ_USER_IMPORT_QUEUE | 1 | userAuthorization.processUserCsv() |
RABBITMQ_BREAK_IMPORT_QUEUE | 1 | breakMaintenance.processBreakCsv() |
RABBITMQ_TAG_IMPORT_QUEUE | 1 | tagMaintenance.processTagCsv() |
Reconexão automática: 5000ms (RABBIT_RECONNECT_INTERVAL_MS)
service/import-contact
Responsabilidade: Fase 2 da nova importação + importações legadas de contato individual
| Fila consumida | Prefetch | Handler |
|---|---|---|
RABBITMQ_QUEUE (legado) | 20 | bulkImportComponent.importContact() |
contact-import/insert | 20 | contactMaintenance.processContactImportChunk() |
RABBITMQ_BLACKLIST_INSERT_QUEUE | 20 | blacklist.processBlacklistBatch() |
O prefetch de 20 permite processamento paralelo de chunks — 20 chunks simultâneos por instância do worker.
Filas de resiliência — contact-import/insert
O sistema usa três filas configuradas com DLX para implementar retry com delay e dead letter:
| Fila | Descrição |
|---|---|
contact-import/insert | Fila principal de processamento. DLX aponta para contact-import/insert-retry |
contact-import/insert-retry | Fila sem consumer. TTL de 30s — após expirar, a mensagem volta para contact-import/insert via DLX. x-death.count é incrementado a cada ciclo |
contact-import/insert-dead-letter | Destino final de chunks descartados após esgotar tentativas. Usado para inspeção manual |
Fluxo de retry:
Tentativas 1-3:
nack(false, false) → contact-import/insert-retry (30s) → contact-import/insert
x-death.count incrementado a cada ciclo
Tentativa 4 (deathCount >= 3):
discardContactImportChunk() → atualiza banco, finaliza importação se for o último chunk
sendToQueue('contact-import/insert-dead-letter') → mensagem preservada para inspeção
ack()
Tratamento de erros na fila contact-import/insert:
O consumer usa o mecanismo nativo x-death do RabbitMQ (configurado via DLX) para controlar retries:
- Erro crítico — tentativas restantes:
nack(false, false)sem requeue → mensagem vai paracontact-import/insert-retry(TTL 30s) → volta paracontact-import/insertcomx-death.countincrementado - Erro crítico — limite de 3 tentativas atingido: chama
discardContactImportChunk()para registrar o descarte e finalizar a importação, depois publica diretamente emcontact-import/insert-dead-letterpara inspeção →ack - Erro não-crítico:
ack(descartado após log) - JSON inválido (exceção de parse): manda payload bruto para
contact-import/insert-dead-letter→ack - Exceção no catch:
finallycom variávelresolvedgarante que o msg sempre recebenack(false, false)de fallback mesmo que o próprio handler de erro falhe
Método discardContactImportChunk: método público do ContactMaintenance que incrementa error_contacts_amount e chama #finalizeContactImportIfProcessedAllChunks — garantindo que a importação possa finalizar normalmente mesmo quando chunks são descartados por excesso de falhas.
Filas RabbitMQ — resumo completo
| Fila | Direção | Descrição |
|---|---|---|
contact-import/process-csv | Produzido por: rota HTTP → Consumido por: mailing-import | Fase 1: sinaliza para processar CSV de uma intent |
contact-import/insert | Produzido por: mailing-import → Consumido por: import-contact | Fase 2: chunk de contatos mapeados para inserção |
contact-import/insert-retry | Produzido por: import-contact (nack) → sem consumer | Fila de delay 30s para retry da fase 2 |
contact-import/insert-dead-letter | Produzido por: import-contact → sem consumer ativo | Chunks descartados após 3 tentativas, para inspeção |
contact-import/process | Produzido por: mailing-import-api (legado) → Consumido por: import-contact | Legado: importação completa de uma intent |
hsm/importing-intent | Produzido por: import-contact → Consumido por: firing-hsm | Disparo de HSM/Email após importação concluída |
sms/firing-intent | Produzido por: import-contact → Consumido por: firing-sms | Disparo de SMS após importação concluída |
8. Deduplicação e Estratégias de Conflito
A nova importação implementa deduplicação em três camadas sequenciais para garantir integridade mesmo com múltiplos workers processando chunks em paralelo.
Estratégias de conflito
| Valor no formulário | Valor no banco | Comportamento |
|---|---|---|
IGNORE | 1 | Contato com CPF/CNPJ já existente é ignorado |
MERGE | 2 | Contato com CPF/CNPJ já existente tem seus dados mesclados |
Atenção: o formulário envia
"IGNORE"(nomenclatura legada). O novo sistema (ContactImportConflictStrategyById) usa"DISCARD"internamente, mas ambos mapeiam para o ID1.
Camada 1 — Deduplicação intra-chunk (em memória)
Método: #deduplicateMappedContactsInChunk
Remove duplicatas dentro do próprio chunk usando um Map com chave cpf_cnpj.
- MERGE: contatos duplicados dentro do chunk são mesclados — atributos, custom fields e canais são combinados
- DISCARD: o segundo registro com o mesmo CPF/CNPJ é descartado
Camada 2 — Deduplicação inter-chunks (Redis)
Método: #deduplicateMappedContactsAcrossChunks
Garante que o mesmo CPF/CNPJ não seja processado por dois chunks diferentes ao mesmo tempo.
Mecanismo: para cada contato único do chunk, tenta adquirir um lock Redis usando SET NX EX (set if not exists with expiration):
chave: contact-import:intent:{intentId}:dedup:{sha1(intentId|clientId|cpf_cnpj)}
TTL: 120 segundos
O TTL de 120 segundos é suficiente para cobrir a janela de processamento simultâneo de até 20 chunks em paralelo. O banco garante a integridade final via ON CONFLICT — o Redis é uma camada de otimização para evitar roundtrips desnecessários ao banco durante o processamento paralelo. Ao final da importação, todas as chaves são limpas via #clearContactImportDedupKeys.
Um índice de todas as chaves criadas é mantido em:
contact-import:intent:{intentId}:dedup:index (Redis Set)
- Lock adquirido: contato segue para inserção
- Lock não adquirido + MERGE: contato vai para
mailing_import_pending_merge - Lock não adquirido + DISCARD: contato é ignorado
Se um chunk falhar e for requeueado, os locks adquiridos por ele são removidos imediatamente no catch do processContactImportChunk, evitando que o retry do mesmo chunk encontre seus próprios locks ainda ativos e encaminhe os contatos incorretamente para pending_merge.
Camada 3 — Verificação no banco (PostgreSQL)
Método: #splitMappedContactsByExistingCpfCnpj
Após as duas primeiras camadas, consulta o banco para verificar se o CPF/CNPJ já existe na tabela contact:
SELECT id, cpf_cnpj, name, street_address
FROM contact
WHERE client_id = $1
AND deleted_at IS NULL
AND regexp_replace(cpf_cnpj::text, '\D', '', 'g') = ANY($2)
- Não existe: vai para
toInsert— será inserido viabulkUpsertContacts - Existe + MERGE: vai para
toMerge— dados são mesclados com o existente - Existe + DISCARD: descartado
Escrita em lote (#applyContactImportBatchWrite)
Após a deduplicação, a escrita no banco é feita em lote:
-
contactRepository.bulkUpsertContacts— INSERT em lote emcontactcomON CONFLICT (client_id, cpf_cnpj) DO UPDATEouDO NOTHINGconforme a estratégia. Retorna[{ id, cpf_cnpj }]para mapeamento dos IDs. -
contactRepository.bulkUpsertContactCustomFields— UPDATE + INSERT em transação explícita emcontact_custom_fieldspara todos os campos personalizados do lote. -
contactRepository.bulkInsertContactDetails— INSERT emcontact_detailsevitando duplicatas via LEFT JOIN, para todos os canais do lote. -
journeyMaintenance.bulkCreateJourneyContacts— INSERT em lote emjourney_contactpara contatos com canal JORNADA, caso a importação esteja vinculada a uma jornada.
Pending merges (finalização)
Método: #processContactImportPendingMerges
Executado somente quando todos os chunks foram processados. Recupera os registros da tabela mailing_import_pending_merge via DELETE ... RETURNING (atômico), resolve os contatos existentes no banco e aplica as mesclagens pendentes.
Finalização atômica
Método: #finalizeContactImportIfProcessedAllChunks
O contador processed_chunks é incrementado atomicamente via UPDATE ... RETURNING. Quando processed_chunks >= total_chunks, o worker que fez o contador bater é responsável pela finalização.
Para evitar que dois workers finalizem ao mesmo tempo (race condition), a atualização final usa:
UPDATE mailing_import_intent
SET status = 'DONE', process_ended_at = $2
WHERE id = $1 AND status NOT IN ('DONE', 'FAILED')
RETURNING id
O process_ended_at é passado como parâmetro (new Date() do Node.js) em vez de usar now() SQL, garantindo consistência de timezone independente da configuração do banco.
Se o RETURNING não retornar nenhuma linha, outro worker já finalizou — o bloco de finalização é ignorado silenciosamente.
9. Integração com Jornadas
Quando a importação está vinculada a uma jornada (journey_id informado), o fluxo adiciona etapas específicas:
Durante a Fase 1 (processContactImportFile)
Antes de iniciar a leitura do CSV, cria uma journey_classification_contact com status AWAITING_IMPORT:
await this.journeyMaintenance.createJourneyContactClassification({
journeyId,
description, // nome do arquivo CSV
clientId,
mailingImportIntentId: intentId,
autoStartJourneyMailing,
status: 'AWAITING_IMPORT'
});
O journeyClassificationContactId gerado é incluído em cada chunk publicado na fila.
Durante a Fase 2 (#applyContactImportBatchWrite)
Para cada contato com canal JORNADA (channel_id = canal de jornada), é criado um registro em journey_contact via bulk insert:
await this.journeyMaintenance.bulkCreateJourneyContacts(
journeyContactsPayload, // [{ detailId, rawRow }]
classificationId,
clientId,
columnNames
);
O bulkCreateJourneyContacts no journey.js:
- Busca todos os
contact_idviagetContactIdsByContactDetailIds(uma query só) - Insere todos os
journey_contactem um únicoINSERT ... VALUES (...), (...)
Na finalização
Quando todos os chunks são processados, a classificação de jornada tem seu status atualizado:
auto_start_journey_mailing = true→ statusACTIVEauto_start_journey_mailing = false→ statusPAUSED
A quantidade total de contatos (contact_amount) é atualizada na journey_classification_contact.
Status enum jourey_classification_contact_status_type:
| Status | Descrição |
|---|---|
AWAITING_IMPORT | Novo — Aguardando conclusão da importação do arquivo |
ACTIVE | Jornada ativa e processando contatos |
PAUSED | Jornada pausada aguardando ação manual |
DONE | Jornada concluída |
CANCELED | Jornada cancelada |
EXPIRED | Jornada expirada |
10. Integração com Disparos Massivos
Rota unificada (nova)
Quando a importação é feita via POST /hsm-api/hsm/intent/import, o disparo é criado junto com o mailing import mas com skipQueue: true:
const firingResult = await this.prepareAndSendHsm({
...intentParams,
mailingFileId: importData.mailing_file_id,
skipQueue: true // não publica na fila agora
});
A intenção de disparo (hsm_intent, sms_intent ou hsm_intent_email) fica registrada no banco com mailing_import_intent_id preenchido, mas sem ser enfileirada.
Disparo automático após importação
Quando todos os chunks são processados, o método #triggerFiringProcesses é chamado. Ele consulta as intenções de disparo associadas via getIntentWithAssociatedFiringIntents e publica nas filas correspondentes:
| Intenção | Fila publicada |
|---|---|
hsm_intent_id | hsm/importing-intent com { id } |
email_intent_id | hsm/importing-intent com { id, channelId: EMAIL } |
sms_intent_id | sms/firing-intent com { id } |
11. Migrações de Banco de Dados
Todas as migrações novas ficam em migrations/main/ com timestamp no nome.
Novas tabelas e índices
| Migração | Operação |
|---|---|
1777312831150_create-mailing-import-pending-merge-table | Cria mailing_import_pending_merge com FK CASCADE para mailing_import_intent |
1777320448002_add-unique-index-contact-cpf-cnpj | Cria índice único (client_id, cpf_cnpj) em contact WHERE cpf_cnpj IS NOT NULL AND deleted_at IS NULL |
Alterações em tabelas existentes
| Migração | Operação |
|---|---|
1779993009297_add-mailing-import-intent-counters | Adiciona total_chunks, processed_chunks, error_contacts_amount em mailing_import_intent |
1780433681830_add_awaiting_import_to_jourey_classification_contact_status_type | Adiciona valor AWAITING_IMPORT ao enum jourey_classification_contact_status_type |
Tabelas removidas (simplificação de tags)
| Migração | Tabela removida | Motivo |
|---|---|---|
1780422922776_fix-contact-tag-contact-fk | contact_tag_contact | Tags passaram a ser JSON inline |
1780429134129_drop-contact-tag-sms-intent-table | contact_tag_sms_intent | Idem |
1780433643708_drop-hsm-intent-contact-tag-table | hsm_intent_contact_tag | Idem |
1780433665880_drop-mailing-import-intent-contact-tag-table | mailing_import_intent_contact_tag | Idem |
1780433681821_drop-contact-tag-table | contact_tag | Idem |
12. Variáveis de Ambiente Relevantes
| Variável | Serviço | Descrição | Padrão |
|---|---|---|---|
CONTACT_IMPORT_CHUNK_SIZE | mailing-import | Tamanho de cada chunk publicado na fila | 500 |
AWS_BUCKET_NAME | mailing-import, import-contact | Bucket S3 para leitura dos CSVs | — |
AWS_REGION | mailing-import, import-contact | Região AWS do S3 | us-east-1 |
RABBITMQ_QUEUE | import-contact | Fila legada de importação de contatos | — |
RABBITMQ_BLACKLIST_INSERT_QUEUE | import-contact | Fila de inserção em blacklist | — |
13. Diagrama de Fluxo Comparativo
Fluxo Legado
Fluxo Novo (refact-massive)
Resumo das diferenças principais
| Aspecto | Fluxo Legado | Fluxo Novo |
|---|---|---|
| Modelo | Síncrono, contato a contato | Assíncrono, chunks em paralelo |
| Paralelismo | Nenhum | Até 20 chunks simultâneos por worker |
| Progresso | Não disponível em tempo real | Contadores parciais no banco |
| Deduplicação | Apenas no banco (query por contato) | 3 camadas: Map + Redis + banco |
| Tags de contato | Obrigatórias, tabelas relacionais | Removidas das tabelas, JSON inline |
| Disparo vinculado | Processo separado, polling | Automático ao fim da importação |
| Entrada HTTP | 2 chamadas (upload + configuração) | 1 chamada multipart unificada |
| Validação de entrada | Manual (if/else no component) | Schema Zod na camada HTTP |
| Jornadas | Processamento posterior | Status AWAITING_IMPORT durante importação |
| Rollback S3 | Não implementado | Arquivo removido do S3 em caso de falha |