瀏覽代碼

tool&rag 上下文管理

Grizzly 5 月之前
父節點
當前提交
7c23be521b

+ 35 - 12
src/main/java/edu/nju/software/aipaasagent/agent/core/base/BaseAgent.java

@@ -3,8 +3,10 @@ package edu.nju.software.aipaasagent.agent.core.base;
 import com.fasterxml.jackson.annotation.JsonIgnore;
 import edu.nju.software.aipaasagent.agent.config.AgentConfiguration;
 import edu.nju.software.aipaasagent.agent.manager.AgentState;
+import edu.nju.software.aipaasagent.dto.StreamEvent;
 import edu.nju.software.aipaasagent.mcp.manage.ToolRegister;
-import edu.nju.software.aipaasagent.service.RagService;
+import edu.nju.software.aipaasagent.rag.dto.RagResult;
+import edu.nju.software.aipaasagent.rag.service.RagService;
 import lombok.Data;
 import lombok.Getter;
 import lombok.Setter;
@@ -81,10 +83,13 @@ public abstract class BaseAgent {
     public String doChat(String message, String chatId, boolean isReact) {
         log.info("Processing Chat [Sync] - chatId: {}, mode: {}", chatId, isReact ? "ReAct" : "Native");
 
-        String finalUserMessage = preparePrompt(message, chatId);
+        // 对于同步对话,在这里执行 RAG
+        RagResult ragResult = preparePrompt(message, chatId, null);
+        String finalUserMessage = ragResult.getEnhancedUserMessage();
+        Message ragKeyInfoMessage = ragResult.getRagKeyInfoMessage();
 
         return isReact ?
-                executeReActStrategy(finalUserMessage, chatId) :
+                executeReActStrategy(finalUserMessage, chatId, ragKeyInfoMessage) :
                 executeNativeStrategy(finalUserMessage, chatId);
     }
 
@@ -94,23 +99,24 @@ public abstract class BaseAgent {
     public Flux<String> doStreamChat(String message, String chatId, boolean isReact) {
         log.info("Processing Chat [Stream] - chatId: {}, mode: {}", chatId, isReact ? "ReAct" : "Native");
 
-        String finalUserMessage = preparePrompt(message, chatId);
-
+        // 对于流式对话,不在这里执行 RAG,移到 ToolCallAgent 里执行
+        // 这样可以在创建事件发射器后再发射 RAG 事件
         return isReact ?
-                executeReActStreamStrategy(finalUserMessage, chatId) :
-                executeNativeStreamStrategy(finalUserMessage, chatId);
+                executeReActStreamStrategy(message, chatId, null) :
+                executeNativeStreamStrategy(preparePrompt(message, chatId, null).getEnhancedUserMessage(), chatId);
     }
 
     /**
      * 前置处理逻辑:RAG 判定与 Prompt 增强
      */
-    private String preparePrompt(String message, String chatId) {
+    protected RagResult preparePrompt(String message, String chatId, 
+                                       java.util.function.Consumer<edu.nju.software.aipaasagent.dto.StreamEvent> eventEmitter) {
         if (isRagEnabled()) {
             log.info("[RAG] 触发检索增强 - chatId: {}", chatId);
             List<Message> history = getHistoryFromMemory(chatId);
-            return ragService.executeRagAndBuildMessage(message, chatId, agentConfiguration.getRag(), history);
+            return ragService.executeRagAndBuildMessage(message, chatId, agentConfiguration.getRag(), history, eventEmitter);
         }
-        return message;
+        return new RagResult(message, null);
     }
 
     /**
@@ -189,9 +195,26 @@ public abstract class BaseAgent {
 
     protected abstract String executeNativeStrategy(String message, String chatId);
 
-    protected abstract String executeReActStrategy(String message, String chatId);
+    protected abstract String executeReActStrategy(String message, String chatId, Message ragKeyInfoMessage);
 
     protected abstract Flux<String> executeNativeStreamStrategy(String message, String chatId);
 
-    protected abstract Flux<String> executeReActStreamStrategy(String message, String chatId);
+    protected abstract Flux<String> executeReActStreamStrategy(String message, String chatId, Message ragKeyInfoMessage);
+    
+    /**
+     * 统一流式对话入口(返回 SSE 事件格式
+     * 用于 stream=true + isThink=true 的情况
+     */
+    public Flux<StreamEvent> doStreamChatSSE(String message, String chatId, boolean isReact) {
+        log.info("Processing Chat [SSE Stream] - chatId: {}, mode: {}", chatId, isReact ? "ReAct" : "Native");
+        // 只有 ReAct 模式才支持 SSE 事件格式
+        if (isReact) {
+            return executeReActStreamStrategySSE(message, chatId, null);
+        } else {
+            // 普通 Agent 不支持 SSE 事件格式,返回错误
+            return Flux.error(new UnsupportedOperationException("SSE format only supported for ReAct mode (isThink=true)"));
+        }
+    }
+    
+    protected abstract Flux<StreamEvent> executeReActStreamStrategySSE(String message, String chatId, Message ragKeyInfoMessage);
 }

+ 16 - 142
src/main/java/edu/nju/software/aipaasagent/agent/core/reactagent/ReActAgent.java

@@ -1,8 +1,10 @@
 package edu.nju.software.aipaasagent.agent.core.reactagent;
 
 import edu.nju.software.aipaasagent.agent.config.AgentConfiguration;
+import edu.nju.software.aipaasagent.dto.StreamEvent;
 import edu.nju.software.aipaasagent.mcp.manage.ToolRegister;
-import edu.nju.software.aipaasagent.service.RagService;
+import edu.nju.software.aipaasagent.rag.service.RagService;
+import reactor.core.publisher.Flux;
 import lombok.Data;
 import lombok.Getter;
 import lombok.Setter;
@@ -78,150 +80,22 @@ public abstract class ReActAgent extends BaseAgent {
     }
 
     @Override
-    protected String executeReActStrategy(String userPrompt, String chatId) {
-        //todo:用户工具注册
-        //List<ToolCallback> currentTools =  toolManager.
-        if (this.state != AgentState.IDLE) {
-            throw new RuntimeException("Cannot run agent from state: " + this.state);
-        }
-        if (userPrompt == null || userPrompt.trim().isEmpty()) {
-            throw new RuntimeException("Cannot run agent with empty user prompt");
-        }
-        state = AgentState.RUNNING;
-        //messageList.add(new UserMessage(userPrompt));
-        List<String> results = new ArrayList<>();
-        try {
-            for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
-                currentStep = i + 1;
-                log.info("[{}] 执行第 {} 步", name, currentStep);
-                String stepResult = step(userPrompt,chatId);
-                results.add("Step " + currentStep + ": " + stepResult);
-            }
-            if (currentStep >= maxSteps) {
-                state = AgentState.FINISHED;
-                results.add("Terminated: Reached max steps (" + maxSteps + ")");
-            }
-            return String.join("\n", results);
-        } catch (Exception e) {
-            state = AgentState.ERROR;
-            log.error("[{}] 执行错误: {}", name, e.getMessage(), e);
-            return "执行错误: " + e.getMessage();
-        } finally {
-            this.cleanup();
-        }
+    protected String executeReActStrategy(String userPrompt, String chatId, Message ragKeyInfoMessage) {
+        // 这个方法由子类 ToolCallAgent 实现
+        throw new UnsupportedOperationException("ReActAgent.executeReActStrategy should be implemented by subclass");
     }
 
     @Override
-    protected Flux<String> executeReActStreamStrategy(String userPrompt, String chatId) {
-        // 1. 入参校验(和同步版保持一致)
-        if (this.state != AgentState.IDLE) {
-            return Flux.error(new RuntimeException("Cannot run agent from state: " + this.state));
-        }
-        if (userPrompt == null || userPrompt.trim().isEmpty()) {
-            return Flux.error(new RuntimeException("Cannot run agent with empty user prompt"));
-        }
-
-        // 2. 创建 Sink(用于异步发射流式数据)
-        Sinks.Many<String> resultSink = Sinks.many().unicast().onBackpressureBuffer();
-        // 3. 标记 Agent 状态为运行中
-        this.state = AgentState.RUNNING;
-        this.currentStep = 0;
-
-        // 4. 异步执行 ReAct 步骤循环(避免阻塞主线程)
-        Schedulers.boundedElastic().schedule(() -> {
-            try {
-                // 步骤循环(和同步版逻辑完全对齐)
-                for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
-                    currentStep = i + 1;
-                    log.info("[{}] 流式执行第 {} 步", name, currentStep);
-
-                    // 执行单步逻辑,获取步骤结果
-                    String stepResult = step(userPrompt, chatId);
-                    String streamMsg = "Step " + currentStep + ": " + stepResult;
-
-                    // 发射流式数据(每一步结果实时推送)
-                    Sinks.EmitResult emitResult = resultSink.tryEmitNext(streamMsg);
-                    if (emitResult.isFailure()) {
-                        log.warn("[{}] 第 {} 步流式数据发射失败: {}", name, currentStep, emitResult);
-                        break; // 客户端断开连接,终止循环
-                    }
-                }
-
-                // 处理最大步骤终止逻辑
-                if (currentStep >= maxSteps) {
-                    state = AgentState.FINISHED;
-                    String finishMsg = "Terminated: Reached max steps (" + maxSteps + ")";
-                    resultSink.tryEmitNext(finishMsg);
-                }
-
-                // 流式响应完成
-                resultSink.tryEmitComplete();
-
-            } catch (Exception e) {
-                // 异常处理:发射错误信号,标记状态为ERROR
-                state = AgentState.ERROR;
-                log.error("[{}] 流式执行错误: {}", name, e.getMessage(), e);
-                resultSink.tryEmitError(new RuntimeException("执行错误: " + e.getMessage()));
-            } finally {
-                // 清理资源(和同步版保持一致)
-                this.cleanup();
-            }
-        });
-
-        // 5. 返回流式响应(转换 Sink 为 Flux)
-        return resultSink.asFlux()
-                // 响应式生命周期管理:取消订阅时清理状态
-                .doOnCancel(() -> {
-                    log.warn("[{}] 客户端取消流式订阅,chatId={}", name, chatId);
-                    this.state = AgentState.IDLE;
-                    this.cleanup();
-                })
-                // 异常兜底:返回友好提示
-                .onErrorResume(e -> {
-                    log.error("[{}] 流式响应异常", name, e);
-                    return Flux.just("流式响应异常: " + e.getMessage());
-                });
+    protected Flux<String> executeReActStreamStrategy(String userPrompt, String chatId, Message ragKeyInfoMessage) {
+        // 这个方法由子类 ToolCallAgent 实现
+        return Flux.error(new UnsupportedOperationException("ReActAgent.executeReActStreamStrategy should be implemented by subclass"));
+    }
+    
+    @Override
+    protected Flux<StreamEvent> executeReActStreamStrategySSE(String userPrompt, String chatId, Message ragKeyInfoMessage) {
+        // 这个方法由子类 ToolCallAgent 实现
+        return Flux.error(new UnsupportedOperationException("ReActAgent.executeReActStreamStrategySSE should be implemented by subclass"));
     }
-    /**
-     * 运行 Agent
-     * 
-     * @param userPrompt 用户输入
-     * @return 执行结果
-     */
-//    public String run(String userPrompt) {
-//        if (this.state != AgentState.IDLE) {
-//            throw new RuntimeException("Cannot run agent from state: " + this.state);
-//        }
-//        if (userPrompt == null || userPrompt.trim().isEmpty()) {
-//            throw new RuntimeException("Cannot run agent with empty user prompt");
-//        }
-//
-//        state = AgentState.RUNNING;
-//        messageList.add(new UserMessage(userPrompt));
-//        List<String> results = new ArrayList<>();
-//
-//        try {
-//            for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
-//                currentStep = i + 1;
-//                log.info("[{}] 执行第 {} 步", name, currentStep);
-//                String stepResult = step();
-//                results.add("Step " + currentStep + ": " + stepResult);
-//            }
-//
-//            if (currentStep >= maxSteps) {
-//                state = AgentState.FINISHED;
-//                results.add("Terminated: Reached max steps (" + maxSteps + ")");
-//            }
-//
-//            return String.join("\n", results);
-//        } catch (Exception e) {
-//            state = AgentState.ERROR;
-//            log.error("[{}] 执行错误: {}", name, e.getMessage(), e);
-//            return "执行错误: " + e.getMessage();
-//        } finally {
-//            this.cleanup();
-//        }
-//    }
 
     /**
      * 单步执行
@@ -229,7 +103,7 @@ public abstract class ReActAgent extends BaseAgent {
      * 
      * @return 步骤执行结果
      */
-    public String step(String prompt,String chatID) {
+    public String step(String prompt, String chatID) {
         try {
             // 思考阶段:决定是否需要调用工具
             boolean shouldAct = think(prompt, chatID);

+ 525 - 295
src/main/java/edu/nju/software/aipaasagent/agent/core/reactagent/ToolCallAgent.java

@@ -2,11 +2,13 @@ package edu.nju.software.aipaasagent.agent.core.reactagent;
 
 import cn.hutool.core.collection.CollUtil;
 
+import com.fasterxml.jackson.databind.ObjectMapper;
 import edu.nju.software.aipaasagent.agent.config.AgentConfiguration;
 import edu.nju.software.aipaasagent.agent.manager.AgentState;
+import edu.nju.software.aipaasagent.dto.StreamEvent;
 import edu.nju.software.aipaasagent.mcp.manage.ToolRegister;
-import edu.nju.software.aipaasagent.memory.chat.RedisBasedChatMemory;
-import edu.nju.software.aipaasagent.service.RagService;
+import edu.nju.software.aipaasagent.rag.dto.RagResult;
+import edu.nju.software.aipaasagent.rag.service.RagService;
 import lombok.Getter;
 import lombok.Setter;
 import lombok.extern.slf4j.Slf4j;
@@ -15,185 +17,470 @@ import org.springframework.ai.chat.memory.ChatMemory;
 import org.springframework.ai.chat.messages.AssistantMessage;
 import org.springframework.ai.chat.messages.Message;
 import org.springframework.ai.chat.messages.ToolResponseMessage;
+import org.springframework.ai.chat.messages.UserMessage;
 import org.springframework.ai.chat.model.ChatModel;
 import org.springframework.ai.chat.model.ChatResponse;
 import org.springframework.ai.chat.prompt.ChatOptions;
-import org.springframework.ai.chat.prompt.DefaultChatOptions;
 import org.springframework.ai.chat.prompt.Prompt;
 import org.springframework.ai.model.tool.ToolCallingManager;
 import org.springframework.ai.model.tool.ToolExecutionResult;
 import org.springframework.ai.tool.ToolCallback;
-
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Sinks;
+import reactor.core.scheduler.Schedulers;
+import java.util.ArrayList;
 import java.util.List;
 import java.util.stream.Collectors;
-import reactor.core.publisher.Flux;
 
 /**
  * 工具调用 Agent 类
- * 继承 ReActAgent,核心实现思考和行动逻辑
- * 重点:禁用 Spring AI 内置工具托管,手动控制流程
+ * - 循环中保存历史上下文的问题
+ * - 添加 SSE 事件流支持
+ * - 优化提示词,确保 AI 能正确使用 terminate 工具
  */
 @Setter
 @Getter
 @Slf4j
 public class ToolCallAgent extends ReActAgent {
 
-    /**
-     * 可用的工具列表
-     */
     protected ToolCallback[] availableTools;
-
-    /**
-     * ChatClient 实例(大模型客户端)
-     */
     protected ChatClient chatClient;
+    protected ToolCallingManager toolCallingManager;
+    protected ChatOptions chatOptions;
+    protected ChatResponse toolCallChatResponse;
 
     /**
-     * 工具调用管理器
+     * 当前对话的消息列表(内存中维护完整上下文)
+     * 这个列表用于每次 think() 时传递完整的对话历史
      */
-    protected ToolCallingManager toolCallingManager;
+    protected List<Message> currentConversationMessages = new ArrayList<>();
 
     /**
-     * 聊天选项配置
+     * SSE 事件发射器(用于流式输出)
      */
-    protected ChatOptions chatOptions;
+    protected Sinks.Many<StreamEvent> eventSink;
+
+    private static final ObjectMapper objectMapper = new ObjectMapper();
+
+    public ToolCallAgent(ChatModel chatModel, ChatClient chatClient, 
+                         AgentConfiguration agentConfiguration, RagService ragService, 
+                         ChatMemory chatMemory, ToolRegister toolRegister) {
+        super(chatModel, agentConfiguration, chatMemory, ragService, toolRegister);
+        this.chatClient = chatClient;
+        this.toolCallingManager = ToolCallingManager.builder().build();
+        this.chatOptions = getDefaultOptions();
+        
+        log.info("[ToolCallAgent] 初始化完成,可用工具数量:{}", 
+                availableTools != null ? availableTools.length : 0);
+    }
+
+    @Override
+    protected String executeReActStrategy(String userPrompt, String chatId, Message ragKeyInfoMessage) {
+        log.info("[ToolCallAgent] 开始 ReAct 策略 - chatId: {}", chatId);
+        
+        if (this.state != AgentState.IDLE) {
+            throw new RuntimeException("Cannot run agent from state: " + this.state);
+        }
+        if (userPrompt == null || userPrompt.trim().isEmpty()) {
+            throw new RuntimeException("Cannot run agent with empty user prompt");
+        }
+
+        state = AgentState.RUNNING;
+        List<String> results = new ArrayList<>();
+        
+        try {
+            // 1. 从 Redis 加载历史对话!(非常重要)
+            log.info("[ToolCallAgent] 从 Redis 加载历史对话 - chatId: {}", chatId);
+            List<Message> historyMessages = chatMemory.get(chatId);
+            
+            // 2. 初始化对话上下文:历史 + 当前用户消息 + RAG关键信息
+            currentConversationMessages.clear();
+            if (historyMessages != null && !historyMessages.isEmpty()) {
+                currentConversationMessages.addAll(historyMessages);
+                log.info("[ToolCallAgent] 已加载历史消息 {} 条", historyMessages.size());
+            }
+            
+            // 【重要】记录初始大小,后面只收集本轮新增的消息
+            int initialSize = currentConversationMessages.size();
+            
+            // 3. 添加 RAG 关键信息(如果有)
+            if (ragKeyInfoMessage != null) {
+                currentConversationMessages.add(ragKeyInfoMessage);
+                log.info("[ToolCallAgent] 添加 RAG 关键信息到上下文");
+            }
+            
+            // 4. 添加当前用户消息(已经过 BaseAgent 的 RAG 增强)
+            currentConversationMessages.add(new UserMessage(userPrompt));
+            log.info("[ToolCallAgent] 添加当前用户消息到上下文");
+            
+            // 5. 【重要】存储这一轮的开始消息到 Redis
+            List<Message> roundStartMessages = new ArrayList<>();
+            if (ragKeyInfoMessage != null) {
+                roundStartMessages.add(ragKeyInfoMessage);
+            }
+            roundStartMessages.add(new UserMessage(userPrompt));
+            
+            log.info("[ToolCallAgent] 存储本轮开始消息到 Redis - chatId: {}, messagesCount: {}", 
+                    chatId, roundStartMessages.size());
+            chatMemory.add(chatId, roundStartMessages);
+
+            // 3. 执行 ReAct 循环
+            for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
+                currentStep = i + 1;
+                log.info("[ToolCallAgent] 执行第 {} 步", currentStep);
+                
+                String stepResult = step(userPrompt, chatId);
+                results.add("Step " + currentStep + ": " + stepResult);
+                
+                // 检查是否检测到 terminate 工具
+                if (state == AgentState.FINISHED) {
+                    log.info("[ToolCallAgent] 检测到 terminate,提前结束");
+                    break;
+                }
+            }
+
+            if (currentStep >= maxSteps) {
+                state = AgentState.FINISHED;
+                results.add("Terminated: Reached max steps (" + maxSteps + ")");
+            }
+
+            String finalResult = String.join("\n", results);
+            
+            // 4. 收集本轮对话的所有新增消息(跳过初始历史)
+            // 注意:roundStartMessages 已经存了 RAG 关键信息和用户消息
+            List<Message> roundMessages = new ArrayList<>();
+            for (int i = initialSize + roundStartMessages.size(); i < currentConversationMessages.size(); i++) {
+                roundMessages.add(currentConversationMessages.get(i));
+            }
+            
+            // 5. 【重要】批量存储这一轮的剩余消息到 Redis
+            log.info("[ToolCallAgent] 存储本轮剩余消息到 Redis - chatId: {}, messagesCount: {}", 
+                    chatId, roundMessages.size());
+            if (!roundMessages.isEmpty()) {
+                chatMemory.add(chatId, roundMessages);
+            }
+
+            return finalResult;
+
+        } catch (Exception e) {
+            state = AgentState.ERROR;
+            log.error("[ToolCallAgent] 执行错误: {}", e.getMessage(), e);
+            return "执行错误: " + e.getMessage();
+        } finally {
+            this.cleanup();
+        }
+    }
+
+    @Override
+    protected Flux<String> executeReActStreamStrategy(String originalUserPrompt, String chatId, Message unusedRagKeyInfoMessage) {
+        log.info("[ToolCallAgent] 开始流式 ReAct 策略 - chatId: {}", chatId);
+        
+        if (this.state != AgentState.IDLE) {
+            return Flux.error(new RuntimeException("Cannot run agent from state: " + this.state));
+        }
+        if (originalUserPrompt == null || originalUserPrompt.trim().isEmpty()) {
+            return Flux.error(new RuntimeException("Cannot run agent with empty user prompt"));
+        }
+
+        // 创建 SSE 事件发射器
+        eventSink = Sinks.many().unicast().onBackpressureBuffer();
+        state = AgentState.RUNNING;
+
+        // 异步执行
+        Sinks.Many<StreamEvent> sink = eventSink;
+        Schedulers.boundedElastic().schedule(() -> {
+            try {
+                // 1. 【重要】先调用 RAG(在创建了事件发射器之后!)
+                // 这样可以发射 RAG 相关的事件给前端
+                log.info("[ToolCallAgent] 流式 - 执行 RAG 检索 - chatId: {}", chatId);
+                RagResult ragResult;
+                if (isRagEnabled()) {
+                    List<Message> history = getHistoryFromMemory(chatId);
+                    ragResult = ragService.executeRagAndBuildMessage(
+                        originalUserPrompt, chatId, agentConfiguration.getRag(), history, this::emitEvent
+                    );
+                } else {
+                    ragResult = new RagResult(originalUserPrompt, null);
+                }
+                String finalUserMessage = ragResult.getEnhancedUserMessage();
+                Message ragKeyInfoMessage = ragResult.getRagKeyInfoMessage();
+                
+                // 2. 从 Redis 加载历史对话!(非常重要)
+                log.info("[ToolCallAgent] 流式 - 从 Redis 加载历史对话 - chatId: {}", chatId);
+                List<Message> historyMessages = chatMemory.get(chatId);
+                
+                // 3. 初始化对话上下文:历史 + RAG关键信息 + 当前用户消息
+                currentConversationMessages.clear();
+                if (historyMessages != null && !historyMessages.isEmpty()) {
+                    currentConversationMessages.addAll(historyMessages);
+                    log.info("[ToolCallAgent] 流式 - 已加载历史消息 {} 条", historyMessages.size());
+                }
+                
+                // 【重要】记录初始大小,后面只收集本轮新增的消息
+                int initialSize = currentConversationMessages.size();
+                
+                // 4. 添加 RAG 关键信息(如果有)
+                if (ragKeyInfoMessage != null) {
+                    currentConversationMessages.add(ragKeyInfoMessage);
+                    log.info("[ToolCallAgent] 流式 - 添加 RAG 关键信息到上下文");
+                }
+                
+                // 5. 添加当前用户消息(已经过 RAG 增强)
+                currentConversationMessages.add(new UserMessage(finalUserMessage));
+                log.info("[ToolCallAgent] 流式 - 添加当前用户消息到上下文");
+                
+                // 6. 【重要】存储这一轮的开始消息到 Redis
+                List<Message> roundStartMessages = new ArrayList<>();
+                if (ragKeyInfoMessage != null) {
+                    roundStartMessages.add(ragKeyInfoMessage);
+                }
+                roundStartMessages.add(new UserMessage(finalUserMessage));
+                
+                log.info("[ToolCallAgent] 流式 - 存储本轮开始消息到 Redis - chatId: {}, messagesCount: {}", 
+                        chatId, roundStartMessages.size());
+                chatMemory.add(chatId, roundStartMessages);
+                
+                // 发射 thinking_start 事件
+                emitEvent(StreamEvent.thinkingStart());
+
+                // 2. 执行 ReAct 循环
+                for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
+                    currentStep = i + 1;
+                    log.info("[ToolCallAgent] 流式执行第 {} 步", currentStep);
+                    
+                    // 执行单步
+                    boolean shouldContinue = executeStreamStep(finalUserMessage, chatId, sink);
+                    if (!shouldContinue) {
+                        break;
+                    }
+                }
+
+                if (currentStep >= maxSteps) {
+                    state = AgentState.FINISHED;
+                    emitEvent(StreamEvent.contentChunk("Terminated: Reached max steps (" + maxSteps + ")"));
+                }
+
+                // 存储本轮对话的所有新增消息(跳过初始历史)
+                // 注意:roundStartMessages 已经存了 RAG 关键信息和用户消息
+                List<Message> roundMessages = new ArrayList<>();
+                for (int i = initialSize + roundStartMessages.size(); i < currentConversationMessages.size(); i++) {
+                    roundMessages.add(currentConversationMessages.get(i));
+                }
+                log.info("[ToolCallAgent] 流式 - 存储本轮剩余消息到 Redis - chatId: {}, messagesCount: {}", 
+                        chatId, roundMessages.size());
+                if (!roundMessages.isEmpty()) {
+                    chatMemory.add(chatId, roundMessages);
+                }
+
+                // 发射 done 事件
+                emitEvent(StreamEvent.done());
+
+            } catch (Exception e) {
+                state = AgentState.ERROR;
+                log.error("[ToolCallAgent] 流式执行错误: {}", e.getMessage(), e);
+                emitEvent(StreamEvent.error(e.getMessage()));
+            } finally {
+                this.cleanup();
+            }
+        });
+
+        // 返回事件流(转换为 JSON 字符串)
+        return sink.asFlux()
+                .map(this::eventToSse)
+                .doOnCancel(() -> {
+                    log.warn("[ToolCallAgent] 客户端取消订阅 - chatId: {}", chatId);
+                    state = AgentState.IDLE;
+                });
+    }
+
+    @Override
+    protected Flux<StreamEvent> executeReActStreamStrategySSE(String originalUserPrompt, String chatId, Message unusedRagKeyInfoMessage) {
+        log.info("[ToolCallAgent] 开始 SSE 流式 ReAct 策略 - chatId: {}", chatId);
+        
+        if (this.state != AgentState.IDLE) {
+            return Flux.error(new RuntimeException("Cannot run agent from state: " + this.state));
+        }
+        if (originalUserPrompt == null || originalUserPrompt.trim().isEmpty()) {
+            return Flux.error(new RuntimeException("Cannot run agent with empty user prompt"));
+        }
+
+        // 创建 SSE 事件发射器
+        eventSink = Sinks.many().unicast().onBackpressureBuffer();
+        state = AgentState.RUNNING;
+
+        // 异步执行
+        Sinks.Many<StreamEvent> sink = eventSink;
+        Schedulers.boundedElastic().schedule(() -> {
+            try {
+                // 1. 【重要】先调用 RAG(在创建了事件发射器之后!)
+                // 这样可以发射 RAG 相关的事件给前端
+                log.info("[ToolCallAgent] SSE流式 - 执行 RAG 检索 - chatId: {}", chatId);
+                RagResult ragResult;
+                if (isRagEnabled()) {
+                    List<Message> history = getHistoryFromMemory(chatId);
+                    ragResult = ragService.executeRagAndBuildMessage(
+                        originalUserPrompt, chatId, agentConfiguration.getRag(), history, this::emitEvent
+                    );
+                } else {
+                    ragResult = new RagResult(originalUserPrompt, null);
+                }
+                String finalUserMessage = ragResult.getEnhancedUserMessage();
+                Message ragKeyInfoMessage = ragResult.getRagKeyInfoMessage();
+                
+                // 2. 从 Redis 加载历史对话!(非常重要)
+                log.info("[ToolCallAgent] SSE流式 - 从 Redis 加载历史对话 - chatId: {}", chatId);
+                List<Message> historyMessages = chatMemory.get(chatId);
+                
+                // 3. 初始化对话上下文:历史 + RAG关键信息 + 当前用户消息
+                currentConversationMessages.clear();
+                if (historyMessages != null && !historyMessages.isEmpty()) {
+                    currentConversationMessages.addAll(historyMessages);
+                    log.info("[ToolCallAgent] SSE流式 - 已加载历史消息 {} 条", historyMessages.size());
+                }
+                
+                // 【重要】记录初始大小,后面只收集本轮新增的消息
+                int initialSize = currentConversationMessages.size();
+                
+                // 4. 添加 RAG 关键信息(如果有)
+                if (ragKeyInfoMessage != null) {
+                    currentConversationMessages.add(ragKeyInfoMessage);
+                    log.info("[ToolCallAgent] SSE流式 - 添加 RAG 关键信息到上下文");
+                }
+                
+                // 5. 添加当前用户消息(已经过 RAG 增强)
+                currentConversationMessages.add(new UserMessage(finalUserMessage));
+                log.info("[ToolCallAgent] SSE流式 - 添加当前用户消息到上下文");
+                
+                // 6. 【重要】存储这一轮的开始消息到 Redis
+                List<Message> roundStartMessages = new ArrayList<>();
+                if (ragKeyInfoMessage != null) {
+                    roundStartMessages.add(ragKeyInfoMessage);
+                }
+                roundStartMessages.add(new UserMessage(finalUserMessage));
+                
+                log.info("[ToolCallAgent] SSE流式 - 存储本轮开始消息到 Redis - chatId: {}, messagesCount: {}", 
+                        chatId, roundStartMessages.size());
+                chatMemory.add(chatId, roundStartMessages);
+                
+                // 发射 thinking_start 事件
+                emitEvent(StreamEvent.thinkingStart());
+
+                // 2. 执行 ReAct 循环
+                for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
+                    currentStep = i + 1;
+                    log.info("[ToolCallAgent] SSE流式执行第 {} 步", currentStep);
+                    
+                    // 执行单步
+                    boolean shouldContinue = executeStreamStep(finalUserMessage, chatId, sink);
+                    if (!shouldContinue) {
+                        break;
+                    }
+                }
+
+                if (currentStep >= maxSteps) {
+                    state = AgentState.FINISHED;
+                    emitEvent(StreamEvent.contentChunk("Terminated: Reached max steps (" + maxSteps + ")"));
+                }
+
+                // 存储本轮对话的所有新增消息(跳过初始历史)
+                // 注意:roundStartMessages 已经存了 RAG 关键信息和用户消息
+                List<Message> roundMessages = new ArrayList<>();
+                for (int i = initialSize + roundStartMessages.size(); i < currentConversationMessages.size(); i++) {
+                    roundMessages.add(currentConversationMessages.get(i));
+                }
+                log.info("[ToolCallAgent] SSE流式 - 存储本轮剩余消息到 Redis - chatId: {}, messagesCount: {}", 
+                        chatId, roundMessages.size());
+                if (!roundMessages.isEmpty()) {
+                    chatMemory.add(chatId, roundMessages);
+                }
+
+                // 发射 done 事件
+                emitEvent(StreamEvent.done());
+
+            } catch (Exception e) {
+                state = AgentState.ERROR;
+                log.error("[ToolCallAgent] SSE流式执行错误: {}", e.getMessage(), e);
+                emitEvent(StreamEvent.error(e.getMessage()));
+            } finally {
+                this.cleanup();
+            }
+        });
+
+        // 直接返回 StreamEvent 流
+        return sink.asFlux()
+                .doOnCancel(() -> {
+                    log.warn("[ToolCallAgent] SSE流式 - 客户端取消订阅 - chatId: {}", chatId);
+                    state = AgentState.IDLE;
+                });
+    }
 
     /**
-     * 工具调用响应(用于 act() 阶段)
+     * 执行流式单步
      */
-    protected ChatResponse toolCallChatResponse;
+    private boolean executeStreamStep(String userPrompt, String chatId, Sinks.Many<StreamEvent> sink) {
+        try {
+            // Think 阶段
+            boolean shouldAct = think(userPrompt, chatId);
+            if (!shouldAct) {
+                log.info("[ToolCallAgent] 无需行动,直接返回结果");
+                emitEvent(StreamEvent.thinkingEnd());
+                // 获取最后一条消息作为结果
+                Message lastMsg = CollUtil.getLast(currentConversationMessages);
+                if (lastMsg instanceof AssistantMessage) {
+                    emitEvent(StreamEvent.contentChunk(((AssistantMessage) lastMsg).getText()));
+                }
+                state = AgentState.FINISHED;
+                return false;
+            }
 
-    // 注意:agentConfiguration 继承自 BaseAgent,不要在这里重新定义
-
-//    /**
-//     * 执行对话
-//     * 重写父类方法,提供 ToolCallAgent 特有的对话实现
-//     * 支持 RAG 检索
-//     *
-//     * @param message 用户消息
-//     * @param chatId  会话 ID(用于记忆管理)
-//     * @return AI 回复
-//     */
-//    @Override
-//    public String doChat(String message, String chatId) {
-//        log.info("ToolCallAgent 收到对话请求 - chatId: {}, message: {}", chatId, message);
-//
-//        if (chatClient == null) {
-//            throw new IllegalStateException("ChatClient 未初始化");
-//        }
-//
-//        // 1. 如果启用 RAG,先执行检索
-//        String finalUserMessage = message;
-//        if (isRagEnabled()) {
-//            List<Message> history = getHistoryFromMemory(chatId);
-//            finalUserMessage = ragService.executeRagAndBuildMessage(message, chatId, agentConfiguration.getRag(), history);
-//        }
-//
-//        // 2. 使用 ReAct 模式执行对话
-//        return run(finalUserMessage);
-//    }
-
-//    /**
-//     * 执行流式对话
-//     * 重写父类方法,提供 ToolCallAgent 特有的流式对话实现
-//     * 支持 RAG 检索
-//     *
-//     * @param message 用户消息
-//     * @param chatId  会话 ID(用于记忆管理)
-//     * @return 流式 AI 回复
-//     */
-//    @Override
-//    public Flux<String> doStreamChat(String message, String chatId) {
-//        log.info("ToolCallAgent 收到流式对话请求 - chatId: {}, message: {}", chatId, message);
-//
-//        if (chatClient == null) {
-//            throw new IllegalStateException("ChatClient 未初始化");
-//        }
-//
-//        // 1. 如果启用 RAG,先执行检索
-//        String finalUserMessage = message;
-//        if (isRagEnabled()) {
-//            List<Message> history = getHistoryFromMemory(chatId);
-//            finalUserMessage = ragService.executeRagAndBuildMessage(message, chatId, agentConfiguration.getRag(), history);
-//        }
-//
-//        // 2. 直接使用 ChatClient 进行流式对话
-//        return chatClient
-//                .prompt()
-//                .system(getSystemPrompt())
-//                .user(finalUserMessage)
-//                .options(chatOptions)
-//                .tools(availableTools)
-//                .stream()
-//                .content()
-//                .doOnNext(chunk -> log.debug("ToolCallAgent 流式输出 - chatId: {}, chunk: {}", chatId, chunk))
-//                .doOnComplete(() -> log.info("ToolCallAgent 流式对话完成 - chatId: {}", chatId));
-//    }
+            // Act 阶段
+            String actResult = act(userPrompt, chatId);
+            emitEvent(StreamEvent.contentChunk("Step " + currentStep + ": " + actResult));
 
+            return state != AgentState.FINISHED;
 
-    public String getSystemPrompt() {
-        return agentConfiguration.getSystemPrompt();
+        } catch (Exception e) {
+            log.error("[ToolCallAgent] 流式步骤执行失败: {}", e.getMessage(), e);
+            emitEvent(StreamEvent.error("步骤执行失败: " + e.getMessage()));
+            return false;
+        }
     }
 
-
     /**
-     * 构造函数
-     *
-     * @param chatClient         ChatClient 实例
-     * @param agentConfiguration Agent 配置
-     * @param ragService         RAG 服务
-     * @param chatMemory         Redis 对话记忆
+     * 发射 SSE 事件
      */
-    public ToolCallAgent(ChatModel chatModel, ChatClient chatClient, AgentConfiguration agentConfiguration, RagService ragService, ChatMemory chatMemory, ToolRegister toolRegister) {
-        super(chatModel, agentConfiguration, chatMemory, ragService, toolRegister);
-        this.chatClient = chatClient;
-        this.toolCallingManager = ToolCallingManager.builder().build();
-        this.chatOptions=getDefaultOptions();
-        // 根据模型配置创建聊天选项
-        log.info("[{}] ToolCallAgent 初始化完成,可用工具数量:{}, RAG: {}", name,
-                availableTools != null ? availableTools.length : 0,
-                agentConfiguration != null && agentConfiguration.getRag() != null && Boolean.TRUE.equals(agentConfiguration.getRag().getEnabled()) ? "enabled" : "disabled");
+    protected void emitEvent(StreamEvent event) {
+        if (eventSink != null) {
+            Sinks.EmitResult result = eventSink.tryEmitNext(event);
+            if (result.isFailure()) {
+                log.warn("[ToolCallAgent] 事件发射失败: {}", result);
+            }
+        }
     }
 
     /**
-     * 思考阶段实现
-     * 发送消息给大模型,判断是否需要调用工具
-     * 模型通过工具的 description 来判断应该使用哪个工具
-     *
-     * @return true - 需要调用工具,false - 无需调用工具
+     * 将事件转换为 SSE 格式字符串
      */
+    private String eventToSse(StreamEvent event) {
+        try {
+            return "data: " + objectMapper.writeValueAsString(event) + "\n\n";
+        } catch (Exception e) {
+            log.error("[ToolCallAgent] 事件序列化失败", e);
+            return "data: {\"type\":\"error\",\"data\":\"序列化失败\"}\n\n";
+        }
+    }
+
     @Override
     public boolean think(String userPrompt, String chatID) {
-        ChatOptions manualOptions= ChatOptions.builder().build().copy();
-        Prompt prompt = new Prompt(userPrompt,manualOptions);
+        log.info("[ToolCallAgent] 开始思考阶段 - chatID: {}, 历史消息数: {}", chatID, currentConversationMessages.size());
 
         try {
-            //log.info("[{}] 思考阶段 - 发送消息到模型,当前消息数:{}", name, getMessageList().size());
-
-            // 打印上一轮消息的关键信息(用于调试)
-//            if (getMessageList().size() > 1) {
-//                Message lastMessage = getMessageList().get(getMessageList().size() - 1);
-//                log.info("[{}] 上一轮消息类型:{}", name, lastMessage.getMessageType());
-//                if (lastMessage instanceof ToolResponseMessage) {
-//                    ToolResponseMessage toolResp = (ToolResponseMessage) lastMessage;
-//                    log.info("[{}] 上一轮工具响应数:{}", name, toolResp.getResponses().size());
-//                } else if (lastMessage instanceof AssistantMessage) {
-//                    AssistantMessage assistantResp = (AssistantMessage) lastMessage;
-//                    String text = assistantResp.getText();
-//                    log.info("[{}] 上一轮助手消息:{}", name, text != null && text.length() > 100 ? text.substring(0, 100) + "..." : text);
-//                }
-//            }
-
-            // 构建工具定义的 JSON 描述
-            //String toolsJson = buildToolsJson();
-
-            // 构建增强的系统提示词,
             String enhancedSystemPrompt = buildEnhancedSystemPrompt();
 
-            //TODO:注入工具
-            // 调用大模型,传入工具定义(模型通过工具描述判断是否需要调用)
-            ChatResponse chatResponse = chatClient.prompt(prompt)
+            // 使用完整的对话历史
+            ChatResponse chatResponse = chatClient.prompt()
                     .system(enhancedSystemPrompt)
+                    .messages(currentConversationMessages)
                     .toolCallbacks(toolRegister.getTools().toArray(new ToolCallback[0]))
                     .call()
                     .chatResponse();
@@ -202,202 +489,143 @@ public class ToolCallAgent extends ReActAgent {
             AssistantMessage assistantMessage = chatResponse.getResult().getOutput();
             List<AssistantMessage.ToolCall> toolCallList = assistantMessage.getToolCalls();
 
-            log.info("[{}] 思考阶段 - 模型回复:{}, 工具调用数:{}",
-                    name,
-                    assistantMessage.getText() != null ? assistantMessage.getText().substring(0, Math.min(50, assistantMessage.getText().length())) + "..." : "null",
+            log.info("[ToolCallAgent] 思考阶段 - 模型回复:{}, 工具调用数:{}",
+                    assistantMessage.getText() != null ? 
+                        assistantMessage.getText().substring(0, Math.min(50, assistantMessage.getText().length())) + "..." : "null",
                     toolCallList != null ? toolCallList.size() : 0);
 
+            // 将 AI 的回复添加到对话历史
+            currentConversationMessages.add(assistantMessage);
 
-            // 如果没有工具调用
-            if (toolCallList == null || toolCallList.isEmpty()) {
-                // 检查是否有 JSON 格式的工具调用
-                String content = assistantMessage.getText();
-                if (content != null && content.contains("\"tool_calls\"")) {
-                    log.info("[{}] 检测到 JSON 格式工具调用,进入行动阶段", name);
-                    // 将 Assistant 消息添加到列表,然后进入行动阶段
-                    //getMessageList().add(assistantMessage);
-                    return true; // 有 JSON 格式工具调用,需要执行
-                }
-
-                // 检查是否包含工具调用标记
-                boolean hasToolCall = content.contains("```tool_code") ||
-                        content.contains("```json") ||
-                        (content.contains("\"name\"") && content.contains("\"input\""));
-
-                if (hasToolCall) {
-                    log.info("[{}] 检测到工具调用,进入行动阶段", name);
-                    //  getMessageList().add(assistantMessage);
-                    return true;
-                }
+            // 检查是否有工具调用
+            if (toolCallList != null && !toolCallList.isEmpty()) {
+                log.info("[ToolCallAgent] 检测到 {} 个工具调用", toolCallList.size());
+                return true;
+            }
 
-                // 其他情况:直接返回文本
-                log.info("[{}] 模型返回文本,无需调用工具", name);
-                //getMessageList().add(assistantMessage);
+            // 检查是否有 terminate 意图
+            String content = assistantMessage.getText();
+            if (content != null && isTerminateIntent(content)) {
+                log.info("[ToolCallAgent] 检测到终止意图");
+                state = AgentState.FINISHED;
                 return false;
-
-                // force 模式:强制工具调用,返回需要调用工具
-                // auto 模式:将模型回复添加到消息列表,返回无需调用工具
-            } else {
-                // 有工具调用,将包含工具调用的消息添加到列表
-                log.info("[{}] 检测到 {} 个工具调用,进入行动阶段", name, toolCallList.size());
-                //getMessageList().add(assistantMessage);
-                return true;
             }
+
+            // 没有工具调用,直接返回
+            log.info("[ToolCallAgent] 无需调用工具");
+            return false;
+
         } catch (Exception e) {
-            log.error("[{}] 思考异常:{}", name, e.getMessage(), e);
-            //getMessageList().add(new AssistantMessage("处理时遇到错误:" + e.getMessage()));
+            log.error("[ToolCallAgent] 思考异常:{}", e.getMessage(), e);
             return false;
         }
     }
 
-
-    /**
-     * 构建增强的系统提示词
-     */
-    private String buildEnhancedSystemPrompt() {
-        return getSystemPrompt() + "\n\n" +
-                "## 使用说明\n" +
-                "1. 如果你需要调用工具,请直接回复工具调用请求\n" +
-                "2. 工具调用完成后,你会看到工具执行结果\n" +
-                "3. 当你已经完成任务或不需要再调用工具时,请使用 terminate 终止工具终止\n" +
-                "4. 如果你不需要调用工具,直接回复文本即可\n" +
-                "5. 重要:一旦你得到工具执行结果并且可以回答用户问题,立即返回最终答案,不要再调用工具";
-    }
-
-    /**
-     * 行动阶段实现
-     * 执行工具调用并处理结果
-     *
-     * @return 行动结果
-     */
     @Override
     public String act(String userPrompt, String chartID) {
         if (toolCallChatResponse == null) {
             return "没有工具调用";
         }
 
-        // 检查是否有 Spring AI 原生工具调用
         AssistantMessage assistantMessage = toolCallChatResponse.getResult().getOutput();
         List<AssistantMessage.ToolCall> toolCallList = assistantMessage.getToolCalls();
 
-        // 如果有原生工具调用,使用 ToolCallingManager 执行
-        if (toolCallList != null && !toolCallList.isEmpty()) {
-            try {
-                log.info("[{}] 行动阶段 - 执行原生工具调用,数量:{}", name, toolCallList.size());
-
-                // 创建 Prompt
-                Prompt prompt = new Prompt(userPrompt, chatOptions);
-
-                // 执行工具调用
-                ToolExecutionResult toolExecutionResult = toolCallingManager.executeToolCalls(prompt, toolCallChatResponse);
-
-                // 更新消息列表(包含工具调用结果)
-                // setMessageList(toolExecutionResult.conversationHistory());
-
-                // 获取最后一个工具响应消息
-                Message lastMessage = CollUtil.getLast(toolExecutionResult.conversationHistory());
+        if (toolCallList == null || toolCallList.isEmpty()) {
+            return "没有工具调用";
+        }
 
-                if (lastMessage instanceof ToolResponseMessage) {
-                    ToolResponseMessage toolResponseMessage = (ToolResponseMessage) lastMessage;
+        try {
+            log.info("[ToolCallAgent] 行动阶段 - 执行 {} 个工具调用", toolCallList.size());
 
-                    // 构建结果字符串
-                    String results = toolResponseMessage.getResponses().stream()
-                            .map(response -> "工具 " + response.name() + " 完成了它的任务!结果:" + response.responseData())
-                            .collect(Collectors.joining("\n"));
+            // 发射工具调用开始事件
+            for (AssistantMessage.ToolCall toolCall : toolCallList) {
+                log.info("[ToolCallAgent] 调用工具: {}", toolCall.name());
+                if (eventSink != null) {
+                    emitEvent(StreamEvent.toolCallStart(toolCall.name(), toolCall.arguments()));
+                }
+            }
 
-                    log.info("[{}] 行动阶段 - 工具执行完成", name);
-                    log.info("[{}] 工具响应数:{}", name, toolResponseMessage.getResponses().size());
+            // 执行工具调用
+            Prompt prompt = new Prompt(currentConversationMessages, chatOptions);
+            ToolExecutionResult toolExecutionResult = toolCallingManager.executeToolCalls(prompt, toolCallChatResponse);
 
-                    // 关键:添加一条助手消息,告诉模型任务已完成,请返回最终答案
-                    String assistantMessages = "工具调用已完成。\n" + results + "\n\n请根据以上结果,直接返回最终答案给用户。不需要再调用工具。";
-                    // getMessageList().add(new AssistantMessage(assistantMessages));
+            // 获取工具响应消息
+            Message lastMessage = CollUtil.getLast(toolExecutionResult.conversationHistory());
+            String results = "";
 
-                    log.info("[{}] 已添加助手消息到对话上下文,引导模型返回最终答案", name);
-                    // log.info("[{}] 下一轮对话消息数:{}", name, getMessageList().size());
+            if (lastMessage instanceof ToolResponseMessage) {
+                ToolResponseMessage toolResponseMessage = (ToolResponseMessage) lastMessage;
+                
+                results = toolResponseMessage.getResponses().stream()
+                        .map(response -> "工具 " + response.name() + " 完成!结果:" + 
+                            (response.responseData() != null ? response.responseData().toString().substring(0, Math.min(100, response.responseData().toString().length())) + "..." : "无结果"))
+                        .collect(Collectors.joining("\n"));
 
-                    // 检测终止工具调用,更新状态
-                    boolean terminateToolCalled = toolResponseMessage.getResponses().stream()
-                            .anyMatch(response -> "terminate".equals(response.name()) || "doTerminate".equals(response.name()));
+                log.info("[ToolCallAgent] 工具执行完成");
 
-                    if (terminateToolCalled) {
-                        setState(AgentState.FINISHED);
-                        log.info("[{}] 检测到终止工具调用,任务结束", name);
+                // 发射工具调用结束事件
+                for (ToolResponseMessage.ToolResponse response : toolResponseMessage.getResponses()) {
+                    if (eventSink != null) {
+                        emitEvent(StreamEvent.toolCallEnd(response.name(), response.responseData()));
+                    }
+                    
+                    // 检查是否是 terminate 工具
+                    if ("terminate".equals(response.name()) || "doTerminate".equals(response.name())) {
+                        log.info("[ToolCallAgent] 检测到 terminate 工具调用");
+                        state = AgentState.FINISHED;
                     }
-
-                    return results;
                 }
 
-                return "工具执行完成";
-            } catch (Exception e) {
-                log.error("[{}] 行动阶段异常:{}", name, e.getMessage(), e);
-                return "工具执行失败:" + e.getMessage();
+                // 将工具响应添加到对话历史
+                currentConversationMessages.add(lastMessage);
             }
-        }
 
-//        // 没有原生工具调用,尝试解析 JSON 格式的工具调用请求
-//        String content = assistantMessage.getText();
-//        if (content != null && content.contains("\"tool_calls\"")) {
-//            try {
-//                log.info("[{}] 行动阶段 - 解析 JSON 格式工具调用", name);
-//                return parseAndExecuteToolCalls(content);
-//            } catch (Exception e) {
-//                log.error("[{}] 解析 JSON 工具调用失败:{}", name, e.getMessage(), e);
-//            }
-//        }
-
-        // 检查是否为 force 模式
-        boolean forceMode = false;
-        if (agentConfiguration != null && agentConfiguration.getMcp() != null
-                && agentConfiguration.getMcp().getPolicy() != null
-                && "force".equals(agentConfiguration.getMcp().getPolicy().getMode())) {
-            forceMode = true;
+            return results;
+
+        } catch (Exception e) {
+            log.error("[ToolCallAgent] 行动阶段异常:{}", e.getMessage(), e);
+            return "工具执行失败:" + e.getMessage();
         }
+    }
 
-        if (forceMode) {
-            // force 模式:使用第一个可用的 MCP 工具进行调用
-            if (availableTools != null && availableTools.length > 0) {
-                // 找到第一个 MCP 工具
-                for (ToolCallback tool : availableTools) {
-                    //tool instanceof edu.nju.software.aipaasagent.client.mcp.McpTool
-                    if (true) {
-                        String toolName = tool.getToolDefinition().name();
-                        log.info("[{}] force 模式:使用 MCP 工具 {} 进行调用", name, toolName);
-
-                        try {
-                            // 执行工具调用
-                            String result = tool.call("{}");
-                            log.info("[{}] force 模式:工具 {} 执行结果:{}", name, toolName, result);
-                            return "force 模式:工具 " + toolName + " 执行结果:" + result;
-                        } catch (Exception e) {
-                            log.error("[{}] force 模式:工具调用失败:{}", name, e.getMessage(), e);
-                            return "force 模式:工具调用失败:" + e.getMessage();
-                        }
-                    }
-                }
+    /**
+     * 构建优化的系统提示词
+     * 明确告诉 AI 如何使用工具和 terminate
+     */
+    private String buildEnhancedSystemPrompt() {
+        return getSystemPrompt() + "\n\n" +
+                "## 工具使用说明\n" +
+                "1. 当你需要获取信息或执行操作时,请调用相应的工具\n" +
+                "2. 工具执行后,你会收到工具的执行结果\n" +
+                "3. 【重要】当你已经获得足够信息可以回答用户问题时,**必须调用 terminate 工具**来结束对话\n" +
+                "4. 【重要】不要忘记调用 terminate 工具!这是结束对话的唯一方式\n" +
+                "5. 如果你不需要调用工具,直接回答用户即可\n" +
+                "6. 只在以下情况调用 terminate:\n" +
+                "   - 你已经可以完整回答用户的问题\n" +
+                "   - 任务已经完成\n" +
+                "   - 你确定不需要再调用任何工具";
+    }
 
-                // 如果没有 MCP 工具,使用第一个可用工具
-                ToolCallback firstTool = availableTools[0];
-                String toolName = firstTool.getToolDefinition().name();
-                log.info("[{}] force 模式:使用第一个可用工具 {} 进行调用", name, toolName);
-
-                try {
-                    String result = firstTool.call("{}");
-                    log.info("[{}] force 模式:工具 {} 执行结果:{}", name, toolName, result);
-                    return "force 模式:工具 " + toolName + " 执行结果:" + result;
-                } catch (Exception e) {
-                    log.error("[{}] force 模式:工具调用失败:{}", name, e.getMessage(), e);
-                    return "force 模式:工具调用失败:" + e.getMessage();
-                }
-            } else {
-                return "force 模式:无可用工具";
-            }
-        } else {
-            return "没有工具调用";
-        }
+    /**
+     * 检查是否有终止意图(文本中包含结束信号)
+     */
+    private boolean isTerminateIntent(String content) {
+        if (content == null) return false;
+        String lowerContent = content.toLowerCase();
+        return lowerContent.contains("terminate") || 
+               lowerContent.contains("结束") || 
+               lowerContent.contains("完成") ||
+               lowerContent.contains("任务完成");
+    }
+
+    public String getSystemPrompt() {
+        return agentConfiguration.getSystemPrompt();
     }
 
     @Override
     protected String executeNativeStrategy(String message, String chatId) {
+        log.info("[ToolCallAgent] 执行 Native 策略 - chatId: {}", chatId);
+        
         ChatResponse response = chatClient
                 .prompt()
                 .user(message)
@@ -407,12 +635,15 @@ public class ToolCallAgent extends ReActAgent {
                         .param(ChatMemory.CONVERSATION_ID, chatId))
                 .call()
                 .chatResponse();
+        
         assert response != null;
         return response.toString();
     }
 
     @Override
     protected Flux<String> executeNativeStreamStrategy(String message, String chatId) {
+        log.info("[ToolCallAgent] 执行 Native 流式策略 - chatId: {}", chatId);
+        
         return chatClient
                 .prompt()
                 .user(message)
@@ -422,8 +653,7 @@ public class ToolCallAgent extends ReActAgent {
                         .param(ChatMemory.CONVERSATION_ID, chatId))
                 .stream()
                 .content()
-                .doOnNext(chunk -> log.debug("流式输出 - chatId: {}, chunk: {}", chatId, chunk))
-                .doOnComplete(() -> log.info("流式对话完成 - chatId: {}", chatId));
-
+                .doOnNext(chunk -> log.debug("[ToolCallAgent] 流式输出 - chatId: {}, chunk: {}", chatId, chunk))
+                .doOnComplete(() -> log.info("[ToolCallAgent] 流式对话完成 - chatId: {}", chatId));
     }
-}
+}

+ 26 - 7
src/main/java/edu/nju/software/aipaasagent/agent/manager/AgentFactory.java

@@ -7,7 +7,7 @@ import edu.nju.software.aipaasagent.agent.core.base.BaseAgent;
 import edu.nju.software.aipaasagent.agent.core.reactagent.ToolCallAgent;
 import edu.nju.software.aipaasagent.mcp.manage.ToolRegister;
 import edu.nju.software.aipaasagent.memory.chat.RedisBasedChatMemory;
-import edu.nju.software.aipaasagent.service.RagService;
+import edu.nju.software.aipaasagent.rag.service.RagService;
 import jakarta.annotation.Resource;
 import lombok.Data;
 import lombok.extern.slf4j.Slf4j;
@@ -243,21 +243,40 @@ public class AgentFactory {
      */
     private ChatMemory createChatMemory(AgentConfiguration config) {
         String strategy = config.getMemory().getStrategy();
+        ChatMemory baseMemory;
 
         switch (strategy.toLowerCase()) {
             case "in-memory":
                 // 使用内存记忆策略
-                MessageWindowChatMemory chatMemory= MessageWindowChatMemory.builder().
-                        chatMemoryRepository(new InMemoryChatMemoryRepository())
-                        .maxMessages(agentConfiguration.getMemory().getSize()).build();
+                MessageWindowChatMemory inMemoryChatMemory = MessageWindowChatMemory.builder()
+                        .chatMemoryRepository(new InMemoryChatMemoryRepository())
+                        .maxMessages(config.getMemory().getSize())
+                        .build();
                 log.info("使用内存记忆策略");
-                return chatMemory;
+                baseMemory = inMemoryChatMemory;
+                break;
             case "redis":
                 log.info("使用 Redis 记忆策略");
-                return redisBasedChatMemory;
+                baseMemory = redisBasedChatMemory;
+                break;
             default:
                 log.warn("未知的记忆策略:{},使用默认的 Redis 记忆策略", strategy);
-                return redisBasedChatMemory;
+                baseMemory = redisBasedChatMemory;
         }
+
+        // 包装一层 SummaryChatMemory,添加记忆压缩功能
+        int maxRounds = config.getMemory().getSize();
+        int compressRounds = Math.max(1, maxRounds / 2); // 压缩前一半
+        
+        // 选择用于压缩的模型(优先使用配置的模型)
+        ChatModel compressModel = selectChatModel(config.getModel().getProvider());
+        
+        log.info("[AgentFactory] 启用记忆压缩 - maxRounds: {}, compressRounds: {}", maxRounds, compressRounds);
+        return new edu.nju.software.aipaasagent.memory.chat.SummaryChatMemory(
+                baseMemory,
+                compressModel,
+                maxRounds,
+                compressRounds
+        );
     }
 }

+ 25 - 6
src/main/java/edu/nju/software/aipaasagent/controller/ChatCompletionController.java

@@ -3,6 +3,7 @@ package edu.nju.software.aipaasagent.controller;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import edu.nju.software.aipaasagent.dto.ChatCompletionRequest;
 import edu.nju.software.aipaasagent.dto.ChatCompletionResponse;
+import edu.nju.software.aipaasagent.dto.StreamEvent;
 import edu.nju.software.aipaasagent.service.ChatCompletionService;
 import io.swagger.v3.oas.annotations.Operation;
 import io.swagger.v3.oas.annotations.tags.Tag;
@@ -11,6 +12,7 @@ import lombok.extern.slf4j.Slf4j;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.http.MediaType;
 import org.springframework.http.ResponseEntity;
+import org.springframework.http.codec.ServerSentEvent;
 import org.springframework.web.bind.annotation.PostMapping;
 import org.springframework.web.bind.annotation.RequestBody;
 import org.springframework.web.bind.annotation.RequestMapping;
@@ -41,21 +43,38 @@ public class ChatCompletionController {
      * 聊天完成接口
      * 兼容 OpenAI /v1/chat/completions
      * 通过请求体中的 stream 参数区分流式和非流式
+     * 通过 isThink 参数区分是否返回 SSE 事件格式(包含 RAG/思考/工具调用事件)
      *
      * @param request 聊天完成请求
      * @return 聊天完成响应(非流式)或 SSE 流(流式)
      */
     @PostMapping(value = "/chat/completions")
-    @Operation(summary = "聊天接口", description = "OpenAI 兼容的聊天完成接口,支持流式和非流式输出")
+    @Operation(summary = "聊天接口", description = "OpenAI 兼容的聊天完成接口,支持流式和非流式输出," +
+            "当 stream=true 且 isThink=true 时返回 SSE 事件格式(包含 RAG/思考/工具调用事件)")
     public Object chatCompletions(@Valid @RequestBody ChatCompletionRequest request) {
-        log.info("收到聊天完成请求 - stream: {}, agentId: {}", request.getStream(), request.getAgentId());
+        log.info("收到聊天完成请求 - stream: {}, isThink: {}, agentId: {}", 
+                request.getStream(), request.getIsThink(), request.getAgentId());
 
-        // 根据 stream 参数决定返回类型
-        if (Boolean.TRUE.equals(request.getStream())) {
-            // 流式输出 - 使用 SseEmitter
+        // 判断是否需要返回 SSE 事件格式
+        boolean isSSEFormat = Boolean.TRUE.equals(request.getStream()) && 
+                              Boolean.TRUE.equals(request.getIsThink());
+
+        if (isSSEFormat) {
+            // 流式 + ReAct 模式:返回 SSE 事件格式(StreamEvent)
+            log.info("返回 SSE 事件格式(包含 RAG/思考/工具调用事件)");
+            Flux<StreamEvent> stream = chatCompletionService.streamChatCompletionSSE(request);
+            return stream.map(event -> 
+                ServerSentEvent.<StreamEvent>builder()
+                    .data(event)
+                    .build()
+            );
+        } else if (Boolean.TRUE.equals(request.getStream())) {
+            // 流式 + 普通模式:返回 OpenAI 兼容格式
+            log.info("返回 OpenAI 兼容流式格式");
             return streamChatCompletion(request);
         } else {
-            // 非流式输出
+            // 非流式:返回 OpenAI 兼容格式
+            log.info("返回 OpenAI 兼容非流式格式");
             ChatCompletionResponse response = chatCompletionService.chatCompletion(request);
             return ResponseEntity.ok()
                 .contentType(MediaType.APPLICATION_JSON)

+ 150 - 0
src/main/java/edu/nju/software/aipaasagent/dto/StreamEvent.java

@@ -0,0 +1,150 @@
+package edu.nju.software.aipaasagent.dto;
+
+import com.fasterxml.jackson.annotation.JsonInclude;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+/**
+ * 流式事件 DTO
+ * 用于 SSE 输出结构化事件
+ */
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+@JsonInclude(JsonInclude.Include.NON_NULL)
+public class StreamEvent {
+
+    /**
+     * 事件类型
+     */
+    private String type;
+
+    /**
+     * 事件数据
+     */
+    private Object data;
+
+    /**
+     * 时间戳
+     */
+    private Long timestamp;
+
+    // ==================== 工厂方法 ====================
+
+    public static StreamEvent thinkingStart() {
+        return StreamEvent.builder()
+                .type("thinking_start")
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent thinkingEnd() {
+        return StreamEvent.builder()
+                .type("thinking_end")
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent toolCallStart(String toolName, Object args) {
+        return StreamEvent.builder()
+                .type("tool_call_start")
+                .data(ToolCallData.builder()
+                        .tool(toolName)
+                        .args(args)
+                        .build())
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent toolCallEnd(String toolName, Object result) {
+        return StreamEvent.builder()
+                .type("tool_call_end")
+                .data(ToolCallData.builder()
+                        .tool(toolName)
+                        .result(result)
+                        .build())
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent contentChunk(String content) {
+        return StreamEvent.builder()
+                .type("content_chunk")
+                .data(content)
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent done() {
+        return StreamEvent.builder()
+                .type("done")
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent error(String message) {
+        return StreamEvent.builder()
+                .type("error")
+                .data(message)
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    // ==================== RAG 事件 ====================
+
+    public static StreamEvent ragStart() {
+        return StreamEvent.builder()
+                .type("rag_start")
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent ragRetrieve(int resultCount) {
+        return StreamEvent.builder()
+                .type("rag_retrieve")
+                .data(RagData.builder()
+                        .resultCount(resultCount)
+                        .build())
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent ragKeyInfo(String keyInfo) {
+        return StreamEvent.builder()
+                .type("rag_key_info")
+                .data(RagData.builder()
+                        .keyInfo(keyInfo)
+                        .build())
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    public static StreamEvent ragEnd() {
+        return StreamEvent.builder()
+                .type("rag_end")
+                .timestamp(System.currentTimeMillis())
+                .build();
+    }
+
+    @Data
+    @Builder
+    @NoArgsConstructor
+    @AllArgsConstructor
+    public static class ToolCallData {
+        private String tool;
+        private Object args;
+        private Object result;
+    }
+
+    @Data
+    @Builder
+    @NoArgsConstructor
+    @AllArgsConstructor
+    public static class RagData {
+        private Integer resultCount;
+        private String keyInfo;
+    }
+}

+ 6 - 3
src/main/java/edu/nju/software/aipaasagent/mcp/manage/McpClientService.java

@@ -111,7 +111,10 @@ public class McpClientService {
     public void init() {
         rebuild();
     }
-
+    public void getOrRebuild(){
+        //mcpAsyncClientMap 检查服务是否正常
+        //不正常则 rebuild
+    }
 
     // 2. 增加同步锁,防止多线程同时重构
     public synchronized void rebuild() {
@@ -197,8 +200,8 @@ public class McpClientService {
                 log.info("📡 节点 [{}] 请求工具列表...", name);
                 McpSchema.ListToolsResult toolResponse = client.listTools().block();
 
-                log.info("🔍 节点 [{}] 工具列表响应:{} 个工具", name,
-                        toolResponse != null && toolResponse.tools() != null ? toolResponse.tools().size() : "null");
+                    log.info("🔍 节点 [{}] 工具列表响应:{} 个工具", name,
+                            toolResponse != null && toolResponse.tools() != null ? toolResponse.tools().size() : "null");
 
                 if (toolResponse != null && toolResponse.tools() != null) {
                     List<McpSchema.Tool> tools = toolResponse.tools();

+ 226 - 0
src/main/java/edu/nju/software/aipaasagent/memory/chat/SummaryChatMemory.java

@@ -0,0 +1,226 @@
+package edu.nju.software.aipaasagent.memory.chat;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.ai.chat.messages.AssistantMessage;
+import org.springframework.ai.chat.messages.Message;
+import org.springframework.ai.chat.messages.SystemMessage;
+import org.springframework.ai.chat.messages.UserMessage;
+import org.springframework.ai.chat.memory.ChatMemory;
+import org.springframework.ai.chat.model.ChatModel;
+import org.springframework.ai.chat.model.ChatResponse;
+import org.springframework.ai.chat.prompt.Prompt;
+
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * 带记忆压缩功能的 ChatMemory 包装类
+ * 当对话轮数达到配置阈值时,自动压缩最早的对话
+ */
+@Slf4j
+public class SummaryChatMemory implements ChatMemory {
+
+    private final ChatMemory delegate;
+    private final ChatModel chatModel;
+    private final int maxRounds;
+    private final int compressRounds;
+
+    /**
+     * 构造函数
+     *
+     * @param delegate     底层的 ChatMemory 实现(如 RedisBasedChatMemory)
+     * @param chatModel    用于压缩记忆的 ChatModel
+     * @param maxRounds    最大对话轮数,达到此值触发压缩
+     * @param compressRounds 需要压缩的轮数(从开头数)
+     */
+    public SummaryChatMemory(ChatMemory delegate, ChatModel chatModel, int maxRounds, int compressRounds) {
+        this.delegate = delegate;
+        this.chatModel = chatModel;
+        this.maxRounds = maxRounds;
+        this.compressRounds = compressRounds;
+        log.info("[SummaryChatMemory] 初始化 - maxRounds: {}, compressRounds: {}", maxRounds, compressRounds);
+    }
+
+    @Override
+    public void add(String conversationId, List<Message> messages) {
+        log.info("[SummaryChatMemory] 添加消息 - conversationId: {}, messagesCount: {}", conversationId, messages.size());
+        
+        // 先添加新消息
+        delegate.add(conversationId, messages);
+        
+        // 检查是否需要压缩
+        checkAndCompress(conversationId);
+    }
+
+    @Override
+    public List<Message> get(String conversationId) {
+        return delegate.get(conversationId);
+    }
+
+    @Override
+    public void clear(String conversationId) {
+        delegate.clear(conversationId);
+        log.info("[SummaryChatMemory] 清除记忆 - conversationId: {}", conversationId);
+    }
+
+    /**
+     * 检查并执行记忆压缩
+     */
+    private void checkAndCompress(String conversationId) {
+        try {
+            List<Message> allMessages = delegate.get(conversationId);
+            
+            if (allMessages == null || allMessages.isEmpty()) {
+                return;
+            }
+
+            // 估算当前有多少轮对话
+            int currentRounds = estimateRounds(allMessages);
+            
+            log.info("[SummaryChatMemory] 检查压缩 - conversationId: {}, currentRounds: {}, maxRounds: {}", 
+                    conversationId, currentRounds, maxRounds);
+
+            if (currentRounds >= maxRounds) {
+                log.info("[SummaryChatMemory] 触发压缩 - conversationId: {}", conversationId);
+                compressMemory(conversationId, allMessages);
+            }
+        } catch (Exception e) {
+            log.error("[SummaryChatMemory] 检查压缩失败 - conversationId: {}", conversationId, e);
+        }
+    }
+
+    /**
+     * 估算对话轮数
+     * 每一轮大致包含:UserMessage + (若干工具调用) + AssistantMessage
+     */
+    private int estimateRounds(List<Message> messages) {
+        int rounds = 0;
+        for (Message msg : messages) {
+            if (msg instanceof UserMessage) {
+                rounds++;
+            }
+        }
+        return rounds;
+    }
+
+    /**
+     * 执行记忆压缩
+     */
+    private void compressMemory(String conversationId, List<Message> allMessages) {
+        try {
+            log.info("[SummaryChatMemory] 开始压缩 - conversationId: {}, totalMessages: {}", 
+                    conversationId, allMessages.size());
+
+            // 1. 找到需要压缩的消息范围(前 compressRounds 轮)
+            List<Message> messagesToCompress = new ArrayList<>();
+            List<Message> messagesToKeep = new ArrayList<>();
+            
+            int userMessageCount = 0;
+            boolean startKeeping = false;
+            
+            for (Message msg : allMessages) {
+                if (msg instanceof UserMessage) {
+                    userMessageCount++;
+                    if (userMessageCount > compressRounds) {
+                        startKeeping = true;
+                    }
+                }
+                
+                if (startKeeping) {
+                    messagesToKeep.add(msg);
+                } else {
+                    messagesToCompress.add(msg);
+                }
+            }
+
+            if (messagesToCompress.isEmpty()) {
+                log.info("[SummaryChatMemory] 没有需要压缩的消息");
+                return;
+            }
+
+            log.info("[SummaryChatMemory] 压缩范围 - compressCount: {}, keepCount: {}", 
+                    messagesToCompress.size(), messagesToKeep.size());
+
+            // 2. 调用大模型压缩
+            String summary = compressWithModel(messagesToCompress);
+            
+            if (summary == null || summary.trim().isEmpty()) {
+                log.warn("[SummaryChatMemory] 压缩结果为空,跳过压缩");
+                return;
+            }
+
+            // 3. 构建压缩后的消息列表
+            List<Message> newMessages = new ArrayList<>();
+            
+            // 添加压缩总结
+            newMessages.add(new SystemMessage("【历史对话总结】\n" + summary));
+            
+            // 添加保留的消息
+            newMessages.addAll(messagesToKeep);
+
+            // 4. 替换原有的记忆
+            // 先清除再添加
+            delegate.clear(conversationId);
+            delegate.add(conversationId, newMessages);
+
+            log.info("[SummaryChatMemory] 压缩完成 - conversationId: {}, newMessageCount: {}", 
+                    conversationId, newMessages.size());
+
+        } catch (Exception e) {
+            log.error("[SummaryChatMemory] 压缩失败 - conversationId: {}", conversationId, e);
+        }
+    }
+
+    /**
+     * 使用大模型压缩对话历史
+     */
+    private String compressWithModel(List<Message> messagesToCompress) {
+        try {
+            // 构建压缩提示词
+            String prompt = buildCompressPrompt(messagesToCompress);
+            
+            log.info("[SummaryChatMemory] 调用模型压缩 - promptLength: {}", prompt.length());
+
+            // 调用大模型
+            ChatResponse response = chatModel.call(new Prompt(prompt));
+            String summary = response.getResult().getOutput().getText();
+
+            log.info("[SummaryChatMemory] 压缩结果 - summaryLength: {}", summary != null ? summary.length() : 0);
+            log.debug("[SummaryChatMemory] 压缩结果内容: {}", summary);
+
+            return summary;
+
+        } catch (Exception e) {
+            log.error("[SummaryChatMemory] 模型压缩失败", e);
+            return null;
+        }
+    }
+
+    /**
+     * 构建压缩提示词
+     */
+    private String buildCompressPrompt(List<Message> messages) {
+        StringBuilder sb = new StringBuilder();
+        sb.append("请将以下对话历史压缩成一个简洁的总结。要求:\n");
+        sb.append("1. 保留所有关键信息和上下文\n");
+        sb.append("2. 去除冗余细节\n");
+        sb.append("3. 保持对话的逻辑连贯性\n");
+        sb.append("4. 只返回总结内容,不要有其他说明\n\n");
+        sb.append("【对话历史】\n");
+
+        for (Message msg : messages) {
+            String role;
+            if (msg instanceof UserMessage) {
+                role = "用户";
+            } else if (msg instanceof AssistantMessage) {
+                role = "助手";
+            } else {
+                role = "系统";
+            }
+            sb.append(role).append(": ").append(msg.getText()).append("\n");
+        }
+
+        sb.append("\n【总结】");
+        return sb.toString();
+    }
+}

+ 1 - 86
src/main/java/edu/nju/software/aipaasagent/rag/client/RagClient.java

@@ -1,14 +1,9 @@
 package edu.nju.software.aipaasagent.rag.client;
 
-import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
-import com.fasterxml.jackson.annotation.JsonProperty;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import edu.nju.software.aipaasagent.rag.dto.RagRetrieveRequest;
+import edu.nju.software.aipaasagent.rag.dto.RagRetrieveResponse;
 import jakarta.annotation.PostConstruct;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.beans.factory.annotation.Value;
@@ -18,8 +13,6 @@ import org.springframework.stereotype.Component;
 import org.springframework.web.reactive.function.client.WebClient;
 import reactor.core.publisher.Mono;
 
-import java.util.List;
-
 /**
  * RAG 客户端
  * 调用同学的 RAG 检索接口
@@ -75,82 +68,4 @@ public class RagClient {
                     log.error("RAG 检索请求失败", error);
                 });
     }
-
-    /**
-     * RAG 检索响应 DTO
-     */
-    @Data
-    @Builder
-    @NoArgsConstructor
-    @AllArgsConstructor
-    @JsonIgnoreProperties(ignoreUnknown = true)
-    public static class RagRetrieveResponse {
-
-        /**
-         * 状态:ok/error
-         */
-        private String status;
-
-        /**
-         * 类型:agent_retrieve
-         */
-        private String type;
-
-        /**
-         * 检索 ID
-         */
-        @JsonProperty("retrieval_id")
-        private String retrievalId;
-
-        /**
-         * 查询词
-         */
-        private String query;
-
-        /**
-         * 检索结果列表
-         */
-        private List<Result> results;
-
-        /**
-         * 错误信息
-         */
-        private String error;
-
-        /**
-         * 检索结果项
-         */
-        @Data
-        @Builder
-        @NoArgsConstructor
-        @AllArgsConstructor
-        @JsonIgnoreProperties(ignoreUnknown = true)
-        public static class Result {
-            /**
-             * 文档 ID
-             */
-            private String id;
-
-            /**
-             * 文档内容
-             */
-            private String text;
-
-            /**
-             * 来源
-             */
-            private String source;
-
-            /**
-             * 相似度分数
-             */
-            private Double score;
-
-            /**
-             * 重排序分数
-             */
-            @JsonProperty("rerank_score")
-            private Double rerankScore;
-        }
-    }
 }

+ 25 - 0
src/main/java/edu/nju/software/aipaasagent/rag/dto/RagResult.java

@@ -0,0 +1,25 @@
+package edu.nju.software.aipaasagent.rag.dto;
+
+import lombok.AllArgsConstructor;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+import org.springframework.ai.chat.messages.Message;
+
+/**
+ * RAG执行结果
+ * 包含:增强后的用户消息 + RAG关键信息(可选)
+ */
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+public class RagResult {
+    /**
+     * 增强后的用户消息(包含RAG上下文)
+     */
+    private String enhancedUserMessage;
+    
+    /**
+     * RAG关键信息(可为null)
+     */
+    private Message ragKeyInfoMessage;
+}

+ 78 - 24
src/main/java/edu/nju/software/aipaasagent/service/RagService.java → src/main/java/edu/nju/software/aipaasagent/rag/service/RagService.java

@@ -1,9 +1,12 @@
-package edu.nju.software.aipaasagent.service;
+package edu.nju.software.aipaasagent.rag.service;
 
 import edu.nju.software.aipaasagent.agent.config.AgentConfiguration;
+import edu.nju.software.aipaasagent.dto.StreamEvent;
+import edu.nju.software.aipaasagent.memory.chat.RedisBasedChatMemory;
 import edu.nju.software.aipaasagent.rag.client.RagClient;
 import edu.nju.software.aipaasagent.rag.dto.RagRetrieveRequest;
-import edu.nju.software.aipaasagent.memory.chat.RedisBasedChatMemory;
+import edu.nju.software.aipaasagent.rag.dto.RagRetrieveResponse;
+import edu.nju.software.aipaasagent.rag.dto.RagResult;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.ai.chat.messages.Message;
 import org.springframework.ai.chat.messages.UserMessage;
@@ -17,6 +20,7 @@ import reactor.core.publisher.Mono;
 
 import java.util.Comparator;
 import java.util.List;
+import java.util.function.Consumer;
 import java.util.stream.Collectors;
 
 /**
@@ -43,7 +47,7 @@ public class RagService {
      * @param ragConfig RAG 配置
      * @return 检索结果
      */
-    public Mono<RagClient.RagRetrieveResponse> retrieve(String query, AgentConfiguration.RagConfig ragConfig) {
+    public Mono<RagRetrieveResponse> retrieve(String query, AgentConfiguration.RagConfig ragConfig) {
         if (ragConfig == null || !Boolean.TRUE.equals(ragConfig.getEnabled())) {
             log.debug("RAG 未启用,跳过检索");
             return Mono.empty();
@@ -123,8 +127,8 @@ public class RagService {
      * @param maxResults 最大结果数
      * @return 处理后的结果
      */
-    public List<RagClient.RagRetrieveResponse.Result> sortAndLimitResults(
-            RagClient.RagRetrieveResponse response,
+    public List<RagRetrieveResponse.Result> sortAndLimitResults(
+            RagRetrieveResponse response,
             Integer maxResults) {
         
         if (response == null || response.getResults() == null) {
@@ -134,8 +138,8 @@ public class RagService {
 
         log.info("[RAG-Sort] 开始排序,原始结果数:{},最大结果数:{}", response.getResults().size(), maxResults);
         
-        List<RagClient.RagRetrieveResponse.Result> sorted = response.getResults().stream()
-                .sorted(Comparator.comparing(RagClient.RagRetrieveResponse.Result::getScore).reversed())
+        List<RagRetrieveResponse.Result> sorted = response.getResults().stream()
+                .sorted(Comparator.comparing(RagRetrieveResponse.Result::getScore).reversed())
                 .limit(maxResults)
                 .collect(Collectors.toList());
         
@@ -149,7 +153,7 @@ public class RagService {
      * @param results 检索结果
      * @return 拼接后的上下文
      */
-    public String buildRagContext(List<RagClient.RagRetrieveResponse.Result> results) {
+    public String buildRagContext(List<RagRetrieveResponse.Result> results) {
         if (results == null || results.isEmpty()) {
             return "";
         }
@@ -209,40 +213,70 @@ public class RagService {
     }
 
     /**
-     * 执行 RAG 检索并构建最终用户消息
+     * 执行RAG检索并构建最终用户消息
+     * 包含:生成查询词 -> 调用检索接口 -> 排序过滤 -> 拼接上下文 -> 同步提取关键信息
      * 
-     * @param userMessage 用户消息
-     * @param chatId 会话 ID
-     * @param ragConfig RAG 配置
+     * @param userMessage 用户原始问题
+     * @param chatId 会话ID
+     * @param ragConfig RAG配置
      * @param history 历史对话
-     * @return 包含 RAG 上下文的最终用户消息
+     * @return RAG执行结果(包含增强后的用户消息 + RAG关键信息)
      */
-    public String executeRagAndBuildMessage(String userMessage, String chatId, AgentConfiguration.RagConfig ragConfig, List<Message> history) {
+    public RagResult executeRagAndBuildMessage(String userMessage, String chatId, AgentConfiguration.RagConfig ragConfig, List<Message> history) {
+        return executeRagAndBuildMessage(userMessage, chatId, ragConfig, history, null);
+    }
+
+    /**
+     * 执行RAG检索并构建最终用户消息(支持事件回调)
+     * 
+     * @param userMessage 用户原始问题
+     * @param chatId 会话ID
+     * @param ragConfig RAG配置
+     * @param history 历史对话
+     * @param eventEmitter 事件发射器(可为null)
+     * @return RAG执行结果(包含增强后的用户消息 + RAG关键信息)
+     */
+    public RagResult executeRagAndBuildMessage(String userMessage, String chatId, AgentConfiguration.RagConfig ragConfig, 
+                                                 List<Message> history, 
+                                                 Consumer<StreamEvent> eventEmitter) {
         try {
             log.info("[RAG] ========== 开始执行 RAG 检索 ==========");
             log.info("[RAG] 会话 ID: {}", chatId);
             log.info("[RAG] 用户原始问题: {}", userMessage);
 
+            // 发射 RAG 开始事件
+            if (eventEmitter != null) {
+                eventEmitter.accept(StreamEvent.ragStart());
+            }
+
             // 1. 生成检索查询词
             String query = generateRetrievalQuery(userMessage, history);
             log.info("[RAG] 生成的检索查询词:{}", query);
 
             // 2. 执行 RAG 检索
             log.info("[RAG] 调用 RAG 检索接口...");
-            RagClient.RagRetrieveResponse ragResponse = retrieve(query, ragConfig).block();
+            RagRetrieveResponse ragResponse = retrieve(query, ragConfig).block();
             
             if (ragResponse == null || ragResponse.getResults() == null || ragResponse.getResults().isEmpty()) {
                 log.warn("[RAG] 检索结果为空,使用原始问题");
-                return userMessage;
+                if (eventEmitter != null) {
+                    eventEmitter.accept(StreamEvent.ragEnd());
+                }
+                return new RagResult(userMessage, null);
             }
 
             log.info("[RAG] 检索成功!");
             log.info("[RAG] 检索 ID: {}", ragResponse.getRetrievalId());
             log.info("[RAG] 结果数量:{}", ragResponse.getResults().size());
             
+            // 发射 RAG 检索结果事件
+            if (eventEmitter != null) {
+                eventEmitter.accept(StreamEvent.ragRetrieve(ragResponse.getResults().size()));
+            }
+            
             // 打印所有检索结果
             for (int i = 0; i < ragResponse.getResults().size(); i++) {
-                RagClient.RagRetrieveResponse.Result result = ragResponse.getResults().get(i);
+                RagRetrieveResponse.Result result = ragResponse.getResults().get(i);
                 log.info("[RAG] --- 结果 {} ---", i + 1);
                 log.info("[RAG]   文档 ID: {}", result.getId());
                 log.info("[RAG]   来源:{}", result.getSource());
@@ -265,14 +299,14 @@ public class RagService {
             }
             log.info("[RAG] 使用配置参数 - topK: {}", topK);
             
-            List<RagClient.RagRetrieveResponse.Result> filteredResults = sortAndLimitResults(
+            List<RagRetrieveResponse.Result> filteredResults = sortAndLimitResults(
                     ragResponse, 
                     topK
             );
 
             if (filteredResults.isEmpty()) {
                 log.warn("[RAG] 结果为空,使用原始问题");
-                return userMessage;
+                return new RagResult(userMessage, null);
             }
 
             log.info("[RAG] 保留 {} 条结果", filteredResults.size());
@@ -281,20 +315,40 @@ public class RagService {
             String ragContext = buildRagContext(filteredResults);
             log.info("[RAG] 拼接后的 RAG 上下文:\n{}", ragContext);
 
-            // 5. 触发异步总结和存储关键信息(不阻塞主流程)
-            log.info("[RAG] 触发异步关键信息总结和存储...");
-            summarizeAndStoreKeyInfoAsync(ragContext, userMessage, chatId);
+            // 5. 【重要】同步提取关键信息(不再异步!)
+            Message ragKeyInfoMessage = null;
+            String keyInfoJson = extractKeyInfo(ragContext, userMessage);
+            if (keyInfoJson != null && needStore(keyInfoJson)) {
+                log.info("[RAG] 提取的关键信息 JSON: {}", keyInfoJson);
+                // 使用 AssistantMessage 存储 RAG 关键信息,和普通 UserMessage 区分开
+                ragKeyInfoMessage = new org.springframework.ai.chat.messages.AssistantMessage(
+                    "[RAG 关键信息] " + keyInfoJson
+                );
+                // 发射 RAG 关键信息事件
+                if (eventEmitter != null) {
+                    eventEmitter.accept(StreamEvent.ragKeyInfo(keyInfoJson));
+                }
+            }
 
             // 6. 构建最终用户消息
             String finalMessage = buildFinalUserMessage(userMessage, ragContext);
             log.info("[RAG] 构建的最终用户消息:{}", finalMessage);
             log.info("[RAG] ========== RAG 检索完成 ==========");
 
-            return finalMessage;
+            // 发射 RAG 结束事件
+            if (eventEmitter != null) {
+                eventEmitter.accept(StreamEvent.ragEnd());
+            }
+
+            return new RagResult(finalMessage, ragKeyInfoMessage);
 
         } catch (Exception e) {
             log.error("[RAG] 执行检索失败,使用原始问题", e);
-            return userMessage;
+            if (eventEmitter != null) {
+                eventEmitter.accept(StreamEvent.ragEnd());
+                eventEmitter.accept(StreamEvent.error(e.getMessage()));
+            }
+            return new RagResult(userMessage, null);
         }
     }
 

+ 29 - 0
src/main/java/edu/nju/software/aipaasagent/service/ChatCompletionService.java

@@ -6,6 +6,7 @@ import edu.nju.software.aipaasagent.agent.config.AgentConfigRegistry;
 import edu.nju.software.aipaasagent.agent.config.AgentConfiguration;
 import edu.nju.software.aipaasagent.dto.ChatCompletionRequest;
 import edu.nju.software.aipaasagent.dto.ChatCompletionResponse;
+import edu.nju.software.aipaasagent.dto.StreamEvent;
 import lombok.extern.slf4j.Slf4j;
 
 import org.springframework.beans.factory.annotation.Autowired;
@@ -140,6 +141,34 @@ public class ChatCompletionService {
                     .build()
             ));
     }
+    
+    /**
+     * SSE 流式聊天完成(返回 StreamEvent 格式)
+     * 用于 stream=true + isThink=true 的情况
+     * 
+     * @param request 聊天完成请求
+     * @return StreamEvent 流
+     */
+    public Flux<StreamEvent> streamChatCompletionSSE(ChatCompletionRequest request) {
+        String threadId = getOrCreateThreadId(request.getThreadId());
+        String agentId = getOrCreateAgentId(request.getAgentId());
+        
+        log.info("SSE流式聊天完成 - threadId: {}, agentId: {}", threadId, agentId);
+
+        // 获取或创建 Agent
+        BaseAgent agent = getOrCreateAgent();
+        
+        // 转换消息格式并处理多模态文件
+        String userMessage = extractLastUserMessage(request.getMessages());
+        
+        // 处理多模态文件(如果提供了 file_ids)
+        if (request.getFileIds() != null && !request.getFileIds().isEmpty()) {
+            userMessage = processMultimodalContent(userMessage, request.getFileIds());
+        }
+
+        // 调用 SSE 流式方法
+        return agent.doStreamChatSSE(userMessage, threadId, Boolean.TRUE.equals(request.getIsThink()));
+    }
 
     /**
      * 获取或创建 Agent

+ 1 - 10
src/main/resources/agent-config.yml

@@ -38,9 +38,6 @@ react:
 mcp:
   # SeeCoder MCP 工具列表
   onlineTool:
-    SeeCoder-mcp:
-      - everything_get_sum
-      - chat_improve_prompt
   #本地的
   localTools: #plugin
     - terminate
@@ -56,15 +53,9 @@ app:
       nodes:
         SeeCoder-mcp:
           url: https://ai-paas-mcp-endpoint.njuu.top
-          endpoint: /mcp/airouting
+          endpoint: /mcp
           headers:
             Authorization: "sqGYuMvKgdxmzmTM5lNBgLdVpl6XNnPX"
-#        map-mcp:
-#          url: https://dashscope.aliyuncs.com/api/v1/mcps/amap-maps
-#          endpoint: /mcp
-#          headers:
-#            Authorization: "sk-63f4b5a5f7ab42e78843c26a89c377ec"
-
 # RAG 配置
 rag:
   enabled: false

+ 4 - 3
src/test/java/edu/nju/software/aipaasagent/client/rag/RagClientTest.java

@@ -2,6 +2,7 @@ package edu.nju.software.aipaasagent.client.rag;
 
 import edu.nju.software.aipaasagent.rag.client.RagClient;
 import edu.nju.software.aipaasagent.rag.dto.RagRetrieveRequest;
+import edu.nju.software.aipaasagent.rag.dto.RagRetrieveResponse;
 import lombok.extern.slf4j.Slf4j;
 import org.junit.jupiter.api.Test;
 import org.springframework.beans.factory.annotation.Autowired;
@@ -40,8 +41,8 @@ public class RagClientTest {
 
         // 调用接口
         try {
-            Mono<RagClient.RagRetrieveResponse> responseMono = ragClient.retrieve(request);
-            RagClient.RagRetrieveResponse response = responseMono.block();
+            Mono<RagRetrieveResponse> responseMono = ragClient.retrieve(request);
+            RagRetrieveResponse response = responseMono.block();
 
             if (response != null) {
                 log.info("RAG 检索成功!");
@@ -52,7 +53,7 @@ public class RagClientTest {
 
                 if (response.getResults() != null) {
                     for (int i = 0; i < response.getResults().size(); i++) {
-                        RagClient.RagRetrieveResponse.Result result = response.getResults().get(i);
+                        RagRetrieveResponse.Result result = response.getResults().get(i);
                         log.info("\n--- 结果 {} ---", i + 1);
                         log.info("文档 ID: {}", result.getId());
                         log.info("来源:{}", result.getSource());

+ 607 - 0
前端技术方案.md

@@ -0,0 +1,607 @@
+# 前端技术方案:流式对话与工具调用交互
+
+## 一、整体架构
+
+### 后端 SSE 事件流实现(已完成)
+
+---
+
+## 二、HTTP 接口文档(完整)
+
+### 1. 接口地址
+```
+POST /v1/chat/completions
+Content-Type: application/json
+```
+
+### 2. 请求参数(完整)
+
+| 参数名 | 类型 | 必填 | 说明 |
+|--------|------|------|------|
+| `messages` | Array | 是 | OpenAI 兼容的消息数组,例如 `[{"role": "user", "content": "你好"}]` |
+| `stream` | Boolean | 否 | 是否流式输出,默认 `false` |
+| `is_think` | Boolean | 否 | 是否启用 ReAct 思考模式,默认 `false` |
+| `thread_id` | String | 否 | 对话 ID,用于保持会话记忆,不传则自动生成 |
+| `agent_id` | String | 否 | Agent ID,不传则使用默认 Agent |
+| `model` | String | 否 | 模型名称,优先级:接口参数 > 配置文件 > 默认值 |
+| `file_ids` | Array | 否 | 多模态文件 ID 列表(可选) |
+
+### 3. 四种组合的返回格式
+
+后端会根据 `stream` 和 `is_think` 参数自动判断返回格式:
+
+| 组合 | 说明 | 返回格式 |
+|------|------|---------|
+| `stream=false` + `is_think=false` | 非流式 + 普通 Agent | OpenAI 兼容 JSON |
+| `stream=false` + `is_think=true` | 非流式 + ReAct Agent | OpenAI 兼容 JSON |
+| `stream=true` + `is_think=false` | 流式 + 普通 Agent | OpenAI 兼容 SSE 流 |
+| `stream=true` + `is_think=true` | 流式 + ReAct Agent | **SSE 事件流(包含 RAG/思考/工具调用事件)** ⭐ |
+
+---
+
+## 三、SSE 事件格式(stream=true + is_think=true)
+
+### 1. 事件类型定义
+
+当 `stream=true` 且 `is_think=true` 时,后端返回以下事件:
+
+| 事件类型 | 说明 | 数据字段 |
+|---------|------|---------|
+| `rag_start` | RAG 检索开始 | - |
+| `rag_retrieve` | RAG 检索到结果 | `{ resultCount: 5 }` |
+| `rag_key_info` | 提取到 RAG 关键信息 | `{ keyInfo: "{\"key_points\": [...], \"summary\": \"...\"}" }` |
+| `rag_end` | RAG 处理完成 | - |
+| `thinking_start` | 开始思考阶段 | - |
+| `thinking_end` | 结束思考阶段 | - |
+| `tool_call_start` | 开始调用工具 | `{ tool: "工具名", args: {...} }` |
+| `tool_call_end` | 工具调用完成 | `{ tool: "工具名", result: {...} }` |
+| `content_chunk` | 文本内容块 | `"文本内容"` |
+| `done` | 对话完成 | - |
+| `error` | 发生错误 | `"错误信息"` |
+
+### 2. SSE 输出格式示例
+
+```
+data: {"type":"rag_start","timestamp":1712456789000}
+
+data: {"type":"rag_retrieve","data":{"resultCount":5},"timestamp":1712456789001}
+
+data: {"type":"rag_key_info","data":{"keyInfo":"{\"key_points\": [...], \"summary\": \"...\"}"},"timestamp":1712456789002}
+
+data: {"type":"rag_end","timestamp":1712456789003}
+
+data: {"type":"thinking_start","timestamp":1712456789004}
+
+data: {"type":"tool_call_start","data":{"tool":"search_web","args":{"query":"今天天气"}},"timestamp":1712456789005}
+
+data: {"type":"tool_call_end","data":{"tool":"search_web","result":{"temperature":"25°C"}},"timestamp":1712456789006}
+
+data: {"type":"content_chunk","data":"今天","timestamp":1712456789007}
+
+data: {"type":"content_chunk","data":"天气","timestamp":1712456789008}
+
+data: {"type":"content_chunk","data":"很好","timestamp":1712456789009}
+
+data: {"type":"done","timestamp":1712456789010}
+
+```
+
+---
+
+## 四、前端实现
+
+### 技术栈
+- **框架**: React 18+
+- **状态管理**: Zustand / React Context
+- **流式通信**: fetch ReadableStream(比 EventSource 更灵活)
+- **UI组件**: shadcn/ui + Tailwind CSS
+
+---
+
+## 五、核心功能设计
+
+### 1. 打字机效果(流式文本输出)
+
+#### 实现思路
+```tsx
+// 使用 fetch + ReadableStream 实现
+async function streamChat(messages: ChatMessage[]) {
+  const response = await fetch('/api/chat/stream', {
+    method: 'POST',
+    headers: { 'Content-Type': 'application/json' },
+    body: JSON.stringify({ messages })
+  });
+  
+  const reader = response.body!.getReader();
+  const decoder = new TextDecoder();
+  let buffer = '';
+  
+  while (true) {
+    const { done, value } = await reader.read();
+    if (done) break;
+    
+    buffer += decoder.decode(value, { stream: true });
+    
+    // 解析 SSE 格式或直接处理文本
+    processBuffer(buffer);
+  }
+}
+```
+
+#### 打字机效果 Hook
+```tsx
+function useTypewriter(content: string, speed: number = 10) {
+  const [displayText, setDisplayText] = useState('');
+  const [index, setIndex] = useState(0);
+  
+  useEffect(() => {
+    if (index < content.length) {
+      const timer = setTimeout(() => {
+        setDisplayText(prev => prev + content[index]);
+        setIndex(prev => prev + 1);
+      }, speed);
+      return () => clearTimeout(timer);
+    }
+  }, [index, content, speed]);
+  
+  return { displayText, isTyping: index < content.length };
+}
+```
+
+---
+
+### 2. RAG 检索状态展示
+
+#### 组件结构
+```
+┌─────────────────────────────────────┐
+│                                     │
+│  🔍 检索知识库中...                  │
+│                                     │
+│  ┌─────────────────────────────┐   │
+│  │ ✨ 找到 5 条相关文档        │   │
+│  │ 📝 提取关键信息中...        │   │
+│  └─────────────────────────────┘   │
+│                                     │
+└─────────────────────────────────────┘
+```
+
+#### 状态管理
+```typescript
+type RagState = {
+  isRetrieving: boolean;
+  resultCount: number | null;
+  keyInfo: string | null;
+};
+```
+
+#### 组件实现
+```tsx
+function RagStatus({ ragState }: { ragState: RagState }) {
+  if (!ragState.isRetrieving && !ragState.resultCount) {
+    return null;
+  }
+
+  return (
+    <div className="flex items-center gap-3 p-3 rounded-lg bg-purple-50 dark:bg-purple-900/20 border border-purple-200 dark:border-purple-800">
+      {/* RAG 图标 */}
+      {ragState.isRetrieving ? (
+        <div className="animate-spin w-4 h-4 border-2 border-purple-500 border-t-transparent rounded-full" />
+      ) : (
+        <Database className="w-4 h-4 text-purple-500" />
+      )}
+      
+      {/* 状态文本 */}
+      <div className="text-sm">
+        {ragState.isRetrieving ? (
+          <span className="text-purple-700 dark:text-purple-300">
+            检索知识库中...
+          </span>
+        ) : (
+          <span className="text-purple-700 dark:text-purple-300">
+            找到 <span className="font-bold">{ragState.resultCount}</span> 条相关文档
+          </span>
+        )}
+        
+        {/* 关键信息提示 */}
+        {ragState.keyInfo && (
+          <div className="mt-1 text-xs text-purple-600 dark:text-purple-400">
+            已提取关键信息
+          </div>
+        )}
+      </div>
+    </div>
+  );
+}
+```
+
+---
+
+### 3. 工具调用小窗口(思考动画)
+
+#### 组件结构
+```
+┌─────────────────────────────────────┐
+│                                     │
+│  🤔 思考中...                       │
+│                                     │
+│  ┌─────────────────────────────┐   │
+│  │ ✨ 正在调用 get_weather     │   │
+│  │    Loading spinner          │   │
+│  └─────────────────────────────┘   │
+│                                     │
+└─────────────────────────────────────┘
+```
+
+#### 状态管理
+```typescript
+type ToolCallState = {
+  isThinking: boolean;
+  currentTool: string | null;
+  toolCalls: Array<{
+    id: string;
+    name: string;
+    status: 'pending' | 'running' | 'success' | 'error';
+    result?: any;
+    timestamp: number;
+  }>;
+};
+```
+
+#### 动画实现
+```tsx
+function ToolCallCard({ toolCall }: { toolCall: ToolCall }) {
+  return (
+    <div className="flex items-center gap-3 p-3 rounded-lg bg-gray-50 dark:bg-gray-800">
+      {/* 状态图标 */}
+      {toolCall.status === 'running' && (
+        <div className="animate-spin w-4 h-4 border-2 border-blue-500 border-t-transparent rounded-full" />
+      )}
+      {toolCall.status === 'success' && <CheckCircle2 className="w-4 h-4 text-green-500" />}
+      {toolCall.status === 'error' && <XCircle className="w-4 h-4 text-red-500" />}
+      
+      {/* 工具名称 */}
+      <span className="font-mono text-sm">{toolCall.name}</span>
+      
+      {/* 结果预览 */}
+      {toolCall.status === 'success' && toolCall.result && (
+        <span className="text-xs text-gray-500 truncate">
+          {JSON.stringify(toolCall.result).slice(0, 50)}...
+        </span>
+      )}
+    </div>
+  );
+}
+```
+
+---
+
+## 六、后端事件类型(建议)
+
+### SSE 事件格式
+```typescript
+// 事件类型定义
+type StreamEvent = 
+  | { type: 'rag_start' }
+  | { type: 'rag_retrieve'; data: { resultCount: number } }
+  | { type: 'rag_key_info'; data: { keyInfo: string } }
+  | { type: 'rag_end' }
+  | { type: 'thinking_start' }
+  | { type: 'thinking_end' }
+  | { type: 'tool_call_start'; data: { tool: string; args: any } }
+  | { type: 'tool_call_end'; data: { tool: string; result: any } }
+  | { type: 'content_chunk'; data: string }
+  | { type: 'done' }
+  | { type: 'error'; data: string };
+```
+
+### 后端输出示例
+```
+event: rag_start
+data: {}
+
+event: rag_retrieve
+data: {"resultCount": 5}
+
+event: rag_key_info
+data: {"keyInfo": "{\"key_points\": [...], \"summary\": \"...\"}"}
+
+event: rag_end
+data: {}
+
+event: thinking_start
+data: {}
+
+event: tool_call_start
+data: {"tool": "search_web", "args": {"query": "今天天气"}}
+
+event: tool_call_end
+data: {"tool": "search_web", "result": {"temperature": "25°C"}}
+
+event: content_chunk
+data: "今天"
+
+event: content_chunk
+data: "天气"
+
+event: content_chunk
+data: "很好"
+
+event: done
+data: {}
+```
+
+---
+
+## 七、前端 SSE 客户端实现
+
+### 使用 fetch + ReadableStream
+
+```typescript
+async function* streamChat(messages: ChatMessage[], threadId?: string) {
+  const response = await fetch('/api/chat/completions', {
+    method: 'POST',
+    headers: { 'Content-Type': 'application/json' },
+    body: JSON.stringify({
+      messages,
+      thread_id: threadId,
+      is_think: true,
+      stream: true
+    })
+  });
+
+  if (!response.ok) {
+    throw new Error('Request failed');
+  }
+
+  const reader = response.body!.getReader();
+  const decoder = new TextDecoder();
+  let buffer = '';
+
+  while (true) {
+    const { done, value } = await reader.read();
+    if (done) break;
+
+    buffer += decoder.decode(value, { stream: true });
+    
+    // 解析 SSE 格式
+    const lines = buffer.split('\n');
+    buffer = lines.pop() || ''; // 保留不完整的行
+
+    for (const line of lines) {
+      if (line.startsWith('data: ')) {
+        const jsonStr = line.slice(6);
+        if (jsonStr === '[DONE]') continue;
+        
+        try {
+          const event = JSON.parse(jsonStr);
+          yield event;
+        } catch (e) {
+          console.error('Parse error:', e);
+        }
+      }
+    }
+  }
+}
+```
+
+### React Hook 封装
+
+```typescript
+function useChatStream() {
+  const [messages, setMessages] = useState<ChatMessage[]>([]);
+  const [isLoading, setIsLoading] = useState(false);
+  const [toolCalls, setToolCalls] = useState<ToolCall[]>([]);
+  const [isThinking, setIsThinking] = useState(false);
+  // RAG 状态
+  const [ragState, setRagState] = useState<RagState>({
+    isRetrieving: false,
+    resultCount: null,
+    keyInfo: null,
+  });
+
+  const sendMessage = async (content: string) => {
+    const userMessage: ChatMessage = { role: 'user', content };
+    setMessages(prev => [...prev, userMessage]);
+    setIsLoading(true);
+    setToolCalls([]);
+    // 重置 RAG 状态
+    setRagState({
+      isRetrieving: false,
+      resultCount: null,
+      keyInfo: null,
+    });
+
+    const assistantMessage: ChatMessage = { 
+      role: 'assistant', 
+      content: '', 
+      toolCalls: [] 
+    };
+
+    try {
+      for await (const event of streamChat([...messages, userMessage])) {
+        switch (event.type) {
+          // ========== RAG 事件处理 ==========
+          case 'rag_start':
+            setRagState(prev => ({ ...prev, isRetrieving: true }));
+            break;
+          case 'rag_retrieve':
+            setRagState(prev => ({ ...prev, resultCount: event.data.resultCount }));
+            break;
+          case 'rag_key_info':
+            setRagState(prev => ({ ...prev, keyInfo: event.data.keyInfo }));
+            break;
+          case 'rag_end':
+            setRagState(prev => ({ ...prev, isRetrieving: false }));
+            break;
+          
+          // ========== 思考事件处理 ==========
+          case 'thinking_start':
+            setIsThinking(true);
+            break;
+          case 'thinking_end':
+            setIsThinking(false);
+            break;
+          
+          // ========== 工具调用事件处理 ==========
+          case 'tool_call_start':
+            setToolCalls(prev => [...prev, {
+              id: Date.now().toString(),
+              name: event.data.tool,
+              status: 'running',
+              args: event.data.args,
+              timestamp: event.timestamp
+            }]);
+            break;
+          case 'tool_call_end':
+            setToolCalls(prev => prev.map(tc => 
+              tc.name === event.data.tool 
+                ? { ...tc, status: 'success', result: event.data.result }
+                : tc
+            ));
+            break;
+          
+          // ========== 内容事件处理 ==========
+          case 'content_chunk':
+            assistantMessage.content += event.data;
+            setMessages(prev => [
+              ...prev.slice(0, -1),
+              { ...assistantMessage }
+            ]);
+            break;
+          
+          // ========== 结束/错误事件处理 ==========
+          case 'done':
+            setIsLoading(false);
+            break;
+          case 'error':
+            console.error('Stream error:', event.data);
+            setIsLoading(false);
+            setRagState({ isRetrieving: false, resultCount: null, keyInfo: null });
+            break;
+        }
+      }
+    } catch (error) {
+      console.error('Chat error:', error);
+      setIsLoading(false);
+      setRagState({ isRetrieving: false, resultCount: null, keyInfo: null });
+    }
+  };
+
+  return { messages, isLoading, isThinking, toolCalls, ragState, sendMessage };
+}
+```
+
+---
+
+## 八、完整对话组件示例
+
+```tsx
+function ChatPage() {
+  const { messages, isLoading, isThinking, toolCalls, ragState, sendMessage } = useChatStream();
+
+  return (
+    <div className="flex flex-col h-screen">
+      {/* 聊天消息区域 */}
+      <div className="flex-1 overflow-y-auto p-4 space-y-4">
+        {messages.map((message, index) => (
+          <ChatMessage key={index} message={message} />
+        ))}
+        
+        {/* 加载状态 - RAG + 思考 + 工具调用 */}
+        {isLoading && (
+          <div className="flex gap-4 justify-start">
+            <div className="w-8 h-8 rounded-full bg-blue-500 flex items-center justify-center">
+              🤖
+            </div>
+            <div className="flex-1 space-y-3">
+              {/* RAG 状态组件 */}
+              <RagStatus ragState={ragState} />
+              
+              {/* 工具调用状态组件 */}
+              {toolCalls.length > 0 && (
+                <div className="space-y-2">
+                  {toolCalls.map(tool => (
+                    <ToolCallCard key={tool.id} toolCall={tool} />
+                  ))}
+                </div>
+              )}
+              
+              {/* 思考中状态 */}
+              {isThinking && (
+                <div className="flex items-center gap-2 p-3 rounded-lg bg-gray-50 dark:bg-gray-800">
+                  <div className="animate-pulse w-2 h-2 rounded-full bg-blue-500" />
+                  <span className="text-sm text-gray-500">思考中...</span>
+                </div>
+              )}
+            </div>
+          </div>
+        )}
+      </div>
+      
+      {/* 输入框 */}
+      <ChatInput onSend={sendMessage} disabled={isLoading} />
+    </div>
+  );
+}
+
+function ChatMessage({ message }: { message: ChatMessage }) {
+  const [showTools, setShowTools] = useState(false);
+  
+  return (
+    <div className={`flex gap-4 ${message.role === 'user' ? 'justify-end' : 'justify-start'}`}>
+      <div className={`max-w-[80%] ${message.role === 'user' ? 'order-2' : 'order-1'}`}>
+        {/* 头像 */}
+        <div className="w-8 h-8 rounded-full bg-blue-500 flex items-center justify-center">
+          {message.role === 'user' ? '👤' : '🤖'}
+        </div>
+      </div>
+      
+      <div className={`flex-1 ${message.role === 'user' ? 'order-1' : 'order-2'}`}>
+        {/* 工具调用区域 */}
+        {message.toolCalls && message.toolCalls.length > 0 && (
+          <div className="mb-2">
+            <button 
+              onClick={() => setShowTools(!showTools)}
+              className="text-xs text-blue-500 hover:underline"
+            >
+              {showTools ? '隐藏' : '显示'} 工具调用 ({message.toolCalls.length})
+            </button>
+            
+            {showTools && (
+              <div className="mt-2 space-y-2">
+                {message.toolCalls.map(tool => (
+                  <ToolCallCard key={tool.id} toolCall={tool} />
+                ))}
+              </div>
+            )}
+          </div>
+        )}
+        
+        {/* 消息内容 - 打字机效果 */}
+        <div className="bg-white dark:bg-gray-800 rounded-lg p-4 shadow-sm">
+          <TypewriterText content={message.content} />
+        </div>
+      </div>
+    </div>
+  );
+}
+```
+
+---
+
+## 九、优化建议
+
+1. **虚拟滚动**: 对话历史很长时使用
+2. **Markdown渲染**: 使用 `react-markdown` 渲染富文本
+3. **代码高亮**: 使用 `prismjs` 或 `shiki`
+4. **暂停/恢复**: 支持暂停和恢复流式输出
+5. **错误重试**: 网络断开时自动重连
+
+---
+
+## 十、参考实现库
+
+- [Vercel AI SDK](https://sdk.vercel.ai/) - 完整的 AI 聊天 SDK
+- [LangChain UI](https://js.langchain.com/docs/modules/chains/popular/chat_vector_db) - LangChain 官方 UI 组件
+- [Chatbot UI](https://github.com/mckaywrigley/chatbot-ui) - 开源聊天界面参考