Compare commits
4 Commits
v0.12.0-be
...
v0.12.2-be
| Author | SHA1 | Date | |
|---|---|---|---|
| 7f519588b6 | |||
| 0799c850ec | |||
| 113df6a391 | |||
| a1f3b9b796 |
@@ -5,6 +5,7 @@ port (LLM, embeddings, extracteur PDF), en fonction des Settings — modifiables
|
|||||||
à chaud depuis l'écran Paramètres de l'UI. Les routers ne connaissent que les
|
à chaud depuis l'écran Paramètres de l'UI. Les routers ne connaissent que les
|
||||||
ports et les use cases, jamais Ollama/Mistral/etc.
|
ports et les use cases, jamais Ollama/Mistral/etc.
|
||||||
"""
|
"""
|
||||||
|
import logging
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
||||||
from fastapi import Depends, HTTPException
|
from fastapi import Depends, HTTPException
|
||||||
@@ -29,11 +30,39 @@ from app.infrastructure.onemin_adapter import OneMinAiLLMProvider
|
|||||||
from app.infrastructure.openrouter_adapter import OpenRouterLLMProvider
|
from app.infrastructure.openrouter_adapter import OpenRouterLLMProvider
|
||||||
from app.infrastructure.pdf_extractor import PyMuPdfTextExtractor
|
from app.infrastructure.pdf_extractor import PyMuPdfTextExtractor
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Extracteur PDF partagé : la détection OCR (version Tesseract) a un coût
|
# Extracteur PDF partagé : la détection OCR (version Tesseract) a un coût
|
||||||
# (subprocess) qu'on ne veut pas payer à chaque requête → singleton module.
|
# (subprocess) qu'on ne veut pas payer à chaque requête → singleton module.
|
||||||
_PDF_EXTRACTOR = PyMuPdfTextExtractor()
|
_PDF_EXTRACTOR = PyMuPdfTextExtractor()
|
||||||
|
|
||||||
|
|
||||||
|
def _effective_import_chunk_tokens(settings: Settings) -> int:
|
||||||
|
"""Taille de morceau réellement utilisable pour l'import.
|
||||||
|
|
||||||
|
Avec Ollama, le morceau (entrée) ET sa réécriture en sections (sortie ≈ même
|
||||||
|
taille) doivent tenir ensemble dans `num_ctx` — sinon Ollama remplit la fenêtre
|
||||||
|
avec le prompt et la génération s'arrête après quelques tokens (JSON coupé net,
|
||||||
|
morceau perdu). Budget : entrée×~1.3 (les morceaux sont mesurés en tokens
|
||||||
|
cl100k, plus compacts que les tokenizers locaux) + consignes + sortie×~1.4
|
||||||
|
≤ num_ctx → morceau ≤ (num_ctx − 800) / 2.7. On plafonne, avec un log pour
|
||||||
|
rester transparent. Les providers cloud (gros contexte) ne sont pas plafonnés.
|
||||||
|
"""
|
||||||
|
requested = settings.import_chunk_tokens
|
||||||
|
if settings.llm_provider != "ollama":
|
||||||
|
return requested
|
||||||
|
cap = max(1000, int((settings.llm_num_ctx - 800) / 2.7))
|
||||||
|
if requested > cap:
|
||||||
|
logger.warning(
|
||||||
|
"Taille de morceau d'import réduite de %s à %s tokens : avec num_ctx=%s, "
|
||||||
|
"un morceau plus gros ne laisserait pas la place à la sortie du modèle "
|
||||||
|
"(génération coupée). Augmentez num_ctx pour utiliser de plus gros morceaux.",
|
||||||
|
requested, cap, settings.llm_num_ctx,
|
||||||
|
)
|
||||||
|
return cap
|
||||||
|
return requested
|
||||||
|
|
||||||
|
|
||||||
def get_llm_provider(
|
def get_llm_provider(
|
||||||
settings: Annotated[Settings, Depends(get_settings)],
|
settings: Annotated[Settings, Depends(get_settings)],
|
||||||
) -> LLMProvider:
|
) -> LLMProvider:
|
||||||
@@ -82,8 +111,14 @@ def get_import_rules_use_case(
|
|||||||
settings: Annotated[Settings, Depends(get_settings)],
|
settings: Annotated[Settings, Depends(get_settings)],
|
||||||
) -> ImportRulesUseCase:
|
) -> ImportRulesUseCase:
|
||||||
"""Factory du use case d'import de règles PDF (extraction + structuration)."""
|
"""Factory du use case d'import de règles PDF (extraction + structuration)."""
|
||||||
|
# Modèle LOCAL → mode segmentation : le LLM ne renvoie que les frontières des
|
||||||
|
# sections (~200 tokens) et le texte original est découpé localement. Réécrire
|
||||||
|
# tout le contenu à ~100 tokens/s prendrait des dizaines de minutes par livre.
|
||||||
|
# Les providers cloud (rapides, grand contexte) gardent la réécriture nettoyée.
|
||||||
return ImportRulesUseCase(
|
return ImportRulesUseCase(
|
||||||
llm=llm, extractor=_PDF_EXTRACTOR, chunk_target_tokens=settings.import_chunk_tokens)
|
llm=llm, extractor=_PDF_EXTRACTOR,
|
||||||
|
chunk_target_tokens=_effective_import_chunk_tokens(settings),
|
||||||
|
segment_only=settings.llm_provider == "ollama")
|
||||||
|
|
||||||
|
|
||||||
def get_import_campaign_use_case(
|
def get_import_campaign_use_case(
|
||||||
@@ -94,7 +129,7 @@ def get_import_campaign_use_case(
|
|||||||
return ImportCampaignUseCase(
|
return ImportCampaignUseCase(
|
||||||
llm=llm,
|
llm=llm,
|
||||||
extractor=_PDF_EXTRACTOR,
|
extractor=_PDF_EXTRACTOR,
|
||||||
chunk_target_tokens=settings.import_chunk_tokens,
|
chunk_target_tokens=_effective_import_chunk_tokens(settings),
|
||||||
map_concurrency=settings.llm_map_concurrency,
|
map_concurrency=settings.llm_map_concurrency,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -29,7 +29,12 @@ from app.domain.models import (
|
|||||||
RoomProposal,
|
RoomProposal,
|
||||||
SceneProposal,
|
SceneProposal,
|
||||||
)
|
)
|
||||||
from app.domain.ports import LLMProvider, LLMProviderError, PdfTextExtractor
|
from app.domain.ports import (
|
||||||
|
LLMGenerationTimeout,
|
||||||
|
LLMProvider,
|
||||||
|
LLMProviderError,
|
||||||
|
PdfTextExtractor,
|
||||||
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -107,6 +112,84 @@ Format de réponse :
|
|||||||
- N'invente pas de contenu : tu réorganises et recopies ce qui est présent dans l'extrait.
|
- N'invente pas de contenu : tu réorganises et recopies ce qui est présent dans l'extrait.
|
||||||
- Si l'extrait ne contient aucune matière narrative, renvoie {{"arcs": []}}."""
|
- Si l'extrait ne contient aucune matière narrative, renvoie {{"arcs": []}}."""
|
||||||
|
|
||||||
|
# Schéma de l'arbre attendu, passé aux providers à sorties structurées (Ollama
|
||||||
|
# contraint la grammaire : un modèle local ne PEUT plus produire de clés
|
||||||
|
# inventées, d'objets bavards type "thought" ni de texte hors JSON). Les
|
||||||
|
# adapters cloud le traduisent en mode JSON natif. Seuls les "name" sont
|
||||||
|
# requis : le _TreeMerger tolère déjà tous les champs absents.
|
||||||
|
_TREE_SCHEMA: dict = {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"arcs": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"name": {"type": "string"},
|
||||||
|
"description": {"type": "string"},
|
||||||
|
"type": {"type": "string", "enum": ["LINEAR", "HUB"]},
|
||||||
|
"chapters": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"name": {"type": "string"},
|
||||||
|
"description": {"type": "string"},
|
||||||
|
"scenes": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"name": {"type": "string"},
|
||||||
|
"description": {"type": "string"},
|
||||||
|
"player_narration": {"type": "string"},
|
||||||
|
"gm_notes": {"type": "string"},
|
||||||
|
"rooms": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"name": {"type": "string"},
|
||||||
|
"description": {"type": "string"},
|
||||||
|
"enemies": {"type": "string"},
|
||||||
|
"loot": {"type": "string"},
|
||||||
|
},
|
||||||
|
"required": ["name"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"required": ["name"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"required": ["name"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"required": ["name"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"npcs": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"name": {"type": "string"},
|
||||||
|
"description": {"type": "string"},
|
||||||
|
},
|
||||||
|
"required": ["name"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"required": ["arcs"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
}
|
||||||
|
|
||||||
# Bloc TOC injecté quand le PDF a des bookmarks : les morceaux étant traités
|
# Bloc TOC injecté quand le PDF a des bookmarks : les morceaux étant traités
|
||||||
# séparément, c'est CE référentiel commun qui garantit que tous nomment les
|
# séparément, c'est CE référentiel commun qui garantit que tous nomment les
|
||||||
# mêmes chapitres à l'identique → la fusion par nom du _TreeMerger recolle
|
# mêmes chapitres à l'identique → la fusion par nom du _TreeMerger recolle
|
||||||
@@ -480,6 +563,17 @@ class ImportCampaignUseCase:
|
|||||||
f"Dernier message : {last_error or 'inconnu'}"}
|
f"Dernier message : {last_error or 'inconnu'}"}
|
||||||
return
|
return
|
||||||
|
|
||||||
|
if total > 0 and merger.counts()[0] == 0 and not merger.npcs():
|
||||||
|
# Le texte a été extrait mais le modèle n'a produit AUCUNE structure
|
||||||
|
# exploitable : sans ce signal, l'UI reçoit un `done` vide et
|
||||||
|
# l'utilisateur conclut à tort que le PDF est illisible.
|
||||||
|
yield {"type": "error",
|
||||||
|
"message": "Le texte du PDF a été extrait, mais le modèle n'a produit "
|
||||||
|
"aucune structure exploitable (réponses JSON vides ou coupées). "
|
||||||
|
"Réduisez la taille des morceaux d'import, augmentez la fenêtre "
|
||||||
|
"de contexte (num_ctx) ou essayez un autre modèle."}
|
||||||
|
return
|
||||||
|
|
||||||
# Consolidation finale : fusion des quasi-doublons inter-morceaux
|
# Consolidation finale : fusion des quasi-doublons inter-morceaux
|
||||||
# (best-effort, voir _consolidate). Inutile sur un import mono-morceau.
|
# (best-effort, voir _consolidate). Inutile sur un import mono-morceau.
|
||||||
if total > 1:
|
if total > 1:
|
||||||
@@ -557,8 +651,26 @@ class ImportCampaignUseCase:
|
|||||||
+ f"\n\n--- EXTRAIT {index + 1}/{total} ---\n{text}\n\n"
|
+ f"\n\n--- EXTRAIT {index + 1}/{total} ---\n{text}\n\n"
|
||||||
"Renvoie maintenant le JSON de l'arborescence."
|
"Renvoie maintenant le JSON de l'arborescence."
|
||||||
)
|
)
|
||||||
|
try:
|
||||||
raw = await generate_with_retry(
|
raw = await generate_with_retry(
|
||||||
self._llm, prompt, output_format="json", temperature=_TEMPERATURE)
|
self._llm, prompt, output_format=_TREE_SCHEMA, temperature=_TEMPERATURE)
|
||||||
|
except LLMGenerationTimeout:
|
||||||
|
# Génération trop lente pour la taille demandée (fréquent en local /
|
||||||
|
# tier gratuit) : même remède que la troncature, deux moitiés →
|
||||||
|
# sortie 2× plus courte. Re-lever si plus découpable.
|
||||||
|
if depth >= _MAX_SPLIT_DEPTH:
|
||||||
|
raise
|
||||||
|
left, right = split_in_half(text)
|
||||||
|
if not left or not right:
|
||||||
|
raise
|
||||||
|
logger.info(
|
||||||
|
"Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).",
|
||||||
|
index, depth + 1)
|
||||||
|
a = await self._extract_payload(
|
||||||
|
left, index=index, total=total, depth=depth + 1, toc_block=toc_block)
|
||||||
|
b = await self._extract_payload(
|
||||||
|
right, index=index, total=total, depth=depth + 1, toc_block=toc_block)
|
||||||
|
return {"arcs": a["arcs"] + b["arcs"], "npcs": a["npcs"] + b["npcs"]}
|
||||||
payload, truncated = self._parse_payload(raw, index=index)
|
payload, truncated = self._parse_payload(raw, index=index)
|
||||||
|
|
||||||
if truncated and depth < _MAX_SPLIT_DEPTH:
|
if truncated and depth < _MAX_SPLIT_DEPTH:
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ Ne dépend que des abstractions du domaine (ports LLMProvider + PdfTextExtractor
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
import re
|
||||||
|
|
||||||
from app.application.chunking import CHUNK_TARGET_TOKENS, chunk_text, split_in_half
|
from app.application.chunking import CHUNK_TARGET_TOKENS, chunk_text, split_in_half
|
||||||
from app.application.llm_json import load_json_object, looks_like_truncated_json
|
from app.application.llm_json import load_json_object, looks_like_truncated_json
|
||||||
@@ -25,7 +26,12 @@ from app.application.streaming import with_heartbeat
|
|||||||
# 1-2 niveaux suffisent en pratique, le reste est un garde-fou).
|
# 1-2 niveaux suffisent en pratique, le reste est un garde-fou).
|
||||||
_MAX_SPLIT_DEPTH = 3
|
_MAX_SPLIT_DEPTH = 3
|
||||||
from app.domain.models import RulesImportResult
|
from app.domain.models import RulesImportResult
|
||||||
from app.domain.ports import LLMProvider, LLMProviderError, PdfTextExtractor
|
from app.domain.ports import (
|
||||||
|
LLMGenerationTimeout,
|
||||||
|
LLMProvider,
|
||||||
|
LLMProviderError,
|
||||||
|
PdfTextExtractor,
|
||||||
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -34,6 +40,16 @@ logger = logging.getLogger(__name__)
|
|||||||
# Plus la valeur est haute, plus le modèle "brode" (invente du contenu absent).
|
# Plus la valeur est haute, plus le modèle "brode" (invente du contenu absent).
|
||||||
_TEMPERATURE = 0.1
|
_TEMPERATURE = 0.1
|
||||||
|
|
||||||
|
# Schéma de la sortie attendue : objet PLAT {titre: markdown}. Passé tel quel à
|
||||||
|
# Ollama (structured outputs : la grammaire interdit physiquement les objets
|
||||||
|
# imbriqués, les clés "thought" à valeur non-string, le bavardage hors JSON…
|
||||||
|
# indispensable pour les petits modèles locaux qui ne suivent pas les consignes).
|
||||||
|
# Les adapters cloud le traduisent en mode JSON natif (json_object).
|
||||||
|
_SECTIONS_SCHEMA: dict = {
|
||||||
|
"type": "object",
|
||||||
|
"additionalProperties": {"type": "string"},
|
||||||
|
}
|
||||||
|
|
||||||
# Taxonomie canonique suggérée au modèle pour homogénéiser les titres entre
|
# Taxonomie canonique suggérée au modèle pour homogénéiser les titres entre
|
||||||
# morceaux (sinon "Combat" / "Le combat" / "Règles de combat" se dispersent).
|
# morceaux (sinon "Combat" / "Le combat" / "Règles de combat" se dispersent).
|
||||||
# Le modèle reste libre d'en créer d'autres si rien ne correspond.
|
# Le modèle reste libre d'en créer d'autres si rien ne correspond.
|
||||||
@@ -56,9 +72,13 @@ On te donne un EXTRAIT brut d'un PDF de règles (texte parfois mal coupé par la
|
|||||||
|
|
||||||
Ta tâche : répartir le contenu de cet extrait dans des SECTIONS THÉMATIQUES.
|
Ta tâche : répartir le contenu de cet extrait dans des SECTIONS THÉMATIQUES.
|
||||||
|
|
||||||
|
Format EXACT attendu — un objet JSON plat {{titre de section: contenu markdown}} :
|
||||||
|
{{"Combat": "## Initiative\\n\\nChaque participant lance 1d20...", "Magie et sorts": "## Sorts\\n\\n..."}}
|
||||||
|
|
||||||
Règles impératives :
|
Règles impératives :
|
||||||
- Tu réponds UNIQUEMENT par un objet JSON valide, sans markdown ni commentaire autour.
|
- Tu réponds UNIQUEMENT par cet objet JSON, sans texte avant ni après.
|
||||||
- Les CLÉS sont des titres de section (texte court). Les VALEURS sont le contenu de la règle en markdown.
|
- Les CLÉS sont des titres de section (texte court). Les VALEURS sont le contenu de la règle en markdown (chaîne de caractères, jamais un objet ou une liste).
|
||||||
|
- INTERDIT : des clés génériques comme "title", "content", "sections", "thought" ou "notes" ; des objets imbriqués ; tout commentaire sur ta démarche ou ton raisonnement.
|
||||||
- Utilise EN PRIORITÉ ces titres canoniques quand le contenu y correspond :
|
- Utilise EN PRIORITÉ ces titres canoniques quand le contenu y correspond :
|
||||||
{canonical}
|
{canonical}
|
||||||
- Si un contenu ne rentre dans aucun, crée un titre clair et concis (en français).
|
- Si un contenu ne rentre dans aucun, crée un titre clair et concis (en français).
|
||||||
@@ -67,6 +87,55 @@ Règles impératives :
|
|||||||
- N'INVENTE AUCUNE règle, ne résume pas abusivement : tu réorganises, tu ne réécris pas le fond.
|
- N'INVENTE AUCUNE règle, ne résume pas abusivement : tu réorganises, tu ne réécris pas le fond.
|
||||||
- Ignore les pages de garde, sommaires, crédits, pages vides (renvoie {{}} si l'extrait n'a aucune règle)."""
|
- Ignore les pages de garde, sommaires, crédits, pages vides (renvoie {{}} si l'extrait n'a aucune règle)."""
|
||||||
|
|
||||||
|
# --- Mode SEGMENTATION (modèles locaux) --------------------------------------
|
||||||
|
# Réécrire tout le texte en JSON impose une SORTIE ≈ taille de l'ENTRÉE : à
|
||||||
|
# ~100 tokens/s en local, un livre = des dizaines de minutes et des troncatures
|
||||||
|
# en cascade. Ici le modèle ne renvoie que les FRONTIÈRES des sections (titre +
|
||||||
|
# premiers mots exacts) — ~200 tokens quel que soit le morceau — et c'est NOUS
|
||||||
|
# qui découpons le texte original. ~50× plus rapide, fidélité parfaite du
|
||||||
|
# contenu (texte source intact), plus de troncature possible.
|
||||||
|
|
||||||
|
_SEGMENT_SYSTEM = """Tu analyses un EXTRAIT brut d'un livre de règles de jeu de rôle.
|
||||||
|
Ta tâche : repérer où COMMENCENT les sections thématiques. Tu ne réécris RIEN.
|
||||||
|
|
||||||
|
Format EXACT attendu :
|
||||||
|
{{"sections": [{{"titre": "Combat", "debut": "Le combat se déroule en tours de"}}, ...]}}
|
||||||
|
|
||||||
|
Règles impératives :
|
||||||
|
- "debut" = les 5 à 10 PREMIERS MOTS du passage où la section commence, COPIÉS À L'IDENTIQUE
|
||||||
|
depuis l'extrait (même orthographe, même ponctuation, même langue). JAMAIS un résumé.
|
||||||
|
- La PREMIÈRE entrée commence aux tout premiers mots de l'extrait (même si le contenu
|
||||||
|
poursuit une section entamée avant cet extrait).
|
||||||
|
- Les entrées suivent l'ordre du texte. Vise des sections LARGES (un thème), pas un titre
|
||||||
|
par paragraphe : un extrait contient typiquement 1 à 6 sections.
|
||||||
|
- Titres : EN PRIORITÉ parmi :
|
||||||
|
{canonical}
|
||||||
|
sinon un titre court et clair en français.
|
||||||
|
- Pages de garde, sommaires, crédits : n'en fais pas des sections. Si l'extrait n'est que ça,
|
||||||
|
renvoie {{"sections": []}}."""
|
||||||
|
|
||||||
|
# Schéma passé à Ollama (structured outputs) : un objet {"sections": [...]}.
|
||||||
|
# Racine objet (pas tableau) car l'extraction côté Brain repère le premier {…}.
|
||||||
|
_ANCHORS_SCHEMA: dict = {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"sections": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"titre": {"type": "string"},
|
||||||
|
"debut": {"type": "string"},
|
||||||
|
},
|
||||||
|
"required": ["titre", "debut"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"required": ["sections"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
class _SectionMerger:
|
class _SectionMerger:
|
||||||
"""Fusionne les sections issues des différents morceaux, ordre préservé.
|
"""Fusionne les sections issues des différents morceaux, ordre préservé.
|
||||||
@@ -102,6 +171,84 @@ class _SectionMerger:
|
|||||||
return {title: "\n\n".join(parts) for title, parts in self._merged.items()}
|
return {title: "\n\n".join(parts) for title, parts in self._merged.items()}
|
||||||
|
|
||||||
|
|
||||||
|
# Clés "méta" que certains modèles glissent dans le JSON (fuite de raisonnement,
|
||||||
|
# schéma title/content inventé…) : jamais des titres de section voulus.
|
||||||
|
_META_KEYS = frozenset({
|
||||||
|
"thought", "thoughts", "thinking", "reasoning", "raisonnement",
|
||||||
|
"comment", "commentaire", "commentaires", "note", "notes", "explanation",
|
||||||
|
})
|
||||||
|
|
||||||
|
|
||||||
|
def _normalize_sections(parsed: dict) -> dict:
|
||||||
|
"""Ramène les formes déviantes courantes au format attendu {titre: contenu}.
|
||||||
|
|
||||||
|
Observé sur les petits modèles locaux (gemma 12b) malgré les consignes :
|
||||||
|
- enveloppe {"sections": {...}} ou {"règles": {...}} autour du vrai contenu ;
|
||||||
|
- schéma inventé {"title": "...", "content": "...", "thought": "..."} →
|
||||||
|
une seule section dont le titre est la valeur de "title" ;
|
||||||
|
- clés méta ("thought", "notes"…) mêlées aux vraies sections → retirées.
|
||||||
|
"""
|
||||||
|
by_lower = {str(k).strip().lower(): k for k in parsed}
|
||||||
|
# Enveloppe : un unique conteneur connu dont la valeur est l'objet attendu.
|
||||||
|
if len(parsed) == 1:
|
||||||
|
only_key, only_val = next(iter(parsed.items()))
|
||||||
|
if (isinstance(only_val, dict)
|
||||||
|
and str(only_key).strip().lower() in {"sections", "règles", "regles", "rules"}):
|
||||||
|
return _normalize_sections(only_val)
|
||||||
|
# Schéma {"title": ..., "content": ...} : le titre est une VALEUR, pas une clé.
|
||||||
|
if "title" in by_lower and "content" in by_lower:
|
||||||
|
title = str(parsed[by_lower["title"]]).strip()
|
||||||
|
content = parsed[by_lower["content"]]
|
||||||
|
if title and not isinstance(content, dict):
|
||||||
|
return {title: content}
|
||||||
|
return {k: v for k, v in parsed.items()
|
||||||
|
if str(k).strip().lower() not in _META_KEYS}
|
||||||
|
|
||||||
|
|
||||||
|
def _coerce_markdown(value: object) -> str:
|
||||||
|
"""Convertit une valeur de section renvoyée par le LLM en markdown plat.
|
||||||
|
|
||||||
|
Malgré la consigne « valeurs = markdown », certains modèles nichent des
|
||||||
|
sous-sections ({titre: {sous-titre: contenu}}) ou des listes. Un `str(v)`
|
||||||
|
naïf produirait du repr Python ({'k': 'v'}) ; on aplatit récursivement à la
|
||||||
|
place pour ne perdre aucun contenu.
|
||||||
|
"""
|
||||||
|
if isinstance(value, str):
|
||||||
|
return value
|
||||||
|
if isinstance(value, dict):
|
||||||
|
parts = []
|
||||||
|
for k, v in value.items():
|
||||||
|
content = _coerce_markdown(v)
|
||||||
|
# Clé = sous-titre (cas normal) ; si la "valeur" est vide, la clé
|
||||||
|
# elle-même porte le contenu (dérive observée sur certains modèles).
|
||||||
|
parts.append(f"{k}\n\n{content}".strip() if content else str(k))
|
||||||
|
return "\n\n".join(parts)
|
||||||
|
if isinstance(value, list):
|
||||||
|
return "\n\n".join(_coerce_markdown(v) for v in value)
|
||||||
|
return "" if value is None else str(value)
|
||||||
|
|
||||||
|
|
||||||
|
def _find_anchor(text: str, anchor: str, start: int) -> int | None:
|
||||||
|
"""Position de `anchor` dans `text` à partir de `start`, ou None.
|
||||||
|
|
||||||
|
Le modèle recopie les premiers mots d'un passage, mais le texte extrait du
|
||||||
|
PDF contient des sauts de ligne/espaces multiples au même endroit, et le
|
||||||
|
modèle normalise parfois la casse. Trois passes, de la plus stricte à la
|
||||||
|
plus tolérante : exacte → espaces≈\\s+ → idem insensible à la casse."""
|
||||||
|
pos = text.find(anchor, start)
|
||||||
|
if pos != -1:
|
||||||
|
return pos
|
||||||
|
words = anchor.split()
|
||||||
|
if not words:
|
||||||
|
return None
|
||||||
|
pattern = r"\s+".join(re.escape(w) for w in words)
|
||||||
|
match = re.compile(pattern).search(text, start)
|
||||||
|
if match:
|
||||||
|
return match.start()
|
||||||
|
match = re.compile(pattern, re.IGNORECASE).search(text, start)
|
||||||
|
return match.start() if match else None
|
||||||
|
|
||||||
|
|
||||||
def _combine_sections(a: dict[str, str], b: dict[str, str]) -> dict[str, str]:
|
def _combine_sections(a: dict[str, str], b: dict[str, str]) -> dict[str, str]:
|
||||||
"""Fusionne deux dicts de sections (issus des 2 moitiés d'un morceau re-découpé).
|
"""Fusionne deux dicts de sections (issus des 2 moitiés d'un morceau re-découpé).
|
||||||
|
|
||||||
@@ -128,10 +275,16 @@ class ImportRulesUseCase:
|
|||||||
llm: LLMProvider,
|
llm: LLMProvider,
|
||||||
extractor: PdfTextExtractor,
|
extractor: PdfTextExtractor,
|
||||||
chunk_target_tokens: int = CHUNK_TARGET_TOKENS,
|
chunk_target_tokens: int = CHUNK_TARGET_TOKENS,
|
||||||
|
segment_only: bool = False,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
"""`segment_only=True` (modèles locaux) : le LLM ne renvoie que les
|
||||||
|
frontières des sections (titre + premiers mots) et le texte original est
|
||||||
|
découpé localement — sortie minuscule, pas de réécriture. False (cloud) :
|
||||||
|
le LLM réécrit le contenu en sections markdown nettoyées."""
|
||||||
self._llm = llm
|
self._llm = llm
|
||||||
self._extractor = extractor
|
self._extractor = extractor
|
||||||
self._chunk_target_tokens = chunk_target_tokens
|
self._chunk_target_tokens = chunk_target_tokens
|
||||||
|
self._segment_only = segment_only
|
||||||
|
|
||||||
async def execute(self, pdf_bytes: bytes) -> RulesImportResult:
|
async def execute(self, pdf_bytes: bytes) -> RulesImportResult:
|
||||||
"""Variante non-streamée : traite tout puis renvoie le résultat complet."""
|
"""Variante non-streamée : traite tout puis renvoie le résultat complet."""
|
||||||
@@ -215,9 +368,21 @@ class ImportRulesUseCase:
|
|||||||
f"Dernier message : {last_error or 'inconnu'}"}
|
f"Dernier message : {last_error or 'inconnu'}"}
|
||||||
return
|
return
|
||||||
|
|
||||||
|
sections = merger.result()
|
||||||
|
if total > 0 and not sections:
|
||||||
|
# Le texte a bien été extrait mais AUCUN morceau n'a produit de JSON
|
||||||
|
# exploitable (sorties coupées/illisibles). Sans ce signal, l'UI reçoit
|
||||||
|
# un `done` vide et l'utilisateur conclut à tort que le PDF est illisible.
|
||||||
|
yield {"type": "error",
|
||||||
|
"message": "Le texte du PDF a été extrait, mais le modèle n'a produit "
|
||||||
|
"aucune section exploitable (réponses JSON vides ou coupées). "
|
||||||
|
"Réduisez la taille des morceaux d'import, augmentez la fenêtre "
|
||||||
|
"de contexte (num_ctx) ou essayez un autre modèle."}
|
||||||
|
return
|
||||||
|
|
||||||
yield {
|
yield {
|
||||||
"type": "done",
|
"type": "done",
|
||||||
"sections": merger.result(),
|
"sections": sections,
|
||||||
"page_count": doc.page_count,
|
"page_count": doc.page_count,
|
||||||
"ocr_page_count": doc.ocr_page_count,
|
"ocr_page_count": doc.ocr_page_count,
|
||||||
"skipped": skipped,
|
"skipped": skipped,
|
||||||
@@ -234,15 +399,36 @@ class ImportRulesUseCase:
|
|||||||
"""Extrait les sections d'un texte. Si la SORTIE est tronquée, retraite le
|
"""Extrait les sections d'un texte. Si la SORTIE est tronquée, retraite le
|
||||||
texte en DEUX moitiés (chacune produit une réponse complète) et fusionne —
|
texte en DEUX moitiés (chacune produit une réponse complète) et fusionne —
|
||||||
ainsi aucune section n'est perdue, quel que soit le plafond de sortie."""
|
ainsi aucune section n'est perdue, quel que soit le plafond de sortie."""
|
||||||
|
system = _SEGMENT_SYSTEM if self._segment_only else _MAP_SYSTEM
|
||||||
|
schema = _ANCHORS_SCHEMA if self._segment_only else _SECTIONS_SCHEMA
|
||||||
prompt = (
|
prompt = (
|
||||||
_MAP_SYSTEM.format(
|
system.format(
|
||||||
canonical="\n".join(f" - {s}" for s in _CANONICAL_SECTIONS)
|
canonical="\n".join(f" - {s}" for s in _CANONICAL_SECTIONS)
|
||||||
)
|
)
|
||||||
+ f"\n\n--- EXTRAIT {index + 1}/{total} ---\n{text}\n\n"
|
+ f"\n\n--- EXTRAIT {index + 1}/{total} ---\n{text}\n\n"
|
||||||
"Renvoie maintenant le JSON des sections."
|
"Renvoie maintenant le JSON des sections."
|
||||||
)
|
)
|
||||||
|
try:
|
||||||
raw = await generate_with_retry(
|
raw = await generate_with_retry(
|
||||||
self._llm, prompt, output_format="json", temperature=_TEMPERATURE)
|
self._llm, prompt, output_format=schema, temperature=_TEMPERATURE)
|
||||||
|
except LLMGenerationTimeout:
|
||||||
|
# Le modèle générait mais trop lentement pour réécrire tout le morceau
|
||||||
|
# dans le temps imparti (fréquent sur tier gratuit + gros morceaux).
|
||||||
|
# Même remède que la troncature : deux moitiés → sortie 2× plus courte.
|
||||||
|
if depth >= _MAX_SPLIT_DEPTH:
|
||||||
|
raise
|
||||||
|
left, right = split_in_half(text)
|
||||||
|
if not left or not right:
|
||||||
|
raise
|
||||||
|
logger.info(
|
||||||
|
"Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).",
|
||||||
|
index, depth + 1)
|
||||||
|
a = await self._extract_sections(left, index=index, total=total, depth=depth + 1)
|
||||||
|
b = await self._extract_sections(right, index=index, total=total, depth=depth + 1)
|
||||||
|
return _combine_sections(a, b)
|
||||||
|
if self._segment_only:
|
||||||
|
sections, truncated = self._parse_anchors(raw, text, index=index)
|
||||||
|
else:
|
||||||
sections, truncated = self._parse_sections(raw, index=index)
|
sections, truncated = self._parse_sections(raw, index=index)
|
||||||
|
|
||||||
if truncated and depth < _MAX_SPLIT_DEPTH:
|
if truncated and depth < _MAX_SPLIT_DEPTH:
|
||||||
@@ -259,6 +445,68 @@ class ImportRulesUseCase:
|
|||||||
"Morceau %s : sortie tronquée, profondeur max atteinte — partiel conservé.", index)
|
"Morceau %s : sortie tronquée, profondeur max atteinte — partiel conservé.", index)
|
||||||
return sections
|
return sections
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _parse_anchors(raw: str, text: str, *, index: int) -> tuple[dict[str, str], bool]:
|
||||||
|
"""Mode segmentation : réponse {"sections": [{titre, debut}, …]} → on localise
|
||||||
|
chaque `debut` dans le texte ORIGINAL et on découpe entre les ancres.
|
||||||
|
|
||||||
|
Une ancre introuvable est abandonnée (son contenu reste dans la section
|
||||||
|
précédente — aucun texte n'est perdu). Le texte avant la première ancre
|
||||||
|
trouvée est rattaché à la première section (le prompt demande au modèle de
|
||||||
|
faire démarrer la première entrée aux premiers mots de l'extrait)."""
|
||||||
|
parsed, recovered = load_json_object(raw)
|
||||||
|
if parsed is None:
|
||||||
|
truncated = looks_like_truncated_json(raw)
|
||||||
|
if not truncated:
|
||||||
|
logger.warning(
|
||||||
|
"Morceau %s : aucun objet JSON exploitable (segmentation), ignoré. "
|
||||||
|
"Début de la réponse du modèle : %r",
|
||||||
|
index, (raw or "").strip()[:300] or "(réponse VIDE)")
|
||||||
|
return {}, truncated
|
||||||
|
entries = parsed.get("sections") if isinstance(parsed, dict) else None
|
||||||
|
if not isinstance(entries, list):
|
||||||
|
logger.warning("Morceau %s : pas de liste 'sections' exploitable, ignoré.", index)
|
||||||
|
return {}, False
|
||||||
|
|
||||||
|
# Localisation séquentielle : chaque ancre est cherchée APRÈS la précédente
|
||||||
|
# (préserve l'ordre du texte, évite qu'une phrase répétée matche trop tôt).
|
||||||
|
located: list[tuple[str, int]] = []
|
||||||
|
cursor = 0
|
||||||
|
dropped = 0
|
||||||
|
for entry in entries:
|
||||||
|
if not isinstance(entry, dict):
|
||||||
|
continue
|
||||||
|
title = str(entry.get("titre") or "").strip()
|
||||||
|
anchor = str(entry.get("debut") or "").strip()
|
||||||
|
if not title or not anchor:
|
||||||
|
continue
|
||||||
|
pos = _find_anchor(text, anchor, cursor)
|
||||||
|
if pos is None:
|
||||||
|
dropped += 1
|
||||||
|
continue
|
||||||
|
located.append((title, pos))
|
||||||
|
cursor = pos + 1
|
||||||
|
if dropped:
|
||||||
|
logger.info(
|
||||||
|
"Morceau %s : %s ancre(s) de section introuvable(s) — contenu rattaché "
|
||||||
|
"à la section précédente.", index, dropped)
|
||||||
|
if not located:
|
||||||
|
return {}, False
|
||||||
|
|
||||||
|
# Découpe entre ancres ; le préambule éventuel rejoint la première section.
|
||||||
|
located[0] = (located[0][0], 0)
|
||||||
|
sections: dict[str, str] = {}
|
||||||
|
for i, (title, start) in enumerate(located):
|
||||||
|
end = located[i + 1][1] if i + 1 < len(located) else len(text)
|
||||||
|
content = text[start:end].strip()
|
||||||
|
if not content:
|
||||||
|
continue
|
||||||
|
if title in sections:
|
||||||
|
sections[title] = f"{sections[title]}\n\n{content}"
|
||||||
|
else:
|
||||||
|
sections[title] = content
|
||||||
|
return sections, recovered
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _parse_sections(raw: str, *, index: int) -> tuple[dict[str, str], bool]:
|
def _parse_sections(raw: str, *, index: int) -> tuple[dict[str, str], bool]:
|
||||||
"""Parse robuste → (sections, tronqué). `tronqué`=True si récupération partielle."""
|
"""Parse robuste → (sections, tronqué). `tronqué`=True si récupération partielle."""
|
||||||
@@ -276,4 +524,5 @@ class ImportRulesUseCase:
|
|||||||
if not isinstance(parsed, dict):
|
if not isinstance(parsed, dict):
|
||||||
logger.warning("Morceau %s : le LLM n'a pas renvoyé un objet, ignoré.", index)
|
logger.warning("Morceau %s : le LLM n'a pas renvoyé un objet, ignoré.", index)
|
||||||
return {}, False
|
return {}, False
|
||||||
return {str(k): str(v) for k, v in parsed.items()}, recovered
|
normalized = _normalize_sections(parsed)
|
||||||
|
return {str(k): _coerce_markdown(v) for k, v in normalized.items()}, recovered
|
||||||
|
|||||||
@@ -38,13 +38,16 @@ def load_json_object(raw: str) -> tuple[object | None, bool]:
|
|||||||
obj = extract_json_object(raw)
|
obj = extract_json_object(raw)
|
||||||
if obj is not None:
|
if obj is not None:
|
||||||
try:
|
try:
|
||||||
return json.loads(obj), False
|
# strict=False : tolère les caractères de contrôle BRUTS (retours à la
|
||||||
|
# ligne non échappés…) dans les chaînes — erreur fréquente des LLM hors
|
||||||
|
# mode JSON natif, qui invalidait toute la réponse.
|
||||||
|
return json.loads(obj, strict=False), False
|
||||||
except json.JSONDecodeError:
|
except json.JSONDecodeError:
|
||||||
pass
|
pass
|
||||||
repaired = repair_truncated_json(raw)
|
repaired = repair_truncated_json(raw)
|
||||||
if repaired is not None:
|
if repaired is not None:
|
||||||
try:
|
try:
|
||||||
return json.loads(repaired), True
|
return json.loads(repaired, strict=False), True
|
||||||
except json.JSONDecodeError:
|
except json.JSONDecodeError:
|
||||||
pass
|
pass
|
||||||
return None, False
|
return None, False
|
||||||
@@ -54,12 +57,20 @@ def looks_like_truncated_json(raw: str) -> bool:
|
|||||||
"""La sortie ressemble-t-elle à un JSON COUPÉ (accolades/crochets non refermés)
|
"""La sortie ressemble-t-elle à un JSON COUPÉ (accolades/crochets non refermés)
|
||||||
plutôt qu'à de la prose ? Sert à déclencher un re-découpage même quand RIEN n'a
|
plutôt qu'à de la prose ? Sert à déclencher un re-découpage même quand RIEN n'a
|
||||||
pu être récupéré (cas où le 1er contenu est si long qu'il est coupé avant toute
|
pu être récupéré (cas où le 1er contenu est si long qu'il est coupé avant toute
|
||||||
sous-structure complète). On exige un contenu substantiel pour éviter les
|
sous-structure complète).
|
||||||
faux positifs sur une courte réponse non-JSON."""
|
|
||||||
s = (raw or "").strip()
|
Une réponse qui COMMENCE par `{` est jugée sur le seul équilibre des accolades,
|
||||||
if "{" not in s or len(s) < 100:
|
même très courte : en mode JSON un `{"` de 2 caractères est une génération
|
||||||
|
interrompue net (contexte plein, plafond de sortie), pas de la prose — c'est le
|
||||||
|
signal de re-découpage. Pour le reste (prose contenant des accolades), on exige
|
||||||
|
un contenu substantiel pour éviter les faux positifs."""
|
||||||
|
s = _strip_reasoning(raw or "").strip()
|
||||||
|
if "{" not in s:
|
||||||
return False
|
return False
|
||||||
return s.count("{") > s.count("}") or s.count("[") > s.count("]")
|
unbalanced = s.count("{") > s.count("}") or s.count("[") > s.count("]")
|
||||||
|
if s.startswith("{"):
|
||||||
|
return unbalanced
|
||||||
|
return len(s) >= 100 and unbalanced
|
||||||
|
|
||||||
|
|
||||||
def extract_json_object(raw: str) -> str | None:
|
def extract_json_object(raw: str) -> str | None:
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ import asyncio
|
|||||||
import logging
|
import logging
|
||||||
import re
|
import re
|
||||||
|
|
||||||
from app.domain.ports import LLMProvider, LLMProviderError
|
from app.domain.ports import LLMGenerationTimeout, LLMProvider, LLMProviderError
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -60,7 +60,7 @@ async def generate_with_retry(
|
|||||||
llm: LLMProvider,
|
llm: LLMProvider,
|
||||||
prompt: str,
|
prompt: str,
|
||||||
*,
|
*,
|
||||||
output_format: str | None = None,
|
output_format: str | dict | None = None,
|
||||||
temperature: float | None = None,
|
temperature: float | None = None,
|
||||||
) -> str:
|
) -> str:
|
||||||
"""Comme `llm.generate`, mais réessaie les erreurs transitoires (backoff).
|
"""Comme `llm.generate`, mais réessaie les erreurs transitoires (backoff).
|
||||||
@@ -74,6 +74,12 @@ async def generate_with_retry(
|
|||||||
for attempt in range(_ATTEMPTS):
|
for attempt in range(_ATTEMPTS):
|
||||||
try:
|
try:
|
||||||
return await llm.generate(prompt, output_format=output_format, temperature=temperature)
|
return await llm.generate(prompt, output_format=output_format, temperature=temperature)
|
||||||
|
except LLMGenerationTimeout:
|
||||||
|
# Timeout de DÉBIT (génération trop lente pour la sortie demandée) :
|
||||||
|
# rejouer le même prompt re-timeoutera à l'identique — on a déjà perdu
|
||||||
|
# `timeout` secondes. On remonte tout de suite : l'appelant (import)
|
||||||
|
# sait re-découper le morceau en deux pour réduire la sortie.
|
||||||
|
raise
|
||||||
except LLMProviderError as exc:
|
except LLMProviderError as exc:
|
||||||
last_error = exc
|
last_error = exc
|
||||||
# Quota JOURNALIER épuisé : inutile d'insister, on remonte tout de suite
|
# Quota JOURNALIER épuisé : inutile d'insister, on remonte tout de suite
|
||||||
|
|||||||
@@ -24,17 +24,20 @@ class LLMProvider(Protocol):
|
|||||||
self,
|
self,
|
||||||
prompt: str,
|
prompt: str,
|
||||||
*,
|
*,
|
||||||
output_format: str | None = None,
|
output_format: str | dict | None = None,
|
||||||
temperature: float | None = None,
|
temperature: float | None = None,
|
||||||
) -> str:
|
) -> str:
|
||||||
"""Génère une réponse textuelle à partir d'un prompt donné.
|
"""Génère une réponse textuelle à partir d'un prompt donné.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
prompt: le texte envoyé au modèle.
|
prompt: le texte envoyé au modèle.
|
||||||
output_format: contrainte de format optionnelle. Exemple : "json"
|
output_format: contrainte de format optionnelle. "json" pour forcer
|
||||||
pour forcer le modèle à renvoyer du JSON valide. Les
|
un JSON valide ; un dict = SCHÉMA JSON décrivant la structure
|
||||||
fournisseurs qui ne supportent pas une valeur donnée doivent
|
attendue (les fournisseurs qui supportent les sorties
|
||||||
l'ignorer silencieusement ou la traduire au mieux.
|
structurées — ex. Ollama — contraignent la génération au schéma,
|
||||||
|
les autres retombent sur leur mode JSON natif). Les fournisseurs
|
||||||
|
qui ne supportent pas une valeur donnée doivent l'ignorer
|
||||||
|
silencieusement ou la traduire au mieux.
|
||||||
temperature: créativité du modèle, 0.0 (déterministe/factuel) à
|
temperature: créativité du modèle, 0.0 (déterministe/factuel) à
|
||||||
1.0+ (très créatif, hallucine plus facilement). None =
|
1.0+ (très créatif, hallucine plus facilement). None =
|
||||||
valeur par défaut de l'adapter. Recommandation LoreMind :
|
valeur par défaut de l'adapter. Recommandation LoreMind :
|
||||||
@@ -113,3 +116,14 @@ class LLMProviderError(Exception):
|
|||||||
Définie dans le domaine (pas dans l'infra) pour que les couches
|
Définie dans le domaine (pas dans l'infra) pour que les couches
|
||||||
supérieures puissent l'attraper sans connaître l'adapter concret.
|
supérieures puissent l'attraper sans connaître l'adapter concret.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
class LLMGenerationTimeout(LLMProviderError):
|
||||||
|
"""La génération a démarré mais n'a pas FINI dans le temps imparti.
|
||||||
|
|
||||||
|
Cas distinct d'un échec transitoire (file d'attente, 503) : le modèle
|
||||||
|
produisait des tokens mais trop lentement pour la taille de sortie demandée.
|
||||||
|
Réessayer à l'identique est inutile (même entrée → même lenteur) ; la bonne
|
||||||
|
réaction est de RÉDUIRE la sortie demandée (ex. import : re-découper le
|
||||||
|
morceau en deux moitiés).
|
||||||
|
"""
|
||||||
|
|||||||
@@ -20,7 +20,7 @@ import httpx
|
|||||||
|
|
||||||
from app.core.config import Settings
|
from app.core.config import Settings
|
||||||
from app.domain.models import ChatMessage
|
from app.domain.models import ChatMessage
|
||||||
from app.domain.ports import LLMProviderError
|
from app.domain.ports import LLMGenerationTimeout, LLMProviderError
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -95,7 +95,7 @@ class GeminiLLMProvider:
|
|||||||
try:
|
try:
|
||||||
return await asyncio.wait_for(_collect(), timeout=self._timeout)
|
return await asyncio.wait_for(_collect(), timeout=self._timeout)
|
||||||
except asyncio.TimeoutError as exc:
|
except asyncio.TimeoutError as exc:
|
||||||
raise LLMProviderError(
|
raise LLMGenerationTimeout(
|
||||||
f"Erreur Gemini : génération non terminée en {self._timeout}s. Réduisez la "
|
f"Erreur Gemini : génération non terminée en {self._timeout}s. Réduisez la "
|
||||||
"taille des morceaux d'import ou augmentez le timeout."
|
"taille des morceaux d'import ou augmentez le timeout."
|
||||||
) from exc
|
) from exc
|
||||||
@@ -130,6 +130,12 @@ class GeminiLLMProvider:
|
|||||||
}
|
}
|
||||||
if temperature is not None:
|
if temperature is not None:
|
||||||
body["temperature"] = temperature
|
body["temperature"] = temperature
|
||||||
|
# Mode JSON natif (supporté par l'endpoint OpenAI-compatible de Gemini) :
|
||||||
|
# supprime fences ```json et JSON invalide, principale cause de morceaux
|
||||||
|
# ignorés. Un SCHÉMA (dict) est traduit en json_object : suffisant, les
|
||||||
|
# grands modèles cloud respectent la structure demandée par le prompt.
|
||||||
|
if output_format is not None:
|
||||||
|
body["response_format"] = {"type": "json_object"}
|
||||||
|
|
||||||
async with httpx.AsyncClient(timeout=self._timeout) as client:
|
async with httpx.AsyncClient(timeout=self._timeout) as client:
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -20,7 +20,7 @@ import httpx
|
|||||||
|
|
||||||
from app.core.config import Settings
|
from app.core.config import Settings
|
||||||
from app.domain.models import ChatMessage
|
from app.domain.models import ChatMessage
|
||||||
from app.domain.ports import LLMProviderError
|
from app.domain.ports import LLMGenerationTimeout, LLMProviderError
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -101,7 +101,7 @@ class MistralLLMProvider:
|
|||||||
try:
|
try:
|
||||||
return await asyncio.wait_for(_collect(), timeout=self._timeout)
|
return await asyncio.wait_for(_collect(), timeout=self._timeout)
|
||||||
except asyncio.TimeoutError as exc:
|
except asyncio.TimeoutError as exc:
|
||||||
raise LLMProviderError(
|
raise LLMGenerationTimeout(
|
||||||
f"Erreur Mistral : génération non terminée en {self._timeout}s. Réduisez la "
|
f"Erreur Mistral : génération non terminée en {self._timeout}s. Réduisez la "
|
||||||
"taille des morceaux d'import, augmentez le timeout, ou changez de modèle."
|
"taille des morceaux d'import, augmentez le timeout, ou changez de modèle."
|
||||||
) from exc
|
) from exc
|
||||||
@@ -136,6 +136,13 @@ class MistralLLMProvider:
|
|||||||
}
|
}
|
||||||
if temperature is not None:
|
if temperature is not None:
|
||||||
body["temperature"] = temperature
|
body["temperature"] = temperature
|
||||||
|
# Mode JSON natif : TOUS les modèles Mistral le supportent → plus de fences
|
||||||
|
# ```json ni de JSON invalide (retours à la ligne bruts dans les chaînes),
|
||||||
|
# principale cause de morceaux d'import ignorés. Un SCHÉMA (dict) est
|
||||||
|
# traduit en json_object : suffisant ici, les grands modèles cloud
|
||||||
|
# respectent la structure demandée par le prompt.
|
||||||
|
if output_format is not None:
|
||||||
|
body["response_format"] = {"type": "json_object"}
|
||||||
|
|
||||||
async with httpx.AsyncClient(timeout=self._timeout) as client:
|
async with httpx.AsyncClient(timeout=self._timeout) as client:
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -5,13 +5,16 @@ Isole le reste de l'application des spécificités du protocole Ollama
|
|||||||
demain, on écrit un nouvel adapter sans toucher au reste du code.
|
demain, on écrit un nouvel adapter sans toucher au reste du code.
|
||||||
"""
|
"""
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
from typing import AsyncIterator
|
from typing import AsyncIterator
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
|
|
||||||
from app.core.config import Settings
|
from app.core.config import Settings
|
||||||
from app.domain.models import ChatMessage
|
from app.domain.models import ChatMessage
|
||||||
from app.domain.ports import LLMProviderError
|
from app.domain.ports import LLMGenerationTimeout, LLMProviderError
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class OllamaLLMProvider:
|
class OllamaLLMProvider:
|
||||||
@@ -45,7 +48,7 @@ class OllamaLLMProvider:
|
|||||||
self,
|
self,
|
||||||
prompt: str,
|
prompt: str,
|
||||||
*,
|
*,
|
||||||
output_format: str | None = None,
|
output_format: str | dict | None = None,
|
||||||
temperature: float | None = None,
|
temperature: float | None = None,
|
||||||
) -> str:
|
) -> str:
|
||||||
url = f"{self._base_url}/api/generate"
|
url = f"{self._base_url}/api/generate"
|
||||||
@@ -55,6 +58,10 @@ class OllamaLLMProvider:
|
|||||||
"stream": False,
|
"stream": False,
|
||||||
"options": self._build_options(temperature),
|
"options": self._build_options(temperature),
|
||||||
}
|
}
|
||||||
|
# "json" (mode JSON simple) ou un SCHÉMA JSON complet (structured outputs) :
|
||||||
|
# Ollama contraint alors la grammaire de génération au schéma — un petit
|
||||||
|
# modèle local ne PEUT physiquement plus produire d'objets imbriqués, de
|
||||||
|
# clés "thought" bavardes ou de texte hors JSON.
|
||||||
if output_format is not None:
|
if output_format is not None:
|
||||||
payload["format"] = output_format
|
payload["format"] = output_format
|
||||||
|
|
||||||
@@ -71,12 +78,45 @@ class OllamaLLMProvider:
|
|||||||
raise LLMProviderError(
|
raise LLMProviderError(
|
||||||
f"Ollama HTTP {response.status_code} : {err_msg.strip()[:500]}"
|
f"Ollama HTTP {response.status_code} : {err_msg.strip()[:500]}"
|
||||||
)
|
)
|
||||||
|
except httpx.ConnectTimeout as exc:
|
||||||
|
# Serveur injoignable : erreur d'infrastructure, pas de lenteur.
|
||||||
|
raise LLMProviderError(
|
||||||
|
f"Erreur lors de l'appel à Ollama : {exc}"
|
||||||
|
) from exc
|
||||||
|
except httpx.TimeoutException as exc:
|
||||||
|
# `stream: False` → le read-timeout court jusqu'à la réponse COMPLÈTE,
|
||||||
|
# donc le dépasser = génération trop lente pour la sortie demandée
|
||||||
|
# (fréquent : modèle local modeste + gros morceau d'import à réécrire).
|
||||||
|
# Type dédié → pas de retry à l'identique ; l'import re-découpe le
|
||||||
|
# morceau en deux moitiés (sortie 2× plus courte) à la place.
|
||||||
|
raise LLMGenerationTimeout(
|
||||||
|
f"Erreur Ollama : génération non terminée en {self._timeout}s. Réduisez "
|
||||||
|
"la taille des morceaux d'import, augmentez le timeout, ou utilisez un "
|
||||||
|
"modèle plus rapide."
|
||||||
|
) from exc
|
||||||
except httpx.HTTPError as exc:
|
except httpx.HTTPError as exc:
|
||||||
raise LLMProviderError(
|
raise LLMProviderError(
|
||||||
f"Erreur lors de l'appel à Ollama : {exc}"
|
f"Erreur lors de l'appel à Ollama : {exc}"
|
||||||
) from exc
|
) from exc
|
||||||
|
|
||||||
return response.json()["response"]
|
data = response.json()
|
||||||
|
# Diagnostic crucial pour les imports : `done_reason` != "stop" signifie que
|
||||||
|
# la génération a été INTERROMPUE (fenêtre de contexte pleine, num_predict…)
|
||||||
|
# et non terminée par le modèle. Sans ce log, on ne voit qu'un JSON coupé
|
||||||
|
# en aval, sans la cause. `prompt_eval_count` révèle aussi la VRAIE taille
|
||||||
|
# du prompt en tokens du modèle (les morceaux sont mesurés en tokens
|
||||||
|
# cl100k, ~20-40% plus compacts que les tokenizers locaux).
|
||||||
|
done_reason = data.get("done_reason")
|
||||||
|
if done_reason and done_reason != "stop":
|
||||||
|
logger.warning(
|
||||||
|
"Ollama a interrompu la génération (done_reason=%s) : prompt=%s tokens, "
|
||||||
|
"sortie=%s tokens, num_ctx demandé=%s. Si prompt+sortie ≈ num_ctx, la "
|
||||||
|
"fenêtre de contexte est pleine : réduisez la taille des morceaux "
|
||||||
|
"d'import ou augmentez num_ctx (Paramètres).",
|
||||||
|
done_reason, data.get("prompt_eval_count"),
|
||||||
|
data.get("eval_count"), self._num_ctx,
|
||||||
|
)
|
||||||
|
return data["response"]
|
||||||
|
|
||||||
async def stream_chat(
|
async def stream_chat(
|
||||||
self,
|
self,
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ from app.core.config import Settings
|
|||||||
from app.domain.models import ChatMessage
|
from app.domain.models import ChatMessage
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
from app.domain.ports import LLMProviderError
|
from app.domain.ports import LLMGenerationTimeout, LLMProviderError
|
||||||
|
|
||||||
_API_URL = "https://openrouter.ai/api/v1/chat/completions"
|
_API_URL = "https://openrouter.ai/api/v1/chat/completions"
|
||||||
|
|
||||||
@@ -113,7 +113,7 @@ class OpenRouterLLMProvider:
|
|||||||
try:
|
try:
|
||||||
return await asyncio.wait_for(_collect(), timeout=self._timeout)
|
return await asyncio.wait_for(_collect(), timeout=self._timeout)
|
||||||
except asyncio.TimeoutError as exc:
|
except asyncio.TimeoutError as exc:
|
||||||
raise LLMProviderError(
|
raise LLMGenerationTimeout(
|
||||||
f"Erreur {provider} : génération non terminée en {self._timeout}s. Réduisez la "
|
f"Erreur {provider} : génération non terminée en {self._timeout}s. Réduisez la "
|
||||||
"taille des morceaux d'import, augmentez le timeout, ou changez de modèle."
|
"taille des morceaux d'import, augmentez le timeout, ou changez de modèle."
|
||||||
) from exc
|
) from exc
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ from app.infrastructure.ollama_model_installer import ensure_ollama_embedding_mo
|
|||||||
app = FastAPI(
|
app = FastAPI(
|
||||||
title="LoreMind Brain",
|
title="LoreMind Brain",
|
||||||
description="Backend IA pour la génération de contenu narratif.",
|
description="Backend IA pour la génération de contenu narratif.",
|
||||||
version="0.12.0-beta",
|
version="0.12.2-beta",
|
||||||
)
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|||||||
@@ -14,7 +14,7 @@
|
|||||||
|
|
||||||
<groupId>com.loremind</groupId>
|
<groupId>com.loremind</groupId>
|
||||||
<artifactId>loremind-core</artifactId>
|
<artifactId>loremind-core</artifactId>
|
||||||
<version>0.12.0-beta</version>
|
<version>0.12.2-beta</version>
|
||||||
<name>LoreMind Core</name>
|
<name>LoreMind Core</name>
|
||||||
<description>Backend Core - Architecture Hexagonale</description>
|
<description>Backend Core - Architecture Hexagonale</description>
|
||||||
|
|
||||||
|
|||||||
@@ -64,9 +64,11 @@ public class CampaignImportService {
|
|||||||
byte[] pdfBytes,
|
byte[] pdfBytes,
|
||||||
String filename,
|
String filename,
|
||||||
Consumer<CampaignImportProgress> onProgress,
|
Consumer<CampaignImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<CampaignImportProposal> onDone,
|
Consumer<CampaignImportProposal> onDone,
|
||||||
Consumer<Throwable> onError) {
|
Consumer<Throwable> onError) {
|
||||||
campaignPdfImporter.importCampaignStreaming(pdfBytes, filename, onProgress, onDone, onError);
|
campaignPdfImporter.importCampaignStreaming(
|
||||||
|
pdfBytes, filename, onProgress, onHeartbeat, onDone, onError);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -39,9 +39,11 @@ public class GameSystemService {
|
|||||||
byte[] pdfBytes,
|
byte[] pdfBytes,
|
||||||
String filename,
|
String filename,
|
||||||
java.util.function.Consumer<com.loremind.domain.gamesystemcontext.RulesImportProgress> onProgress,
|
java.util.function.Consumer<com.loremind.domain.gamesystemcontext.RulesImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
java.util.function.Consumer<RulesImportResult> onDone,
|
java.util.function.Consumer<RulesImportResult> onDone,
|
||||||
java.util.function.Consumer<Throwable> onError) {
|
java.util.function.Consumer<Throwable> onError) {
|
||||||
rulesPdfImporter.importRulesStreaming(pdfBytes, filename, onProgress, onDone, onError);
|
rulesPdfImporter.importRulesStreaming(
|
||||||
|
pdfBytes, filename, onProgress, onHeartbeat, onDone, onError);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -16,6 +16,9 @@ public interface CampaignPdfImporter {
|
|||||||
* l'avancement au fil de l'eau, puis la proposition finale.
|
* l'avancement au fil de l'eau, puis la proposition finale.
|
||||||
*
|
*
|
||||||
* @param onProgress invoqué à chaque étape (extraction, puis par morceau).
|
* @param onProgress invoqué à chaque étape (extraction, puis par morceau).
|
||||||
|
* @param onHeartbeat invoqué périodiquement pendant un appel LLM long (aucune
|
||||||
|
* avancée à afficher, mais le canal SSE vers le navigateur
|
||||||
|
* doit rester actif — sinon un proxy intermédiaire le coupe).
|
||||||
* @param onDone invoqué une fois avec l'arbre proposé (non persisté).
|
* @param onDone invoqué une fois avec l'arbre proposé (non persisté).
|
||||||
* @param onError invoqué si l'extraction/structuration échoue.
|
* @param onError invoqué si l'extraction/structuration échoue.
|
||||||
*/
|
*/
|
||||||
@@ -23,6 +26,7 @@ public interface CampaignPdfImporter {
|
|||||||
byte[] pdfBytes,
|
byte[] pdfBytes,
|
||||||
String filename,
|
String filename,
|
||||||
Consumer<CampaignImportProgress> onProgress,
|
Consumer<CampaignImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<CampaignImportProposal> onDone,
|
Consumer<CampaignImportProposal> onDone,
|
||||||
Consumer<Throwable> onError);
|
Consumer<Throwable> onError);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -27,6 +27,9 @@ public interface RulesPdfImporter {
|
|||||||
* d'exécution de l'adapter (synchrone jusqu'à {@code onDone}/{@code onError}).
|
* d'exécution de l'adapter (synchrone jusqu'à {@code onDone}/{@code onError}).
|
||||||
*
|
*
|
||||||
* @param onProgress invoqué à chaque étape (extraction, puis par morceau).
|
* @param onProgress invoqué à chaque étape (extraction, puis par morceau).
|
||||||
|
* @param onHeartbeat invoqué périodiquement pendant un appel LLM long (aucune
|
||||||
|
* avancée à afficher, mais le canal SSE vers le navigateur
|
||||||
|
* doit rester actif — sinon un proxy intermédiaire le coupe).
|
||||||
* @param onDone invoqué une fois avec le résultat final.
|
* @param onDone invoqué une fois avec le résultat final.
|
||||||
* @param onError invoqué si l'extraction/structuration échoue.
|
* @param onError invoqué si l'extraction/structuration échoue.
|
||||||
*/
|
*/
|
||||||
@@ -34,6 +37,7 @@ public interface RulesPdfImporter {
|
|||||||
byte[] pdfBytes,
|
byte[] pdfBytes,
|
||||||
String filename,
|
String filename,
|
||||||
Consumer<RulesImportProgress> onProgress,
|
Consumer<RulesImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<RulesImportResult> onDone,
|
Consumer<RulesImportResult> onDone,
|
||||||
Consumer<Throwable> onError);
|
Consumer<Throwable> onError);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -60,6 +60,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
|
|||||||
byte[] pdfBytes,
|
byte[] pdfBytes,
|
||||||
String filename,
|
String filename,
|
||||||
Consumer<CampaignImportProgress> onProgress,
|
Consumer<CampaignImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<CampaignImportProposal> onDone,
|
Consumer<CampaignImportProposal> onDone,
|
||||||
Consumer<Throwable> onError) {
|
Consumer<Throwable> onError) {
|
||||||
|
|
||||||
@@ -83,7 +84,8 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
|
|||||||
flux
|
flux
|
||||||
.timeout(Duration.ofSeconds(importTimeoutSeconds))
|
.timeout(Duration.ofSeconds(importTimeoutSeconds))
|
||||||
.doOnNext(sse -> handleEvent(
|
.doOnNext(sse -> handleEvent(
|
||||||
sse, pageCount, ocrPageCount, terminated, onProgress, onDone, onError))
|
sse, pageCount, ocrPageCount, terminated,
|
||||||
|
onProgress, onHeartbeat, onDone, onError))
|
||||||
.blockLast();
|
.blockLast();
|
||||||
if (!terminated[0]) {
|
if (!terminated[0]) {
|
||||||
onError.accept(new CampaignImportException(
|
onError.accept(new CampaignImportException(
|
||||||
@@ -107,12 +109,19 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
|
|||||||
int[] ocrPageCount,
|
int[] ocrPageCount,
|
||||||
boolean[] terminated,
|
boolean[] terminated,
|
||||||
Consumer<CampaignImportProgress> onProgress,
|
Consumer<CampaignImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<CampaignImportProposal> onDone,
|
Consumer<CampaignImportProposal> onDone,
|
||||||
Consumer<Throwable> onError) {
|
Consumer<Throwable> onError) {
|
||||||
|
|
||||||
String event = sse.event();
|
String event = sse.event();
|
||||||
String data = sse.data() == null ? "" : sse.data();
|
String data = sse.data() == null ? "" : sse.data();
|
||||||
|
|
||||||
|
if ("heartbeat".equals(event)) {
|
||||||
|
// Keep-alive du Brain pendant un appel LLM long : à PROPAGER jusqu'au
|
||||||
|
// navigateur, sinon nginx (proxy_read_timeout) coupe le SSE Core→front.
|
||||||
|
onHeartbeat.run();
|
||||||
|
return;
|
||||||
|
}
|
||||||
if ("error".equals(event)) {
|
if ("error".equals(event)) {
|
||||||
terminated[0] = true;
|
terminated[0] = true;
|
||||||
onError.accept(new CampaignImportException(
|
onError.accept(new CampaignImportException(
|
||||||
|
|||||||
@@ -114,6 +114,7 @@ public class BrainRulesImportClient implements RulesPdfImporter {
|
|||||||
byte[] pdfBytes,
|
byte[] pdfBytes,
|
||||||
String filename,
|
String filename,
|
||||||
Consumer<RulesImportProgress> onProgress,
|
Consumer<RulesImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<RulesImportResult> onDone,
|
Consumer<RulesImportResult> onDone,
|
||||||
Consumer<Throwable> onError) {
|
Consumer<Throwable> onError) {
|
||||||
|
|
||||||
@@ -139,7 +140,8 @@ public class BrainRulesImportClient implements RulesPdfImporter {
|
|||||||
flux
|
flux
|
||||||
.timeout(Duration.ofSeconds(importTimeoutSeconds))
|
.timeout(Duration.ofSeconds(importTimeoutSeconds))
|
||||||
.doOnNext(sse -> handleEvent(
|
.doOnNext(sse -> handleEvent(
|
||||||
sse, pageCount, ocrPageCount, terminated, onProgress, onDone, onError))
|
sse, pageCount, ocrPageCount, terminated,
|
||||||
|
onProgress, onHeartbeat, onDone, onError))
|
||||||
.blockLast();
|
.blockLast();
|
||||||
// Flux terminé sans event done/error (ex: connexion coupée) → on signale.
|
// Flux terminé sans event done/error (ex: connexion coupée) → on signale.
|
||||||
if (!terminated[0]) {
|
if (!terminated[0]) {
|
||||||
@@ -165,12 +167,20 @@ public class BrainRulesImportClient implements RulesPdfImporter {
|
|||||||
int[] ocrPageCount,
|
int[] ocrPageCount,
|
||||||
boolean[] terminated,
|
boolean[] terminated,
|
||||||
Consumer<RulesImportProgress> onProgress,
|
Consumer<RulesImportProgress> onProgress,
|
||||||
|
Runnable onHeartbeat,
|
||||||
Consumer<RulesImportResult> onDone,
|
Consumer<RulesImportResult> onDone,
|
||||||
Consumer<Throwable> onError) {
|
Consumer<Throwable> onError) {
|
||||||
|
|
||||||
String event = sse.event();
|
String event = sse.event();
|
||||||
String data = sse.data() == null ? "" : sse.data();
|
String data = sse.data() == null ? "" : sse.data();
|
||||||
|
|
||||||
|
if ("heartbeat".equals(event)) {
|
||||||
|
// Keep-alive du Brain pendant un appel LLM long : à PROPAGER jusqu'au
|
||||||
|
// navigateur, sinon nginx (proxy_read_timeout) coupe le SSE Core→front
|
||||||
|
// resté silencieux pendant tout le traitement du morceau.
|
||||||
|
onHeartbeat.run();
|
||||||
|
return;
|
||||||
|
}
|
||||||
if ("error".equals(event)) {
|
if ("error".equals(event)) {
|
||||||
terminated[0] = true;
|
terminated[0] = true;
|
||||||
onError.accept(new RulesImportException(
|
onError.accept(new RulesImportException(
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import org.springframework.http.converter.HttpMessageNotReadableException;
|
|||||||
import org.springframework.web.bind.MethodArgumentNotValidException;
|
import org.springframework.web.bind.MethodArgumentNotValidException;
|
||||||
import org.springframework.web.bind.annotation.ExceptionHandler;
|
import org.springframework.web.bind.annotation.ExceptionHandler;
|
||||||
import org.springframework.web.bind.annotation.RestControllerAdvice;
|
import org.springframework.web.bind.annotation.RestControllerAdvice;
|
||||||
|
import org.springframework.web.context.request.async.AsyncRequestNotUsableException;
|
||||||
|
|
||||||
import java.util.LinkedHashMap;
|
import java.util.LinkedHashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
@@ -70,6 +71,18 @@ public class GlobalExceptionHandler {
|
|||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Client HTTP parti pendant une reponse asynchrone (SSE) : le navigateur a ferme
|
||||||
|
* la connexion (onglet ferme, proxy coupe...), la reponse n'est plus utilisable.
|
||||||
|
* Ce n'est PAS une erreur serveur -> pas de log ERROR + stack trace (bruit),
|
||||||
|
* et aucune reponse a renvoyer (le canal est mort).
|
||||||
|
*/
|
||||||
|
@ExceptionHandler(AsyncRequestNotUsableException.class)
|
||||||
|
public void handleClientDisconnected(HttpServletRequest request, AsyncRequestNotUsableException ex) {
|
||||||
|
log.debug("Client deconnecte pendant la reponse asynchrone sur {} {} : {}",
|
||||||
|
request.getMethod(), request.getRequestURI(), ex.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Fallback : tout ce qui n'a pas ete catche au-dessus -> 500, mais avec
|
* Fallback : tout ce qui n'a pas ete catche au-dessus -> 500, mais avec
|
||||||
* un log ERROR explicite (path + stack trace) et un body JSON debuggable
|
* un log ERROR explicite (path + stack trace) et un body JSON debuggable
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
|
|||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* REST Controller pour l'import d'un PDF de campagne → arbre arc/chapitre/scène.
|
* REST Controller pour l'import d'un PDF de campagne → arbre arc/chapitre/scène.
|
||||||
@@ -30,8 +31,13 @@ public class CampaignImportController {
|
|||||||
|
|
||||||
private static final Logger log = LoggerFactory.getLogger(CampaignImportController.class);
|
private static final Logger log = LoggerFactory.getLogger(CampaignImportController.class);
|
||||||
|
|
||||||
/** Timeout SSE généreux : un import de livre entier peut durer plusieurs minutes. */
|
/**
|
||||||
private static final long IMPORT_SSE_TIMEOUT_MS = 15 * 60 * 1000L;
|
* Timeout SSE = durée TOTALE maximale de l'import (pas un timeout d'inactivité :
|
||||||
|
* les heartbeats ne le réarment pas). Un livre entier sur un modèle local peut
|
||||||
|
* largement dépasser 15 min → 60 min. La déconnexion du client reste détectée
|
||||||
|
* immédiatement par ailleurs (échec d'envoi → interruption de l'import).
|
||||||
|
*/
|
||||||
|
private static final long IMPORT_SSE_TIMEOUT_MS = 60 * 60 * 1000L;
|
||||||
|
|
||||||
private final CampaignImportService campaignImportService;
|
private final CampaignImportService campaignImportService;
|
||||||
private final TaskExecutor taskExecutor;
|
private final TaskExecutor taskExecutor;
|
||||||
@@ -54,33 +60,63 @@ public class CampaignImportController {
|
|||||||
@RequestParam("file") MultipartFile file) throws IOException {
|
@RequestParam("file") MultipartFile file) throws IOException {
|
||||||
SseEmitter emitter = new SseEmitter(IMPORT_SSE_TIMEOUT_MS);
|
SseEmitter emitter = new SseEmitter(IMPORT_SSE_TIMEOUT_MS);
|
||||||
if (file == null || file.isEmpty()) {
|
if (file == null || file.isEmpty()) {
|
||||||
sendError(emitter, "Fichier PDF vide.");
|
sendError(emitter, new AtomicBoolean(false), "Fichier PDF vide.");
|
||||||
return emitter;
|
return emitter;
|
||||||
}
|
}
|
||||||
byte[] bytes = file.getBytes();
|
byte[] bytes = file.getBytes();
|
||||||
String filename = file.getOriginalFilename();
|
String filename = file.getOriginalFilename();
|
||||||
|
|
||||||
|
// Suivi de la déconnexion du navigateur : dès qu'un envoi échoue (ou que
|
||||||
|
// l'emitter se termine), on cesse d'envoyer ET on interrompt le streaming
|
||||||
|
// amont (ClientGoneException remonte dans le doOnNext du WebClient →
|
||||||
|
// annule la souscription → le Brain voit la coupure et stoppe le LLM).
|
||||||
|
AtomicBoolean clientGone = new AtomicBoolean(false);
|
||||||
|
emitter.onTimeout(() -> {
|
||||||
|
// Timeout = durée totale dépassée, mais la connexion est encore vivante :
|
||||||
|
// on envoie une vraie erreur au navigateur AVANT de fermer (sinon le flux
|
||||||
|
// se termine en silence et l'UI reste figée sur la barre de progression).
|
||||||
|
sendError(emitter, clientGone,
|
||||||
|
"L'import a dépassé la durée maximale autorisée et a été interrompu. "
|
||||||
|
+ "Réessayez avec un modèle plus rapide ou un PDF plus petit.");
|
||||||
|
clientGone.set(true);
|
||||||
|
});
|
||||||
|
emitter.onError(e -> clientGone.set(true));
|
||||||
|
|
||||||
taskExecutor.execute(() -> {
|
taskExecutor.execute(() -> {
|
||||||
try {
|
try {
|
||||||
campaignImportService.importStructureStreaming(
|
campaignImportService.importStructureStreaming(
|
||||||
bytes, filename,
|
bytes, filename,
|
||||||
progress -> sendEvent(emitter, "progress", progress),
|
progress -> sendEvent(emitter, clientGone, "progress", progress),
|
||||||
|
() -> sendHeartbeat(emitter, clientGone),
|
||||||
proposal -> {
|
proposal -> {
|
||||||
sendEvent(emitter, "done", proposal);
|
sendEvent(emitter, clientGone, "done", proposal);
|
||||||
emitter.complete();
|
emitter.complete();
|
||||||
},
|
},
|
||||||
error -> {
|
error -> {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
log.info("Import campagne (stream) interrompu : client déconnecté.");
|
||||||
|
return;
|
||||||
|
}
|
||||||
log.warn("Import campagne (stream) échoué : {}", error.getMessage());
|
log.warn("Import campagne (stream) échoué : {}", error.getMessage());
|
||||||
sendError(emitter, error.getMessage());
|
sendError(emitter, clientGone, error.getMessage());
|
||||||
});
|
});
|
||||||
|
} catch (ClientGoneException e) {
|
||||||
|
log.info("Import campagne (stream) interrompu : client déconnecté.");
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
log.warn("Import campagne (stream) échoué : {}", e.getMessage());
|
log.warn("Import campagne (stream) échoué : {}", e.getMessage());
|
||||||
sendError(emitter, e.getMessage());
|
sendError(emitter, clientGone, e.getMessage());
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
return emitter;
|
return emitter;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Signale que le navigateur a fermé le flux SSE : inutile de continuer l'import. */
|
||||||
|
private static final class ClientGoneException extends RuntimeException {
|
||||||
|
ClientGoneException(Throwable cause) {
|
||||||
|
super("Client SSE déconnecté.", cause);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@PostMapping(value = "/apply", consumes = MediaType.APPLICATION_JSON_VALUE)
|
@PostMapping(value = "/apply", consumes = MediaType.APPLICATION_JSON_VALUE)
|
||||||
public ResponseEntity<CampaignImportService.ApplyResult> apply(
|
public ResponseEntity<CampaignImportService.ApplyResult> apply(
|
||||||
@PathVariable String campaignId,
|
@PathVariable String campaignId,
|
||||||
@@ -96,23 +132,52 @@ public class CampaignImportController {
|
|||||||
|
|
||||||
// --- Helpers SSE ---------------------------------------------------------
|
// --- Helpers SSE ---------------------------------------------------------
|
||||||
|
|
||||||
private void sendEvent(SseEmitter emitter, String eventName, Object payload) {
|
private void sendEvent(
|
||||||
|
SseEmitter emitter, AtomicBoolean clientGone, String eventName, Object payload) {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
throw new ClientGoneException(null);
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
emitter.send(SseEmitter.event().name(eventName).data(
|
emitter.send(SseEmitter.event().name(eventName).data(
|
||||||
objectMapper.writeValueAsString(payload), MediaType.APPLICATION_JSON));
|
objectMapper.writeValueAsString(payload), MediaType.APPLICATION_JSON));
|
||||||
} catch (IOException e) {
|
} catch (Exception e) {
|
||||||
|
// IOException OU IllegalStateException (emitter déjà terminé) : le client
|
||||||
|
// est parti — on interrompt le pipeline amont au lieu de rejouer l'échec.
|
||||||
|
clientGone.set(true);
|
||||||
emitter.completeWithError(e);
|
emitter.completeWithError(e);
|
||||||
|
throw new ClientGoneException(e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void sendError(SseEmitter emitter, String message) {
|
/**
|
||||||
|
* Keep-alive vers le navigateur pendant un appel LLM long : un commentaire SSE
|
||||||
|
* (ignoré par le front) suffit à réarmer le {@code proxy_read_timeout} de nginx.
|
||||||
|
*/
|
||||||
|
private void sendHeartbeat(SseEmitter emitter, AtomicBoolean clientGone) {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
throw new ClientGoneException(null);
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
emitter.send(SseEmitter.event().comment("keepalive"));
|
||||||
|
} catch (Exception e) {
|
||||||
|
clientGone.set(true);
|
||||||
|
emitter.completeWithError(e);
|
||||||
|
throw new ClientGoneException(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void sendError(SseEmitter emitter, AtomicBoolean clientGone, String message) {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
return; // le client n'est plus là pour lire le message d'erreur.
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
emitter.send(SseEmitter.event().name("error").data(
|
emitter.send(SseEmitter.event().name("error").data(
|
||||||
objectMapper.writeValueAsString(Map.of(
|
objectMapper.writeValueAsString(Map.of(
|
||||||
"message", message != null ? message : "Erreur inconnue.")),
|
"message", message != null ? message : "Erreur inconnue.")),
|
||||||
MediaType.APPLICATION_JSON));
|
MediaType.APPLICATION_JSON));
|
||||||
emitter.complete();
|
emitter.complete();
|
||||||
} catch (IOException e) {
|
} catch (Exception e) {
|
||||||
|
clientGone.set(true);
|
||||||
emitter.completeWithError(e);
|
emitter.completeWithError(e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,7 +3,6 @@ package com.loremind.infrastructure.web.controller;
|
|||||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import com.loremind.application.gamesystemcontext.GameSystemService;
|
import com.loremind.application.gamesystemcontext.GameSystemService;
|
||||||
import com.loremind.domain.gamesystemcontext.GameSystem;
|
import com.loremind.domain.gamesystemcontext.GameSystem;
|
||||||
import com.loremind.domain.gamesystemcontext.RulesImportProgress;
|
|
||||||
import com.loremind.domain.gamesystemcontext.RulesImportResult;
|
import com.loremind.domain.gamesystemcontext.RulesImportResult;
|
||||||
import com.loremind.domain.gamesystemcontext.ports.RulesImportException;
|
import com.loremind.domain.gamesystemcontext.ports.RulesImportException;
|
||||||
import com.loremind.domain.shared.template.TemplateField;
|
import com.loremind.domain.shared.template.TemplateField;
|
||||||
@@ -27,6 +26,7 @@ import java.io.IOException;
|
|||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
@RestController
|
@RestController
|
||||||
@@ -35,8 +35,13 @@ public class GameSystemController {
|
|||||||
|
|
||||||
private static final Logger log = LoggerFactory.getLogger(GameSystemController.class);
|
private static final Logger log = LoggerFactory.getLogger(GameSystemController.class);
|
||||||
|
|
||||||
/** Timeout SSE généreux : un import de livre entier peut durer plusieurs minutes. */
|
/**
|
||||||
private static final long IMPORT_SSE_TIMEOUT_MS = 15 * 60 * 1000L;
|
* Timeout SSE = durée TOTALE maximale de l'import (pas un timeout d'inactivité :
|
||||||
|
* les heartbeats ne le réarment pas). Un livre entier sur un modèle local peut
|
||||||
|
* largement dépasser 15 min → 60 min. La déconnexion du client reste détectée
|
||||||
|
* immédiatement par ailleurs (échec d'envoi → interruption de l'import).
|
||||||
|
*/
|
||||||
|
private static final long IMPORT_SSE_TIMEOUT_MS = 60 * 60 * 1000L;
|
||||||
|
|
||||||
private final GameSystemService gameSystemService;
|
private final GameSystemService gameSystemService;
|
||||||
private final GameSystemMapper gameSystemMapper;
|
private final GameSystemMapper gameSystemMapper;
|
||||||
@@ -135,7 +140,7 @@ public class GameSystemController {
|
|||||||
public SseEmitter importRulesStream(@RequestParam("file") MultipartFile file) throws IOException {
|
public SseEmitter importRulesStream(@RequestParam("file") MultipartFile file) throws IOException {
|
||||||
SseEmitter emitter = new SseEmitter(IMPORT_SSE_TIMEOUT_MS);
|
SseEmitter emitter = new SseEmitter(IMPORT_SSE_TIMEOUT_MS);
|
||||||
if (file == null || file.isEmpty()) {
|
if (file == null || file.isEmpty()) {
|
||||||
sendImportError(emitter, "Fichier PDF vide.");
|
sendImportError(emitter, new AtomicBoolean(false), "Fichier PDF vide.");
|
||||||
return emitter;
|
return emitter;
|
||||||
}
|
}
|
||||||
// Les octets sont lus sur le thread servlet (le MultipartFile n'est plus
|
// Les octets sont lus sur le thread servlet (le MultipartFile n'est plus
|
||||||
@@ -143,46 +148,108 @@ public class GameSystemController {
|
|||||||
byte[] bytes = file.getBytes();
|
byte[] bytes = file.getBytes();
|
||||||
String filename = file.getOriginalFilename();
|
String filename = file.getOriginalFilename();
|
||||||
|
|
||||||
|
// Suivi de la déconnexion du navigateur : dès qu'un envoi échoue (ou que
|
||||||
|
// l'emitter se termine), on cesse d'envoyer ET on interrompt le streaming
|
||||||
|
// amont (l'exception ClientGone remonte dans le doOnNext du WebClient →
|
||||||
|
// annule la souscription → le Brain voit la coupure et stoppe le LLM).
|
||||||
|
AtomicBoolean clientGone = new AtomicBoolean(false);
|
||||||
|
emitter.onTimeout(() -> {
|
||||||
|
// Timeout = durée totale dépassée, mais la connexion est encore vivante :
|
||||||
|
// on envoie une vraie erreur au navigateur AVANT de fermer (sinon le flux
|
||||||
|
// se termine en silence et l'UI reste figée sur la barre de progression).
|
||||||
|
sendImportError(emitter, clientGone,
|
||||||
|
"L'import a dépassé la durée maximale autorisée et a été interrompu. "
|
||||||
|
+ "Réessayez avec un modèle plus rapide ou un PDF plus petit.");
|
||||||
|
clientGone.set(true);
|
||||||
|
});
|
||||||
|
emitter.onError(e -> clientGone.set(true));
|
||||||
|
|
||||||
taskExecutor.execute(() -> {
|
taskExecutor.execute(() -> {
|
||||||
try {
|
try {
|
||||||
gameSystemService.importRulesFromPdfStreaming(
|
gameSystemService.importRulesFromPdfStreaming(
|
||||||
bytes, filename,
|
bytes, filename,
|
||||||
progress -> sendImportEvent(emitter, "progress", progress),
|
progress -> sendImportEvent(emitter, clientGone, "progress", progress),
|
||||||
|
() -> sendImportHeartbeat(emitter, clientGone),
|
||||||
result -> {
|
result -> {
|
||||||
sendImportEvent(emitter, "done", result);
|
sendImportEvent(emitter, clientGone, "done", result);
|
||||||
emitter.complete();
|
emitter.complete();
|
||||||
},
|
},
|
||||||
error -> {
|
error -> {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
// La "panne" amont n'est que l'écho de la déconnexion
|
||||||
|
// du navigateur : pas un échec d'import.
|
||||||
|
log.info("Import de règles (stream) interrompu : client déconnecté.");
|
||||||
|
return;
|
||||||
|
}
|
||||||
log.warn("Import de règles (stream) échoué : {}", error.getMessage());
|
log.warn("Import de règles (stream) échoué : {}", error.getMessage());
|
||||||
sendImportError(emitter, error.getMessage());
|
sendImportError(emitter, clientGone, error.getMessage());
|
||||||
});
|
});
|
||||||
|
} catch (ClientGoneException e) {
|
||||||
|
log.info("Import de règles (stream) interrompu : client déconnecté.");
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
log.warn("Import de règles (stream) échoué : {}", e.getMessage());
|
log.warn("Import de règles (stream) échoué : {}", e.getMessage());
|
||||||
sendImportError(emitter, e.getMessage());
|
sendImportError(emitter, clientGone, e.getMessage());
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
return emitter;
|
return emitter;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Signale que le navigateur a fermé le flux SSE : inutile de continuer l'import. */
|
||||||
|
private static final class ClientGoneException extends RuntimeException {
|
||||||
|
ClientGoneException(Throwable cause) {
|
||||||
|
super("Client SSE déconnecté.", cause);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** Sérialise `payload` en JSON et l'envoie comme évènement SSE nommé. */
|
/** Sérialise `payload` en JSON et l'envoie comme évènement SSE nommé. */
|
||||||
private void sendImportEvent(SseEmitter emitter, String eventName, Object payload) {
|
private void sendImportEvent(
|
||||||
|
SseEmitter emitter, AtomicBoolean clientGone, String eventName, Object payload) {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
throw new ClientGoneException(null);
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
emitter.send(SseEmitter.event().name(eventName).data(
|
emitter.send(SseEmitter.event().name(eventName).data(
|
||||||
objectMapper.writeValueAsString(payload), MediaType.APPLICATION_JSON));
|
objectMapper.writeValueAsString(payload), MediaType.APPLICATION_JSON));
|
||||||
} catch (IOException e) {
|
} catch (Exception e) {
|
||||||
|
// IOException OU IllegalStateException (emitter déjà terminé) : le client
|
||||||
|
// est parti. On marque l'état et on INTERROMPT le pipeline amont — sinon
|
||||||
|
// chaque évènement suivant rejouerait l'échec (bruit de logs + LLM gaspillé).
|
||||||
|
clientGone.set(true);
|
||||||
emitter.completeWithError(e);
|
emitter.completeWithError(e);
|
||||||
|
throw new ClientGoneException(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Keep-alive vers le navigateur pendant un appel LLM long : un commentaire SSE
|
||||||
|
* (ignoré par le front) suffit à réarmer le {@code proxy_read_timeout} de nginx.
|
||||||
|
*/
|
||||||
|
private void sendImportHeartbeat(SseEmitter emitter, AtomicBoolean clientGone) {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
throw new ClientGoneException(null);
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
emitter.send(SseEmitter.event().comment("keepalive"));
|
||||||
|
} catch (Exception e) {
|
||||||
|
clientGone.set(true);
|
||||||
|
emitter.completeWithError(e);
|
||||||
|
throw new ClientGoneException(e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Envoie un évènement `error` {message} puis termine le flux. */
|
/** Envoie un évènement `error` {message} puis termine le flux. */
|
||||||
private void sendImportError(SseEmitter emitter, String message) {
|
private void sendImportError(SseEmitter emitter, AtomicBoolean clientGone, String message) {
|
||||||
|
if (clientGone.get()) {
|
||||||
|
return; // le client n'est plus là pour lire le message d'erreur.
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
emitter.send(SseEmitter.event().name("error").data(
|
emitter.send(SseEmitter.event().name("error").data(
|
||||||
objectMapper.writeValueAsString(Map.of(
|
objectMapper.writeValueAsString(Map.of(
|
||||||
"message", message != null ? message : "Erreur inconnue.")),
|
"message", message != null ? message : "Erreur inconnue.")),
|
||||||
MediaType.APPLICATION_JSON));
|
MediaType.APPLICATION_JSON));
|
||||||
emitter.complete();
|
emitter.complete();
|
||||||
} catch (IOException e) {
|
} catch (Exception e) {
|
||||||
|
clientGone.set(true);
|
||||||
emitter.completeWithError(e);
|
emitter.completeWithError(e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
4
web/package-lock.json
generated
4
web/package-lock.json
generated
@@ -1,12 +1,12 @@
|
|||||||
{
|
{
|
||||||
"name": "loremind-web",
|
"name": "loremind-web",
|
||||||
"version": "0.12.0-beta",
|
"version": "0.12.2-beta",
|
||||||
"lockfileVersion": 3,
|
"lockfileVersion": 3,
|
||||||
"requires": true,
|
"requires": true,
|
||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"name": "loremind-web",
|
"name": "loremind-web",
|
||||||
"version": "0.12.0-beta",
|
"version": "0.12.2-beta",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@angular/animations": "^21.2.16",
|
"@angular/animations": "^21.2.16",
|
||||||
"@angular/common": "^21.2.16",
|
"@angular/common": "^21.2.16",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "loremind-web",
|
"name": "loremind-web",
|
||||||
"version": "0.12.0-beta",
|
"version": "0.12.2-beta",
|
||||||
"description": "LoreMind Frontend - Angular",
|
"description": "LoreMind Frontend - Angular",
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"ng": "ng",
|
"ng": "ng",
|
||||||
|
|||||||
@@ -61,17 +61,24 @@ export class CampaignImportService {
|
|||||||
let buffer = '';
|
let buffer = '';
|
||||||
let currentEvent: string | null = null;
|
let currentEvent: string | null = null;
|
||||||
let currentData = '';
|
let currentData = '';
|
||||||
|
// Le flux s'est-il terminé PROPREMENT (évènement done ou error reçu) ?
|
||||||
|
// Sans ce suivi, une connexion coupée en plein import (timeout serveur,
|
||||||
|
// proxy, Core redémarré) terminait l'Observable en silence : barre de
|
||||||
|
// progression figée et aucun message pour l'utilisateur.
|
||||||
|
let terminated = false;
|
||||||
|
|
||||||
const dispatch = () => {
|
const dispatch = () => {
|
||||||
const name = currentEvent ?? 'message';
|
const name = currentEvent ?? 'message';
|
||||||
if (name === 'error') {
|
if (name === 'error') {
|
||||||
let message = 'Échec de l\'import.';
|
let message = 'Échec de l\'import.';
|
||||||
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* défaut */ }
|
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* défaut */ }
|
||||||
|
terminated = true;
|
||||||
subscriber.error(new Error(message));
|
subscriber.error(new Error(message));
|
||||||
} else if (name === 'progress' || name === 'done') {
|
} else if (name === 'progress' || name === 'done') {
|
||||||
try {
|
try {
|
||||||
const obj = JSON.parse(currentData);
|
const obj = JSON.parse(currentData);
|
||||||
if (name === 'done') {
|
if (name === 'done') {
|
||||||
|
terminated = true;
|
||||||
subscriber.next({ type: 'done', arcs: obj.arcs ?? [], npcs: obj.npcs ?? [] });
|
subscriber.next({ type: 'done', arcs: obj.arcs ?? [], npcs: obj.npcs ?? [] });
|
||||||
subscriber.complete();
|
subscriber.complete();
|
||||||
} else {
|
} else {
|
||||||
@@ -105,9 +112,12 @@ export class CampaignImportService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (currentEvent !== null || currentData !== '') dispatch();
|
if (currentEvent !== null || currentData !== '') dispatch();
|
||||||
subscriber.complete();
|
if (!terminated) {
|
||||||
|
subscriber.error(new Error(
|
||||||
|
'L\'import s\'est interrompu avant la fin (connexion coupée ou délai dépassé). Réessayez.'));
|
||||||
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
subscriber.error(err);
|
if (!terminated) subscriber.error(err);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -92,17 +92,24 @@ export class GameSystemService {
|
|||||||
let buffer = '';
|
let buffer = '';
|
||||||
let currentEvent: string | null = null;
|
let currentEvent: string | null = null;
|
||||||
let currentData = '';
|
let currentData = '';
|
||||||
|
// Le flux s'est-il terminé PROPREMENT (évènement done ou error reçu) ?
|
||||||
|
// Sans ce suivi, une connexion coupée en plein import (timeout serveur,
|
||||||
|
// proxy, Core redémarré) terminait l'Observable en silence : barre de
|
||||||
|
// progression figée et aucun message pour l'utilisateur.
|
||||||
|
let terminated = false;
|
||||||
|
|
||||||
const dispatch = () => {
|
const dispatch = () => {
|
||||||
const name = currentEvent ?? 'message';
|
const name = currentEvent ?? 'message';
|
||||||
if (name === 'error') {
|
if (name === 'error') {
|
||||||
let message = 'Échec de l\'import.';
|
let message = 'Échec de l\'import.';
|
||||||
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* garde le défaut */ }
|
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* garde le défaut */ }
|
||||||
|
terminated = true;
|
||||||
subscriber.error(new Error(message));
|
subscriber.error(new Error(message));
|
||||||
} else if (name === 'progress' || name === 'done') {
|
} else if (name === 'progress' || name === 'done') {
|
||||||
try {
|
try {
|
||||||
const obj = JSON.parse(currentData);
|
const obj = JSON.parse(currentData);
|
||||||
if (name === 'done') {
|
if (name === 'done') {
|
||||||
|
terminated = true;
|
||||||
subscriber.next({ type: 'done', ...obj });
|
subscriber.next({ type: 'done', ...obj });
|
||||||
subscriber.complete();
|
subscriber.complete();
|
||||||
} else {
|
} else {
|
||||||
@@ -136,9 +143,12 @@ export class GameSystemService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (currentEvent !== null || currentData !== '') dispatch();
|
if (currentEvent !== null || currentData !== '') dispatch();
|
||||||
subscriber.complete();
|
if (!terminated) {
|
||||||
|
subscriber.error(new Error(
|
||||||
|
'L\'import s\'est interrompu avant la fin (connexion coupée ou délai dépassé). Réessayez.'));
|
||||||
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
subscriber.error(err);
|
if (!terminated) subscriber.error(err);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user