Aller au contenu

Runtime — workers & exécution

Trois workers portent le cycle de vie d'un P2P. Chacun est un service Railway dédié (même image que l'API, startCommand: node dist/worker/main.js), activé par un flag d'environnement, et tourne en boucle continue via safeInfiniteLoop (heartbeat loggé à chaque itération). Pause idle par défaut : 3 s. Workers P2P : PENDING_LEMONWAY_P2P_WORKER_SLEEP_TIME_MS / PLAYED_LEMONWAY_P2P_WORKER_SLEEP_TIME_MS (défaut 3 s).

Worker Flag d'activation Timeout / itération Rôle
bricks-purchase-p2p IS_BRICKS_P2P_WORKER 10 s Crée les P2P des achats de bricks
pending-lemonway-p2p IS_PENDING_LEMONWAY_P2P_WORKER 60 s Joue les P2P pending chez Lemonway
played-lemonway-p2p IS_PLAYED_LEMONWAY_P2P_WORKER 60 s Réconcilie : tire les conséquences métier des P2P joués

Les fichiers worker/*.worker.ts exposent une fonction start* appelée par worker/main.ts ; la logique vit dans les state-watchers et le module lemonway-p2p.

1. bricks-purchase-p2p — création des P2P d'achat

State-watcher : purchase-waiting-for-p2p.state-watcher.ts. À chaque itération, dans une seule transaction Postgres :

  1. Lock d'un batch (≤ 1000) d'achats en waiting_for_p2p_creation.
  2. Groupage par propertyId, résolution du wallet Lemonway de la SPV.
  3. PrimaryPurchaseConfirmationService.setP2PForPurchases — attache à chaque WT d'achat un P2P pending (investisseur → SPV).
  4. Passage des achats en waiting_for_contract_creation.

Les autres flux métier (gift card, referral, checkout…) n'ont pas de worker de création : leur service attache le P2P à la WT au fil de l'eau, via la même factory WalletTransaction.createPendingP2P.

2. pending-lemonway-p2p — exécution chez Lemonway

State-watcher : pending-p2p-to-play.state-watcher.ts.

Sélection priorisée — l'ordre des débits est sacré

La requête getWtWithPendingLemonwayP2PIdsToPlayWithPriority ne prend pas « tous les pending » : elle applique des contraintes d'ordre pour que le solde Lemonway d'un customer ne soit jamais consommé dans le désordre.

  • Éligibilité : WT waiting, P2P pending, attemptAfter dépassé.
  • Comptes bloqués exclus : si le wallet débité ou crédité appartient à un customer bloqué (transactionRights ≠ 'all', blockedTemporarilyAtLemonway, lemonwayStatus ≠ '6', ou investorIdentityStatus.status = 'outdated'), le P2P et les retraits waiting de ce customer sont invisibles pour le ranking — ils attendent le déblocage / la recertif. verified + update n'est pas exclu. Sinon un withdraw devient rn = 1 pendant qu'un P2P en attente (ex. gift-card 167) est filtré. Un P2P exclu sort du ranking : il ne tient plus le rn = 1, les débits suivants d'un débiteur verified peuvent passer devant.
  • Un seul débit à la fois par compte customer : les débits P2P sont rankés par compte débiteur (ROW_NUMBER() … PARTITION BY debitAccountId ORDER BY index), unifiés avec les retraits waiting dans la même CTE. Seul le plus ancien (rn = 1) est sélectionnable : un retrait en attente passe donc avant les P2P débit suivants du même customer, et vice-versa (FIFO strict par compte).
  • Comptes non-customer en parallèle : les wallets techniques et SPV n'ont pas de contrainte de solde côté Bricks — tous leurs P2P sont sélectionnables d'emblée (utile pour les vagues de remboursements de capital).
  • ORDER BY index ASC, LIMIT concurrencyLimit (PENDING_LEMONWAY_P2P_WORKER_CONCURRENCY_LIMIT) : les plus anciens d'abord, débit ou crédit confondus.

Recevoir des crédits dans le désordre est sans risque (pas de contrainte de solde) — seuls les débits exigent l'ordre.

Jeu d'un P2P

Pour chaque WT sélectionnée (concurrence bornée par p-limit) :

  1. Transaction + lock : FOR UPDATE OF wallet_transactions SKIP LOCKED, en revalidant toutes les conditions d'éligibilité — un P2P déjà pris par une autre itération/réplique est simplement sauté.
  2. Appel Lemonway : attemptP2PAtLemonway → POST /p2p (play-p2p.lemonway.endpoint.ts) avec reference = walletTransactionId.
  3. Transition selon le résultat — voir Statuts & transitions pour l'arbre de décision complet.
  4. Persistance : saveLemonwayP2PForWt réécrit le JSON ; si succeeded, wallet_transactions.lemonwayTransactionId reçoit l'id Lemonway.

Idempotence côté Lemonway

Garantie par la reference envoyée à Lemonway : c'est l'id de la WT, invariant entre les retries.

Sur toute erreur du POST, le service interroge GET /p2p?reference=… (get-p2p.endpoint.ts) avant de classer l'échec : si un P2P existe déjà pour cette référence, l'appel est considéré réussi avec l'id existant. Ce check systématique (et pas seulement sur le code 348 « duplicate ») vient d'un incident réel : un déploiement entre l'appel Lemonway et la sauvegarde en DB avait fait rejouer un P2P, et Lemonway avait répondu « solde insuffisant » au lieu de « duplicate ».

Le filet de sécurité complet : SKIP LOCKED (pas deux replicas sur la même WT) + référence invariante (pas de double exécution chez Lemonway) + re-vérification au retry (un crash entre le POST et le save se résout tout seul à la tentative suivante).

Second consommateur de ce GET par référence : isP2PPlayedAtLw (remboursement de carte cadeau). Ok(true) veut dire que l'argent a bougé chez Lemonway, pas que le P2P est résolu au sens du worker played-lemonway-p2p (statut succeeded en base et conséquences métier appliquées). Ok(false) est une liste vide. Tout autre échec du GET est lemonway-p2p.lw-check-failed.

3. played-lemonway-p2p — conséquences métier

Worker : played-lemonway-p2p.worker.ts. Il sélectionne les WT waiting dont le P2P est joué :

  • crédits (totalValue_view > 0) : uniquement si P2P succeeded (un crédit ne passe jamais failed, cf. Statuts) ;
  • débits (totalValue_view < 0) : dès que le P2P n'est plus pending (succeeded → confirmer, failed → rembourser/décliner).

Une seule WT par customer et par itération, pour appliquer les conséquences dans l'ordre des exécutions Lemonway. La file ne rank plus tout le stock waiting : scan indexé ORDER BY "lemonwayTransactionId" borné à 200 lignes (prédicat aligné sur idx_wallet_transactions_p2p_played_ready, y compris "lemonwayP2P" IS NOT NULL), puis DISTINCT ON ("customerId") sur cette fenêtre, puis LIMIT PLAYED_LEMONWAY_P2P_WORKER_CONCURRENCY_LIMIT (défaut 2). Un ROW_NUMBER() OVER (PARTITION BY "customerId") forcerait le tri de toute la file — mesuré à 1,4 M lignes / ~4–6 s par appel quand les coupons du 08/09 saturent l'index. Chaque traitement locke la WT (SKIP LOCKED) et la ligne customer (cohérence des balances), dans un safePgTransaction.

Conséquences métier par kind

P2P succeeded → handleSucceededLemonwayP2P :

WT kind Action
PRIMARY_PURCHASE PrimaryPurchaseConfirmationService.confirmPrimaryPurchase puis WT → confirmed
PRIMARY_PURCHASE_REFUND PrimaryPurchaseRefundService.handleSucceededLemonwayP2PRefund
BUYING_DEAL MarketDealService.confirmDealPurchase
GIFT_CARD_PURCHASE / GIFT_CARD_CREDIT confirmGiftCardCreation / confirmGiftCardCredit. La carte cadeau est émise à l'achat, pas ici : confirmGiftCardCreation ne fait que passer la WT en confirmed, sauf sur les WT antérieures à ce changement (pas de giftcardId en contexte) qu'il émet encore
CREDIT_FOR_WALLET_ASSET_TRANSFER AssetTransferService.confirmAssetWalletTransfer
REVENUE_OBLIGATION_COUPON, WITHHOLDING_TAX, autres WT → confirmed

P2P failed → handleFailedLemonwayP2P :

WT kind Action
PRIMARY_PURCHASE refundPrimaryPurchase (reason p2pFailed)
BUYING_DEAL MarketDealService.declineDealPurchase
GIFT_CARD_PURCHASE declineGiftCardPurchase, sauf si le donneur vient d'être gelé (debit_account_blocked) : P2P remis en pending (FailedLemonwayP2P.retry, sans nouveau délai) et WT laissée waiting
WITHHOLDING_TAX P2P remis en pending (FailedLemonwayP2P.retry, J+1) + alerte Slack — les débits suivants du customer restent bloqués entre-temps
autres WT → declined + alerte Slack transactionBlocked (cas non prévu)

En amont du match par kind : si businessErrorType = debit_account_blocked (codes 111/167) et que le customer est le compte débité, le worker gèle le customer (CustomerFreezeService.freezeCustomer) — son identité n'est plus considérée vérifiée chez Lemonway.

Dans ce cas précis, un GIFT_CARD_PURCHASE n'est pas décliné : son P2P repart en pending et la WT reste waiting, donc les fonds restent réservés. Le worker pending-lemonway-p2p exclut les P2P dont le compte débité est gelé, le P2P dort donc jusqu'au dégel — quel que soit le chemin de dégel — puis se rejoue tout seul. Les autres kinds gardent leur refund/decline.

Au premier revenu encaissé (kinds de revenus du portfolio), le worker pose aussi customer.firstRevenueDate.

Monitoring & exploitation

Observabilité :

  • Spans dd-trace pending-p2p-to-play-state-watcher et played-lemonway-p2p-watcher (APM), heartbeats logués par itération.
  • Monitor Datadog sur les P2P à plus de 5 tentatives techniques.
  • Alertes Slack (InternalAlertSlackService) : debit-p2p-failed-insufficient-balance-no-pending-credits, withholding-tax-p2p-failed-will-retry, transactionBlocked (WT sans handler), échec de refund d'achat.

Performance : les requêtes de sélection s'appuient sur des index partiels dédiés — 2025-12-24-pending-p2p-worker-index-optimization.sql et 2025-12-24-played-p2p-worker-index-optimization.sql (clauses WHERE alignées sur les conditions exactes des workers).

Ces deux index bloatent par construction. Une WT entre dans status = 'waiting' puis en sort : le churn tombe exactement dans leur prédicat, et chaque transition laisse une entrée morte. Constaté en prod : 2 560 MB pour 53 026 entrées vivantes sur l'index pending, et un getWaitingWtIdsPlayedAtLw qui lit 38 406 pages d'index pour renvoyer 0 ligne (8,5 s par appel, 41,9 % du temps DB). Le remède est en deux parties, les deux nécessaires : 2026-08-30-p2p-worker-indexes-reindex-concurrent.sql rend le disque déjà perdu (mesuré en prod le 2026-09-01 : requête à 8,5 s → 0,634 ms et 708 227 blocs → 311 buffers, index 2 560 MB → 7 896 kB et 452 MB → 48 kB, ~2,9 GB rendus, sans interruption des workers), V202608301613 abaisse le seuil d'autovacuum de wallet_transactions pour que le bloat ne revienne pas.

La cadence n'est pas un intervalle. Les 60 s de la table ci-dessus sont un timeout par itération, pas une période. safeInfiniteLoop enchaîne les itérations et ne dort que 3 s quand la requête ne renvoie rien (PENDING_LEMONWAY_P2P_WORKER_SLEEP_TIME_MS / PLAYED_LEMONWAY_P2P_WORKER_SLEEP_TIME_MS, défaut 3 s). La période réelle est donc durée de la requête + 3 s, sur un seul réplica (railway.played-lemonway-p2p.worker.json, numReplicas: 1). Conséquence contre-intuitive : accélérer la requête augmente le nombre d'appels — une hausse du compteur calls de pg_stat_statements après un REINDEX est attendue, c'est total_exec_time et les blocs lus qu'il faut suivre.

Débloquer un P2P failed : pas d'endpoint admin générique. Après correction de la cause (compte débloqué, solde reconstitué), un script ponctuel remet le P2P en pending via FailedLemonwayP2P.retry — le pipeline reprend tout seul. Exemple : V202604211541__replay-failed-referral-p2p-for-3-customers.sql.

Scénarios de test bout en bout : projects/api/scripts/tests/lemonway-p2p/ (interactions P2P × retraits, multi-customers, idempotence).