From 83be40aeb577a59c5722401cdf1ed705ac1a9aa6 Mon Sep 17 00:00:00 2001 From: KDB <937925477@qq.com> Date: Thu, 13 Aug 2026 01:53:22 +0800 Subject: [PATCH] fix SSE trace correlation and cancellation --- AGENTS.md | 12 ++ .../ai/ragent/rag/dto/MetaPayload.java | 2 +- .../handler/StreamChatEventHandler.java | 26 +++- .../service/handler/StreamTaskManager.java | 54 ++++++-- .../service/pipeline/StreamChatPipeline.java | 55 ++++++-- .../service/ratelimit/ChatQueueLimiter.java | 4 +- .../rag/trace/StreamChatTraceRunner.java | 46 +++++-- .../rag/trace/StreamChatTraceRunnerTest.java | 122 ++++++++++++++++++ docs/external-governance-operations.md | 16 +++ docs/harness-remediation-tracker.md | 6 +- docs/verification-routing.md | 24 ++++ .../hooks/__tests__/useStreamResponse.test.ts | 6 +- frontend/src/stores/chatStore.ts | 5 +- frontend/src/types/index.ts | 1 + .../infra/chat/ForwardingStreamCallback.java | 5 + .../ai/ragent/infra/chat/StreamCallback.java | 8 ++ 16 files changed, 350 insertions(+), 42 deletions(-) create mode 100644 bootstrap/src/test/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunnerTest.java create mode 100644 docs/verification-routing.md diff --git a/AGENTS.md b/AGENTS.md index ba80bb1ea..1130a43dd 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -37,6 +37,18 @@ npm run build npm run dev # 开发服务器 5173,/api 代理到 localhost:9090 ``` +### 最小相关检查路由 + +局部改动先跑最小相关检查获得快速反馈,提交前仍按影响范围升级到完整门禁。完整映射与升级条件见 `docs/verification-routing.md`。 + +```bash +# 后端单个测试类;-am 场景必须关闭“未找到指定测试即失败” +./mvnw -B -ntp -pl bootstrap -am -Dtest=StreamChatTraceRunnerTest -Dsurefire.failIfNoSpecifiedTests=false test + +# 前端单个测试文件 +cd frontend && npm run test -- src/hooks/__tests__/useStreamResponse.test.ts +``` + ## 模块分层 ```text diff --git a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/dto/MetaPayload.java b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/dto/MetaPayload.java index 1a83bf4d8..4d69c298b 100644 --- a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/dto/MetaPayload.java +++ b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/dto/MetaPayload.java @@ -17,5 +17,5 @@ package com.nageoffer.ai.ragent.rag.dto; -public record MetaPayload(String conversationId, String taskId) { +public record MetaPayload(String conversationId, String taskId, String traceId) { } diff --git a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamChatEventHandler.java b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamChatEventHandler.java index c2039cc3c..482de8efd 100644 --- a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamChatEventHandler.java +++ b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamChatEventHandler.java @@ -79,18 +79,24 @@ public StreamChatEventHandler(StreamChatHandlerParams params) { this.messageChunkSize = resolveMessageChunkSize(params.getModelProperties()); this.sendTitleOnComplete = shouldSendTitle(); - // 初始化(发送初始事件、注册任务) + // 先返回 taskId,保证排队期间也可取消;Trace 建立后会发送一次带 traceId 的 META 更新。 initialize(); } /** - * 初始化:发送元数据事件并注册任务 + * 初始化:发送可取消所需的元数据并注册任务 */ private void initialize() { - sender.sendEvent(SSEEventType.META.value(), new MetaPayload(conversationId, taskId)); + sender.sendEvent(SSEEventType.META.value(), new MetaPayload(conversationId, taskId, null)); taskManager.register(taskId, sender, this::buildCompletionPayloadOnCancel); } + @Override + public void onTraceStarted(String traceId) { + taskManager.attachTrace(taskId, traceId); + sender.sendEvent(SSEEventType.META.value(), new MetaPayload(conversationId, taskId, traceId)); + } + /** * 解析消息块大小 */ @@ -127,7 +133,9 @@ private CompletionPayload buildCompletionPayloadOnCancel() { message.setMessageStatus(ChatMessage.MessageStatus.INTERRUPTED); messageId = memoryService.append(conversationId, userId, message); } catch (Exception e) { - log.error("取消时持久化消息失败,conversationId:{}", conversationId, e); + log.error("Failed to persist cancelled SSE message: traceId={}, taskId={}, errorType={}", + StreamTaskManager.safeCorrelationId(taskManager.traceId(taskId)), + StreamTaskManager.safeCorrelationId(taskId), e.getClass().getSimpleName()); } } String title = resolveTitleForEvent(); @@ -209,13 +217,18 @@ public void onComplete() { message.setMessageStatus(ChatMessage.MessageStatus.NORMAL); messageId = memoryService.append(conversationId, userId, message); } catch (Exception e) { - log.error("对话完成时持久化消息失败,conversationId:{}", conversationId, e); + log.error("Failed to persist completed SSE message: traceId={}, taskId={}, errorType={}", + StreamTaskManager.safeCorrelationId(taskManager.traceId(taskId)), + StreamTaskManager.safeCorrelationId(taskId), e.getClass().getSimpleName()); } String title = resolveTitleForEvent(); String messageIdText = StrUtil.isBlank(messageId) ? null : messageId; sender.sendEvent(SSEEventType.FINISH.value(), new CompletionPayload(messageIdText, title, sources, ChatMessage.MessageStatus.NORMAL)); sender.sendEvent(SSEEventType.DONE.value(), "[DONE]"); + log.info("SSE stream completed: traceId={}, taskId={}", + StreamTaskManager.safeCorrelationId(taskManager.traceId(taskId)), + StreamTaskManager.safeCorrelationId(taskId)); taskManager.unregister(taskId); sender.complete(); } @@ -225,6 +238,9 @@ public void onError(Throwable t) { if (taskManager.isCancelled(taskId)) { return; } + log.warn("SSE stream failed: traceId={}, taskId={}", + StreamTaskManager.safeCorrelationId(taskManager.traceId(taskId)), + StreamTaskManager.safeCorrelationId(taskId)); taskManager.unregister(taskId); sender.fail(t); } diff --git a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamTaskManager.java b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamTaskManager.java index 54dfa6a6a..ff0264c22 100644 --- a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamTaskManager.java +++ b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/handler/StreamTaskManager.java @@ -87,6 +87,23 @@ public void register(String taskId, SseEmitterSender sender, Supplier bucket = redissonClient.getBucket(cancelKey(taskId)); bucket.set(Boolean.TRUE, CANCEL_TTL); @@ -139,15 +157,21 @@ private void cancelLocal(String taskId) { return; } - if (taskInfo.handle != null) { - taskInfo.handle.cancel(); - } + try { + if (taskInfo.handle != null) { + taskInfo.handle.cancel(); + } - // 在取消时执行回调,保存已累积的内容 - if (taskInfo.sender != null) { - CompletionPayload payload = taskInfo.onCancelSupplier.get(); - sendCancelAndDone(taskInfo.sender, payload); - taskInfo.sender.complete(); + // 在取消时执行回调,保存已累积的内容 + if (taskInfo.sender != null) { + CompletionPayload payload = taskInfo.onCancelSupplier.get(); + sendCancelAndDone(taskInfo.sender, payload); + taskInfo.sender.complete(); + } + } finally { + notifyCancellationObserver(taskInfo); + log.info("SSE cancellation applied: traceId={}, taskId={}", + safeCorrelationId(taskInfo.traceId), safeCorrelationId(taskId)); } } @@ -169,6 +193,17 @@ private void sendCancelAndDone(SseEmitterSender sender, CompletionPayload payloa sender.sendEvent(SSEEventType.DONE.value(), "[DONE]"); } + private void notifyCancellationObserver(StreamTaskInfo taskInfo) { + Runnable observer = taskInfo.cancellationObserver; + if (observer != null && taskInfo.cancellationObserved.compareAndSet(false, true)) { + observer.run(); + } + } + + public static String safeCorrelationId(String value) { + return value != null && value.matches("[A-Za-z0-9_-]{1,64}") ? value : ""; + } + @SneakyThrows private StreamTaskInfo getOrCreate(String taskId) { return tasks.get(taskId, StreamTaskInfo::new); @@ -176,8 +211,11 @@ private StreamTaskInfo getOrCreate(String taskId) { private static final class StreamTaskInfo { private final AtomicBoolean cancelled = new AtomicBoolean(false); + private final AtomicBoolean cancellationObserved = new AtomicBoolean(false); private volatile StreamCancellationHandle handle; private volatile SseEmitterSender sender; private volatile Supplier onCancelSupplier; + private volatile Runnable cancellationObserver; + private volatile String traceId; } } diff --git a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/pipeline/StreamChatPipeline.java b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/pipeline/StreamChatPipeline.java index 974b6b649..83dd12746 100644 --- a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/pipeline/StreamChatPipeline.java +++ b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/pipeline/StreamChatPipeline.java @@ -22,6 +22,7 @@ import com.nageoffer.ai.ragent.framework.convention.ChatMessage; import com.nageoffer.ai.ragent.framework.convention.ChatRequest; import com.nageoffer.ai.ragent.framework.convention.SourceRef; +import com.nageoffer.ai.ragent.framework.trace.RagTraceContext; import com.nageoffer.ai.ragent.infra.chat.LLMService; import com.nageoffer.ai.ragent.infra.chat.StreamCallback; import com.nageoffer.ai.ragent.infra.chat.StreamCancellationHandle; @@ -80,23 +81,51 @@ public class StreamChatPipeline { * 执行流式对话管道 */ public void execute(StreamChatContext ctx) { - loadMemory(ctx); - rewriteQuery(ctx); - resolveIntents(ctx); + String stage = "memory"; + try { + logStage(stage, ctx); + loadMemory(ctx); + stage = "rewrite"; + logStage(stage, ctx); + rewriteQuery(ctx); + stage = "intent"; + logStage(stage, ctx); + resolveIntents(ctx); - if (handleGuidance(ctx)) { - return; - } - if (handleSystemOnly(ctx)) { - return; - } + stage = "guidance"; + logStage(stage, ctx); + if (handleGuidance(ctx)) { + return; + } + stage = "system-response"; + logStage(stage, ctx); + if (handleSystemOnly(ctx)) { + return; + } - RetrievalContext retrievalCtx = retrieve(ctx); - if (handleEmptyRetrieval(ctx, retrievalCtx)) { - return; + stage = "retrieval"; + logStage(stage, ctx); + RetrievalContext retrievalCtx = retrieve(ctx); + if (handleEmptyRetrieval(ctx, retrievalCtx)) { + return; + } + + stage = "llm-stream"; + logStage(stage, ctx); + streamRagResponse(ctx, retrievalCtx); + } catch (RuntimeException ex) { + log.warn("SSE pipeline failed: traceId={}, taskId={}, stage={}, errorType={}", + StreamTaskManager.safeCorrelationId(RagTraceContext.getTraceId()), + StreamTaskManager.safeCorrelationId(ctx.getTaskId()), stage, + ex.getClass().getSimpleName()); + throw ex; } + } - streamRagResponse(ctx, retrievalCtx); + private void logStage(String stage, StreamChatContext ctx) { + log.debug("SSE pipeline stage: traceId={}, taskId={}, stage={}", + StreamTaskManager.safeCorrelationId(RagTraceContext.getTraceId()), + StreamTaskManager.safeCorrelationId(ctx.getTaskId()), stage); } // ==================== 流水线阶段 ==================== diff --git a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/ratelimit/ChatQueueLimiter.java b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/ratelimit/ChatQueueLimiter.java index 46e079147..875b6dec5 100644 --- a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/ratelimit/ChatQueueLimiter.java +++ b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/service/ratelimit/ChatQueueLimiter.java @@ -145,7 +145,9 @@ private String buildFallbackTitle(String question) { private void sendRejectEvents(SseEmitter emitter, RejectedContext rejectedContext) { SseEmitterSender sender = new SseEmitterSender(emitter); if (rejectedContext != null) { - sender.sendEvent(SSEEventType.META.value(), new MetaPayload(rejectedContext.conversationId, rejectedContext.taskId)); + // 限流拒绝发生在 Trace 建立前,不能返回一个没有对应运行记录的伪 traceId。 + sender.sendEvent(SSEEventType.META.value(), + new MetaPayload(rejectedContext.conversationId, rejectedContext.taskId, null)); sender.sendEvent(SSEEventType.REJECT.value(), new MessageDelta(RESPONSE_TYPE, REJECT_MESSAGE)); sender.sendEvent(SSEEventType.FINISH.value(), new CompletionPayload(String.valueOf(rejectedContext.messageId), rejectedContext.title, diff --git a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunner.java b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunner.java index e18b18097..f57ec27ec 100644 --- a/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunner.java +++ b/bootstrap/src/main/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunner.java @@ -28,11 +28,13 @@ import com.nageoffer.ai.ragent.rag.dao.entity.RagTraceNodeDO; import com.nageoffer.ai.ragent.rag.dao.entity.RagTraceRunDO; import com.nageoffer.ai.ragent.rag.service.RagTraceRecordService; +import com.nageoffer.ai.ragent.rag.service.handler.StreamTaskManager; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import java.util.Date; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; /** @@ -48,11 +50,13 @@ public class StreamChatTraceRunner { private static final String STATUS_RUNNING = "RUNNING"; private static final String STATUS_SUCCESS = "SUCCESS"; private static final String STATUS_ERROR = "ERROR"; + private static final String STATUS_CANCELLED = "CANCELLED"; private static final String USER_TTFT_NODE_NAME = "user-first-packet"; private static final String USER_TTFT_NODE_TYPE = "USER_TTFT"; private final RagTraceProperties traceProperties; private final RagTraceRecordService traceRecordService; + private final StreamTaskManager taskManager; /** * @param businessLogic 接收 trace 增强后的 callback:onComplete / onError 会触发 finishRun @@ -70,6 +74,7 @@ public void run(String question, String traceId = IdUtil.getSnowflakeNextIdStr(); long startMillis = System.currentTimeMillis(); + AtomicBoolean finished = new AtomicBoolean(false); traceRecordService.startRun(RagTraceRunDO.builder() .traceId(traceId) .traceName(TRACE_NAME) @@ -94,16 +99,27 @@ protected void onFirstContent() { @Override protected void onFinish(boolean success, Throwable error) { - finishRun(traceId, success, error, startMillis); + finishRunOnce(finished, traceId, taskId, + success ? STATUS_SUCCESS : STATUS_ERROR, error, startMillis); } }; RagTraceContext.setTraceId(traceId); RagTraceContext.setTaskId(taskId); try { + taskManager.bindCancellationObserver(taskId, + () -> finishRunOnce(finished, traceId, taskId, STATUS_CANCELLED, null, startMillis)); + traceAwareCallback.onTraceStarted(traceId); + log.info("SSE stream started: traceId={}, taskId={}", + StreamTaskManager.safeCorrelationId(traceId), StreamTaskManager.safeCorrelationId(taskId)); + if (taskManager.isCancelled(taskId)) { + return; + } businessLogic.accept(traceAwareCallback); } catch (Throwable ex) { - log.warn("执行流式对话失败(同步阶段),会话ID:{},任务ID:{}", conversationId, taskId, ex); + log.warn("SSE stream failed synchronously: traceId={}, taskId={}, errorType={}", + StreamTaskManager.safeCorrelationId(traceId), StreamTaskManager.safeCorrelationId(taskId), + ex.getClass().getSimpleName()); // 走 traceAwareCallback.onError 以复用其内部 CAS,避免与 pipeline 内已触发的终态重复收尾 try { traceAwareCallback.onError(ex); @@ -137,21 +153,34 @@ private void recordUserTtft(String traceId, Date runStartTime, long startMillis) .build()); traceRecordService.finishNode(traceId, nodeId, STATUS_SUCCESS, null, new Date(now), durationMs); } catch (Exception e) { - log.warn("写入 user-first-packet 节点失败,traceId:{}", traceId, e); + log.warn("Failed to record SSE first packet: traceId={}, errorType={}", + StreamTaskManager.safeCorrelationId(traceId), e.getClass().getSimpleName()); } } - private void finishRun(String traceId, boolean success, Throwable error, long startMillis) { + private void finishRunOnce(AtomicBoolean finished, + String traceId, + String taskId, + String status, + Throwable error, + long startMillis) { + if (!finished.compareAndSet(false, true)) { + return; + } try { traceRecordService.finishRun( traceId, - success ? STATUS_SUCCESS : STATUS_ERROR, - success ? null : truncateError(error), + status, + STATUS_ERROR.equals(status) ? truncateError(error) : null, new Date(), System.currentTimeMillis() - startMillis ); + log.info("SSE stream finished: traceId={}, taskId={}, status={}", + StreamTaskManager.safeCorrelationId(traceId), + StreamTaskManager.safeCorrelationId(taskId), status); } catch (Exception e) { - log.warn("finishRun 失败,traceId:{}", traceId, e); + log.warn("Failed to finish SSE trace: traceId={}, errorType={}", + StreamTaskManager.safeCorrelationId(traceId), e.getClass().getSimpleName()); } } @@ -162,7 +191,8 @@ private void runWithoutTrace(String conversationId, try { businessLogic.accept(callback); } catch (Throwable ex) { - log.warn("执行流式对话失败,会话ID:{},任务ID:{}", conversationId, taskId, ex); + log.warn("SSE stream failed without trace: taskId={}, errorType={}", + StreamTaskManager.safeCorrelationId(taskId), ex.getClass().getSimpleName()); callback.onError(ex); } } diff --git a/bootstrap/src/test/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunnerTest.java b/bootstrap/src/test/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunnerTest.java new file mode 100644 index 000000000..9b7c2e73a --- /dev/null +++ b/bootstrap/src/test/java/com/nageoffer/ai/ragent/rag/trace/StreamChatTraceRunnerTest.java @@ -0,0 +1,122 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.nageoffer.ai.ragent.rag.trace; + +import com.nageoffer.ai.ragent.framework.trace.RagTraceContext; +import com.nageoffer.ai.ragent.infra.chat.StreamCallback; +import com.nageoffer.ai.ragent.rag.config.RagTraceProperties; +import com.nageoffer.ai.ragent.rag.dao.entity.RagTraceRunDO; +import com.nageoffer.ai.ragent.rag.service.RagTraceRecordService; +import com.nageoffer.ai.ragent.rag.service.handler.StreamTaskManager; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import java.util.Date; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class StreamChatTraceRunnerTest { + + @Mock + private RagTraceProperties traceProperties; + + @Mock + private RagTraceRecordService traceRecordService; + + @Mock + private StreamTaskManager taskManager; + + @Mock + private StreamCallback callback; + + @Captor + private ArgumentCaptor runCaptor; + + @Captor + private ArgumentCaptor cancellationObserverCaptor; + + @InjectMocks + private StreamChatTraceRunner traceRunner; + + @AfterEach + void clearTraceContext() { + RagTraceContext.clear(); + } + + @Test + void shouldExposeSameTraceIdAndFinishCancelledRunExactlyOnce() { + when(traceProperties.isEnabled()).thenReturn(true); + AtomicReference traceAware = new AtomicReference<>(); + + traceRunner.run("question", "conversation-1", "task-1", callback, traceAware::set); + + verify(traceRecordService).startRun(runCaptor.capture()); + String traceId = runCaptor.getValue().getTraceId(); + verify(taskManager).bindCancellationObserver(eq("task-1"), cancellationObserverCaptor.capture()); + verify(callback).onTraceStarted(traceId); + assertEquals("task-1", runCaptor.getValue().getTaskId()); + assertNull(RagTraceContext.getTraceId()); + assertNull(RagTraceContext.getTaskId()); + + cancellationObserverCaptor.getValue().run(); + traceAware.get().onComplete(); + + verify(traceRecordService, times(1)) + .finishRun(eq(traceId), eq("CANCELLED"), isNull(), any(Date.class), anyLong()); + verify(callback).onComplete(); + } + + @Test + void shouldPropagateTraceContextDuringBusinessLogicAndClearAfterwards() { + when(traceProperties.isEnabled()).thenReturn(true); + AtomicReference observedTraceId = new AtomicReference<>(); + AtomicReference observedTaskId = new AtomicReference<>(); + AtomicReference observedCallback = new AtomicReference<>(); + + traceRunner.run("question", "conversation-1", "task-1", callback, enhanced -> { + observedTraceId.set(RagTraceContext.getTraceId()); + observedTaskId.set(RagTraceContext.getTaskId()); + observedCallback.set(enhanced); + }); + + verify(traceRecordService).startRun(runCaptor.capture()); + assertEquals(runCaptor.getValue().getTraceId(), observedTraceId.get()); + assertEquals("task-1", observedTaskId.get()); + observedCallback.get().onComplete(); + verify(callback).onComplete(); + assertNull(RagTraceContext.getTraceId()); + assertNull(RagTraceContext.getTaskId()); + verify(taskManager).bindCancellationObserver(eq("task-1"), any(Runnable.class)); + } +} diff --git a/docs/external-governance-operations.md b/docs/external-governance-operations.md index 00970d7b6..a3726423d 100644 --- a/docs/external-governance-operations.md +++ b/docs/external-governance-operations.md @@ -90,6 +90,8 @@ 仓库内已修复当前 5 个 Critical SSRF、1 个 High 鉴权绕过和 1 个 High 前端不完整转义,并增加定向回归测试;前端本地 npm audit 已无 Critical/High。PR #9 的 CodeQL 重扫已通过且无本 PR 新增告警;默认分支存量 CodeQL/Dependabot 告警仍需在合并后按最新 API 快照逐条 triage,不要为清零数字批量 dismiss。 +2026-08-13 再次分页读取默认分支得到 99 个 open(2 Critical / 65 Medium / 32 未分级)。SSE 诊断后续分支同时消除其所触及文件中的 3 个存量 `java/log-injection` 数据流;是否正式关闭以该 PR CodeQL 扫描和合并后的默认分支快照为准。 + 完成证据:修复 PR、重新扫描结果、关闭或带理由 dismiss 的告警记录。 ## 6. 真实集成测试验收 @@ -139,3 +141,17 @@ | `integration-isolation-gap` | 隔离环境连续两次真实集成测试通过,且清理 postcondition 有证据 | | `dependency-supply-chain-gap` | Dependency Review required;Code Scanning merge protection 生效;Critical/High 存量完成修复或有依据的判定 | | `production-default-credential-boundary` | 代码层已闭合;staging 正/负/混合 profile 验收作为发布接受证据 | + +## 10. Better Harness 复评后的个人待办(2026-08-13) + +以下事项不能由无真实凭据、基础设施或仓库管理员最终判断的代理代做: + +1. 在隔离环境连续运行两次 protected integration workflow,保存 doctor、测试、cleanup/reset postcondition 和脱敏日志。 +2. 在 staging 执行生产凭据守卫的负向、正向、混合 profile 验收。 +3. 轮换曾暴露于本机配置的 API key,并在验证新 key 后吊销旧 key。 +4. 备份后执行 v1.1.0 数据库升级,保存升级和回滚验证证据。 +5. PR 合并后复核默认分支的 CodeQL/Dependabot 最新分页快照;逐条修复或给出技术判定,不批量 dismiss。 +6. 若 GitHub 套餐支持,在 `main-protection` 增加 Code Scanning 结果严重度规则;否则继续保留双语言 required jobs 和人工告警验收。 +7. 在未来至少两个真实开发任务中保存“目标 → 改动 → 最小检查 → 完整检查 → PR 接受结果”,再做纵向 Better Harness 复评。 + +仓库内已经能先完成的部分是 SSE 关联诊断和 affected-check 路由;这不替代第 1、5、7 项的真实运行与长期证据。 diff --git a/docs/harness-remediation-tracker.md b/docs/harness-remediation-tracker.md index 744218f7a..dd351fa28 100644 --- a/docs/harness-remediation-tracker.md +++ b/docs/harness-remediation-tracker.md @@ -58,7 +58,7 @@ - [ ] **轮换 mygpt API key**:曾明文存于 opencode 配置(已迁移 auth.json;fork/上游全历史扫描零泄露),供应商侧轮换一次收尾 - [ ] **执行 v1.1.0 SQL 升级**:本地与部署库均需执行(先备份,按 `docs/v1.1.0-upgrade-guide.md`) - [ ] **消化 Dependabot 漏洞告警**:2026-08-11 API 快照为 58 个 open(1 critical / 23 high / 32 medium / 2 low);按 critical→high 优先处理,数量变化时以新 API 快照为准 -- [ ] **消化 CodeQL 告警**:2026-08-11 API 快照为 104 个 open(security severity:5 critical / 2 high / 65 medium / 32 未分级),先完成 triage、去重与误报处置,再确定阻断阈值 +- [ ] **消化 CodeQL 告警**:2026-08-13 分页 API 快照为 99 个 open(2 critical / 65 medium / 32 未分级);SSE 后续分支已处理其触及文件中的 3 个 log-injection 数据流,最终关闭以 PR 扫描为准;其余继续按严重度 triage - [x] **收紧 main bypass**:个人 bypass 已移除;单人维护模式下不强制 approval - [ ] **定期上游同步**(建议月例行):按 `docs/upstream-sync.md` SOP @@ -73,8 +73,10 @@ ### 4.3 可观测性阶段(31–60 天,含 P3) - [x] SSE 诊断脚本增强:捕获 stderr/headers/raw SSE、correlation ID、curl 与业务完成断言;真实端到端运行待外部验证 +- [x] SSE 仓库内关联诊断:正常流将同一 `traceId`/`taskId` 贯穿 Trace 记录、SSE META、阶段日志和取消终态;取消以 `CANCELLED` 一次性收尾;真实端到端运行仍待外部验证 +- [x] affected-check 路由:`docs/verification-routing.md` 提供后端模块/类与前端单文件快速检查,并明确升级到完整或真实集成验证的条件 - [x] coverage ratchet:后端 JaCoCo 报告 + 前端 Vitest V8 低基线硬阈值,见 `docs/coverage-baseline.md` -- [ ] 可观测性:结构化日志、metrics(模型首包/检索通道/降级次数/SSE 断开原因)、OpenTelemetry tracing、health/liveness、Trace ID 跨前后端 +- [ ] 可观测性后续:metrics(模型首包/检索通道/降级次数/SSE 断开原因)、OpenTelemetry、health/liveness;基础 Trace ID 跨前后端已实现 - [ ] release 流程:SemVer/changelog/GitHub Release/制品 checksum - [ ] Docker/部署 ADR:环境分层、DB migration、备份恢复、RTO/RPO - [x] Actions 升级:checkout v7、setup-java v5、CodeQL Action v4(全部固定完整 SHA) diff --git a/docs/verification-routing.md b/docs/verification-routing.md new file mode 100644 index 000000000..f7e50201e --- /dev/null +++ b/docs/verification-routing.md @@ -0,0 +1,24 @@ +# 改动到验证的精确路由 + +本文给代理和维护者提供局部改动的最小反馈入口。局部检查用于尽早发现问题,不替代提交前或 PR 上的完整 CI。 + +## 路由表 + +| 改动范围 | 首个最小检查 | 提交前升级 | +|---|---|---| +| `framework/` 通用上下文、Web/SSE 工具 | `./mvnw -B -ntp -pl framework test` | 根目录 `./mvnw -B -ntp verify` | +| `infra-ai/` 模型回调、路由、容错 | `./mvnw -B -ntp -pl infra-ai -am test` | 根目录 `./mvnw -B -ntp verify` | +| SSE Trace、pipeline、取消 | `./mvnw -B -ntp -pl bootstrap -am -Dtest=StreamChatTraceRunnerTest,StreamChatPipelineTest -Dsurefire.failIfNoSpecifiedTests=false test` | 根目录 `./mvnw -B -ntp verify`;真实依赖行为另走 protected integration workflow | +| 生产凭据守卫 | `./mvnw -B -ntp -pl bootstrap -am -Dtest=ProductionCredentialGuardTest -Dsurefire.failIfNoSpecifiedTests=false test` | 根目录 `./mvnw -B -ntp verify`;发布前需 staging 验收 | +| 前端单个 hook/component | `cd frontend && npm run test -- ` | `npm run lint && npm run test:coverage && npm run build` | +| Maven/npm 依赖或 workflow | 受影响模块构建、本地 audit/语法检查 | PR 上等待 Dependency Review、CodeQL 和全部 required checks | +| `@Tag("integration")` 测试或外部适配器 | 不在无真实依赖的本机伪运行 | 手动 protected integration workflow,保存清理 postcondition | + +## 升级规则 + +- 跨模块接口、根 `pom.xml`、共享类型或公共回调发生变化:直接升级到根目录 `verify`。 +- 前后端协议字段发生变化:后端定向测试和前端对应解析测试都要运行,再跑两端完整门禁。 +- 触及数据库、Redis、Milvus、对象存储或模型 API 的真实交互:单元测试通过仍只能说明代码边界;必须在隔离环境跑 opt-in integration。 +- 触及认证、授权、URL 获取、反序列化、日志或凭据:除功能测试外,必须等待 PR CodeQL/Dependency Review;扫描 job 成功不等于存量告警已处置。 + +PowerShell 调用 Maven 时,可将 `-Dtest=...` 和 `-Dsurefire.failIfNoSpecifiedTests=false` 分别放在双引号内,避免参数被 shell 误解析。 diff --git a/frontend/src/hooks/__tests__/useStreamResponse.test.ts b/frontend/src/hooks/__tests__/useStreamResponse.test.ts index 4e62d3bcc..17788345d 100644 --- a/frontend/src/hooks/__tests__/useStreamResponse.test.ts +++ b/frontend/src/hooks/__tests__/useStreamResponse.test.ts @@ -30,7 +30,7 @@ describe("createStreamResponse", () => { const fetchMock = vi.fn().mockResolvedValue( okResponse( sseBody( - ["meta", '{"conversationId":"conv-1","taskId":"task-1"}'], + ["meta", '{"conversationId":"conv-1","taskId":"task-1","traceId":"trace-1"}'], ["message", '{"type":"response","delta":"你好"}'], ["finish", '{"messageId":"m-1","title":"会话标题"}'], ["done", "[DONE]"] @@ -68,7 +68,9 @@ describe("createStreamResponse", () => { ).start(); expect(events).toEqual(["meta", "message", "finish", "done"]); - expect(metaPayloads).toEqual([{ conversationId: "conv-1", taskId: "task-1" }]); + expect(metaPayloads).toEqual([ + { conversationId: "conv-1", taskId: "task-1", traceId: "trace-1" } + ]); expect(messagePayloads).toEqual([{ type: "response", delta: "你好" }]); expect(finishPayloads).toEqual([{ messageId: "m-1", title: "会话标题" }]); expect(done).toHaveBeenCalledTimes(1); diff --git a/frontend/src/stores/chatStore.ts b/frontend/src/stores/chatStore.ts index beb00705b..2dfc68052 100644 --- a/frontend/src/stores/chatStore.ts +++ b/frontend/src/stores/chatStore.ts @@ -6,7 +6,8 @@ import type { FeedbackValue, Message, MessageDeltaPayload, - Session + Session, + StreamMetaPayload } from "@/types"; import { listMessages, @@ -306,7 +307,7 @@ export const useChatStore = create((set, get) => ({ const token = storage.getToken(); const handlers = { - onMeta: (payload: { conversationId: string; taskId: string }) => { + onMeta: (payload: StreamMetaPayload) => { if (get().streamingMessageId !== assistantId) return; const nextId = payload.conversationId || get().currentSessionId; if (!nextId) return; diff --git a/frontend/src/types/index.ts b/frontend/src/types/index.ts index 0209b6be6..6613a190a 100644 --- a/frontend/src/types/index.ts +++ b/frontend/src/types/index.ts @@ -60,6 +60,7 @@ export interface RecommendedQuestionsPayload { export interface StreamMetaPayload { conversationId: string; taskId: string; + traceId?: string | null; } export interface MessageDeltaPayload { diff --git a/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/ForwardingStreamCallback.java b/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/ForwardingStreamCallback.java index bb986dc0c..eed734b77 100644 --- a/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/ForwardingStreamCallback.java +++ b/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/ForwardingStreamCallback.java @@ -39,6 +39,11 @@ protected ForwardingStreamCallback(StreamCallback delegate) { this.delegate = delegate; } + @Override + public final void onTraceStarted(String traceId) { + delegate.onTraceStarted(traceId); + } + @Override public final void onContent(String content) { if (firstContentSeen.compareAndSet(false, true)) { diff --git a/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/StreamCallback.java b/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/StreamCallback.java index 820cce173..9caec29fd 100644 --- a/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/StreamCallback.java +++ b/infra-ai/src/main/java/com/nageoffer/ai/ragent/infra/chat/StreamCallback.java @@ -42,6 +42,14 @@ */ public interface StreamCallback { + /** + * Trace 已建立、流式业务即将开始。 + * + * @param traceId 启用 Trace 时的关联 ID;未启用时为 {@code null} + */ + default void onTraceStarted(String traceId) { + } + /** * 记录当前回答对应的用户消息 ID *