Aller au contenu

Réconciliation des compteurs de funding (Redis ↔ PG)

La task de réconciliation des compteurs protège les collectes (fundings) contre les incohérences entre le compteur rapide et la réalité des achats.

Elle s'exécute chaque nuit et garantit que, quand une propriété approche des 100 % vendus, le nombre de bricks "comptées comme vendues" correspond exactement aux achats primaires réellement enregistrés et actifs en base.

Pourquoi cette task existe

Pendant une collecte, deux systèmes comptent les bricks vendues :

  • Redis (compteur rapide) : mis à jour en temps réel à chaque création d'achat, réservation, decline ou refund — sauf les declines no-bricks-left du worker d'assignation (compteur volontairement non décrémenté).
  • Postgres (source de vérité) : les lignes de la table primary_purchase. Seuls les achats dont le statut n'est pas declined ou refunded comptent vraiment.

Dans la très grande majorité des cas, les deux sont synchronisés grâce au "dual-write" (on incrémente Redis en même temps qu'on écrit dans PG, avec compensation en cas de rollback).

Mais des écarts (drift) peuvent apparaître :

  • L'incrément Redis vit dans la transaction PG : un crash ou un timeout entre l'incrément et le commit laisse Redis sur-compté (le hook de compensation ne tourne pas).
  • Des rollbacks qui compensent le compteur.
  • Des workers concurrents.
  • Des scénarios de blocage (assignation de briques) qui laissent le compteur dans un état incohérent.

Tôt dans la collecte, un écart de quelques bricks est sans gravité.
Près de la fin (99.9 % ou plus), le même écart peut faire vendre plus de bricks que la propriété n'en contient → sur-vente, funding bloqué, et travail manuel pour réparer.

La task de réconciliation agit comme un audit + réalignement automatique juste avant la ligne d'arrivée.

Les deux compteurs et la vérité

Compteur Stockage Rôle principal Comment on le lit
purchased_brick_count Redis (property_{id}_purchased_brick_count) Disponible bricks, near-completion, checks rapides PrimaryPurchaseRepository.getPropertyPurchasedBricksCount
auto_invest_purchased_brick_count Redis (variante) Suivi séparé des achats auto-invest ...getPropertyAutoInvestPurchasedBricksCount
activePurchasedBrickCount + autoInvestBrickCount Postgres (somme sur primary_purchase) Vérité officielle FundingCounterReconciliationRepository.getPgPurchasedBrickCounts

La requête PG filtre explicitement :

where "propertyId_view" = $1
  and "status_view" not in ('declined', 'refunded')

Seuls les achats "vivants" comptent. Les declined et refunded ont en principe déjà fait l'objet d'un décrément du compteur Redis au moment de leur traitement — sauf les declines no-bricks-left du worker d'assignation (declinePrimaryPurchasesForMissingBricks), qui ne décrémente volontairement pas Redis (source principale du drift sur les fundings bloqués en bricks-oversell).

Comment les compteurs se désalignent (dual-write)

Le flux normal d'un achat direct ou auto-invest :

  1. Lecture des bricks disponibles (compteur Redis) ; la balance WT a été vérifiée en amont (routage achat direct/réservation, distribution auto-invest) et n'est re-contrainte que par l'échec du débit WT en cas de race condition sur le solde (un autre débit concurrent entre la vérif et le débit). Distinct de la race condition sur le compteur Redis (étapes 3–5).
  2. Débit de la WT d'achat puis création de la primary_purchase, dans la même transaction Postgres (la réservation, elle, ne crée pas de WT à ce stade).
  3. Incrément atomique du compteur Redis (incrementAndGet).
  4. Si la transaction PG échoue → setOnRollback déclenche un décrément compensatoire.
  5. Après l'incrément, si newBrickCount > property.totalBrickCountErr('primary-purchase.not-enough-bricks-left') : rollback complet (purchase + WT) + décrément compensatoire via setOnRollback (protection nominale contre la concurrence ; la réconciliation n'est que le filet du dernier kilomètre).

Ce mécanisme est très robuste, mais pas infaillible à 100 % en présence de :

  • Concurrence forte (plusieurs achats simultanés sur la même propriété).
  • Timing entre l'incrément Redis et le commit PG visible par les autres workers.
  • Blocs posés par le watcher d'assignation de briques (qui peut décliner des achats en excès).

C'est pourquoi on ne corrige pas au premier signe d'écart.

Déclenchement et périmètre

La task tourne via un cron Graphile :

  • Nom : funding-counter-reconciliation
  • Horaire : 30 2 * * * (02:30 UTC tous les jours)
  • Position dans la journée : après l'expiration des auto-invest et avant la clôture auto des fundings.
  • Retries Graphile : maxAttempts: 3 (1 run + 2 retries rapides). Si les compteurs restent instables, on n'enchaîne pas les retries pendant des heures (défaut Graphile = 25) ; le cron du lendemain reprend les fundings encore bloqués.

Elle ne touche que les propriétés : - publishStatus = 'published' - Funding démarré (startedAt <= now) - Funding non clôturé (ended est null)

Parmi elles, elle traite :

  1. Celles qui sont near completion selon le compteur Redis :
    redisPurchasedBrickCount / fundingCapacityBrickCount >= 0.999
    
  2. Celles qui sont déjà bloquées pour réessayer la correction :
  3. blocked.reason === 'counter-reconciliation'
  4. blocked.reason === 'bricks-oversell'

Les fundings clôturés ou en brouillon sont ignorés.

Le déroulé de la réconciliation (pas à pas)

flowchart TD
    Start([Cron quotidien]) --> Load[Charger fundings actifs<br/>published + started + non clôturés]
    Load --> Identify[Near-completion ≥99.9%<br/>ou déjà bloqués ?]
    Identify -->|non| Ignore[Ignorer]
    Identify -->|oui| Loop[Pour chaque propriété]

    Loop --> Block["1. ensureBlocked"]
    Block --> Drain["2. waitUntilAssignationDrained<br/>max 3 × 5s"]
    Drain -->|timeout| FailDrain[Alerte drain-timeout<br/>funding reste bloqué]
    Drain -->|waiting = 0| Stabilize["3. stabilizeCounters<br/>max 3 × 5s"]
    Stabilize -->|unstable| FailUnstable[Alerte failed<br/>funding reste bloqué]
    Stabilize -->|synced / needs-fix| Oversell{PG active > capacity ?}
    Oversell -->|oui| FailOversell[Alerte pg-oversell-after-drain<br/>funding reste bloqué]
    Oversell -->|non| Fix["4. applyOrPreviewFix<br/>cap → write Redis sauf skip<br/>si needs-fix"]
    Fix --> Unblock["5. unblock"]
    FailDrain --> Next
    FailUnstable --> Next
    FailOversell --> Next
    Unblock --> Next
    Next --> Loop

À l'intérieur de reconcileOneProjectInFunding :

  1. ensureBlocked — posé dès la détection near-completion (avant le reconcile lent). Si déjà bricks-oversell, on garde ce reason.
  2. waitUntilAssignationDrained — poll COUNT(waiting_for_assignation) jusqu'à 0 (max 3 lectures espacées de stabilizationDelayMs). Timeout → funding reste bloqué + alerte assignation-drain-timeout.
  3. stabilizeCounters — jusqu'à 3 lectures Redis+PG. synced → rien à écrire ; needs-fix → overwrite ; unstable → reste bloqué.
  4. applyOrPreviewFixcapRedisCounterWriteValues puis writeRedisCounters (sauf skipRedisCounterWrite). Si PG active > capacity après drain → pas de fix, pas d'unblock, alerte pg-oversell-after-drain.
  5. unblock — seul ce job lève les blocs. Flag déjà disparu → succès idempotent (pas d'échec cron).

Ce job est le seul autorisé à lever les blocs (counter-reconciliation et bricks-oversell). Le worker d'assignation continue sous toute raison de blocage (drain : assignation partielle + declines) — requis car la task pose counter-reconciliation avant d'attendre waiting_for_assignation = 0. Les watchers expiration / confirmation excluent toute raison de blocage : le path d'échec de confirmation décrémente Redis (not-enough-balance), l'expiration aussi.

La stabilisation : attendre que ça se calme

C'est le cœur de la logique anti-faux-positif.

// pseudo
previous = undefined
current = readBoth()

for read = 1 to 3:
  decision = evaluateCountersOnTwoIntervals(previous, current)
  if decision == 'synced'  OK, rien à faire
  if decision == 'needs-fix'  on a deux lectures identiques  on peut corriger
  if decision == 'retry' et read < 3  sleep(5s), previous = current, current = read()

si toujours pas stable  'unstable'

Pourquoi deux lectures identiques ?

Après ensureBlocked + drain assignation, plus de nouveaux achats ni de declines waiting_for_assignation. Un écart à la première lecture peut encore être un dual-write in-flight. Si deux lectures successives (espacées de 5s) donnent exactement les mêmes valeurs → l'écart est réel et stable.

Si les lectures continuent de bouger → quelque chose mute encore (bug, worker qui n'a pas vu le blocage). On abandonne et on alerte.

Le blocage du funding comme "fenêtre de gel"

Le funding peut être bloqué pour deux raisons (champ funding.blocked depuis la migration de juillet 2026) :

Raison Qui pose le blocage Effet
bricks-oversell Watcher d'assignation quand il n'y a pas assez de bricks physiques libres Sur-vente détectée → achats en excès déclinés, funding gelé en attendant investigation
counter-reconciliation La task elle-même (ou retry sur un précédent blocage) Gel temporaire pour stabilisation + correction des compteurs

Effet du blocage selon l'opération :

  1. Increments refusés via checkInvestorPropertyAccessprimary-purchase.not-enough-bricks-left : purchaseBricks, reserveBricks.
  2. Mutations refusées via assertFundingNotBlockedprimary-purchase.funding-blocked : refund, annulation de réservation (setExpirationToNowForReservation), confirmation de réservation (carte ou admin).
  3. Watchers batch : expiration / confirmation ne retiennent que funding->'blocked' is null (évite busy-loop) ; assignation continue sous toute raison (drain pendant la réconciliation).
  4. Defense in-TX : expiration skip silencieux (warn + return) et assert confirm/refund/setExpiration restent pour la race batch→TX (block posé entre le SELECT et la mutation) et les chemins HTTP.

Le blocage est posé avant drain + stabilisation : il coupe les nouveaux achats et les watchers qui peuvent décrémenter Redis. L'assignation continue jusqu'à waiting_for_assignation = 0 (y compris sous counter-reconciliation).

Après correction (ou constat d'alignement), et seulement si PG active ≤ capacity, on retire le flag.

Ce que fait vraiment la correction

Quand on décide de fixer :

const targets = capRedisCounterWriteValues({
  targets: {
    purchasedBrickCount: pg.activePurchasedBrickCount,
    autoInvestPurchasedBrickCount: pg.autoInvestBrickCount,
  },
  fundingCapacityBrickCount,
})
// always capped; Redis write skipped when skipRedisCounterWrite
await writeRedisCounters({ projectId, targets })

C'est un overwrite pur (setPrimitiveValue), pas un incrément/décrément.

Postgres n'est jamais modifié sur les compteurs (primary_purchase) par la task ; la task écrit toutefois properties.funding.blocked pour poser/lever le gel. PG reste la référence pour les compteurs — le decline d'excès reste uniquement dans le worker d'assignation.

Alertes et modes

Deux types d'alertes Slack via InternalAlertSlackService :

  • funding-counter-reconciliation-fixed : écart détecté et corrigé (ou "would fix" si skipRedisCounterWrite). Contient le delta, la décomposition PG, le nombre de lectures de stabilisation, etc.
  • funding-counter-reconciliation-failed :
  • counter-unstable : les compteurs continuaient de bouger après blocage + drain.
  • assignation-drain-timeout : des waiting_for_assignation restent après ~15s — l'assignation n'a pas fini le drain.
  • pg-oversell-after-drain : PG active > capacity après drain — intégrité cassée, pas d'unblock.
  • Erreur inattendue.

Mode skipRedisCounterWrite (FUNDING_COUNTER_RECONCILIATION_SKIP_REDIS_COUNTER_WRITE=true) :

  • La task pose/lève toujours le blocage PG (funding.blocked) pour figer les mutations d'achat pendant l'observation.
  • Elle détecte les écarts et envoie les alertes "would fix" ([SKIP REDIS WRITE]).
  • Elle ne modifie jamais les compteurs Redis.
  • Très utile en production pour un premier passage sans risque de correction Redis.

Config : - FUNDING_COUNTER_RECONCILIATION_STABILIZATION_DELAY_MS (défaut : 5 000 ms) — utilisé pour le drain et la stabilisation (budget ~15s chacun). - FUNDING_COUNTER_RECONCILIATION_SKIP_REDIS_COUNTER_WRITE (défaut : false).

Exemples de scénarios (tirés des tests)

  • Redis sur-compte de 1 sur un funding à 9/10 → la task remet Redis à 9, débloque, alerte.
  • Funding bloqué "bricks-oversell" avec Redis à 12 alors que PG dit 7 → correction vers le bas.
  • Seul le compteur auto-invest est décalé → seule cette valeur est corrigée.
  • Deux fundings : seul celui ≥ 99.9 % est traité.
  • Funding déjà clôturé → ignoré même si compteurs désalignés.
  • Lecture qui bouge encore après blocage → unstable + funding reste bloqué + alerte.
  • waiting_for_assignation encore présents → drain-timeout, Redis inchangé, funding reste bloqué.
  • waiting_for_assignation drainés mid-wait (worker assignation simulé via hook) → fix Redis + unblock.
  • PG active > capacity après drain → pg-oversell-after-drain, pas d'unblock.

Références code & tests

  • Logique d'évaluation : src/__new/modules/primary-purchase/business/evaluate-counter-reconciliation.ts
  • Service : .../services/funding-counter-reconciliation.service.ts (reconcileOneProjectInFunding = pipeline block → drain → stabilize → fix → unblock)
  • Repository (recon reads / Redis set) : .../repository/funding-counter-reconciliation.repository.ts
  • Repository (active fundings / lock / update / isBlocked) : .../repository/primary-purchase-funding.repository.ts
  • Service (block / unblock multi-step) : .../services/primary-purchase-funding.service.ts
  • Cron : projects/api/cron-task/funding/funding-counter-reconciliation.cron-task.ts
  • Alertes : .../services/funding-counter-alert.service.ts
  • Blocage & compteurs : primary-purchase.service.ts + PrimaryPurchaseRepository (inc/dec) + PrimaryPurchaseFundingRepository (assertFundingNotBlocked)
  • Watcher d'assignation (autre source de blocage) : state-watcher/purchase-waiting-for-assignation.state-watcher.ts
  • Tests unitaires : __tests__/evaluate-counter-reconciliation.unit-test.ts
  • Tests d'intégration : cron-task/funding/__tests__/funding-counter-reconciliation.integration-test.ts + fixtures tests/integration/fixtures/primary-purchase/

Liens rapides (dans la codebase) :

  • isFundingNearCompletion / capRedisCounterWriteValues / isPgOversellingCapacity
  • evaluateCountersOnTwoIntervals
  • countWaitingForAssignation
  • PrimaryPurchaseFundingService.blockProjectFundingIfNotAlreadyBlocked / unblockProjectFundingIfBlocked
  • getPgPurchasedBrickCounts

Cette task est le dernier filet de sécurité avant la clôture d'une collecte. Elle transforme un risque de sur-vente en un processus observable, réparable automatiquement, et tracé par des alertes claires.