L6.1 : dedupliquer les chargements paresseux en vol (single-flight)

Le serveur traite les tools/call en concurrence. Les trois services a cache
chargent paresseusement sans se coordonner : le premier appelant qui trouve le
cache invalide lance le fetch, et tous ceux qui arrivent pendant ce fetch le
trouvent *encore* invalide et lancent le leur. Une rafale de 6 appels
identiques declenchait donc 6 chargements complets pour une seule cle.

Ce n'est pas qu'un gaspillage : la duplication surcharge l'API AD au point de
la faire echouer. Rafale mixte de 14 appels, avant correction — les 4 appels
EasyWMS (~4000 workflows) reviennent en erreur, les memes passent en
sequentiel :

  [Workflow] Error fetching workflows for "EasyWMS": POST
  https://10.255.255.2/AD/api/Workflow/GetByApplication failed (HTTP 500)
  fetching from API: 8 | EntityResolver Cache expired or empty: 3

Motif commun extrait dans src/services/single-flight.js — une Map de
promesses, pas de dependance externe. Une cle par entree de cache
(workflows::<app>, applications, <app>::<type>, metadata) : deux cles
distinctes se chargent toujours en parallele, aucun prechargement (D26
intact). La promesse est retiree au reglement, succes *ou* echec, pour qu'un
fetch en erreur ne reste pas coince.

Le log de fetch reste l'observable (un par chargement reel) ; les appelants
joints emettent une ligne distincte "Fetch already in flight ... joining it".

--- Verifications (LIMAGRAIN), rafales rejouees 3 fois ---

Phase 0, reproduction avant correction :
  6 x search_workflows CustomApp  -> count 44 x6, 'fetching from API' : 6
  6 x query_wms_entities Container -> 6 succes, 'EntityResolver] Cache
                                      expired or empty' : 6

Rafale de 6 search_workflows {"query":"CST_","application":"CustomApp"} :

  ===== RUN 1 =====            ===== RUN 2 =====            ===== RUN 3 =====
  id 10 success=true application=CustomApp count=44   (idem RUN 2 et RUN 3,
  id 11 success=true application=CustomApp count=44    les 6 reponses a 44)
  id 12 success=true application=CustomApp count=44
  id 13 success=true application=CustomApp count=44
  id 14 success=true application=CustomApp count=44
  id 15 success=true application=CustomApp count=44
  -- 'fetching from API' : 1  | 'joining it' : 5   [RUN 1]
  -- 'fetching from API' : 1  | 'joining it' : 5   [RUN 2]
  -- 'fetching from API' : 1  | 'joining it' : 5   [RUN 3]

Rafale de 6 query_wms_entities {"entity_type":"Container","limit":1} :

  RUN 1/2/3 : id 10..15 success=true count=1 (6/6)
  -- 'EntityResolver] Cache expired or empty' : 1 | joins : 5 | GET
     Metadata/Entities : 5   [identique RUN 1, RUN 2, RUN 3]
  (avant : 6 chargements, soit 30 GET Metadata)

Rafale mixte EasyWMS + CustomApp (3 + 3) — un fetch par application :

  RUN 1/2/3 : CustomApp count=44 x3, EasyWMS count=50 x3
  [Workflow] Cache expired or empty for "CustomApp", fetching from API...
  [Workflow] Cache expired or empty for "EasyWMS", fetching from API...
     total fetch=2 joins=4   [identique RUN 1, RUN 2, RUN 3]
  Plus aucun HTTP 500 : un seul fetch EasyWMS concurrent au lieu de 4.

Chemin sequentiel nominal, strictement inchange (driver sequentiel) :

  [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

Liberation de la Map sur echec (test direct, apiService.post substitue :
echoue au 1er appel, reussit ensuite) :

  [Workflow] Cache expired or empty for "TestApp", fetching from API...
  [Workflow] Fetch already in flight for "workflows::TestApp", joining it  (x2)
  --- rafale de 3 sur un fetch en echec :
    appelant 0/1/2: rejected - Failed to fetch workflows ... panne reseau simulee
    appels reseau reels: 1 (attendu 1 : les 3 partagent le meme fetch)
    cache pose ? {} (attendu {} : rien en cache sur echec)
  --- appel suivant (la Map doit avoir ete liberee) :
    resultat: 1 workflow(s), appels reseau cumules: 2

Baseline : tools/list 23, resources/list 6 ; npm test 4/4 exit 0.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Arthur Ria
2026-08-25 15:43:56 +02:00
parent 8a631f5ea4
commit b7151b3bc7
7 changed files with 173 additions and 10 deletions
+4
View File
@@ -199,6 +199,10 @@ par un appel réseau (D26).
Les tailles de page par type viennent de l'observation des timeouts serveur — Les tailles de page par type viennent de l'observation des timeouts serveur —
ne les augmentez pas à l'aveugle. 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.
`get_application_summary` expose l'état des caches sans redémarrage. `get_application_summary` expose l'état des caches sans redémarrage.
--- ---
+48
View File
@@ -645,3 +645,51 @@ passé), et ajoutent un `hint` quand la recherche revient **vide** :
- Il nomme `CustomApp` en clair, sauf quand c'est déjà l'application - Il nomme `CustomApp` en clair, sauf quand c'est déjà l'application
interrogée : c'est une connaissance statique, déjà portée par les interrogée : c'est une connaissance statique, déjà portée par les
descriptions d'outils, pas une donnée à aller chercher. descriptions d'outils, pas une donnée à aller chercher.
---
## D27 — Chargements paresseux sous concurrence : single-flight par clé
**Contexte (mesuré le 25/08/2026, `LIMAGRAI2512`).** Le serveur traite les
`tools/call` **en concurrence** : une rafale d'appels dans une même session
s'exécute en parallèle. Les trois services à cache chargeaient paresseusement
sans se coordonner — le premier appelant qui trouve le cache invalide lance le
fetch, et tous ceux qui arrivent pendant ce fetch trouvent le cache **encore**
invalide et lancent le leur. Mesures avant correction :
| Rafale | Résultat |
|---|---|
| 6 × `search_workflows` (CustomApp) | 6 × `fetching from API` pour une seule clé |
| 6 × `query_wms_entities` (Container) | 6 × `[EntityResolver] Cache expired or empty` — soit 30 GET Metadata |
| 14 appels mixtes | l'API AD répond **HTTP 500** sur `Workflow/GetByApplication` (EasyWMS, ~4 000 workflows) — les 4 appels EasyWMS échouent, les mêmes passent en séquentiel |
La dernière ligne est le vrai coût : la duplication ne gaspille pas seulement
des appels, elle **surcharge l'API AD au point de la faire échouer**.
**Décision.** Un motif unique, `src/services/single-flight.js`, partagé par
`workflow-service`, `ad-service` et `entity-resolver` — une `Map` de promesses,
pas une dépendance externe :
- **Une clé de single-flight par entrée de cache** : `workflows::<app>` et
`applications` pour les workflows, `<app>::<type>` pour l'AD, une clé unique
pour le resolver. Deux clés distinctes se chargent toujours **en parallèle**
le single-flight ne sérialise rien au-delà de la clé demandée, et
n'introduit aucun préchargement (D26 intact).
- **La promesse est retirée au règlement, succès *ou* échec.** Un fetch en
erreur ne reste pas coincé dans la Map : l'appel suivant refetche. Les
appelants joints reçoivent la même erreur, et rien n'est mis en cache.
- **Le log de fetch reste l'observable** (D6) : une ligne `fetching from API`
/ `Cache expired or empty` par chargement **réel**. Les appelants joints
émettent une ligne distincte (`Fetch already in flight … joining it`) — ne
fusionnez pas les deux, c'est ce qui rend la déduplication vérifiable depuis
stderr.
**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
total, plus aucun HTTP 500. Le chemin séquentiel nominal est inchangé : 1 fetch
puis 1 `Using cached data`, 0 join.
**Ce que cette décision ne couvre pas.** La bascule de profil concurrente aux
appels en vol (un `switch_wms_profile` qui redirige des requêtes déjà parties)
reste un point ouvert de la ROADMAP, distinct.
-6
View File
@@ -95,12 +95,6 @@ Rejoués en séquentiel, les mêmes appels sont corrects.
Trois défauts structurels dans les services à cache (`workflow-service`, Trois défauts structurels dans les services à cache (`workflow-service`,
`ad-service`, `entity-resolver`) : `ad-service`, `entity-resolver`) :
### L6.1 — Dédupliquer les fetchs en vol
Aucun service ne mémorise la promesse de chargement en cours : N appelants
concurrents sur la même clé déclenchent N chaînes de fetch complètes.
Single-flight par clé de cache (promesse partagée, libérée au règlement).
### L6.2 — Protéger l'invalidation contre les fetchs en vol ### L6.2 — Protéger l'invalidation contre les fetchs en vol
Un fetch parti avant `clearCache()`/`invalidateCache()` (bascule de profil, Un fetch parti avant `clearCache()`/`invalidateCache()` (bascule de profil,
+13
View File
@@ -6,12 +6,16 @@
const apiService = require('./api-service').getInstance(); const apiService = require('./api-service').getInstance();
const profileManager = require('../config/profile-manager'); const profileManager = require('../config/profile-manager');
const { createSingleFlight } = require('./single-flight');
// Cache state - one cache per (application, element type) (D26) // Cache state - one cache per (application, element type) (D26)
const cache = {}; const cache = {};
const cacheTimestamps = {}; const cacheTimestamps = {};
const CACHE_TTL = parseInt(process.env.WORKFLOW_CACHE_TTL) || 3600000; // 1 hour const CACHE_TTL = parseInt(process.env.WORKFLOW_CACHE_TTL) || 3600000; // 1 hour
// Déduplication des chargements concurrents, une clé par (application, type) (D27).
const singleFlight = createSingleFlight('AD');
/** /**
* Application effective : celle demandée, sinon celle du profil actif. * Application effective : celle demandée, sinon celle du profil actif.
*/ */
@@ -96,6 +100,15 @@ async function getElements(elementType, application) {
return cache[key]; return cache[key];
} }
// Un seul chargement par (application, type), même sous rafale (D27).
return singleFlight.run(key, () => loadElements(app, elementType, 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).
*/
async function loadElements(app, elementType, key) {
console.error(`[AD] Cache expired or empty, fetching ${key}...`); console.error(`[AD] Cache expired or empty, fetching ${key}...`);
try { try {
+14
View File
@@ -11,6 +11,7 @@
const apiService = require('./api-service').getInstance(); const apiService = require('./api-service').getInstance();
const profileManager = require('../config/profile-manager'); const profileManager = require('../config/profile-manager');
const { createSingleFlight } = require('./single-flight');
// Cache state — même TTL que les autres caches (D10) // Cache state — même TTL que les autres caches (D10)
let resolutionMap = null; // Map lower(Name | TableName) -> TableName let resolutionMap = null; // Map lower(Name | TableName) -> TableName
@@ -18,6 +19,10 @@ let tableNames = null; // TableName[] triés (suggestions + comptage)
let cacheTimestamp = null; let cacheTimestamp = null;
const CACHE_TTL = parseInt(process.env.WORKFLOW_CACHE_TTL) || 3600000; const CACHE_TTL = parseInt(process.env.WORKFLOW_CACHE_TTL) || 3600000;
// Table unique : une seule clé de single-flight (D27).
const singleFlight = createSingleFlight('EntityResolver');
const METADATA_KEY = 'metadata';
// La table de résolution est par tenant — invalidée à chaque bascule (D8). // La table de résolution est par tenant — invalidée à chaque bascule (D8).
profileManager.onSwitch(() => invalidateCache()); profileManager.onSwitch(() => invalidateCache());
@@ -37,6 +42,15 @@ function isCacheValid() {
async function loadResolutionMap() { async function loadResolutionMap() {
if (isCacheValid()) return; 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);
}
/**
* Chargement réel de la table de résolution (D27).
*/
async function fetchResolutionMap() {
console.error('[EntityResolver] Cache expired or empty, fetching Metadata...'); console.error('[EntityResolver] Cache expired or empty, fetching Metadata...');
const apps = await apiService.get('/configuration/applications'); const apps = await apiService.get('/configuration/applications');
+70
View File
@@ -0,0 +1,70 @@
/**
* Single-flight déduplication des chargements paresseux en vol (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).
*
* 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).
*
* @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
/**
* Exécute `fetcher` pour cette clé, ou rejoint le chargement déjà en vol.
* 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.
*
* @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
* @returns {Promise<any>} le résultat du chargement (partagé par les joignants)
*/
function run(key, fetcher) {
const pending = inFlight.get(key);
if (pending) {
console.error(`[${label}] Fetch already in flight for "${key}", joining it`);
return pending;
}
const promise = (async () => fetcher())();
inFlight.set(key, promise);
// Libération au règlement. Le test d'identité évite qu'une promesse
// périmée (Map vidée par une invalidation, puis nouveau fetch démarré)
// supprime l'entrée de son successeur.
const release = () => {
if (inFlight.get(key) === promise) inFlight.delete(key);
};
promise.then(release, release);
return promise;
}
/**
* 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.
*/
function clear() {
inFlight.clear();
}
/**
* Nombre de chargements en vol diagnostic seulement.
*/
function pendingCount() {
return inFlight.size;
}
return { run, clear, pendingCount };
}
module.exports = { createSingleFlight };
+24 -4
View File
@@ -9,6 +9,7 @@
const apiService = require('./api-service').getInstance(); const apiService = require('./api-service').getInstance();
const profileManager = require('../config/profile-manager'); const profileManager = require('../config/profile-manager');
const { createSingleFlight } = require('./single-flight');
// Cache state — un cache de workflows par application (D26) // Cache state — un cache de workflows par application (D26)
let workflowCaches = {}; // application -> workflows[] let workflowCaches = {}; // application -> workflows[]
@@ -17,6 +18,10 @@ let applicationsCache = null; // liste allégée de POST /Application/GetAll
let applicationsTimestamp = null; let applicationsTimestamp = null;
const CACHE_TTL = parseInt(process.env.WORKFLOW_CACHE_TTL) || 3600000; // 1 hour in milliseconds const CACHE_TTL = parseInt(process.env.WORKFLOW_CACHE_TTL) || 3600000; // 1 hour in milliseconds
// Déduplication des chargements concurrents, par clé de cache (D27). Deux
// clés distinctes ici : une par application, plus la liste d'applications.
const singleFlight = createSingleFlight('Workflow');
// Clear cache when profile changes — workflows are per-tenant, so the previous // Clear cache when profile changes — workflows are per-tenant, so the previous
// profile's cache is meaningless after a switch. // profile's cache is meaningless after a switch.
profileManager.onSwitch(() => clearCache()); profileManager.onSwitch(() => clearCache());
@@ -56,6 +61,15 @@ async function fetchAllWorkflows(application) {
return workflowCaches[app]; return workflowCaches[app];
} }
// Un seul chargement par application, même sous rafale concurrente (D27).
return singleFlight.run(`workflows::${app}`, () => loadWorkflows(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).
*/
async function loadWorkflows(app) {
console.error(`[Workflow] Cache expired or empty for "${app}", fetching from API...`); console.error(`[Workflow] Cache expired or empty for "${app}", fetching from API...`);
try { try {
@@ -113,11 +127,17 @@ async function fetchAllWorkflows(application) {
* @returns {Promise<Array<{name: string, id: string, version: number}>>} * @returns {Promise<Array<{name: string, id: string, version: number}>>}
*/ */
async function fetchApplications() { async function fetchApplications() {
if (applicationsCache && applicationsTimestamp && const cached = getCachedApplications();
Date.now() - applicationsTimestamp < CACHE_TTL) { if (cached) return cached;
return applicationsCache;
}
// Même déduplication que les workflows, sur sa propre clé (D27).
return singleFlight.run('applications', loadApplications);
}
/**
* Chargement réel de la liste d'applications (D27).
*/
async function loadApplications() {
console.error('[Workflow] Fetching application list (Application/GetAll)...'); console.error('[Workflow] Fetching application list (Application/GetAll)...');
const response = await apiService.post('/Application/GetAll', null, true); const response = await apiService.post('/Application/GetAll', null, true);
const entities = response?.entities || []; const entities = response?.entities || [];