package com.xly.web; import com.fasterxml.jackson.databind.ObjectMapper; import com.xly.agent.ReActAgent; import dev.langchain4j.service.TokenStream; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.io.IOException; import java.util.LinkedHashMap; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; /** * 单 agent 对话入口(M1)。 * *

{@code POST /xlyAi/api/agent/chat} 以 SSE(text/event-stream)流式返回助手的 token。 * 每帧是一条 JSON:{@code {"type":"token|done|error","content":"..."}},前端逐帧渲染。 * *

会话记忆按 {@code conversationId} 隔离(多条命名会话);缺省用 {@code userid:default}。 * 说明:token 流本身是异步的({@link TokenStream#start()} 立即返回,回调在模型线程触发), * 这里用一个线程池提交,避免占用请求线程。 */ @RestController @RequestMapping("/api/agent") public class AgentChatController { private static final Logger log = LoggerFactory.getLogger(AgentChatController.class); private final ReActAgent agent; private final ObjectMapper mapper; private final ExecutorService exec = Executors.newCachedThreadPool(); public AgentChatController(ReActAgent agent, ObjectMapper mapper) { this.agent = agent; this.mapper = mapper; } /** 前端请求体:身份字段透传(M1 只用 userid + conversationId + text)。 */ public static class ChatReq { public String text; public String userid; public String conversationId; } @PostMapping(value = "/chat", produces = "text/event-stream;charset=UTF-8") public SseEmitter chat(@RequestBody ChatReq req) { SseEmitter emitter = new SseEmitter(180_000L); final String userInput = req.text == null ? "" : req.text; final String convId = (req.conversationId != null && !req.conversationId.isBlank()) ? req.conversationId : ((req.userid == null ? "anon" : req.userid) + ":default"); exec.submit(() -> { try { TokenStream ts = agent.chat(convId, userInput); ts.onPartialResponse(token -> send(emitter, "token", token)) // 工具执行 = 一轮结束:此前流出的是模型的“调用旁白/思考”,让前端清空, // 只保留工具执行之后的最终答复(既去掉噪声、又保留流式)。 .onToolExecuted(te -> send(emitter, "reset", "")) .onCompleteResponse(resp -> { send(emitter, "done", ""); emitter.complete(); }) .onError(err -> { log.warn("agent chat stream error (conv={})", convId, err); send(emitter, "error", err.getMessage() == null ? "服务异常" : err.getMessage()); emitter.complete(); }) .start(); } catch (Exception e) { log.error("agent invoke failed (conv={})", convId, e); send(emitter, "error", "服务异常:" + e.getMessage()); emitter.complete(); } }); return emitter; } private void send(SseEmitter emitter, String type, String content) { try { Map m = new LinkedHashMap<>(); m.put("type", type); m.put("content", content); emitter.send(SseEmitter.event().data(mapper.writeValueAsString(m), MediaType.APPLICATION_JSON)); } catch (IOException | IllegalStateException e) { // 客户端已断开或 emitter 已结束——忽略。 } } }