Skip to content
bitzorcas
中EN

Reference

AIManage 流式响应、客户端缓存与用量

NDJSON 帧合同、增量流式、强化错误帧、幂等请求、取消与部分持久化、Semantic Kernel 客户端缓存,以及用量治理。

Last updated

AIManage 把 Provider 的增量分块实时下发到客户端。流不是端到端缓冲:分块一到就 yield,持久化作为有界、至多一次的工作并行进行。线上格式是带 kind 判别字段的类型化 NDJSON 帧协议,每次失败只产生一个稳定错误帧,绝不泄漏原始异常文本。

1. NDJSON 帧合同

每帧是独占一行的 JSON 对象(application/x-ndjson; charset=utf-8)。kind 字段始终存在并决定帧形状;值为 null 的字段被省略。

成功流:data 帧序列,最后是 done 帧
{"kind":"data","text":"第一段"}
{"kind":"data","text":"第二段"}
{"kind":"done"}
失败流:data 帧后恰好一个终态 error 帧
{"kind":"data","text":"部分内容"}
{"kind":"error","errorCode":"AI.Conversation.StreamFailed","errorType":"Unexpected","detail":"...","traceId":"...","correlationId":"...","requestId":"..."}

帧类型是 NdjsonTextStreamFrame(src/Framework/BitzOrcas.Framework.AspNetCore/Results/NdjsonTextStreamFrame.cs)。三个工厂产生唯一合法的形状:

工厂kind写入字段
Data(text)datakind、text
Failure(errorCode, errorType, detail, traceId, correlationId, requestId)errorkind + 全部六个错误字段
Done()done仅 kind

序列化使用源生成(NdjsonTextStreamJsonContext),兼容 Native AOT。每帧写入后立即 Flush,首字节延迟跟随 Provider 而非缓冲区。

浏览器消费 NDJSON
const response = await fetch(
`/api/ai/chat/conversations/${conversationId}/stream`,
{
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
requestId: "turn-42",
content: prompt,
modelId: null,
}),
signal: abortController.signal,
},
);
// 许可/授权拒绝在流开始前返回标准 JSON ProblemDetails(403),
// 此时响应体不是 NDJSON。
if (!response.ok) {
const problem = await response.json();
showError(problem.errorCode ?? response.status);
return;
}
const reader = response.body.pipeThrough(new TextDecoderStream()).getReader();
let pending = "";
// 网络 chunk 不保证与 NDJSON 行对齐,必须保留尚未结束的一行。
for (;;) {
const { value, done } = await reader.read();
if (done) break;
pending += value;
const lines = pending.split("\n");
pending = lines.pop() ?? "";
for (const line of lines) {
if (!line) continue;
const frame = JSON.parse(line);
if (frame.kind === "data") renderIncrementally(frame.text);
else if (frame.kind === "error") {
showError(frame.errorCode, frame.errorType);
return; // 终态——流在此停止
}
// kind === "done" 是空终态标记,循环自然结束。
}
}

生产客户端还需检查 HTTP 状态(流开始前的 403 返回 ProblemDetails 而非 NDJSON)、限制单帧大小、做 schema 校验,并处理跨 chunk 的 UTF-8/换行边界。

2. 增量流式与预取边界

ClientNdjsonTextStreamResultStreamMessage handlerProviderClientNdjsonTextStreamResultStreamMessage handlerProvider分块 1yield 分块 1flush {"kind":"data"}分块 2..N逐块 yield逐块 flush枚举完成{"kind":"done"}

NdjsonTextStreamResult.CreateAsync 在返回 HTTP 结果前预取第一项。这是刻意的:许可或授权拒绝发生在响应头发送之前,因此可以被异常中间件投影为标准 403 ProblemDetails,而非流内错误帧。拿到第一项后设置头(no-cache、X-Accel-Buffering: no),随后的增量分块一到就 yield 并 flush。Handler 不先收集完整响应再下发——首字节延迟≈Provider 首 token 延迟,而非完整推理时间。

3. 强化的流式错误合同

流式失败绝不泄漏内部细节。该合同由 NdjsonTextStreamResult.WriteFramesAsync 强制,并被 contract test 钉死:

  • Result<string>.IsFailure 项,或响应开始后的未处理异常,写入恰好一个终态错误帧后停止。只有枚举无失败完成时才写 done。
  • 错误 detail 仅来自 L.GetString(error.Code)——本地化 i18n key。面向开发者的 Error.WithDescription(...) 文本永不序列化。包含 “database shard secret-diagnostic” 的诊断描述不得出现在响应体中。
  • 响应开始后的未处理异常映射到调用方声明的稳定错误(AIManageErrors.ConversationStreamFailed = AI.Conversation.StreamFailed),绝不映射原始异常。原始异常文本不会到达客户端。
  • 即使 DisposeAsync 随后抛出,也最多写一个终态帧。
  • 客户端取消不写错误帧,也不发 LogLevel.Error。
端点声明稳定的未处理流错误
// AIManageEndpoints.cs —— 流路由把 mediator 流包装进结果,
// 并命名在开始后逃逸的未处理异常所用的稳定错误。
return NdjsonTextStreamResult.CreateAsync(
mediator.CreateStream(new StreamMessage.Command(conversationId, request.RequestId, request.Content, request.ModelId)),
HttpContext,
AIManageErrors.ConversationStreamFailed,
logger,
streamName: "ai-conversation");

4. 许可与授权在流开始前门控

流管线由源生成器按固定顺序闭合注册:LoggingStreamPipelineBehavior → RuntimeLicenseStreamPipelineBehavior → AuthorizationStreamPipelineBehavior。两个门控行为都在 Handler 枚举之前抛出:

  • RuntimeLicenseStreamPipelineBehavior 拒绝时抛 RuntimeLicenseExecutionDeniedException。
  • AuthorizationStreamPipelineBehavior 授权决策拒绝时抛 AuthorizationExecutionDeniedException。

因为第一项被预取(第 2 节),这些拒绝发生在响应开始之前,由异常中间件映射为带 errorCode(Licensing.Runtime.Denied 或 Authorization.Denied)的标准 403 ProblemDetails。拒绝请求从不调用对话存储。

5. 幂等请求与 claim 状态机

SendMessage.Command 和 StreamMessage.Command 都要求客户端提供 RequestId。IAIMessageRequestCoordinator 在任何 Provider 调用之前通过数据库唯一约束原子占用该键,因此重试和并发重复都解析为一次 Provider 调用。

StreamMessage.Command 携带幂等键
public sealed record Command(
string ConversationId,
string RequestId, // 必填;重试间保持不变
string Content,
string? ModelId) : IStreamCommand<Result<string>>, IAuthorizedRequest;

协调器的 ClaimAsync 返回 AIMessageRequestClaim,其 State 决定响应:

AIMessageRequestState含义结果
Claimed当前调用方占有键;可恰好调用一次 Provider继续
Pending另一调用方占有且未完成AI.Message.RequestInProgress(Conflict)
Completed已完成;重放已持久化的 assistant 消息重放,不调用 Provider
Failed此前同键请求已失败AI.Message.PreviouslyFailed(Conflict)
Conflict同键携带不同内容或模型AI.Message.RequestConflict(Conflict)

SendMessage.Command 还实现 INonTransactionalCommand:它跳过通用长事务,协调器在网络调用前后用有界短事务提交用户消息 + claim(及随后的 assistant 消息)。Completed 重放已持久化的 assistant 消息而不调用 Provider——成功 turn 之后的重试免费且幂等。

6. 对话所有权被强制执行

两个 Handler 都注入 ICurrentUser,在任何 Provider 调用或消息写入之前,拒绝任何 TenantId 或 UserId 不匹配调用方的对话,返回 AIManageErrors.ConversationNotOwner(AI.Conversation.NotOwner,Forbidden)。

StreamMessage.Handler 所有权守卫
// 身份只能来自已认证调用上下文,不能接受请求体代填。
var caller = currentUser.User;
var userId = currentUser.RequireUserId(AIManageErrors.ConversationUserRequired).GetValueOrThrow();
// 调用 Provider 前同时约束租户和会话 owner。
if (!string.Equals(conversation.TenantId, caller.TenantId, StringComparison.Ordinal)
|| !string.Equals(conversation.UserId, userId, StringComparison.Ordinal))
{
yield return Result.Failure<string>(AIManageErrors.ConversationNotOwner);
yield break;
}

两个 Command 还实现 IAuthorizedRequest,声明 ResourceDescriptor(AIManagePermissions.Module, AIManagePermissions.ConversationResource) 与 AuthorizationAction.Use,因此授权管线在 Handler 运行前求值。

7. 取消、部分持久化与至多一次不变量

Handler 跟踪单个 AssistantPersistenceAttempt(NotStarted → InProgress → Succeeded/Failed),确保每个流枚举器至多一次 assistant 写入,即使跨越取消与 dispose 路径。

  • 有非空部分内容的客户端取消:在有界 5 秒 token 下持久化部分文本,随后重抛 OperationCanceledException 使取消仍传播到客户端。正常情况下部分内容在取消前已 yield。
  • Provider 失败:向已交付的部分追加 \n\n[流式中断] 标记并持久化;若无任何交付则标记 attempt 失败,然后 yield AIManageErrors.ChatRequestFailed。
  • Provider 内部取消被当作 Provider 失败处理,而非客户端取消。 Provider 内部取消的 token 不走静默客户端取消路径。
  • 模糊的 CompleteAsync 失败绝不盲目重试。 持久化 attempt 保留在其记录状态,而非猜测。
  • DisposeAsync 路径独立地(尽力而为)持久化任何已交付的部分内容,但绝不产生第二个终态帧。

8. 跨模块缝:IAIConversationPort

其他模块通过稳定的对话端口调用 AIManage,而非内部 Command。AIConversationMediatorPort 把每次调用都路由进 Mediator,因此授权、运行时许可、日志全部生效——没有绕过路径。

IAIConversationPort —— 稳定的跨模块缝
public interface IAIConversationPort
{
// 创建与非流式回合通过 Result 返回已持久化的摘要。
Task<Result<ConversationSummary>> CreateConversationAsync(
AIConversationCreateRequest request, CancellationToken cancellationToken = default);
Task<Result<MessageDto>> SendMessageAsync(
AIConversationMessageRequest request, CancellationToken cancellationToken = default);
// 流式调用保留逐块业务失败,不暴露内部 Handler。
IAsyncEnumerable<Result<string>> StreamMessageAsync(
AIConversationMessageRequest request, CancellationToken cancellationToken = default);
}

AIConversationCreateRequest(WorkspaceId, Title?, SystemPrompt?) 与 AIConversationMessageRequest(ConversationId, RequestId, Content, ModelId?) 是请求记录。租户/用户身份不由调用方传入——由统一管线内部从可信调用方上下文派生。

9. ChatClientFactory

工厂缓存 IChatClient,最大 32 个条目。键是 ProviderId、ModelId 和 Endpoint/ModelId/ApiKey 哈希;命中时更新 LastAccessed,满容量时找最旧项并 Dispose。具体构建由 IChatClientBuilder 注入,生产实现返回 SemanticKernelAdapter。

通过工厂解析聊天客户端
// 缓存键同时覆盖 Provider、模型、Endpoint 与凭据身份。
IChatClient client = chatClientFactory.GetOrCreate(provider, modelId);
var response = await client.GetResponseAsync(messages, new ChatOptions
{
MaxOutputTokens = 4096,
Temperature = 0.7f
}, cancellationToken);

10. Semantic Kernel 边界

Adapter 注册 AddOpenAIChatCompletion(modelId, apiKey, endpoint),把 MEAI ChatMessage 映射为 SK ChatHistory,仅映射 MaxTokens/Temperature。它不注册 plugin/function,也不启用自动函数调用。商业支持矩阵必须由每种 Provider 的 contract test 证明,包括 Endpoint 形态、认证、模型名、流式、usage 与错误映射。

11. Token、费用与运维治理

非流式 Assistant 消息的 TokenCount 保存 response.Usage.TotalTokenCount;流式保存 null。DTO 不返回 usage 拆分。不可变 Usage Ledger(租户/工作区/对话/turn/Provider/模型、input/output/cached token、价格版本、估算/结算费用、状态)与租户预算强制仍是治理目标而非当前行为。

运维上至少采集:first-byte/full duration、chunk gap、chunks/bytes、取消时机、partial persisted、provider error、cache 活动、并发流、token、费用、预算拒绝。告警应区分 Provider 延迟、代理缓冲、客户端断开。

12. 流式行为故障排查

症状可能层证据
客户端收到 403 JSON 而非 NDJSON首帧前的许可/授权拒绝响应状态 + errorCode(Licensing.Runtime.Denied / Authorization.Denied)
部分文本后一个 error 帧流中 Provider 失败终态帧的 errorCode;持久化的 [流式中断] 标记
流结束但无 done 也无 error 帧客户端取消取消已传播;无错误日志
重试时 AI.Message.RequestInProgress幂等协调器发现重复RequestId 已 Claimed/Pending
AI.Conversation.NotOwner调用方非对话所有者所有权守卫在 Provider 调用前触发

13. 验证命令

Terminal window
# 帧合同、错误映射、预取与管线顺序。
rg -n "NdjsonTextStreamFrame|CreateAsync\(|ConversationStreamFailed|IAIMessageRequestCoordinator|ConversationNotOwner" \
src/Platform/AIManage src/Framework/BitzOrcas.Framework.AspNetCore -g '*.cs'

AIManage 总览 · 对话与安全 · 运行时许可

100%

滚轮或按钮缩放 · 放大后拖动画面 · 双击切换 100% / 200%