Files
MailUI4Agents/gateway/internal/handler/agents.go
JianFeeeee 89356d4a9b feat: 每平台可用模型范围 + 降级尝试 + 失败回报
配置页为每个 Agent 平台划定「邮件场景下可用的模型」,插件按顺序逐个尝试,
全部失败把原因封装成邮件回复。目录由插件上报、管理员只做勾选 —— 手打模型名
会打错,而打错的后果要到真发邮件时才暴露成一次失败。

## 目录上报走心跳,不另设端点

模型清单会在运行中变(换 provider 配置、上游上下线、换 API key)。
只在注册时报一次的话目录会静静变陈,管理员在配置页选中一个平台其实调不到的
模型。心跳本来就是 30 秒一次的现成通道;另设一个 POST 等于给「目录是谁写的」
留两个答案,排查时要同时看两处。

心跳响应回传 `allowed_models`,因此管理员改了范围后最多一个周期生效,
不必重启插件。

与 platform_sessions 同一约定:拉不到目录时**省略字段**(保留现有目录),
传空数组会把配置页清成空白。

## 目录与选择分两张表

模型会从平台目录里消失(上游临时下线、换了 provider 配置)。合成一张带
allowed 标记的表时,整行被删就连带把管理员的选择也删了,模型回来还得重配一遍。
分开存之后「选了什么」是持久的,目录只决定「这一项现在是否可用」;
已选但不在目录里的标为 stale 显示出来 —— 不显示会让人以为自己没选过它。

## 最难的一点:模型失败不是同步抛出的

两个平台都踩了。`promptAsync()` 立即返回、`ctx.agents.create()` 不校验模型,
只包 try/catch 的话第二个模型永远不会被试到 —— 第一个无效模型会被判成成功。

必须等异步结论:
- opencode → `session.error` 事件(event 钩子在 deliverMail 之外,
  因此用 turnWatchers 表把两者接起来)
- DSH → `turn/end` 的 `reason.kind === 'error'`

DSH 还有个陷阱:**`assistant/chunk` 不能当成功信号**,它的 `finish` 子类型
也带错误 —— `{chunk:{type:'finish',reason:{kind:'error',failure:{code:'NO_ADAPTER'}}}}`。
实测「无效 provider 却判成功」正是因为把任意 chunk 当成了走通。判据要落在
chunk 的类型上:finish 看 reason,其余才意味着模型真的在产出。

超时按成功处理(60 秒窗口):模型可能只是很慢,把慢当成失败会在换模型的同时
把已经在跑的那一轮丢掉。

DSH 换模型要换会话 id(`<原 id>-r1`)并 dispose 失败那个 agent:复用同一个 id
会让重试接在一条已经出错的会话后面,不 dispose 则 agent/status 还会为那个
死会话触发一次自动转发。

## 其他决策

- **范围优先于环境变量**:范围是运行时可改的策略,`AGENTMAIL_REPLY_*` 是部署时
  的兜底。反过来的话管理员在配置页改了却不生效,得去改 service 文件重启
- **范围为空返回 `[undefined]` 而非 `[]`**:空数组会让调用方一次都不试,
  而「管理员没配」的正确含义是不限定,不是「一个都不许用」
- **上限 10 个**:降级是串行的,选 50 个意味着最坏情况下一封邮件要等 50 次超时
- 前端 key 按**第一个** `/` 切分 provider/model:model id 可能含 `/`
  (如 `org/model-name`),按最后一个切会把 provider 切错
- 保存后用服务端返回的结果刷新界面而非回显入参:repo 层会跳过重复与空字段

## 验证

- Go 10 个新测试(含「模型从目录消失后选择必须留存」的直接回归)
- 两插件各 18 个模型范围测试,共 180 个
- 端到端四轮:正常路由 → 全部无效(收到失败回报邮件,used_rounds 保持 0
  确认走了免配额通道)→ DSH 降级(fake-a 失败 → llmsproxy/AUTO 成功)→
  opencode 降级(nonexistent/bad 失败 → AUTO 成功,日志确认「前 1 个失败」)
- 生产已部署,前端「模型范围」页可用
2026-09-02 21:34:55 +08:00

212 lines
7.1 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package handler
import (
"net/http"
"github.com/agentmail/gateway/internal/middleware"
"github.com/agentmail/gateway/internal/models"
"github.com/agentmail/gateway/internal/repo"
)
// ---------- Agent ----------
type registerRequest struct {
Name string `json:"name"`
Secret string `json:"secret"`
Workspaces []models.Workspace `json:"workspaces"`
Platform string `json:"platform"`
}
// heartbeatRequest 是心跳可选带的上报体。
//
// 字段全可省:旧插件发空心跳,不能因为新增了上报就把它们报错。
type heartbeatRequest struct {
// PlatformSessions 是平台侧当前的会话快照(按最近活跃排序)。
//
// 为什么让插件上报而不是 Gateway 反向拉取:当前架构是单向的
// Agent 持密钥主动连 GatewayGateway 从不外呼)。反向拉取需要 Gateway
// 保存各平台的地址与凭证,那是另一套信任模型。
//
// nil 与空数组语义不同nil = 本次不上报(保留现有镜像),
// 空数组 = 平台侧确实一条会话都没有(清空镜像)。
// 拿不到会话列表的插件应当省略该字段,而不是传空数组把镜像抹掉。
PlatformSessions []repo.PlatformSession `json:"platform_sessions"`
// Models 是平台当前看得见的模型目录,供配置页勾选。
//
// 随心跳上报而不是只在注册时上报:模型清单会在运行中变
// (换 provider 配置、上游上下线、换了 API key。只在注册时报一次的话
// 目录会静静变陈,而管理员在配置页上看到的是上次重启时的快照 ——
// 选中一个平台已经调不到的模型,失败要到真发邮件时才暴露。
//
// 与 PlatformSessions 同一约定nil = 本次不上报(保留现有目录),
// 空数组 = 平台确实一个模型都拿不到。拿不到目录时必须省略:
// 清空目录会让配置页变成空白,管理员以为该平台没有任何可用模型。
Models []repo.CatalogModel `json:"models"`
}
// POST /api/v1/agent/register
//
// 两种认证方式:
// 1. Authorization: Bearer <agent_key_token> —— 密钥认证(推荐)。
// 密钥未绑定时用本请求的 name 落定;已绑定时 name 必须与之一致,
// 否则等于拿别人的密钥冒充新身份。
// 2. body 里带 secret —— 旧方式,兼容保留。
func RegisterAgent(w http.ResponseWriter, r *http.Request) {
var req registerRequest
if err := Decode(r, &req); err != nil {
Error(w, http.StatusBadRequest, "Invalid JSON")
return
}
if req.Name == "" {
Error(w, http.StatusBadRequest, "Missing name")
return
}
keyToken := middleware.BearerToken(r)
if keyToken == "" && req.Secret == "" {
Error(w, http.StatusBadRequest, "需要 Authorization: Bearer <密钥> 或 body 里的 secret")
return
}
if keyToken != "" {
bound, err := repo.VerifyAgentKey(r.Context(), keyToken)
if err != nil {
writeKeyErr(w, err)
return
}
if bound != "" && bound != req.Name {
Error(w, http.StatusForbidden,
"该密钥已绑定到 Agent \""+bound+"\",不能用于注册 \""+req.Name+"\"")
return
}
}
if req.Platform == "" {
req.Platform = "pi"
}
// 三维地址的 name 位与人类用户名共用命名空间,不得重名
if ok, err := repo.AgentNameAvailable(r.Context(), req.Name); err != nil {
Error(w, http.StatusInternalServerError, "Failed to validate agent name")
return
} else if !ok {
Error(w, http.StatusConflict, "该名称已被人类用户占用")
return
}
if req.Name == "human" {
Error(w, http.StatusBadRequest, "human 是保留别名,不能作为 Agent 名")
return
}
// 密钥认证时不需要 secret但 agents.secret 非空约束仍在;
// 存密钥本身作占位,旧的 name/secret 路径不受影响。
secret := req.Secret
if secret == "" {
secret = keyToken
}
if err := repo.CreateOrUpdateAgent(r.Context(), req.Name, secret, req.Platform, req.Workspaces); err != nil {
Error(w, http.StatusInternalServerError, "Failed to register agent")
return
}
// 待绑定密钥在首次注册成功后落定到该 Agent
if keyToken != "" {
if err := repo.ClaimAgentKey(r.Context(), keyToken, req.Name); err != nil {
Error(w, http.StatusInternalServerError, "Failed to bind key")
return
}
}
JSON(w, http.StatusOK, map[string]string{
"status": "registered",
"agent_name": req.Name,
})
}
// POST /api/v1/agent/heartbeat
func HeartbeatAgent(w http.ResponseWriter, r *http.Request) {
agentName := middleware.GetAgentName(r)
if agentName == "" {
Error(w, http.StatusUnauthorized, "Unauthorized")
return
}
pending, err := repo.HeartbeatAgent(r.Context(), agentName)
if err != nil {
Error(w, http.StatusInternalServerError, "Failed to heartbeat")
return
}
// 可选的平台会话快照。解不开就当作没带:心跳的主职责是「我还活着」,
// 不该因为上报体格式不对就把 Agent 判成离线。
var req heartbeatRequest
if r.ContentLength > 0 {
_ = Decode(r, &req)
}
syncedSessions := -1 // -1 = 本次未上报
if req.PlatformSessions != nil {
if err := repo.ReplacePlatformSessions(r.Context(), agentName, req.PlatformSessions); err != nil {
// 镜像写失败只影响候选补全,不影响投递,因此不报错
syncedSessions = -1
} else {
syncedSessions = len(req.PlatformSessions)
}
}
// 模型目录同理:写失败只让配置页看到的目录陈一轮,下一次心跳会补上。
syncedModels := -1
if req.Models != nil {
if err := repo.ReplaceModelCatalog(r.Context(), agentName, req.Models); err == nil {
syncedModels = len(req.Models)
}
}
// 心跳回传该 Agent 的累计统计与新任务默认预算。
//
// 不再回传「剩余额度」:额度属于具体任务(会话)而不属于 Agent
// 剩余往返随每次发信响应budget_remaining回传在那里才有意义。
stats, sErr := repo.GetAgentStats(r.Context(), agentName)
if sErr != nil {
// 统计读不到不影响心跳本身
stats = repo.AgentStats{AgentName: agentName}
}
resp := map[string]interface{}{
"status": "ok",
"pending_mails": pending,
"stats": stats,
}
if syncedSessions >= 0 {
resp["platform_sessions_synced"] = syncedSessions
}
if syncedModels >= 0 {
resp["models_synced"] = syncedModels
}
// 回传当前生效的模型范围,插件无需另起一个请求去读。
//
// 随心跳回传而不是让插件自己轮询:管理员在配置页改了范围后,
// 插件最多一个心跳周期30 秒)就能看到新值,不需要重启。
if allowed, aErr := repo.ListAllowedModels(r.Context(), agentName); aErr == nil {
resp["allowed_models"] = allowed
resp["models_unrestricted"] = len(allowed) == 0
}
JSON(w, http.StatusOK, resp)
}
// GET /api/v1/agents
func ListAgents(w http.ResponseWriter, r *http.Request) {
statusFilter := r.URL.Query().Get("status")
agents, err := repo.ListAgents(r.Context(), statusFilter)
if err != nil {
Error(w, http.StatusInternalServerError, "Failed to list agents")
return
}
JSON(w, http.StatusOK, map[string]interface{}{
"agents": emptySlice(agents),
})
}