L6.2 : garde de generation contre les ecritures post-invalidation

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 <noreply@anthropic.com>
This commit is contained in:
Arthur Ria
2026-08-25 15:49:39 +02:00
parent b7151b3bc7
commit 097b76c7ef
7 changed files with 137 additions and 55 deletions
+4 -1
View File
@@ -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.
+24 -1
View File
@@ -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
-7
View File
@@ -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
+14 -7
View File
@@ -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');
}
}
+18 -7
View File
@@ -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');
}
+53 -17
View File
@@ -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<any>} 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<any>} 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 };
+24 -15
View File
@@ -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');
}