feat(swarm): #60 模型 key 注入(方案A)+ pool_terminated 吊销握手
对接 agent_swarm PR#43 定死的参数,实现 HM 侧的 per-user 模型 key 注入:
- A.1 粒度/命名:为认证用户 mint 一把 per-user sk-(getOrMintSwarmModelToken),
跨该用户所有 swarm run 复用;KV 密文名 swarm-model-key-<user_id>。
- A.2 KV value:JSON {"openai_api_key":"sk-..."},对齐 Swarm
orchestrator/agent_launcher._extract_model_key 解析字段。
- A.4 OPENAI_API_BASE 为 Swarm 部署常量,HM 不经 billing_context 下发。
- A.5 吊销:事件驱动。注册 swarm.pool_terminated(swarm_lifecycle 类),
回调 handleSwarmPoolTerminated 删 KV 密文 + 软删 token;user_id 优先取
payload,缺失回退 deployment 上下文。
createAgentDeploymentFromPlan 在 swarm 模式且 billing_context 未自带 secret_ref
时调 provisionSwarmModelKey,把 azkv:// 引用写进 billing_context.secret_ref 下发。
Key Vault 未配置/不可达时降级:记日志、secret_ref 留空,不阻断 create(联调前可用)。
明文 sk- 仅经 secret_ref 服务端解析,绝不入代码/日志/事件/argv。
测试:swarm_model_key_test.go 覆盖 per-user 复用、吊销重 mint、pool_terminated
回调(payload/上下文两路 user_id)、非目标事件不误吊销、事件已注册。
影响面:agent_swarm(契约消费侧)、密钥(per-user sk- mint+吊销)、计费
(走 User.Quota)、审计(回调入库+审计事件)。不影响 Client/CodeGW/发布链路。
未决依赖:A.3(Swarm 读 heicode-kv 的 RBAC,运维授权 pending)是端到端联调前置,
不阻塞本 PR。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -487,6 +487,10 @@ var agentCallbackEventRequiredFields = map[string][]string{
|
||||
"swarm.completed": {},
|
||||
"swarm.failed": {"reason"},
|
||||
"swarm.stopped": {},
|
||||
// #60 A.5:用户全部 run 停机时 Swarm 恰好发一次,payload {user_id, secret_ref(azkv://)}。
|
||||
// 控制面生命周期事件,HM 据此吊销该用户 swarm 模型 key。最小必填集避免误拒;user_id 缺失
|
||||
// 时回退到 deployment 上下文。
|
||||
"swarm.pool_terminated": {},
|
||||
"artifact.created": {"artifact_id"},
|
||||
"timeline.updated": {"title"},
|
||||
"sk_tool.called": {"tool_name", "tool_invocation_id"},
|
||||
@@ -519,6 +523,7 @@ var agentCallbackEventCategories = map[string]string{
|
||||
"swarm.completed": "swarm_lifecycle",
|
||||
"swarm.failed": "swarm_lifecycle",
|
||||
"swarm.stopped": "swarm_lifecycle",
|
||||
"swarm.pool_terminated": "swarm_lifecycle",
|
||||
"artifact.created": "artifact",
|
||||
"timeline.updated": "timeline",
|
||||
"sk_tool.called": "sk",
|
||||
@@ -660,6 +665,26 @@ func persistAgentApprovalFromCallback(payload agentCallbackEnvelope, record agen
|
||||
return nil
|
||||
}
|
||||
|
||||
// handleSwarmPoolTerminated 处理 #60 A.5 吊销握手:收到 swarm.pool_terminated 时,
|
||||
// 删除该用户的 swarm 模型 key(KV 密文 + token)。user_id 优先取 payload,缺失时回退到
|
||||
// deployment 上下文。仅在事件首次入库(inserted)时调用,天然幂等。
|
||||
func handleSwarmPoolTerminated(payload agentCallbackEnvelope, record agentDeploymentRecord) {
|
||||
if payload.EventType != "swarm.pool_terminated" {
|
||||
return
|
||||
}
|
||||
userID := strings.TrimSpace(callbackStringValue(payload.Payload, "user_id"))
|
||||
if userID == "" {
|
||||
userID = strings.TrimSpace(record.Plan.UserContext.UserID)
|
||||
}
|
||||
uid, err := strconv.Atoi(userID)
|
||||
if err != nil || uid <= 0 {
|
||||
common.SysLog("swarm.pool_terminated: missing/invalid user_id; skip swarm model key revocation")
|
||||
return
|
||||
}
|
||||
revokeSwarmModelKey(uid)
|
||||
common.SysLog(fmt.Sprintf("swarm.pool_terminated: revoked swarm model key for user %d", uid))
|
||||
}
|
||||
|
||||
func AgentReceiveRuntimeEventCallback(c *gin.Context) {
|
||||
rawBody, err := io.ReadAll(io.LimitReader(c.Request.Body, 1<<20))
|
||||
if err != nil {
|
||||
@@ -748,6 +773,7 @@ func AgentReceiveRuntimeEventCallback(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
record, _ = applyAgentCallbackDeploymentState(payload, record)
|
||||
handleSwarmPoolTerminated(payload, record)
|
||||
recordAgentAuditEvent(agentEvent{
|
||||
EventID: "evt_" + common.GetUUID()[:12],
|
||||
Event: "callback." + payload.EventType,
|
||||
|
||||
@@ -1097,6 +1097,20 @@ func createAgentDeploymentFromPlan(c *gin.Context, plan agentOrchestrationPlan,
|
||||
plan.Metadata.RuntimeMode = agentRuntimeModeAgent
|
||||
}
|
||||
|
||||
// #60 模型 key 注入(方案A):swarm 模式下,若 billing_context 未自带 secret_ref,
|
||||
// 则为认证用户 mint/复用 per-user sk- 写入 Key Vault,把 azkv:// 引用放进
|
||||
// billing_context.secret_ref 下发给 Swarm(对接 agent_swarm PR#43)。Key Vault
|
||||
// 未配置/不可达时降级:记日志、secret_ref 留空,不阻断 create(联调前可用)。
|
||||
if plan.Metadata.RuntimeMode == agentRuntimeModeSwarm && strings.TrimSpace(plan.BillingContext.SecretRef) == "" {
|
||||
if uid := c.GetInt("id"); uid > 0 {
|
||||
if secretRef, err := provisionSwarmModelKey(uid); err != nil {
|
||||
common.SysLog("provisionSwarmModelKey (swarm create): " + err.Error())
|
||||
} else {
|
||||
plan.BillingContext.SecretRef = secretRef
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
now := agentNow()
|
||||
deploymentID := "dep_" + common.GetUUID()[:12]
|
||||
record := agentDeploymentRecord{
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"gorm.io/gorm"
|
||||
|
||||
"github.com/heicode/manager/common"
|
||||
"github.com/heicode/manager/model"
|
||||
)
|
||||
|
||||
// #60 模型 key 注入(方案A) —— 对接 agent_swarm PR#43 定死的参数:
|
||||
//
|
||||
// - A.1 粒度/命名:HM 为每个用户 mint 一把 per-user `sk-`,跨该用户所有 swarm run 复用;
|
||||
// KV 密文名 = `swarm-model-key-<user_id>`。
|
||||
// - A.2 KV value:JSON `{"openai_api_key":"sk-..."}`,对齐 Swarm
|
||||
// `orchestrator/agent_launcher._extract_model_key` 的解析字段。
|
||||
// - A.4 OPENAI_API_BASE:Swarm 部署常量(=HM 网关 /v1),HM **不**经 billing_context 下发。
|
||||
// - A.5 吊销:事件驱动 —— 用户全部 run 被 stop 时 Swarm 恰好发一次 `swarm.pool_terminated`,
|
||||
// HM 据此删 KV 密文 + 删 token(见 agent_callback.go handleSwarmPoolTerminated)。
|
||||
//
|
||||
// create 时把 putSecret 返回的 azkv:// secret_ref 放进 billing_context.secret_ref 下发;
|
||||
// Swarm 侧凭 KV 读权限(A.3,运维授权 pending)解析后注入 agent 环境。
|
||||
// 明文 `sk-` 仅经 secret_ref 服务端解析,绝不入代码/日志/事件/argv。
|
||||
|
||||
func swarmModelKeyTokenName(userID int) string {
|
||||
return fmt.Sprintf("swarm:user:%d", userID)
|
||||
}
|
||||
|
||||
func swarmModelKeySecretName(userID int) string {
|
||||
return fmt.Sprintf("swarm-model-key-%d", userID)
|
||||
}
|
||||
|
||||
// getOrMintSwarmModelToken 返回该用户长存的 swarm 模型 token 的 `sk-` bearer。
|
||||
// 已存在(未软删)则复用(A.1:跨 run 复用、稳定 key),否则新建一把隐藏、无限额度、
|
||||
// 不自然过期的系统托管 token(计费直接走 User.Quota,与 mintAgentModelToken 一致)。
|
||||
func getOrMintSwarmModelToken(userID int) (string, error) {
|
||||
if userID <= 0 || model.DB == nil {
|
||||
return "", errors.New("invalid user for swarm model key")
|
||||
}
|
||||
name := swarmModelKeyTokenName(userID)
|
||||
var tok model.Token
|
||||
err := model.DB.Where("user_id = ? AND name = ?", userID, name).First(&tok).Error
|
||||
if err == nil && strings.TrimSpace(tok.Key) != "" {
|
||||
return "sk-" + tok.Key, nil
|
||||
}
|
||||
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return "", err
|
||||
}
|
||||
rawKey, err := common.GenerateKey()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
now := common.GetTimestamp()
|
||||
tok = model.Token{
|
||||
UserId: userID,
|
||||
Name: name,
|
||||
Key: rawKey,
|
||||
Status: common.TokenStatusEnabled,
|
||||
CreatedTime: now,
|
||||
AccessedTime: now,
|
||||
ExpiredTime: -1, // never naturally
|
||||
UnlimitedQuota: true, // bills straight from User.Quota
|
||||
HideFromUserUI: true, // 系统托管令牌,不在用户令牌列表展示
|
||||
}
|
||||
if err := tok.Insert(); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return "sk-" + rawKey, nil
|
||||
}
|
||||
|
||||
// provisionSwarmModelKey 确保用户的 swarm 模型 key 已就位(token + KV 密文),返回 azkv:// secret_ref。
|
||||
// 仅在 Azure Key Vault 已配置时可用;未配置/不可达返回 error,由调用方按需降级
|
||||
// (联调前 secret_ref 可留空,不阻断 swarm create)。
|
||||
func provisionSwarmModelKey(userID int) (string, error) {
|
||||
store, err := newSecretStoreClientFromEnv()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
bearer, err := getOrMintSwarmModelToken(userID)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
// A.2:KV value = JSON {"openai_api_key":"sk-..."}。putSecret 幂等(同名新建版本)。
|
||||
secretRef, err := store.putSecret(swarmModelKeySecretName(userID), map[string]any{
|
||||
"openai_api_key": bearer,
|
||||
})
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return secretRef, nil
|
||||
}
|
||||
|
||||
// revokeSwarmModelKey 吊销用户的 swarm 模型 key:删 KV 密文 + 软删 token(均尽力而为)。
|
||||
// 由 swarm.pool_terminated 回调触发(A.5)。token 软删后,下次 swarm create 会重新 mint 一把新 key。
|
||||
func revokeSwarmModelKey(userID int) {
|
||||
if userID <= 0 {
|
||||
return
|
||||
}
|
||||
if store, err := newSecretStoreClientFromEnv(); err == nil {
|
||||
if err := store.deleteSecret(store.secretRef(swarmModelKeySecretName(userID))); err != nil {
|
||||
common.SysLog("revokeSwarmModelKey delete KV secret: " + err.Error())
|
||||
}
|
||||
}
|
||||
if model.DB != nil {
|
||||
if err := model.DB.Where("user_id = ? AND name = ?", userID, swarmModelKeyTokenName(userID)).Delete(&model.Token{}).Error; err != nil {
|
||||
common.SysLog("revokeSwarmModelKey delete token: " + err.Error())
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,139 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/glebarez/sqlite"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/gorm"
|
||||
|
||||
"github.com/heicode/manager/common"
|
||||
"github.com/heicode/manager/model"
|
||||
)
|
||||
|
||||
func setupSwarmModelKeyTestDB(t *testing.T) {
|
||||
t.Helper()
|
||||
common.UsingSQLite = true
|
||||
common.UsingMySQL = false
|
||||
common.UsingPostgreSQL = false
|
||||
common.RedisEnabled = false
|
||||
|
||||
dsn := fmt.Sprintf("file:%s?mode=memory&cache=shared", strings.ReplaceAll(t.Name(), "/", "_"))
|
||||
db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{})
|
||||
require.NoError(t, err)
|
||||
model.DB = db
|
||||
require.NoError(t, db.AutoMigrate(&model.Token{}))
|
||||
t.Cleanup(func() {
|
||||
if sqlDB, err := db.DB(); err == nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
model.DB = nil
|
||||
})
|
||||
}
|
||||
|
||||
func TestSwarmModelKeyNames(t *testing.T) {
|
||||
require.Equal(t, "swarm:user:42", swarmModelKeyTokenName(42))
|
||||
require.Equal(t, "swarm-model-key-42", swarmModelKeySecretName(42))
|
||||
}
|
||||
|
||||
func TestGetOrMintSwarmModelTokenReusesPerUser(t *testing.T) {
|
||||
setupSwarmModelKeyTestDB(t)
|
||||
|
||||
// A.1: first call mints, second call reuses the SAME sk- across runs.
|
||||
first, err := getOrMintSwarmModelToken(7001)
|
||||
require.NoError(t, err)
|
||||
require.True(t, strings.HasPrefix(first, "sk-"))
|
||||
second, err := getOrMintSwarmModelToken(7001)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, first, second, "per-user swarm key must be reused across runs")
|
||||
|
||||
// only one token persisted for the user, and it is hidden + system-managed.
|
||||
var toks []model.Token
|
||||
require.NoError(t, model.DB.Where("user_id = ?", 7001).Find(&toks).Error)
|
||||
require.Len(t, toks, 1)
|
||||
require.Equal(t, swarmModelKeyTokenName(7001), toks[0].Name)
|
||||
require.True(t, toks[0].HideFromUserUI)
|
||||
require.True(t, toks[0].UnlimitedQuota)
|
||||
require.EqualValues(t, -1, toks[0].ExpiredTime)
|
||||
|
||||
// distinct users get distinct keys.
|
||||
other, err := getOrMintSwarmModelToken(7002)
|
||||
require.NoError(t, err)
|
||||
require.NotEqual(t, first, other)
|
||||
|
||||
// invalid user rejected.
|
||||
_, err = getOrMintSwarmModelToken(0)
|
||||
require.Error(t, err)
|
||||
}
|
||||
|
||||
func TestRevokeSwarmModelKeyReMintsFresh(t *testing.T) {
|
||||
setupSwarmModelKeyTestDB(t)
|
||||
|
||||
first, err := getOrMintSwarmModelToken(7003)
|
||||
require.NoError(t, err)
|
||||
|
||||
// revoke soft-deletes the token (KV not configured -> delete logged, non-fatal).
|
||||
revokeSwarmModelKey(7003)
|
||||
var live []model.Token
|
||||
require.NoError(t, model.DB.Where("user_id = ?", 7003).Find(&live).Error)
|
||||
require.Len(t, live, 0, "revoked swarm token must be soft-deleted")
|
||||
|
||||
// next provision mints a brand-new key (A.5: stop -> revoke -> new key on re-run).
|
||||
second, err := getOrMintSwarmModelToken(7003)
|
||||
require.NoError(t, err)
|
||||
require.NotEqual(t, first, second)
|
||||
}
|
||||
|
||||
func TestHandleSwarmPoolTerminatedRevokesByPayloadUserID(t *testing.T) {
|
||||
setupSwarmModelKeyTestDB(t)
|
||||
_, err := getOrMintSwarmModelToken(7004)
|
||||
require.NoError(t, err)
|
||||
|
||||
payload := agentCallbackEnvelope{
|
||||
EventType: "swarm.pool_terminated",
|
||||
Payload: map[string]any{"user_id": "7004", "secret_ref": "azkv://heicode-kv.vault.azure.net/secrets/swarm-model-key-7004"},
|
||||
}
|
||||
handleSwarmPoolTerminated(payload, agentDeploymentRecord{})
|
||||
|
||||
var live []model.Token
|
||||
require.NoError(t, model.DB.Where("user_id = ?", 7004).Find(&live).Error)
|
||||
require.Len(t, live, 0)
|
||||
}
|
||||
|
||||
func TestHandleSwarmPoolTerminatedFallsBackToDeploymentUser(t *testing.T) {
|
||||
setupSwarmModelKeyTestDB(t)
|
||||
_, err := getOrMintSwarmModelToken(7005)
|
||||
require.NoError(t, err)
|
||||
|
||||
// no user_id in payload -> fall back to deployment user context.
|
||||
payload := agentCallbackEnvelope{EventType: "swarm.pool_terminated", Payload: map[string]any{}}
|
||||
record := agentDeploymentRecord{}
|
||||
record.Plan.UserContext.UserID = "7005"
|
||||
handleSwarmPoolTerminated(payload, record)
|
||||
|
||||
var live []model.Token
|
||||
require.NoError(t, model.DB.Where("user_id = ?", 7005).Find(&live).Error)
|
||||
require.Len(t, live, 0)
|
||||
}
|
||||
|
||||
func TestHandleSwarmPoolTerminatedIgnoresOtherEvents(t *testing.T) {
|
||||
setupSwarmModelKeyTestDB(t)
|
||||
_, err := getOrMintSwarmModelToken(7006)
|
||||
require.NoError(t, err)
|
||||
|
||||
// wrong event type -> no revocation.
|
||||
handleSwarmPoolTerminated(agentCallbackEnvelope{EventType: "swarm.stopped", Payload: map[string]any{"user_id": "7006"}}, agentDeploymentRecord{})
|
||||
var live []model.Token
|
||||
require.NoError(t, model.DB.Where("user_id = ?", 7006).Find(&live).Error)
|
||||
require.Len(t, live, 1)
|
||||
}
|
||||
|
||||
func TestSwarmPoolTerminatedEventRegistered(t *testing.T) {
|
||||
// control-plane lifecycle event must be a known callback type so the schema
|
||||
// validator accepts it and routes it to the swarm_lifecycle category.
|
||||
_, hasFields := agentCallbackEventRequiredFields["swarm.pool_terminated"]
|
||||
require.True(t, hasFields)
|
||||
require.Equal(t, "swarm_lifecycle", agentCallbackEventCategories["swarm.pool_terminated"])
|
||||
}
|
||||
Reference in New Issue
Block a user