|
|
@ -38,7 +38,9 @@ import java.util.Collections; |
|
|
import java.util.LinkedHashMap; |
|
|
import java.util.LinkedHashMap; |
|
|
import java.util.List; |
|
|
import java.util.List; |
|
|
import java.util.Map; |
|
|
import java.util.Map; |
|
|
|
|
|
import java.util.UUID; |
|
|
import java.util.concurrent.CopyOnWriteArrayList; |
|
|
import java.util.concurrent.CopyOnWriteArrayList; |
|
|
|
|
|
import java.util.concurrent.atomic.AtomicBoolean; |
|
|
import java.util.concurrent.atomic.AtomicInteger; |
|
|
import java.util.concurrent.atomic.AtomicInteger; |
|
|
import java.util.concurrent.atomic.AtomicReference; |
|
|
import java.util.concurrent.atomic.AtomicReference; |
|
|
|
|
|
|
|
|
@ -389,6 +391,169 @@ public class AssistantApp { |
|
|
.onErrorResume(e -> Flux.just("抱歉,AI 服务调用失败:" + e.getMessage())); |
|
|
.onErrorResume(e -> Flux.just("抱歉,AI 服务调用失败:" + e.getMessage())); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
/** |
|
|
|
|
|
* 流式对话(OpenAI Chat Completions 标准 SSE 格式)。 |
|
|
|
|
|
* <p> |
|
|
|
|
|
* 复用 {@link #chatStream(ChatContext)} 的完整编排逻辑(熔断早退 / FAQ 命中早退 / |
|
|
|
|
|
* 正常流式调用 / 空白缓冲 / 埋点),差异在于把每个文本片段包装为 OpenAI 标准 JSON chunk: |
|
|
|
|
|
* 首片 delta 携带 role=assistant,流结束时追加 finish_reason=stop 的 chunk 与 [DONE]。 |
|
|
|
|
|
* <p> |
|
|
|
|
|
* 每个 Flux 元素即一个完整 JSON 字符串,Spring WebFlux 自动加 data: 前缀。 |
|
|
|
|
|
* |
|
|
|
|
|
* @param ctx 对话上下文 |
|
|
|
|
|
* @return OpenAI 标准格式的流式回答 |
|
|
|
|
|
*/ |
|
|
|
|
|
public Flux<String> 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<ToolCallEvent> 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<Usage> usageRef = new AtomicReference<>(); |
|
|
|
|
|
AtomicReference<String> errorTypeRef = new AtomicReference<>(); |
|
|
|
|
|
AtomicReference<String> errorMessageRef = new AtomicReference<>(); |
|
|
|
|
|
Flux<ChatResponse> responseFlux = spec.stream().chatResponse(); |
|
|
|
|
|
Flux<String> 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 字符串)。 |
|
|
|
|
|
* <p> |
|
|
|
|
|
* 字段顺序固定为 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<String, Object> delta = new LinkedHashMap<>(); |
|
|
|
|
|
if (first) { |
|
|
|
|
|
delta.put("role", "assistant"); |
|
|
|
|
|
} |
|
|
|
|
|
if (content != null && !content.isEmpty()) { |
|
|
|
|
|
delta.put("content", content); |
|
|
|
|
|
} |
|
|
|
|
|
Map<String, Object> choice = new LinkedHashMap<>(); |
|
|
|
|
|
choice.put("index", 0); |
|
|
|
|
|
choice.put("delta", delta); |
|
|
|
|
|
choice.put("finish_reason", finishReason); |
|
|
|
|
|
Map<String, Object> 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<String> 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 合并后再发出。 |
|
|
* 缓冲以空白字符结尾的 chunk,将其与下一个 chunk 合并后再发出。 |
|
|
* <p> |
|
|
* <p> |
|
|
|