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
|
||||
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,
|
||||
)
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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).
|
||||
"""
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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__)
|
||||
|
||||
@@ -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>
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
4
web/package-lock.json
generated
@@ -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",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "loremind-web",
|
||||
"version": "0.12.0-beta",
|
||||
"version": "0.12.2-beta",
|
||||
"description": "LoreMind Frontend - Angular",
|
||||
"scripts": {
|
||||
"ng": "ng",
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user