From 097b76c7efd731346f64346391d103113b768026 Mon Sep 17 00:00:00 2001 From: Arthur Ria Date: Tue, 25 Aug 2026 15:49:39 +0200 Subject: [PATCH] L6.2 : garde de generation contre les ecritures post-invalidation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Un fetch parti AVANT une invalidation (clearCache/invalidateCache, declenchees par onSwitch a la bascule de profil, D8) terminait APRES elle et ecrivait quand meme son resultat : le cache repartait peuple avec les donnees de l'ancien tenant, timestamp neuf, valid: true. Defaut latent avant L6.1 ; deterministe apres, puisque la promesse en vol survit desormais a l'invalidation. Compteur de generation par service, incremente a chaque invalidation. Le fetch capture la generation au depart et, au moment de publier, jette son resultat si elle a bouge. Surtout : les fonctions de chargement n'ecrivent plus rien en cache — la publication est un `commit` passe a singleFlight.run, appele seulement si la generation n'a pas change. La garde devient structurelle, pas conventionnelle : un chargement ne peut plus publier par inadvertance. L'invalidation vide aussi la Map des promesses en vol. L'appelant deja en attente recoit quand meme son resultat — il l'a demande avant la bascule ; c'est sa mise en cache qui est refusee. --- Verifications (LIMAGRAIN) --- Rafale [search_workflows(EasyWMS), switch_wms_profile(EUROTRAFIC)] envoyee d'un bloc, puis get_application_summary SEQUENTIEL apres les deux reponses. AVANT la garde (HEAD = L6.1), 3 rejeux — le cache survit a la bascule : ===== RUN 1 ===== search_workflows : success=true count=50 switch_wms_profile : success=true get_application_summary -> workflowCachesByApplication = {"EasyWMS":{"cached":true,"count":3944,"timestamp":1787665677594, "age":0,"valid":true}} adElementsByApplication = {} => caches peuples : 1 (attendu 0) ===== RUN 2 ===== idem, count 3944, timestamp 1787665687744, valid true ===== RUN 3 ===== idem, count 3944, timestamp 1787665703945, valid true APRES la garde, 3 rejeux — aucun cache peuple : ===== RUN 1 ===== search_workflows : success=true count=50 switch_wms_profile : success=true get_application_summary -> workflowCachesByApplication = {} adElementsByApplication = {} => caches peuples : 0 (attendu 0) -- stderr : 1 resultat(s) jete(s) | 2 'Cache cleared' ===== RUN 2 ===== identique : caches peuples 0, 1 resultat jete ===== RUN 3 ===== identique : caches peuples 0, 1 resultat jete Test direct et deterministe (invalidation declenchee pendant un fetch tenu ouvert par une porte) : [Workflow] Cache expired or empty for "TestApp", fetching from API... --- invalidation PENDANT le fetch (clearCache) --- [Workflow] Cache cleared [Workflow] Fetched 1 workflows (total: 1) [Workflow] Result for "workflows::TestApp" discarded, not cached: cache invalidated during fetch (generation 0 -> 1) l'appelant recoit bien son resultat : 1 workflow(s) cache apres coup : {} (attendu {}) Non-regression L6.1, 3 rejeux de la rafale de 6 (CustomApp) : == RUN 1 : 6 reponses a 44 | fetch=1 joins=5 caches=1 == RUN 2 : 6 reponses a 44 | fetch=1 joins=5 caches=1 == RUN 3 : 6 reponses a 44 | fetch=1 joins=5 caches=1 Chemin sequentiel nominal, inchange : [search_workflows] success=true application=EasyWMS count=50 len=14220 [search_workflows] success=true application=EasyWMS count=50 len=14220 fetch=1 cached=1 joins=0 Co-Authored-By: Claude Opus 5 --- CLAUDE.md | 5 ++- DECISIONS.md | 25 +++++++++++- ROADMAP.md | 7 ---- src/services/ad-service.js | 21 ++++++---- src/services/entity-resolver.js | 25 ++++++++---- src/services/single-flight.js | 70 ++++++++++++++++++++++++-------- src/services/workflow-service.js | 39 +++++++++++------- 7 files changed, 137 insertions(+), 55 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 70ab689..70a3ffd 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -201,7 +201,10 @@ ne les augmentez pas à l'aveugle. Les chargements sont **dédupliqués par clé de cache** : sous appels concurrents, une seule chaîne de fetch part par clé et les autres appelants la rejoignent -(D27) — deux applications différentes se chargent toujours en parallèle. +(D27) — deux applications différentes se chargent toujours en parallèle. Un +fetch parti avant une invalidation ne repeuple plus le cache après elle : la +publication passe par un `commit` gardé par un compteur de génération. **Ne +remettez jamais d'écriture de cache dans une fonction de chargement.** `get_application_summary` expose l'état des caches sans redémarrage. diff --git a/DECISIONS.md b/DECISIONS.md index a9a07cd..f21728c 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -648,7 +648,7 @@ passé), et ajoutent un `hint` quand la recherche revient **vide** : --- -## D27 — Chargements paresseux sous concurrence : single-flight par clé +## D27 — Chargements paresseux sous concurrence : single-flight + génération **Contexte (mesuré le 25/08/2026, `LIMAGRAI2512`).** Le serveur traite les `tools/call` **en concurrence** : une rafale d'appels dans une même session @@ -684,6 +684,29 @@ pas une dépendance externe : fusionnez pas les deux, c'est ce qui rend la déduplication vérifiable depuis stderr. +**Garde de génération.** Un fetch parti *avant* une invalidation terminait +*après* elle et écrivait quand même son résultat : le cache repartait peuplé +avec les données de l'ancien tenant, timestamp neuf, `valid: true`. Défaut +latent avant le single-flight, **déterministe après** (la promesse en vol +survit à l'invalidation). D'où : + +- Un **compteur de génération par service**, incrémenté à chaque invalidation + (`clearCache()` / `invalidateCache()`, toujours déclenchées par + `onSwitch()` — D8 inchangé). Le fetch capture la génération au départ. +- **Les fonctions de chargement n'écrivent plus rien en cache** : la + publication est un `commit` passé à `singleFlight.run`, appelé *seulement* + si la génération n'a pas bougé. C'est structurel, pas conventionnel — un + fetch ne peut plus publier par inadvertance. +- L'invalidation vide aussi la Map des promesses en vol. **L'appelant reçoit + quand même son résultat** — il l'a demandé avant la bascule ; c'est sa mise + en cache qui est refusée, tracée par + `Result for "…" discarded, not cached`. + +Mesuré sur `[search_workflows(EasyWMS), switch_wms_profile(EUROTRAFIC)]` envoyé +d'un bloc, puis `get_application_summary` en séquentiel : **3/3 avant**, le +cache EasyWMS de LIMAGRAIN (3 944 workflows) survit à la bascule avec un +timestamp neuf ; **3/3 après**, aucun cache peuplé. + **Mesures après.** Rafale de 6 (CustomApp) → 1 fetch + 5 joins, les 6 réponses à `count: 44`. Rafale de 6 (resolver) → 1 chargement, 5 GET Metadata au lieu de 30. Rafale mixte EasyWMS + CustomApp → **un fetch par application**, deux au diff --git a/ROADMAP.md b/ROADMAP.md index 801cac9..a562482 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -95,13 +95,6 @@ Rejoués en séquentiel, les mêmes appels sont corrects. Trois défauts structurels dans les services à cache (`workflow-service`, `ad-service`, `entity-resolver`) : -### L6.2 — Protéger l'invalidation contre les fetchs en vol - -Un fetch parti avant `clearCache()`/`invalidateCache()` (bascule de profil, -D8) écrit son résultat **après** l'invalidation : cache repeuplé avec les -données de l'ancien tenant. Garde de génération (epoch) : un résultat issu -d'une génération antérieure est jeté, pas écrit. - ### L6.3 — Ne pas encaisser un vide anormal `fetchAllWorkflows` fait `response?.entities || []` puis met en cache le diff --git a/src/services/ad-service.js b/src/services/ad-service.js index 9ed4ac2..3b80a4b 100644 --- a/src/services/ad-service.js +++ b/src/services/ad-service.js @@ -100,13 +100,23 @@ async function getElements(elementType, application) { return cache[key]; } - // Un seul chargement par (application, type), même sous rafale (D27). - return singleFlight.run(key, () => loadElements(app, elementType, key)); + // Un seul chargement par (application, type), même sous rafale, et + // publication refusée si le cache a été invalidé pendant le fetch (D27). + return singleFlight.run( + key, + () => loadElements(app, elementType, key), + (elements) => { + cache[key] = elements; + cacheTimestamps[key] = Date.now(); + console.error(`[AD] Successfully cached ${elements.length} ${key}`); + } + ); } /** * Chargement réel d'un (application, type) (pagination complète). * Appelé au plus une fois par clé tant qu'il est en vol (D27). + * N'écrit RIEN en cache : la publication est le `commit` de singleFlight.run. */ async function loadElements(app, elementType, key) { console.error(`[AD] Cache expired or empty, fetching ${key}...`); @@ -146,11 +156,6 @@ async function loadElements(app, elementType, key) { offset += pageSize; } - // Update cache - cache[key] = allElements; - cacheTimestamps[key] = Date.now(); - - console.error(`[AD] Successfully cached ${allElements.length} ${key}`); return allElements; } catch (error) { console.error(`[AD] Error fetching ${key}:`, error.message); @@ -251,12 +256,14 @@ function invalidateCache(elementType = null) { delete cache[k]; delete cacheTimestamps[k]; }); + singleFlight.invalidate(); console.error(`[AD] Cache invalidated: ${elementType}`); } else { Object.keys(cache).forEach(k => { delete cache[k]; delete cacheTimestamps[k]; }); + singleFlight.invalidate(); console.error('[AD] All caches invalidated'); } } diff --git a/src/services/entity-resolver.js b/src/services/entity-resolver.js index 50a5751..f30455d 100644 --- a/src/services/entity-resolver.js +++ b/src/services/entity-resolver.js @@ -43,12 +43,20 @@ async function loadResolutionMap() { if (isCacheValid()) return; // Un seul chargement Metadata, même sous rafale concurrente (D27) : sans - // lui, 6 appels concurrents déclenchaient 6 chargements complets. - await singleFlight.run(METADATA_KEY, fetchResolutionMap); + // lui, 6 appels concurrents déclenchaient 6 chargements complets. La + // publication est refusée si le cache a été invalidé pendant le fetch. + await singleFlight.run(METADATA_KEY, fetchResolutionMap, (loaded) => { + resolutionMap = loaded.map; + tableNames = loaded.names; + cacheTimestamp = Date.now(); + console.error(`[EntityResolver] Cached ${tableNames.length} entities from ${loaded.applicationCount} application(s)`); + }); } /** - * Chargement réel de la table de résolution (D27). + * Chargement réel de la table de résolution (D27). N'écrit rien en cache : + * la publication est le `commit` de singleFlight.run. + * @returns {Promise<{map: Map, names: string[], applicationCount: number}>} */ async function fetchResolutionMap() { console.error('[EntityResolver] Cache expired or empty, fetching Metadata...'); @@ -81,10 +89,11 @@ async function fetchResolutionMap() { throw new Error('Metadata API returned no entity for any application'); } - resolutionMap = map; - tableNames = Array.from(names).sort(); - cacheTimestamp = Date.now(); - console.error(`[EntityResolver] Cached ${tableNames.length} entities from ${appNames.length} application(s)`); + return { + map, + names: Array.from(names).sort(), + applicationCount: appNames.length, + }; } /** @@ -193,6 +202,8 @@ function invalidateCache() { resolutionMap = null; tableNames = null; cacheTimestamp = null; + // Les fetchs déjà partis ne repeupleront pas la table (D27). + singleFlight.invalidate(); console.error('[EntityResolver] Cache cleared'); } diff --git a/src/services/single-flight.js b/src/services/single-flight.js index 6fc6922..5705666 100644 --- a/src/services/single-flight.js +++ b/src/services/single-flight.js @@ -1,40 +1,71 @@ /** - * Single-flight — déduplication des chargements paresseux en vol (D27) + * Single-flight + génération de cache — chargements paresseux sous + * concurrence (D27) * * Les services à cache (workflow, AD, resolver) chargent paresseusement : le * premier appelant qui trouve le cache invalide déclenche le fetch. Le serveur - * traitant les `tools/call` en concurrence, N appelants arrivés pendant ce - * fetch trouvaient tous le cache invalide et lançaient N chaînes complètes - * (mesuré : 6 chargements Metadata en parallèle pour une seule table). + * traitant les `tools/call` en concurrence, deux défauts en découlaient, et ce + * module porte les deux : * - * Le motif est le même partout : une Map `clé de cache -> promesse en vol`, - * pas de dépendance externe. Le single-flight est **par clé** — deux - * applications différentes se chargent toujours en parallèle (D26 : rien n'est - * préchargé, et rien n'est sérialisé au-delà de la clé demandée). + * 1. **Duplication** — N appelants arrivés pendant un fetch trouvaient tous le + * cache invalide et lançaient N chaînes complètes (mesuré : 6 chargements + * Metadata en parallèle pour une seule table). Une Map + * `clé de cache -> promesse en vol` les fait rejoindre le fetch en cours. + * 2. **Écriture post-invalidation** — un fetch parti avant une bascule de + * profil (D8) terminait après elle et repeuplait le cache avec les données + * de l'ancien tenant, timestamp neuf. Un compteur de génération, incrémenté + * à chaque invalidation, fait **jeter** un résultat d'une génération + * périmée au lieu de l'écrire. + * + * Le single-flight est **par clé** — deux applications différentes se chargent + * toujours en parallèle (D26 : rien n'est préchargé, rien n'est sérialisé + * au-delà de la clé demandée). Pas de dépendance externe : une Map. * * @param {string} label - préfixe de log du service appelant (D6, MONITORING §2) */ function createSingleFlight(label) { const inFlight = new Map(); // clé de cache -> promesse du chargement en cours + let generation = 0; // incrémenté à chaque invalidation /** - * Exécute `fetcher` pour cette clé, ou rejoint le chargement déjà en vol. + * Exécute `fetcher` pour cette clé, ou rejoint le chargement déjà en vol, + * puis publie le résultat via `commit` **si la génération n'a pas changé**. + * * La promesse est retirée de la Map au règlement, succès **ou** échec : un * fetch en erreur ne reste pas coincé, l'appel suivant refetche. * + * L'appelant reçoit toujours le résultat de son fetch, même périmé — c'est + * sa **mise en cache** qui est refusée, pas sa réponse : il a demandé ces + * données avant l'invalidation, il les obtient. + * * @param {string} key - clé de cache (une par entrée de cache indépendante) * @param {() => Promise} fetcher - le chargement réel, appelé au plus - * une fois tant qu'il est en vol + * une fois tant qu'il est en vol ; il ne doit **rien** écrire en cache + * @param {(value: any) => void} [commit] - publication en cache, appelée + * seulement si aucune invalidation n'est survenue pendant le fetch * @returns {Promise} le résultat du chargement (partagé par les joignants) */ - function run(key, fetcher) { + function run(key, fetcher, commit) { const pending = inFlight.get(key); if (pending) { console.error(`[${label}] Fetch already in flight for "${key}", joining it`); return pending; } - const promise = (async () => fetcher())(); + const startGeneration = generation; + const promise = (async () => { + const value = await fetcher(); + if (generation !== startGeneration) { + console.error( + `[${label}] Result for "${key}" discarded, not cached: ` + + `cache invalidated during fetch (generation ${startGeneration} -> ${generation})` + ); + return value; + } + if (commit) commit(value); + return value; + })(); + inFlight.set(key, promise); // Libération au règlement. Le test d'identité évite qu'une promesse @@ -49,11 +80,16 @@ function createSingleFlight(label) { } /** - * Oublie les promesses en vol (invalidation de cache). Les appelants déjà - * en attente reçoivent bien leur résultat — c'est la mise en cache de ce - * résultat qui est refusée, ailleurs. + * Marque toutes les données en vol comme périmées : la génération avance et + * la Map est vidée. À appeler depuis l'invalidation du service (l'abonnement + * `onSwitch()` reste le seul déclencheur, D8). + * + * Vider la Map ne coupe personne : les appelants déjà en attente gardent + * leur référence à la promesse et reçoivent son résultat — simplement, ce + * résultat ne sera pas mis en cache. */ - function clear() { + function invalidate() { + generation++; inFlight.clear(); } @@ -64,7 +100,7 @@ function createSingleFlight(label) { return inFlight.size; } - return { run, clear, pendingCount }; + return { run, invalidate, pendingCount }; } module.exports = { createSingleFlight }; diff --git a/src/services/workflow-service.js b/src/services/workflow-service.js index 4488355..e9cae83 100644 --- a/src/services/workflow-service.js +++ b/src/services/workflow-service.js @@ -61,13 +61,25 @@ async function fetchAllWorkflows(application) { return workflowCaches[app]; } - // Un seul chargement par application, même sous rafale concurrente (D27). - return singleFlight.run(`workflows::${app}`, () => loadWorkflows(app)); + // Un seul chargement par application, même sous rafale concurrente, et + // publication en cache seulement si aucune invalidation n'est survenue + // pendant le fetch (D27). + return singleFlight.run( + `workflows::${app}`, + () => loadWorkflows(app), + (workflows) => { + workflowCaches[app] = workflows; + cacheTimestamps[app] = Date.now(); + console.error(`[Workflow] Successfully cached ${workflows.length} workflows for "${app}"`); + } + ); } /** * Chargement réel des workflows d'une application (pagination complète). * Appelé au plus une fois par application tant qu'il est en vol (D27). + * N'écrit RIEN en cache : la publication est le `commit` de singleFlight.run, + * qui la refuse si le cache a été invalidé entre-temps. */ async function loadWorkflows(app) { console.error(`[Workflow] Cache expired or empty for "${app}", fetching from API...`); @@ -107,11 +119,6 @@ async function loadWorkflows(app) { offset += pageSize; } - // Update cache - workflowCaches[app] = allWorkflows; - cacheTimestamps[app] = Date.now(); - - console.error(`[Workflow] Successfully cached ${allWorkflows.length} workflows for "${app}"`); return allWorkflows; } catch (error) { console.error(`[Workflow] Error fetching workflows for "${app}":`, error.message); @@ -130,27 +137,27 @@ async function fetchApplications() { const cached = getCachedApplications(); if (cached) return cached; - // Même déduplication que les workflows, sur sa propre clé (D27). - return singleFlight.run('applications', loadApplications); + // Même déduplication et même garde de génération, sur sa propre clé (D27). + return singleFlight.run('applications', loadApplications, (applications) => { + applicationsCache = applications; + applicationsTimestamp = Date.now(); + console.error(`[Workflow] Cached ${applications.length} application(s)`); + }); } /** - * Chargement réel de la liste d'applications (D27). + * Chargement réel de la liste d'applications (D27). N'écrit rien en cache. */ async function loadApplications() { console.error('[Workflow] Fetching application list (Application/GetAll)...'); const response = await apiService.post('/Application/GetAll', null, true); const entities = response?.entities || []; - applicationsCache = entities.map(a => ({ + return entities.map(a => ({ name: a.name || a.Name, id: a.id || a.Id, version: a.version ?? a.Version, })).filter(a => a.name); - applicationsTimestamp = Date.now(); - - console.error(`[Workflow] Cached ${applicationsCache.length} application(s)`); - return applicationsCache; } /** @@ -290,6 +297,8 @@ function clearCache() { cacheTimestamps = {}; applicationsCache = null; applicationsTimestamp = null; + // Les fetchs déjà partis ne repeupleront pas ce cache (D27). + singleFlight.invalidate(); console.error('[Workflow] Cache cleared'); }