|
|
@@ -44,7 +44,7 @@ import java.util.*;
|
|
|
import java.util.concurrent.ConcurrentHashMap;
|
|
|
import java.util.concurrent.atomic.AtomicBoolean;
|
|
|
import java.util.concurrent.atomic.AtomicReference;
|
|
|
-import java.util.concurrent.locks.ReentrantLock;
|
|
|
+import java.util.concurrent.Semaphore;
|
|
|
import java.util.stream.Collectors;
|
|
|
|
|
|
/**
|
|
|
@@ -70,14 +70,19 @@ public class AgentChatServiceImpl extends ServiceImpl<AgentChatSessionMapper, Ag
|
|
|
private static final String DEFAULT_TITLE = "新对话";
|
|
|
|
|
|
/**
|
|
|
- * 会话级流互斥锁:同一会话同一时刻只允许一条 AI 流(对照研判链路 {@code InsightServiceImpl.sessionLocks})。
|
|
|
+ * 会话级流互斥闸门:同一会话同一时刻只允许一条 AI 流(对照研判链路 {@code InsightServiceImpl.sessionLocks})。
|
|
|
*
|
|
|
* <p>双标签页/双击重试对同一会话并发发消息时,两条流会交错写同一 Agent 记忆槽
|
|
|
* {@code ("default", sessionId)} 与同一工作区 transcript 文件,产生记忆污染。
|
|
|
- * 非阻塞 tryLock:抢不到立即返回错误帧,不排队不阻塞请求线程。
|
|
|
+ * 非阻塞 tryAcquire:抢不到立即返回错误帧,不排队不阻塞请求线程。
|
|
|
* 进程内闸门(单实例部署假设);多实例部署时需再补 DB 状态位。
|
|
|
+ *
|
|
|
+ * <p>必须用 {@link Semaphore} 而不是 {@code ReentrantLock}:抢占发生在订阅线程(Tomcat 请求线程),
|
|
|
+ * 释放发生在 {@code doFinally}(Reactor 在上游事件线程回调),ReentrantLock 跨线程 unlock 会抛
|
|
|
+ * IllegalMonitorStateException,闸门永久不释放 —— 表现为第一条对话结束后,
|
|
|
+ * 之后每条消息都误报「该会话正在回复中」。
|
|
|
*/
|
|
|
- private final ConcurrentHashMap<Long, ReentrantLock> streamLocks = new ConcurrentHashMap<>();
|
|
|
+ private final ConcurrentHashMap<Long, Semaphore> streamLocks = new ConcurrentHashMap<>();
|
|
|
|
|
|
private final AgentChatSessionMapper chatSessionMapper;
|
|
|
private final AgentMessageMapper messageMapper;
|
|
|
@@ -301,14 +306,17 @@ public class AgentChatServiceImpl extends ServiceImpl<AgentChatSessionMapper, Ag
|
|
|
return Flux.just(sse("error", Map.of("type", "error", "error", "消息内容不能为空")));
|
|
|
}
|
|
|
// defer:订阅时才解析会话/构建 agent,异常统一转 error 帧
|
|
|
- return Flux.defer(() -> doStream(req))
|
|
|
+ // 过程累积器随请求创建并下传:出错时错误消息要把「深度思考 + 工具调用 + 产物元数据」一并落库,
|
|
|
+ // 否则出错轮次刷新后只剩一行错误提示,本轮跑过的思考/工具/产物全部消失
|
|
|
+ StreamProcess process = new StreamProcess();
|
|
|
+ return Flux.defer(() -> doStream(req, process))
|
|
|
.onErrorResume(ex -> {
|
|
|
String errorMsg = ex.getMessage() != null ? ex.getMessage() : "未知错误";
|
|
|
log.warn("Chat stream error: sessionKey={}, ", req.getSessionKey(), ex);
|
|
|
- // 错误信息持久化:确保刷新后仍可见
|
|
|
+ // 错误信息持久化:确保刷新后仍可见(含本轮已产生的过程数据)
|
|
|
Long sessionId = parseSessionId(req.getSessionKey());
|
|
|
if (sessionId != null) {
|
|
|
- Mono.fromRunnable(() -> persistErrorMessage(sessionId, errorMsg))
|
|
|
+ Mono.fromRunnable(() -> persistErrorMessage(sessionId, errorMsg, process))
|
|
|
.subscribeOn(Schedulers.boundedElastic())
|
|
|
.subscribe();
|
|
|
}
|
|
|
@@ -324,28 +332,28 @@ public class AgentChatServiceImpl extends ServiceImpl<AgentChatSessionMapper, Ag
|
|
|
* 视为动态切换模型并回写会话。
|
|
|
* 记忆槽位 = (运行时用户标识, sessionId),与 agent 实例/模型无关,切换模型后记忆自动延续。
|
|
|
*/
|
|
|
- private Flux<ServerSentEvent<String>> doStream(ChatRequestDTO req) {
|
|
|
+ private Flux<ServerSentEvent<String>> doStream(ChatRequestDTO req, StreamProcess process) {
|
|
|
Long sessionId = parseSessionId(req.getSessionKey());
|
|
|
if (sessionId == null) {
|
|
|
throw new ServerException(400, "缺少有效会话ID(sessionKey)");
|
|
|
}
|
|
|
// 会话级防并发双流:同一会话同时只允许一条 AI 流。
|
|
|
- // 非阻塞抢锁,抢不到立即错误帧;锁在流终态(完成/异常/取消)由 doFinally 释放,
|
|
|
+ // 非阻塞抢闸门,抢不到立即错误帧;闸门在流终态(完成/异常/取消)由 doFinally 释放,
|
|
|
// doStreamLocked 同步抛异常(404/400 等)时由 catch 释放
|
|
|
- ReentrantLock lock = streamLocks.computeIfAbsent(sessionId, k -> new ReentrantLock());
|
|
|
- if (!lock.tryLock()) {
|
|
|
+ Semaphore gate = streamLocks.computeIfAbsent(sessionId, k -> new Semaphore(1));
|
|
|
+ if (!gate.tryAcquire()) {
|
|
|
return Flux.just(sse("error", Map.of("type", "error", "error", "该会话正在回复中,请等待当前回复完成后再发送")));
|
|
|
}
|
|
|
try {
|
|
|
- return doStreamLocked(sessionId, req)
|
|
|
- .doFinally(signal -> lock.unlock());
|
|
|
+ return doStreamLocked(sessionId, req, process)
|
|
|
+ .doFinally(signal -> gate.release());
|
|
|
} catch (RuntimeException e) {
|
|
|
- lock.unlock();
|
|
|
+ gate.release();
|
|
|
throw e;
|
|
|
}
|
|
|
}
|
|
|
|
|
|
- private Flux<ServerSentEvent<String>> doStreamLocked(Long sessionId, ChatRequestDTO req) {
|
|
|
+ private Flux<ServerSentEvent<String>> doStreamLocked(Long sessionId, ChatRequestDTO req, StreamProcess process) {
|
|
|
AgentChatSession session = requireSessionInOpenCase(requireSession(sessionId));
|
|
|
|
|
|
// 附件:归属校验 + 注入块(内联全文 / 表格 Python 直读指引 / RAG 片段)。
|
|
|
@@ -406,16 +414,17 @@ public class AgentChatServiceImpl extends ServiceImpl<AgentChatSessionMapper, Ag
|
|
|
long startMs = System.currentTimeMillis();
|
|
|
String sessionKey = String.valueOf(sessionId);
|
|
|
|
|
|
- // 工具调用状态聚合(流内累积,随流结束释放)
|
|
|
- Map<String, StringBuilder> toolInputs = new ConcurrentHashMap<>();
|
|
|
- Map<String, StringBuilder> toolResults = new ConcurrentHashMap<>();
|
|
|
- Map<String, String> toolNames = new ConcurrentHashMap<>();
|
|
|
+ // 工具调用状态聚合(流内累积,随流结束释放)—— 统一放在 StreamProcess 里,
|
|
|
+ // 出错时 persistErrorMessage 要用同一份数据落库,保住本轮的工具调用与产物元数据
|
|
|
+ Map<String, StringBuilder> toolInputs = process.toolInputs;
|
|
|
+ Map<String, StringBuilder> toolResults = process.toolResults;
|
|
|
+ Map<String, String> toolNames = process.toolNames;
|
|
|
|
|
|
// 流中断追踪:用户主动中止时,若 AgentResultEvent 尚未到达,将已累积的文本保存
|
|
|
- AtomicReference<StringBuilder> accumulatedText = new AtomicReference<>(new StringBuilder());
|
|
|
+ AtomicReference<StringBuilder> accumulatedText = process.text;
|
|
|
// 模型思考过程累积(reasoning_content / thinking):随流结束后写入消息 metadata
|
|
|
- AtomicReference<StringBuilder> accumulatedThinking = new AtomicReference<>(new StringBuilder());
|
|
|
- AtomicBoolean resultPersisted = new AtomicBoolean(false);
|
|
|
+ AtomicReference<StringBuilder> accumulatedThinking = process.thinking;
|
|
|
+ AtomicBoolean resultPersisted = process.resultPersisted;
|
|
|
|
|
|
return events.concatMap(event -> {
|
|
|
List<ServerSentEvent<String>> frames = new ArrayList<>();
|
|
|
@@ -651,16 +660,28 @@ public class AgentChatServiceImpl extends ServiceImpl<AgentChatSessionMapper, Ag
|
|
|
|
|
|
/**
|
|
|
* 持久化错误信息为 assistant 消息(messageType=error),确保刷新后仍可见。
|
|
|
+ *
|
|
|
+ * <p>本轮已产生的「深度思考 + 工具调用(含产物元数据)」必须一并落库:
|
|
|
+ * 出错前的过程是真实发生过的,只存一行错误提示会让刷新后的时间线
|
|
|
+ * 「深度思考后调用工具」整段消失、产物文件也无从回溯(历史会话看不到)。
|
|
|
*/
|
|
|
- private void persistErrorMessage(Long sessionId, String errorMsg) {
|
|
|
+ private void persistErrorMessage(Long sessionId, String errorMsg, StreamProcess process) {
|
|
|
try {
|
|
|
+ // 完整回答已落库时(AgentResultEvent 先到、后置环节才出错)不再重复携带过程数据,
|
|
|
+ // 避免同轮的时间线/产物在历史里出现两份
|
|
|
+ boolean processLost = process != null && !process.resultPersisted.get();
|
|
|
+ String toolEventsJson = processLost
|
|
|
+ ? buildToolEventsJson(process.toolInputs, process.toolResults, process.toolNames) : null;
|
|
|
+ String thinkingJson = processLost
|
|
|
+ ? buildThinkingMetadata(process.thinking.get().toString()) : null;
|
|
|
AgentMessage msg = new AgentMessage();
|
|
|
msg.setSessionId(sessionId);
|
|
|
msg.setRole("assistant");
|
|
|
msg.setContent("⚠️ " + errorMsg);
|
|
|
msg.setContentType("text");
|
|
|
msg.setMessageType("error");
|
|
|
- msg.setMetadata("{\"error\": true}");
|
|
|
+ msg.setMetadata(mergeThinkingMetadata("{\"error\": true}", thinkingJson));
|
|
|
+ msg.setToolEvents(toolEventsJson);
|
|
|
msg.setStarred(false);
|
|
|
msg.setCreateAt(LocalDateTime.now());
|
|
|
messageMapper.insert(msg);
|
|
|
@@ -871,4 +892,21 @@ public class AgentChatServiceImpl extends ServiceImpl<AgentChatSessionMapper, Ag
|
|
|
return ServerSentEvent.<String>builder().event(eventType).data(json).build();
|
|
|
}
|
|
|
|
|
|
+ /**
|
|
|
+ * 一轮流式对话的过程累积器:工具调用聚合 + 正文/思考累积。
|
|
|
+ *
|
|
|
+ * <p>生命周期 = 一次 {@link #stream} 请求:正常结束随 assistant 消息落库(AgentResultEvent),
|
|
|
+ * 用户中止随 interrupted 消息落库,<b>中途出错也要随错误消息落库</b> ——
|
|
|
+ * 三处共用同一份数据,任何终态都不丢过程记录。条目极少、随请求释放,不额外清理。
|
|
|
+ */
|
|
|
+ private static final class StreamProcess {
|
|
|
+ final Map<String, StringBuilder> toolInputs = new ConcurrentHashMap<>();
|
|
|
+ final Map<String, StringBuilder> toolResults = new ConcurrentHashMap<>();
|
|
|
+ final Map<String, String> toolNames = new ConcurrentHashMap<>();
|
|
|
+ final AtomicReference<StringBuilder> text = new AtomicReference<>(new StringBuilder());
|
|
|
+ final AtomicReference<StringBuilder> thinking = new AtomicReference<>(new StringBuilder());
|
|
|
+ /** 完整回答是否已落库(落库后出错不再重复携带过程数据) */
|
|
|
+ final AtomicBoolean resultPersisted = new AtomicBoolean(false);
|
|
|
+ }
|
|
|
+
|
|
|
}
|