Compare commits
6 Commits
1851b38c62
...
5eadc1f6a6
| Author | SHA1 | Date | |
|---|---|---|---|
| 5eadc1f6a6 | |||
| b758a0e09d | |||
| 242b0c0f1c | |||
| 097b76c7ef | |||
| b7151b3bc7 | |||
| 8a631f5ea4 |
@@ -64,6 +64,8 @@ src/
|
|||||||
│ ├── workflow-service.js Workflows, lazy loading + cache
|
│ ├── workflow-service.js Workflows, lazy loading + cache
|
||||||
│ ├── ad-service.js Application Dictionary, 20 types, cache par type
|
│ ├── ad-service.js Application Dictionary, 20 types, cache par type
|
||||||
│ ├── wms-query-service.js Construction d'expressions LINQ
|
│ ├── wms-query-service.js Construction d'expressions LINQ
|
||||||
|
│ ├── single-flight.js Déduplication des chargements + génération (D27)
|
||||||
|
│ ├── ad-envelope.js Enveloppe { entities } des API AD : vide anormal = erreur (D27)
|
||||||
│ ├── log-service.js Lecture et recherche dans les fichiers de logs
|
│ ├── log-service.js Lecture et recherche dans les fichiers de logs
|
||||||
│ └── response-limit.js Plafond de taille commun aux outils de requête (D24)
|
│ └── response-limit.js Plafond de taille commun aux outils de requête (D24)
|
||||||
└── tools/ 23 outils MCP
|
└── tools/ 23 outils MCP
|
||||||
@@ -199,6 +201,18 @@ 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. 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.**
|
||||||
|
|
||||||
|
Une réponse d'API AD **hors enveloppe** `{ entities: [...] }` lève au lieu de
|
||||||
|
passer pour un tableau vide : sinon un cache vide s'installe pour tout le TTL
|
||||||
|
(D27). Un `entities: []` **réel** reste cachable — des applications sont
|
||||||
|
légitimement vides.
|
||||||
|
|
||||||
`get_application_summary` expose l'état des caches sans redémarrage.
|
`get_application_summary` expose l'état des caches sans redémarrage.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|||||||
@@ -645,3 +645,100 @@ 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 + 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
|
||||||
|
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.
|
||||||
|
|
||||||
|
**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é.
|
||||||
|
|
||||||
|
**Vide anormal ≠ vide réel.** Les services lisaient `response?.entities || []`
|
||||||
|
sur les réponses de l'API AD (enveloppe `{ entities: [...] }`, D4). Toute
|
||||||
|
réponse d'une **autre forme** devenait donc un tableau vide, indistinguable
|
||||||
|
d'une page finale légitime — et mise en cache avec un timestamp valide : un
|
||||||
|
cache vide empoisonné pour tout le TTL, sans message. C'est la cause probable
|
||||||
|
du `count: 0` mesuré sous rafale, et le mode d'échec le plus coûteux du lot :
|
||||||
|
il se lit comme une réponse.
|
||||||
|
|
||||||
|
`src/services/ad-envelope.js` porte le contrat pour les trois sites
|
||||||
|
(`Workflow/GetByApplication`, `Application/GetAll`, `<Type>/GetByApplication`) :
|
||||||
|
|
||||||
|
| Réponse | Traitement |
|
||||||
|
|---|---|
|
||||||
|
| `{ entities: [...] }`, `[]` réel compris | rendue telle quelle — une application peut être légitimement vide (`SmartUI` : 0 workflow ; 3 types AD valides mais vides, D17) |
|
||||||
|
| toute autre forme | **lève** — l'appel échoue, rien n'est mis en cache, l'appel suivant refetche |
|
||||||
|
|
||||||
|
**Pas de retry, pas de résilience.** L'anomalie doit être **visible et non
|
||||||
|
persistante** ; la rattraper la rendrait invisible, ce qui est exactement le
|
||||||
|
défaut corrigé. `entity-resolver` était déjà conforme : il lève déjà si
|
||||||
|
`/configuration/applications` ou le Metadata ne rendent aucune entité.
|
||||||
|
|
||||||
|
Mesures : `search_workflows` sur `SmartUI` → `count: 0`, `success: true`,
|
||||||
|
cache posé (`count: 0`, `valid: true`) et hint L5.4 présent. Les trois sites
|
||||||
|
face à une réponse `{}` → erreur levée, `{}` en cache, et le fetch suivant
|
||||||
|
repart normalement.
|
||||||
|
|
||||||
|
**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.
|
||||||
|
|||||||
+18
-16
@@ -105,23 +105,25 @@ plus riche que `/AD/api/Application/GetAll`), `GET /healthcheck?tenantCode=` et
|
|||||||
test.
|
test.
|
||||||
- **Bascule de profil concurrente aux appels en vol** (mesuré le 25/08/2026,
|
- **Bascule de profil concurrente aux appels en vol** (mesuré le 25/08/2026,
|
||||||
révision du lot 2). Le serveur traite les `tools/call` **en concurrence** :
|
révision du lot 2). Le serveur traite les `tools/call` **en concurrence** :
|
||||||
en envoyant une rafale de requêtes dans une même session, les réponses
|
un `switch_wms_profile` émis pendant que des requêtes sont en vol les fait
|
||||||
reviennent dans le désordre, et un `switch_wms_profile` émis pendant que des
|
partir sur le nouveau profil (observé : une requête destinée à EUROTRAFIC
|
||||||
requêtes sont en vol les fait partir sur le nouveau profil (observé : une
|
exécutée sur l'hôte `10.255.255.2` après la bascule suivante). Conséquence du
|
||||||
requête destinée à EUROTRAFIC exécutée sur l'hôte `10.255.255.2` après la
|
singleton d'état global (D8). À traiter si un cas réel de mélange de profils
|
||||||
bascule suivante). Conséquence du singleton d'état global (D8). Sans gravité
|
|
||||||
pour un usage conversationnel séquentiel, mais Claude peut émettre des appels
|
|
||||||
d'outils **en parallèle** : à traiter si un cas réel de mélange de profils
|
|
||||||
est observé (piste : sérialiser les `tools/call` ou figer le profil résolu au
|
est observé (piste : sérialiser les `tools/call` ou figer le profil résolu au
|
||||||
début de chaque appel). **Seconde manifestation mesurée (25/08/2026, révision
|
début de chaque appel). La manifestation « caches » de la même concurrence
|
||||||
du lot 5)** : sans aucune bascule de profil, une rafale d'appels concurrents
|
est **traitée** (D27 : single-flight par clé, garde de génération) ; celle-ci
|
||||||
pendant le chargement paresseux du cache workflow renvoie des résultats
|
ne l'est pas — D27 borne les chargements paresseux, pas le routage d'une
|
||||||
faux en silence — `search_workflows("CST_", application: "CustomApp")` a
|
requête déjà partie.
|
||||||
répondu `count: 0` (contre 44 en séquentiel), et `get_workflow_details` une
|
- **L'API AD échoue sous appels concurrents nombreux** (mesuré le 25/08/2026,
|
||||||
réponse anormale — les chargements concurrents du même cache ne sont pas
|
lot 6). Avant D27, une rafale de 14 `tools/call` faisait répondre **HTTP 500**
|
||||||
synchronisés (pas de déduplication de fetch en vol). Rejoués en séquentiel,
|
à `POST /AD/api/Workflow/GetByApplication` pour `EasyWMS` (~4 000 workflows) —
|
||||||
les mêmes appels sont corrects. Piste supplémentaire : mémoriser la promesse
|
les 4 appels concernés en erreur, les mêmes corrects en séquentiel. C'est une
|
||||||
de fetch en cours par clé de cache et la partager entre appelants.
|
limite du serveur AD, pas du MCP. D27 l'atténue fortement (un seul fetch par
|
||||||
|
clé au lieu de N, et le 500 n'a pas reparu depuis), sans la supprimer : des
|
||||||
|
clés **différentes** se chargent toujours en parallèle. À reconsidérer si le
|
||||||
|
500 réapparaît — piste : plafonner le nombre de chargements simultanés, tous
|
||||||
|
clés confondues. N'implémentez rien avant d'avoir une mesure : brider les
|
||||||
|
chargements parallèles coûte de la latence sur le chemin nominal.
|
||||||
- **`select_expression`** : les projections via le paramètre `Select` provoquent
|
- **`select_expression`** : les projections via le paramètre `Select` provoquent
|
||||||
des erreurs de compilation côté serveur (D13). Irritant principal restant.
|
des erreurs de compilation côté serveur (D13). Irritant principal restant.
|
||||||
- **Déploiement SSH sur la VM** : l'exécutable est validé, la configuration SSH
|
- **Déploiement SSH sur la VM** : l'exécutable est validé, la configuration SSH
|
||||||
|
|||||||
@@ -0,0 +1,59 @@
|
|||||||
|
/**
|
||||||
|
* Enveloppe des API AD — un vide anormal n'est pas un vide (D27)
|
||||||
|
*
|
||||||
|
* Les API AD renvoient `{ entities: [...] }` (D4). Les services lisaient
|
||||||
|
* `response?.entities || []` : toute réponse d'une **autre forme** (pas de
|
||||||
|
* champ `entities`, corps vide, objet d'erreur) devenait un tableau vide,
|
||||||
|
* indistinguable d'une page finale légitime — donc mise en cache avec un
|
||||||
|
* timestamp valide. Un cache vide empoisonné pour tout le TTL, sans le
|
||||||
|
* moindre message.
|
||||||
|
*
|
||||||
|
* Deux cas, deux traitements :
|
||||||
|
*
|
||||||
|
* | Réponse | Traitement |
|
||||||
|
* |---|---|
|
||||||
|
* | `{ entities: [...] }`, y compris `[]` réel | rendue telle quelle — une application peut être légitimement vide (`SmartUI` : 0 workflow, D26) |
|
||||||
|
* | tout le reste | **lève** — l'appel échoue, rien n'est mis en cache, l'appel suivant refetche |
|
||||||
|
*
|
||||||
|
* Volontairement sans retry ni résilience : le but est de rendre l'anomalie
|
||||||
|
* **visible et non persistante**, pas de la rattraper.
|
||||||
|
*/
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Décrit la forme reçue, pour un message d'erreur exploitable (convention 4).
|
||||||
|
*/
|
||||||
|
function describeShape(response) {
|
||||||
|
if (response === null) return 'null';
|
||||||
|
if (response === undefined) return 'undefined';
|
||||||
|
if (Array.isArray(response)) return `un tableau nu de ${response.length} élément(s)`;
|
||||||
|
if (typeof response !== 'object') return `un ${typeof response}`;
|
||||||
|
const keys = Object.keys(response);
|
||||||
|
if (keys.length === 0) return 'un objet vide';
|
||||||
|
return `un objet sans champ "entities" (champs reçus : ${keys.slice(0, 10).join(', ')})`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Extrait le tableau `entities` d'une réponse d'API AD, ou lève.
|
||||||
|
*
|
||||||
|
* @param {any} response - la réponse brute de `apiService.post(..., true)`
|
||||||
|
* @param {string} context - l'appel concerné, pour le message d'erreur
|
||||||
|
* (ex. `Workflow/GetByApplication (application "EasyWMS", offset 0)`)
|
||||||
|
* @returns {Array} le tableau `entities`, éventuellement vide
|
||||||
|
* @throws {Error} si la réponse n'a pas la forme `{ entities: [...] }`
|
||||||
|
*/
|
||||||
|
function requireEntities(response, context) {
|
||||||
|
const entities = response ? response.entities : undefined;
|
||||||
|
|
||||||
|
if (!Array.isArray(entities)) {
|
||||||
|
throw new Error(
|
||||||
|
`Réponse inattendue de l'API AD sur ${context} : ${describeShape(response)}, ` +
|
||||||
|
`au lieu de l'enveloppe attendue { entities: [...] }. ` +
|
||||||
|
`Rien n'a été mis en cache — relancez l'appel. ` +
|
||||||
|
`Si l'erreur persiste, l'API AD est en défaut (elle échoue notamment sous appels concurrents nombreux).`
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
return entities;
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { requireEntities };
|
||||||
@@ -6,12 +6,17 @@
|
|||||||
|
|
||||||
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');
|
||||||
|
const { requireEntities } = require('./ad-envelope');
|
||||||
|
|
||||||
// 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 +101,25 @@ async function getElements(elementType, application) {
|
|||||||
return cache[key];
|
return cache[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}...`);
|
console.error(`[AD] Cache expired or empty, fetching ${key}...`);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
@@ -112,11 +136,15 @@ async function getElements(elementType, application) {
|
|||||||
// Use AD API (useAdApi=true)
|
// Use AD API (useAdApi=true)
|
||||||
const response = await apiService.post(`/${elementType}/GetByApplication`, body, true);
|
const response = await apiService.post(`/${elementType}/GetByApplication`, body, true);
|
||||||
|
|
||||||
// Extract entities array from response
|
// Une réponse hors enveloppe { entities: [...] } lève au lieu de se
|
||||||
const elements = response?.entities || [];
|
// faire passer pour une page vide (D27).
|
||||||
|
const elements = requireEntities(
|
||||||
|
response,
|
||||||
|
`${elementType}/GetByApplication (application "${app}", offset ${offset})`
|
||||||
|
);
|
||||||
|
|
||||||
// Check if response is valid
|
// Vide réel : fin de pagination (3 types sont valides mais vides, D17).
|
||||||
if (!elements || elements.length === 0) {
|
if (elements.length === 0) {
|
||||||
console.error(`[AD] No more ${elementType} to fetch`);
|
console.error(`[AD] No more ${elementType} to fetch`);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -133,11 +161,6 @@ async function getElements(elementType, application) {
|
|||||||
offset += pageSize;
|
offset += pageSize;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Update cache
|
|
||||||
cache[key] = allElements;
|
|
||||||
cacheTimestamps[key] = Date.now();
|
|
||||||
|
|
||||||
console.error(`[AD] Successfully cached ${allElements.length} ${key}`);
|
|
||||||
return allElements;
|
return allElements;
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[AD] Error fetching ${key}:`, error.message);
|
console.error(`[AD] Error fetching ${key}:`, error.message);
|
||||||
@@ -238,12 +261,14 @@ function invalidateCache(elementType = null) {
|
|||||||
delete cache[k];
|
delete cache[k];
|
||||||
delete cacheTimestamps[k];
|
delete cacheTimestamps[k];
|
||||||
});
|
});
|
||||||
|
singleFlight.invalidate();
|
||||||
console.error(`[AD] Cache invalidated: ${elementType}`);
|
console.error(`[AD] Cache invalidated: ${elementType}`);
|
||||||
} else {
|
} else {
|
||||||
Object.keys(cache).forEach(k => {
|
Object.keys(cache).forEach(k => {
|
||||||
delete cache[k];
|
delete cache[k];
|
||||||
delete cacheTimestamps[k];
|
delete cacheTimestamps[k];
|
||||||
});
|
});
|
||||||
|
singleFlight.invalidate();
|
||||||
console.error('[AD] All caches invalidated');
|
console.error('[AD] All caches invalidated');
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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,23 @@ 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. 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). 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...');
|
console.error('[EntityResolver] Cache expired or empty, fetching Metadata...');
|
||||||
|
|
||||||
const apps = await apiService.get('/configuration/applications');
|
const apps = await apiService.get('/configuration/applications');
|
||||||
@@ -67,10 +89,11 @@ async function loadResolutionMap() {
|
|||||||
throw new Error('Metadata API returned no entity for any application');
|
throw new Error('Metadata API returned no entity for any application');
|
||||||
}
|
}
|
||||||
|
|
||||||
resolutionMap = map;
|
return {
|
||||||
tableNames = Array.from(names).sort();
|
map,
|
||||||
cacheTimestamp = Date.now();
|
names: Array.from(names).sort(),
|
||||||
console.error(`[EntityResolver] Cached ${tableNames.length} entities from ${appNames.length} application(s)`);
|
applicationCount: appNames.length,
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -179,6 +202,8 @@ function invalidateCache() {
|
|||||||
resolutionMap = null;
|
resolutionMap = null;
|
||||||
tableNames = null;
|
tableNames = null;
|
||||||
cacheTimestamp = null;
|
cacheTimestamp = null;
|
||||||
|
// Les fetchs déjà partis ne repeupleront pas la table (D27).
|
||||||
|
singleFlight.invalidate();
|
||||||
console.error('[EntityResolver] Cache cleared');
|
console.error('[EntityResolver] Cache cleared');
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,106 @@
|
|||||||
|
/**
|
||||||
|
* 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, deux défauts en découlaient, et ce
|
||||||
|
* module porte les deux :
|
||||||
|
*
|
||||||
|
* 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,
|
||||||
|
* 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 ; 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, commit) {
|
||||||
|
const pending = inFlight.get(key);
|
||||||
|
if (pending) {
|
||||||
|
console.error(`[${label}] Fetch already in flight for "${key}", joining it`);
|
||||||
|
return pending;
|
||||||
|
}
|
||||||
|
|
||||||
|
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
|
||||||
|
// 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;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 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 invalidate() {
|
||||||
|
generation++;
|
||||||
|
inFlight.clear();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Nombre de chargements en vol — diagnostic seulement.
|
||||||
|
*/
|
||||||
|
function pendingCount() {
|
||||||
|
return inFlight.size;
|
||||||
|
}
|
||||||
|
|
||||||
|
return { run, invalidate, pendingCount };
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { createSingleFlight };
|
||||||
@@ -9,6 +9,8 @@
|
|||||||
|
|
||||||
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');
|
||||||
|
const { requireEntities } = require('./ad-envelope');
|
||||||
|
|
||||||
// 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 +19,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 +62,27 @@ async function fetchAllWorkflows(application) {
|
|||||||
return workflowCaches[app];
|
return workflowCaches[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...`);
|
console.error(`[Workflow] Cache expired or empty for "${app}", fetching from API...`);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
@@ -72,11 +99,17 @@ async function fetchAllWorkflows(application) {
|
|||||||
// Use AD API (useAdApi=true)
|
// Use AD API (useAdApi=true)
|
||||||
const response = await apiService.post('/Workflow/GetByApplication', body, true);
|
const response = await apiService.post('/Workflow/GetByApplication', body, true);
|
||||||
|
|
||||||
// Extract entities array from response
|
// Une réponse hors enveloppe { entities: [...] } lève au lieu de se
|
||||||
const workflows = response?.entities || [];
|
// faire passer pour une page vide (D27) : un cache vide empoisonné
|
||||||
|
// durerait tout le TTL.
|
||||||
|
const workflows = requireEntities(
|
||||||
|
response,
|
||||||
|
`Workflow/GetByApplication (application "${app}", offset ${offset})`
|
||||||
|
);
|
||||||
|
|
||||||
// Check if response is valid
|
// Vide réel : fin de pagination (une application peut n'avoir aucun
|
||||||
if (!workflows || workflows.length === 0) {
|
// workflow — SmartUI, D26).
|
||||||
|
if (workflows.length === 0) {
|
||||||
console.error('[Workflow] No more workflows to fetch');
|
console.error('[Workflow] No more workflows to fetch');
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -93,11 +126,6 @@ async function fetchAllWorkflows(application) {
|
|||||||
offset += pageSize;
|
offset += pageSize;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Update cache
|
|
||||||
workflowCaches[app] = allWorkflows;
|
|
||||||
cacheTimestamps[app] = Date.now();
|
|
||||||
|
|
||||||
console.error(`[Workflow] Successfully cached ${allWorkflows.length} workflows for "${app}"`);
|
|
||||||
return allWorkflows;
|
return allWorkflows;
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[Workflow] Error fetching workflows for "${app}":`, error.message);
|
console.error(`[Workflow] Error fetching workflows for "${app}":`, error.message);
|
||||||
@@ -113,24 +141,30 @@ 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 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). N'écrit rien en cache.
|
||||||
|
*/
|
||||||
|
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 = requireEntities(response, 'Application/GetAll');
|
||||||
|
|
||||||
applicationsCache = entities.map(a => ({
|
return entities.map(a => ({
|
||||||
name: a.name || a.Name,
|
name: a.name || a.Name,
|
||||||
id: a.id || a.Id,
|
id: a.id || a.Id,
|
||||||
version: a.version ?? a.Version,
|
version: a.version ?? a.Version,
|
||||||
})).filter(a => a.name);
|
})).filter(a => a.name);
|
||||||
applicationsTimestamp = Date.now();
|
|
||||||
|
|
||||||
console.error(`[Workflow] Cached ${applicationsCache.length} application(s)`);
|
|
||||||
return applicationsCache;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -270,6 +304,8 @@ function clearCache() {
|
|||||||
cacheTimestamps = {};
|
cacheTimestamps = {};
|
||||||
applicationsCache = null;
|
applicationsCache = null;
|
||||||
applicationsTimestamp = null;
|
applicationsTimestamp = null;
|
||||||
|
// Les fetchs déjà partis ne repeupleront pas ce cache (D27).
|
||||||
|
singleFlight.invalidate();
|
||||||
console.error('[Workflow] Cache cleared');
|
console.error('[Workflow] Cache cleared');
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user