Compare commits

...

1 Commits

Author SHA1 Message Date
8cc90bd24d Amélioration du feedback pendant les imports sur les PDF
All checks were successful
Build & Push Images / build (brain) (push) Successful in 1m38s
Build & Push Images / build (core) (push) Successful in 1m59s
Build & Push Images / build-switcher (push) Successful in 17s
Build & Push Images / build (web) (push) Successful in 1m54s
passage en 0.12.4-beta
2026-06-12 14:35:10 +02:00
28 changed files with 364 additions and 115 deletions

View File

@@ -14,6 +14,11 @@ import asyncio
import logging import logging
from app.application.chunking import chunk_text, split_in_half 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_json import load_json_object, looks_like_truncated_json
from app.application.llm_retry import generate_with_retry from app.application.llm_retry import generate_with_retry
from app.application.streaming import with_heartbeat from app.application.streaming import with_heartbeat
@@ -510,6 +515,11 @@ class ImportCampaignUseCase:
skipped = 0 skipped = 0
last_error: str | None = None last_error: str | None = None
done_count = 0 done_count = 0
# 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` # 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 # appels simultanés. L'ordre narratif est préservé : la fusion se fait
# vague par vague, dans l'ordre du livre. # vague par vague, dans l'ordre du livre.
@@ -526,9 +536,12 @@ class ImportCampaignUseCase:
return_exceptions=True, return_exceptions=True,
) )
results: list | None = None results: list | None = None
async for kind, payload in with_heartbeat(gathered): async for kind, payload in with_heartbeat(gathered, status_queue=status_queue):
if kind == "heartbeat": if kind == "heartbeat":
yield {"type": "heartbeat", "current": done_count + 1, "total": total} yield {"type": "heartbeat", "current": done_count + 1, "total": total}
elif kind == "status":
yield {"type": "status", "message": payload,
"current": done_count + 1, "total": total}
else: else:
results = payload results = payload
for (i, _), res in zip(wave, results or []): for (i, _), res in zip(wave, results or []):
@@ -578,9 +591,14 @@ class ImportCampaignUseCase:
# (best-effort, voir _consolidate). Inutile sur un import mono-morceau. # (best-effort, voir _consolidate). Inutile sur un import mono-morceau.
if total > 1: if total > 1:
yield {"type": "consolidating", "total": total} yield {"type": "consolidating", "total": total}
async for kind, _ in with_heartbeat(self._consolidate(merger)): async for kind, payload in with_heartbeat(
self._consolidate(merger), status_queue=status_queue
):
if kind == "heartbeat": if kind == "heartbeat":
yield {"type": "heartbeat", "current": total, "total": total} yield {"type": "heartbeat", "current": total, "total": total}
elif kind == "status":
yield {"type": "status", "message": payload,
"current": total, "total": total}
yield { yield {
"type": "done", "type": "done",
@@ -590,6 +608,8 @@ class ImportCampaignUseCase:
"ocr_page_count": doc.ocr_page_count, "ocr_page_count": doc.ocr_page_count,
"skipped": skipped, "skipped": skipped,
} }
finally:
reset_status_queue(status_token)
# --- Consolidation finale (fusion des quasi-doublons) --------------------- # --- Consolidation finale (fusion des quasi-doublons) ---------------------
@@ -666,6 +686,9 @@ class ImportCampaignUseCase:
logger.info( logger.info(
"Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).", "Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).",
index, depth + 1) 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( a = await self._extract_payload(
left, index=index, total=total, depth=depth + 1, toc_block=toc_block) left, index=index, total=total, depth=depth + 1, toc_block=toc_block)
b = await self._extract_payload( b = await self._extract_payload(
@@ -679,6 +702,9 @@ class ImportCampaignUseCase:
logger.info( logger.info(
"Morceau %s : sortie tronquée → re-découpage en 2 moitiés (niveau %s).", "Morceau %s : sortie tronquée → re-découpage en 2 moitiés (niveau %s).",
index, depth + 1) 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( a = await self._extract_payload(
left, index=index, total=total, depth=depth + 1, toc_block=toc_block) left, index=index, total=total, depth=depth + 1, toc_block=toc_block)
b = await self._extract_payload( b = await self._extract_payload(

View File

@@ -15,7 +15,14 @@ from __future__ import annotations
import logging import logging
import re import re
import asyncio
from app.application.chunking import CHUNK_TARGET_TOKENS, chunk_text, split_in_half from app.application.chunking import CHUNK_TARGET_TOKENS, chunk_text, split_in_half
from app.application.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_json import load_json_object, looks_like_truncated_json
from app.application.llm_retry import generate_with_retry from app.application.llm_retry import generate_with_retry
from app.application.streaming import with_heartbeat from app.application.streaming import with_heartbeat
@@ -332,6 +339,11 @@ class ImportRulesUseCase:
merger = _SectionMerger() merger = _SectionMerger()
skipped = 0 skipped = 0
last_error: str | None = None last_error: str | None = None
# 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): for i, chunk in enumerate(chunks):
# RÉSILIENCE : un morceau qui échoue est SAUTÉ, l'import continue. # RÉSILIENCE : un morceau qui échoue est SAUTÉ, l'import continue.
# Abandon seulement si AUCUN morceau ne passe (cf. après la boucle). # Abandon seulement si AUCUN morceau ne passe (cf. après la boucle).
@@ -341,10 +353,14 @@ class ImportRulesUseCase:
try: try:
sections: dict[str, str] | None = None sections: dict[str, str] | None = None
async for kind, payload in with_heartbeat( async for kind, payload in with_heartbeat(
self._map_chunk(chunk, index=i, total=total) self._map_chunk(chunk, index=i, total=total),
status_queue=status_queue,
): ):
if kind == "heartbeat": if kind == "heartbeat":
yield {"type": "heartbeat", "current": i + 1, "total": total} yield {"type": "heartbeat", "current": i + 1, "total": total}
elif kind == "status":
yield {"type": "status", "message": payload,
"current": i + 1, "total": total}
else: else:
sections = payload sections = payload
new_titles = merger.add(sections or {}) new_titles = merger.add(sections or {})
@@ -361,6 +377,8 @@ class ImportRulesUseCase:
"new_sections": new_titles, "new_sections": new_titles,
"skipped": skipped, "skipped": skipped,
} }
finally:
reset_status_queue(status_token)
if total > 0 and skipped == total: if total > 0 and skipped == total:
yield {"type": "error", yield {"type": "error",
@@ -423,6 +441,9 @@ class ImportRulesUseCase:
logger.info( logger.info(
"Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).", "Morceau %s : timeout de génération → re-découpage en 2 moitiés (niveau %s).",
index, depth + 1) 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) 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) b = await self._extract_sections(right, index=index, total=total, depth=depth + 1)
return _combine_sections(a, b) return _combine_sections(a, b)
@@ -437,6 +458,9 @@ class ImportRulesUseCase:
logger.info( logger.info(
"Morceau %s : sortie tronquée → re-découpage en 2 moitiés (niveau %s).", "Morceau %s : sortie tronquée → re-découpage en 2 moitiés (niveau %s).",
index, depth + 1) 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) 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) b = await self._extract_sections(right, index=index, total=total, depth=depth + 1)
return _combine_sections(a, b) return _combine_sections(a, b)

View File

@@ -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)

View File

@@ -14,6 +14,7 @@ import asyncio
import logging import logging
import re import re
from app.application.import_status import notify_status
from app.domain.ports import LLMGenerationTimeout, LLMProvider, LLMProviderError from app.domain.ports import LLMGenerationTimeout, LLMProvider, LLMProviderError
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -103,6 +104,14 @@ async def generate_with_retry(
attempt + 1, _ATTEMPTS, " [rate limit]" if _is_rate_limit(exc) else "", attempt + 1, _ATTEMPTS, " [rate limit]" if _is_rate_limit(exc) else "",
exc, wait, 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) await asyncio.sleep(wait)
assert last_error is not None assert last_error is not None
raise last_error raise last_error

View File

@@ -25,21 +25,43 @@ async def with_heartbeat(
coro: Awaitable[Any], coro: Awaitable[Any],
*, *,
interval: float = HEARTBEAT_INTERVAL_SECONDS, interval: float = HEARTBEAT_INTERVAL_SECONDS,
status_queue: "asyncio.Queue | None" = None,
) -> AsyncIterator[tuple[str, Any]]: ) -> AsyncIterator[tuple[str, Any]]:
"""Exécute `coro` en émettant ('heartbeat', None) toutes les `interval`s tant """Exécute `coro` en émettant ('heartbeat', None) toutes les `interval`s tant
qu'elle n'est pas terminée, puis ('result', valeur). 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()`), 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 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. (client déconnecté), la tâche sous-jacente est annulée.
""" """
task: asyncio.Task = asyncio.ensure_future(coro) task: asyncio.Task = asyncio.ensure_future(coro)
getter: asyncio.Task | None = None
try: try:
while not task.done(): 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: if not done:
yield ("heartbeat", None) 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()) yield ("result", task.result())
finally: finally:
if getter is not None and not getter.done():
getter.cancel()
if not task.done(): if not task.done():
task.cancel() task.cancel()

View File

@@ -144,6 +144,16 @@ class GeminiLLMProvider:
) as response: ) as response:
if response.status_code >= 400: if response.status_code >= 400:
detail = (await response.aread()).decode("utf-8", "replace").strip() 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( raise LLMProviderError(
f"Erreur Gemini (HTTP {response.status_code})" f"Erreur Gemini (HTTP {response.status_code})"
+ (f" : {detail[:500]}" if detail else "") + (f" : {detail[:500]}" if detail else "")

View File

@@ -26,7 +26,7 @@ from app.infrastructure.ollama_model_installer import ensure_ollama_embedding_mo
app = FastAPI( app = FastAPI(
title="LoreMind Brain", title="LoreMind Brain",
description="Backend IA pour la génération de contenu narratif.", description="Backend IA pour la génération de contenu narratif.",
version="0.12.3-beta", version="0.12.4-beta",
) )
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)

View File

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

View File

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

View File

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

View File

@@ -19,6 +19,10 @@ public interface CampaignPdfImporter {
* @param onHeartbeat invoqué périodiquement pendant un appel LLM long (aucune * @param onHeartbeat invoqué périodiquement pendant un appel LLM long (aucune
* avancée à afficher, mais le canal SSE vers le navigateur * avancée à afficher, mais le canal SSE vers le navigateur
* doit rester actif — sinon un proxy intermédiaire le coupe). * 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 onDone invoqué une fois avec l'arbre proposé (non persisté).
* @param onError invoqué si l'extraction/structuration échoue. * @param onError invoqué si l'extraction/structuration échoue.
*/ */
@@ -27,6 +31,7 @@ public interface CampaignPdfImporter {
String filename, String filename,
Consumer<CampaignImportProgress> onProgress, Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat, Runnable onHeartbeat,
Consumer<String> onStatus,
Consumer<CampaignImportProposal> onDone, Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError); Consumer<Throwable> onError);
} }

View File

@@ -30,6 +30,10 @@ public interface RulesPdfImporter {
* @param onHeartbeat invoqué périodiquement pendant un appel LLM long (aucune * @param onHeartbeat invoqué périodiquement pendant un appel LLM long (aucune
* avancée à afficher, mais le canal SSE vers le navigateur * avancée à afficher, mais le canal SSE vers le navigateur
* doit rester actif — sinon un proxy intermédiaire le coupe). * 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 onDone invoqué une fois avec le résultat final.
* @param onError invoqué si l'extraction/structuration échoue. * @param onError invoqué si l'extraction/structuration échoue.
*/ */
@@ -38,6 +42,7 @@ public interface RulesPdfImporter {
String filename, String filename,
Consumer<RulesImportProgress> onProgress, Consumer<RulesImportProgress> onProgress,
Runnable onHeartbeat, Runnable onHeartbeat,
Consumer<String> onStatus,
Consumer<RulesImportResult> onDone, Consumer<RulesImportResult> onDone,
Consumer<Throwable> onError); Consumer<Throwable> onError);
} }

View File

@@ -61,6 +61,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
String filename, String filename,
Consumer<CampaignImportProgress> onProgress, Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat, Runnable onHeartbeat,
Consumer<String> onStatus,
Consumer<CampaignImportProposal> onDone, Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError) { Consumer<Throwable> onError) {
@@ -85,7 +86,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
.timeout(Duration.ofSeconds(importTimeoutSeconds)) .timeout(Duration.ofSeconds(importTimeoutSeconds))
.doOnNext(sse -> handleEvent( .doOnNext(sse -> handleEvent(
sse, pageCount, ocrPageCount, terminated, sse, pageCount, ocrPageCount, terminated,
onProgress, onHeartbeat, onDone, onError)) onProgress, onHeartbeat, onStatus, onDone, onError))
.blockLast(); .blockLast();
if (!terminated[0]) { if (!terminated[0]) {
onError.accept(new CampaignImportException( onError.accept(new CampaignImportException(
@@ -110,6 +111,7 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
boolean[] terminated, boolean[] terminated,
Consumer<CampaignImportProgress> onProgress, Consumer<CampaignImportProgress> onProgress,
Runnable onHeartbeat, Runnable onHeartbeat,
Consumer<String> onStatus,
Consumer<CampaignImportProposal> onDone, Consumer<CampaignImportProposal> onDone,
Consumer<Throwable> onError) { Consumer<Throwable> onError) {
@@ -122,6 +124,22 @@ public class BrainCampaignImportClient implements CampaignPdfImporter {
onHeartbeat.run(); onHeartbeat.run();
return; 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)) { if ("error".equals(event)) {
terminated[0] = true; terminated[0] = true;
onError.accept(new CampaignImportException( onError.accept(new CampaignImportException(

View File

@@ -115,6 +115,7 @@ public class BrainRulesImportClient implements RulesPdfImporter {
String filename, String filename,
Consumer<RulesImportProgress> onProgress, Consumer<RulesImportProgress> onProgress,
Runnable onHeartbeat, Runnable onHeartbeat,
Consumer<String> onStatus,
Consumer<RulesImportResult> onDone, Consumer<RulesImportResult> onDone,
Consumer<Throwable> onError) { Consumer<Throwable> onError) {
@@ -141,7 +142,7 @@ public class BrainRulesImportClient implements RulesPdfImporter {
.timeout(Duration.ofSeconds(importTimeoutSeconds)) .timeout(Duration.ofSeconds(importTimeoutSeconds))
.doOnNext(sse -> handleEvent( .doOnNext(sse -> handleEvent(
sse, pageCount, ocrPageCount, terminated, sse, pageCount, ocrPageCount, terminated,
onProgress, onHeartbeat, onDone, onError)) onProgress, onHeartbeat, onStatus, onDone, onError))
.blockLast(); .blockLast();
// Flux terminé sans event done/error (ex: connexion coupée) → on signale. // Flux terminé sans event done/error (ex: connexion coupée) → on signale.
if (!terminated[0]) { if (!terminated[0]) {
@@ -168,6 +169,7 @@ public class BrainRulesImportClient implements RulesPdfImporter {
boolean[] terminated, boolean[] terminated,
Consumer<RulesImportProgress> onProgress, Consumer<RulesImportProgress> onProgress,
Runnable onHeartbeat, Runnable onHeartbeat,
Consumer<String> onStatus,
Consumer<RulesImportResult> onDone, Consumer<RulesImportResult> onDone,
Consumer<Throwable> onError) { Consumer<Throwable> onError) {
@@ -181,6 +183,22 @@ public class BrainRulesImportClient implements RulesPdfImporter {
onHeartbeat.run(); onHeartbeat.run();
return; 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)) { if ("error".equals(event)) {
terminated[0] = true; terminated[0] = true;
onError.accept(new RulesImportException( onError.accept(new RulesImportException(

View File

@@ -88,6 +88,8 @@ public class CampaignImportController {
bytes, filename, bytes, filename,
progress -> sendEvent(emitter, clientGone, "progress", progress), progress -> sendEvent(emitter, clientGone, "progress", progress),
() -> sendHeartbeat(emitter, clientGone), () -> sendHeartbeat(emitter, clientGone),
status -> sendEvent(emitter, clientGone, "status",
Map.of("message", status != null ? status : "")),
proposal -> { proposal -> {
sendEvent(emitter, clientGone, "done", proposal); sendEvent(emitter, clientGone, "done", proposal);
emitter.complete(); emitter.complete();

View File

@@ -170,6 +170,8 @@ public class GameSystemController {
bytes, filename, bytes, filename,
progress -> sendImportEvent(emitter, clientGone, "progress", progress), progress -> sendImportEvent(emitter, clientGone, "progress", progress),
() -> sendImportHeartbeat(emitter, clientGone), () -> sendImportHeartbeat(emitter, clientGone),
status -> sendImportEvent(emitter, clientGone, "status",
Map.of("message", status != null ? status : "")),
result -> { result -> {
sendImportEvent(emitter, clientGone, "done", result); sendImportEvent(emitter, clientGone, "done", result);
emitter.complete(); emitter.complete();

4
web/package-lock.json generated
View File

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

View File

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

View File

@@ -31,6 +31,9 @@
</div> </div>
</div> </div>
} }
@if (importStatus) {
<p class="import-status" role="status">{{ importStatus }}</p>
}
@if (importCounts) { @if (importCounts) {
<p class="import-counts"> <p class="import-counts">
Trouvé jusqu'ici : {{ importCounts.arcs }} arc(s) · {{ importCounts.chapters }} chapitre(s) · {{ importCounts.scenes }} scène(s) · {{ importCounts.npcs }} PNJ Trouvé jusqu'ici : {{ importCounts.arcs }} arc(s) · {{ importCounts.chapters }} chapitre(s) · {{ importCounts.scenes }} scène(s) · {{ importCounts.npcs }} PNJ

View File

@@ -87,6 +87,16 @@
.import-counts { margin: 0.55rem 0 0; color: #9ca3af; font-size: 0.8rem; } .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, .import-error,
.apply-error { .apply-error {
margin: 1rem 0 0; margin: 1rem 0 0;

View File

@@ -74,6 +74,12 @@ export class CampaignImportComponent implements OnInit {
importing = false; importing = false;
importPhase = ''; importPhase = '';
importProgress: { current: number; total: number } | null = null; 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; importCounts: { arcs: number; chapters: number; scenes: number; npcs: number } | null = null;
importError: string | null = null; importError: string | null = null;
/** Vrai une fois la proposition reçue (on affiche l'arbre éditable). */ /** 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.applyError = null;
this.importPhase = 'Extraction du texte…'; this.importPhase = 'Extraction du texte…';
this.importProgress = null; this.importProgress = null;
this.importStatus = null;
this.importCounts = null; this.importCounts = null;
this.tree = []; this.tree = [];
this.service.importStructureStream(this.campaignId, file).subscribe({ this.service.importStructureStream(this.campaignId, file).subscribe({
next: (ev) => { next: (ev) => {
if (ev.type === 'progress') { if (ev.type === 'progress') {
// Un morceau vient d'aboutir : le message d'attente est obsolète.
this.importStatus = null;
if (ev.total === 0) { if (ev.total === 0) {
this.importPhase = 'Extraction du texte…'; this.importPhase = 'Extraction du texte…';
this.importProgress = null; this.importProgress = null;
@@ -149,10 +158,13 @@ export class CampaignImportComponent implements OnInit {
scenes: ev.sceneCount, npcs: ev.npcCount ?? 0 scenes: ev.sceneCount, npcs: ev.npcCount ?? 0
}; };
} }
} else if (ev.type === 'status') {
this.importStatus = ev.message;
} else if (ev.type === 'done') { } else if (ev.type === 'done') {
this.importing = false; this.importing = false;
this.importPhase = ''; this.importPhase = '';
this.importProgress = null; this.importProgress = null;
this.importStatus = null;
if ((ev.arcs ?? []).length === 0 && (ev.npcs ?? []).length === 0) { if ((ev.arcs ?? []).length === 0 && (ev.npcs ?? []).length === 0) {
this.importError = "Aucune structure narrative détectée dans ce PDF."; this.importError = "Aucune structure narrative détectée dans ce PDF.";
this.reviewing = false; this.reviewing = false;

View File

@@ -60,6 +60,9 @@
</div> </div>
</div> </div>
} }
@if (importStatus) {
<p class="import-status" role="status">{{ importStatus }}</p>
}
@if (importFound.length) { @if (importFound.length) {
<p class="import-found"> <p class="import-found">
Sections trouvées : {{ importFound.join(' · ') }} Sections trouvées : {{ importFound.join(' · ') }}

View File

@@ -163,6 +163,16 @@
line-height: 1.4; 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 { .import-note {
margin: -0.4rem 0 1rem; margin: -0.4rem 0 1rem;
padding: 0.55rem 0.8rem; padding: 0.55rem 0.8rem;

View File

@@ -69,6 +69,12 @@ export class GameSystemEditComponent implements OnInit {
importProgress: { current: number; total: number } | null = null; importProgress: { current: number; total: number } | null = null;
/** Titres de sections trouvés au fil de l'eau (affichage live). */ /** Titres de sections trouvés au fil de l'eau (affichage live). */
importFound: string[] = []; 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 = ''; name = '';
description = ''; description = '';
@@ -139,10 +145,13 @@ export class GameSystemEditComponent implements OnInit {
this.importPhase = 'Extraction du texte…'; this.importPhase = 'Extraction du texte…';
this.importProgress = null; this.importProgress = null;
this.importFound = []; this.importFound = [];
this.importStatus = null;
this.service.importRulesStream(file).subscribe({ this.service.importRulesStream(file).subscribe({
next: (ev) => { next: (ev) => {
if (ev.type === 'progress') { if (ev.type === 'progress') {
// Un morceau vient d'aboutir : le message d'attente est obsolète.
this.importStatus = null;
if (ev.total === 0) { if (ev.total === 0) {
// Phase d'extraction (total encore inconnu). // Phase d'extraction (total encore inconnu).
this.importPhase = 'Extraction du texte…'; this.importPhase = 'Extraction du texte…';
@@ -154,6 +163,8 @@ export class GameSystemEditComponent implements OnInit {
if (!this.importFound.includes(t)) this.importFound.push(t); if (!this.importFound.includes(t)) this.importFound.push(t);
} }
} }
} else if (ev.type === 'status') {
this.importStatus = ev.message;
} else if (ev.type === 'done') { } else if (ev.type === 'done') {
this.finishImport(ev.sections, ev.pageCount, ev.ocrPageCount); this.finishImport(ev.sections, ev.pageCount, ev.ocrPageCount);
} }

View File

@@ -63,10 +63,13 @@ export interface CampaignImportApplyResult {
/** /**
* Évènements du flux SSE d'import streamé. * Évènements du flux SSE d'import streamé.
* - progress : avancement (total=0 ⇒ extraction en cours). * - 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). * - done : arbre proposé (à réviser).
* - error : message d'erreur côté serveur. * - error : message d'erreur côté serveur.
*/ */
export type CampaignImportStreamEvent = export type CampaignImportStreamEvent =
| { type: 'status'; message: string }
| { | {
type: 'progress'; type: 'progress';
current: number; current: number;

View File

@@ -74,6 +74,12 @@ export class CampaignImportService {
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* défaut */ } try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* défaut */ }
terminated = true; terminated = true;
subscriber.error(new Error(message)); 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') { } else if (name === 'progress' || name === 'done') {
try { try {
const obj = JSON.parse(currentData); const obj = JSON.parse(currentData);

View File

@@ -32,11 +32,14 @@ export interface RulesImportResponse {
/** /**
* Évènements du flux SSE d'import streamé. * Évènements du flux SSE d'import streamé.
* - progress : avancement (total=0 ⇒ phase d'extraction en cours). * - 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). * - done : résultat final (sections proposées).
* - error : message d'erreur côté serveur. * - error : message d'erreur côté serveur.
*/ */
export type RulesImportStreamEvent = export type RulesImportStreamEvent =
| { type: 'progress'; current: number; total: number; pageCount: number; ocrPageCount: number; newSectionTitles: string[] } | { type: 'progress'; current: number; total: number; pageCount: number; ocrPageCount: number; newSectionTitles: string[] }
| { type: 'status'; message: string }
| { type: 'done'; sections: Record<string, string>; pageCount: number; ocrPageCount: number } | { type: 'done'; sections: Record<string, string>; pageCount: number; ocrPageCount: number }
| { type: 'error'; message: string }; | { type: 'error'; message: string };

View File

@@ -105,6 +105,12 @@ export class GameSystemService {
try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* garde le défaut */ } try { message = (JSON.parse(currentData) as { message?: string }).message ?? message; } catch { /* garde le défaut */ }
terminated = true; terminated = true;
subscriber.error(new Error(message)); 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') { } else if (name === 'progress' || name === 'done') {
try { try {
const obj = JSON.parse(currentData); const obj = JSON.parse(currentData);