Chat atelier : evenement SSE sources + affichage des pages utilisees
Le Brain emet un evenement sources (source_id, page, score des passages retenus) AVANT le premier token ; le Core le relaie tel quel (JSON brut) ; l UI affiche une ligne discrete sous la reponse (ex: 12, 47, 103), prefixee du nom de fichier si plusieurs sources. Transparence pour le MJ et diagnostic immediat quand le RAG repond a cote. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -17,9 +17,10 @@ public interface NotebookChatStreamer {
|
||||
|
||||
/**
|
||||
* Streame la réponse ancrée sur les sources. Les callbacks sont invoqués au fil
|
||||
* de l'eau : {@code onToken} par fragment, {@code onProgress} (mode approfondi
|
||||
* uniquement) pendant la lecture du document, {@code onDone} à la fin,
|
||||
* {@code onError} en cas d'échec.
|
||||
* de l'eau : {@code onSourcesJson} UNE fois avant le premier token (JSON brut des
|
||||
* passages utilisés — transparence UI), {@code onToken} par fragment,
|
||||
* {@code onProgress} (mode approfondi uniquement) pendant la lecture du document,
|
||||
* {@code onDone} à la fin, {@code onError} en cas d'échec.
|
||||
*
|
||||
* @param deep true = analyse approfondie (map-reduce sur tout le document) ;
|
||||
* false = chat RAG (top-k).
|
||||
@@ -29,6 +30,7 @@ public interface NotebookChatStreamer {
|
||||
List<Msg> messages,
|
||||
String context,
|
||||
boolean deep,
|
||||
Consumer<String> onSourcesJson,
|
||||
Consumer<String> onToken,
|
||||
Consumer<Progress> onProgress,
|
||||
Runnable onDone,
|
||||
|
||||
@@ -51,6 +51,7 @@ public class BrainNotebookChatClient implements NotebookChatStreamer {
|
||||
List<Msg> messages,
|
||||
String context,
|
||||
boolean deep,
|
||||
Consumer<String> onSourcesJson,
|
||||
Consumer<String> onToken,
|
||||
Consumer<Progress> onProgress,
|
||||
Runnable onDone,
|
||||
@@ -75,7 +76,8 @@ public class BrainNotebookChatClient implements NotebookChatStreamer {
|
||||
try {
|
||||
flux
|
||||
.timeout(Duration.ofSeconds(timeoutSeconds))
|
||||
.doOnNext(sse -> handleEvent(sse, terminated, onToken, onProgress, onDone, onError))
|
||||
.doOnNext(sse -> handleEvent(
|
||||
sse, terminated, onSourcesJson, onToken, onProgress, onDone, onError))
|
||||
.blockLast();
|
||||
if (!terminated[0]) {
|
||||
onDone.run(); // flux terminé sans event done explicite
|
||||
@@ -93,6 +95,7 @@ public class BrainNotebookChatClient implements NotebookChatStreamer {
|
||||
private void handleEvent(
|
||||
ServerSentEvent<String> sse,
|
||||
boolean[] terminated,
|
||||
Consumer<String> onSourcesJson,
|
||||
Consumer<String> onToken,
|
||||
Consumer<Progress> onProgress,
|
||||
Runnable onDone,
|
||||
@@ -103,6 +106,10 @@ public class BrainNotebookChatClient implements NotebookChatStreamer {
|
||||
if ("token".equals(event)) {
|
||||
String token = readField(data, "token");
|
||||
if (token != null && !token.isEmpty()) onToken.accept(token);
|
||||
} else if ("sources".equals(event)) {
|
||||
// Passages utilisés par le RAG : relayés tels quels (JSON brut) — le
|
||||
// Core n'a pas besoin de les comprendre, seulement de les transmettre.
|
||||
onSourcesJson.accept(data);
|
||||
} else if ("progress".equals(event)) {
|
||||
onProgress.accept(new Progress(readInt(data, "current"), readInt(data, "total")));
|
||||
} else if ("done".equals(event)) {
|
||||
|
||||
@@ -132,6 +132,7 @@ public class NotebookController {
|
||||
StringBuilder assistant = new StringBuilder();
|
||||
chatStreamer.stream(
|
||||
sourceIds, history, context, deep,
|
||||
sourcesJson -> sendSources(emitter, sourcesJson),
|
||||
token -> { assistant.append(token); sendToken(emitter, token); },
|
||||
progress -> sendProgress(emitter, progress),
|
||||
() -> {
|
||||
@@ -161,6 +162,15 @@ public class NotebookController {
|
||||
}
|
||||
}
|
||||
|
||||
private void sendSources(SseEmitter emitter, String sourcesJson) {
|
||||
try {
|
||||
// JSON brut du Brain ({"sources":[{source_id,page,score},…]}), relayé tel quel.
|
||||
emitter.send(SseEmitter.event().name("sources").data(sourcesJson));
|
||||
} catch (IOException | IllegalStateException e) {
|
||||
// flux fermé/expiré : on cesse d'écrire
|
||||
}
|
||||
}
|
||||
|
||||
private void sendProgress(SseEmitter emitter, NotebookChatStreamer.Progress p) {
|
||||
try {
|
||||
emitter.send(SseEmitter.event().name("progress")
|
||||
|
||||
Reference in New Issue
Block a user