Skip to content

fix(segments): exclusão de contato invalida o cache de deletados em todos os processos (CRM-245) - #129

Merged
gomessguii merged 7 commits into
developfrom
feat/CRM-245-sinal-deletados-entre-processos
Oct 9, 2026
Merged

gomessguii merged 7 commits into
developfrom
feat/CRM-245-sinal-deletados-entre-processos

Conversation

@nickoliveira23

@nickoliveira23 nickoliveira23 commented Oct 9, 2026 •

Copy link
Copy Markdown
Contributor

Summary

  • O sinal de contato excluído saía só pelo EventEmitter2 do processo que recebe o /events/identify. Com API e workers em processos separados (RUN_MODE=api e segment-worker), o worker que recalcula os segmentos nunca o recebia: o contato excluído seguia contando por até 5 minutos, até o TTL do cache vencer.
  • Novo DeletedContactsSignalRelay, com uma instância por processo:
    • publica o sinal num canal Redis evo-flow:db<REDIS_DB>:segments:contact-deleted (pub/sub ignora o índice do banco);
    • reemite no EventEmitter2 local, marcado como relayed, o que chega de outros processos, e ignora a própria origem;
    • Redis fora do ar não segura o boot nem a ingestão: os caches voltam ao TTL;
    • loga uma linha por queda de conexão e usa disconnect() no shutdown.
  • DeletedContactsCacheService:
    • chamadas concorrentes compartilham uma única busca ao ClickHouse; antes, na janela de 15 s depois de uma exclusão, cada chamada disparava a sua;
    • cada invalidação incrementa uma geração, então uma busca iniciada antes da exclusão não grava o conjunto velho;
    • a janela continua valendo pelo início da busca.

Security

  • O payload do canal leva só o contactId e um id de processo. O canal apenas invalida cache, não transporta dado.
  • Sem endpoint novo. Mensagem malformada ou de outro canal é ignorada.

Test plan

  • npx jest --ci --maxWorkers=2 → 1150 passed, 27 skipped.
  • Spec do relay com dois processos sobre um Redis em memória: a exclusão ingerida na API invalida o cache do worker; a origem não recebe duas vezes; o que chega de fora não é republicado; falha de publish não lança.
  • Teste com TestingModule e o EventEmitterModule real: o @OnEvent publica a exclusão local uma vez e não republica a reemitida.
  • Cache:
    • concorrência: uma busca só, dentro e fora da janela;
    • busca anterior à invalidação: não grava;
    • busca que começa dentro da janela e responde depois dela: não é servida do cache.
  • Mutações em cada uma dessas peças deixam os testes vermelhos.
  • Ao vivo no compose, com um container RUN_MODE=segment-worker extra:
    • os dois processos assinaram o canal (PUBSUB NUMSUB = 2);
    • um contact.deleted em POST /api/v1/events/identify publicou {origin, contactId} no canal.

Changed Files

  • src/modules/segments/services/deleted-contacts-signal.relay.ts (novo)
  • src/modules/segments/services/deleted-contacts-cache.service.ts
  • src/modules/segments/segments-cache.module.ts
  • src/modules/segments/services/deleted-contacts-signal.relay.spec.ts (novo)
  • src/modules/segments/services/segment-canonical-event-names.spec.ts

Linked Issue

  • CRM-245

🤖 Generated with Claude Code

Summary by Sourcery

Synchronize deleted-contact cache invalidation across processes and harden cache refresh behavior against concurrent and in-flight queries.

New Features:

  • Propagate contact-deletion signals across API and segment-worker processes through Redis so segment caches are invalidated promptly.

Bug Fixes:

  • Prevent stale deleted-contact data from being re-cached after an invalidation and avoid duplicate local processing of relayed signals.

Enhancements:

  • Coalesce concurrent deleted-contact lookups into a single ClickHouse query while preserving invalidation-generation and fetch-start timing guarantees.
  • Make Redis relay failures non-blocking, with resilient connection handling, scoped channels, malformed-message filtering, and clean shutdown.

Tests:

  • Add coverage for cross-process signal delivery, loop prevention, Redis failure handling, lifecycle behavior, and concurrent cache-fetch consistency.

…sos via Redis (CRM-245)

O sinal de contato excluído saía só pelo EventEmitter2 do processo que
recebe o /events/identify. No split api/worker, o segment-worker, que
recalcula os segmentos, nunca o recebia e seguia com a lista velha até o
TTL de 5 minutos.

O DeletedContactsSignalRelay publica o sinal num canal Redis
(evo-flow:db<REDIS_DB>:segments:contact-deleted; pub/sub ignora o índice
do banco) e reemite como evento local, marcado como relayed, o que chega
dos outros processos. Cada processo ignora a própria origem. Redis fora
do ar não segura o boot nem a ingestão: os caches voltam ao TTL.

Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
… velho (CRM-245)

Na janela de 15 s depois de uma exclusão o cache não responde, e cada
chamada ia ao ClickHouse: um recálculo de N segmentos disparava N buscas
iguais. Chamadas concorrentes agora compartilham a busca em andamento.

Cada invalidação incrementa uma geração e descarta a busca em voo, então
um resultado que começou antes da exclusão volta para quem pediu mas não
é gravado por 5 minutos. O expiresAt passa a ser calculado quando a
busca termina.

Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
…de deletados (CRM-245)

O spec do relay simula dois processos sobre um Redis em memória: a
exclusão ingerida na api invalida o cache do worker, a origem não recebe
o sinal duas vezes, o que chega de fora não é republicado, mensagem
malformada ou de outro canal é ignorada e falha de publish não lança.

No cache: chamadas concorrentes fazem uma só busca, dentro e fora da
janela de bypass, e uma busca iniciada antes do sinal não grava o
conjunto velho.

Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
… fim (CRM-245)

Achado da review: com o expiresAt calculado ao fim da busca, uma busca
que começava dentro da janela de 15 s e respondia depois dela gravava por
5 minutos um conjunto lido dentro da janela, possivelmente sem a
exclusão. É a corrida que a janela existe para evitar, e a develop já
decidia pelo início.

Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
…dis (CRM-245)

Com o Redis fora, o quit() rejeitava e os clientes seguiam reconectando
depois do shutdown; disconnect() encerra na hora. O erro de conexão era
logado a cada tentativa de reconexão (40 linhas em 10 s): agora sai uma
linha por queda e outra quando a conexão volta.

Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
…y (CRM-245)

Cobre os achados da review: a busca que começa dentro da janela e
responde depois dela não é servida do cache; um TestingModule com o
EventEmitterModule real prova que o @onevent publica a exclusão local
uma vez e não republica a reemitida; shutdown usa disconnect; erro de
conexão sai uma vez por queda.

Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
@sourcery-ai

sourcery-ai Bot commented Oct 9, 2026 •

Copy link
Copy Markdown
Contributor

Reviewer's Guide

Fixes stale deleted-contact counts across separate API and segment-worker processes by relaying invalidation signals over Redis pub/sub, while hardening cache refresh behavior against concurrent requests and races with invalidation. The relay is deliberately best effort—Redis outages preserve the existing TTL fallback—and the new tests cover cross-process delivery, loop prevention, malformed input, lifecycle behavior, and cache concurrency/generation semantics.

Flow diagram for generation-safe deleted-contacts cache refresh

flowchart TD
    Request["getDeletedContacts()"] --> Hit{"Valid cached set?"}
    Hit -- Yes --> ReturnCache["Return cached set"]
    Hit -- No --> InFlight{"inFlight fetch exists?"}
    InFlight -- Yes --> Share["Share existing Promise"]
    InFlight -- No --> Fetch["fetchDeletedContactsFromClickHouse()"]
    Invalidate["invalidateCache()"] --> Generation["Increment generation"]
    Generation --> Clear["Clear cached set and expiry"]
    Fetch --> Compare{"Fetch generation matches?"}
    Compare -- Yes --> Store["Store result with fetch-start expiry"]
    Compare -- No --> ReturnFresh["Return result without caching"]
    Store --> ReturnFresh
    Share --> ReturnFresh
Loading

File-Level Changes

Change Details Files
Propagate deleted-contact cache invalidation between API and segment-worker processes through Redis pub/sub while preserving local event handling.
  • Create one publisher and one subscriber per process on a Redis-DB-namespaced channel.
  • Publish only locally ingested deletions with a process origin identifier.
  • Re-emit valid messages locally with a relayed marker, ignoring self-originated, malformed, and unrelated-channel messages.
  • Treat Redis connectivity and publish failures as best effort, with outage-aware logging and non-blocking shutdown.
src/modules/segments/services/deleted-contacts-signal.relay.ts
src/modules/segments/segments-cache.module.ts
src/modules/segments/services/deleted-contacts-signal.relay.spec.ts
Make deleted-contact cache refreshes safe and efficient under concurrent requests and invalidations.
  • Share concurrent ClickHouse fetches through a single in-flight promise.
  • Track invalidation generations so pre-invalidation results cannot repopulate the cache.
  • Keep bypass-window expiry based on fetch start time and force subsequent refreshes after the shared request settles.
src/modules/segments/services/deleted-contacts-cache.service.ts
src/modules/segments/services/segment-canonical-event-names.spec.ts

Tips and commands

Interacting with Sourcery

  • Trigger a new review: Comment @sourcery-ai review on the pull request.
  • Continue discussions: Reply directly to Sourcery's review comments.
  • Generate a GitHub issue from a review comment: Ask Sourcery to create an
    issue from a review comment by replying to it. You can also reply to a
    review comment with @sourcery-ai issue to create an issue from it.
  • Generate a pull request title: Write @sourcery-ai anywhere in the pull
    request title to generate a title at any time. You can also comment
    @sourcery-ai title on the pull request to (re-)generate the title at any time.
  • Generate a pull request summary: Write @sourcery-ai summary anywhere in
    the pull request body to generate a PR summary at any time exactly where you
    want it. You can also comment @sourcery-ai summary on the pull request to
    (re-)generate the summary at any time.
  • Generate reviewer's guide: Comment @sourcery-ai guide on the pull
    request to (re-)generate the reviewer's guide at any time.
  • Resolve all Sourcery comments: Comment @sourcery-ai resolve on the
    pull request to resolve all Sourcery comments. Useful if you've already
    addressed all the comments and don't want to see them anymore.
  • Dismiss all Sourcery reviews: Comment @sourcery-ai dismiss on the pull
    request to dismiss all existing Sourcery reviews. Especially useful if you
    want to start fresh with a new review - don't forget to comment
    @sourcery-ai review to trigger a new review!

Customizing Your Experience

Access your dashboard to:

  • Enable or disable review features such as the Sourcery-generated pull request
    summary, the reviewer's guide, and others.
  • Change the review language.
  • Add, remove or edit custom review instructions.
  • Adjust other review settings.

Getting Help

@sourcery-ai sourcery-ai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hey - I've found 3 issues

Prompt for AI Agents
Please address the comments from this code review:

## Individual Comments

### Comment 1
<location path="src/modules/segments/services/deleted-contacts-cache.service.ts" line_range="71" />
<code_context>
+            : startedAt + this.CACHE_TTL;
+        this.logger.debug(`Cached ${deletedContacts.size} deleted contacts`);
+      }
       return deletedContacts;
     } catch (error) {
       this.logger.error('Failed to fetch deleted contacts:', error);
</code_context>
<issue_to_address>
**Deleted contacts remain in segments**

When a deletion signal invalidates the cache while a segment computation is awaiting a ClickHouse fetch, `loadDeletedContacts` returns the fetched set after its generation becomes stale, so the segment computation uses the pre-deletion set and can retain the deleted contact for that run.

Check the generation before returning the set and refetch or discard it when it has become stale.

Also at `src/modules/segments/services/deleted-contacts-cache.service.ts:72`.
</issue_to_address>

### Comment 2
<location path="src/modules/segments/services/deleted-contacts-signal.relay.ts" line_range="69" />
<code_context>
+    });
+
+    for (const client of [this.publisher, subscriber]) {
+      client.connect().catch(() => undefined);
+    }
+  }
</code_context>
<issue_to_address>
**Deletion signals are dropped**

When a deletion event arrives while the publisher is connecting or reconnecting, `onContactDeletedIngested` publishes only once, and `enableOfflineQueue: false` makes `publish` reject before the publisher is ready; the relay logs and drops the signal, leaving other processes with stale deleted-contact data until their cache TTL expires.

Wait for publisher readiness or retry/queue failed publishes so each deletion signal is delivered.

Also at `src/modules/segments/services/deleted-contacts-signal.relay.ts:131`.
</issue_to_address>

### Comment 3
<location path="src/modules/segments/services/deleted-contacts-cache.service.ts" line_range="40-41" />
<code_context>
       return this.cached;
     }

+    if (this.inFlight) {
+      return this.inFlight;
+    }
+
</code_context>
<issue_to_address>
**Bypass fetch is reused**

When a fetch started during the 15-second bypass window remains pending after the window closes and another caller requests deleted contacts, `getDeletedContacts` returns the existing `inFlight` promise without checking when that fetch started, so the caller receives a result that may not include the deletion and can treat the deleted contact as active.

Check the in-flight fetch's start time against the bypass window and start a fresh fetch for callers arriving after the window closes.
</issue_to_address>

Sourcery assessment

Needs a human reviewer. 3 findings to address first, and if the relay or generation handling is wrong, deleted-contact results can remain stale or cause extra ClickHouse queries across processes, but the impact is bounded to in-memory cache state and can be repaired by clearing or allowing the cache to expire. Reverting removes the new behavior, although cache state already created before the revert does not disappear immediately.

Blocking findings: src/modules/segments/services/deleted-contacts-cache.service.ts:71, src/modules/segments/services/deleted-contacts-signal.relay.ts:69, src/modules/segments/services/deleted-contacts-cache.service.ts:41


Sourcery is free for open source - if you like our reviews please consider sharing them ✨

Comment thread src/modules/segments/services/deleted-contacts-cache.service.ts
Comment thread src/modules/segments/services/deleted-contacts-signal.relay.ts
Comment thread src/modules/segments/services/deleted-contacts-cache.service.ts
… fetch

The bypass-window comment still said every call queries ClickHouse;
concurrent calls now share one fetch. Also drops a history note from
the relay header.
@gomessguii
gomessguii merged commit 4a8c627 into develop Oct 9, 2026
7 checks passed
@gomessguii
gomessguii deleted the feat/CRM-245-sinal-deletados-entre-processos branch October 9, 2026 20:20
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants