fix(sse): SSE 改为按会话隔离 + 修 trace 租户 bug + 清理冗余

SSE 串台修复:
- SseEmitterManager 新增按 sessionId 维度的 connect/sendEvent/disconnect,
  每个会话一个 SSE 连接,替代原 userId+token 维度(同用户多会话串台)
- SseMessageDto 增加 sessionId + eventDto 字段,跨实例按会话路由
- SseTopicListener 优先按 sessionId 路由,回退原 userId/群发逻辑
- SseMessageUtils 增加 sessionId 重载,对话链路全部切换
- ChatServiceFacade / MyMcpClientListener / WorkflowStarter 切到 sessionId
- 通知/全局 /sse 端点保留 userId 模式(通知按用户)

trace 租户 bug 修复:
- trace_run / trace_node 加入 tenant.excludes,监控表跨租户全局可见,
  绕过异步线程租户上下文不传播导致 node 查不到的问题

冗余清理:
- 删除零使用的 @TraceNode 注解 + TraceNodeAspect 切面 + 对应单测
- TraceProperties 删除 recordDetail/maxInputLength/maxOutputLength
  (仅对已删切面生效,编程式埋点用 RagTracePayloadBuilder 不受影响)
- TracePayloadUtils 删除 input/output/asString 方法
- TraceConstants 删除无引用的 NODE_METHOD/HTTP/DB/CACHE/TASK/STREAM
- TraceContext 删除预留的 copyNodeStack
- TraceNodeVo 删除未用的 children 字段
- TraceRecordServiceImpl.getDetail 删除重复 enrich
- application.yml 删除 trace.payload 无效配置项
- RagTraceNodeTypes 删除无引用的 NODE_STREAM

Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
ageerle
2026-07-23 15:35:10 +08:00
parent 11bb1dba0f
commit d1a820728f
21 changed files with 259 additions and 405 deletions

View File

@@ -32,6 +32,16 @@ public class SseEmitterManager {
*/
private final static String SSE_TOPIC = "global:sse";
/**
* 按会话维度管理:每个会话一个 SSE 连接,用于对话流式响应
* Key: sessionId
*/
private final static Map<String, SseEmitter> SESSION_EMITTERS = new ConcurrentHashMap<>();
/**
* 按用户维度管理:全局通知、站内信等场景,一个用户可有多个连接(按 token 区分)
* Key: userId, Value: token -> SseEmitter
*/
private final static Map<Long, Map<String, SseEmitter>> USER_TOKEN_EMITTERS = new ConcurrentHashMap<>();
public SseEmitterManager() {
@@ -40,6 +50,98 @@ public class SseEmitterManager {
.scheduleWithFixedDelay(this::sseMonitor, 60L, 60L, TimeUnit.SECONDS);
}
// ======================== 会话维度(对话流式响应) ========================
/**
* 建立与指定会话的 SSE 连接,每个会话仅保留一个连接,重复建连会替换旧连接。
*
* @param sessionId 会话ID
* @return SseEmitter 实例
*/
public SseEmitter connect(String sessionId) {
if (sessionId == null) {
throw new IllegalArgumentException("sessionId 不能为空");
}
// 关闭已存在的 SseEmitter保证每个会话只有一个活跃连接
SseEmitter oldEmitter = SESSION_EMITTERS.remove(sessionId);
if (oldEmitter != null) {
oldEmitter.complete();
}
SseEmitter emitter = new SseEmitter(86400000L);
SESSION_EMITTERS.put(sessionId, emitter);
emitter.onCompletion(() -> removeSessionEmitter(sessionId, emitter));
emitter.onTimeout(() -> removeSessionEmitter(sessionId, emitter));
emitter.onError(e -> removeSessionEmitter(sessionId, emitter));
try {
emitter.send(SseEmitter.event().comment("connected"));
} catch (IOException e) {
SESSION_EMITTERS.remove(sessionId);
}
return emitter;
}
/**
* 断开指定会话的 SSE 连接
*
* @param sessionId 会话ID
*/
public void disconnect(String sessionId) {
if (sessionId == null) {
return;
}
SseEmitter emitter = SESSION_EMITTERS.remove(sessionId);
if (emitter != null) {
try {
emitter.send(SseEmitter.event().comment("disconnected"));
} catch (Exception exception) {
log.error(exception.getMessage());
}
emitter.complete();
}
}
/**
* 向指定会话发送结构化事件
*
* @param sessionId 会话ID
* @param eventDto SSE事件对象
*/
public void sendEvent(String sessionId, SseEventDto eventDto) {
if (sessionId == null || eventDto == null) {
return;
}
SseEmitter emitter = SESSION_EMITTERS.get(sessionId);
if (emitter == null) {
log.warn("【SSE发送失败】sessionId: {} 没有活跃的SSE连接", sessionId);
return;
}
try {
log.debug("【SSE发送】sessionId: {}, event: {}", sessionId, eventDto.getEvent());
emitter.send(SseEmitter.event()
.name(eventDto.getEvent())
.data(JSONUtil.toJsonStr(eventDto)));
} catch (Exception e) {
log.error("【SSE发送失败】sessionId: {}, error: {}", sessionId, e.getMessage());
removeSessionEmitter(sessionId, emitter);
}
}
private void removeSessionEmitter(String sessionId, SseEmitter emitter) {
boolean removed = SESSION_EMITTERS.remove(sessionId, emitter);
if (removed) {
try {
emitter.complete();
} catch (Exception ignore) {
// 忽略重复关闭异常
}
}
}
// ======================== 用户维度(全局通知) ========================
/**
* 建立与指定用户的 SSE 连接
*
@@ -154,6 +256,23 @@ public class SseEmitterManager {
// 循环结束后统一清理空用户,避免并发修改异常
toRemoveUsers.forEach(USER_TOKEN_EMITTERS::remove);
// 会话维度心跳:发送失败的连接移除
if (!SESSION_EMITTERS.isEmpty()) {
SESSION_EMITTERS.entrySet().removeIf(entry -> {
try {
entry.getValue().send(heartbeat);
return false;
} catch (Exception ex) {
try {
entry.getValue().complete();
} catch (Exception ignore) {
// 忽略重复关闭异常
}
return true;
}
});
}
}
/**
@@ -243,9 +362,11 @@ public class SseEmitterManager {
SseMessageDto broadcastMessage = new SseMessageDto();
broadcastMessage.setMessage(sseMessageDto.getMessage());
broadcastMessage.setUserIds(sseMessageDto.getUserIds());
broadcastMessage.setSessionId(sseMessageDto.getSessionId());
broadcastMessage.setEventDto(sseMessageDto.getEventDto());
RedisUtils.publish(SSE_TOPIC, broadcastMessage, consumer -> {
log.info("SSE发送主题订阅消息topic:{} session keys:{} message:{}",
SSE_TOPIC, sseMessageDto.getUserIds(), sseMessageDto.getMessage());
log.info("SSE发送主题订阅消息topic:{} session:{} session keys:{} message:{}",
SSE_TOPIC, sseMessageDto.getSessionId(), sseMessageDto.getUserIds(), sseMessageDto.getMessage());
});
}

View File

@@ -26,4 +26,14 @@ public class SseMessageDto implements Serializable {
* 需要发送的消息
*/
private String message;
/**
* 按会话定向推送的会话ID非空时优先按会话路由忽略 userIds
*/
private String sessionId;
/**
* 结构化事件按会话定向推送时使用message 为兼容旧逻辑保留)
*/
private SseEventDto eventDto;
}

View File

@@ -1,8 +1,10 @@
package org.ruoyi.common.sse.listener;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.util.StrUtil;
import lombok.extern.slf4j.Slf4j;
import org.ruoyi.common.sse.core.SseEmitterManager;
import org.ruoyi.common.sse.dto.SseMessageDto;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
@@ -28,8 +30,20 @@ public class SseTopicListener implements ApplicationRunner, Ordered {
@Override
public void run(ApplicationArguments args) throws Exception {
sseEmitterManager.subscribeMessage((message) -> {
log.info("SSE主题订阅收到消息session keys={} message={}", message.getUserIds(), message.getMessage());
// 如果key不为空就按照key发消息 如果为空就群发
log.info("SSE主题订阅收到消息session:{} session keys={} message={}",
message.getSessionId(), message.getUserIds(), message.getMessage());
// 优先按会话路由(对话流式响应)
if (StrUtil.isNotBlank(message.getSessionId())) {
if (message.getEventDto() != null) {
sseEmitterManager.sendEvent(message.getSessionId(), message.getEventDto());
} else if (message.getMessage() != null) {
// 兼容按会话发纯文本的场景
sseEmitterManager.sendEvent(message.getSessionId(),
org.ruoyi.common.sse.dto.SseEventDto.content(message.getMessage()));
}
return;
}
// 否则按用户/群发路由(全局通知)
if (CollUtil.isNotEmpty(message.getUserIds())) {
message.getUserIds().forEach(key -> {
sseEmitterManager.sendMessage(key, message.getMessage());

View File

@@ -93,6 +93,15 @@ public class SseMessageUtils {
MANAGER.disconnect(userId, tokenValue);
}
/**
* 完成指定会话的SSE连接
*
* @param sessionId 会话ID
*/
public static void completeConnection(String sessionId) {
MANAGER.disconnect(sessionId);
}
/**
* 向指定的SSE会话发送结构化事件
*
@@ -106,6 +115,22 @@ public class SseMessageUtils {
MANAGER.sendEvent(userId, eventDto);
}
/**
* 向指定会话发送结构化事件(通过 Redis 广播,跨实例可达)
*
* @param sessionId 会话ID
* @param eventDto SSE事件对象
*/
public static void sendEvent(String sessionId, SseEventDto eventDto) {
if (!isEnable() || sessionId == null) {
return;
}
SseMessageDto dto = new SseMessageDto();
dto.setSessionId(sessionId);
dto.setEventDto(eventDto);
MANAGER.publishMessage(dto);
}
/**
* 发送内容事件
*
@@ -116,6 +141,16 @@ public class SseMessageUtils {
sendEvent(userId, SseEventDto.content(content));
}
/**
* 向指定会话发送内容事件
*
* @param sessionId 会话ID
* @param content 内容
*/
public static void sendContent(String sessionId, String content) {
sendEvent(sessionId, SseEventDto.content(content));
}
/**
* 发送推理内容事件
*
@@ -126,6 +161,16 @@ public class SseMessageUtils {
sendEvent(userId, SseEventDto.reasoning(reasoningContent));
}
/**
* 向指定会话发送推理内容事件
*
* @param sessionId 会话ID
* @param reasoningContent 推理内容
*/
public static void sendReasoning(String sessionId, String reasoningContent) {
sendEvent(sessionId, SseEventDto.reasoning(reasoningContent));
}
/**
* 发送完成事件
*
@@ -135,6 +180,15 @@ public class SseMessageUtils {
sendEvent(userId, SseEventDto.done());
}
/**
* 向指定会话发送完成事件
*
* @param sessionId 会话ID
*/
public static void sendDone(String sessionId) {
sendEvent(sessionId, SseEventDto.done());
}
/**
* 发送错误事件
*
@@ -145,6 +199,16 @@ public class SseMessageUtils {
sendEvent(userId, SseEventDto.error(error));
}
/**
* 向指定会话发送错误事件
*
* @param sessionId 会话ID
* @param error 错误信息
*/
public static void sendError(String sessionId, String error) {
sendEvent(sessionId, SseEventDto.error(error));
}
/**
* 是否开启
*/