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); } /**