|
|
@ -44,7 +44,6 @@ import java.util.List; |
|
|
import java.util.Map; |
|
|
import java.util.Map; |
|
|
import java.util.UUID; |
|
|
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; |
|
|
|
|
|
|
|
|
@ -519,11 +518,9 @@ public class AssistantApp { |
|
|
}); |
|
|
}); |
|
|
// 聚合所有分片用于埋点(在 doFinally 时取完整回复文本) |
|
|
// 聚合所有分片用于埋点(在 doFinally 时取完整回复文本) |
|
|
StringBuilder aggregated = new StringBuilder(); |
|
|
StringBuilder aggregated = new StringBuilder(); |
|
|
// 首片标记:第一片 delta 需带 role=assistant,后续片仅含 content |
|
|
|
|
|
AtomicBoolean first = new AtomicBoolean(true); |
|
|
|
|
|
return preserveTrailingWhitespace(rawStream) |
|
|
return preserveTrailingWhitespace(rawStream) |
|
|
.doOnNext(aggregated::append) |
|
|
.doOnNext(aggregated::append) |
|
|
.map(chunk -> buildOpenAiChunk(completionId, model, created, chunk, first.getAndSet(false), null)) |
|
|
|
|
|
|
|
|
.map(chunk -> buildOpenAiChunk(completionId, model, created, chunk, false, null)) |
|
|
.doOnComplete(() -> aiCircuitBreaker.recordSuccess(AI_CIRCUIT_KEY)) |
|
|
.doOnComplete(() -> aiCircuitBreaker.recordSuccess(AI_CIRCUIT_KEY)) |
|
|
.doOnError(e -> { |
|
|
.doOnError(e -> { |
|
|
aiCircuitBreaker.recordFailure(AI_CIRCUIT_KEY); |
|
|
aiCircuitBreaker.recordFailure(AI_CIRCUIT_KEY); |
|
|
@ -543,13 +540,17 @@ public class AssistantApp { |
|
|
usage != null ? usage.getTotalTokens() : null, |
|
|
usage != null ? usage.getTotalTokens() : null, |
|
|
events)); |
|
|
events)); |
|
|
}) |
|
|
}) |
|
|
|
|
|
// 首片(仅 role=assistant、无 content)在流订阅时立即发出,确保 SSE 响应头/首字节及时 flush。 |
|
|
|
|
|
// 推理模型(如 doubao-seed)思考阶段 delta.content 为空、被 preserveTrailingWhitespace 吞掉, |
|
|
|
|
|
// 若不提前发首片,思考阶段将无任何字节输出,前端等待首字节会触发 60s 超时。 |
|
|
|
|
|
.startWith(buildOpenAiChunk(completionId, model, created, "", true, null)) |
|
|
// 流正常结束时追加 finish_reason=stop 的 chunk 与 [DONE] |
|
|
// 流正常结束时追加 finish_reason=stop 的 chunk 与 [DONE] |
|
|
.concatWith(Flux.just( |
|
|
.concatWith(Flux.just( |
|
|
buildOpenAiChunk(completionId, model, created, "", false, "stop"), |
|
|
buildOpenAiChunk(completionId, model, created, "", false, "stop"), |
|
|
"[DONE]")) |
|
|
"[DONE]")) |
|
|
// 错误兜底:脱敏错误信息,避免泄露内部细节;role 是否出现取决于此前是否已发出过内容片 |
|
|
|
|
|
|
|
|
// 错误兜底:脱敏错误信息,避免泄露内部细节(首片 role 已提前发出,此处不再带 role) |
|
|
.onErrorResume(e -> openAiFallbackStream(completionId, model, created, |
|
|
.onErrorResume(e -> openAiFallbackStream(completionId, model, created, |
|
|
"抱歉,AI 服务调用失败:" + maskError(e.getMessage()), first.get())); |
|
|
|
|
|
|
|
|
"抱歉,AI 服务调用失败:" + maskError(e.getMessage()), false)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
/** |
|
|
/** |
|
|
|