From c494e9576e4618deeca58521fc97f933ed2743bc Mon Sep 17 00:00:00 2001 From: chenchen Date: Thu, 11 Jun 2026 18:30:32 +0800 Subject: [PATCH] =?UTF-8?q?feat(swarm):=20#60=20=E6=A8=A1=E5=9E=8B=20key?= =?UTF-8?q?=20=E6=B3=A8=E5=85=A5(=E6=96=B9=E6=A1=88A)+=20pool=5Fterminated?= =?UTF-8?q?=20=E5=90=8A=E9=94=80=E6=8F=A1=E6=89=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 对接 agent_swarm PR#43 定死的参数,实现 HM 侧的 per-user 模型 key 注入: - A.1 粒度/命名:为认证用户 mint 一把 per-user sk-(getOrMintSwarmModelToken), 跨该用户所有 swarm run 复用;KV 密文名 swarm-model-key-。 - 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) --- heicode/controller/agent_callback.go | 26 ++++ heicode/controller/agent_control_plane.go | 14 +++ heicode/controller/swarm_model_key.go | 112 +++++++++++++++++ heicode/controller/swarm_model_key_test.go | 139 +++++++++++++++++++++ 4 files changed, 291 insertions(+) create mode 100644 heicode/controller/swarm_model_key.go create mode 100644 heicode/controller/swarm_model_key_test.go diff --git a/heicode/controller/agent_callback.go b/heicode/controller/agent_callback.go index b476502..12d32be 100644 --- a/heicode/controller/agent_callback.go +++ b/heicode/controller/agent_callback.go @@ -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, diff --git a/heicode/controller/agent_control_plane.go b/heicode/controller/agent_control_plane.go index c1942e4..109a603 100644 --- a/heicode/controller/agent_control_plane.go +++ b/heicode/controller/agent_control_plane.go @@ -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{ diff --git a/heicode/controller/swarm_model_key.go b/heicode/controller/swarm_model_key.go new file mode 100644 index 0000000..463fcec --- /dev/null +++ b/heicode/controller/swarm_model_key.go @@ -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-`。 +// - 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()) + } + } +} diff --git a/heicode/controller/swarm_model_key_test.go b/heicode/controller/swarm_model_key_test.go new file mode 100644 index 0000000..339e1ef --- /dev/null +++ b/heicode/controller/swarm_model_key_test.go @@ -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"]) +}