package com.pharmacopoeia.controller; import com.pharmacopoeia.config.QwenProperties; import com.pharmacopoeia.dto.ChatRequest; import com.pharmacopoeia.dto.FeedbackRequest; import com.pharmacopoeia.dto.ImageChatRequest; import com.pharmacopoeia.dto.MultimodalChatRequest; import com.pharmacopoeia.service.*; import jakarta.servlet.http.HttpServletRequest; import lombok.extern.slf4j.Slf4j; import org.jetbrains.annotations.NotNull; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.http.codec.ServerSentEvent; import org.springframework.security.core.context.SecurityContextHolder; import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.web.multipart.MultipartFile; import java.util.*; import java.util.stream.Collectors; @RestController @RequestMapping("/api/v1/chat") @Slf4j public class ChatController { // 复用 PromptService.SECTION_DISPLAY 统一权威映射,避免两处重复定义导致不一致 private final RetrieverService retrieverService; private final LLMService llmService; private final PromptService promptService; private final ChatPersistenceService persistenceService; private final RerankerService rerankerService; private final QACacheService qaCache; private final JdbcTemplate jdbc; private final QwenProperties props; private final HttpServletRequest request; public ChatController(RetrieverService rs, LLMService ls, PromptService ps, ChatPersistenceService cps, RerankerService rrs, QACacheService qaCache, JdbcTemplate jdbc, QwenProperties props, HttpServletRequest request) { this.retrieverService = rs; this.llmService = ls; this.promptService = ps; this.persistenceService = cps; this.rerankerService = rrs; this.qaCache = qaCache; this.jdbc = jdbc; this.props = props; this.request = request; } @PostMapping("/ask") public ResponseEntity> chatAsk(@RequestBody ChatRequest request) { final String userKey = getCurrentUserKey(); String query = request.getMessage(); String normalized = qaCache.normalize(query); String cid = request.getConversationId() != null && !request.getConversationId().isBlank() ? request.getConversationId() : UUID.randomUUID().toString(); // 检查全局缓存(24 小时有效,不区分用户) Map cached = qaCache.get(normalized); if (cached != null) { String answer = (String) cached.get("answer"); String intent = (String) cached.getOrDefault("intent", ""); @SuppressWarnings("unchecked") List> sources = (List>) cached.getOrDefault("sources", List.of()); persistenceService.saveMessage(userKey, cid, "user", query, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", answer, intent, sources); return ResponseEntity.ok(Map.of( "answer", answer, "sources", sources, "intent", intent, "conversation_id", cid, "cached", true )); } // 未命中缓存:尝试抢占处理权,避免并发重复调 LLM if (qaCache.tryMarkPending(normalized)) { Map waited = qaCache.waitForCache(normalized); if (waited != null) { String answer = (String) waited.get("answer"); String intent = (String) waited.getOrDefault("intent", ""); @SuppressWarnings("unchecked") List> sources = (List>) waited.getOrDefault("sources", List.of()); persistenceService.saveMessage(userKey, cid, "user", query, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", answer, intent, sources); return ResponseEntity.ok(Map.of( "answer", answer, "sources", sources, "intent", intent, "conversation_id", cid, "cached", true )); } } String intent = retrieverService.classifyIntent(query); List> docs; String llmAnswer; List> sources; try { docs = retrieverService.search(query, intent, 20); docs = rerankerService.rerank(docs, query, 5); // 统一:LLM 回答 + 原文对照 List> messages = promptService.buildPrompt(query, docs, intent); llmAnswer = cleanAnswer(llmService.chat(messages)); sources = buildSources(docs); } catch (Exception e) { // 异常时清除 PENDING 标记,避免后续同问题请求被锁死 qaCache.removePending(normalized); throw e; } String answer = llmAnswer; persistenceService.saveMessage(userKey, cid, "user", query, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", answer, intent, sources); // 写入全局缓存 qaCache.put(normalized, answer, intent, sources); return ResponseEntity.ok(Map.of( "answer", answer, "sources", sources, "intent", intent, "conversation_id", cid )); } @PostMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux> chatStream(@RequestBody ChatRequest request) { final String userKey = getCurrentUserKey(); String query = request.getMessage(); String normalized = qaCache.normalize(query); final String cid = request.getConversationId() != null && !request.getConversationId().isBlank() ? request.getConversationId() : UUID.randomUUID().toString(); // 检查全局缓存 Map cached = qaCache.get(normalized); if (cached != null) { return streamCached(cached, cid, query, userKey); } // 未命中缓存:尝试抢占处理权,避免并发重复调 LLM if (qaCache.tryMarkPending(normalized)) { Map waited = qaCache.waitForCache(normalized); if (waited != null) { return streamCached(waited, cid, query, userKey); } } final String intent = retrieverService.classifyIntent(query); final long t0 = System.currentTimeMillis(); log.info("[chatStream] intent={}, query={}", intent, query.substring(0, Math.min(50, query.length()))); // 先发射 intent/status 事件,再用 flatMapMany 接回管道保持取消链完整。 // 纯 Reactor 管道(零裸 subscribe),连接断开时整条链路自动取消到百炼。 return Flux.just( ServerSentEvent.builder().event("intent").data(intent).build(), ServerSentEvent.builder().event("status").data("Retrieving...").build() ).concatWith(retrieverService.searchReactive(query, intent, 20) .map(docs -> rerankerService.rerank(docs, query, 5)) .flatMapMany(docs -> { long t2 = System.currentTimeMillis(); log.info("[chatStream] search+rerank done, docs={}, elapsed={}ms", docs.size(), t2 - t0); final List> messages = promptService.buildPrompt(query, docs, intent); final long t3 = System.currentTimeMillis(); log.info("[chatStream] prompt built, elapsed={}ms, chars={}", t3 - t0, messages.stream().mapToInt(m -> m.get("content").length()).sum()); // 用 StringBuilder 攒完整答案(Flux 内部同步操作,无线程安全问题) var fullAnswerBuf = new StringBuilder(); // 纯管道拼接,无裸 subscribe: // [status] → [token1, token2, ...] → [meta] → complete Flux> contentFlux = llmService.chatStream(messages) .map(token -> { fullAnswerBuf.append(token); return ServerSentEvent.builder().data(token).build(); }); Flux> tailFlux = Flux.defer(() -> { String finalAnswer = cleanAnswer(fullAnswerBuf.toString()); final List> sources = buildSources(docs); String meta; try { meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of( "intent", intent, "sources", sources, "conversation_id", cid )); } catch (Exception e) { meta = "{}"; } long t4 = System.currentTimeMillis(); log.info("[chatStream] chatStream done, llmElapsed={}ms, total={}ms", t4 - t3, t4 - t0); persistenceService.saveMessage(userKey, cid, "user", query, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", finalAnswer, intent, sources); qaCache.put(normalized, finalAnswer, intent, sources); return Flux.just( ServerSentEvent.builder().event("meta").data(meta).build() ); }); Flux> headFlux = Flux.just( ServerSentEvent.builder().event("status") .data("Matched " + docs.size() + " records, generating...").build() ); return Flux.concat(headFlux, contentFlux, tailFlux) .doOnError(e -> { long t4 = System.currentTimeMillis(); log.error("[chatStream] error, totalElapsed={}ms, error={}", t4 - t0, e.getMessage()); qaCache.removePending(normalized); }); })); } /** 将缓存命中结果以流式 SSE 形式返回 */ private Flux> streamCached(Map cached, String cid, String query, String userKey) { String answer = (String) cached.get("answer"); String intent = (String) cached.getOrDefault("intent", ""); @SuppressWarnings("unchecked") List> sources = (List>) cached.getOrDefault("sources", List.of()); persistenceService.saveMessage(userKey, cid, "user", query, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", answer, intent, sources); return Flux.create(sink -> { sink.next(ServerSentEvent.builder().event("intent").data(intent).build()); sink.next(ServerSentEvent.builder().event("status").data("命中缓存,直接返回...").build()); // 将缓存答案按段落拆分发送,模拟流式体验 String[] chunks = answer.split("(?<=\\n)"); for (String chunk : chunks) { sink.next(ServerSentEvent.builder().data(chunk).build()); } try { String meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of( "intent", intent, "sources", sources, "conversation_id", cid, "cached", true )); sink.next(ServerSentEvent.builder().event("meta").data(meta).build()); } catch (Exception ignored) {} sink.complete(); }); } // ============================================================ // 图片对话 API(Qwen VL 分析 + OCR → RAG 检索 → 联网搜索) // ============================================================ @PostMapping("/ask-image") public ResponseEntity> chatAskImage(@RequestBody ImageChatRequest request) { final String userKey = getCurrentUserKey(); String cid = request.getConversationId() != null && !request.getConversationId().isBlank() ? request.getConversationId() : UUID.randomUUID().toString(); // Step 1: Qwen VL 分析图片 + OCR 提取文字 String ocrText = llmService.analyzeImage( request.getImageBase64(), request.getMimeType(), "请分析这张图片,提取其中所有文字信息(OCR),特别是药品名称、成分、用法用量等关键药学信息。简要输出即可。"); // Step 2: 拼接查询 → RAG 检索 String query = (!request.getMessage().isBlank()) ? request.getMessage() + "\n\n(图片OCR提取内容:" + ocrText + ")" : ocrText; String intent = retrieverService.classifyIntent(query); List> docs = retrieverService.search(query, intent, 20); docs = rerankerService.rerank(docs, query, 5); // Step 3: 构建 Prompt(含图片分析上下文)+ 联网搜索 List> messages = promptService.buildPrompt(query, docs, intent); String imageContext = "\n\n【图片分析结果】\n" + ocrText + "\n"; messages.getFirst().put("content", messages.getFirst().get("content") + imageContext); String answer = cleanAnswer(llmService.chat(messages, true)); List> sources = buildSources(docs); persistenceService.saveMessage(userKey, cid, "user", request.getMessage().isBlank() ? "[图片]" : "[图片] " + request.getMessage(), intent, null); persistenceService.saveMessage(userKey, cid, "assistant", answer, intent, sources); return ResponseEntity.ok(Map.of( "answer", answer, "sources", sources, "intent", intent, "conversation_id", cid )); } @PostMapping(value = "/stream-image", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux> chatStreamImage(@RequestBody ImageChatRequest request) { final String userKey = getCurrentUserKey(); final String cid = request.getConversationId() != null && !request.getConversationId().isBlank() ? request.getConversationId() : UUID.randomUUID().toString(); // 校验图片数据 if (request.getImageBase64() == null || request.getImageBase64().isBlank()) { return Flux.just( ServerSentEvent.builder().event("status").data("图片数据为空,请重新上传").build() ); } // 先发状态,不等图片分析完成 return Flux.just( ServerSentEvent.builder().event("status").data("正在分析图片(OCR 文字识别)...").build() ).concatWith( llmService.analyzeImageStream(request.getImageBase64(), request.getMimeType(), "请分析这张图片,提取其中所有文字信息(OCR),特别是药品名称、成分、用法用量等。简要输出。") .collectList() .flatMapMany(tokens -> { String ocrText = String.join("", tokens); String query = (!request.getMessage().isBlank()) ? request.getMessage() + "\n\n(图片OCR提取内容:" + ocrText + ")" : ocrText; final String intent = retrieverService.classifyIntent(query); return Flux.just( ServerSentEvent.builder().event("status").data("图片分析完成,正在检索药典知识库...").build(), ServerSentEvent.builder().event("intent").data(intent).build() ).concatWith(retrieverService.searchReactive(query, intent, 20) .map(docs -> rerankerService.rerank(docs, query, 5)) .flatMapMany(docs -> { final List> messages = promptService.buildPrompt(query, docs, intent); String imageContext = "\n\n【图片分析结果】\n" + ocrText + "\n"; messages.getFirst().put("content", messages.getFirst().get("content") + imageContext); var fullAnswerBuf = new StringBuilder(); Flux> contentFlux = llmService.chatStream(messages, true) .map(token -> { fullAnswerBuf.append(token); return ServerSentEvent.builder().data(token).build(); }); Flux> tailFlux = Flux.defer(() -> { String finalAnswer = cleanAnswer(fullAnswerBuf.toString()); final List> sources = buildSources(docs); String meta; try { meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of( "intent", intent, "sources", sources, "conversation_id", cid, "ocr_text", ocrText )); } catch (Exception e) { meta = "{}"; } persistenceService.saveMessage(userKey, cid, "user", request.getMessage().isBlank() ? "[图片]" : "[图片] " + request.getMessage(), intent, null); persistenceService.saveMessage(userKey, cid, "assistant", finalAnswer, intent, sources); return Flux.just( ServerSentEvent.builder().event("meta").data(meta).build() ); }); Flux> headFlux = Flux.just( ServerSentEvent.builder().event("status") .data("已匹配 " + docs.size() + " 条药典资料,生成回答中(已启用联网搜索)...").build() ); return Flux.concat(headFlux, contentFlux, tailFlux); })); }) ); } // ============================================================ // 统一多模态对话 API(文本 + 图片 + 视频) // ============================================================ @PostMapping("/ask-multimodal") public ResponseEntity> chatAskMultimodal(@RequestBody MultimodalChatRequest request) { final String userKey = getCurrentUserKey(); String cid = request.getConversationId() != null && !request.getConversationId().isBlank() ? request.getConversationId() : UUID.randomUUID().toString(); // Step 1: 媒体分析(如果有附件) String ocrText = ""; String mediaLabel = ""; if (request.getMediaBase64() != null && !request.getMediaBase64().isBlank() && request.getMediaType() != null && !request.getMediaType().isBlank()) { mediaLabel = "video".equals(request.getMediaType()) ? "视频" : "图片"; ocrText = llmService.analyzeMedia( request.getMediaBase64(), request.getMediaType(), request.getMediaMime(), ""); } // Step 2: 拼接查询 String query = request.getMessage() != null ? request.getMessage().trim() : ""; if (!query.isEmpty() && !ocrText.isEmpty()) { query = query + "\n\n(" + mediaLabel + "OCR提取内容:" + ocrText + ")"; } else if (!ocrText.isEmpty()) { query = ocrText; } else if (query.isEmpty()) { query = "请介绍一下自己"; } // Step 3: RAG 检索 String intent = retrieverService.classifyIntent(query); List> docs = retrieverService.search(query, intent, 20); docs = rerankerService.rerank(docs, query, 5); // Step 4: 构建 Prompt + 联网搜索 List> messages = promptService.buildPrompt(query, docs, intent); if (!ocrText.isEmpty()) { messages.getFirst().put("content", messages.getFirst().get("content") + "\n\n【" + mediaLabel + "分析结果】\n" + ocrText + "\n"); } boolean enableSearch = !ocrText.isEmpty() || props.isEnableWebSearch(); String answer = cleanAnswer(llmService.chat(messages, enableSearch)); List> sources = buildSources(docs); String userMsg = !request.getMessage().isBlank() ? request.getMessage() : !ocrText.isEmpty() ? "[" + mediaLabel + "]" : request.getMessage(); persistenceService.saveMessage(userKey, cid, "user", userMsg, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", answer, intent, sources); return ResponseEntity.ok(Map.of( "answer", answer, "sources", sources, "intent", intent, "conversation_id", cid )); } @PostMapping(value = "/stream-multimodal", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux> chatStreamMultimodal(@RequestBody MultimodalChatRequest request) { final String userKey = getCurrentUserKey(); final String cid = request.getConversationId() != null && !request.getConversationId().isBlank() ? request.getConversationId() : UUID.randomUUID().toString(); final boolean hasMedia = request.getMediaBase64() != null && !request.getMediaBase64().isBlank() && request.getMediaType() != null && !request.getMediaType().isBlank(); final String mediaLabel = hasMedia && "video".equals(request.getMediaType()) ? "视频" : "图片"; if (hasMedia) { // 先发 OCR section header return Flux.just( ServerSentEvent.builder().event("status") .data("🔍 正在分析" + mediaLabel + "...").build(), ServerSentEvent.builder().data("【📷 " + mediaLabel + "分析】\n\n").build() ).concatWith( // 用 collectList 收集流式结果,同时攒 OCR 文本,全程在管道内无阻塞 llmService.analyzeMediaStream(request.getMediaBase64(), request.getMediaType(), request.getMediaMime(), "") .flatMap(token -> Flux.just( ServerSentEvent.builder().data(token).build() // 前端实时看到 )) .concatWith(Flux.just(ServerSentEvent.builder().data("\n\n").build())) // collectList 后再 flatMapMany 接 RAG——Reactor 取消时自动 abort .collectList() .flatMapMany(events -> { // 从已发送的事件中拼回 OCR 文本(不额外调 API) StringBuilder ocrBuilder = new StringBuilder(); for (ServerSentEvent e : events) { String d = e.data(); if (d != null && !d.equals("\n\n")) ocrBuilder.append(d); } String mediaOcr = ocrBuilder.toString().trim(); // 不重复发送 OCR 事件(已经通过 analyzeMediaStream 实时发送过了) return buildRagPipeline(cid, request.getMessage(), mediaOcr, userKey, mediaLabel); }) ); } else { return buildRagPipeline(cid, request.getMessage(), "", userKey, ""); } } /** 构建 RAG → 百炼流式管道(纯 Reactor,零裸 subscribe) */ private Flux> buildRagPipeline(String cid, String rawMsg, String ocrText, String userKey, String mediaLabel) { String rq = rawMsg != null ? rawMsg.trim() : ""; if (!rq.isEmpty() && ocrText != null && !ocrText.isEmpty()) { rq = rq + "\n\n(" + mediaLabel + "OCR提取内容:" + ocrText + ")"; } else if (ocrText != null && !ocrText.isEmpty()) { rq = ocrText; } else if (rq.isEmpty()) { rq = "请介绍一下自己"; } final String query = rq; final String intent = retrieverService.classifyIntent(query); final boolean enableSearch = (ocrText != null && !ocrText.isEmpty()) || props.isEnableWebSearch(); // 如果等待检索,先发状态 Flux> prefixFlux = (ocrText != null && !ocrText.isEmpty()) ? Flux.just(ServerSentEvent.builder().event("status") .data("📚 检索药典知识库...").build()) : Flux.just(ServerSentEvent.builder().event("intent").data(intent).build()); Flux> intentFlux = (ocrText != null && !ocrText.isEmpty()) ? Flux.just(ServerSentEvent.builder().event("intent").data(intent).build()) : Flux.empty(); return prefixFlux.concatWith(intentFlux) .concatWith(retrieverService.searchReactive(query, intent, 20) .map(docs -> rerankerService.rerank(docs, query, 5)) .flatMapMany(docs -> { final List> messages = promptService.buildPrompt(query, docs, intent); if (ocrText != null && !ocrText.isEmpty()) { messages.getFirst().put("content", messages.getFirst().get("content") + "\n\n【" + mediaLabel + "分析结果】\n" + ocrText + "\n"); } var fullAnswerBuf = new StringBuilder(); Flux> contentFlux = llmService.chatStream(messages, enableSearch) .map(token -> { fullAnswerBuf.append(token); return ServerSentEvent.builder().data(token).build(); }); Flux> tailFlux = Flux.defer(() -> { final List> sources = buildSources(docs); String meta; try { meta = new com.fasterxml.jackson.databind.ObjectMapper().writeValueAsString(Map.of( "intent", intent, "sources", sources, "conversation_id", cid, "ocr_text", ocrText != null ? ocrText : "")); } catch (Exception e) { meta = "{}"; } String finalAnswer = cleanAnswer(fullAnswerBuf.toString()); String userMsg = (rawMsg != null && !rawMsg.isBlank()) ? rawMsg : (ocrText != null && !ocrText.isEmpty()) ? "[" + mediaLabel + "]" : ""; persistenceService.saveMessage(userKey, cid, "user", userMsg, intent, null); persistenceService.saveMessage(userKey, cid, "assistant", finalAnswer, intent, sources); return Flux.just( ServerSentEvent.builder().event("meta").data(meta).build() ); }); Flux> headFlux = Flux.just( ServerSentEvent.builder().data("\n【📚 药典参考回答】\n\n").build(), ServerSentEvent.builder().event("status") .data("已匹配 " + docs.size() + " 条药典资料,生成回答中" + (enableSearch ? "(已启用联网搜索)" : "") + "...").build() ); return Flux.concat(headFlux, contentFlux, tailFlux); })); } // ============================================================ // 文件上传 API(multipart → base64 → 复用已有对话管线) // ============================================================ @PostMapping("/upload-image") public ResponseEntity> uploadImage( @RequestParam("file") MultipartFile file, @RequestParam(defaultValue = "") String message, @RequestParam(defaultValue = "") String conversationId) { // 校验 MIME 类型 Set allowed = Set.of("image/jpeg", "image/png", "image/webp", "image/bmp"); String contentType = file.getContentType(); if (contentType == null || !allowed.contains(contentType)) { throw new IllegalArgumentException( "不支持的图片格式: " + contentType + ",支持 jpg/png/webp/bmp"); } // 校验大小 ≤ 10MB if (file.getSize() > 10 * 1024 * 1024) { throw new IllegalArgumentException("图片大小不能超过 10MB"); } // 转 base64 → 委托给 ask-image String base64; try { base64 = Base64.getEncoder().encodeToString(file.getBytes()); } catch (Exception e) { throw new RuntimeException("读取上传文件失败", e); } ImageChatRequest req = new ImageChatRequest(); req.setImageBase64(base64); req.setMimeType(contentType); req.setMessage(message); req.setConversationId( conversationId.isBlank() ? UUID.randomUUID().toString() : conversationId); return chatAskImage(req); } @PostMapping("/upload-media") public ResponseEntity> uploadMedia( @RequestParam("file") MultipartFile file, @RequestParam(defaultValue = "") String message, @RequestParam(defaultValue = "") String conversationId) { String contentType = file.getContentType(); if (contentType == null) { throw new IllegalArgumentException("无法识别的媒体类型"); } String mediaType; long maxSize; if (contentType.startsWith("image/")) { mediaType = "image"; maxSize = 10 * 1024 * 1024; // 10MB } else if (contentType.startsWith("video/")) { mediaType = "video"; maxSize = 50 * 1024 * 1024; // 50MB } else { throw new IllegalArgumentException( "不支持的媒体格式: " + contentType + ",支持 jpg/png/webp/bmp/mp4/mov/avi/webm"); } if (file.getSize() > maxSize) { throw new IllegalArgumentException( "文件大小不能超过 " + (maxSize / 1024 / 1024) + "MB"); } String base64; try { base64 = Base64.getEncoder().encodeToString(file.getBytes()); } catch (Exception e) { throw new RuntimeException("读取上传文件失败", e); } MultimodalChatRequest req = new MultimodalChatRequest(); req.setMessage(message); req.setMediaType(mediaType); req.setMediaBase64(base64); req.setMediaMime(contentType); req.setConversationId( conversationId.isBlank() ? UUID.randomUUID().toString() : conversationId); return chatAskMultimodal(req); } @GetMapping("/history") public ResponseEntity> getHistory( @RequestParam(defaultValue = "1") int page, @RequestParam(defaultValue = "20") int pageSize) { final String userKey = getCurrentUserKey(); // getHistory 现在直接返回包含 items/page/page_size/total/total_pages 的 Map var result = persistenceService.getHistory(userKey, page, pageSize); return ResponseEntity.ok(result); } @GetMapping("/history/{cid}") public ResponseEntity> getConversationDetail(@PathVariable String cid) { var msgs = persistenceService.getConversationDetail(cid); return ResponseEntity.ok(Map.of("conversation_id", cid, "messages", msgs)); } /** 返回最近 N 条消息,供前端恢复对话(微信 WebView 等 IndexedDB 不可用场景) */ @GetMapping("/recent-messages") public ResponseEntity> getRecentMessages( @RequestParam(defaultValue = "50") int limit) { final String userKey = getCurrentUserKey(); var msgs = persistenceService.getRecentMessages(userKey, Math.min(limit, 200)); return ResponseEntity.ok(Map.of("messages", msgs)); } @PostMapping("/feedback") public ResponseEntity> submitFeedback(@RequestBody FeedbackRequest request) { persistenceService.updateFeedback(request.getMessageId(), request.getFeedback()); return ResponseEntity.ok(Map.of("status", "ok")); } @GetMapping("/admin/conversations") public ResponseEntity> adminListConversations( @RequestParam(defaultValue = "1") int page, @RequestParam(defaultValue = "20") int pageSize, @RequestParam(required = false) String keyword) { int offset = (page - 1) * pageSize; StringBuilder sql = new StringBuilder(""" SELECT DISTINCT ON (c.conversation_id) c.conversation_id, c.title, c.created_at, m.content AS last_msg, m.role FROM conversations c JOIN messages m ON m.conversation_id = c.conversation_id """); List params = new ArrayList<>(); if (keyword != null && !keyword.isBlank()) { sql.append("WHERE m.content ILIKE ? "); params.add("%" + keyword + "%"); } sql.append(""" ORDER BY c.conversation_id, m.created_at DESC LIMIT ? OFFSET ? """); params.add(pageSize); params.add(offset); List> items = jdbc.queryForList( sql.toString(), params.toArray()); return ResponseEntity.ok(Map.of( "items", items, "page", page, "page_size", pageSize )); } private List> buildSources(List> docs) { Set rawSeen = new HashSet<>(); Set seen = new HashSet<>(); return docs.stream() // 第一层:按原始 name|section 去重,消除同一栏目的多个分块 .filter(d -> { String key = d.getOrDefault("name", "") + "|" + d.getOrDefault("section", ""); return rawSeen.add(key); }) .map(d -> { String content = (String) d.getOrDefault("content", ""); String drugName = (String) d.getOrDefault("name", ""); String storedSection = (String) d.getOrDefault("section", ""); String sourceVersion = (String) d.getOrDefault("source_version", ""); String sourceVolume = (String) d.getOrDefault("source_volume", ""); String category = (String) d.getOrDefault("category", ""); // 优先用 DB 元数据,回退到内容解析 if (drugName == null || drugName.isEmpty()) { drugName = extractDrugName(content); } String sectionDisplay = PromptService.SECTION_DISPLAY.getOrDefault(storedSection, storedSection); if (sectionDisplay == null || sectionDisplay.isEmpty()) { sectionDisplay = realSection(content, storedSection); } // 跳过"正文"类无意义栏目 if ("正文".equals(sectionDisplay)) { return null; } // 构建完整来源引用 String fullSource = getFullSource(sourceVersion, sourceVolume); content = content.replaceAll("\\s*来源:.*$", ""); content = content.replaceAll("[\\r\\n]+", " ").trim(); String excerpt = content.length() > 500 ? content.substring(0, 500) + "…" : content; return Map.of( "drug_id", d.getOrDefault("drug_id", ""), "name", drugName, "section", sectionDisplay, "category", category != null ? category : "", "source", fullSource, "excerpt", excerpt ); }) .filter(Objects::nonNull) // 第二层:按显示名去重,避免别名映射(如"功能""主治"→"功能与主治")导致重复 .filter(m -> { String key = m.get("name") + "|" + m.get("section"); return seen.add(key); }) // 每种药最多 5 个栏目,总数最多 8 条 .collect(Collectors.groupingBy(m -> (String) m.get("name"), LinkedHashMap::new, Collectors.toList())) .values().stream() .flatMap(list -> list.stream().limit(5)) .limit(8) .collect(Collectors.toList()); } @NotNull private static String getFullSource(String sourceVersion, String sourceVolume) { StringBuilder sourceBuilder = new StringBuilder(); if (!sourceVersion.isEmpty()) { sourceBuilder.append(sourceVersion); } if (!sourceVolume.isEmpty()) { if (!sourceBuilder.isEmpty()) { sourceBuilder.append(" "); } sourceBuilder.append(sourceVolume); } return sourceBuilder.toString(); } private String extractDrugName(String content) { if (content == null) { return ""; } int start = content.indexOf("【"); int end = content.indexOf(" - "); if (start >= 0 && end > start) { return content.substring(start + 1, end); } return content.length() > 20 ? content.substring(0, 20) : content; } /** 从 content 文本中提取真实 section(兜底"正文") */ private String realSection(String content, String storedSection) { if (!"正文".equals(storedSection) || content == null) { return storedSection; } int sep = content.indexOf(" - "); if (sep < 0) { return storedSection; } int end = content.indexOf("】", sep); if (end > sep) { return content.substring(sep + 3, end).trim(); } return storedSection; } /** 从 SecurityContext 获取当前用户标识(JWT subject),未登录则用 IP 隔离 */ private String getCurrentUserKey() { var auth = SecurityContextHolder.getContext().getAuthentication(); if (auth != null && auth.isAuthenticated() && !"anonymousUser".equals(auth.getPrincipal())) { return auth.getName(); } // 未登录用 IP 隔离,避免不同手机会话串了 String ip = request.getRemoteAddr(); String forwarded = request.getHeader("X-Forwarded-For"); if (forwarded != null && !forwarded.isBlank()) { ip = forwarded.split(",")[0].trim(); } return "ip:" + ip; } private String cleanAnswer(String text) { if (text == null) { return ""; } // 去除多余空白行(保留单个换行),修复 Qwen 常见格式问题 return text .replace("\r\n", "\n") .replaceAll("\\n{3,}", "\n\n") .trim(); } }