AIManage 把 Provider 的增量分块实时下发到客户端。流不是端到端缓冲:分块一到就 yield,持久化作为有界、至多一次的工作并行进行。线上格式是带 kind 判别字段的类型化 NDJSON 帧协议,每次失败只产生一个稳定错误帧,绝不泄漏原始异常文本。
1. NDJSON 帧合同
每帧是独占一行的 JSON 对象(application/x-ndjson; charset=utf-8)。kind 字段始终存在并决定帧形状;值为 null 的字段被省略。
{"kind":"data","text":"第一段"}{"kind":"data","text":"第二段"}{"kind":"done"}{"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) | data | kind、text |
Failure(errorCode, errorType, detail, traceId, correlationId, requestId) | error | kind + 全部六个错误字段 |
Done() | done | 仅 kind |
序列化使用源生成(NdjsonTextStreamJsonContext),兼容 Native AOT。每帧写入后立即 Flush,首字节延迟跟随 Provider 而非缓冲区。
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. 增量流式与预取边界
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 调用。
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)。
// 身份只能来自已认证调用上下文,不能接受请求体代填。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 失败,然后 yieldAIManageErrors.ChatRequestFailed。 - Provider 内部取消被当作 Provider 失败处理,而非客户端取消。 Provider 内部取消的 token 不走静默客户端取消路径。
- 模糊的
CompleteAsync失败绝不盲目重试。 持久化 attempt 保留在其记录状态,而非猜测。 DisposeAsync路径独立地(尽力而为)持久化任何已交付的部分内容,但绝不产生第二个终态帧。
8. 跨模块缝:IAIConversationPort
其他模块通过稳定的对话端口调用 AIManage,而非内部 Command。AIConversationMediatorPort 把每次调用都路由进 Mediator,因此授权、运行时许可、日志全部生效——没有绕过路径。
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. 验证命令
# 帧合同、错误映射、预取与管线顺序。rg -n "NdjsonTextStreamFrame|CreateAsync\(|ConversationStreamFailed|IAIMessageRequestCoordinator|ConversationNotOwner" \ src/Platform/AIManage src/Framework/BitzOrcas.Framework.AspNetCore -g '*.cs'AIManage 总览 · 对话与安全 · 运行时许可