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'); }