import { Inject, Injectable, Logger } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { createHash } from 'crypto'; import { CACHE_PORT, CachePort } from '@domain/ports/out/cache.port'; import { TRADE_EMBEDDINGS, TradeEmbeddingPort, TradePassage, TradeRetrievalPort, } from '@domain/ports/out/trade-assistant.port'; import corpus from './knowledge/wiki-corpus.json'; /** * Recherche documentaire sur le wiki du site. * * Le corpus est celui de `dashboard.wikiPages` (voir * `scripts/setup/build-knowledge-corpus.js`) : l'assistant repond donc a partir * de la documentation que l'utilisateur peut ouvrir, et chaque reponse peut * citer la page correspondante. * * Deux caches evitent de refaire le meme calcul : * * 1. **L'index** est vectorise une seule fois. La cle Redis contient l'empreinte * du corpus : tant que le wiki ne change pas, aucun appel d'embedding n'est * refait, meme apres un redemarrage ou un deploiement. Il est aussi garde en * memoire du processus, donc une recherche ne touche Redis qu'au demarrage. * 2. **Les questions** sont vectorisees une fois par formulation, pour tous les * utilisateurs : reposer la meme question ne coute pas un second appel. * * Sans cle OpenAI — ou si l'API echoue — la recherche bascule sur un score * lexical. Une reponse moins bien documentee vaut mieux qu'une page en erreur. */ interface CorpusDocument { id: string; locale: string; topic: string; title: string; section: string; href: string; text: string; } const DOCUMENTS = (corpus as { documents: CorpusDocument[] }).documents; /** Empreinte du corpus : change des que la documentation du site change. */ const CORPUS_FINGERPRINT = createHash('sha256') .update(DOCUMENTS.map(d => `${d.id}${d.text}`).join('')) .digest('hex') .slice(0, 16); const QUESTION_TTL_SECONDS = 7 * 24 * 3600; const INDEX_TTL_SECONDS = 90 * 24 * 3600; const DEFAULT_LIMIT = 4; /** * Seuils de pertinence, mesures le 2026-09-06 sur le corpus reel. * * Une question metier (« quels sont les regimes douaniers ») remonte des * extraits entre 0,57 et 0,71. Une question sans rapport avec le wiki * (« combien de reservations ai-je ») plafonne a 0,40. La coupure a 0,45 separe * les deux avec de la marge des deux cotes ; a 0,28, l'assistant citait des * pages de calcul du fret sous une reponse sur le nombre de reservations. * * Le repli lexical ne vit pas sur la meme echelle — c'est une proportion de * mots couverts — et garde donc son propre seuil. */ const MIN_VECTOR_SCORE = 0.45; const MIN_LEXICAL_SCORE = 0.3; interface IndexedDocument { document: CorpusDocument; vector: Float32Array | null; /** Sac de mots, utilise par le repli lexical. */ terms: Set; } @Injectable() export class WikiRetriever implements TradeRetrievalPort { private readonly logger = new Logger(WikiRetriever.name); /** Index par langue, construit paresseusement puis garde en memoire. */ private readonly indexes = new Map>(); constructor( @Inject(TRADE_EMBEDDINGS) private readonly embeddings: TradeEmbeddingPort, @Inject(CACHE_PORT) private readonly cache: CachePort, private readonly config: ConfigService ) {} async search(question: string, language: string, limit = DEFAULT_LIMIT): Promise { const locale = DOCUMENTS.some(d => d.locale === language) ? language : 'fr'; const index = await this.index(locale); const vector = await this.questionVector(question); const terms = tokenize(question); const scored = index .map(entry => { const vectorised = Boolean(vector && entry.vector); return { entry, vectorised, score: vectorised ? dot(vector!, entry.vector!) : lexicalScore(entry.terms, terms), }; }) .filter(row => row.score >= (row.vectorised ? MIN_VECTOR_SCORE : MIN_LEXICAL_SCORE)) .sort((a, b) => b.score - a.score) .slice(0, limit); return scored.map(({ entry, score }) => ({ id: entry.document.id, title: entry.document.title, section: entry.document.section, href: entry.document.href, text: entry.document.text, score: Number(score.toFixed(4)), })); } /* ------------------------------------------------------------------------ */ /* Index */ /* ------------------------------------------------------------------------ */ private index(locale: string): Promise { const existing = this.indexes.get(locale); if (existing) return existing; // La promesse est memorisee avant d'etre resolue : deux requetes simultanees // au demarrage ne doivent pas vectoriser le corpus deux fois. const building = this.buildIndex(locale).catch(error => { this.logger.warn(`Falling back to lexical search: ${asMessage(error)}`); return documentsFor(locale).map(document => ({ document, vector: null, terms: tokenize(`${document.title} ${document.section} ${document.text}`), })); }); this.indexes.set(locale, building); return building; } private async buildIndex(locale: string): Promise { const documents = documentsFor(locale); const lexical = documents.map(document => ({ document, vector: null, terms: tokenize(`${document.title} ${document.section} ${document.text}`), })); if (!this.embeddings.isAvailable()) return lexical; const key = `trade:kb:${this.model()}:${CORPUS_FINGERPRINT}:${locale}`; const cached = await this.readCache(key); const vectors = cached?.length === documents.length ? cached : await this.vectorizeCorpus(key, documents, locale); return documents.map((document, i) => ({ ...lexical[i], document, vector: vectors[i] })); } private async vectorizeCorpus( key: string, documents: CorpusDocument[], locale: string ): Promise { this.logger.log(`Building the ${locale} knowledge index (${documents.length} passages)`); const raw = await this.embeddings.embed( documents.map(d => `${d.title} — ${d.section}\n${d.text}`) ); const vectors = raw.map(toFloat32); await this.writeCache(key, vectors, INDEX_TTL_SECONDS); return vectors; } /* ------------------------------------------------------------------------ */ /* Question */ /* ------------------------------------------------------------------------ */ private async questionVector(question: string): Promise { if (!this.embeddings.isAvailable()) return null; const key = `trade:q:${this.model()}:${hash(normalizeQuestion(question))}`; const cached = await this.readCache(key); if (cached?.length) return cached[0]; try { const [vector] = await this.embeddings.embed([question]); if (!vector) return null; const typed = toFloat32(vector); await this.writeCache(key, [typed], QUESTION_TTL_SECONDS); return typed; } catch (error) { this.logger.warn(`Question embedding failed: ${asMessage(error)}`); return null; } } /* ------------------------------------------------------------------------ */ /* Cache */ /* ------------------------------------------------------------------------ */ private model(): string { return this.config.get('OPENAI_EMBEDDING_MODEL', 'text-embedding-3-small'); } /** * Les vecteurs transitent en base64 de `Float32Array`. En JSON, l'index * francais pesait 2,6 Mo par langue pour la seule mise en forme des flottants. */ private async readCache(key: string): Promise { try { const packed = await this.cache.get(key); return packed?.map(unpack) ?? null; } catch (error) { this.logger.warn(`Knowledge cache unreadable: ${asMessage(error)}`); return null; } } private async writeCache(key: string, vectors: Float32Array[], ttl: number): Promise { try { await this.cache.set(key, vectors.map(pack), ttl); } catch (error) { // Un cache indisponible coute un recalcul, pas une panne. this.logger.warn(`Knowledge cache unwritable: ${asMessage(error)}`); } } } /* -------------------------------------------------------------------------- */ /* Fonctions pures */ /* -------------------------------------------------------------------------- */ const documentsFor = (locale: string) => DOCUMENTS.filter(d => d.locale === locale); const asMessage = (error: unknown) => (error instanceof Error ? error.message : String(error)); const hash = (value: string) => createHash('sha256').update(value).digest('hex').slice(0, 32); /** Les vecteurs etant normes, le produit scalaire est la similarite cosinus. */ export function dot(a: Float32Array, b: Float32Array): number { let sum = 0; for (let i = 0; i < a.length && i < b.length; i++) sum += a[i] * b[i]; return sum; } export const toFloat32 = (vector: number[]) => Float32Array.from(vector); export const pack = (vector: Float32Array) => Buffer.from(vector.buffer, vector.byteOffset, vector.byteLength).toString('base64'); export function unpack(packed: string): Float32Array { const buffer = Buffer.from(packed, 'base64'); // `Buffer` est alloue dans un pool partage : son `byteOffset` n'est pas // garanti multiple de 4, ce qu'exige une vue `Float32Array`. On copie. const copy = new ArrayBuffer(buffer.byteLength); new Uint8Array(copy).set(buffer); return new Float32Array(copy); } /** * Deux formulations identiques a la ponctuation et aux accents pres partagent * leur vecteur : « Quels documents ? » et « quels documents » ne sont pas deux * questions. */ export function normalizeQuestion(question: string): string { return question .normalize('NFD') .replace(/[\u0300-\u036f]/g, '') .toLowerCase() .replace(/[^a-z0-9]+/g, ' ') .trim(); } /** Mots de liaison : presents partout, ils ne discriminent rien. */ const STOP_WORDS = new Set( ( 'le la les un une des du de d au aux et ou a en dans pour par sur avec sans que qui quoi quel ' + 'quelle quels quelles est sont ete etre ai as avons avez ont mon ma mes votre vos notre nos ce ' + 'cet cette ces il elle ils elles je tu nous vous on se sa son ses plus moins tres bien ' + 'the a an of to in on for with and or is are be been what which how why when where do does my your' ).split(' ') ); export function tokenize(text: string): Set { return new Set( normalizeQuestion(text) .split(' ') .filter(word => word.length > 2 && !STOP_WORDS.has(word)) ); } /** * Repli lexical : proportion des mots de la question couverts par l'extrait. * Le score vise la meme echelle que la similarite cosinus pour que `MIN_SCORE` * garde un sens dans les deux modes. */ export function lexicalScore(documentTerms: Set, questionTerms: Set): number { if (!questionTerms.size) return 0; let matched = 0; for (const term of questionTerms) if (documentTerms.has(term)) matched++; return matched / questionTerms.size; }