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 :
- Lock d'un batch (≤ 1000) d'achats en
waiting_for_p2p_creation. - Groupage par
propertyId, résolution du wallet Lemonway de la SPV. PrimaryPurchaseConfirmationService.setP2PForPurchases— attache à chaque WT d'achat un P2Ppending(investisseur → SPV).- 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, P2Ppending,attemptAfterdépassé. - Comptes bloqués exclus : si le wallet débité ou crédité appartient à un customer bloqué (
transactionRights ≠ 'all',blockedTemporarilyAtLemonway,lemonwayStatus ≠ '6', ouinvestorIdentityStatus.status = 'outdated'), le P2P et les retraitswaitingde ce customer sont invisibles pour le ranking — ils attendent le déblocage / la recertif.verified+updaten'est pas exclu. Sinon un withdraw devientrn = 1pendant qu'un P2P en attente (ex. gift-card 167) est filtré. Un P2P exclu sort du ranking : il ne tient plus lern = 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 retraitswaitingdans 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) :
- 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é. - Appel Lemonway :
attemptP2PAtLemonway→POST /p2p(play-p2p.lemonway.endpoint.ts) avecreference = walletTransactionId. - Transition selon le résultat — voir Statuts & transitions pour l'arbre de décision complet.
- Persistance :
saveLemonwayP2PForWtréécrit le JSON ; sisucceeded,wallet_transactions.lemonwayTransactionIdreç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 P2Psucceeded(un crédit ne passe jamaisfailed, cf. Statuts) ; - débits (
totalValue_view < 0) : dès que le P2P n'est pluspending(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-watcheretplayed-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).