Aller au contenu

Module leaderboard — recompute stateless

/internal/leaderboard/* sert le top investisseurs et les key-numbers du site vitrine bricks.co — endpoints authentifiés par x-api-key (consommateur site-vitrine), exposés via le SDK OpenAPI des routes /internal. Aucune donnée sensible : publicInvestorId opaque, jamais le customerId. Les coupons viennent de wallet_transactions (~37,5 M lignes, 53 GB), les achats de primary_purchase (table dédiée, source de vérité des achats — bien plus petite que wallet_transactions) : agréger en live = ~6 min par requête, intenable. Le module precompute une fois par jour vers une table dérivée — les endpoints servent la table en quelques ms. Coupons & achats comptés dès qu'ils sont utilisables par l'investisseur : wallet_transactions.status IN (confirmed, waiting) et primary_purchase.status sécurisé (confirmed, waiting_for_p2p_creation, waiting_for_contract_creation).

Surface du module

Trois endpoints, tous POST, authentifiés x-api-key (consommateur site-vitrine) : la data alimente la vitrine bricks.co et n'expose aucune donnée sensible (pas de customerId ni de referralCode vanity). Deux identifiants par investisseur : un publicInvestorId opaque (hash salé du customerId) = clé unique/stable de lookup, et un investorDisplayName cosmétique (Brickers_ + 4 derniers chars du referralCode) = libellé affichable. Lectures sur la replica (readPgAndValidate). Montants en cents ; le type branded cents porte l'unité (pas de suffixe Cents dans le naming).

Endpoint Source Forme
POST /internal/leaderboard/coupon/top leaderboard_investor_live + RANK() OVER (…) { entries: [{ rank, publicInvestorId, investorDisplayName, currentMonthCouponRevenue, totalCouponRevenue, investedObligation, obligationProjectCount }], totalCount, avgInvestedObligation, avgObligationProjectCount, currentMonthCouponRevenueTotal, obligationProjectCountTotal } — paginé { limit, offset }, tri { sort } (currentMonthCouponRevenue_* | investedObligation_*, défaut currentMonthCouponRevenue_desc) — oblig + prêt
POST /internal/leaderboard/statistics Sub-selects sur leaderboard_investor_live + COUNT(properties) + COUNT(customers) + AVG(yearlyRate) oblig+prêt live { investorCount, avgObligationProjectCount, avgInvestedObligation, computedAt, projectCount, investedAmount, revenue, investedObligation, revenueCoupon, avgYieldCoupon } — racine = plateforme (oblig+prêt+royalty), champs investedObligation/revenueCoupon + avgYieldCoupon = oblig+prêt
POST /internal/leaderboard/coupon/investor 1 row de leaderboard_investor_live (dont metricsLast12Months) + RANK + détail par projet { publicInvestorId, investorDisplayName, rank, registeredAt, currentMonthCouponRevenue, totalCouponRevenue, investedObligation, obligationProjectCount, monthlyHistory: [{ month, revenue, invested, projectCount, investmentCount, averageInvestment }], investmentsByProject: [{ projectId, contractType, amount }] }404 si publicInvestorId inconnu ou investisseur en visibilité privateoblig + prêt

Les trois endpoints servent la table dérivée → réponse en quelques ms. La courbe 12 mois est désormais par investisseur : le cron précalcule metricsLast12Months (revenue coupon + invested oblig+prêt + projectCount + investmentCount par mois) dans une colonne JSONB du snapshot — une donnée différente par ligne, zéro duplication (avant, une courbe plateforme unique était stampée à l'identique sur les ~138 k lignes). /investor sert sa monthlyHistory directement depuis cette colonne (plus de query live) et y dérive l'averageInvestment (invested / investmentCount). /top ne renvoie plus de courbe. registeredAt = date de création du compte (customers.createdAt). Les month des séries sont au format YearMonthDate (YYYY-MM) mais suivent le cycle de versement du 8 au 8 (les coupons tombent le 8), pas le mois calendaire — décalage de 7 jours avant date_trunc. Idem pour currentMonthCouponRevenue (la clé de tri) : « mois courant » = depuis le dernier 8, sinon le classement afficherait 0€ partout du 1er au 8. /coupon/top et /coupon/investor ne servent que les produits à coupon (obligation + prêt, royalty exclu) : totalCouponRevenue = revenue_obligation_coupon, investedObligation/obligationProjectCount = achats dont primary_purchase.json.contractType = 'obligation' (les prêts y sont enregistrés obligation). /statistics compare cette base coupon (investedObligation, revenueCoupon, avgYieldCoupon) à la plateforme entière : investedAmount = tous contrats, revenue = revenue_obligation_coupon + revenue_royalty, projectCount = toutes les properties.

investmentsByProject agrège les achats primaires confirmés ou en cours de settlement (primary_purchase.statusconfirmed, waiting_for_p2p_creation, waiting_for_contract_creation) par projectId, avec le contractType lu sur la property (obligation | loan uniquement — pas les royalties). Montants en cents, triés par montant décroissant. Query live sur primary_purchase (pas précalculée dans le snapshot) — acceptable car filtrée par un seul investorId.

Visibilité investisseur (leaderboardVisibility)

Chaque investisseur peut se retirer du classement public via PATCH /customers/leaderboard/visibility (valeurs public | private, défaut public — modèle opt-out). Stocké sur customers.leaderboardVisibility.

Visibilité /top /investor Snapshot cron
public Inclus (rang calculé parmi les publics uniquement) 200 + détail Row toujours présente
private Exclu (totalCount et agrégats /top ne le comptent pas) 404 investor-not-found Row toujours présente (filtre lecture uniquement)

Le rank est global parmi les investisseurs publics : un privé avec le plus gros revenu n'apparaît pas et ne décale pas le rang des autres. Le cron ne filtre pas — basculer en private prend effet immédiatement sans attendre le refresh quotidien.

/statistics n'applique pas ce filtre sur investedAmount / revenue / moyennes (agrégats plateforme marketing). investorCount = COUNT(customers) (tous les comptes, pas seulement le top public).

Clients : toggle dans le profil web (EditLeaderboardVisibility), champ exposé dans GET /customers/account (leaderboardVisibility). Voir account.md.

Le cron refresh-leaderboard

Schedule 0 3 * * * UTC (3h du matin) — déclaré dans graphile-cron.worker.ts. Durée mesurée sur fork prod (16/06/2026, 37,5 M wallet_transactions) : ~6 min 36 s. Le scan master (CTE sur referral-link + wallet_transactions + properties) en consomme la quasi-totalité ; le reste — bulk INSERT de ~101 k rows via unnest($1::uuid[], …) puis les 3 RENAME en transaction — tient en quelques secondes.

La fréquence est un levier facile

Le coût quasi-intégral du cron est un scan read-only sur le master (~6 min) — l'écriture (INSERT + RENAME) ne pèse que ~3 s. Si le métier veut des stats plus fraîches, on peut basculer à 0 */6 * * * (toutes les 6 h) ou même 0 */1 * * * (toutes les heures) en changeant la seule ligne match du cron dans graphile-cron.worker.ts. Le scan ne prend aucun lock et les endpoints lisent la replica → aucun impact lecture. Le plafond pratique reste durée du scan < intervalle (aujourd'hui ~7 min < 1 h, large marge).

flowchart TD
    Cron(["0 3 * * * UTC<br/>graphile-worker"]):::endpoint
    Read["readAggregateForRefresh()<br/>master · hors transaction"]:::master
    Rows[["~101 k rows<br/>buffer Node ~40 MB"]]:::data
    Tx["safePgTransaction (master)"]:::master
    Trunc["TRUNCATE leaderboard_investor_live_next"]:::master
    Insert["INSERT INTO leaderboard_investor_live_next<br/>SELECT FROM unnest($1::uuid[], …)"]:::master
    R1["RENAME live → swap_tmp"]:::master
    R2["RENAME next → live"]:::master
    R3["RENAME swap_tmp → next"]:::master
    Done(["snapshot publié<br/>endpoints voient la nouvelle version"]):::success

    Cron --> Read --> Rows --> Tx
    Tx --> Trunc --> Insert --> R1 --> R2 --> R3 --> Done

    classDef endpoint fill:#4f46e5,color:#fff,stroke:#3730a3
    classDef master fill:#ef4444,color:#fff,stroke:#b91c1c
    classDef data fill:#f59e0b,color:#fff,stroke:#b45309
    classDef success fill:#10b981,color:#fff,stroke:#047857

Pourquoi les lectures vivent hors safePgTransaction

Les deux scans (readAggregateForRefresh(), readInvestorMetricsLast12Months()) tournent sur le master via queryPgAndValidate (getAppDataSource().query) : sur la replica, un scan de ~6 min se fait canceller par un conflit de recovery WAL (canceling statement due to conflict with recovery), inexistant sur le master qui n'a pas de replay. Chacun reste un statement unique hors transaction : pgbouncer est en mode transaction côté Bricks, et une connexion master idle au sein d'une transaction le temps du scan serait releasée par le pooler → QueryRunnerProviderAlreadyReleasedError. Le cron lance les deux scans en parallèle (Promise.all, fail-fast) puis writeAggregateAndSwap(rows, metricsRows) ouvre sa propre safePgTransaction (master) pour l'INSERT + les 3 RENAME.

Pourquoi deux tables ?

flowchart LR
    subgraph before ["Avant le cron"]
        A1[("leaderboard_investor_live<br/>snapshot J-1 — lue par /top, /statistics, /investor")]:::live
        A2[("leaderboard_investor_live_next<br/>vide")]:::staging
    end
    subgraph during ["Pendant le cron"]
        B1[("leaderboard_investor_live<br/>snapshot J-1 — toujours lue")]:::live
        B2[("leaderboard_investor_live_next<br/>en cours de remplissage")]:::staging
    end
    subgraph swap ["Swap atomique (1 transaction, 3 RENAME)"]
        C1[("leaderboard_investor_live<br/>snapshot J — lue dès le COMMIT")]:::live
        C2[("leaderboard_investor_live_next<br/>vide à nouveau (ancien J-1)")]:::staging
    end

    before --> during --> swap

    classDef live fill:#10b981,color:#fff,stroke:#047857
    classDef staging fill:#94a3b8,color:#fff,stroke:#475569
  • leaderboard_investor_live — celle lue par les endpoints. Toujours pleine, toujours cohérente.
  • leaderboard_investor_live_next — staging. Vide hors cron, remplie pendant, swappée à la fin.
  • Swap atomique : 3 ALTER TABLE … RENAME dans la même transaction Postgres. Aucun endpoint ne voit d'état intermédiaire — la requête courante lit l'ancien snapshot complet, la suivante lit le nouveau snapshot complet (jamais un mix moitié J-1 / moitié J).
  • LIKE INCLUDING ALL à la création de _next : PK + indexes clonés. Ils survivent au RENAME, donc pas de reconstruction à chaque cron — c'est essentiel pour garder l'INSERT à ~3 s.

Colonnes split coupon vs all (BRI-1256)

Le snapshot porte, en plus des totaux historiques, trois colonnes par-investisseur pour la comparaison oblig-vs-royalty : totalRevenue (revenue_obligation_coupon + revenue_royalty), investedObligation et obligationProjectCount (achats contractType = 'obligation'). Elles sont calculées dans les scans existants de readAggregateForRefresh (filtres conditionnels, aucun scan supplémentaire) → durée du cron inchangée. Le _next doit recevoir ces colonnes via ALTER explicite (le LIKE n'est posé qu'à la création), cf. migration V202606221200.

Le IF NOT EXISTS partout dans la migration

Permet de re-run Flyway sans dommage si on doit dropper / recréer manuellement. Le cron est idempotent : deux passes successives produisent exactement le même snapshot.

Replica vs master

Trois flux, deux pools : le cron (scan + écriture) vit entièrement sur le master, les endpoints lisent la replica.

Phase Connexion Pourquoi
Scan agrégat wallet_transactions + primary_purchase (cron) Master (TypeORM DataSource.query, queryPgAndValidate) Sur la replica, le scan ~6 min se fait canceller par un conflit de recovery WAL (canceling statement due to conflict with recovery) ; le master n'a pas de replay → aucun conflit. À 3h UTC il absorbe sans peine un scan read-only.
TRUNCATE + INSERT + RENAME (cron) Master (TypeORM EntityManager dans safePgTransaction) Écritures + DDL → master obligatoire, et la transaction unique garantit le swap atomique.
POST /coupon/top · /statistics · /coupon/investor Replica (readPgAndValidate) Lecture seule, stale-tolerant. Le pool replica absorbe la charge vitrine sans peser sur le master.

Le pool replica = même socle que les segments SQL opérateur

getReadOnlyPgPool + le rôle bricks_ro sont mutualisés avec /internal/segments. Si tu touches au pool, vérifie l'impact sur les deux.

Limites & roadmap

Aujourd'hui on tient sans effort. Voici les limites connues et ce qu'on fera quand elles deviendront gênantes — pas avant.

Limite Aujourd'hui Devient un problème quand… Mitigation envisagée
Fraîcheur 1 refresh / jour (3h UTC) Demande métier "stats < 1 h" Passer le schedule à 0 */1 * * *. Le cron dure ~7 min, compatible avec un slot horaire.
Durée du cron ~6 min 36 s pour 37,5 M wallet_transactions > 30 min → pression sur le scheduling, fenêtre de maintenance qui rétrécit (a) Scan incrémental : ne lire que wt."createdAt" > last_run puis merger sur l'agrégat précédent. (b) pg-cursor côté Node si le buffer transient (~40 MB) devient bloquant (>1 M investisseurs actifs).
Volume table dérivée ~101 k rows / ~10 MB > 10 M rows Idem : incrémental + colonne updatedAt pour détecter les lignes modifiées.
Précision globale /statistics Undercount des investisseurs sans referral-link (filtre du CTE latest_referral_per_investor) Demande comptable ou réglementaire Hors scope marketing — la query AMF/MiFID live reste la source officielle. Ne pas industrialiser.

On a le temps

La vitrine bricks.co tolère 1 refresh / jour et il n'y a aucun SLO sur la fraîcheur. Cette page existe pour que le prochain dev (a) ne refasse pas le diagnostic pgbouncer en panique à 3h du matin et (b) sache où chercher quand on lui demandera "on peut avoir des stats en temps quasi-réel ?".

Décisions & FAQ

Pourquoi pas une MATERIALIZED VIEW Postgres ?

REFRESH MATERIALIZED VIEW CONCURRENTLY exige un index unique et lock la vue le temps de la passe ; REFRESH sans CONCURRENTLY lock en exclusif → endpoints bloqués pendant la durée du refresh. Le swap par RENAME nous donne le même résultat (atomic update) sans aucun lock côté lecture et nous laisse le contrôle fin du SQL (CTE multi-table, casting, sources mixtes replica/master).

Pourquoi unnest($1::uuid[], $2::uuid[], …) plutôt qu'un INSERT … VALUES multi-row ?

101 k rows en un seul round-trip réseau, payload ~25 MB en arrays binaires (vs ~80 MB en CSV multi-VALUES + parsing PG plus coûteux). unnest reste l'idiome standard côté driver pg pour les bulk inserts massifs.

Pourquoi RANK() OVER (…) dans une sous-SELECT, et pas un ROW_NUMBER() ?

Ties autorisés : deux investisseurs à 0 € de revenu mensuel doivent avoir le même rang. ROW_NUMBER() casserait la cohérence visible vitrine (rang 12 != rang 12). Le rang reste global (pas relatif à la page) en calculant RANK() avant LIMIT/OFFSET.

Comment récupère-t-on un investisseur, et que se passe-t-il s'il est inconnu ?

Par son publicInvestorId (hash salé du customerId, ex. a3f90c2b1d…) : clé opaque, unique et stable, calculée par le cron et stockée en colonne indexée. C'est lui qu'on passe en body et que le /top renvoie par entrée → le front navigue vers la fiche par cet id. On expose aussi un investorDisplayName purement cosmétique (Brickers_ + 4 derniers chars du referralCode) pour l'affichage — il peut collisionner, donc jamais utilisé comme clé. On n'expose jamais le referralCode vanity (souvent un prénom/pseudo → leak) ni le customerId/UUID interne. Inconnu → 404 (le contrat n'est plus .nullable()).

Pourquoi compter les waiting et pas seulement les confirmed ?

Une transaction waiting (coupon en cours de settlement Lemonway) ou un achat primary_purchase dont les bricks sont déjà sécurisées (waiting_for_p2p_creation, waiting_for_contract_creation) est immédiatement utilisable par l'investisseur — pas la peine d'attendre la confirmation finale pour la compter. Mois courant, all-time et monthlyHistory partagent ce même filtre → cohérence sur toute la réponse. On exclut seulement les états annulables/refusés (declined, canceled, refunded).

Le tri /coupon/top (sort) re-calcule-t-il le rank ? Faut-il un index ?

Non au rang : il reste global par revenu mensuel (RANK() OVER (…)), indépendant du sort — trier par investedObligation réordonne la page mais garde les badges de rang cohérents. Et pas d'index dédié : la window function trie déjà toute la table à chaque appel, donc un index n'accélérerait pas le ORDER BY … LIMIT. Un index ne deviendrait utile que si on stockait le rank en colonne au refresh (à reconsidérer si la table dépasse ~1 M lignes).