Compare commits

...

4 Commits

Author SHA1 Message Date
7f519588b6 Amélioration de l'exploitation des PDF par l'IA
All checks were successful
Build & Push Images / build (brain) (push) Successful in 1m44s
Build & Push Images / build (core) (push) Successful in 2m0s
Build & Push Images / build-switcher (push) Successful in 23s
Build & Push Images / build (web) (push) Successful in 1m50s
Amélioration des feedbacks en cas d'erreur d'exploitation des PDF
2026-06-12 01:28:45 +02:00
0799c850ec On essai de contraindre le modèle utilisé par ollama à répondre dans un certain format et ne plus mettre à l'interieur sa "réflexion"
All checks were successful
Build & Push Images / build (brain) (push) Successful in 1m30s
Build & Push Images / build (core) (push) Successful in 1m56s
Build & Push Images / build-switcher (push) Successful in 26s
Build & Push Images / build (web) (push) Successful in 1m44s
2026-06-11 15:51:22 +02:00
113df6a391 amélioration import ollama
All checks were successful
Build & Push Images / build (brain) (push) Successful in 1m59s
Build & Push Images / build (core) (push) Successful in 1m59s
Build & Push Images / build-switcher (push) Successful in 29s
Build & Push Images / build (web) (push) Successful in 1m59s
2026-06-11 15:29:28 +02:00
a1f3b9b796 Améliorations sur l'utilisation de l'IA pour l'exploitation des PDF, que ce soit la partie cloud ou la partie ollama + montée en version
All checks were successful
Build & Push Images / build (brain) (push) Successful in 1m43s
Build & Push Images / build (core) (push) Successful in 1m51s
Build & Push Images / build-switcher (push) Successful in 26s
Build & Push Images / build (web) (push) Successful in 1m46s
2026-06-11 01:31:24 +02:00
25 changed files with 755 additions and 79 deletions

View File

@@ -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
ports et les use cases, jamais Ollama/Mistral/etc.
"""
import logging
from typing import Annotated
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.pdf_extractor import PyMuPdfTextExtractor
logger = logging.getLogger(__name__)
# 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.
_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(
settings: Annotated[Settings, Depends(get_settings)],
) -> LLMProvider:
@@ -82,8 +111,14 @@ def get_import_rules_use_case(
settings: Annotated[Settings, Depends(get_settings)],
) -> ImportRulesUseCase:
"""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(
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(
@@ -94,7 +129,7 @@ def get_import_campaign_use_case(
return ImportCampaignUseCase(
llm=llm,
extractor=_PDF_EXTRACTOR,
chunk_target_tokens=settings.import_chunk_tokens,
chunk_target_tokens=_effective_import_chunk_tokens(settings),
map_concurrency=settings.llm_map_concurrency,
)

View File

@@ -29,7 +29,12 @@ from app.domain.models import (
RoomProposal,
SceneProposal,
)
from app.domain.ports import LLMProvider, LLMProviderError, PdfTextExtractor
from app.domain.ports import (
LLMGenerationTimeout,
LLMProvider,
LLMProviderError,
PdfTextExtractor,
)
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.
- 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
# 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
@@ -480,6 +563,17 @@ class ImportCampaignUseCase:
f"Dernier message : {last_error or 'inconnu'}"}
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
# (best-effort, voir _consolidate). Inutile sur un import mono-morceau.
if total > 1:
@@ -557,8 +651,26 @@ class ImportCampaignUseCase:
+ f"\n\n--- EXTRAIT {index + 1}/{total} ---\n{text}\n\n"
"Renvoie maintenant le JSON de l'arborescence."
)
try:
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)
if truncated and depth < _MAX_SPLIT_DEPTH:

View File

@@ -13,6 +13,7 @@ Ne dépend que des abstractions du domaine (ports LLMProvider + PdfTextExtractor
from __future__ import annotations
import logging
import re
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
@@ -25,7 +26,12 @@ from app.application.streaming import with_heartbeat
# 1-2 niveaux suffisent en pratique, le reste est un garde-fou).
_MAX_SPLIT_DEPTH = 3
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__)
@@ -34,6 +40,16 @@ logger = logging.getLogger(__name__)
# Plus la valeur est haute, plus le modèle "brode" (invente du contenu absent).
_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
# 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.
@@ -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.
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 :
- Tu réponds UNIQUEMENT par un objet JSON valide, sans markdown ni commentaire autour.
- Les CLÉS sont des titres de section (texte court). Les VALEURS sont le contenu de la règle en markdown.
- 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 (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 :
{canonical}
- 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.
- 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:
"""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()}
# 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]:
"""Fusionne deux dicts de sections (issus des 2 moitiés d'un morceau re-découpé).
@@ -128,10 +275,16 @@ class ImportRulesUseCase:
llm: LLMProvider,
extractor: PdfTextExtractor,
chunk_target_tokens: int = CHUNK_TARGET_TOKENS,
segment_only: bool = False,
) -> 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._extractor = extractor
self._chunk_target_tokens = chunk_target_tokens
self._segment_only = segment_only
async def execute(self, pdf_bytes: bytes) -> RulesImportResult:
"""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'}"}
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 {
"type": "done",
"sections": merger.result(),
"sections": sections,
"page_count": doc.page_count,
"ocr_page_count": doc.ocr_page_count,
"skipped": skipped,
@@ -234,15 +399,36 @@ class ImportRulesUseCase:
"""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 —
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 = (
_MAP_SYSTEM.format(
system.format(
canonical="\n".join(f" - {s}" for s in _CANONICAL_SECTIONS)
)
+ f"\n\n--- EXTRAIT {index + 1}/{total} ---\n{text}\n\n"
"Renvoie maintenant le JSON des sections."
)
try:
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)
if truncated and depth < _MAX_SPLIT_DEPTH:
@@ -259,6 +445,68 @@ class ImportRulesUseCase:
"Morceau %s : sortie tronquée, profondeur max atteinte — partiel conservé.", index)
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
def _parse_sections(raw: str, *, index: int) -> tuple[dict[str, str], bool]:
"""Parse robuste → (sections, tronqué). `tronqué`=True si récupération partielle."""
@@ -276,4 +524,5 @@ class ImportRulesUseCase:
if not isinstance(parsed, dict):
logger.warning("Morceau %s : le LLM n'a pas renvoyé un objet, ignoré.", index)
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

View File

@@ -38,13 +38,16 @@ def load_json_object(raw: str) -> tuple[object | None, bool]:
obj = extract_json_object(raw)
if obj is not None:
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:
pass
repaired = repair_truncated_json(raw)
if repaired is not None:
try:
return json.loads(repaired), True
return json.loads(repaired, strict=False), True
except json.JSONDecodeError:
pass
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)
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
sous-structure complète). On exige un contenu substantiel pour éviter les
faux positifs sur une courte réponse non-JSON."""
s = (raw or "").strip()
if "{" not in s or len(s) < 100:
sous-structure complète).
Une réponse qui COMMENCE par `{` est jugée sur le seul équilibre des accolades,
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 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:

View File

@@ -14,7 +14,7 @@ import asyncio
import logging
import re
from app.domain.ports import LLMProvider, LLMProviderError
from app.domain.ports import LLMGenerationTimeout, LLMProvider, LLMProviderError
logger = logging.getLogger(__name__)
@@ -60,7 +60,7 @@ async def generate_with_retry(
llm: LLMProvider,
prompt: str,
*,
output_format: str | None = None,
output_format: str | dict | None = None,
temperature: float | None = None,
) -> str:
"""Comme `llm.generate`, mais réessaie les erreurs transitoires (backoff).
@@ -74,6 +74,12 @@ async def generate_with_retry(
for attempt in range(_ATTEMPTS):
try:
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:
last_error = exc
# Quota JOURNALIER épuisé : inutile d'insister, on remonte tout de suite

View File

@@ -24,17 +24,20 @@ class LLMProvider(Protocol):
self,
prompt: str,
*,
output_format: str | None = None,
output_format: str | dict | None = None,
temperature: float | None = None,
) -> str:
"""Génère une réponse textuelle à partir d'un prompt donné.
Args:
prompt: le texte envoyé au modèle.
output_format: contrainte de format optionnelle. Exemple : "json"
pour forcer le modèle à renvoyer du JSON valide. Les
fournisseurs qui ne supportent pas une valeur donnée doivent
l'ignorer silencieusement ou la traduire au mieux.
output_format: contrainte de format optionnelle. "json" pour forcer
un JSON valide ; un dict = SCHÉMA JSON décrivant la structure
attendue (les fournisseurs qui supportent les sorties
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) à
1.0+ (très créatif, hallucine plus facilement). None =
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
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).
"""

View File

@@ -20,7 +20,7 @@ import httpx
from app.core.config import Settings
from app.domain.models import ChatMessage
from app.domain.ports import LLMProviderError
from app.domain.ports import LLMGenerationTimeout, LLMProviderError
logger = logging.getLogger(__name__)
@@ -95,7 +95,7 @@ class GeminiLLMProvider:
try:
return await asyncio.wait_for(_collect(), timeout=self._timeout)
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 "
"taille des morceaux d'import ou augmentez le timeout."
) from exc
@@ -130,6 +130,12 @@ class GeminiLLMProvider:
}
if temperature is not None:
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:
try:

View File

@@ -20,7 +20,7 @@ import httpx
from app.core.config import Settings
from app.domain.models import ChatMessage
from app.domain.ports import LLMProviderError
from app.domain.ports import LLMGenerationTimeout, LLMProviderError
logger = logging.getLogger(__name__)
@@ -101,7 +101,7 @@ class MistralLLMProvider:
try:
return await asyncio.wait_for(_collect(), timeout=self._timeout)
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 "
"taille des morceaux d'import, augmentez le timeout, ou changez de modèle."
) from exc
@@ -136,6 +136,13 @@ class MistralLLMProvider:
}
if temperature is not None:
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:
try:

View File

@@ -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.
"""
import json
import logging
from typing import AsyncIterator
import httpx
from app.core.config import Settings
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:
@@ -45,7 +48,7 @@ class OllamaLLMProvider:
self,
prompt: str,
*,
output_format: str | None = None,
output_format: str | dict | None = None,
temperature: float | None = None,
) -> str:
url = f"{self._base_url}/api/generate"
@@ -55,6 +58,10 @@ class OllamaLLMProvider:
"stream": False,
"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:
payload["format"] = output_format
@@ -71,12 +78,45 @@ class OllamaLLMProvider:
raise LLMProviderError(
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:
raise LLMProviderError(
f"Erreur lors de l'appel à Ollama : {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(
self,

View File

@@ -22,7 +22,7 @@ from app.core.config import Settings
from app.domain.models import ChatMessage
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"
@@ -113,7 +113,7 @@ class OpenRouterLLMProvider:
try:
return await asyncio.wait_for(_collect(), timeout=self._timeout)
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 "
"taille des morceaux d'import, augmentez le timeout, ou changez de modèle."
) from exc

View File

@@ -26,7 +26,7 @@ from app.infrastructure.ollama_model_installer import ensure_ollama_embedding_mo
app = FastAPI(
title="LoreMind Brain",
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__)

View File

@@ -14,7 +14,7 @@
<groupId>com.loremind</groupId>
<artifactId>loremind-core</artifactId>
<version>0.12.0-beta</version>
<version>0.12.2-beta</version>
<name>LoreMind Core</name>
<description>Backend Core - Architecture Hexagonale</description>

View File

@@ -64,9 +64,11 @@ public class CampaignImportService {
byte[] pdfBytes,
String filename,
Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError) {
campaignPdfImporter.importCampaignStreaming(pdfBytes, filename, onProgress, onDone, onError);
campaignPdfImporter.importCampaignStreaming(
pdfBytes, filename, onProgress, onHeartbeat, onDone, onError);
}
/**

View File

@@ -39,9 +39,11 @@ public class GameSystemService {
byte[] pdfBytes,
String filename,
java.util.function.Consumer<com.loremind.domain.gamesystemcontext.RulesImportProgress> onProgress,
Runnable onHeartbeat,
java.util.function.Consumer<RulesImportResult> onDone,
java.util.function.Consumer<Throwable> onError) {
rulesPdfImporter.importRulesStreaming(pdfBytes, filename, onProgress, onDone, onError);
rulesPdfImporter.importRulesStreaming(
pdfBytes, filename, onProgress, onHeartbeat, onDone, onError);
}
/**

View File

@@ -16,6 +16,9 @@ public interface CampaignPdfImporter {
* l'avancement au fil de l'eau, puis la proposition finale.
*
* @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 onError invoqué si l'extraction/structuration échoue.
*/
@@ -23,6 +26,7 @@ public interface CampaignPdfImporter {
byte[] pdfBytes,
String filename,
Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError);
}

View File

@@ -27,6 +27,9 @@ public interface RulesPdfImporter {
* d'exécution de l'adapter (synchrone jusqu'à {@code onDone}/{@code onError}).
*
* @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 onError invoqué si l'extraction/structuration échoue.
*/
@@ -34,6 +37,7 @@ public interface RulesPdfImporter {
byte[] pdfBytes,
String filename,
Consumer<RulesImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<RulesImportResult> onDone,
Consumer<Throwable> onError);
}

View File

@@ -60,6 +60,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
byte[] pdfBytes,
String filename,
Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError) {
@@ -83,7 +84,8 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
flux
.timeout(Duration.ofSeconds(importTimeoutSeconds))
.doOnNext(sse -> handleEvent(
sse, pageCount, ocrPageCount, terminated, onProgress, onDone, onError))
sse, pageCount, ocrPageCount, terminated,
onProgress, onHeartbeat, onDone, onError))
.blockLast();
if (!terminated[0]) {
onError.accept(new CampaignImportException(
@@ -107,12 +109,19 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
int[] ocrPageCount,
boolean[] terminated,
Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError) {
String event = sse.event();
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)) {
terminated[0] = true;
onError.accept(new CampaignImportException(

View File

@@ -114,6 +114,7 @@ public class BrainRulesImportClient implements RulesPdfImporter {
byte[] pdfBytes,
String filename,
Consumer<RulesImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<RulesImportResult> onDone,
Consumer<Throwable> onError) {
@@ -139,7 +140,8 @@ public class BrainRulesImportClient implements RulesPdfImporter {
flux
.timeout(Duration.ofSeconds(importTimeoutSeconds))
.doOnNext(sse -> handleEvent(
sse, pageCount, ocrPageCount, terminated, onProgress, onDone, onError))
sse, pageCount, ocrPageCount, terminated,
onProgress, onHeartbeat, onDone, onError))
.blockLast();
// Flux terminé sans event done/error (ex: connexion coupée) → on signale.
if (!terminated[0]) {
@@ -165,12 +167,20 @@ public class BrainRulesImportClient implements RulesPdfImporter {
int[] ocrPageCount,
boolean[] terminated,
Consumer<RulesImportProgress> onProgress,
Runnable onHeartbeat,
Consumer<RulesImportResult> onDone,
Consumer<Throwable> onError) {
String event = sse.event();
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)) {
terminated[0] = true;
onError.accept(new RulesImportException(

View File

@@ -10,6 +10,7 @@ import org.springframework.http.converter.HttpMessageNotReadableException;
import org.springframework.web.bind.MethodArgumentNotValidException;
import org.springframework.web.bind.annotation.ExceptionHandler;
import org.springframework.web.bind.annotation.RestControllerAdvice;
import org.springframework.web.context.request.async.AsyncRequestNotUsableException;
import java.util.LinkedHashMap;
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
* un log ERROR explicite (path + stack trace) et un body JSON debuggable

View File

@@ -15,6 +15,7 @@ import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
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.
@@ -30,8 +31,13 @@ public class CampaignImportController {
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 TaskExecutor taskExecutor;
@@ -54,33 +60,63 @@ public class CampaignImportController {
@RequestParam("file") MultipartFile file) throws IOException {
SseEmitter emitter = new SseEmitter(IMPORT_SSE_TIMEOUT_MS);
if (file == null || file.isEmpty()) {
sendError(emitter, "Fichier PDF vide.");
sendError(emitter, new AtomicBoolean(false), "Fichier PDF vide.");
return emitter;
}
byte[] bytes = file.getBytes();
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(() -> {
try {
campaignImportService.importStructureStreaming(
bytes, filename,
progress -> sendEvent(emitter, "progress", progress),
progress -> sendEvent(emitter, clientGone, "progress", progress),
() -> sendHeartbeat(emitter, clientGone),
proposal -> {
sendEvent(emitter, "done", proposal);
sendEvent(emitter, clientGone, "done", proposal);
emitter.complete();
},
error -> {
if (clientGone.get()) {
log.info("Import campagne (stream) interrompu : client déconnecté.");
return;
}
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) {
log.warn("Import campagne (stream) échoué : {}", e.getMessage());
sendError(emitter, e.getMessage());
sendError(emitter, clientGone, e.getMessage());
}
});
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)
public ResponseEntity<CampaignImportService.ApplyResult> apply(
@PathVariable String campaignId,
@@ -96,23 +132,52 @@ public class CampaignImportController {
// --- 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 {
emitter.send(SseEmitter.event().name(eventName).data(
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);
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 {
emitter.send(SseEmitter.event().name("error").data(
objectMapper.writeValueAsString(Map.of(
"message", message != null ? message : "Erreur inconnue.")),
MediaType.APPLICATION_JSON));
emitter.complete();
} catch (IOException e) {
} catch (Exception e) {
clientGone.set(true);
emitter.completeWithError(e);
}
}

View File

@@ -3,7 +3,6 @@ package com.loremind.infrastructure.web.controller;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.loremind.application.gamesystemcontext.GameSystemService;
import com.loremind.domain.gamesystemcontext.GameSystem;
import com.loremind.domain.gamesystemcontext.RulesImportProgress;
import com.loremind.domain.gamesystemcontext.RulesImportResult;
import com.loremind.domain.gamesystemcontext.ports.RulesImportException;
import com.loremind.domain.shared.template.TemplateField;
@@ -27,6 +26,7 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
@RestController
@@ -35,8 +35,13 @@ public class GameSystemController {
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 GameSystemMapper gameSystemMapper;
@@ -135,7 +140,7 @@ public class GameSystemController {
public SseEmitter importRulesStream(@RequestParam("file") MultipartFile file) throws IOException {
SseEmitter emitter = new SseEmitter(IMPORT_SSE_TIMEOUT_MS);
if (file == null || file.isEmpty()) {
sendImportError(emitter, "Fichier PDF vide.");
sendImportError(emitter, new AtomicBoolean(false), "Fichier PDF vide.");
return emitter;
}
// Les octets sont lus sur le thread servlet (le MultipartFile n'est plus
@@ -143,46 +148,108 @@ public class GameSystemController {
byte[] bytes = file.getBytes();
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(() -> {
try {
gameSystemService.importRulesFromPdfStreaming(
bytes, filename,
progress -> sendImportEvent(emitter, "progress", progress),
progress -> sendImportEvent(emitter, clientGone, "progress", progress),
() -> sendImportHeartbeat(emitter, clientGone),
result -> {
sendImportEvent(emitter, "done", result);
sendImportEvent(emitter, clientGone, "done", result);
emitter.complete();
},
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());
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) {
log.warn("Import de règles (stream) échoué : {}", e.getMessage());
sendImportError(emitter, e.getMessage());
sendImportError(emitter, clientGone, e.getMessage());
}
});
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é. */
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 {
emitter.send(SseEmitter.event().name(eventName).data(
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);
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. */
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 {
emitter.send(SseEmitter.event().name("error").data(
objectMapper.writeValueAsString(Map.of(
"message", message != null ? message : "Erreur inconnue.")),
MediaType.APPLICATION_JSON));
emitter.complete();
} catch (IOException e) {
} catch (Exception e) {
clientGone.set(true);
emitter.completeWithError(e);
}
}

4
web/package-lock.json generated
View File

@@ -1,12 +1,12 @@
{
"name": "loremind-web",
"version": "0.12.0-beta",
"version": "0.12.2-beta",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "loremind-web",
"version": "0.12.0-beta",
"version": "0.12.2-beta",
"dependencies": {
"@angular/animations": "^21.2.16",
"@angular/common": "^21.2.16",

View File

@@ -1,6 +1,6 @@
{
"name": "loremind-web",
"version": "0.12.0-beta",
"version": "0.12.2-beta",
"description": "LoreMind Frontend - Angular",
"scripts": {
"ng": "ng",

View File

@@ -61,17 +61,24 @@ export class CampaignImportService {
let buffer = '';
let currentEvent: string | null = null;
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 name = currentEvent ?? 'message';
if (name === 'error') {
let message = 'Échec de l\'import.';
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* défaut */ }
terminated = true;
subscriber.error(new Error(message));
} else if (name === 'progress' || name === 'done') {
try {
const obj = JSON.parse(currentData);
if (name === 'done') {
terminated = true;
subscriber.next({ type: 'done', arcs: obj.arcs ?? [], npcs: obj.npcs ?? [] });
subscriber.complete();
} else {
@@ -105,9 +112,12 @@ export class CampaignImportService {
}
}
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) {
subscriber.error(err);
if (!terminated) subscriber.error(err);
}
}
}

View File

@@ -92,17 +92,24 @@ export class GameSystemService {
let buffer = '';
let currentEvent: string | null = null;
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 name = currentEvent ?? 'message';
if (name === 'error') {
let message = 'Échec de l\'import.';
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* garde le défaut */ }
terminated = true;
subscriber.error(new Error(message));
} else if (name === 'progress' || name === 'done') {
try {
const obj = JSON.parse(currentData);
if (name === 'done') {
terminated = true;
subscriber.next({ type: 'done', ...obj });
subscriber.complete();
} else {
@@ -136,9 +143,12 @@ export class GameSystemService {
}
}
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) {
subscriber.error(err);
if (!terminated) subscriber.error(err);
}
}
}