From 8cc90bd24d9a14ae4c68d6fb154d8ee5ad7a28e2 Mon Sep 17 00:00:00 2001 From: "IETM_FIXE\\ietm6" Date: Fri, 12 Jun 2026 14:35:10 +0200 Subject: [PATCH] =?UTF-8?q?Am=C3=A9lioration=20du=20feedback=20pendant=20l?= =?UTF-8?q?es=20imports=20sur=20les=20PDF=20passage=20en=200.12.4-beta?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- brain/app/application/import_campaign.py | 178 ++++++++++-------- brain/app/application/import_rules.py | 82 +++++--- brain/app/application/import_status.py | 39 ++++ brain/app/application/llm_retry.py | 9 + brain/app/application/streaming.py | 24 ++- brain/app/infrastructure/gemini_adapter.py | 10 + brain/app/main.py | 2 +- core/pom.xml | 2 +- .../CampaignImportService.java | 3 +- .../gamesystemcontext/GameSystemService.java | 3 +- .../ports/CampaignPdfImporter.java | 5 + .../ports/RulesPdfImporter.java | 5 + .../ai/BrainCampaignImportClient.java | 20 +- .../ai/BrainRulesImportClient.java | 20 +- .../controller/CampaignImportController.java | 2 + .../web/controller/GameSystemController.java | 2 + web/package-lock.json | 4 +- web/package.json | 2 +- .../campaign-import.component.html | 3 + .../campaign-import.component.scss | 10 + .../campaign-import.component.ts | 12 ++ .../game-system-edit.component.html | 3 + .../game-system-edit.component.scss | 10 + .../game-system-edit.component.ts | 11 ++ web/src/app/services/campaign-import.model.ts | 3 + .../app/services/campaign-import.service.ts | 6 + web/src/app/services/game-system.model.ts | 3 + web/src/app/services/game-system.service.ts | 6 + 28 files changed, 364 insertions(+), 115 deletions(-) create mode 100644 brain/app/application/import_status.py diff --git a/brain/app/application/import_campaign.py b/brain/app/application/import_campaign.py index dd3cbea..00dee07 100644 --- a/brain/app/application/import_campaign.py +++ b/brain/app/application/import_campaign.py @@ -14,6 +14,11 @@ import asyncio import logging from app.application.chunking import chunk_text, split_in_half +from app.application.import_status import ( + notify_status, + reset_status_queue, + set_status_queue, +) from app.application.llm_json import load_json_object, looks_like_truncated_json from app.application.llm_retry import generate_with_retry from app.application.streaming import with_heartbeat @@ -510,86 +515,101 @@ class ImportCampaignUseCase: skipped = 0 last_error: str | None = None done_count = 0 - # PARALLÉLISME : les morceaux sont traités par VAGUES de `map_concurrency` - # appels simultanés. L'ordre narratif est préservé : la fusion se fait - # vague par vague, dans l'ordre du livre. - # RÉSILIENCE : un morceau qui échoue (provider saturé, quota, etc.) est - # SAUTÉ — on ne perd pas tout l'import pour autant. On n'abandonne que - # si AUCUN morceau ne passe (cf. après la boucle). - # HEARTBEAT : keep-alive pendant la vague d'appels LLM pour ne jamais - # laisser le flux SSE silencieux (sinon le Core coupe sur inactivité). - for start in range(0, total, self._map_concurrency): - wave = list(enumerate(chunks))[start:start + self._map_concurrency] - gathered = asyncio.gather( - *(self._map_chunk(c, index=i, total=total, toc_block=toc_block) - for i, c in wave), - return_exceptions=True, - ) - results: list | None = None - async for kind, payload in with_heartbeat(gathered): - if kind == "heartbeat": - yield {"type": "heartbeat", "current": done_count + 1, "total": total} - else: - results = payload - for (i, _), res in zip(wave, results or []): - done_count += 1 - if isinstance(res, LLMProviderError): - skipped += 1 - last_error = str(res) - logger.warning("Morceau %s/%s ignoré (échec LLM) : %s", i + 1, total, res) - yield {"type": "chunk_failed", "current": i + 1, "total": total, - "message": str(res)[:300]} - elif isinstance(res, BaseException): - raise res # bug inattendu : ne pas l'avaler en silence - else: - merger.add((res or {}).get("arcs") or []) - merger.add_npcs((res or {}).get("npcs") or []) - arcs, chapters, scenes = merger.counts() - yield { - "type": "progress", - "current": done_count, - "total": total, - "arc_count": arcs, - "chapter_count": chapters, - "scene_count": scenes, - "npc_count": len(merger.npcs()), - "skipped": skipped, - } + # Canal de statut : les couches profondes (retry LLM, re-découpage) y + # publient des messages destinés à l'UI — cf. import_status.notify_status. + status_queue: asyncio.Queue = asyncio.Queue() + status_token = set_status_queue(status_queue) + try: + # PARALLÉLISME : les morceaux sont traités par VAGUES de `map_concurrency` + # appels simultanés. L'ordre narratif est préservé : la fusion se fait + # vague par vague, dans l'ordre du livre. + # RÉSILIENCE : un morceau qui échoue (provider saturé, quota, etc.) est + # SAUTÉ — on ne perd pas tout l'import pour autant. On n'abandonne que + # si AUCUN morceau ne passe (cf. après la boucle). + # HEARTBEAT : keep-alive pendant la vague d'appels LLM pour ne jamais + # laisser le flux SSE silencieux (sinon le Core coupe sur inactivité). + for start in range(0, total, self._map_concurrency): + wave = list(enumerate(chunks))[start:start + self._map_concurrency] + gathered = asyncio.gather( + *(self._map_chunk(c, index=i, total=total, toc_block=toc_block) + for i, c in wave), + return_exceptions=True, + ) + results: list | None = None + async for kind, payload in with_heartbeat(gathered, status_queue=status_queue): + if kind == "heartbeat": + yield {"type": "heartbeat", "current": done_count + 1, "total": total} + elif kind == "status": + yield {"type": "status", "message": payload, + "current": done_count + 1, "total": total} + else: + results = payload + for (i, _), res in zip(wave, results or []): + done_count += 1 + if isinstance(res, LLMProviderError): + skipped += 1 + last_error = str(res) + logger.warning("Morceau %s/%s ignoré (échec LLM) : %s", i + 1, total, res) + yield {"type": "chunk_failed", "current": i + 1, "total": total, + "message": str(res)[:300]} + elif isinstance(res, BaseException): + raise res # bug inattendu : ne pas l'avaler en silence + else: + merger.add((res or {}).get("arcs") or []) + merger.add_npcs((res or {}).get("npcs") or []) + arcs, chapters, scenes = merger.counts() + yield { + "type": "progress", + "current": done_count, + "total": total, + "arc_count": arcs, + "chapter_count": chapters, + "scene_count": scenes, + "npc_count": len(merger.npcs()), + "skipped": skipped, + } - if total > 0 and skipped == total: - # Tout a échoué : "done" vide serait trompeur → erreur explicite. - yield {"type": "error", - "message": "Tous les morceaux ont échoué auprès du fournisseur IA. " - f"Dernier message : {last_error or 'inconnu'}"} - return + if total > 0 and skipped == total: + # Tout a échoué : "done" vide serait trompeur → erreur explicite. + yield {"type": "error", + "message": "Tous les morceaux ont échoué auprès du fournisseur IA. " + 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 + 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: - yield {"type": "consolidating", "total": total} - async for kind, _ in with_heartbeat(self._consolidate(merger)): - if kind == "heartbeat": - yield {"type": "heartbeat", "current": total, "total": total} + # Consolidation finale : fusion des quasi-doublons inter-morceaux + # (best-effort, voir _consolidate). Inutile sur un import mono-morceau. + if total > 1: + yield {"type": "consolidating", "total": total} + async for kind, payload in with_heartbeat( + self._consolidate(merger), status_queue=status_queue + ): + if kind == "heartbeat": + yield {"type": "heartbeat", "current": total, "total": total} + elif kind == "status": + yield {"type": "status", "message": payload, + "current": total, "total": total} - yield { - "type": "done", - "arcs": _serialize_arcs(merger.result()), - "npcs": [{"name": n.name, "description": n.description} for n in merger.npcs()], - "page_count": doc.page_count, - "ocr_page_count": doc.ocr_page_count, - "skipped": skipped, - } + yield { + "type": "done", + "arcs": _serialize_arcs(merger.result()), + "npcs": [{"name": n.name, "description": n.description} for n in merger.npcs()], + "page_count": doc.page_count, + "ocr_page_count": doc.ocr_page_count, + "skipped": skipped, + } + finally: + reset_status_queue(status_token) # --- Consolidation finale (fusion des quasi-doublons) --------------------- @@ -666,6 +686,9 @@ class ImportCampaignUseCase: logger.info( "Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).", index, depth + 1) + notify_status( + f"Le modèle est trop lent sur le morceau {index + 1} : " + "re-découpage en 2 moitiés plus digestes…") a = await self._extract_payload( left, index=index, total=total, depth=depth + 1, toc_block=toc_block) b = await self._extract_payload( @@ -679,6 +702,9 @@ class ImportCampaignUseCase: logger.info( "Morceau %s : sortie tronquée → re-découpage en 2 moitiés (niveau %s).", index, depth + 1) + notify_status( + f"Réponse du modèle coupée sur le morceau {index + 1} : " + "re-découpage en 2 moitiés plus digestes…") a = await self._extract_payload( left, index=index, total=total, depth=depth + 1, toc_block=toc_block) b = await self._extract_payload( diff --git a/brain/app/application/import_rules.py b/brain/app/application/import_rules.py index 2f5424f..62b4007 100644 --- a/brain/app/application/import_rules.py +++ b/brain/app/application/import_rules.py @@ -15,7 +15,14 @@ from __future__ import annotations import logging import re +import asyncio + from app.application.chunking import CHUNK_TARGET_TOKENS, chunk_text, split_in_half +from app.application.import_status import ( + notify_status, + reset_status_queue, + set_status_queue, +) from app.application.llm_json import load_json_object, looks_like_truncated_json from app.application.llm_retry import generate_with_retry from app.application.streaming import with_heartbeat @@ -332,35 +339,46 @@ class ImportRulesUseCase: merger = _SectionMerger() skipped = 0 last_error: str | None = None - for i, chunk in enumerate(chunks): - # RÉSILIENCE : un morceau qui échoue est SAUTÉ, l'import continue. - # Abandon seulement si AUCUN morceau ne passe (cf. après la boucle). - # HEARTBEAT : on émet des keep-alive pendant l'appel LLM (long sur un - # provider lent) pour que le flux SSE ne soit jamais coupé par le Core. - new_titles: list[str] = [] - try: - sections: dict[str, str] | None = None - async for kind, payload in with_heartbeat( - self._map_chunk(chunk, index=i, total=total) - ): - if kind == "heartbeat": - yield {"type": "heartbeat", "current": i + 1, "total": total} - else: - sections = payload - new_titles = merger.add(sections or {}) - except LLMProviderError as exc: - skipped += 1 - last_error = str(exc) - logger.warning("Morceau %s/%s ignoré (échec LLM) : %s", i + 1, total, exc) - yield {"type": "chunk_failed", "current": i + 1, "total": total, - "message": str(exc)[:300]} - yield { - "type": "progress", - "current": i + 1, - "total": total, - "new_sections": new_titles, - "skipped": skipped, - } + # Canal de statut : les couches profondes (retry LLM, re-découpage) y + # publient des messages destinés à l'UI — cf. import_status.notify_status. + status_queue: asyncio.Queue = asyncio.Queue() + status_token = set_status_queue(status_queue) + try: + for i, chunk in enumerate(chunks): + # RÉSILIENCE : un morceau qui échoue est SAUTÉ, l'import continue. + # Abandon seulement si AUCUN morceau ne passe (cf. après la boucle). + # HEARTBEAT : on émet des keep-alive pendant l'appel LLM (long sur un + # provider lent) pour que le flux SSE ne soit jamais coupé par le Core. + new_titles: list[str] = [] + try: + sections: dict[str, str] | None = None + async for kind, payload in with_heartbeat( + self._map_chunk(chunk, index=i, total=total), + status_queue=status_queue, + ): + if kind == "heartbeat": + yield {"type": "heartbeat", "current": i + 1, "total": total} + elif kind == "status": + yield {"type": "status", "message": payload, + "current": i + 1, "total": total} + else: + sections = payload + new_titles = merger.add(sections or {}) + except LLMProviderError as exc: + skipped += 1 + last_error = str(exc) + logger.warning("Morceau %s/%s ignoré (échec LLM) : %s", i + 1, total, exc) + yield {"type": "chunk_failed", "current": i + 1, "total": total, + "message": str(exc)[:300]} + yield { + "type": "progress", + "current": i + 1, + "total": total, + "new_sections": new_titles, + "skipped": skipped, + } + finally: + reset_status_queue(status_token) if total > 0 and skipped == total: yield {"type": "error", @@ -423,6 +441,9 @@ class ImportRulesUseCase: logger.info( "Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).", index, depth + 1) + notify_status( + f"Le modèle est trop lent sur le morceau {index + 1} : " + "re-découpage en 2 moitiés plus digestes…") 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) @@ -437,6 +458,9 @@ class ImportRulesUseCase: logger.info( "Morceau %s : sortie tronquée → re-découpage en 2 moitiés (niveau %s).", index, depth + 1) + notify_status( + f"Réponse du modèle coupée sur le morceau {index + 1} : " + "re-découpage en 2 moitiés plus digestes…") 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) diff --git a/brain/app/application/import_status.py b/brain/app/application/import_status.py new file mode 100644 index 0000000..b8eba44 --- /dev/null +++ b/brain/app/application/import_status.py @@ -0,0 +1,39 @@ +"""Canal de statut des imports : remonte à l'UI ce qui n'existait qu'en logs. + +Problème résolu : pendant un import, les événements internes (retry parce que +le fournisseur IA est saturé, re-découpage d'un morceau trop gros…) n'étaient +visibles que dans les logs Docker. L'utilisateur voyait une barre de +progression figée sans explication. + +Mécanisme : le flux d'import (use case `stream()`) installe une Queue dans une +ContextVar ; les couches profondes (retry LLM, re-découpage) y publient des +messages via `notify_status()` sans connaître le flux SSE. La ContextVar est +propagée automatiquement aux tâches asyncio enfants → chaque import concurrent +a SA queue, sans couplage ni paramètre à faire transiter partout. +""" +from __future__ import annotations + +import asyncio +from contextvars import ContextVar, Token + +_QUEUE: ContextVar[asyncio.Queue | None] = ContextVar("import_status_queue", default=None) + + +def set_status_queue(queue: asyncio.Queue | None) -> Token: + """Installe la queue de statut pour le contexte courant (et ses tâches filles). + + Renvoie le token à passer à `reset_status_queue` en fin d'import. + """ + return _QUEUE.set(queue) + + +def reset_status_queue(token: Token) -> None: + _QUEUE.reset(token) + + +def notify_status(message: str) -> None: + """Publie un message de statut si un import écoute. No-op sinon (appels + LLM hors import : chat, génération de page…).""" + queue = _QUEUE.get() + if queue is not None: + queue.put_nowait(message) diff --git a/brain/app/application/llm_retry.py b/brain/app/application/llm_retry.py index 773d609..9fb059e 100644 --- a/brain/app/application/llm_retry.py +++ b/brain/app/application/llm_retry.py @@ -14,6 +14,7 @@ import asyncio import logging import re +from app.application.import_status import notify_status from app.domain.ports import LLMGenerationTimeout, LLMProvider, LLMProviderError logger = logging.getLogger(__name__) @@ -103,6 +104,14 @@ async def generate_with_retry( attempt + 1, _ATTEMPTS, " [rate limit]" if _is_rate_limit(exc) else "", exc, wait, ) + # Remonte aussi l'info à l'UI (flux d'import) : sans ça l'utilisateur + # voit une barre figée sans savoir que le fournisseur est saturé. + notify_status( + ("Fournisseur IA saturé (rate limit)" if _is_rate_limit(exc) + else "Appel IA échoué") + + f" — tentative {attempt + 1}/{_ATTEMPTS}, nouvel essai dans {int(wait)}s. " + + str(exc)[:160] + ) await asyncio.sleep(wait) assert last_error is not None raise last_error diff --git a/brain/app/application/streaming.py b/brain/app/application/streaming.py index 2d8d7db..8062812 100644 --- a/brain/app/application/streaming.py +++ b/brain/app/application/streaming.py @@ -25,21 +25,43 @@ async def with_heartbeat( coro: Awaitable[Any], *, interval: float = HEARTBEAT_INTERVAL_SECONDS, + status_queue: "asyncio.Queue | None" = None, ) -> AsyncIterator[tuple[str, Any]]: """Exécute `coro` en émettant ('heartbeat', None) toutes les `interval`s tant qu'elle n'est pas terminée, puis ('result', valeur). + Si `status_queue` est fournie, les messages qui y sont publiés pendant + l'exécution (cf. import_status.notify_status : retry LLM, re-découpage…) + sont émis AU FIL DE L'EAU sous forme ('status', message) — c'est ce qui + permet à l'UI d'expliquer une attente au lieu d'une barre figée. + L'exception éventuelle de `coro` est propagée (re-levée par `task.result()`), donc l'appelant peut l'attraper normalement. Si l'itération est abandonnée (client déconnecté), la tâche sous-jacente est annulée. """ task: asyncio.Task = asyncio.ensure_future(coro) + getter: asyncio.Task | None = None try: while not task.done(): - done, _ = await asyncio.wait({task}, timeout=interval) + waiters: set[asyncio.Task] = {task} + if status_queue is not None and getter is None: + getter = asyncio.ensure_future(status_queue.get()) + if getter is not None: + waiters.add(getter) + done, _ = await asyncio.wait( + waiters, timeout=interval, return_when=asyncio.FIRST_COMPLETED) + if getter is not None and getter in done: + yield ("status", getter.result()) + getter = None # un nouveau get() sera créé au tour suivant if not done: yield ("heartbeat", None) + # Vide les statuts restés en file (publiés juste avant la fin de la tâche). + if status_queue is not None: + while not status_queue.empty(): + yield ("status", status_queue.get_nowait()) yield ("result", task.result()) finally: + if getter is not None and not getter.done(): + getter.cancel() if not task.done(): task.cancel() diff --git a/brain/app/infrastructure/gemini_adapter.py b/brain/app/infrastructure/gemini_adapter.py index 0a2489c..e6d7a2f 100644 --- a/brain/app/infrastructure/gemini_adapter.py +++ b/brain/app/infrastructure/gemini_adapter.py @@ -144,6 +144,16 @@ class GeminiLLMProvider: ) as response: if response.status_code >= 400: detail = (await response.aread()).decode("utf-8", "replace").strip() + # 401/403 = clé rejetée par GOOGLE (pas un problème LoreMind) : + # message actionnable plutôt que le JSON brut de l'API. + if response.status_code in (401, 403): + raise LLMProviderError( + "Erreur Gemini : clé API refusée par Google " + f"(HTTP {response.status_code}). Vérifiez que la clé vient bien " + "de aistudio.google.com (« Get API key ») et qu'elle n'a pas de " + "restrictions (API ou adresse IP) dans la Google Cloud Console. " + f"Détail : {detail[:300]}" + ) raise LLMProviderError( f"Erreur Gemini (HTTP {response.status_code})" + (f" : {detail[:500]}" if detail else "") diff --git a/brain/app/main.py b/brain/app/main.py index 6ee4ce6..4856d19 100644 --- a/brain/app/main.py +++ b/brain/app/main.py @@ -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.3-beta", + version="0.12.4-beta", ) logger = logging.getLogger(__name__) diff --git a/core/pom.xml b/core/pom.xml index bd3be9b..57f5b95 100644 --- a/core/pom.xml +++ b/core/pom.xml @@ -14,7 +14,7 @@ com.loremind loremind-core - 0.12.3-beta + 0.12.4-beta LoreMind Core Backend Core - Architecture Hexagonale diff --git a/core/src/main/java/com/loremind/application/campaigncontext/CampaignImportService.java b/core/src/main/java/com/loremind/application/campaigncontext/CampaignImportService.java index ced84df..80cb078 100644 --- a/core/src/main/java/com/loremind/application/campaigncontext/CampaignImportService.java +++ b/core/src/main/java/com/loremind/application/campaigncontext/CampaignImportService.java @@ -65,10 +65,11 @@ public class CampaignImportService { String filename, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError) { campaignPdfImporter.importCampaignStreaming( - pdfBytes, filename, onProgress, onHeartbeat, onDone, onError); + pdfBytes, filename, onProgress, onHeartbeat, onStatus, onDone, onError); } /** diff --git a/core/src/main/java/com/loremind/application/gamesystemcontext/GameSystemService.java b/core/src/main/java/com/loremind/application/gamesystemcontext/GameSystemService.java index d75c429..0d6f211 100644 --- a/core/src/main/java/com/loremind/application/gamesystemcontext/GameSystemService.java +++ b/core/src/main/java/com/loremind/application/gamesystemcontext/GameSystemService.java @@ -40,10 +40,11 @@ public class GameSystemService { String filename, java.util.function.Consumer onProgress, Runnable onHeartbeat, + java.util.function.Consumer onStatus, java.util.function.Consumer onDone, java.util.function.Consumer onError) { rulesPdfImporter.importRulesStreaming( - pdfBytes, filename, onProgress, onHeartbeat, onDone, onError); + pdfBytes, filename, onProgress, onHeartbeat, onStatus, onDone, onError); } /** diff --git a/core/src/main/java/com/loremind/domain/campaigncontext/ports/CampaignPdfImporter.java b/core/src/main/java/com/loremind/domain/campaigncontext/ports/CampaignPdfImporter.java index ffb8e02..e11ef1d 100644 --- a/core/src/main/java/com/loremind/domain/campaigncontext/ports/CampaignPdfImporter.java +++ b/core/src/main/java/com/loremind/domain/campaigncontext/ports/CampaignPdfImporter.java @@ -19,6 +19,10 @@ public interface CampaignPdfImporter { * @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 onStatus invoqué avec un message lisible quand quelque chose se + * passe pendant l'attente (fournisseur saturé → retry, + * morceau re-découpé, morceau ignoré…) — affiché par l'UI + * pour que l'utilisateur n'ait pas à lire les logs. * @param onDone invoqué une fois avec l'arbre proposé (non persisté). * @param onError invoqué si l'extraction/structuration échoue. */ @@ -27,6 +31,7 @@ public interface CampaignPdfImporter { String filename, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError); } diff --git a/core/src/main/java/com/loremind/domain/gamesystemcontext/ports/RulesPdfImporter.java b/core/src/main/java/com/loremind/domain/gamesystemcontext/ports/RulesPdfImporter.java index 828a83a..6ca54ec 100644 --- a/core/src/main/java/com/loremind/domain/gamesystemcontext/ports/RulesPdfImporter.java +++ b/core/src/main/java/com/loremind/domain/gamesystemcontext/ports/RulesPdfImporter.java @@ -30,6 +30,10 @@ public interface RulesPdfImporter { * @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 onStatus invoqué avec un message lisible quand quelque chose se + * passe pendant l'attente (fournisseur saturé → retry, + * morceau re-découpé, morceau ignoré…) — affiché par l'UI + * pour que l'utilisateur n'ait pas à lire les logs. * @param onDone invoqué une fois avec le résultat final. * @param onError invoqué si l'extraction/structuration échoue. */ @@ -38,6 +42,7 @@ public interface RulesPdfImporter { String filename, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError); } diff --git a/core/src/main/java/com/loremind/infrastructure/ai/BrainCampaignImportClient.java b/core/src/main/java/com/loremind/infrastructure/ai/BrainCampaignImportClient.java index 6870c9d..3c11180 100644 --- a/core/src/main/java/com/loremind/infrastructure/ai/BrainCampaignImportClient.java +++ b/core/src/main/java/com/loremind/infrastructure/ai/BrainCampaignImportClient.java @@ -61,6 +61,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter { String filename, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError) { @@ -85,7 +86,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter { .timeout(Duration.ofSeconds(importTimeoutSeconds)) .doOnNext(sse -> handleEvent( sse, pageCount, ocrPageCount, terminated, - onProgress, onHeartbeat, onDone, onError)) + onProgress, onHeartbeat, onStatus, onDone, onError)) .blockLast(); if (!terminated[0]) { onError.accept(new CampaignImportException( @@ -110,6 +111,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter { boolean[] terminated, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError) { @@ -122,6 +124,22 @@ public class BrainCampaignImportClient implements CampaignPdfImporter { onHeartbeat.run(); return; } + if ("status".equals(event)) { + // Message d'attente lisible (retry sur fournisseur saturé, morceau + // re-découpé…) : affiché par l'UI au lieu de n'exister qu'en logs. + onStatus.accept(readMessage(data)); + return; + } + if ("chunk_failed".equals(event)) { + JsonNode node = readJson(data); + String msg = node != null && node.hasNonNull("message") + ? node.get("message").asText() : ""; + int current = node != null ? node.path("current").asInt() : 0; + int total = node != null ? node.path("total").asInt() : 0; + onStatus.accept("Morceau " + current + "/" + total + " ignoré" + + (msg.isEmpty() ? "." : " : " + msg)); + return; + } if ("error".equals(event)) { terminated[0] = true; onError.accept(new CampaignImportException( diff --git a/core/src/main/java/com/loremind/infrastructure/ai/BrainRulesImportClient.java b/core/src/main/java/com/loremind/infrastructure/ai/BrainRulesImportClient.java index 84831ff..436d19f 100644 --- a/core/src/main/java/com/loremind/infrastructure/ai/BrainRulesImportClient.java +++ b/core/src/main/java/com/loremind/infrastructure/ai/BrainRulesImportClient.java @@ -115,6 +115,7 @@ public class BrainRulesImportClient implements RulesPdfImporter { String filename, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError) { @@ -141,7 +142,7 @@ public class BrainRulesImportClient implements RulesPdfImporter { .timeout(Duration.ofSeconds(importTimeoutSeconds)) .doOnNext(sse -> handleEvent( sse, pageCount, ocrPageCount, terminated, - onProgress, onHeartbeat, onDone, onError)) + onProgress, onHeartbeat, onStatus, onDone, onError)) .blockLast(); // Flux terminé sans event done/error (ex: connexion coupée) → on signale. if (!terminated[0]) { @@ -168,6 +169,7 @@ public class BrainRulesImportClient implements RulesPdfImporter { boolean[] terminated, Consumer onProgress, Runnable onHeartbeat, + Consumer onStatus, Consumer onDone, Consumer onError) { @@ -181,6 +183,22 @@ public class BrainRulesImportClient implements RulesPdfImporter { onHeartbeat.run(); return; } + if ("status".equals(event)) { + // Message d'attente lisible (retry sur fournisseur saturé, morceau + // re-découpé…) : affiché par l'UI au lieu de n'exister qu'en logs. + onStatus.accept(readMessage(data)); + return; + } + if ("chunk_failed".equals(event)) { + JsonNode node = readJson(data); + String msg = node != null && node.hasNonNull("message") + ? node.get("message").asText() : ""; + int current = node != null ? node.path("current").asInt() : 0; + int total = node != null ? node.path("total").asInt() : 0; + onStatus.accept("Morceau " + current + "/" + total + " ignoré" + + (msg.isEmpty() ? "." : " : " + msg)); + return; + } if ("error".equals(event)) { terminated[0] = true; onError.accept(new RulesImportException( diff --git a/core/src/main/java/com/loremind/infrastructure/web/controller/CampaignImportController.java b/core/src/main/java/com/loremind/infrastructure/web/controller/CampaignImportController.java index 085b30e..e0d7c2e 100644 --- a/core/src/main/java/com/loremind/infrastructure/web/controller/CampaignImportController.java +++ b/core/src/main/java/com/loremind/infrastructure/web/controller/CampaignImportController.java @@ -88,6 +88,8 @@ public class CampaignImportController { bytes, filename, progress -> sendEvent(emitter, clientGone, "progress", progress), () -> sendHeartbeat(emitter, clientGone), + status -> sendEvent(emitter, clientGone, "status", + Map.of("message", status != null ? status : "")), proposal -> { sendEvent(emitter, clientGone, "done", proposal); emitter.complete(); diff --git a/core/src/main/java/com/loremind/infrastructure/web/controller/GameSystemController.java b/core/src/main/java/com/loremind/infrastructure/web/controller/GameSystemController.java index 8adf53b..c1a7c36 100644 --- a/core/src/main/java/com/loremind/infrastructure/web/controller/GameSystemController.java +++ b/core/src/main/java/com/loremind/infrastructure/web/controller/GameSystemController.java @@ -170,6 +170,8 @@ public class GameSystemController { bytes, filename, progress -> sendImportEvent(emitter, clientGone, "progress", progress), () -> sendImportHeartbeat(emitter, clientGone), + status -> sendImportEvent(emitter, clientGone, "status", + Map.of("message", status != null ? status : "")), result -> { sendImportEvent(emitter, clientGone, "done", result); emitter.complete(); diff --git a/web/package-lock.json b/web/package-lock.json index d97f91c..d3dfe14 100644 --- a/web/package-lock.json +++ b/web/package-lock.json @@ -1,12 +1,12 @@ { "name": "loremind-web", - "version": "0.12.3-beta", + "version": "0.12.4-beta", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "loremind-web", - "version": "0.12.3-beta", + "version": "0.12.4-beta", "dependencies": { "@angular/animations": "^21.2.16", "@angular/common": "^21.2.16", diff --git a/web/package.json b/web/package.json index ac7e471..8d3d99e 100644 --- a/web/package.json +++ b/web/package.json @@ -1,6 +1,6 @@ { "name": "loremind-web", - "version": "0.12.3-beta", + "version": "0.12.4-beta", "description": "LoreMind Frontend - Angular", "scripts": { "ng": "ng", diff --git a/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.html b/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.html index 70d498c..5b686e9 100644 --- a/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.html +++ b/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.html @@ -31,6 +31,9 @@ } + @if (importStatus) { +

{{ importStatus }}

+ } @if (importCounts) {

Trouvé jusqu'ici : {{ importCounts.arcs }} arc(s) · {{ importCounts.chapters }} chapitre(s) · {{ importCounts.scenes }} scène(s) · {{ importCounts.npcs }} PNJ diff --git a/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.scss b/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.scss index 82fc027..aef67f9 100644 --- a/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.scss +++ b/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.scss @@ -87,6 +87,16 @@ .import-counts { margin: 0.55rem 0 0; color: #9ca3af; font-size: 0.8rem; } +// Message d'attente live (fournisseur saturé → retry, morceau re-découpé…) : +// ambre pour signaler « ça travaille, mais il se passe quelque chose ». +.import-status { + margin: 0.55rem 0 0; + color: #fbbf24; + font-size: 0.8rem; + font-style: italic; + line-height: 1.4; +} + .import-error, .apply-error { margin: 1rem 0 0; diff --git a/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.ts b/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.ts index 3b26920..b791005 100644 --- a/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.ts +++ b/web/src/app/campaigns/campaign/campaign-import/campaign-import.component.ts @@ -74,6 +74,12 @@ export class CampaignImportComponent implements OnInit { importing = false; importPhase = ''; importProgress: { current: number; total: number } | null = null; + /** + * Dernier message de statut du flux (fournisseur saturé → retry, morceau + * re-découpé/ignoré…). Effacé à chaque progression : il explique l'ATTENTE + * en cours, pas l'historique. + */ + importStatus: string | null = null; importCounts: { arcs: number; chapters: number; scenes: number; npcs: number } | null = null; importError: string | null = null; /** Vrai une fois la proposition reçue (on affiche l'arbre éditable). */ @@ -132,12 +138,15 @@ export class CampaignImportComponent implements OnInit { this.applyError = null; this.importPhase = 'Extraction du texte…'; this.importProgress = null; + this.importStatus = null; this.importCounts = null; this.tree = []; this.service.importStructureStream(this.campaignId, file).subscribe({ next: (ev) => { if (ev.type === 'progress') { + // Un morceau vient d'aboutir : le message d'attente est obsolète. + this.importStatus = null; if (ev.total === 0) { this.importPhase = 'Extraction du texte…'; this.importProgress = null; @@ -149,10 +158,13 @@ export class CampaignImportComponent implements OnInit { scenes: ev.sceneCount, npcs: ev.npcCount ?? 0 }; } + } else if (ev.type === 'status') { + this.importStatus = ev.message; } else if (ev.type === 'done') { this.importing = false; this.importPhase = ''; this.importProgress = null; + this.importStatus = null; if ((ev.arcs ?? []).length === 0 && (ev.npcs ?? []).length === 0) { this.importError = "Aucune structure narrative détectée dans ce PDF."; this.reviewing = false; diff --git a/web/src/app/game-systems/game-system-edit/game-system-edit.component.html b/web/src/app/game-systems/game-system-edit/game-system-edit.component.html index 442e943..b03d50d 100644 --- a/web/src/app/game-systems/game-system-edit/game-system-edit.component.html +++ b/web/src/app/game-systems/game-system-edit/game-system-edit.component.html @@ -60,6 +60,9 @@ } + @if (importStatus) { +

{{ importStatus }}

+ } @if (importFound.length) {

Sections trouvées : {{ importFound.join(' · ') }} diff --git a/web/src/app/game-systems/game-system-edit/game-system-edit.component.scss b/web/src/app/game-systems/game-system-edit/game-system-edit.component.scss index 5ec3253..5835413 100644 --- a/web/src/app/game-systems/game-system-edit/game-system-edit.component.scss +++ b/web/src/app/game-systems/game-system-edit/game-system-edit.component.scss @@ -163,6 +163,16 @@ line-height: 1.4; } +// Message d'attente live (fournisseur saturé → retry, morceau re-découpé…) : +// ambre pour signaler « ça travaille, mais il se passe quelque chose ». +.import-status { + margin: 0.55rem 0 0; + color: #fbbf24; + font-size: 0.8rem; + font-style: italic; + line-height: 1.4; +} + .import-note { margin: -0.4rem 0 1rem; padding: 0.55rem 0.8rem; diff --git a/web/src/app/game-systems/game-system-edit/game-system-edit.component.ts b/web/src/app/game-systems/game-system-edit/game-system-edit.component.ts index 434e0fc..8afecba 100644 --- a/web/src/app/game-systems/game-system-edit/game-system-edit.component.ts +++ b/web/src/app/game-systems/game-system-edit/game-system-edit.component.ts @@ -69,6 +69,12 @@ export class GameSystemEditComponent implements OnInit { importProgress: { current: number; total: number } | null = null; /** Titres de sections trouvés au fil de l'eau (affichage live). */ importFound: string[] = []; + /** + * Dernier message de statut du flux (fournisseur saturé → retry, morceau + * re-découpé/ignoré…). Effacé à chaque progression : il explique l'ATTENTE + * en cours, pas l'historique. + */ + importStatus: string | null = null; name = ''; description = ''; @@ -139,10 +145,13 @@ export class GameSystemEditComponent implements OnInit { this.importPhase = 'Extraction du texte…'; this.importProgress = null; this.importFound = []; + this.importStatus = null; this.service.importRulesStream(file).subscribe({ next: (ev) => { if (ev.type === 'progress') { + // Un morceau vient d'aboutir : le message d'attente est obsolète. + this.importStatus = null; if (ev.total === 0) { // Phase d'extraction (total encore inconnu). this.importPhase = 'Extraction du texte…'; @@ -154,6 +163,8 @@ export class GameSystemEditComponent implements OnInit { if (!this.importFound.includes(t)) this.importFound.push(t); } } + } else if (ev.type === 'status') { + this.importStatus = ev.message; } else if (ev.type === 'done') { this.finishImport(ev.sections, ev.pageCount, ev.ocrPageCount); } diff --git a/web/src/app/services/campaign-import.model.ts b/web/src/app/services/campaign-import.model.ts index 84adb95..7b4202d 100644 --- a/web/src/app/services/campaign-import.model.ts +++ b/web/src/app/services/campaign-import.model.ts @@ -63,10 +63,13 @@ export interface CampaignImportApplyResult { /** * Évènements du flux SSE d'import streamé. * - progress : avancement (total=0 ⇒ extraction en cours). + * - status : message d'attente lisible (fournisseur saturé → retry, morceau + * re-découpé, morceau ignoré…) — feedback live pour l'utilisateur. * - done : arbre proposé (à réviser). * - error : message d'erreur côté serveur. */ export type CampaignImportStreamEvent = + | { type: 'status'; message: string } | { type: 'progress'; current: number; diff --git a/web/src/app/services/campaign-import.service.ts b/web/src/app/services/campaign-import.service.ts index f5cd6b1..d6f8888 100644 --- a/web/src/app/services/campaign-import.service.ts +++ b/web/src/app/services/campaign-import.service.ts @@ -74,6 +74,12 @@ export class CampaignImportService { try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* défaut */ } terminated = true; subscriber.error(new Error(message)); + } else if (name === 'status') { + // Message d'attente lisible (fournisseur saturé, morceau re-découpé…). + try { + const obj = JSON.parse(currentData) as { message?: string }; + if (obj.message) subscriber.next({ type: 'status', message: obj.message }); + } catch { /* bloc malformé ignoré */ } } else if (name === 'progress' || name === 'done') { try { const obj = JSON.parse(currentData); diff --git a/web/src/app/services/game-system.model.ts b/web/src/app/services/game-system.model.ts index fbb5810..b1ea099 100644 --- a/web/src/app/services/game-system.model.ts +++ b/web/src/app/services/game-system.model.ts @@ -32,11 +32,14 @@ export interface RulesImportResponse { /** * Évènements du flux SSE d'import streamé. * - progress : avancement (total=0 ⇒ phase d'extraction en cours). + * - status : message d'attente lisible (fournisseur saturé → retry, morceau + * re-découpé, morceau ignoré…) — feedback live pour l'utilisateur. * - done : résultat final (sections proposées). * - error : message d'erreur côté serveur. */ export type RulesImportStreamEvent = | { type: 'progress'; current: number; total: number; pageCount: number; ocrPageCount: number; newSectionTitles: string[] } + | { type: 'status'; message: string } | { type: 'done'; sections: Record; pageCount: number; ocrPageCount: number } | { type: 'error'; message: string }; diff --git a/web/src/app/services/game-system.service.ts b/web/src/app/services/game-system.service.ts index 8dfb49c..77a50ee 100644 --- a/web/src/app/services/game-system.service.ts +++ b/web/src/app/services/game-system.service.ts @@ -105,6 +105,12 @@ export class GameSystemService { 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 === 'status') { + // Message d'attente lisible (fournisseur saturé, morceau re-découpé…). + try { + const obj = JSON.parse(currentData) as { message?: string }; + if (obj.message) subscriber.next({ type: 'status', message: obj.message }); + } catch { /* bloc malformé ignoré */ } } else if (name === 'progress' || name === 'done') { try { const obj = JSON.parse(currentData);