From 1cd8ed1248bf3387d930a4a2bc4ef9ca384562e3 Mon Sep 17 00:00:00 2001
From: wanghanlin <1533525126@qq.com>
Date: Tue, 18 Aug 2026 09:56:12 +0800
Subject: [PATCH] =?UTF-8?q?feat(chat):=20=E6=B5=81=E5=BC=8F=E5=AF=B9?=
=?UTF-8?q?=E8=AF=9D=E8=BE=93=E5=87=BA=E5=AF=B9=E9=BD=90=20OpenAI=20Chat?=
=?UTF-8?q?=20Completions=20=E6=A0=87=E5=87=86=20SSE=20=E6=A0=BC=E5=BC=8F?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
- 新增 AssistantApp.chatStreamOpenAi,复用 chatStream 编排逻辑,将流式文本分片包装为 data: {"choices":[{"delta":{"content":"..."}}]} 标准 JSON chunk
- 首片 delta 携带 role=assistant,流结束时追加 finish_reason=stop chunk 与 data: [DONE];熔断降级、FAQ 命中、错误兜底统一走 openAiFallbackStream
- AiController、OpenApiController 的流式接口改调 chatStreamOpenAi,旧路径保留纯文本输出向后兼容
---
.../com/wok/supportbot/app/AssistantApp.java | 165 ++++++++++++++++++
.../supportbot/controller/AiController.java | 2 +-
.../controller/OpenApiController.java | 2 +-
3 files changed, 167 insertions(+), 2 deletions(-)
diff --git a/src/main/java/com/wok/supportbot/app/AssistantApp.java b/src/main/java/com/wok/supportbot/app/AssistantApp.java
index 2351775..375999e 100644
--- a/src/main/java/com/wok/supportbot/app/AssistantApp.java
+++ b/src/main/java/com/wok/supportbot/app/AssistantApp.java
@@ -38,7 +38,9 @@ import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.UUID;
import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
@@ -389,6 +391,169 @@ public class AssistantApp {
.onErrorResume(e -> Flux.just("抱歉,AI 服务调用失败:" + e.getMessage()));
}
+ /**
+ * 流式对话(OpenAI Chat Completions 标准 SSE 格式)。
+ *
+ * 复用 {@link #chatStream(ChatContext)} 的完整编排逻辑(熔断早退 / FAQ 命中早退 /
+ * 正常流式调用 / 空白缓冲 / 埋点),差异在于把每个文本片段包装为 OpenAI 标准 JSON chunk:
+ * 首片 delta 携带 role=assistant,流结束时追加 finish_reason=stop 的 chunk 与 [DONE]。
+ *
+ * 每个 Flux 元素即一个完整 JSON 字符串,Spring WebFlux 自动加 data: 前缀。
+ *
+ * @param ctx 对话上下文
+ * @return OpenAI 标准格式的流式回答
+ */
+ public Flux chatStreamOpenAi(ChatContext ctx) {
+ long startNanos = System.nanoTime();
+ // OpenAI 标准 chunk 的公共元信息:同一次流式回答共享 id / created / model
+ String completionId = "chatcmpl-" + UUID.randomUUID().toString().replace("-", "");
+ long created = System.currentTimeMillis() / 1000;
+ // 活跃模型配置可能不存在(返回 null),回退 unknown
+ AiModelConfig cfg = null;
+ try {
+ cfg = aiModelConfigService.getActiveConfigWithFullKey(ctx.appType());
+ } catch (Exception e) {
+ log.warn("获取活跃模型配置失败,model 回退 unknown: chatId={}, error={}", ctx.chatId(), e.getMessage());
+ }
+ String model = (cfg != null && cfg.getModelName() != null) ? cfg.getModelName() : "unknown";
+ // 熔断:全局 AI 调用处于熔断状态(不做 buildRequest,避免熔断期间仍走意图路由/检索)
+ if (aiCircuitBreaker.isOpen(AI_CIRCUIT_KEY)) {
+ log.warn("AI 调用熔断中(OpenAI 流式),返回降级提示");
+ recordTrace(ctx, null, CIRCUIT_OPEN_MESSAGE, 0, "BYPASS",
+ new TraceMeta("CIRCUIT_BREAK", "AI 服务熔断降级", null, null, null, null));
+ return openAiFallbackStream(completionId, model, created, CIRCUIT_OPEN_MESSAGE, true);
+ }
+ ChatRequest req = chatPipeline.buildRequest(ctx);
+ if (req.faqHit()) {
+ // FAQ 命中:整段答案包装为 OpenAI chunk,随后追加 stop + [DONE]
+ String faqAnswer = req.faqAnswer().get();
+ recordTrace(ctx, req, faqAnswer, 0, "FAQ",
+ new TraceMeta(null, null, null, null, null, null));
+ return openAiFallbackStream(completionId, model, created, faqAnswer, true);
+ }
+ // 显式事件收集器 + 轮次计数器,通过 toolContext 传给 McpToolCallback,规避 Reactor 跨线程丢 ThreadLocal 的问题
+ List events = new CopyOnWriteArrayList<>();
+ AtomicInteger rounds = new AtomicInteger(0);
+ ChatClient.ChatClientRequestSpec spec = getChatClient(ctx.appType(), ctx.allowedMcpTools())
+ .prompt()
+ .user(req.finalMessage())
+ .advisors(s -> s.param(CONVERSATION_ID, ctx.chatId()));
+ if (StringUtils.hasText(req.finalSystemPrompt())) {
+ spec = spec.system(req.finalSystemPrompt());
+ }
+ spec = spec.toolContext(Map.of(
+ McpToolCallback.MCP_EVENTS_KEY, events,
+ McpToolCallback.MCP_ROUNDS_KEY, rounds));
+ // 改为 chatResponse 流以采集 token 用量,再映射回纯文本流
+ AtomicReference usageRef = new AtomicReference<>();
+ AtomicReference errorTypeRef = new AtomicReference<>();
+ AtomicReference errorMessageRef = new AtomicReference<>();
+ Flux responseFlux = spec.stream().chatResponse();
+ Flux rawStream = responseFlux
+ .doOnNext(r -> {
+ if (r != null && r.getMetadata() != null && r.getMetadata().getUsage() != null) {
+ usageRef.set(r.getMetadata().getUsage());
+ }
+ })
+ .map(r -> {
+ String out = r != null && r.getResult() != null && r.getResult().getOutput() != null
+ ? r.getResult().getOutput().getText() : "";
+ return out != null ? out : "";
+ });
+ // 聚合所有分片用于埋点(在 doFinally 时取完整回复文本)
+ StringBuilder aggregated = new StringBuilder();
+ // 首片标记:第一片 delta 需带 role=assistant,后续片仅含 content
+ AtomicBoolean first = new AtomicBoolean(true);
+ return preserveTrailingWhitespace(rawStream)
+ .doOnNext(aggregated::append)
+ .map(chunk -> buildOpenAiChunk(completionId, model, created, chunk, first.getAndSet(false), null))
+ .doOnComplete(() -> aiCircuitBreaker.recordSuccess(AI_CIRCUIT_KEY))
+ .doOnError(e -> {
+ aiCircuitBreaker.recordFailure(AI_CIRCUIT_KEY);
+ errorTypeRef.set(classifyError(e));
+ errorMessageRef.set(maskError(e.getMessage()));
+ log.error("AI 流式调用失败(OpenAI): chatId={}, error={}", ctx.chatId(), e.getMessage());
+ })
+ .doFinally(signalType -> {
+ // 流式埋点:按终止信号区分状态,断连/异常也落库(events 由 toolContext 显式收集,跨线程安全)
+ String status = signalType == SignalType.ON_COMPLETE ? "COMPLETE"
+ : signalType == SignalType.ON_ERROR ? "ERROR" : "CANCEL";
+ Usage usage = usageRef.get();
+ recordTrace(ctx, req, aggregated.toString(), elapsedMillis(startNanos), status,
+ new TraceMeta(errorTypeRef.get(), errorMessageRef.get(),
+ usage != null ? usage.getPromptTokens() : null,
+ usage != null ? usage.getCompletionTokens() : null,
+ usage != null ? usage.getTotalTokens() : null,
+ events));
+ })
+ // 流正常结束时追加 finish_reason=stop 的 chunk 与 [DONE]
+ .concatWith(Flux.just(
+ buildOpenAiChunk(completionId, model, created, "", false, "stop"),
+ "[DONE]"))
+ // 错误兜底:脱敏错误信息,避免泄露内部细节;role 是否出现取决于此前是否已发出过内容片
+ .onErrorResume(e -> openAiFallbackStream(completionId, model, created,
+ "抱歉,AI 服务调用失败:" + maskError(e.getMessage()), first.get()));
+ }
+
+ /**
+ * 组装单个 OpenAI Chat Completions 流式 chunk(JSON 字符串)。
+ *
+ * 字段顺序固定为 id/object/created/model/choices,delta 内 role 在 content 前;
+ * choices 使用 LinkedHashMap 以支持 finish_reason=null(Map.of 不允许 null 值)。
+ *
+ * @param id chunk 唯一 ID(chatcmpl-xxx)
+ * @param model 模型名称
+ * @param created 创建时间(epoch 秒)
+ * @param content 文本片段(stop 片传空串)
+ * @param first 是否首片(首片 delta 携带 role=assistant)
+ * @param finishReason 结束原因(中间片为 null,stop 片为 "stop")
+ * @return OpenAI 标准 chunk 的 JSON 字符串
+ */
+ private String buildOpenAiChunk(String id, String model, long created, String content, boolean first, String finishReason) {
+ // delta:首片带 role=assistant(role 在 content 前),后续片仅 content,stop 片为空对象
+ Map delta = new LinkedHashMap<>();
+ if (first) {
+ delta.put("role", "assistant");
+ }
+ if (content != null && !content.isEmpty()) {
+ delta.put("content", content);
+ }
+ Map choice = new LinkedHashMap<>();
+ choice.put("index", 0);
+ choice.put("delta", delta);
+ choice.put("finish_reason", finishReason);
+ Map chunk = new LinkedHashMap<>();
+ chunk.put("id", id);
+ chunk.put("object", "chat.completion.chunk");
+ chunk.put("created", created);
+ chunk.put("model", model);
+ chunk.put("choices", List.of(choice));
+ try {
+ return OBJECT_MAPPER.writeValueAsString(chunk);
+ } catch (Exception e) {
+ log.warn("序列化 OpenAI chunk 失败: {}", e.getMessage());
+ return "{}";
+ }
+ }
+
+ /**
+ * 组装 OpenAI 格式的早退/兜底流:内容 chunk + finish_reason=stop + [DONE]。
+ * 用于熔断降级、FAQ 命中与错误兜底三种场景。
+ *
+ * @param id chunk 唯一 ID
+ * @param model 模型名称
+ * @param created 创建时间(epoch 秒)
+ * @param content 完整回复文本
+ * @param withRole 首片是否携带 role=assistant(熔断/FAQ 早退为 true;错误兜底时取决于此前是否已发出内容片)
+ * @return OpenAI 标准格式的流
+ */
+ private Flux openAiFallbackStream(String id, String model, long created, String content, boolean withRole) {
+ return Flux.just(
+ buildOpenAiChunk(id, model, created, content, withRole, null),
+ buildOpenAiChunk(id, model, created, "", false, "stop"),
+ "[DONE]");
+ }
+
/**
* 缓冲以空白字符结尾的 chunk,将其与下一个 chunk 合并后再发出。
*
diff --git a/src/main/java/com/wok/supportbot/controller/AiController.java b/src/main/java/com/wok/supportbot/controller/AiController.java
index d9ac504..2530e7a 100644
--- a/src/main/java/com/wok/supportbot/controller/AiController.java
+++ b/src/main/java/com/wok/supportbot/controller/AiController.java
@@ -357,7 +357,7 @@ public class AiController {
ctx = new ChatContext(ctx.message(), ctx.chatId(), ctx.appType(), ctx.systemPrompt(),
ctx.allowedMcpTools(), ctx.categoryIds(), ctx.rewriteStrategy(), ctx.enableRag(), true,
ctx.roleId(), ctx.roleName(), ctx.accountId(), ctx.apiKeyId());
- return assistantApp.chatStream(ctx);
+ return assistantApp.chatStreamOpenAi(ctx);
}
/**
diff --git a/src/main/java/com/wok/supportbot/controller/OpenApiController.java b/src/main/java/com/wok/supportbot/controller/OpenApiController.java
index 1326e31..7572618 100644
--- a/src/main/java/com/wok/supportbot/controller/OpenApiController.java
+++ b/src/main/java/com/wok/supportbot/controller/OpenApiController.java
@@ -121,7 +121,7 @@ public class OpenApiController {
ChatContext ctx = buildOpenApiChatContext(message, resolvedChatId, apiKey, roleId,
categoryIds, rewriteStrategy, enableRag, true);
- return assistantApp.chatStream(ctx);
+ return assistantApp.chatStreamOpenAi(ctx);
}
/**