|
@@ -83,46 +83,45 @@ public class ChatController {
|
|
|
: UUID.randomUUID().toString();
|
|
: UUID.randomUUID().toString();
|
|
|
|
|
|
|
|
final String intent = retrieverService.classifyIntent(query);
|
|
final String intent = retrieverService.classifyIntent(query);
|
|
|
- final List<Map<String, Object>> docs = rerankerService.rerank(retrieverService.search(query, intent, 20), query, 5);
|
|
|
|
|
|
|
|
|
|
- Sinks.Many<ServerSentEvent<String>> sink = Sinks.many().unicast().onBackpressureBuffer();
|
|
|
|
|
-
|
|
|
|
|
- try {
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("intent").data(intent).build());
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("Retrieving...").build());
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("Matched " + docs.size() + " records, generating...").build());
|
|
|
|
|
-
|
|
|
|
|
- final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
|
|
- StringBuilder fullAnswer = new StringBuilder();
|
|
|
|
|
-
|
|
|
|
|
- llmService.chatStream(messages)
|
|
|
|
|
- .doOnNext(token -> {
|
|
|
|
|
- fullAnswer.append(token);
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().data(token).build());
|
|
|
|
|
- })
|
|
|
|
|
- .doOnComplete(() -> {
|
|
|
|
|
- final List<Map<String, Object>> sources = buildSources(docs);
|
|
|
|
|
- try {
|
|
|
|
|
- String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of(
|
|
|
|
|
- "intent", intent,
|
|
|
|
|
- "sources", sources,
|
|
|
|
|
- "conversation_id", cid
|
|
|
|
|
- ));
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("meta").data(meta).build());
|
|
|
|
|
- } catch (Exception ignored) {}
|
|
|
|
|
-
|
|
|
|
|
- String finalAnswer = cleanAnswer(fullAnswer.toString());
|
|
|
|
|
- persistenceService.saveMessage(cid, "user", query, intent, null);
|
|
|
|
|
- persistenceService.saveMessage(cid, "assistant", finalAnswer, intent, sources);
|
|
|
|
|
- sink.tryEmitComplete();
|
|
|
|
|
- })
|
|
|
|
|
- .doOnError(e -> sink.tryEmitError(e))
|
|
|
|
|
- .subscribe();
|
|
|
|
|
- } catch (Exception e) {
|
|
|
|
|
- sink.tryEmitError(e);
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- return sink.asFlux();
|
|
|
|
|
|
|
+ return retrieverService.searchReactive(query, intent, 20)
|
|
|
|
|
+ .map(docs -> rerankerService.rerank(docs, query, 5))
|
|
|
|
|
+ .flatMapMany(docs -> {
|
|
|
|
|
+ Sinks.Many<ServerSentEvent<String>> sink = Sinks.many().unicast().onBackpressureBuffer();
|
|
|
|
|
+
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("intent").data(intent).build());
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("Retrieving...").build());
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("Matched " + docs.size() + " records, generating...").build());
|
|
|
|
|
+
|
|
|
|
|
+ final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
|
|
+ StringBuilder fullAnswer = new StringBuilder();
|
|
|
|
|
+
|
|
|
|
|
+ llmService.chatStream(messages)
|
|
|
|
|
+ .doOnNext(token -> {
|
|
|
|
|
+ fullAnswer.append(token);
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().data(token).build());
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnComplete(() -> {
|
|
|
|
|
+ final List<Map<String, Object>> sources = buildSources(docs);
|
|
|
|
|
+ try {
|
|
|
|
|
+ String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of(
|
|
|
|
|
+ "intent", intent,
|
|
|
|
|
+ "sources", sources,
|
|
|
|
|
+ "conversation_id", cid
|
|
|
|
|
+ ));
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("meta").data(meta).build());
|
|
|
|
|
+ } catch (Exception ignored) {}
|
|
|
|
|
+
|
|
|
|
|
+ String finalAnswer = cleanAnswer(fullAnswer.toString());
|
|
|
|
|
+ persistenceService.saveMessage(cid, "user", query, intent, null);
|
|
|
|
|
+ persistenceService.saveMessage(cid, "assistant", finalAnswer, intent, sources);
|
|
|
|
|
+ sink.tryEmitComplete();
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnError(sink::tryEmitError)
|
|
|
|
|
+ .subscribe();
|
|
|
|
|
+
|
|
|
|
|
+ return sink.asFlux();
|
|
|
|
|
+ });
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// ============================================================
|
|
// ============================================================
|
|
@@ -186,68 +185,69 @@ public class ChatController {
|
|
|
return sink.asFlux();
|
|
return sink.asFlux();
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- try {
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("正在分析图片(OCR 文字识别)...").build());
|
|
|
|
|
-
|
|
|
|
|
- // Step 1: Qwen VL 分析图片
|
|
|
|
|
- StringBuilder ocrBuilder = new StringBuilder();
|
|
|
|
|
- llmService.analyzeImageStream(request.getImageBase64(), request.getMimeType(),
|
|
|
|
|
- "请分析这张图片,提取其中所有文字信息(OCR),特别是药品名称、成分、用法用量等。简要输出。")
|
|
|
|
|
- .doOnNext(ocrBuilder::append)
|
|
|
|
|
- .doOnComplete(() -> {
|
|
|
|
|
- String ocrText = ocrBuilder.toString();
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("图片分析完成,正在检索药典知识库...").build());
|
|
|
|
|
-
|
|
|
|
|
- // Step 2: 拼接查询 → RAG
|
|
|
|
|
- String query = (!request.getMessage().isBlank())
|
|
|
|
|
- ? request.getMessage() + "\n\n(图片OCR提取内容:" + ocrText + ")"
|
|
|
|
|
- : ocrText;
|
|
|
|
|
- final String intent = retrieverService.classifyIntent(query);
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("intent").data(intent).build());
|
|
|
|
|
-
|
|
|
|
|
- final List<Map<String, Object>> docs = rerankerService.rerank(retrieverService.search(query, intent, 20), query, 5);
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("status")
|
|
|
|
|
- .data("已匹配 " + docs.size() + " 条药典资料,生成回答中(已启用联网搜索)...").build());
|
|
|
|
|
-
|
|
|
|
|
- // Step 3: 构建 Prompt + 联网搜索流式生成
|
|
|
|
|
- final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
|
|
- String imageContext = "\n\n【图片分析结果】\n" + ocrText + "\n";
|
|
|
|
|
- messages.get(0).put("content", messages.get(0).get("content") + imageContext);
|
|
|
|
|
-
|
|
|
|
|
- StringBuilder fullAnswer = new StringBuilder();
|
|
|
|
|
- llmService.chatStream(messages, true)
|
|
|
|
|
- .doOnNext(token -> {
|
|
|
|
|
- fullAnswer.append(token);
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().data(token).build());
|
|
|
|
|
- })
|
|
|
|
|
- .doOnComplete(() -> {
|
|
|
|
|
-
|
|
|
|
|
- final List<Map<String, Object>> sources = buildSources(docs);
|
|
|
|
|
- try {
|
|
|
|
|
- String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of(
|
|
|
|
|
- "intent", intent,
|
|
|
|
|
- "sources", sources,
|
|
|
|
|
- "conversation_id", cid,
|
|
|
|
|
- "ocr_text", ocrText
|
|
|
|
|
- ));
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("meta").data(meta).build());
|
|
|
|
|
- } catch (Exception ignored) {}
|
|
|
|
|
-
|
|
|
|
|
- String finalAnswer = cleanAnswer(fullAnswer.toString());
|
|
|
|
|
- persistenceService.saveMessage(cid, "user",
|
|
|
|
|
- request.getMessage().isBlank() ? "[图片]" : "[图片] " + request.getMessage(),
|
|
|
|
|
- intent, null);
|
|
|
|
|
- persistenceService.saveMessage(cid, "assistant", finalAnswer, intent, sources);
|
|
|
|
|
- sink.tryEmitComplete();
|
|
|
|
|
- })
|
|
|
|
|
- .doOnError(sink::tryEmitError)
|
|
|
|
|
- .subscribe();
|
|
|
|
|
- })
|
|
|
|
|
- .doOnError(sink::tryEmitError)
|
|
|
|
|
- .subscribe();
|
|
|
|
|
- } catch (Exception e) {
|
|
|
|
|
- sink.tryEmitError(e);
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("正在分析图片(OCR 文字识别)...").build());
|
|
|
|
|
+
|
|
|
|
|
+ // Step 1: Qwen VL 分析图片(流式),收集完整结果
|
|
|
|
|
+ StringBuilder ocrBuilder = new StringBuilder();
|
|
|
|
|
+ llmService.analyzeImageStream(request.getImageBase64(), request.getMimeType(),
|
|
|
|
|
+ "请分析这张图片,提取其中所有文字信息(OCR),特别是药品名称、成分、用法用量等。简要输出。")
|
|
|
|
|
+ .doOnNext(ocrBuilder::append)
|
|
|
|
|
+ .collectList()
|
|
|
|
|
+ .flatMap(tokens -> {
|
|
|
|
|
+ String ocrText = ocrBuilder.toString();
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("status").data("图片分析完成,正在检索药典知识库...").build());
|
|
|
|
|
+
|
|
|
|
|
+ // Step 2: 拼接查询 → RAG(响应式)
|
|
|
|
|
+ String query = (!request.getMessage().isBlank())
|
|
|
|
|
+ ? request.getMessage() + "\n\n(图片OCR提取内容:" + ocrText + ")"
|
|
|
|
|
+ : ocrText;
|
|
|
|
|
+ final String intent = retrieverService.classifyIntent(query);
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("intent").data(intent).build());
|
|
|
|
|
+
|
|
|
|
|
+ return retrieverService.searchReactive(query, intent, 20)
|
|
|
|
|
+ .map(docs -> rerankerService.rerank(docs, query, 5))
|
|
|
|
|
+ .map(docs -> {
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("status")
|
|
|
|
|
+ .data("已匹配 " + docs.size() + " 条药典资料,生成回答中(已启用联网搜索)...").build());
|
|
|
|
|
+
|
|
|
|
|
+ // Step 3: 构建 Prompt + 联网搜索流式生成
|
|
|
|
|
+ final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
|
|
+ String imageContext = "\n\n【图片分析结果】\n" + ocrText + "\n";
|
|
|
|
|
+ messages.get(0).put("content", messages.get(0).get("content") + imageContext);
|
|
|
|
|
+
|
|
|
|
|
+ StringBuilder fullAnswer = new StringBuilder();
|
|
|
|
|
+ llmService.chatStream(messages, true)
|
|
|
|
|
+ .doOnNext(token -> {
|
|
|
|
|
+ fullAnswer.append(token);
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().data(token).build());
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnComplete(() -> {
|
|
|
|
|
+ final List<Map<String, Object>> sources = buildSources(docs);
|
|
|
|
|
+ try {
|
|
|
|
|
+ String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of(
|
|
|
|
|
+ "intent", intent,
|
|
|
|
|
+ "sources", sources,
|
|
|
|
|
+ "conversation_id", cid,
|
|
|
|
|
+ "ocr_text", ocrText
|
|
|
|
|
+ ));
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("meta").data(meta).build());
|
|
|
|
|
+ } catch (Exception ignored) {}
|
|
|
|
|
+
|
|
|
|
|
+ String finalAnswer = cleanAnswer(fullAnswer.toString());
|
|
|
|
|
+ persistenceService.saveMessage(cid, "user",
|
|
|
|
|
+ request.getMessage().isBlank() ? "[图片]" : "[图片] " + request.getMessage(),
|
|
|
|
|
+ intent, null);
|
|
|
|
|
+ persistenceService.saveMessage(cid, "assistant", finalAnswer, intent, sources);
|
|
|
|
|
+ sink.tryEmitComplete();
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnError(sink::tryEmitError)
|
|
|
|
|
+ .subscribe();
|
|
|
|
|
+
|
|
|
|
|
+ return docs; // dummy return for map
|
|
|
|
|
+ });
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnError(sink::tryEmitError)
|
|
|
|
|
+ .subscribe();
|
|
|
|
|
|
|
|
return sink.asFlux();
|
|
return sink.asFlux();
|
|
|
}
|
|
}
|
|
@@ -461,42 +461,49 @@ public class ChatController {
|
|
|
final String intent = retrieverService.classifyIntent(query);
|
|
final String intent = retrieverService.classifyIntent(query);
|
|
|
sink.tryEmitNext(ServerSentEvent.<String>builder().event("intent").data(intent).build());
|
|
sink.tryEmitNext(ServerSentEvent.<String>builder().event("intent").data(intent).build());
|
|
|
|
|
|
|
|
- final List<Map<String, Object>> docs = rerankerService.rerank(retrieverService.search(query, intent, 20), query, 5);
|
|
|
|
|
final boolean enableSearch = !ocrText.isEmpty() || props.isEnableWebSearch();
|
|
final boolean enableSearch = !ocrText.isEmpty() || props.isEnableWebSearch();
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("status")
|
|
|
|
|
- .data("已匹配 " + docs.size() + " 条药典资料,生成回答中"
|
|
|
|
|
- + (enableSearch ? "(已启用联网搜索)" : "") + "...").build());
|
|
|
|
|
-
|
|
|
|
|
- final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
|
|
- if (!ocrText.isEmpty()) {
|
|
|
|
|
- messages.get(0).put("content",
|
|
|
|
|
- messages.get(0).get("content") + "\n\n【" + mediaLabel + "分析结果】\n" + ocrText + "\n");
|
|
|
|
|
- }
|
|
|
|
|
|
|
|
|
|
- // 发送回答 section header
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().data("\n【📚 药典参考回答】\n\n").build());
|
|
|
|
|
|
|
+ retrieverService.searchReactive(query, intent, 20)
|
|
|
|
|
+ .map(docs -> rerankerService.rerank(docs, query, 5))
|
|
|
|
|
+ .subscribe(docs -> {
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("status")
|
|
|
|
|
+ .data("已匹配 " + docs.size() + " 条药典资料,生成回答中"
|
|
|
|
|
+ + (enableSearch ? "(已启用联网搜索)" : "") + "...").build());
|
|
|
|
|
+
|
|
|
|
|
+ final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
|
|
+ if (!ocrText.isEmpty()) {
|
|
|
|
|
+ messages.get(0).put("content",
|
|
|
|
|
+ messages.get(0).get("content") + "\n\n【" + mediaLabel + "分析结果】\n" + ocrText + "\n");
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- StringBuilder fullAnswer = new StringBuilder();
|
|
|
|
|
- llmService.chatStream(messages, enableSearch)
|
|
|
|
|
- .doOnNext(token -> {
|
|
|
|
|
- fullAnswer.append(token);
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().data(token).build());
|
|
|
|
|
- })
|
|
|
|
|
- .doOnComplete(() -> {
|
|
|
|
|
- final List<Map<String, Object>> sources = buildSources(docs);
|
|
|
|
|
- try {
|
|
|
|
|
- String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of(
|
|
|
|
|
- "intent", intent, "sources", sources, "conversation_id", cid,
|
|
|
|
|
- "ocr_text", ocrText));
|
|
|
|
|
- sink.tryEmitNext(ServerSentEvent.<String>builder().event("meta").data(meta).build());
|
|
|
|
|
- } catch (Exception ignored) {}
|
|
|
|
|
- String finalAnswer = cleanAnswer(fullAnswer.toString());
|
|
|
|
|
- String userMsg = !request.getMessage().isBlank() ? request.getMessage()
|
|
|
|
|
- : !ocrText.isEmpty() ? "[" + mediaLabel + "]" : "";
|
|
|
|
|
- persistenceService.saveMessage(cid, "user", userMsg, intent, null);
|
|
|
|
|
- persistenceService.saveMessage(cid, "assistant", finalAnswer, intent, sources);
|
|
|
|
|
- sink.tryEmitComplete();
|
|
|
|
|
- })
|
|
|
|
|
|
|
+ // 发送回答 section header
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().data("\n【📚 药典参考回答】\n\n").build());
|
|
|
|
|
+
|
|
|
|
|
+ StringBuilder fullAnswer = new StringBuilder();
|
|
|
|
|
+ llmService.chatStream(messages, enableSearch)
|
|
|
|
|
+ .doOnNext(token -> {
|
|
|
|
|
+ fullAnswer.append(token);
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().data(token).build());
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnComplete(() -> {
|
|
|
|
|
+ final List<Map<String, Object>> sources = buildSources(docs);
|
|
|
|
|
+ try {
|
|
|
|
|
+ String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of(
|
|
|
|
|
+ "intent", intent, "sources", sources, "conversation_id", cid,
|
|
|
|
|
+ "ocr_text", ocrText));
|
|
|
|
|
+ sink.tryEmitNext(ServerSentEvent.<String>builder().event("meta").data(meta).build());
|
|
|
|
|
+ } catch (Exception ignored) {}
|
|
|
|
|
+ String finalAnswer = cleanAnswer(fullAnswer.toString());
|
|
|
|
|
+ String userMsg = !request.getMessage().isBlank() ? request.getMessage()
|
|
|
|
|
+ : !ocrText.isEmpty() ? "[" + mediaLabel + "]" : "";
|
|
|
|
|
+ persistenceService.saveMessage(cid, "user", userMsg, intent, null);
|
|
|
|
|
+ persistenceService.saveMessage(cid, "assistant", finalAnswer, intent, sources);
|
|
|
|
|
+ sink.tryEmitComplete();
|
|
|
|
|
+ })
|
|
|
|
|
+ .doOnError(sink::tryEmitError)
|
|
|
|
|
+ .subscribe();
|
|
|
|
|
+ }, sink::tryEmitError);
|
|
|
|
|
+ }
|
|
|
.doOnError(sink::tryEmitError)
|
|
.doOnError(sink::tryEmitError)
|
|
|
.subscribe();
|
|
.subscribe();
|
|
|
}
|
|
}
|