从 Beego 到 SSE:WebAgent 多维问答调用链实现

jessy jessy #webagent#backend#architecture#sse#beego#rag#mcp

沿当前 MixClaw 的真实请求时间线,拆解 WebAgent 多入口问答服务中的身份归一化、持久会话、RAG、MCP、本地文件工具、并发控制与 SSE 终态。

从 Beego 到 SSE:WebAgent 多维问答调用链实现

面向开发者设计一个 Agent 问答服务时,真正困难的通常不是“调用一次大模型”,而是如何让不同身份、不同知识范围、不同工具权限和不同会话策略,共用一套稳定的运行时。

本文以 WebAgent 的三个 SSE 接口为例:

  • GET /api/sse/ask
  • GET /api/sse/tenant_ask
  • GET /agent/sse/ask

我们按当前 MixClaw 中一个请求真实发生的时间线,从 Beego HTTP 路由注册讲到身份归一化、会话持久化、缓存、RAG、MCP 与本地文件工具、并发控制和 SSE 终态。为便于阅读,代码片段会省略无关字段和错误包装,但调用顺序、能力边界和事件名称均以当前实现为准。

1. 三个入口,不是三个 Agent

这三个接口名字相近,但服务对象不同:

入口 主要调用方 身份来源 会话模式 知识范围 工具范围
/api/sse/ask 公开页面、轻量问答 匿名或可选 uid 非托管上下文 平台公共知识 仅内建 RAG
/api/sse/tenant_ask 已登录租户用户 可信网关注入的租户和用户身份 MySQL 持久会话 平台知识 + 当前租户知识 文件工具 + 动态授权的 MCP
/agent/sse/ask 服务端集成、自动化任务 API Key MySQL 持久会话 平台知识 + Key 所属租户知识 文件工具 + Key 精确授权的 MCP

一种容易失控的实现,是为每个入口复制一套“检索、模型、工具、流输出”逻辑。这样做很快会出现三个版本的事件协议、缓存规则和异常处理。

更合适的抽象是:

三个 HTTP 入口
-> 三种接入策略
-> 统一 TaskRequest
-> UnifiedAskAgent.run
-> 简单流或托管流
-> 共用 sseTransport 编码事件

也就是说,入口负责回答“你是谁、允许做什么、会话如何管理”,UnifiedAskAgent 负责回答“这次问题应该怎样完成”。同步 Ask 与 SSE SubscribeAsk 最终也进入同一个 run,避免两条调用方式产生不同结果。

2. 一次请求的完整时间线

先建立全局视角:

flowchart TD
    A[应用启动] --> B[装配数据库 Redis 检索器 模型 工具目录]
    B --> C[Beego 注册三个路由]
    C --> D[HTTP 请求进入 Handler]
    D --> E[解析并归一化 TaskRequest]
    E --> F[鉴权 限流 并发配额]
    F --> G{托管会话}
    G -- 否 --> H[读取可选非托管历史]
    G -- 是 --> I[PrepareTurn 幂等检查 加锁 创建消息占位]
    H --> J[UnifiedAskAgent.run]
    I --> J
    J --> K{已有文件任务可恢复}
    K -- 是 --> L[直接返回幂等恢复结果]
    K -- 否 --> M[构造本次授权 RuntimeTool 集合]
    M --> N{是否有文件或 MCP 工具}
    N -- 否 --> O[RAG-only 流式生成]
    N -- 是 --> P[默认 RAG + 统一模型工具循环]
    P --> Q[artifact_generate 或 tool_search/MCP]
    O --> R[SSE 单 writer 输出]
    Q --> R
    L --> R
    R --> S{托管会话}
    S -- 否 --> T[发送 done]
    S -- 是 --> U[持久化模型终态]
    U --> V{是否创建文件任务}
    V -- 否 --> T
    V -- 是 --> W[等待文件投影终态]
    W --> T
    T --> X[释放会话锁和 API Key 并发槽]

这条时间线里有两个关键原则:

  1. 身份、授权和会话准备必须发生在模型调用之前。
  2. 对持久会话来说,done 不是“模型停止生成”,而是“本轮结果已经可靠落库”。

3. 启动阶段:先构造组合根

路由注册之前,应用会把共享依赖装配好。当前调用链中最重要的四个进程级对象是:

orchestrator.Service
- 检索器与模型
- MCP 动态工具和 MCPToolCatalog
- artifact job creator / replay reader / context reader
conversation.Service
- MySQL 会话与消息仓库
- Redis history cache
- Redis conversation lock
apikey.Service
- API Key 鉴权、scope、RPM、并发槽与 usage tracker
ArtifactProjectionObserver
- 文件任务状态、文件结果与 assistant 终态之间的 barrier

这样做有三个收益:

  • Handler 不直接创建数据库、Redis 或模型客户端,便于测试和替换。
  • 三条入口共享同一个 orchestrator.Service,行为差异由 TaskRequest 字段和可信身份表达。
  • 缓存、鉴权、检索、工具调用都有清晰的所有权,不会散落在 HTTP 层。

MCP 工具仍然分为三个阶段:动态工具定义负责注册,MCPToolCatalog 负责 skill 与工具发现,MCPToolAccessChecker 在候选集合构造时执行授权。执行器只会拿到这次请求已授权并已绑定的 RuntimeTool。本地 artifact_generate 不经过 MCP scope,但必须满足可信托管会话和文件运行时可用两个条件。

基础设施也不应在 Handler 中临时创建。数据库连接池、Redis 客户端和模型客户端通常都需要复用;每次请求重新创建不仅增加延迟,还会放大连接数,并让资源关闭时机变得不可控。

4. Beego 路由注册:用 Handler Factory 注入依赖

Beego 可以通过函数式路由注册 HTTP Handler。对于有依赖的接口,可以使用 Handler Factory,把已经构造好的服务闭包进去:

beego.Get("/api/sse/ask", agentapi.SSEHandler(orchestratorSvc))
beego.Get("/api/sse/tenant_ask", agentapi.TenantSSEHandler(
agentapi.TenantSSEHandlerDeps{
Service: orchestratorSvc, Conversation: conversation,
ArtifactObserver: artifactObserver,
}))
beego.Get("/agent/sse/ask", agentapi.AgentSSEHandler(
agentapi.SSEHandlerDeps{
Service: orchestratorSvc, Identity: identityService,
Conversation: conversation, ArtifactObserver: artifactObserver,
}))

注册完成后,请求的第一段链路是:

Beego HTTP Server
-> Router 匹配 method + path
-> 对应 Handler
-> 接入策略
-> SSE Transport
-> UnifiedAskAgent

Handler Factory 的意义不只是少写全局变量。它让“路由绑定”和“业务执行”分离:启动时决定依赖,请求时只处理输入和输出。

这种写法还有一个隐含好处:函数参数形成了显式的能力白名单。公共 Handler 没有拿到 API Key 服务和会话仓库,就不会因为后续改动而意外获得这些能力。测试时也可以只注入小型 fake,而不需要启动完整应用。

5. 请求归一化:让运行时看不到 HTTP 差异

三个接口的 header、query 和鉴权方式不同,但进入 Runtime 后都会变成同一个 domain.TaskRequest

type TaskRequest struct {
RequestID string
ConversationID string
Query string
Language string
TimeZone string
AccountID string
UserID string
TenantAsk bool
AllowMCPTools bool
MCPPrincipal *MCPPrincipal
AssistantMessageID string
PriorMessages []Message
ManagedConversation bool
}
type askStreamer interface {
SubscribeAsk(context.Context, TaskRequest) <-chan AskEvent
}

可信边界由 Handler 写入这些字段:

  • 公共入口不写 AccountIDTenantAskAllowMCPToolsManagedConversation,因此只能走平台 RAG。
  • 租户入口从 X-Auth-Account-UIDX-Auth-UID 和网关 Bearer 构造租户 MCPPrincipal,并设置 TenantAskAllowMCPTools
  • Agent 入口从 API Key 记录构造 AccountID、Key ID 和精确 scopes;只有 Key 至少包含一个合法 MCP scope 时才打开 AllowMCPTools
  • 托管 Transport 在 PrepareTurn 后补齐 AssistantMessageIDPriorMessagesManagedConversation。客户端不能直接提交这些内部字段。

MCPPrincipal 只用于下游 MCP 鉴权和可信身份注入。Runtime 不再读取 HTTP header,也不会采用模型自行生成的租户或用户身份。文件工具的授权不由客户端布尔值决定,而由完整的托管身份、assistant message 和已配置运行时共同决定。

6. 第一条链路:/api/sse/ask

公共入口强调低门槛和低状态。它的处理过程可以拆成六步:

  1. 拒绝已经废弃的 mode 参数并校验 querylanguagetime_zone 缺省为 zh-CNAsia/Shanghai
  2. 接收可选的 request_idconversation_iduid,其中 uid 只做格式校验,不会形成租户身份。
  3. 构造非托管 TaskRequest,不设置 TenantAskAllowMCPTools 和可信文件字段。
  4. UnifiedAskAgent 读取可选的旧会话存储历史;没有 conversation_id 时就是单轮请求。
  5. Runtime 发出 rag_search 候选及工具事件,检索平台知识后流式生成。
  6. 模型结束后由简单流直接发送 done,不创建 MySQL turn。
sequenceDiagram
    autonumber
    participant C as Client
    participant B as Beego Handler
    participant S as Simple SSE
    participant U as UnifiedAskAgent
    participant R as Retriever
    participant M as Chat Model

    C->>B: GET /api/sse/ask?query=...
    B->>B: 拒绝 mode 校验 query/uid
    B->>S: 非托管 TaskRequest
    S-->>C: 200 text/event-stream
    S->>U: SubscribeAsk
    U-->>S: tool_candidates(rag_search)
    U-->>S: tool_start(rag_search)
    U->>R: 检索平台知识
    R-->>U: 证据片段
    U-->>S: tool_done(rag_search)
    U-->>S: generate_start
    U->>M: system + query + evidence
    loop 增量生成
        M-->>U: chunk
        U-->>S: delta
        S-->>C: event: delta
    end
    U-->>S: done
    S-->>C: event: done

公共入口的边界是:

  • 非托管会话存储只保存带 TTL 的历史。
  • 不能依靠客户端提交的租户 ID 扩大检索范围。
  • 不创建权威的 conversation/message 记录。
  • 不提供需要业务身份的 MCP 工具。

7. 第二条链路:/api/sse/tenant_ask

租户入口通常位于登录网关之后。网关完成登录态校验,再向内部服务注入身份,例如:

X-Auth-Account-UID: tenant-id
X-Auth-UID: user-id
Authorization: Bearer <delegated-credential>

安全边界必须明确:公网客户端不能绕过网关直接访问该入口;网关应先删除外部同名 header,再写入经过验证的值。

Handler 收到请求后:

  1. 校验租户和用户标识。
  2. 构造租户 MCPPrincipal,Bearer 只用于后续 MCP 权限检查。
  3. 校验已有会话,或在未传 conversation_id 时创建会话。
  4. 调用 PrepareTurn 做幂等检查、会话加锁和消息占位。
  5. 把已完成的历史消息传给 Runtime。
  6. Runtime 构造本次授权工具集合:运行时可用时加入 artifact_generate,有已授权 MCP 工具时加入 MCP 候选目录。
  7. 首轮向模型绑定 artifact_generate 和/或 tool_search;符合默认检索条件时先执行 rag_search
  8. 模型在最多六轮的统一循环中直接回答、调用文件工具,或先用 tool_search 发现再调用 MCP 工具。
  9. Runtime 的内部 done 到达后,将 assistant 从 streaming 更新为终态;若创建了文件任务,则进入文件投影 barrier。
  10. 普通回答持久化后发送 done;文件回答等待所有相关 job 进入终态、投影出文件结果后再发送 done
sequenceDiagram
    autonumber
    participant C as Gateway / Client
    participant H as Tenant Handler
    participant CS as Conversation Service
    participant DB as MySQL
    participant R as Redis
    participant U as UnifiedAskAgent
    participant T as Runtime Tools
    participant M as Chat Model

    C->>H: tenant_ask + trusted identity
    H->>CS: EnsureConversation
    CS->>DB: 创建或校验会话
    H->>CS: PrepareTurn
    CS->>DB: 按 request_id 查询已有轮次
    CS->>R: 获取 conversation token lock
    CS->>DB: 二次查重
    CS->>DB: 创建 user completed + assistant streaming
    CS->>R: 读取版本化历史缓存
    opt 缓存未命中
        CS->>DB: 读取完整问答对
        CS->>R: 回填历史缓存
    end
    H-->>C: event: turn_start
    H->>U: SubscribeAsk + PriorMessages
    U->>U: 构造授权 RuntimeTool registry
    U-->>C: tool_candidates(artifact_generate/tool_search)
    opt 默认 RAG
        U->>T: rag_search
        T-->>U: evidence
    end
    loop 最多六轮
        U->>M: history + evidence + tool results
        M-->>U: delta 或 tool calls
        U-->>C: delta
        opt 文件或 MCP tool call
            U->>T: artifact_generate 或已发现 MCP
            T-->>U: structured tool result
        end
    end
    U-->>H: runtime done
    alt 普通回答
        H->>CS: CompleteTurn
        CS->>DB: assistant streaming -> completed
    else 已创建文件 job
        H->>H: 记录模型终态并等待 artifact projection
        H->>CS: 投影 assistant processing -> completed
    end
    H-->>C: event: done
    H->>R: token-safe release lock

这里的 delegated credential 只用于后续权限服务验证,不应被当作租户身份本身。租户身份必须来自已经建立信任关系的网关。

8. 第三条链路:/agent/sse/ask

Agent 入口面向服务间调用。它使用 API Key,而不是浏览器登录态:

Authorization: Bearer wa_xxxxxxxxx

当前 API Key 入口按以下顺序执行:

  1. 从 Bearer header 提取原始 Key。
  2. 计算安全哈希,用哈希值查询 Key,数据库不保存可还原明文。
  3. 检查状态、过期时间、所属租户和允许的 scopes。
  4. 用 Redis Lua 脚本原子检查并占用该 Key 的并发槽位。
  5. 用独立的 Redis 窗口计数检查 RPM。
  6. 创建或确认以稳定 Key ID 为 owner 的持久会话。
  7. 构造 API Key MCPPrincipal;只有精确 MCP scopes 才会打开 MCP 工具能力。
  8. 进入与租户入口相同的托管流,并用 usageTrackingStreamer 消费后端 usage 事件。
  9. 托管流写出 done 后,defer 再完成 usage 持久化、token 计数、指标记录和并发槽释放。
sequenceDiagram
    autonumber
    participant C as API Client
    participant H as Agent Handler
    participant K as API Key Service
    participant R as Redis
    participant CS as Conversation Service
    participant U as UnifiedAskAgent
    participant T as MCP Tool
    participant DB as MySQL

    C->>H: agent/sse/ask + Bearer API Key
    H->>K: Authenticate(hash)
    K->>DB: 查询 Key 状态 租户 scopes
    K->>R: 获取并发槽
    K->>R: RPM 检查
    K-->>H: api_key principal
    H->>CS: EnsureConversation + PrepareTurn
    CS->>R: 获取会话锁
    CS->>DB: 创建消息轮次
    H->>U: Subscribe
    U->>U: 绑定 artifact_generate + tool_search
    opt 模型发现并选择 MCP
        U->>T: 调用 Key 精确授权的工具
        T-->>U: result
    end
    U-->>H: SSE events + backend-only usage
    H->>CS: 持久化 assistant 终态
    H-->>C: event: done
    H->>K: Finalize usage 记录 tokens latency status
    H->>R: 释放并发槽

API Key scope 只接受对话入口和精确 MCP 工具两类表达:

dialogue:sse
mcp:crm:customer_search
mcp:ticket:create

rag:tenantartifact:generateagent:all 和通配符 mcp:service:* 都不是当前合法 scope。租户 RAG 由 Key 所属 AccountUID 决定;本地 artifact_generate 由可信托管会话和文件运行时状态决定,不占用 MCP scope。

9. 两种 SSE Transport:简单流与托管流

三条入口共用 sseTransport 负责 header、JSON 编码、写入和 Flush,但简单流与托管流使用不同的生命周期控制函数。

9.1 简单 SSE

公共入口可以使用简单流:

func serveSimpleAskSSE(
ctx *beegoctx.Context,
streamer askStreamer,
req domain.TaskRequest,
) {
transport, ok := startSSETransport(ctx)
if !ok {
return
}
modelCtx, cancel := context.WithTimeout(
context.WithoutCancel(ctx.Request.Context()), 5*time.Minute,
)
defer cancel()
events := streamer.SubscribeAsk(modelCtx, req)
for {
select {
case <-ctx.Request.Context().Done():
return
case event, open := <-events:
if !open || transport.writeAskEvent(event) != nil {
return
}
}
}
}

它不创建数据库消息占位,Runtime 结束即可发送 done。这条路径的优点是延迟低、依赖少,适合公开问答和临时上下文;代价是不能承诺断线重放、严格幂等和会话终态查询。因此不要为了少写几行代码,把需要可靠会话的租户请求也塞进简单流。

9.2 托管 SSE

租户入口和 Agent 入口需要托管流:

func serveTenantConversationSSE(
ctx *beegoctx.Context,
streamer askStreamer,
lifecycle conversationLifecycle,
owner domain.Owner,
req domain.TaskRequest,
observer ArtifactTurnObserver,
) {
prepared, err := lifecycle.PrepareTurn(
ctx.Request.Context(), owner,
req.ConversationID, req.RequestID, req.Query,
)
if err != nil {
writeTenantConversationError(ctx, err)
return
}
defer releasePreparedTurn(ctx.Request.Context(), lifecycle, prepared)
transport, ok := startSSETransport(ctx)
if !ok {
return
}
transport.writeEvent("turn_start", idsFrom(prepared.Turn))
if prepared.Replay {
writeReplay(ctx.Request.Context(), transport, observer, lifecycle, prepared)
return
}
req = attachManagedTurn(req, prepared)
events := streamer.SubscribeAsk(modelContext(ctx), req)
for event := range events {
switch event.Phase {
case "error":
abortOrProjectFailedTurn(ctx, lifecycle, observer, prepared, event)
return
case "done":
// 文件 job 存在时,这个 helper 记录模型终态、等待所有 job,
// 投影 artifact 结果并写出最终 done。
if finishAcceptedArtifactTurn(
ctx.Request.Context(), transport, observer,
lifecycle, prepared, completed, event.Response,
) {
return
}
lifecycle.CompleteTurn(modelContext(ctx), prepared,
event.Response.Answer, metadata(event.Response))
transport.writeAskEvent(event)
return
default:
transport.writeAskEvent(event)
}
}
}

两者差异不在“是否支持流式”,而在谁对会话终态负责。

这段伪代码省略了错误投影和重连细节,但保留了当前实现的五个关键顺序:

  1. PrepareTurn 位于 Runtime 之前,阻止无效请求消耗模型和工具资源。
  2. replay 分支不进入 Runtime,保证重试不会重复产生外部副作用。
  3. Runtime 的“生成完成”只是内部信号,托管 Transport 持久化后才把它转换为对外 done
  4. 终态写入和解锁使用有界 cleanup context。直接使用已经取消的请求 context,往往会让数据库更新和 Redis 解锁立即失败;完全使用 background context 又可能导致清理永久挂住。
  5. 文件 job 一旦被 observer 接受,普通 CompleteTurn 分支就不再接管;artifact barrier 负责最终消息状态、文件结果和 done

10. PrepareTurn:持久会话的事务边界

PrepareTurn 是托管链路的生成前事务边界,当前签名显式接收已经由 Handler 构造的 domain.Owner

func (s *Service) PrepareTurn(
ctx context.Context,
owner domain.Owner,
conversationID, requestID, query string,
) (*PrepareResult, error)

它按以下顺序执行:

  1. 校验 owner、conversation、query 和 request_id;缺少 request id 时由服务端生成。
  2. 第一次调用 FindTurnByRequest。终态或 processing 轮次直接作为 replay 返回;超过 stale 阈值的 streaming 轮次先修复。
  3. 通过 SET NX 获取 mixclaw:convlock:{account}:{user}:{conversation}。锁已占用时返回 conversation_busy,当前实现不会在 Redis 上排队。
  4. 加锁后再次调用 FindTurnByRequest,关闭两个相同请求同时通过第一次查询的竞态窗口。
  5. BeginTurn 在 MySQL 事务中校验归属并创建 user completedassistant streaming 两条消息。
  6. 根据 conversation 的 history_version 读取 Redis 历史缓存;miss 或版本不一致时从 MySQL 读取完整问答对并回填。
  7. 返回带 lock handle、稳定消息 ID 和 HistoryPrepareResult

Redis 锁、数据库事务和 request id 各自解决不同问题:Redis 锁阻止同一会话同时启动昂贵任务,数据库保证权威消息状态,request id 负责跨进程重启仍可识别的幂等与重放。

assistant placeholder 也不是为了提前“占一条空记录”。它为客户端断线、服务重启和后台恢复提供稳定锚点,让系统能区分正在生成、已经失败和可以标记 stale 的轮次。

一轮消息实际使用下面的状态:

消息 初始状态 成功终态 异常终态
user completed completed completed
assistant streaming completed failed / canceled
assistant(文件 job 已接受) processing completed 由 artifact projection 汇总失败结果后完成投影

BeginTurn 的数据库事务内完成:

  1. 校验 conversation 归属。
  2. 锁定或检查会话版本。
  3. 插入已完成的 user message。
  4. 插入 streaming assistant placeholder。
  5. 保存唯一 request_id
  6. 返回 user message ID、assistant message ID 和 history version。

客户端收到的第一个托管事件可以是:

event: turn_start
data: {
"conversation_id": "...",
"request_id": "...",
"user_message_id": "...",
"assistant_message_id": "..."
}

这样客户端能在答案开始前建立稳定的本地状态。

11. 会话持久化在哪里发生

持久化不是最后一次性保存,而是分为三段:

11.1 生成前

  • 创建或确认 conversation。
  • 创建 user message。
  • 创建 assistant placeholder。
  • 保存 request id。
  • 读取历史。

这一步失败时,不应该调用模型。

11.2 生成中

  • 文本增量只通过 SSE 发送。
  • 当前实现不逐 token 写数据库,也不在统一工具循环中周期性保存 checkpoint。
  • provider usage 通过后端专用事件交给 Agent 入口的 usageTrackingStreamer,不会发给 SSE 客户端。
  • 客户端断开时通过 context 取消后续检索、模型和工具任务。

11.3 生成后

  • 普通回答把 assistant 内容和 metadata 写入数据库,将状态改为 completed,再刷新版本化历史缓存。
  • 文件 job 已创建时,先保存模型 metadata,再把模型终态交给 ArtifactProjectionObserver
  • 文件 barrier 会在 job 运行期间把 assistant 投影为 processing,所有 job 终态后写入文件结果并完成 assistant。
  • 上述权威状态成功后才发送 done

如果异常退出,应把 placeholder 更新为 failedcanceled 或可恢复的 stale,避免会话里永久存在“正在生成”的消息。

12. MySQL 与 Redis:权威状态和加速状态要分开

一个稳健的约束是:

MySQL = 会话与消息的权威状态
Redis = 锁、配额、临时上下文和可丢弃缓存

12.1 历史缓存

当前历史缓存 key 包含完整隔离维度:

mixclaw:conv:v2:{account_id}:{actor_id}:{conversation_id}

value 建议包含:

{
"history_version": 42,
"messages": [
{"role": "user", "content": "..."},
{"role": "assistant", "content": "..."}
]
}

缓存只包含完整的 user/assistant 问答对,不包含 streamingfailed 消息。读取时先比对数据库中的 history_version

  • 版本一致:直接使用缓存。
  • 版本不一致或不存在:从数据库重建并回填。
  • 数据库成功、Redis 回填失败:请求仍可成功,只记录缓存降级。
  • Redis 命中但数据库版本不匹配:丢弃缓存,不能相信旧历史。

写入顺序应是“数据库先成功,Redis 后刷新”。反过来会让缓存暴露尚未持久化的消息。

12.2 RAG 检索缓存

RAG retrieval cache 同样隔离租户、查询、检索数量和索引版本。当前 key 的公共部分是:

mixclaw:rag:retrieval:rag-retrieval-v2:{account_id}:k{top_k}:{query_hash}

随后按检索模式追加版本:

  • 平台检索::p{platform_version}
  • 租户检索::t{tenant_version}
  • 平台与租户联合检索::p{platform_version}:t{tenant_version}

空租户归一化为 admin,查询原文使用 SHA-256 哈希,关键词 fallback 结果不会写入缓存。文档索引更新时递增对应版本,旧 key 会自然失效。

不能只用 query 作为 key,否则不同租户可能读到彼此的检索结果。

13. Runtime 如何统一执行 RAG、MCP 和文件任务

当前实现没有先于主模型运行的独立 Planner/Router,也不再用 artifact-first 路由提前终止本轮。同步 Ask 和 SSE SubscribeAsk 都进入 UnifiedAskAgent.run,按同一顺序执行:

读取 PriorMessages 或非托管历史
-> 恢复已存在的同 request_id 文件 job(命中则直接返回)
-> 构造本次授权 RuntimeTool registry
-> 无本地/MCP工具:RAG-only
-> 有工具:默认 RAG + 最多六轮的主模型工具循环
-> 汇总 answer、artifact_status 和 conversation_status
-> 产生内部 done

入口只提供可信身份与能力开关,不直接指定底层工具。服务端决定哪些工具可以进入 registry,主模型再根据完整历史和当前问题决定直接回答或调用哪个已绑定工具。

13.1 RAG 分配

公共入口和任何没有 RuntimeTool 候选的请求走 RAG-only:

tool_candidates = [rag_search]
rag_search -> generate_start -> model stream

托管请求存在文件或 MCP 候选时进入工具循环。对于没有显式非 RAG 工具意图的问题,循环开始前仍会执行一次默认 rag_search,把证据写入当前 user message。检索范围由服务端字段决定:

public -> platform
tenant / agent -> platform + trusted tenant

TenantAsk 当前固定使用 SearchWithTenant。租户范围来自网关或 API Key 解析出的 AccountID,不是客户端工具参数。

13.2 MCP 工具发现

MCP 工具不会在第一轮全部绑定。Runtime 先构造已经授权的 MCP registry,只把 tool_search 暴露给主模型:

MCPPrincipal
-> MCPToolAccessChecker 过滤动态工具
-> MCPToolCatalog 提供 compact skill catalog
-> 首轮绑定 tool_search
-> tool_search 召回最多 7 个工具
-> 下一轮替换上一次发现的 MCP 工具并重新绑定 schema
-> 模型发出 tool call
-> mcpToolCaller 再次检查 scope/permission 并注入可信参数
-> 调用 MCP server

API Key 只允许 scope 精确匹配的 mcp:<service>:<tool>;租户 principal 使用网关 Bearer 调授权服务。候选构造和远程执行各检查一次,防止目录变化、权限撤销或模型构造越权调用。

13.3 可信参数注入

模型生成的 MCP 参数是不可信输入。mcpToolCaller 会复制业务参数,校验模型提供的 account_uid 不与可信身份冲突,然后覆盖注入:

{
"arguments": {
"keyword": "模型提供的业务参数"
},
"principal": {
"account_id": "服务端验证后的租户",
"actor_id": "服务端验证后的用户或 Key"
},
"data_scope": {
"type": "授权服务返回的范围"
}
}

模型在 arguments 中提供另一个租户时会被拒绝。userdata_scope 只能由服务端 principal 与授权结果产生。

13.4 文件任务

文件能力只在可信托管请求且 creator、reader、支持 function calling 的模型均可用时加入 registry。主模型能看到的唯一公开文件工具是无参数的 artifact_generate,不会看到内部 create_artifact

主模型(携带完整 PriorMessages)
-> artifact_generate {}
-> 读取可信历史、已有文件与同 request_id job
-> 结构化意图解析与服务端策略判定
-> 需要补充信息:返回 clarification payload 给主模型
-> 可以执行:加载目标格式 Skill 和可选 RAG 证据
-> 内部强制调用 create_artifact
-> 校验格式专属 schema,必要时做一次修复
-> 创建异步 artifact job
-> 返回 queued/clarification_required/generation_failed 给主模型

澄清 payload 会回到同一个保留历史的主模型,由主模型用用户当前语言继续提问。SSE 对外只展示 artifact_generatetool_start/tool_donecreate_artifact 是内部实现细节。文件 job 创建后,托管 Transport 的 artifact barrier 等待文件投影终态,再发送最终 done

14. SSE 事件协议:只有一个写协程

当前客户端可见事件如下:

事件 含义
turn_start 托管会话已创建本轮消息
tool_candidates 当前轮实际绑定的工具,带 roundorigin 和服务信息
tool_start 开始调用 rag_searchtool_searchartifact_generate 或 MCP 工具
tool_done 工具调用完成,返回脱敏摘要
tool_error 工具参数或执行失败
tool_unavailable 模型不支持 function calling 或工具绑定失败并降级
generate_start 模型开始生成
delta 文本增量
artifact_progress 文件 job 的阶段、进度、ETA 和投影状态
error 本轮错误
done 客户端可见终态

Runtime 内部还会产生带 ProviderUsageusage 事件,但 sseTransport 会把它过滤,只交给 Agent 入口的 usage tracker。当前没有对外 routeretrieval_startretrieval_done;RAG 使用普通 tool_start/tool_done 表达。未知的内部 phase 会映射成 SSE step

无论内部有多少 goroutine,简单流或托管流的事件消费循环都是 http.ResponseWriter 的唯一 owner:

func (t *sseTransport) writeAskEvent(event AskEvent) error {
if event.Usage != nil {
return nil // backend-only
}
data, err := json.Marshal(event)
if err != nil {
return err
}
return t.writeFrame(sseEventType(event.Phase), data)
}
func (t *sseTransport) writeFrame(eventType string, data []byte) error {
frame := []byte(fmt.Sprintf("event: %s\ndata: %s\n\n", eventType, data))
if _, err := t.writer.Write(frame); err != nil {
return err
}
return http.NewResponseController(t.writer).Flush()
}

startSSETransport 会先检查 http.Flusher,设置 Content-Type: text/event-streamCache-Control: no-cacheConnection: keep-aliveX-Accel-Buffering: no,再写入 200

单 writer 可以避免事件交叉、JSON 截断和 Flush 顺序混乱。http.ResponseWriter 通常不承诺并发安全;即使某个实现碰巧没有数据竞争,也无法保证 tool_startdeltadone 的业务顺序。

单 Writer 并不意味着内部只能串行工作。检索器、模型回调和工具协程可以并发生产 AskEvent,channel 负责线程安全地汇合,Writer 再把内部并发收敛为一条有序网络流。

当前背压和取消行为是:

  • UnifiedAskAgent.StreamAsk 使用容量为 10 的 channel;满载时生产者阻塞,事件不会静默丢弃。
  • tool goroutine 只能调用 push,不能直接写 ResponseWriter。
  • 写入或 Flush 失败时消费循环取消模型或中止托管 turn。
  • 托管流使用请求 context 和五分钟生成超时;客户端断开会停止当前观察与模型链路,但不会反向取消已经提交的持久化 artifact job。

15. 并发处理:四个层次分别解决

“支持并发”不是简单地多开 goroutine。一个多租户 Agent 服务至少要处理四类并发。

15.1 API Key 限流与并发槽

/agent/sse/ask 使用两套独立 Redis 计数:

  • 每分钟请求数。
  • 当前活跃请求数。
  • 当前没有在这条 Handler 中额外增加租户级全局上限。

Handler 先调用 CheckAndIncrConcurrent。Lua 脚本执行 INCR,超过上限时立即 DECR 并返回限流;成功后注册 defer,在所有正常退出路径调用 DecrConcurrent。随后 CheckRPM 用另一段 INCR + PEXPIRE 脚本维护一分钟窗口。两项检查各自原子,但不是同一个事务。

TPM 和每日 token 当前不在建流前预扣。流结束后,usage tracker 根据 provider usage 或估算结果调用 RecordTokens,因此这两项属于事后记录。

15.2 同一会话串行,不同会话并行

同一 owner 下的同一 conversation 实际使用以下锁串行处理:

lock key = mixclaw:convlock:{account_id}:{actor_id}:{conversation_id}
value = random token

获取使用带过期时间的 SET NX。释放时不能直接 DEL,而要执行“value 仍等于我的 token 才删除”的 Lua 逻辑,防止旧请求误删已经续租或重新获取的锁。

-- 只有锁仍属于当前请求时才释放。
-- 如果 TTL 已经过期并被另一个请求重新获取,旧 token 不会误删新锁。
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("DEL", KEYS[1])
end
return 0

随机 token 是锁的所有权证明。当前 lock TTL 至少五分钟,非法或过小配置会使用六分钟默认值;托管模型生成超时是五分钟。当前实现没有锁续租逻辑,因此文档不承诺超过该窗口继续持有会话锁。

数据库事务和行锁是第二道防线。即使 Redis 锁因超时失效,数据库中的唯一 request id、会话版本或 FOR UPDATE 仍应阻止重复 turn。

不同 conversation 没有顺序依赖,可以并行执行。

15.3 request id 幂等

客户端可以为每一轮提供稳定的 request_id;托管请求缺失时由服务端生成。消息表使用以下唯一约束:

unique(conversation_id, request_id, role)

相同 request id 再次请求时:

  • completedfailedcanceled:重放已经持久化的 assistant 结果,不再次进入模型。
  • processing:文件任务重连会读取当前 projection,并继续等待 job 更新。
  • 活跃的 streaming:返回进行中冲突;超过 stale 阈值后由会话服务修复为可诊断终态。

文件创建另外使用 request_id + ":artifact:v1" 作为 job 的 ClientRequestID,并在进入模型和工具内部都检查已有 job,避免取消竞态导致重复创作或重复提交。当前通用 MCP 执行器不会自动为所有远程工具注入幂等键,有副作用的 MCP 工具仍需由对应服务契约保证幂等。

15.4 同一轮多个工具并发

模型一次返回多个 tool calls 时,当前实现为每个已验证调用启动一个 goroutine,并按原调用位置回填结果:

func executeToolCallsWithState(
ctx context.Context,
calls []ToolCall,
) []*Message {
results := make([]*Message, len(calls))
completed := make(chan item, len(calls))
var wg sync.WaitGroup
for i, call := range calls {
tool, args, err := validateBoundCall(call)
if err != nil {
results[i] = toolErrorMessage(call, err)
continue
}
push(toolStartEvent(tool, call, args))
wg.Add(1)
go func(pos int, call ToolCall, tool RuntimeTool) {
defer wg.Done()
result, err := executeRuntimeTool(ctx, tool, args)
push(toolTerminalEvent(tool, call, result, err))
completed <- item{pos: pos, message: toolMessage(call, result, err)}
}(i, call, tool)
}
wg.Wait()
close(completed)
for item := range completed {
results[item.pos] = item.message
}
return results
}

不能按完成顺序把 tool result 交回模型,因为模型协议要求结果与原 tool call 一一对应。预分配切片和带 pos 的结果 channel 保证了下一轮消息顺序;SSE 中并发产生的 tool_done 则按实际完成顺序串行写出。

当前没有独立的工具 semaphore,也没有每个工具单独的 context.WithTimeout;所有调用共享本轮请求 context 和托管生成超时。存在依赖的工具应由主模型分到不同轮次串行调用。tool_search 走专门的发现分支,发现结果只影响下一轮绑定,不会与依赖它的业务工具在同一轮越级执行。

16. done 为什么必须晚于持久化终态

简单流中,模型结束后可以直接发送 done。托管流不能这样做。

假设顺序是:

模型结束 -> 发送 done -> 写数据库失败

客户端会把本轮当作成功,但刷新会话时找不到消息,重试又可能重复调用有副作用的工具。

普通托管回答的正确顺序是:

模型结束
-> assistant 内容和 metadata 持久化
-> history version 更新
-> Redis 历史缓存刷新(失败可降级)
-> 发送 done

文件请求多一层 barrier:

模型结束并返回 artifact_status
-> 保存模型 response metadata
-> 记录 artifact model terminal
-> assistant 投影为 processing
-> 持续发送 artifact_progress
-> 所有关联 job 终态
-> 文件结果与 assistant completed 原子投影
-> 发送 done

Redis 历史缓存刷新失败通常可以降级,因为数据库已经是权威状态;但 assistant 或 artifact 投影终态写入失败时,不能发送成功的 done

Agent API 的 usage 收尾不属于 done barrier。客户端收到 done 后 Handler 返回,defer 才执行 usage Finalize、token 计数、指标记录和并发槽释放。

对于客户端来说,done 应具有明确语义:此后立即查询 conversation、message 或 artifact,都能观察到相同终态。

17. 错误、取消和资源释放

SSE 建立之前的错误适合普通 HTTP 状态码,例如:

  • 400:参数错误。
  • 401:API Key 无效。
  • 403:API Key 已禁用、过期或缺少 dialogue:sse
  • 409:同一会话已有活跃 turn。
  • 429:RPM 或并发配额超限。
  • 503:会话锁、模型或必要依赖不可用。

SSE 建立之后,HTTP 状态已经不能修改,应发送结构化 error 事件,并在服务端完成失败终态。

当前资源释放路径是:

请求结束
-> 取消检索 模型 工具 goroutine
-> 关闭或停止事件生产
-> 标记未完成 assistant 状态
-> 释放 conversation token lock
-> 记录 latency usage status
-> 释放 API Key concurrency slot

会话终态和 Redis 解锁使用脱离请求取消信号、但带三秒上限的 cleanup context。文件 job 一旦提交就是独立持久任务,客户端断开只停止当前 SSE observer,不取消 job。API Key usage 和并发槽通过 defer 收尾。

18. 推荐的模块边界

当前调用链的模块边界如下:

HTTP / Beego
- 路由注册
- 参数解析
- SSE header 与事件编码
Access
- 网关身份解析
- API Key 鉴权
- scope 策略
- 限流与并发配额
Conversation
- conversation 归属
- PrepareTurn
- 幂等与重放
- 消息状态机
- history cache 与 lock
Runtime
- UnifiedAskAgent.run
- RuntimeTool registry
- RAG-only 与统一六轮工具循环
- MCP tool_search、授权与执行
- artifact_generate、内部 create_artifact 与状态汇总
Artifact Projection
- 异步 job 生命周期
- artifact_progress
- 模型终态与文件终态 barrier
Infrastructure
- MySQL
- Redis
- 向量与关键词检索
- Model Provider
- MCP Servers

模块之间用小接口连接,尤其避免让 Runtime 直接依赖 Beego request,也避免让 Handler 直接操作消息表。

19. 三条链路的最终对照

阶段 公共 Ask 租户 Ask Agent Ask
路由 Beego 函数式注册 Beego 函数式注册 Beego 函数式注册
身份 匿名或可选 uid 可信网关 MCPPrincipal API Key MCPPrincipal
前置配额 Handler 内无 Handler 内无 Key 并发槽,再检查 RPM
会话 非托管 TTL 历史 MySQL 持久会话,缺 ID 时创建 MySQL 持久会话,缺 ID 时创建
PrepareTurn 不需要 需要 需要
会话锁 不需要 Redis token lock Redis token lock
历史 可选短期上下文 版本化 Redis cache,MySQL 回源 版本化 Redis cache,MySQL 回源
RAG 平台 平台 + 当前租户 平台 + Key 所属租户
首轮工具 rag_search artifact_generate + 可选 tool_search artifact_generate + scope 允许时的 tool_search
MCP 关闭 目录发现 + 用户权限双检 精确 scope 过滤 + 执行时复检
文件 关闭 可信托管请求可用 可信托管请求可用,不使用 artifact scope
多工具并发 无业务工具 同轮 goroutine 并发,按位置回填;无独立 semaphore/超时 同左
客户端事件 RAG 工具事件、delta、done 加 turn_start、工具事件、artifact_progress 同租户;usage 事件仅后端消费
done 时机 模型结束 普通终态或文件投影终态之后 同租户;usage Finalize 在 done 之后

20. 落地检查清单

修改这条调用链后,可以用下面的问题核对文档与实现是否仍同步:

  • 三个入口是否只负责身份和策略,而不是复制 Agent 逻辑?
  • Runtime 是否只接收 Handler 构造的 TaskRequest,不读取原始 header?
  • 租户 ID 是否始终来自网关或 API Key,而不是模型参数?
  • API Key 是否使用哈希存储、精确 scope、RPM 和并发限制?
  • 持久会话是否在模型调用前创建 user/assistant 消息?
  • request_id 是否通过 (conversation_id, request_id, role) 唯一约束并支持重放?
  • 同一会话是否串行、不同会话是否可以并行?
  • Redis 解锁是否校验 token?
  • 历史缓存是否带 owner、conversation 和 history version?
  • RAG 缓存 key 是否包含 account、topK、query hash 与对应索引版本?
  • MCP 工具是否在发现阶段和执行阶段各鉴权一次?
  • 服务端是否覆盖注入可信身份和 data scope?
  • 多工具结果是否按原 call 位置回填,并明确记录当前没有独立 semaphore 和单工具超时?
  • SSE 是否只有一个 writer,并正确处理 Flush、背压和断连?
  • usage 是否继续只在后端消费,客户端事件表是否没有虚构 routeretrieval_*
  • 客户端是否只看到 artifact_generate,而 create_artifact 保持内部私有?
  • assistant、artifact 和 metadata 是否在 done 前进入终态?
  • 所有失败路径是否都会释放锁、配额并写入可诊断状态?

结语

构建多维调用链的关键,不是为每种调用方维护一套 Agent,而是把差异放在正确的位置:

  • 路由层区分入口。
  • 接入层建立可信身份。
  • 策略层限制知识、工具和会话能力。
  • 会话层负责幂等、锁、历史与终态。
  • UnifiedAskAgent.run 统一完成 RAG、MCP、文件任务和模型循环。
  • SSE Transport 把内部并发收敛成有序、可恢复的事件流。

当这些边界稳定后,增加新的租户类型、认证方式、检索源或 MCP 工具,通常只是扩展策略和能力实现,而不需要重写整条问答链路。