diff --git a/src/app/api/v1/grid/route.ts b/src/app/api/v1/grid/route.ts index 93b96ac..ce27dcf 100644 --- a/src/app/api/v1/grid/route.ts +++ b/src/app/api/v1/grid/route.ts @@ -77,7 +77,7 @@ export async function GET(req: NextRequest) { computedOnDemand = true; } - const result = await queryGridRows({ + const filtres = { codeFournisseur: q.fournisseur, gamme: q.gamme, code1: q.code1, @@ -86,16 +86,21 @@ export async function GET(req: NextRequest) { search: q.search, sort: toSortKey(q.sort), order: q.order, - page: q.page, - limit: q.limit, - }); + }; + + // Sans `limit`, la réponse peut porter tout un fournisseur — et certains en + // comptent plus de 130 000. Tout charger puis tout sérialiser d'un bloc ferait + // exploser la mémoire du processus : on diffuse par lots, à mémoire constante. + if (q.limit == null) { + return streamAllRows(req, filtres, q, freshness, computedOnDemand); + } + + const result = await queryGridRows({ ...filtres, page: q.page, limit: q.limit }); // Métriques Qlik et gamme serveur relues au moment de l'appel (voir api-enrich). const rows = q.enrich === "1" ? await enrichRows(result.rows) : result.rows; - // Sans `limit`, tout a été renvoyé : la pagination doit le refléter (une seule - // page couvrant le total), sinon `hasMore` mentirait. - const pagination = buildPagination(q.page, q.limit ?? Math.max(1, result.total), result.total); + const pagination = buildPagination(q.page, q.limit, result.total); return ok( rows.map((r) => pickFields(r, q.fields)), @@ -127,3 +132,91 @@ export async function GET(req: NextRequest) { }, ); } + +/** Lots internes : compromis entre nombre d'allers-retours SQL et mémoire retenue. */ +const STREAM_BATCH = 1000; + +type Filtres = Parameters[0]; + +/** + * Diffuse **toutes** les lignes du fournisseur en une seule réponse JSON, sans + * jamais la construire entièrement en mémoire. + * + * La forme reste celle des autres réponses (`{ data, pagination, meta }`) : un + * client existant ne voit aucune différence, il reçoit simplement les octets au + * fil de l'eau. `pagination` et `meta` sont écrits en dernier, une fois le total + * connu — l'ordre des clés n'a aucune importance en JSON. + * + * Une erreur survenant après les premiers octets ne peut plus changer le code + * HTTP : on ferme alors le document avec `meta.erreur`, pour que l'appelant + * détecte une réponse tronquée au lieu de la croire complète. + */ +function streamAllRows( + req: NextRequest, + filtres: Filtres, + q: { enrich: string; fields?: string; fournisseur: string }, + freshness: { rowCount: number; computedAt: string | null }, + computedOnDemand: boolean, +): Response { + const encoder = new TextEncoder(); + + const stream = new ReadableStream({ + async start(controller) { + const write = (s: string) => controller.enqueue(encoder.encode(s)); + let envoyees = 0; + let total = 0; + let computedAt: string | null = null; + let erreur: string | null = null; + + write('{"data":['); + try { + for (let page = 1; ; page++) { + const lot = await queryGridRows({ ...filtres, page, limit: STREAM_BATCH }); + if (page === 1) { + total = lot.total; + computedAt = lot.computedAt; + } + if (lot.rows.length === 0) break; + + const enrichies = q.enrich === "1" ? await enrichRows(lot.rows) : lot.rows; + for (const r of enrichies) { + write((envoyees === 0 ? "" : ",") + JSON.stringify(pickFields(r, q.fields))); + envoyees++; + } + if (lot.rows.length < STREAM_BATCH) break; + } + } catch (e) { + erreur = e instanceof Error ? e.message : String(e); + console.error(`[api/v1/grid] diffusion interrompue pour ${q.fournisseur} après ${envoyees} lignes:`, erreur); + } + + const complet = erreur === null && envoyees === total; + write("]," + JSON.stringify({ + pagination: { page: 1, limit: envoyees, total, totalPages: 1, hasMore: !complet }, + meta: { + fournisseur: q.fournisseur, + computedAt, + snapshotComputedAt: freshness.computedAt, + snapshotRowCount: freshness.rowCount, + enrichi: q.enrich === "1", + computedOnDemand, + diffuse: true, + complet, + ...(erreur + ? { erreur: `Diffusion interrompue après ${envoyees} lignes sur ${total} : ${erreur}` } + : {}), + }, + }).slice(1)); + controller.close(); + }, + }); + + return new Response(stream, { + headers: { + "Content-Type": "application/json; charset=utf-8", + "Cache-Control": "no-store", + // Empêche un proxy de tamponner la réponse entière avant de la relayer. + "X-Accel-Buffering": "no", + }, + }); +} diff --git a/src/lib/pg-ff-client.ts b/src/lib/pg-ff-client.ts index 617a4f6..052bfeb 100644 --- a/src/lib/pg-ff-client.ts +++ b/src/lib/pg-ff-client.ts @@ -314,22 +314,29 @@ export async function pgGetGammesByCodeins(codeins: string[]): Promise String(c ?? "").trim()).filter(Boolean))]; if (unique.length === 0) return new Map(); - // DISTINCT ON (codein) → 1 gamme par article, saison la plus récente en premier. - const result = await pgNoParallel(sql` - SELECT DISTINCT ON (a.codein) - TRIM(a.codein::text) AS codein, - g.code AS gamme_code - FROM art_gamme_saison ags - JOIN articles a ON a.no_id = ags.artnoid - JOIN gammes g ON g.no_id = ags.idgamme - JOIN saisons s ON s.no_id = ags.idsaison - WHERE TRIM(a.codein::text) IN (${sql.join(unique.map(c => sql`${c}`), sql`, `)}) - ORDER BY a.codein, s.no_id DESC - `); - + // Découpage obligatoire : chaque codein devient un paramètre lié, or PostgreSQL + // en accepte 65535 au maximum. Un gros fournisseur dépasse 130 000 articles — + // la requête échouait alors purement et simplement. + const CHUNK = 2000; const map = new Map(); - for (const row of result.rows as unknown as { codein: string; gamme_code: string }[]) { - if (row.codein && row.gamme_code) map.set(row.codein, String(row.gamme_code).trim()); + + for (let i = 0; i < unique.length; i += CHUNK) { + const batch = unique.slice(i, i + CHUNK); + // DISTINCT ON (codein) → 1 gamme par article, saison la plus récente en premier. + const result = await pgNoParallel(sql` + SELECT DISTINCT ON (a.codein) + TRIM(a.codein::text) AS codein, + g.code AS gamme_code + FROM art_gamme_saison ags + JOIN articles a ON a.no_id = ags.artnoid + JOIN gammes g ON g.no_id = ags.idgamme + JOIN saisons s ON s.no_id = ags.idsaison + WHERE TRIM(a.codein::text) IN (${sql.join(batch.map(c => sql`${c}`), sql`, `)}) + ORDER BY a.codein, s.no_id DESC + `); + for (const row of result.rows as unknown as { codein: string; gamme_code: string }[]) { + if (row.codein && row.gamme_code) map.set(row.codein, String(row.gamme_code).trim()); + } } console.log(`[pg-ff] pgGetGammesByCodeins: ${map.size}/${unique.length} articles avec gamme`); return map; diff --git a/src/lib/qlik-network-cache.ts b/src/lib/qlik-network-cache.ts index 2c3109e..a27b78b 100644 --- a/src/lib/qlik-network-cache.ts +++ b/src/lib/qlik-network-cache.ts @@ -42,10 +42,19 @@ export async function getNetworkMetricsByCodeCentrale( const unique = [...new Set(codes.filter(Boolean))]; if (unique.length === 0) return out; - const rows = await db - .select() - .from(qlikNetworkMetrics) - .where(inArray(qlikNetworkMetrics.codeCentrale, unique)); + // Découpage obligatoire : `inArray` produit un paramètre lié par valeur, or le + // protocole PostgreSQL en accepte 65535 au maximum. Un gros fournisseur peut + // dépasser 130 000 codes — la requête échouait alors purement et simplement. + const CHUNK = 2000; + const rows: Array = []; + for (let i = 0; i < unique.length; i += CHUNK) { + const batch = unique.slice(i, i + CHUNK); + const part = await db + .select() + .from(qlikNetworkMetrics) + .where(inArray(qlikNetworkMetrics.codeCentrale, batch)); + rows.push(...part); + } for (const r of rows) { out.set(r.codeCentrale, {