diff --git a/README.md b/README.md index 50e328634..222f11740 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,7 @@ All contributing instructions are in [CONTRIBUTING](CONTRIBUTING.md). - [authors-tools](./services/authors-tools) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-authors-tools.svg)](https://hub.docker.com/r/cnrsinist/ws-authors-tools/) - [base-line](./services/base-line) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-base-line.svg)](https://hub.docker.com/r/cnrsinist/ws-base-line/) - [base-line-python](./services/base-line-python) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-base-line-python.svg)](https://hub.docker.com/r/cnrsinist/ws-base-line-python/) +- [baseline-dvc](./services/baseline-dvc) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-baseline-dvc.svg)](https://hub.docker.com/r/cnrsinist/ws-baseline-dvc/) - [biblio-ref](./services/biblio-ref) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-biblio-ref.svg)](https://hub.docker.com/r/cnrsinist/ws-biblio-ref/) - [biblio-tools](./services/biblio-tools) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-biblio-tools.svg)](https://hub.docker.com/r/cnrsinist/ws-biblio-tools/) - [chem-ner](./services/chem-ner) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-chem-ner.svg)](https://hub.docker.com/r/cnrsinist/ws-chem-ner/) @@ -51,6 +52,7 @@ All contributing instructions are in [CONTRIBUTING](CONTRIBUTING.md). - [hal-classifier](./services/hal-classifier) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-hal-classifier.svg)](https://hub.docker.com/r/cnrsinist/ws-hal-classifier/) - [hiddentext-detect](./services/hiddentext-detect) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-hiddentext-detect.svg)](https://hub.docker.com/r/cnrsinist/ws-hiddentext-detect/) - [irc3-species](./services/irc3-species) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-irc3-species.svg)](https://hub.docker.com/r/cnrsinist/ws-irc3-species/) +- [loterre-annotate](./services/loterre-annotate) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-loterre-annotate.svg)](https://hub.docker.com/r/cnrsinist/ws-loterre-annotate/) - [loterre-resolvers](./services/loterre-resolvers) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-loterre-resolvers.svg)](https://hub.docker.com/r/cnrsinist/ws-loterre-resolvers/) - [ner-tagger](./services/ner-tagger) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-ner-tagger.svg)](https://hub.docker.com/r/cnrsinist/ws-ner-tagger/) - [nlp-tools2](./services/nlp-tools2) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-nlp-tools2.svg)](https://hub.docker.com/r/cnrsinist/ws-nlp-tools2/) @@ -58,6 +60,7 @@ All contributing instructions are in [CONTRIBUTING](CONTRIBUTING.md). - [openalex-doctype](./services/openalex-doctype) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-openalex-doctype.svg)](https://hub.docker.com/r/cnrsinist/ws-openalex-doctype/) - [pdf-text](./services/pdf-text) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-pdf-text.svg)](https://hub.docker.com/r/cnrsinist/ws-pdf-text/) - [person-ner](./services/person-ner) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-person-ner.svg)](https://hub.docker.com/r/cnrsinist/ws-person-ner/) +- [rag-tools](./services/rag-tools) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-rag-tools.svg)](https://hub.docker.com/r/cnrsinist/ws-rag-tools/) - [sciencemetrix-classification](./services/sciencemetrix-classification) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-sciencemetrix-classification.svg)](https://hub.docker.com/r/cnrsinist/ws-sciencemetrix-classification/) - [software-extract](./services/software-extract) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-software-extract.svg)](https://hub.docker.com/r/cnrsinist/ws-software-extract/) - [tabbr-mine](./services/tabbr-mine) [![Docker Pulls](https://img.shields.io/docker/pulls/cnrsinist/ws-tabbr-mine.svg)](https://hub.docker.com/r/cnrsinist/ws-tabbr-mine/) diff --git a/package.json b/package.json index e810a78a0..0fd2625e6 100644 --- a/package.json +++ b/package.json @@ -79,6 +79,7 @@ "services/openalex-doctype", "services/pdf-text", "services/person-ner", + "services/rag-tools", "services/sciencemetrix-classification", "services/software-extract", "services/tabbr-mine", @@ -100,4 +101,4 @@ "devDependencies": { "@types/node": "24.13.3" } -} +} \ No newline at end of file diff --git a/services/rag-tools/.dockerignore b/services/rag-tools/.dockerignore new file mode 100644 index 000000000..e280cbb78 --- /dev/null +++ b/services/rag-tools/.dockerignore @@ -0,0 +1,7 @@ +# Ignore all files by default +* + +# White list only the required files +!config.json +!v1 +!swagger.json diff --git a/services/rag-tools/Dockerfile b/services/rag-tools/Dockerfile new file mode 100644 index 000000000..b73d17a9e --- /dev/null +++ b/services/rag-tools/Dockerfile @@ -0,0 +1,28 @@ +# syntax=docker/dockerfile:1.2 + +FROM cnrsinist/ezs-python-server:py3.9-no24-1.0.14 + +WORKDIR /app + +# Installe torch en version CPU uniquement : le wheel par défaut de pip +# embarque tout le support CUDA (bibliothèques NVIDIA), qui représente à +# lui seul plusieurs Go, totalement inutile sans GPU. +RUN pip install --no-cache-dir torch --extra-index-url https://download.pytorch.org/whl/cpu +# sentence-transformers réutilise le torch CPU déjà installé ci-dessus +# (huggingface_hub est déjà une dépendance transitive de +# sentence-transformers/transformers, pas besoin de l'installer à part). +RUN pip install --no-cache-dir sentence-transformers + +# Télécharge le modèle, puis nettoie le cache de téléchargement HuggingFace +# dans la MÊME couche : le modèle existe sinon en double (cache + /models), +# soit environ 2,2 Go gaspillés inutilement. +RUN python -c "\ +from sentence_transformers import SentenceTransformer; \ +model = SentenceTransformer('BAAI/bge-m3'); \ +model.save('/models/BAAI/bge-m3')" \ + && rm -rf /root/.cache/huggingface + +WORKDIR /app/public + +COPY --chown=daemon:daemon . /app/public/ +COPY --chown=daemon:daemon ./config.json /app/config.json \ No newline at end of file diff --git a/services/rag-tools/README.md b/services/rag-tools/README.md new file mode 100644 index 000000000..b0105792d --- /dev/null +++ b/services/rag-tools/README.md @@ -0,0 +1,5 @@ +# ws-rag-tools@0.4.0 + +Ws to generate embeddings and use those for rag + +Ws to generate embeddings and use those for rag diff --git a/services/rag-tools/config.json b/services/rag-tools/config.json new file mode 100644 index 000000000..7dd647531 --- /dev/null +++ b/services/rag-tools/config.json @@ -0,0 +1,15 @@ +{ + "environnement": { + "EZS_TITLE": "Ws to generate embeddings and use those for rag", + "EZS_DESCRIPTION": "Ws to generate embeddings and use those for rag", + "EZS_METRICS": true, + "EZS_CONCURRENCY": 1, + "EZS_CONTINUE_DELAY": 60, + "EZS_NSHARDS": 32, + "EZS_CACHE": true, + "EZS_VERBOSE": false, + "NODE_OPTIONS": "--max_old_space_size=1024", + "NODE_ENV": "production", + "ILAAS_API_KEY": "real_api_key" + } +} \ No newline at end of file diff --git a/services/rag-tools/examples.http b/services/rag-tools/examples.http new file mode 100644 index 000000000..e8c09a1cd --- /dev/null +++ b/services/rag-tools/examples.http @@ -0,0 +1,17 @@ +# These examples can be used directly in VSCode, using HTTPYac extension (anweber.vscode-httpyac) +# They are important, because used to generate the tests.hurl file. + +# Décommenter/commenter les lignes voulues pour tester localement +@host=http://localhost:31976 +# @host=https://rag-tools.services.istex.fr + +### +# @name v1routeInCamelCase +# Description de la route +POST {{host}}/v1/route/in/camel/case?indent=true HTTP/1.1 +Content-Type: application/json + +[ + { "value": "une valeur typique" }, + { "value": "en json" } +] diff --git a/services/rag-tools/package.json b/services/rag-tools/package.json new file mode 100644 index 000000000..c91f9f9b0 --- /dev/null +++ b/services/rag-tools/package.json @@ -0,0 +1,37 @@ +{ + "private": true, + "name": "ws-rag-tools", + "version": "0.4.0", + "description": "Ws to generate embeddings and use those for rag", + "repository": { + "type": "git", + "url": "git+https://github.com/Inist-CNRS/web-services.git" + }, + "keywords": [ + "ezmaster" + ], + "author": "Anki Lucas ", + "license": "MIT", + "bugs": { + "url": "https://github.com/Inist-CNRS/web-services/issues" + }, + "homepage": "https://github.com/Inist-CNRS/web-services/#readme", + "scripts": { + "version:insert:readme": "sed -i \"s#\\(${npm_package_name}.\\)\\([\\.a-z0-9]\\+\\)#\\1${npm_package_version}#g\" README.md && git add README.md", + "version:insert:swagger": "sed -i \"s/\\\"version\\\": \\\"[0-9]\\+.[0-9]\\+.[0-9]\\+\\\"/\\\"version\\\": \\\"${npm_package_version}\\\"/g\" swagger.json && git add swagger.json", + "version:insert": "npm run version:insert:readme && npm run version:insert:swagger", + "version:commit": "git commit -a -m \"release ${npm_package_name}@${npm_package_version}\"", + "version:tag": "git tag \"${npm_package_name}@${npm_package_version}\" -m \"${npm_package_name}@${npm_package_version}\"", + "version:push": "git push && git push origin \"${npm_package_name}@${npm_package_version}\"", + "version": "npm run version:insert && npm run version:commit && npm run version:tag", + "postversion": "npm run version:push", + "build:check": "DOCKER_BUILDKIT=1 docker build --check -t cnrsinist/${npm_package_name}:latest . && docker run --rm -i hadolint/hadolint hadolint - < Dockerfile ", + "build:dev": "docker build -t cnrsinist/${npm_package_name}:latest .", + "start:dev": ". ./.env 2> /dev/null; npm run build:dev && docker run -e ILAAS_API_KEY --name dev --rm --detach -p 31976:31976 cnrsinist/${npm_package_name}:latest", + "stop:dev": "docker stop dev", + "build": "docker build -t cnrsinist/${npm_package_name}:${npm_package_version} .", + "start": "docker run --rm -p 31976:31976 cnrsinist/${npm_package_name}:${npm_package_version}", + "publish": "docker push cnrsinist/${npm_package_name}:${npm_package_version}" + }, + "avoid-testing": false +} diff --git a/services/rag-tools/swagger.json b/services/rag-tools/swagger.json new file mode 100644 index 000000000..148e2b83d --- /dev/null +++ b/services/rag-tools/swagger.json @@ -0,0 +1,33 @@ +{ + "openapi": "3.0.0", + "info": { + "title": "rag-tools - Ws to generate embeddings and use those for rag", + "description": "Ws to generate embeddings and use those for rag", + "version": "0.4.0", + "termsOfService": "https://services.istex.fr/", + "contact": { + "name": "Inist-CNRS", + "url": "https://www.inist.fr/nous-contacter/" + } + }, + "servers": [ + { + "x-comment": "Will be automatically completed by the ezs server." + }, + { + "url": "http://vptdmservices.intra.inist.fr:49225/", + "description": "Latest version for production", + "#DISABLED#x-profil": "Standard" + } + ], + "tags": [ + { + "name": "rag-tools", + "description": "Ws to generate embeddings and use those for rag", + "externalDocs": { + "description": "Plus de documentation", + "url": "https://github.com/inist-cnrs/web-services/tree/main/services/rag-tools" + } + } + ] +} \ No newline at end of file diff --git a/services/rag-tools/tests.hurl b/services/rag-tools/tests.hurl new file mode 100644 index 000000000..4c3447595 --- /dev/null +++ b/services/rag-tools/tests.hurl @@ -0,0 +1,159 @@ +# ============================== +# Embedding +# ============================== + +POST {{host}}/v1/embedding +content-type: application/json +[ + { + "id": "embed_001", + "value": "Quelles sont les principales caractéristiques de Jupiter ?" + } +] + +HTTP 200 +Content-Type: application/json +[Asserts] +jsonpath "$[0].id" == "embed_001" +jsonpath "$[0].value.vector" count > 0 + + +# ============================== +# Embedding API +# ============================== + +#POST {{host}}/v1/embedding-api +#content-type: application/json +#[ +# { +# "id": "embed_001", +# "value": "Quelles sont les principales caractéristiques de Jupiter ?" +# } +#] +# +#HTTP 200 +#Content-Type: application/json +#[Asserts] +#jsonpath "$[0].id" == "embed_001" +#jsonpath "$[0].value.vector" count > 0 + + +# ============================== +# Reformulate +# ============================== + +POST {{host}}/v1/reformulate +content-type: application/json +[ + { + "id": "reformulate_001", + "value": { + "question": "Et Saturne, elle a les mêmes caractéristiques ?", + "historique": [ + { + "role": "user", + "content": "Quelles sont les principales caractéristiques de Jupiter ?" + }, + { + "role": "assistant", + "content": "Jupiter est la plus grande planète du système solaire. C'est une géante gazeuse composée principalement d'hydrogène et d'hélium, connue pour sa Grande Tache Rouge [Document 1]." + } + ] + } + } +] + +HTTP 200 +Content-Type: application/json +[Asserts] +jsonpath "$[0].id" == "reformulate_001" +jsonpath "$[0].value" contains "Saturne" +jsonpath "$[0].value" not contains "Et Saturne, elle a les mêmes caractéristiques ?" + + +# ============================== +# Rerank +# ============================== + +POST {{host}}/v1/rerank +content-type: application/json +[ + { + "id": "rerank_001", + "value": { + "question": "Quels articles portent sur la résistance aux antibiotiques chez Escherichia coli ?", + "documents": [ + { + "text": "Mercure est la planète la plus proche du Soleil. Elle effectue une révolution autour du Soleil en environ 88 jours.", + "metadata": { + "titre": "Le système solaire interne", + "date_publication": "2019" + } + }, + { + "text": "Objectives To characterize the naturally occurring expanded-spectrum β-lactamase from an Escherichia coli clinical isolate...", + "metadata": { + "titre": "Extension of the hydrolysis spectrum of AmpC β-lactamase of Escherichia coli due to amino acid insertion in the H-10 helix", + "auteurs": ["Hedi Mammeri", "Laurent Poirel", "Patrice Nordmann"], + "date_publication": "2007", + "journal": "Journal of Antimicrobial Chemotherapy", + "doi": "10.1093/jac/dkm227" + } + }, + { + "text": "Mars est surnommée la planète rouge à cause de l'oxyde de fer présent à sa surface.", + "metadata": { + "titre": "Exploration martienne", + "date_publication": "2021" + } + } + ], + "top_n": 1 + } + } +] + +HTTP 200 +Content-Type: application/json +[Asserts] +jsonpath "$[0].id" == "rerank_001" +jsonpath "$[0].value.documents" count == 1 +jsonpath "$[0].value.documents[0].metadata.doi" == "10.1093/jac/dkm227" + + +# ============================== +# RAG +# ============================== + +POST {{host}}/v1/rag +content-type: application/json +[ + { + "id": "rag_001", + "value": { + "question": "Qui sont les auteurs de l'article sur l'AmpC β-lactamase et dans quelle revue a-t-il été publié ?", + "documents": [ + { + "id": "doi:10.1093/jac/dkm227#abstract", + "text": "Objectives To characterize the naturally occurring expanded-spectrum β-lactamase from an Escherichia coli clinical isolate...", + "metadata": { + "titre": "Extension of the hydrolysis spectrum of AmpC β-lactamase of Escherichia coli due to amino acid insertion in the H-10 helix", + "auteurs": ["Hedi Mammeri", "Laurent Poirel", "Patrice Nordmann"], + "date_publication": "2007", + "journal": "Journal of Antimicrobial Chemotherapy", + "doi": "10.1093/jac/dkm227" + } + } + ], + "historique": [] + } + } +] + +HTTP 200 +Content-Type: application/json +[Asserts] +jsonpath "$[0].id" == "rag_001" +jsonpath "$[0].value.answer" contains "Mammeri" +jsonpath "$[0].value.answer" contains "Journal of Antimicrobial Chemotherapy" +jsonpath "$[0].value.answer" contains "[doi:10.1093/jac/dkm227#abstract]" diff --git a/services/rag-tools/v1/classification-session.ini b/services/rag-tools/v1/classification-session.ini new file mode 100644 index 000000000..4c8ac3fa9 --- /dev/null +++ b/services/rag-tools/v1/classification-session.ini @@ -0,0 +1,34 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = Classifie la question dans le bon domaine pour la suite (version session) +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Vectorise des documents +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/classification-session.py + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/classification-session.py b/services/rag-tools/v1/classification-session.py new file mode 100755 index 000000000..510c5357b --- /dev/null +++ b/services/rag-tools/v1/classification-session.py @@ -0,0 +1,298 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os +import uuid + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID = "classification_template" + +NO_HISTORY_TEXT = "Aucun historique de conversation disponible." + +TMP_DIR = "/tmp" + +# Types de requête reconnus. DEFAULT_TYPE est utilisé en repli si la +# réponse du LLM ne correspond à aucun type valide (réponse mal formée, +# hors-liste, erreur d'appel...). Pour ajouter un nouveau type à l'avenir : +# 1. l'ajouter à VALID_TYPES, 2. mettre à jour le prompt classification_template +# pour lui décrire ce nouveau type. +VALID_TYPES = {"rag", "definition"} +DEFAULT_TYPE = "rag" + + +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == PROMPT_ID: + return prompt["content"] + + raise ValueError(f"Prompt {PROMPT_ID} not found") + + +PROMPT_TEMPLATE = load_prompt() + + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + + +# ============================== +# Gestion de la session (token -> historique persisté sur disque) +# ============================== + +def history_path(token: str) -> str: + return os.path.join(TMP_DIR, token, f"{token}.json") + + +def generate_token() -> str: + return uuid.uuid4().hex + + +def load_historique(token: str) -> list: + """Charge l'historique associé à un token. Renvoie une liste vide si + le fichier n'existe pas encore (token valide mais session sans + historique enregistré, ne devrait normalement pas arriver hors + scénario de token fourni par erreur).""" + path = history_path(token) + if not os.path.exists(path): + print_log(f"Aucun historique trouvé pour le token {token}, historique vide utilisé") + return [] + + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def resolve_session(token_value): + """ + `token_value` est ce qui se trouve dans value["historique"] en entrée : + - absent / None / chaîne vide -> premier appel, on génère un nouveau + token et l'historique de départ est vide. + - une chaîne non vide -> token existant, on charge l'historique + correspondant depuis le disque. + + Renvoie (token, historique). + """ + if not token_value: + new_token = generate_token() + print_log(f"Aucun token fourni, nouvelle session créée : {new_token}") + return new_token, [] + + return token_value, load_historique(token_value) + + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 20 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"].strip() + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Construction de l'historique +# ============================== + +def build_history_text(historique): + """ + Construit le texte de l'historique de la conversation. + + `historique` est attendu comme une liste de tours de dialogue, + chaque tour étant un dict avec les clés "role" et "content". + + Si aucun historique n'est fourni (liste vide), un texte par défaut + est renvoyé. + """ + if not historique: + return NO_HISTORY_TEXT + + lines = [] + + for turn in historique: + role = turn.get("role", "inconnu") + content = turn.get("content", "") + lines.append(f"{role} : {content}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt de classification +# ============================== + +def build_prompt(question, historique): + historique_text = build_history_text(historique) + + return PROMPT_TEMPLATE.format( + question=question, + historique=historique_text + ) + + +# ============================== +# Parsing de la réponse du LLM +# ============================== + +def parse_classification(raw_response: str) -> str: + """ + Normalise et valide la réponse du LLM. Renvoie DEFAULT_TYPE si la + réponse est vide, mal formée, ou ne correspond à aucun type reconnu + (plutôt que de bloquer le pipeline sur un type invalide). + """ + if raw_response == "Error" or not raw_response: + print_log("Classification invalide, fallback sur le type par défaut") + return DEFAULT_TYPE + + cleaned = raw_response.strip().lower() + # Tolère une éventuelle ponctuation ou guillemets résiduels + cleaned = cleaned.strip(" .\"'`") + + if cleaned not in VALID_TYPES: + print_log( + f"Type '{raw_response}' non reconnu, fallback sur '{DEFAULT_TYPE}'" + ) + return DEFAULT_TYPE + + return cleaned + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + token, historique = resolve_session(item["value"].get("historique")) + + prompt = build_prompt( + question, + historique + ) + + raw_response = call_llm(prompt) + rag_type = parse_classification(raw_response) + + output = { + "id": item["id"], + "value": { + "type": rag_type, + "token": token + } + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch) diff --git a/services/rag-tools/v1/classification.ini b/services/rag-tools/v1/classification.ini new file mode 100644 index 000000000..e9404dd1f --- /dev/null +++ b/services/rag-tools/v1/classification.ini @@ -0,0 +1,34 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = Classifie la question dans le bon domaine pour la suite +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Vectorise des documents +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/classification.py + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/classification.py b/services/rag-tools/v1/classification.py new file mode 100755 index 000000000..12b52513e --- /dev/null +++ b/services/rag-tools/v1/classification.py @@ -0,0 +1,249 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID = "classification_template" + +NO_HISTORY_TEXT = "Aucun historique de conversation disponible." + +# Types de requête reconnus. DEFAULT_TYPE est utilisé en repli si la +# réponse du LLM ne correspond à aucun type valide (réponse mal formée, +# hors-liste, erreur d'appel...). Pour ajouter un nouveau type à l'avenir : +# 1. l'ajouter à VALID_TYPES, 2. mettre à jour le prompt classification_template +# pour lui décrire ce nouveau type. +VALID_TYPES = {"rag", "definition"} +DEFAULT_TYPE = "rag" + + +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == PROMPT_ID: + return prompt["content"] + + raise ValueError(f"Prompt {PROMPT_ID} not found") + + +PROMPT_TEMPLATE = load_prompt() + + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + # si gemma-4-31b émet du reasoning_content, les 20 tokens peuvent être consommés avant la réponse → fallback permanent sur rag. À vérifier + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 20 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"].strip() + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Construction de l'historique +# ============================== + +def build_history_text(historique): + """ + Construit le texte de l'historique de la conversation. + + `historique` est attendu comme une liste de tours de dialogue, + chaque tour étant un dict avec les clés "role" et "content". + + Si aucun historique n'est fourni (absent, None ou liste vide), + un texte par défaut est renvoyé. + """ + if not historique: + return NO_HISTORY_TEXT + + lines = [] + + for turn in historique: + role = turn.get("role", "inconnu") + content = turn.get("content", "") + lines.append(f"{role} : {content}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt de classification +# ============================== + +def build_prompt(question, historique): + historique_text = build_history_text(historique) + + return PROMPT_TEMPLATE.format( + question=question, + historique=historique_text + ) + + +# ============================== +# Parsing de la réponse du LLM +# ============================== + +def parse_classification(raw_response: str) -> str: + """ + Normalise et valide la réponse du LLM. Renvoie DEFAULT_TYPE si la + réponse est vide, mal formée, ou ne correspond à aucun type reconnu + (plutôt que de bloquer le pipeline sur un type invalide). + """ + if raw_response == "Error" or not raw_response: + print_log("Classification invalide, fallback sur le type par défaut") + return DEFAULT_TYPE + + cleaned = raw_response.strip().lower() + # Tolère une éventuelle ponctuation ou guillemets résiduels + cleaned = cleaned.strip(" .\"'`") + + if cleaned not in VALID_TYPES: + print_log( + f"Type '{raw_response}' non reconnu, fallback sur '{DEFAULT_TYPE}'" + ) + return DEFAULT_TYPE + + return cleaned + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + historique = item["value"].get("historique", []) + + prompt = build_prompt( + question, + historique + ) + + raw_response = call_llm(prompt) + rag_type = parse_classification(raw_response) + + output = { + "id": item["id"], + "value": rag_type + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch) \ No newline at end of file diff --git a/services/rag-tools/v1/embedding-api.ini b/services/rag-tools/v1/embedding-api.ini new file mode 100644 index 000000000..ed411f186 --- /dev/null +++ b/services/rag-tools/v1/embedding-api.ini @@ -0,0 +1,34 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = Vectorise des documents via Api +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Vectorise des documents +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/embedding-api.py + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/embedding-api.py b/services/rag-tools/v1/embedding-api.py new file mode 100755 index 000000000..875634d7e --- /dev/null +++ b/services/rag-tools/v1/embedding-api.py @@ -0,0 +1,205 @@ +#!/usr/bin/env python3 +""" +Équivalent de embedding.py, mais appelle l'API d'embedding iLaaS au lieu +de charger le modèle localement. + +Entrée/sortie strictement identiques à vectorize.py (même format NDJSON +sur stdin/stdout), pour rester interchangeable avec le reste du pipeline +(insert.py, orchestrateur RAG). +""" + +import os +import sys +import json +import time +import requests +from concurrent.futures import ThreadPoolExecutor + +# ============================== +# Configuration +# ============================== + +API_KEY = os.getenv("ILAAS_API_KEY") +API_URL = "https://rag-api.ilaas.fr/v1/embeddings" +MODEL_NAME = "bge-m3" + +MAX_RETRIES = 4 +RETRY_DELAY = 2 + +# Nombre de documents traités par "lot" (pour le logging et le regroupement +# des appels en parallèle) — contrairement à vectorize.py, il ne s'agit pas +# d'un batch envoyé en un seul appel API, mais de N appels individuels +# exécutés en parallèle au sein du lot. +BATCH_SIZE = 32 + +# Nombre d'appels API simultanés au sein d'un même lot. L'API étant +# distante (I/O-bound), la parallélisation via threads est efficace ici, +# contrairement à un modèle local où le calcul est borné par le CPU/GPU. +MAX_WORKERS = 8 + + +def log(msg): + """Écrit un message de progression sur stderr.""" + sys.stderr.write(msg + "\n") + sys.stderr.flush() + + +def extract_text_and_metadata(value): + """ + `value` peut être : + - une simple chaîne (cas d'une question à vectoriser, ex: orchestrateur RAG) + -> pas de métadonnées associées. + - un objet {"text": ..., "metadata": {...}} (cas d'un chunk de document + à vectoriser pour l'ingestion, ex: script d'insertion Mongo) + -> métadonnées à propager telles quelles. + """ + if isinstance(value, dict): + return value["text"], value.get("metadata") + + return value, None + + +# ============================== +# Appel à l'API d'embedding +# ============================== + +def call_embedding_api(text: str): + """ + Appelle l'API iLaaS pour un texte donné, avec retries en cas d'échec. + Renvoie le vecteur d'embedding (list[float]), ou None si tous les + essais ont échoué. + """ + headers = { + "Content-Type": "application/json", + "Authorization": f"Bearer {API_KEY}", + } + payload = { + "model": MODEL_NAME, + "input": text, + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + API_URL, + headers=headers, + json=payload, + timeout=60, + ) + response.raise_for_status() + result = response.json() + + embedding = result.get("data", [{}])[0].get("embedding") + if embedding is None: + raise ValueError(f"Réponse API sans embedding : {result}") + + return embedding + + except Exception as e: + log( + f"[WARN] Échec de l'appel à l'API d'embedding " + f"(tentative {attempt}/{MAX_RETRIES}) : {e}" + ) + + if attempt == MAX_RETRIES: + return None + + time.sleep(RETRY_DELAY * attempt) + + +# ============================== +# Traitement d'un document +# ============================== + +def process_item(item: dict): + """ + Vectorise un document via l'API et construit sa sortie au même format + que vectorize.py. Renvoie None si l'embedding a échoué (document ignoré, + avec un warning déjà loggé par call_embedding_api). + """ + text, metadata = extract_text_and_metadata(item["value"]) + + embedding = call_embedding_api(text) + if embedding is None: + log(f"[WARN] Document {item['id']} ignoré (embedding échoué après {MAX_RETRIES} tentatives).") + return None + + output_value = { + "vector": embedding, + "text": text, + } + + # On ne rajoute "metadata" que si elle existe, pour ne pas polluer + # la sortie du cas "question simple" (sans métadonnées). + if metadata is not None: + output_value["metadata"] = metadata + + return { + "id": item["id"], + "value": output_value, + } + + +# ============================== +# Traitement par lot (parallélisé) +# ============================== + +def process_batch(batch, batch_num=None, total_batches=None): + if not batch: + return + + _start = time.time() + + with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: + outputs = list(executor.map(process_item, batch)) + + _duration = time.time() - _start + + if batch_num is not None and total_batches is not None: + batch_label = f"Batch {batch_num}/{total_batches}" + else: + batch_label = f"Batch de {len(batch)} document(s)" + + log(f"[INFO] {batch_label} encodé en {_duration:.2f}s") + + for output in outputs: + if output is None: + continue + sys.stdout.write(json.dumps(output, ensure_ascii=False)) + sys.stdout.write("\n") + + sys.stdout.flush() + + +# ============================== +# Main +# ============================== + +if not API_KEY: + log("[ERREUR] Variable d'environnement ILAAS_API_KEY manquante.") + sys.exit(1) + +log("[INFO] Lecture des données...") +all_data = [] +for line in sys.stdin: + line = line.strip() + if not line: + continue + all_data.append(json.loads(line)) + +total = len(all_data) +total_batches = (total + BATCH_SIZE - 1) // BATCH_SIZE +log(f"[INFO] {total} documents à vectoriser ({total_batches} batch(s))") + +if total == 0: + log("[WARN] Aucun document reçu.") + sys.exit(0) + +processed = 0 + +for batch_num, i in enumerate(range(0, total, BATCH_SIZE), start=1): + batch = all_data[i:i + BATCH_SIZE] + process_batch(batch, batch_num=batch_num, total_batches=total_batches) + processed += len(batch) + +log(f"[INFO] Vectorisation terminée : {processed}/{total} documents") \ No newline at end of file diff --git a/services/rag-tools/v1/embedding.ini b/services/rag-tools/v1/embedding.ini new file mode 100644 index 000000000..ceb3ddb7e --- /dev/null +++ b/services/rag-tools/v1/embedding.ini @@ -0,0 +1,34 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = Vectorise des documents +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Vectorise des documents +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/embedding.py + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/embedding.py b/services/rag-tools/v1/embedding.py new file mode 100755 index 000000000..483e7411d --- /dev/null +++ b/services/rag-tools/v1/embedding.py @@ -0,0 +1,176 @@ +#!/usr/bin/env python3 +import os + +# ────────────────────────────────────────────────────────────────────────── +# Optimisations CPU — à faire AVANT tout import de torch/sentence_transformers, +# car OMP_NUM_THREADS/MKL_NUM_THREADS ne sont lus qu'à l'initialisation des +# bibliothèques de calcul (OpenMP/MKL), pas modifiables après coup. +# ────────────────────────────────────────────────────────────────────────── +_NUM_THREADS = str(os.cpu_count() or 1) +os.environ.setdefault("OMP_NUM_THREADS", _NUM_THREADS) +os.environ.setdefault("MKL_NUM_THREADS", _NUM_THREADS) +# Évite les ralentissements liés aux nombres flottants dénormalisés +# (fréquents en sortie de couches d'activation), gain "gratuit" sur CPU. +os.environ.setdefault("KMP_AFFINITY", "granularity=fine,compact,1,0") + +import sys +import json +import time +import logging +import torch +from sentence_transformers import SentenceTransformer + +# Supprime le warning avant tout chargement +logging.getLogger("transformers.tokenization_utils_base").setLevel(logging.ERROR) + +MODEL_PATH = "/models/BAAI/bge-m3" + +# Batch plus large que la valeur par défaut initiale : sur CPU, un batch +# trop petit sous-exploite la vectorisation SIMD/BLAS ; un batch plus +# grand amortit mieux les coûts fixes. À ajuster selon la RAM disponible +# (chaque doc supplémentaire dans le batch consomme de la mémoire). +BATCH_SIZE = 64 + +# bge-m3 supporte nativement jusqu'à 8192 tokens (usage "contexte long"), +# mais nos chunks font ~400 mots (~500-600 tokens). Sans plafond explicite, +# le modèle peut padder/traiter des séquences bien plus longues que +# nécessaire, ce qui coûte cher en calcul (l'attention croît de façon +# quadratique avec la longueur de séquence). +MAX_SEQ_LENGTH = 512 + + +def log(msg): + """Écrit un message de progression sur stderr.""" + sys.stderr.write(msg + "\n") + sys.stderr.flush() + + +# ────────────────────────────────────────────────────────────────────────── +# Configuration explicite de PyTorch pour le calcul CPU pur. +# +# - set_num_threads : parallélisme "intra-op" (une seule opération, ex: +# une multiplication de matrices, répartie sur plusieurs cœurs). C'est +# le levier principal pour accélérer un encode() séquentiel. +# - set_num_interop_threads : parallélisme "inter-op" (plusieurs opérations +# indépendantes en parallèle). Inutile ici puisqu'on fait des appels +# encode() séquentiels un par un ; le laisser à 1 réduit l'overhead de +# coordination entre threads plutôt que de l'aider. Doit être fixé AVANT +# toute opération torch, sous peine d'erreur si déjà initialisé. +# - set_flush_denormal : évite le ralentissement CPU causé par les nombres +# flottants "dénormalisés" (très proches de zéro), qui peuvent être +# traités beaucoup plus lentement par le FPU sur certains processeurs. +# ────────────────────────────────────────────────────────────────────────── +_num_threads = os.cpu_count() or 1 +torch.set_num_threads(_num_threads) +try: + torch.set_num_interop_threads(1) +except RuntimeError: + # Déjà initialisé ailleurs (ex: import précédent) : sans impact bloquant, + # on continue avec la valeur par défaut. + pass +torch.set_flush_denormal(True) + +log("[INFO] Chargement du modèle...") +_load_start = time.time() + +model = SentenceTransformer(MODEL_PATH, device="cpu") +model.max_seq_length = MAX_SEQ_LENGTH + +_load_duration = time.time() - _load_start +log(f"[INFO] Modèle chargé en {_load_duration:.1f}s sur device : {model.device}") +log(f"[INFO] max_seq_length fixé à {MAX_SEQ_LENGTH} tokens") +log(f"[INFO] Threads CPU (intra-op) : {_num_threads} | BATCH_SIZE : {BATCH_SIZE}") + + +def extract_text_and_metadata(value): + """ + `value` peut être : + - une simple chaîne (cas d'une question à vectoriser, ex: orchestrateur RAG) + -> pas de métadonnées associées. + - un objet {"text": ..., "metadata": {...}} (cas d'un chunk de document + à vectoriser pour l'ingestion, ex: script d'insertion Mongo) + -> métadonnées à propager telles quelles. + """ + if isinstance(value, dict): + return value["text"], value.get("metadata") + + return value, None + + +def process_batch(batch, batch_num=None, total_batches=None): + """Vectorise un batch de documents et écrit le résultat sur stdout.""" + if not batch: + return + + texts = [] + metadatas = [] + + for item in batch: + text, metadata = extract_text_and_metadata(item["value"]) + texts.append(text) + metadatas.append(metadata) + + _encode_start = time.time() + with torch.no_grad(): + embeddings = model.encode( + texts, + normalize_embeddings=True, + batch_size=BATCH_SIZE, + convert_to_numpy=True, + ).tolist() + _encode_duration = time.time() - _encode_start + + if batch_num is not None and total_batches is not None: + batch_label = f"Batch {batch_num}/{total_batches}" + else: + batch_label = f"Batch de {len(batch)} document(s)" + + log(f"[INFO] {batch_label} encodé en {_encode_duration:.2f}s") + + for item, text, metadata, embedding in zip(batch, texts, metadatas, embeddings): + output_value = { + "vector": embedding, + "text": text + } + + # On ne rajoute "metadata" que si elle existe, pour ne pas polluer + # la sortie du cas "question simple" (sans métadonnées). + if metadata is not None: + output_value["metadata"] = metadata + + output = { + "id": item["id"], + "value": output_value + } + sys.stdout.write(json.dumps(output, ensure_ascii=False)) + sys.stdout.write("\n") + + sys.stdout.flush() + + +# Lecture complète de stdin pour connaître le total +log("[INFO] Lecture des données...") +all_data = [] +for line in sys.stdin: + line = line.strip() + if not line: + continue + all_data.append(json.loads(line)) + +total = len(all_data) +total_batches = (total + BATCH_SIZE - 1) // BATCH_SIZE +log(f"[INFO] {total} documents à vectoriser ({total_batches} batch(s))") + +if total == 0: + log("[WARN] Aucun document reçu.") + sys.exit(0) + +# Traitement par batch +processed = 0 + +for batch_num, i in enumerate(range(0, total, BATCH_SIZE), start=1): + batch = all_data[i:i + BATCH_SIZE] + process_batch(batch, batch_num=batch_num, total_batches=total_batches) + processed += len(batch) + +log(f"[INFO] Vectorisation terminée : {processed}/{total} documents") \ No newline at end of file diff --git a/services/rag-tools/v1/prompt.json b/services/rag-tools/v1/prompt.json new file mode 100644 index 000000000..a10e24d3e --- /dev/null +++ b/services/rag-tools/v1/prompt.json @@ -0,0 +1,34 @@ +{ + "prompts": [ + { + "id": "classification_template", + "name": "Query type classification for RAG routing", + "content": "Tu es un système de classification de requêtes pour un pipeline RAG (question-réponse basé sur des documents).\n\nTa mission est de déterminer le type de traitement à appliquer à la demande actuelle de l'utilisateur, parmi les catégories suivantes :\n\n- \"rag\" : une question générale sur le contenu des documents (recherche d'information, résumé, comparaison, question factuelle, demande d'explication générale, etc.).\n- \"definition\" : une demande explicite de définition d'un terme, mot, concept ou expression (par exemple : \"définis X\", \"qu'est-ce que X\", \"donne-moi la définition de X\", \"que veut dire X dans ce contexte\", \"X, ça veut dire quoi ?\").\n\nConsignes :\n- Appuie-toi en priorité sur le dernier échange de l'historique si la demande actuelle est elliptique (par exemple \"et pour Y ?\" faisant suite à une demande de définition doit probablement rester classé \"definition\").\n- En cas de doute ou d'ambiguïté entre les deux catégories, choisis \"rag\" par défaut : c'est la catégorie générale, adaptée à toute question qui n'est pas explicitement et clairement une demande de définition.\n- Réponds UNIQUEMENT avec l'un des mots suivants, exactement, sans aucun autre texte, explication, ponctuation ou guillemet autour : rag, definition\n\n====================\nHISTORIQUE DE LA CONVERSATION\n====================\n{historique}\n\n====================\nDEMANDE ACTUELLE\n====================\n{question}\n\n====================\nCLASSIFICATION\n====================" + }, + { + "id": "rag_template", + "name": "RAG question answering with sources and metadata", + "content": "Tu es un assistant spécialisé dans la recherche d'informations à partir de documents.\n\nTa mission est de répondre à la question de l'utilisateur en utilisant uniquement les documents fournis dans le contexte, en tenant compte de l'historique de la conversation pour comprendre le contexte de l'échange.\n\nConsignes :\n- Analyse attentivement les documents avant de répondre, y compris leurs métadonnées (titre, auteurs, date de publication, journal, DOI, mots-clés, résumé) lorsqu'elles sont fournies.\n- Utilise le contenu ET les métadonnées des documents pour répondre. Une question peut porter sur le contenu (\"que dit ce document sur X ?\") ou sur les métadonnées elles-mêmes (\"qui a écrit ce document ?\", \"de quand date-t-il ?\", \"dans quelle revue a-t-il été publié ?\").\n- Ne complète pas avec tes connaissances générales.\n- N'invente jamais une métadonnée absente d'un document (par exemple une date ou un auteur) : si l'information demandée n'est pas présente, indique-le clairement.\n- Utilise l'historique de la conversation uniquement pour mieux comprendre la question posée (contexte, reformulation, référence à un échange précédent), jamais comme source d'information factuelle.\n- Si la réponse n'est pas présente ou n'est pas déductible des documents, indique clairement que les documents ne permettent pas de répondre précisément.\n- Donne une réponse claire, concise et directement liée à la question.\n- Si plusieurs documents contiennent des informations utiles, synthétise-les.\n- Cite systématiquement les sources utilisées dans ta réponse, en utilisant l'identifiant exact de chaque document tel qu'il apparaît dans son en-tête (--- IDENTIFIANT ---), jamais un numéro de position.\n- Pour chaque information importante, indique le ou les documents dont elle provient en utilisant le format [IDENTIFIANT], où IDENTIFIANT est repris tel quel depuis l'en-tête du document concerné (par exemple [doi:10.1093/jac/dkm227#chunk1]).\n- N'invente jamais de référence : cite uniquement les documents réellement utilisés.\n\n====================\nHISTORIQUE DE LA CONVERSATION\n====================\n{historique}\n\n====================\nQUESTION UTILISATEUR\n====================\n{question}\n\n====================\nDOCUMENTS DE CONTEXTE\n====================\n{documents}\n\n====================\nREPONSE\n====================" + }, + { + "id": "reformulation_template", + "name": "Query reformulation with conversation history", + "content": "Tu es un assistant spécialisé dans la reformulation de questions pour un système de recherche documentaire (RAG).\n\nTa mission est de reformuler la question actuelle de l'utilisateur en une question autonome, claire et complète, qui pourra être utilisée telle quelle par un moteur de recherche documentaire, SANS avoir besoin de connaître l'historique de la conversation.\n\nConsignes :\n- Appuie-toi en priorité sur le dernier échange (dernière question et dernière réponse) de l'historique pour comprendre le sujet implicite de la question actuelle. Les échanges plus anciens ne servent que si le dernier échange ne suffit pas à lever l'ambiguïté.\n- Si la question actuelle fait référence à des éléments mentionnés précédemment (pronoms, ellipses, sujets implicites, \"et pour...\", \"et lui ?\", \"c'est pareil ?\", \"le dernier document\", etc.), remplace ces références par les informations explicites correspondantes issues de l'historique.\n- Ne reprends JAMAIS un identifiant de document ou une citation entre crochets présents dans l'historique (par exemple [doi:...], [Document X], [ark:...]) : ces identifiants sont des références techniques inexploitables par le moteur de recherche et par les étapes suivantes du pipeline. Si la question fait référence à \"ce document\", \"le dernier document\" ou équivalent, décris plutôt son sujet ou son contenu en langage naturel, à partir de ce qui en a été dit dans l'historique.\n- Ne reprends QUE les éléments nécessaires à la compréhension de la question actuelle. Ne recopie pas l'intégralité d'une question précédente si elle n'est pas directement liée à la question actuelle.\n- Ne transforme pas une question sur les caractéristiques d'un sujet en question sur un autre attribut non mentionné (par exemple, si la question actuelle porte sur des caractéristiques et non sur un classement ou une comparaison, ne rajoute pas de comparaison qui n'était pas demandée).\n- Si la question actuelle est déjà autonome et compréhensible sans l'historique, ne la modifie pas ou modifie-la le moins possible.\n- Ne réponds jamais à la question, tu dois uniquement la reformuler.\n- Ne rajoute aucune information qui ne provient pas de la question ou de l'historique.\n- Conserve la langue d'origine de la question.\n- Ta réponse doit contenir uniquement la question reformulée, sans aucun texte, explication, préfixe ou guillemets autour.\n\n====================\nHISTORIQUE DE LA CONVERSATION\n====================\n{historique}\n\n====================\nQUESTION ACTUELLE\n====================\n{question}\n\n====================\nQUESTION REFORMULEE\n====================" + }, + { + "id": "rerank_template", + "name": "Document reranking for RAG retrieval", + "content": "Tu es un système de reranking pour un pipeline de recherche documentaire (RAG).\n\nTa mission est de classer les documents candidats ci-dessous par ordre de pertinence décroissante par rapport à la question posée, et de ne retenir que les {top_n} documents les plus pertinents.\n\nConsignes :\n- Analyse chaque document indépendamment de sa position dans la liste.\n- Un document est pertinent s'il contient des informations, dans son contenu OU dans ses métadonnées (titre, auteurs, date de publication, journal, mots-clés, résumé), qui aident, même partiellement, à répondre à la question.\n- Si la question porte explicitement sur une métadonnée (un auteur, une date, une revue...), priorise les documents dont les métadonnées y répondent directement.\n- Classe uniquement les documents réellement pertinents, du plus pertinent au moins pertinent.\n- Ne retiens jamais un document hors-sujet uniquement pour atteindre le nombre {top_n} demandé : si moins de documents sont pertinents, retourne-en moins.\n- Si aucun document n'est pertinent, réponds avec une liste vide.\n- Réponds UNIQUEMENT avec la liste des numéros des documents retenus, séparés par des virgules, du plus pertinent au moins pertinent (exemple : 3,1,5). Aucun texte, aucune explication, aucun autre caractère, aucun guillemet.\n\n====================\nQUESTION\n====================\n{question}\n\n====================\nDOCUMENTS CANDIDATS\n====================\n{documents}\n\n====================\nCLASSEMENT\n====================" + }, + { + "id": "definition_reformulation_template", + "name": "Query reformulation for definition extraction", + "content": "Tu es un assistant spécialisé dans la reformulation de demandes de définition pour un système de recherche documentaire (RAG).\n\nTa mission est de reformuler la demande actuelle de l'utilisateur en une requête de recherche autonome, claire et complète, qui pourra être utilisée telle quelle par un moteur de recherche documentaire, SANS avoir besoin de connaître l'historique de la conversation.\n\nConsignes :\n- La demande porte sur la définition d'un terme, potentiellement accompagné d'un contexte qui en précise le sens (domaine, sujet, phrase, discussion précédente).\n- Identifie clairement le terme à définir et son contexte, en t'appuyant sur l'historique si la demande actuelle y fait référence de façon elliptique (pronoms, \"ce terme\", \"et dans ce cas ?\", \"pareil pour X\", \"le dernier document\", etc.).\n- Le contexte peut provenir d'un échange précédent de la conversation (un sujet déjà discuté) ou être fourni directement dans la demande actuelle : dans les deux cas, rends-le explicite dans la requête reformulée.\n- Ne reprends JAMAIS un identifiant de document ou une citation entre crochets présents dans l'historique (par exemple [doi:...], [Document X], [ark:...]) : ces identifiants sont des références techniques inexploitables par le moteur de recherche et par les étapes suivantes du pipeline. Si la demande fait référence à \"ce document\", \"le dernier document\" ou équivalent, décris plutôt son sujet ou son contenu en langage naturel, à partir de ce qui en a été dit dans l'historique.\n- Ignore toute contrainte de formulation de la réponse finale (longueur, mot imposé en premier, style, etc.) : elles ne concernent pas la recherche documentaire et ne doivent pas être incluses dans la requête reformulée.\n- Si la demande actuelle est déjà autonome et compréhensible sans l'historique, ne la modifie pas ou modifie-la le moins possible.\n- Ne fournis jamais toi-même une définition, tu dois uniquement reformuler la requête de recherche.\n- Ne rajoute aucune information qui ne provient pas de la demande ou de l'historique.\n- Conserve la langue d'origine de la demande.\n- Ta réponse doit contenir uniquement la requête reformulée, sans aucun texte, explication, préfixe ou guillemets autour.\n\n====================\nHISTORIQUE DE LA CONVERSATION\n====================\n{historique}\n\n====================\nDEMANDE ACTUELLE\n====================\n{question}\n\n====================\nREQUETE REFORMULEE\n====================" + }, + { + "id": "definition_extraction_template", + "name": "Definition extraction from corpus with sources and general knowledge", + "content": "Tu es un assistant spécialisé dans la construction de définitions à partir de documents.\n\nRÈGLE ABSOLUE SUR LES CONTRAINTES DE FORMULATION : avant de rédiger la définition, identifie toute contrainte de formulation exprimée par l'utilisateur, que ce soit dans la demande actuelle OU plus tôt dans l'historique de la conversation. Une contrainte de formulation porte sur la FORME de la réponse, pas sur son contenu (exemples : commencer la définition par le terme lui-même, respecter une longueur maximale ou minimale, rédiger en une seule phrase, employer un ton ou un registre particulier, s'adresser à un public donné, utiliser une langue précise, structurer la réponse d'une certaine façon, éviter certains mots...). Applique CHAQUE contrainte identifiée à la lettre, sans exception, même si le résultat te semble moins naturel que ta formulation par défaut. Une fois qu'une contrainte a été exprimée dans la conversation, elle reste valable pour les demandes suivantes tant que l'utilisateur ne l'annule pas explicitement. En l'absence de toute contrainte exprimée, rédige normalement, sans format imposé.\n\nTa mission est de construire une définition claire et cohérente d'un terme, dans le contexte précisé par l'utilisateur, en te basant en priorité sur les documents fournis, complétés si nécessaire par tes connaissances générales.\n\nIMPORTANT : constater que les documents ne donnent pas de définition explicite ou précise du terme N'EST JAMAIS une réponse finale suffisante à elle seule. C'est uniquement le signal qu'il faut passer à l'étape suivante de la procédure ci-dessous. Ta réponse finale doit toujours se terminer par une définition (sauf cas d'échec total prévu à l'étape 4), qui respecte les éventuelles contraintes de formulation identifiées ci-dessus.\n\nProcède dans cet ordre, sans sauter d'étape :\n\n1. Cherche dans les documents fournis toute information exploitable sur le terme dans le contexte demandé : une définition explicite, mais aussi tout usage, fonction décrite, exemple, propriété, ou mention du terme, même indirecte ou incomplète.\n\n2. Si tu trouves des éléments exploitables (même partiels), construis une définition à partir de ceux-ci. Si ces éléments ne couvrent pas entièrement le sens du terme, complète avec tes connaissances générales (voir étape 3) pour obtenir une définition complète, en indiquant clairement quelle partie vient des documents et quelle partie vient de tes connaissances générales.\n\n3. Si les documents mentionnent le terme sans le définir (comme un simple nom, une référence, ou un usage sans explication), et que ce terme a une définition générale bien établie et non ambiguë dans son domaine (ce qui est le cas de la plupart des termes techniques, scientifiques ou mathématiques standards), tu DOIS fournir cette définition générale, explicitement signalée comme telle (par exemple : \"Les documents ne définissent pas explicitement ce terme, mais il s'agit d'un concept standard en [domaine] : [définition]\"), puis la relier au contexte spécifique mentionné dans les documents.\n\n4. Ne conclus à une impossibilité totale de répondre que si les deux conditions suivantes sont réunies : (a) les documents ne fournissent aucun élément exploitable sur le terme dans ce contexte, ET (b) le terme n'a pas de définition générale bien établie que tu pourrais fournir avec confiance (terme trop spécialisé, trop ambigu, ou trop incertain). Dans ce seul cas, indique-le clairement (par exemple : \"Ni les documents fournis ni des connaissances générales fiables ne permettent de construire une définition de ce terme dans ce contexte.\"), en respectant quand même les éventuelles contraintes de formulation si elles restent applicables à ce message d'échec.\n\nAutres consignes :\n- Base la définition en priorité sur les documents fournis. Sépare TOUJOURS explicitement et clairement ce qui provient des documents (avec citation [IDENTIFIANT]) de ce qui provient de tes connaissances générales (sans citation, mais explicitement signalé comme tel). Ne présente jamais une information issue de tes connaissances générales comme si elle provenait des documents, et inversement.\n- Reste prudent avec tes connaissances générales : n'avance que des éléments bien établis et consensuels. N'invente jamais de faits précis et vérifiables (chiffres, dates, noms, formules, résultats spécifiques) dont tu n'es pas certain.\n- En cas de désaccord entre les documents et tes connaissances générales, les documents priment toujours pour le contexte demandé.\n- Si plusieurs documents apportent des éléments complémentaires sur le terme, synthétise-les en une définition unique et cohérente, sans contradiction non résolue.\n- Cite systématiquement les documents utilisés, en utilisant l'identifiant exact de chaque document tel qu'il apparaît dans son en-tête (--- IDENTIFIANT ---), au format [IDENTIFIANT]. N'invente jamais de référence : cite uniquement les documents réellement utilisés. Si une contrainte de formulation entre en tension avec la citation des sources (par exemple \"une seule phrase\"), privilégie quand même une citation minimale plutôt que de l'omettre complètement.\n- Utilise l'historique de la conversation pour identifier le terme, le contexte, et les contraintes de formulation déjà exprimées, jamais comme source d'information factuelle pour le contenu de la définition elle-même.\n\n====================\nHISTORIQUE DE LA CONVERSATION\n====================\n{historique}\n\n====================\nDEMANDE ACTUELLE\n====================\n{question}\n\n====================\nDOCUMENTS DE CONTEXTE\n====================\n{documents}\n\n====================\nDEFINITION\n====================" + } + ] +} \ No newline at end of file diff --git a/services/rag-tools/v1/rag-session.ini b/services/rag-tools/v1/rag-session.ini new file mode 100644 index 000000000..00096a6b5 --- /dev/null +++ b/services/rag-tools/v1/rag-session.ini @@ -0,0 +1,40 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = Rag (version session) +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Rag-session +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result +post.parameters.2.in = query +post.parameters.2.name = type +post.parameters.2.schema.type = string +post.parameters.2.description = Type de traitement RAG (Choix : Rag ou definition) + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/rag-session.py +args = fix('-p') +args = env('type', "rag") + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/rag-session.py b/services/rag-tools/v1/rag-session.py new file mode 100755 index 000000000..8167aba23 --- /dev/null +++ b/services/rag-tools/v1/rag-session.py @@ -0,0 +1,379 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os +import uuid + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID_RAG = "rag_template" +PROMPT_ID_DEFINITION = "definition_extraction_template" + +NO_HISTORY_TEXT = "Aucun historique de conversation disponible." + +TMP_DIR = "/tmp" + + + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + + +# ============================== +# Arguments optionnels +# ============================== + +rag_type = sys.argv[sys.argv.index("-p") + 1] if "-p" in sys.argv else "rag" +rag_type = "rag" if rag_type not in ["definition"] else rag_type +print_log("Rag type : " + rag_type) + +prompt_id = PROMPT_ID_RAG +if rag_type == "definition": + prompt_id = PROMPT_ID_DEFINITION +print_log("Prompt ID : " + prompt_id) +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(prompt_id): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == prompt_id: + return prompt["content"] + + raise ValueError(f"Prompt {prompt_id} not found") + + +PROMPT_TEMPLATE = load_prompt(prompt_id) + + +# ============================== +# Gestion de la session (token -> historique persisté sur disque) +# ============================== + +def session_dir(token: str) -> str: + return os.path.join(TMP_DIR, token) + + +def history_path(token: str) -> str: + return os.path.join(session_dir(token), f"{token}.json") + + +def generate_token() -> str: + return uuid.uuid4().hex + + +def load_historique(token: str) -> list: + """Charge l'historique associé à un token. Renvoie une liste vide si + le fichier n'existe pas encore.""" + path = history_path(token) + if not os.path.exists(path): + print_log(f"Aucun historique trouvé pour le token {token}, historique vide utilisé") + return [] + + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def save_historique(token: str, historique: list) -> None: + """Écrit l'historique mis à jour sur disque, dans tmp//.json. + Crée le dossier de session s'il n'existe pas encore (premier appel).""" + os.makedirs(session_dir(token), exist_ok=True) + path = history_path(token) + + with open(path, "w", encoding="utf-8") as f: + json.dump(historique, f, ensure_ascii=False, indent=2) + + +def resolve_session(token_value): + """ + `token_value` est ce qui se trouve dans value["historique"] en entrée : + - absent / None / chaîne vide -> premier appel, on génère un nouveau + token et l'historique de départ est vide. + - une chaîne non vide -> token existant, on charge l'historique + correspondant depuis le disque. + + Renvoie (token, historique). + """ + if not token_value: + new_token = generate_token() + print_log(f"Aucun token fourni, nouvelle session créée : {new_token}") + return new_token, [] + + return token_value, load_historique(token_value) + + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 10000 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"] + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Construction de l'historique +# ============================== + +def build_history_text(historique): + """ + Construit le texte de l'historique de la conversation. + + `historique` est attendu comme une liste de tours de dialogue, + chaque tour étant un dict avec les clés "role" et "content", + par exemple : + [ + {"role": "user", "content": "..."}, + {"role": "assistant", "content": "..."} + ] + + Si aucun historique n'est fourni (liste vide), un texte par défaut + est renvoyé. + """ + if not historique: + return NO_HISTORY_TEXT + + lines = [] + + for turn in historique: + role = turn.get("role", "inconnu") + content = turn.get("content", "") + lines.append(f"{role} : {content}") + + return "\n".join(lines) + + +# ============================== +# Sérialisation des métadonnées +# ============================== + +def format_metadata_block(metadata: dict) -> str: + """ + Construit un bloc texte lisible à partir des métadonnées d'un document. + Les champs absents sont simplement ignorés (pas de ligne vide/"None"). + """ + if not metadata: + return "" + + lines = [] + + if metadata.get("titre"): + lines.append(f"Titre : {metadata['titre']}") + + auteurs = metadata.get("auteurs") + if auteurs: + lines.append(f"Auteurs : {', '.join(auteurs)}") + + if metadata.get("date_publication"): + lines.append(f"Date de publication : {metadata['date_publication']}") + + if metadata.get("journal"): + journal_line = f"Journal : {metadata['journal']}" + details = [] + if metadata.get("volume"): + details.append(f"volume {metadata['volume']}") + if metadata.get("numero"): + details.append(f"numéro {metadata['numero']}") + if metadata.get("pages"): + details.append(f"pages {metadata['pages']}") + if details: + journal_line += " (" + ", ".join(details) + ")" + lines.append(journal_line) + + if metadata.get("doi"): + lines.append(f"DOI : {metadata['doi']}") + + mots_cles = metadata.get("mots_cles") + if mots_cles: + lines.append(f"Mots-clés : {', '.join(mots_cles)}") + + if metadata.get("resume"): + lines.append(f"Résumé de l'article : {metadata['resume']}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt RAG +# ============================== + +def build_prompt(question, documents, historique): + documents_text = "" + + for i, doc in enumerate(documents, start=1): + # L'identifiant du chunk (ex: "doi:10.xxxx/xxx#chunk1") sert de + # label de citation dans la réponse finale du LLM, plutôt qu'un + # simple numéro de position. Fallback sur "Document i" si l'id + # est absent (ne devrait pas arriver en usage normal). + doc_label = doc.get("id") or f"Document {i}" + + metadata_block = format_metadata_block(doc.get("metadata")) + text = doc.get("text", "") + + block = f"--- {doc_label} ---\n" + if metadata_block: + block += metadata_block + "\n" + block += f"Contenu :\n{text}\n" + block += f"--- Fin du document {doc_label} ---\n\n" + + documents_text += block + + historique_text = build_history_text(historique) + + return PROMPT_TEMPLATE.format( + question=question, + documents=documents_text, + historique=historique_text + ) + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + documents = item["value"]["documents"] + token, historique = resolve_session(item["value"].get("historique")) + + prompt = build_prompt( + question, + documents, + historique + ) + + answer = call_llm(prompt) + + # Seul ce script écrit l'historique : c'est ici que la question + # originale et la réponse finale sont connues simultanément et + # forment un tour de conversation complet à persister. Les autres + # services (classification, reformulation) ne produisent que des + # informations intermédiaires qui ne font pas partie de l'historique. + updated_historique = historique + [ + {"role": "user", "content": question}, + {"role": "assistant", "content": answer}, + ] + save_historique(token, updated_historique) + + output = { + "id": item["id"], + "value": { + "question": question, + "documents": documents, + "token": token, + "answer": answer + } + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch) diff --git a/services/rag-tools/v1/rag.ini b/services/rag-tools/v1/rag.ini new file mode 100644 index 000000000..f0a415095 --- /dev/null +++ b/services/rag-tools/v1/rag.ini @@ -0,0 +1,40 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = Rag +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Rag +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result +post.parameters.2.in = query +post.parameters.2.name = type +post.parameters.2.schema.type = string +post.parameters.2.description = Type de traitement RAG (Choix : Rag ou definition) + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/rag.py +args = fix('-p') +args = env('type', "rag") + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/rag.py b/services/rag-tools/v1/rag.py new file mode 100755 index 000000000..1323045dc --- /dev/null +++ b/services/rag-tools/v1/rag.py @@ -0,0 +1,308 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID_RAG = "rag_template" +PROMPT_ID_DEFINITION = "definition_extraction_template" + +NO_HISTORY_TEXT = "Aucun historique de conversation disponible." + + + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + + +# ============================== +# Arguments optionnels +# ============================== + +rag_type = sys.argv[sys.argv.index("-p") + 1] if "-p" in sys.argv else "rag" +rag_type = "rag" if rag_type not in ["definition"] else rag_type +print_log("Rag type : " + rag_type) + +prompt_id = PROMPT_ID_RAG +if rag_type == "definition": + prompt_id = PROMPT_ID_DEFINITION +print_log("Prompt ID : " + prompt_id) +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(prompt_id): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == prompt_id: + return prompt["content"] + + raise ValueError(f"Prompt {prompt_id} not found") + + +PROMPT_TEMPLATE = load_prompt(prompt_id) + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 10000 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"] + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Construction de l'historique +# ============================== + +def build_history_text(historique): + """ + Construit le texte de l'historique de la conversation. + + `historique` est attendu comme une liste de tours de dialogue, + chaque tour étant un dict avec les clés "role" et "content", + par exemple : + [ + {"role": "user", "content": "..."}, + {"role": "assistant", "content": "..."} + ] + + Si aucun historique n'est fourni (absent, None ou liste vide), + un texte par défaut est renvoyé. + """ + if not historique: + return NO_HISTORY_TEXT + + lines = [] + + for turn in historique: + role = turn.get("role", "inconnu") + content = turn.get("content", "") + lines.append(f"{role} : {content}") + + return "\n".join(lines) + + +# ============================== +# Sérialisation des métadonnées +# ============================== + +def format_metadata_block(metadata: dict) -> str: + """ + Construit un bloc texte lisible à partir des métadonnées d'un document. + Les champs absents sont simplement ignorés (pas de ligne vide/"None"). + """ + if not metadata: + return "" + + lines = [] + + if metadata.get("titre"): + lines.append(f"Titre : {metadata['titre']}") + + auteurs = metadata.get("auteurs") + if auteurs: + lines.append(f"Auteurs : {', '.join(auteurs)}") + + if metadata.get("date_publication"): + lines.append(f"Date de publication : {metadata['date_publication']}") + + if metadata.get("journal"): + journal_line = f"Journal : {metadata['journal']}" + details = [] + if metadata.get("volume"): + details.append(f"volume {metadata['volume']}") + if metadata.get("numero"): + details.append(f"numéro {metadata['numero']}") + if metadata.get("pages"): + details.append(f"pages {metadata['pages']}") + if details: + journal_line += " (" + ", ".join(details) + ")" + lines.append(journal_line) + + if metadata.get("doi"): + lines.append(f"DOI : {metadata['doi']}") + + mots_cles = metadata.get("mots_cles") + if mots_cles: + lines.append(f"Mots-clés : {', '.join(mots_cles)}") + + if metadata.get("resume"): + lines.append(f"Résumé de l'article : {metadata['resume']}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt RAG +# ============================== + +def build_prompt(question, documents, historique): + documents_text = "" + + for i, doc in enumerate(documents, start=1): + # L'identifiant du chunk (ex: "doi:10.xxxx/xxx#chunk1") sert de + # label de citation dans la réponse finale du LLM, plutôt qu'un + # simple numéro de position. Fallback sur "Document i" si l'id + # est absent (ne devrait pas arriver en usage normal). + doc_label = doc.get("id") or f"Document {i}" + + metadata_block = format_metadata_block(doc.get("metadata")) + text = doc.get("text", "") + + block = f"--- {doc_label} ---\n" + if metadata_block: + block += metadata_block + "\n" + block += f"Contenu :\n{text}\n" + block += f"--- Fin du document {doc_label} ---\n\n" + + documents_text += block + + historique_text = build_history_text(historique) + + return PROMPT_TEMPLATE.format( + question=question, + documents=documents_text, + historique=historique_text + ) + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + documents = item["value"]["documents"] + historique = item["value"].get("historique", []) + + prompt = build_prompt( + question, + documents, + historique + ) + + answer = call_llm(prompt) + + output = { + "id": item["id"], + "value": { + "question": question, + "documents": documents, + "historique": historique, + "answer": answer + } + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch) diff --git a/services/rag-tools/v1/reformulate-session.ini b/services/rag-tools/v1/reformulate-session.ini new file mode 100644 index 000000000..c1cdc0dd2 --- /dev/null +++ b/services/rag-tools/v1/reformulate-session.ini @@ -0,0 +1,41 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = reformulate (version session) +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Reformule les questions +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result +post.parameters.2.in = query +post.parameters.2.name = type +post.parameters.2.schema.type = string +post.parameters.2.description = Type de traitement RAG (Choix : Rag ou definition) + + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/reformulate-session.py +args = fix('-p') +args = env('type', "rag") + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/reformulate-session.py b/services/rag-tools/v1/reformulate-session.py new file mode 100755 index 000000000..74836a0b4 --- /dev/null +++ b/services/rag-tools/v1/reformulate-session.py @@ -0,0 +1,284 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os +import uuid + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID_REFORMULATE = "reformulation_template" +PROMPT_ID_DEFINITION = "definition_reformulation_template" + +NO_HISTORY_TEXT = "Aucun historique de conversation disponible." + +TMP_DIR = "/tmp" + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + +# ============================== +# Arguments optionnels +# ============================== + +rag_type = sys.argv[sys.argv.index("-p") + 1] if "-p" in sys.argv else "rag" +rag_type = "rag" if rag_type not in ["definition"] else rag_type +print_log("Rag type : " + rag_type) +prompt_id = PROMPT_ID_REFORMULATE +if rag_type == "definition": + prompt_id = PROMPT_ID_DEFINITION + +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(prompt_id): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == prompt_id: + return prompt["content"] + + raise ValueError(f"Prompt {prompt_id} not found") + + +PROMPT_TEMPLATE = load_prompt(prompt_id) + + +# ============================== +# Gestion de la session (token -> historique persisté sur disque) +# ============================== + +def history_path(token: str) -> str: + return os.path.join(TMP_DIR, token, f"{token}.json") + + +def generate_token() -> str: + return uuid.uuid4().hex + + +def load_historique(token: str) -> list: + """Charge l'historique associé à un token. Renvoie une liste vide si + le fichier n'existe pas encore.""" + path = history_path(token) + if not os.path.exists(path): + print_log(f"Aucun historique trouvé pour le token {token}, historique vide utilisé") + return [] + + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def resolve_session(token_value): + """ + `token_value` est ce qui se trouve dans value["historique"] en entrée : + - absent / None / chaîne vide -> premier appel, on génère un nouveau + token et l'historique de départ est vide. + - une chaîne non vide -> token existant, on charge l'historique + correspondant depuis le disque. + + Renvoie (token, historique). + """ + if not token_value: + new_token = generate_token() + print_log(f"Aucun token fourni, nouvelle session créée : {new_token}") + return new_token, [] + + return token_value, load_historique(token_value) + + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 2000 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"].strip() + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Construction de l'historique +# ============================== + +def build_history_text(historique): + """ + Construit le texte de l'historique de la conversation. + + `historique` est attendu comme une liste de tours de dialogue, + chaque tour étant un dict avec les clés "role" et "content". + + Si aucun historique n'est fourni (liste vide), un texte par défaut + est renvoyé, et on considère qu'il n'y a pas besoin de reformuler. + """ + if not historique: + return NO_HISTORY_TEXT + + lines = [] + + for turn in historique: + role = turn.get("role", "inconnu") + content = turn.get("content", "") + lines.append(f"{role} : {content}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt de reformulation +# ============================== + +def build_prompt(question, historique): + historique_text = build_history_text(historique) + + return PROMPT_TEMPLATE.format( + question=question, + historique=historique_text + ) + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + token, historique = resolve_session(item["value"].get("historique")) + + # Pas d'historique -> pas besoin d'appeler le LLM, + # la question reste inchangée. + if not historique: + reformulated_question = question + else: + prompt = build_prompt( + question, + historique + ) + + reformulated_question = call_llm(prompt) + + if reformulated_question == "Error": + # En cas d'échec, on retombe sur la question d'origine + # plutôt que de bloquer le pipeline. + print_log( + "Reformulation failed, " + "falling back to original question" + ) + reformulated_question = question + + output = { + "id": item["id"], + "value": { + "question": reformulated_question, + "token": token + } + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch) diff --git a/services/rag-tools/v1/reformulate.ini b/services/rag-tools/v1/reformulate.ini new file mode 100644 index 000000000..a921ba3f1 --- /dev/null +++ b/services/rag-tools/v1/reformulate.ini @@ -0,0 +1,41 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = reformulate +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Reformule les questions +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result +post.parameters.2.in = query +post.parameters.2.name = type +post.parameters.2.schema.type = string +post.parameters.2.description = Type de traitement RAG (Choix : Rag ou definition) + + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/reformulate.py +args = fix('-p') +args = env('type', "rag") + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/reformulate.py b/services/rag-tools/v1/reformulate.py new file mode 100755 index 000000000..9b14b9b43 --- /dev/null +++ b/services/rag-tools/v1/reformulate.py @@ -0,0 +1,236 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID_REFORMULATE = "reformulation_template" +PROMPT_ID_DEFINITION = "definition_reformulation_template" + +NO_HISTORY_TEXT = "Aucun historique de conversation disponible." + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + +# ============================== +# Arguments optionnels +# ============================== + +rag_type = sys.argv[sys.argv.index("-p") + 1] if "-p" in sys.argv else "rag" +rag_type = "rag" if rag_type not in ["definition"] else rag_type +print_log("Rag type : " + rag_type) +prompt_id = PROMPT_ID_REFORMULATE +if rag_type == "definition": + prompt_id = PROMPT_ID_DEFINITION + +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(prompt_id): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == prompt_id: + return prompt["content"] + + raise ValueError(f"Prompt {prompt_id} not found") + + +PROMPT_TEMPLATE = load_prompt(prompt_id) + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 2000 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"].strip() + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Construction de l'historique +# ============================== + +def build_history_text(historique): + """ + Construit le texte de l'historique de la conversation. + + `historique` est attendu comme une liste de tours de dialogue, + chaque tour étant un dict avec les clés "role" et "content". + + Si aucun historique n'est fourni (absent, None ou liste vide), + un texte par défaut est renvoyé, et on considère qu'il n'y a + pas besoin de reformuler. + """ + if not historique: + return NO_HISTORY_TEXT + + lines = [] + + for turn in historique: + role = turn.get("role", "inconnu") + content = turn.get("content", "") + lines.append(f"{role} : {content}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt de reformulation +# ============================== + +def build_prompt(question, historique): + historique_text = build_history_text(historique) + + return PROMPT_TEMPLATE.format( + question=question, + historique=historique_text + ) + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + historique = item["value"].get("historique", []) + + # Pas d'historique -> pas besoin d'appeler le LLM, + # la question reste inchangée. + if not historique: + reformulated_question = question + else: + prompt = build_prompt( + question, + historique + ) + + reformulated_question = call_llm(prompt) + + if reformulated_question == "Error": + # En cas d'échec, on retombe sur la question d'origine + # plutôt que de bloquer le pipeline. + print_log( + "Reformulation failed, " + "falling back to original question" + ) + reformulated_question = question + + output = { + "id": item["id"], + "value": reformulated_question + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch) \ No newline at end of file diff --git a/services/rag-tools/v1/rerank.ini b/services/rag-tools/v1/rerank.ini new file mode 100644 index 000000000..fce018ff7 --- /dev/null +++ b/services/rag-tools/v1/rerank.ini @@ -0,0 +1,34 @@ +# OpenAPI Documentation - JSON format (dot notation) +post.responses.default.description = reranking de documents +post.responses.default.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.summary = Rerank les documents +post.requestBody.required = true +post.requestBody.content.application/json.schema.$ref = #/components/schemas/JSONStream +post.parameters.0.in = query +post.parameters.0.name = path +post.parameters.0.schema.type = string +post.parameters.0.description = The path in each object to enrich with a Python script +post.parameters.1.in = query +post.parameters.1.name = indent +post.parameters.1.schema.type = boolean +post.parameters.1.description = Indent or not the JSON Result + +[use] +plugin = @ezs/spawn +plugin = @ezs/basics + +[JSONParse] +separator = * + +[expand] +path = env('path', 'value') +size = 128 +# in production mode, uncomment the following line +# cache = boost + +[expand/exec] +# command should be executable ! +command = ./v1/rerank.py + +[dump] +indent = env('indent', false) \ No newline at end of file diff --git a/services/rag-tools/v1/rerank.py b/services/rag-tools/v1/rerank.py new file mode 100755 index 000000000..8095cde71 --- /dev/null +++ b/services/rag-tools/v1/rerank.py @@ -0,0 +1,326 @@ +#!/usr/bin/env python3 + +import sys +import json +import time +import requests +import os + +# ============================== +# Configuration +# ============================== + + +API_KEY = os.getenv("ILAAS_API_KEY") +MODEL_NAME = "gemma-4-31b" +MAX_RETRIES = 4 +RETRY_DELAY = 2 +BATCH_SIZE = 32 + +PROMPT_PATH = "v1/prompt.json" +PROMPT_ID = "rerank_template" + +DEFAULT_TOP_N = 6 + + +# ============================== +# Chargement du prompt +# ============================== + +def load_prompt(): + with open(PROMPT_PATH, "r", encoding="utf-8") as f: + data = json.load(f) + + for prompt in data["prompts"]: + if prompt["id"] == PROMPT_ID: + return prompt["content"] + + raise ValueError(f"Prompt {PROMPT_ID} not found") + + +PROMPT_TEMPLATE = load_prompt() + + +# ============================== +# Logs +# ============================== + +def print_log(message): + print(message, file=sys.stderr) + + +# ============================== +# Appel LLM +# ============================== + +def call_llm(prompt: str) -> str: + base_url = "https://llm.ilaas.fr/v1" + + headers = { + "Authorization": f"Bearer {API_KEY}", + "Content-Type": "application/json" + } + + payload = { + "model": MODEL_NAME, + "messages": [ + { + "role": "user", + "content": f"{prompt}" + } + ], + "stream": False, + "max_tokens": 200 + } + + for attempt in range(1, MAX_RETRIES + 1): + try: + response = requests.post( + f"{base_url}/chat/completions", + headers=headers, + json=payload, + timeout=60 + ) + + result = response.json() + + print_log( + "LLM result call : " + + result["choices"][0]["message"]["content"] + ) + + print_log( + result["choices"][0]["message"].get( + "reasoning_content", + None + ) + ) + + return result["choices"][0]["message"]["content"].strip() + + except Exception: + print_log( + f"Error while calling LLM " + f"(attempt {attempt}/{MAX_RETRIES})" + ) + + if attempt == MAX_RETRIES: + return "Error" + + time.sleep(RETRY_DELAY * attempt) + print_log( + "Sleeping " + + str(RETRY_DELAY * attempt) + ) + + +# ============================== +# Sérialisation des métadonnées +# ============================== + +def format_metadata_block(metadata: dict) -> str: + """ + Construit un bloc texte lisible à partir des métadonnées d'un document. + Les champs absents sont simplement ignorés (pas de ligne vide/"None"). + """ + if not metadata: + return "" + + lines = [] + + if metadata.get("titre"): + lines.append(f"Titre : {metadata['titre']}") + + auteurs = metadata.get("auteurs") + if auteurs: + lines.append(f"Auteurs : {', '.join(auteurs)}") + + if metadata.get("date_publication"): + lines.append(f"Date de publication : {metadata['date_publication']}") + + if metadata.get("journal"): + journal_line = f"Journal : {metadata['journal']}" + details = [] + if metadata.get("volume"): + details.append(f"volume {metadata['volume']}") + if metadata.get("numero"): + details.append(f"numéro {metadata['numero']}") + if metadata.get("pages"): + details.append(f"pages {metadata['pages']}") + if details: + journal_line += " (" + ", ".join(details) + ")" + lines.append(journal_line) + + if metadata.get("doi"): + lines.append(f"DOI : {metadata['doi']}") + + mots_cles = metadata.get("mots_cles") + if mots_cles: + lines.append(f"Mots-clés : {', '.join(mots_cles)}") + + if metadata.get("resume"): + lines.append(f"Résumé de l'article : {metadata['resume']}") + + return "\n".join(lines) + + +# ============================== +# Construction du prompt de reranking +# ============================== + +def build_documents_text(documents: list[dict]) -> str: + documents_text = "" + + for i, doc in enumerate(documents, start=1): + metadata_block = format_metadata_block(doc.get("metadata")) + text = doc.get("text", "") + + block = f"--- Document {i} ---\n" + if metadata_block: + block += metadata_block + "\n" + block += f"Contenu :\n{text}\n" + block += f"--- Fin du document {i} ---\n\n" + + documents_text += block + + return documents_text + + +def build_prompt(question, documents, top_n): + documents_text = build_documents_text(documents) + + return PROMPT_TEMPLATE.format( + question=question, + documents=documents_text, + top_n=top_n + ) + + +# ============================== +# Parsing de la réponse du LLM +# ============================== + +def parse_ranking(raw_response, documents, top_n): + """ + Attendu : une chaîne du type "3,1,5". + Renvoie la liste des documents originaux (dict {text, metadata}), + réordonnée et tronquée à top_n, dans l'ordre de pertinence donné + par le LLM. + + En cas de réponse invalide ou d'erreur LLM, on retombe sur les + top_n premiers documents dans leur ordre d'origine (fallback + non bloquant plutôt que de casser le pipeline). + """ + if raw_response == "Error" or not raw_response: + print_log("Ranking invalide, fallback sur l'ordre d'origine") + return documents[:top_n] + + cleaned = raw_response.strip() + + if not cleaned: + return documents[:top_n] + + try: + indices = [ + int(x.strip()) + for x in cleaned.split(",") + if x.strip() != "" + ] + except ValueError: + print_log( + f"Impossible de parser le classement '{raw_response}', " + "fallback sur l'ordre d'origine" + ) + return documents[:top_n] + + reranked = [] + seen_indices = set() + + for idx in indices: + if 1 <= idx <= len(documents) and idx not in seen_indices: + reranked.append(documents[idx - 1]) + seen_indices.add(idx) + + if len(reranked) >= top_n: + break + + if not reranked: + return documents[:top_n] + + return reranked + + +# ============================== +# Traitement batch +# ============================== + +def process_batch(batch): + if not batch: + return + + for item in batch: + question = item["value"]["question"] + documents = item["value"]["documents"] + top_n = item["value"].get("top_n", DEFAULT_TOP_N) + + if not documents: + output = { + "id": item["id"], + "value": { + "documents": [] + } + } + else: + prompt = build_prompt( + question, + documents, + top_n + ) + + raw_response = call_llm(prompt) + + reranked_documents = parse_ranking( + raw_response, + documents, + top_n + ) + + output = { + "id": item["id"], + "value": { + "documents": reranked_documents + } + } + + sys.stdout.write( + json.dumps( + output, + ensure_ascii=False + ) + ) + sys.stdout.write("\n") + + +# ============================== +# Main +# ============================== + +batch = [] + +for line in sys.stdin: + line = line.strip() + + if not line: + continue + + data = json.loads(line) + + batch.append(data) + + if len(batch) >= BATCH_SIZE: + process_batch(batch) + batch = [] + + +# Dernier batch incomplet +process_batch(batch)