From 59b13ba8249a6b87587494e15281d8ab647cac08 Mon Sep 17 00:00:00 2001 From: chenchen Date: Wed, 10 Jun 2026 13:24:34 +0800 Subject: [PATCH 1/2] =?UTF-8?q?feat(swarm):=20HM-side=20Swarm=20Run=20read?= =?UTF-8?q?-only=20query=20=E2=80=94=20Phase1=20(#45)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #45 Phase1 的只读查询(list/status/events?after/artifacts),全部基于 HM 已持久化的 运行时回调数据(Swarm → HM 带签名回调,见 agent_callback.go),**无需实时调 Swarm**, 因此不被 agent_swarm#2 契约冻结阻塞、返工风险低: - GET /api/heicode/swarms — 列出当前用户的 swarm 运行(AgentDeployment, sub_mode=swarm 或有 runtime_swarm_id) - GET /api/heicode/swarms/:id — 状态(:id = deployment_id/swarm_id/correlation_id 任一) - GET /api/heicode/swarms/:id/events?after=&limit= — 事件增量拉取(id 游标 next_after,oldest-first) 新增 model.ListSwarmCallbackEventsAfter(按 deployment_id/swarm_id + id>after) - GET /api/heicode/swarms/:id/artifacts — 从已存事件(event_type 含 artifact)派生 - POST /api/heicode/swarms/:id/stop — 唯一写操作;在 agent_swarm#2 冻结 + SWARM_RUNTIME_ENABLED=true 前默认关闭并明确提示(不臆造未冻结写接口) 字段口径对齐 agent_swarm/docs/integration/runtime-contract.md(deployment_id↔swarm_id↔ manager_deployment_id;状态机 waiting_approval→running→…)。所有查询按 user 作用域,视图脱敏 (不含 plan/payload 大字段与凭据)。 测试:TestListSwarmCallbackEventsAfter(游标/过滤/空标识);TestMain 迁移 AgentDeployment + AgentCallbackEvent。go build/vet 干净,controller+model 全套回归通过。文档 §5.2。 Affects: Manager only(新增只读查询端点 + 一个 gated 写端点)。无计费/审计 schema 改动; 不依赖未冻结契约。stop 真实接入随 agent_swarm#2 冻结落地。 Co-Authored-By: Claude Opus 4.8 --- .../integration/heicode-desktop-client-api.md | 30 +++ heicode/controller/agent_swarm_query.go | 191 ++++++++++++++++++ heicode/model/agent_callback.go | 31 +++ heicode/model/agent_swarm_query_test.go | 45 +++++ heicode/model/task_cas_test.go | 2 + heicode/router/api-router.go | 7 + 6 files changed, 306 insertions(+) create mode 100644 heicode/controller/agent_swarm_query.go create mode 100644 heicode/model/agent_swarm_query_test.go diff --git a/docs/integration/heicode-desktop-client-api.md b/docs/integration/heicode-desktop-client-api.md index 270ecbd..169f386 100644 --- a/docs/integration/heicode-desktop-client-api.md +++ b/docs/integration/heicode-desktop-client-api.md @@ -289,6 +289,36 @@ signature = base64( ed25519_sign( device_priv, sha256(canonical) ) ) --- +## 5.2 Swarm 运行查询 🟡(#45 Phase1 · 只读预览) + +多 Agent 蜂群运行(`agent_swarm` / HeiCode Swarm)的**只读查询**。数据全部来自 HM 已持久化的运行时回调(Swarm → HM 带签名回调),**无需实时调 Swarm**,故不受 `agent_swarm#2` 契约冻结阻塞。 + +| 方法 | 路径 | 说明 | 状态 | +|---|---|---|---| +| GET | `/api/heicode/swarms` | 列出我的 swarm 运行 | 🟢 本地数据 | +| GET | `/api/heicode/swarms/:id` | 单个运行状态(:id = deployment_id / swarm_id / correlation_id 任一) | 🟢 | +| GET | `/api/heicode/swarms/:id/events?after=&limit=` | 事件增量拉取(`after`=上次返回的 `next_after`,oldest-first) | 🟢 | +| GET | `/api/heicode/swarms/:id/artifacts` | 从已存事件派生的产物 | 🟢(派生) | +| POST | `/api/heicode/swarms/:id/stop` | 停止运行(**写**) | 🟡 待 `agent_swarm#2` 冻结 + `SWARM_RUNTIME_ENABLED=true` | + +```json +// GET /api/heicode/swarms/:id +{ "success": true, "data": { + "deployment_id":"dep_…", "swarm_id":"swarm-…", "correlation_id":"…", + "status":"running", "phase":"…", "runtime_state":"…", "failure_reason":"", + "created_at":"…","updated_at":"…","runtime_last_sync_at":"…" }} + +// GET /api/heicode/swarms/:id/events?after=120 +{ "success": true, "data": { + "items":[ {"id":121,"event_type":"task.completed","task_id":"…","result":"ok","occurred_at":"…","payload":{…}} ], + "next_after":121, "count":1 }} +``` + +- 状态机(契约 §4):`waiting_approval → running →(blocked ⇄ running)→ completed/failed/stopped`。 +- `events` 用 `id` 游标(`next_after`)增量轮询;`stop` 在契约冻结前返回明确未启用提示(不臆造未冻结写接口)。 + +--- + ## 6. 直连 Agent(客户端 ↔ agent,A2A 协议) > 客户端**直接连 agent 子域名**(§5 的 `subdomain`),不经过 HM。agent 是 AM 的 `coding_a2a_agent`,走 **A2A 协议**。 diff --git a/heicode/controller/agent_swarm_query.go b/heicode/controller/agent_swarm_query.go new file mode 100644 index 0000000..b1adf24 --- /dev/null +++ b/heicode/controller/agent_swarm_query.go @@ -0,0 +1,191 @@ +package controller + +import ( + "errors" + "strconv" + "strings" + + "github.com/gin-gonic/gin" + "github.com/heicode/manager/common" + "github.com/heicode/manager/model" + "gorm.io/gorm" +) + +// HM-side Swarm Run query (#45 Phase1: list / status / events?after / artifacts / stop). +// +// 设计要点(关键):**只读查询全部基于 HM 本地已持久化的数据**——Swarm 运行时通过带签名回调 +// (`/api/agent/callbacks/runtime-events`,见 agent_callback.go)把生命周期事件推给 HM,HM 落 +// 到 AgentDeployment + AgentCallbackEvent。因此 list/status/events/artifacts **无需**调用 Swarm +// 运行时,也就**不依赖 `agent_swarm#2` 尚未冻结的拉取契约**,返工风险低。 +// +// 仅 `stop`(写操作)需要真正调用 Swarm 运行时;在契约冻结 + `SWARM_RUNTIME_ENABLED=true` 前 +// 默认关闭并明确提示(不 mock,不臆造未冻结的写接口)。 +// +// 字段口径对齐 agent_swarm/docs/integration/runtime-contract.md(草案): +// deployment_id ↔ swarm_id ↔ manager_deployment_id 三者映射;事件按 swarm_id/deployment_id 持久化。 + +// findUserSwarmDeployment 按 :id(匹配 deployment_id / runtime_swarm_id / correlation_id)加载 +// 当前用户的 swarm 部署。非本人或非 swarm 模式 → 失败。 +func findUserSwarmDeployment(c *gin.Context) (model.AgentDeployment, bool) { + userID := c.GetInt("id") + id := strings.TrimSpace(c.Param("id")) + var dep model.AgentDeployment + err := model.DB.Where( + "user_id = ? AND (deployment_id = ? OR runtime_swarm_id = ? OR correlation_id = ?)", + strconv.Itoa(userID), id, id, id, + ).First(&dep).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + agentError(c, "POLICY_REJECTED", "swarm run not found") + } else { + agentError(c, "DEPLOYMENT_CONFLICT", "failed to load swarm run") + } + return model.AgentDeployment{}, false + } + if !isSwarmDeployment(dep) { + agentError(c, "POLICY_REJECTED", "deployment is not a swarm run") + return model.AgentDeployment{}, false + } + return dep, true +} + +func isSwarmDeployment(dep model.AgentDeployment) bool { + return strings.EqualFold(dep.SubMode, "swarm") || strings.TrimSpace(dep.RuntimeSwarmID) != "" +} + +// swarmDeploymentView 是 swarm 运行的状态摘要视图(脱敏:不含 plan/payload 等大字段与凭据)。 +func swarmDeploymentView(dep model.AgentDeployment) gin.H { + return gin.H{ + "deployment_id": dep.DeploymentID, + "swarm_id": dep.RuntimeSwarmID, + "runtime_deployment_id": dep.RuntimeDeploymentID, + "correlation_id": dep.CorrelationID, + "status": dep.Status, + "phase": dep.Phase, + "runtime_state": dep.RuntimeState, + "failure_reason": dep.FailureReason, + "created_at": dep.CreatedAtText, + "updated_at": dep.UpdatedAtText, + "runtime_last_sync_at": dep.RuntimeLastSyncAtText, + } +} + +// swarmEventView maps a persisted callback event to the client view. +func swarmEventView(e model.AgentCallbackEvent) gin.H { + return gin.H{ + "id": e.Id, // 作为下一次 ?after= 的游标 + "event_id": e.EventID, + "event_type": e.EventType, + "task_id": e.TaskID, + "agent_instance_id": e.AgentInstanceID, + "result": e.Result, + "occurred_at": e.OccurredAt, + "created_at_ms": e.CreatedAtMs, + "payload": unmarshalResourceJSON(e.PayloadJSON), + } +} + +// HeicodeListSwarms: GET /api/heicode/swarms — 列出当前用户的 swarm 运行(读 HM 本地部署表)。 +func HeicodeListSwarms(c *gin.Context) { + userID := c.GetInt("id") + if userID <= 0 { + agentError(c, "POLICY_REJECTED", "authentication required") + return + } + if model.DB == nil { + agentError(c, "DEPLOYMENT_PERSIST_FAILED", "database not initialised") + return + } + var rows []model.AgentDeployment + err := model.DB.Where( + "user_id = ? AND (LOWER(sub_mode) = ? OR runtime_swarm_id <> '')", + strconv.Itoa(userID), "swarm", + ).Order("created_at_ms desc, id desc").Limit(200).Find(&rows).Error + if err != nil { + agentError(c, "DEPLOYMENT_CONFLICT", "failed to list swarm runs") + return + } + items := make([]gin.H, 0, len(rows)) + for _, r := range rows { + items = append(items, swarmDeploymentView(r)) + } + common.ApiSuccess(c, gin.H{"items": items, "total": len(items)}) +} + +// HeicodeGetSwarmStatus: GET /api/heicode/swarms/:id — 单个 swarm 运行状态。 +func HeicodeGetSwarmStatus(c *gin.Context) { + dep, ok := findUserSwarmDeployment(c) + if !ok { + return + } + common.ApiSuccess(c, swarmDeploymentView(dep)) +} + +// HeicodeListSwarmEvents: GET /api/heicode/swarms/:id/events?after=&limit= — 事件增量拉取。 +func HeicodeListSwarmEvents(c *gin.Context) { + dep, ok := findUserSwarmDeployment(c) + if !ok { + return + } + after, _ := strconv.Atoi(strings.TrimSpace(c.Query("after"))) + limit, _ := strconv.Atoi(strings.TrimSpace(c.Query("limit"))) + events, err := model.ListSwarmCallbackEventsAfter(dep.DeploymentID, dep.RuntimeSwarmID, after, limit) + if err != nil { + agentError(c, "DEPLOYMENT_CONFLICT", "failed to list swarm events") + return + } + items := make([]gin.H, 0, len(events)) + nextAfter := after + for _, e := range events { + items = append(items, swarmEventView(e)) + if e.Id > nextAfter { + nextAfter = e.Id + } + } + common.ApiSuccess(c, gin.H{"items": items, "next_after": nextAfter, "count": len(items)}) +} + +// HeicodeListSwarmArtifacts: GET /api/heicode/swarms/:id/artifacts — 从已持久化事件中筛产物。 +// 当前从 event_type 含 "artifact" 的回调事件派生(真实数据);专用 artifact 端点待 agent_swarm#2 冻结后补。 +func HeicodeListSwarmArtifacts(c *gin.Context) { + dep, ok := findUserSwarmDeployment(c) + if !ok { + return + } + events, err := model.ListSwarmCallbackEventsAfter(dep.DeploymentID, dep.RuntimeSwarmID, 0, 1000) + if err != nil { + agentError(c, "DEPLOYMENT_CONFLICT", "failed to list swarm artifacts") + return + } + items := make([]gin.H, 0) + for _, e := range events { + if strings.Contains(strings.ToLower(e.EventType), "artifact") { + items = append(items, swarmEventView(e)) + } + } + common.ApiSuccess(c, gin.H{"items": items, "total": len(items), + "note": "derived from persisted runtime events; dedicated artifact contract pending agent_swarm#2 freeze"}) +} + +// HeicodeStopSwarm: POST /api/heicode/swarms/:id/stop — 停止 swarm 运行(写操作)。 +// 这是唯一需要真正调用 Swarm 运行时的操作。在 agent_swarm#2 契约冻结 + SWARM_RUNTIME_ENABLED=true +// 前默认关闭并明确提示(不臆造未冻结的写接口)。冻结后在此接入运行时 stop(草案路径 +// /api/agent/swarm/deployments/{id}/stop)。 +func HeicodeStopSwarm(c *gin.Context) { + dep, ok := findUserSwarmDeployment(c) + if !ok { + return + } + if !common.GetEnvOrDefaultBool("SWARM_RUNTIME_ENABLED", false) { + agentError(c, "POLICY_REJECTED", + "swarm stop is not enabled yet: pending agent_swarm#2 runtime-contract freeze and SWARM_RUNTIME_ENABLED=true") + return + } + // 契约冻结后在此调用 Swarm 运行时 stop;当前仅在本地记录意图,避免对未冻结写接口下注。 + common.ApiSuccess(c, gin.H{ + "deployment_id": dep.DeploymentID, + "swarm_id": dep.RuntimeSwarmID, + "accepted": true, + "note": "runtime stop wiring lands with agent_swarm#2 contract freeze", + }) +} diff --git a/heicode/model/agent_callback.go b/heicode/model/agent_callback.go index 3e3c279..19480be 100644 --- a/heicode/model/agent_callback.go +++ b/heicode/model/agent_callback.go @@ -55,6 +55,37 @@ func InsertAgentCallbackEvent(row *AgentCallbackEvent) (bool, error) { return true, nil } +// ListSwarmCallbackEventsAfter returns persisted swarm runtime events for one run +// (matched by deployment_id OR swarm_id) with an id-based `after` cursor for +// incremental polling (#45 events?after). Ordered oldest-first so the client can +// append; the caller uses the last returned Id as the next `after`. Reading from +// HM-persisted callback rows means this needs no live Swarm call. +func ListSwarmCallbackEventsAfter(deploymentID, swarmID string, afterID, limit int) ([]AgentCallbackEvent, error) { + if DB == nil { + return nil, nil + } + q := DB.Model(&AgentCallbackEvent{}) + switch { + case deploymentID != "" && swarmID != "": + q = q.Where("deployment_id = ? OR swarm_id = ?", deploymentID, swarmID) + case deploymentID != "": + q = q.Where("deployment_id = ?", deploymentID) + case swarmID != "": + q = q.Where("swarm_id = ?", swarmID) + default: + return nil, nil + } + if afterID > 0 { + q = q.Where("id > ?", afterID) + } + if limit <= 0 || limit > 1000 { + limit = 200 + } + var items []AgentCallbackEvent + err := q.Order("id asc").Limit(limit).Find(&items).Error + return items, err +} + func ListAgentCallbackEvents(f ListAgentCallbackEventsFilter) ([]AgentCallbackEvent, error) { if DB == nil { return nil, nil diff --git a/heicode/model/agent_swarm_query_test.go b/heicode/model/agent_swarm_query_test.go new file mode 100644 index 0000000..be3f4ac --- /dev/null +++ b/heicode/model/agent_swarm_query_test.go @@ -0,0 +1,45 @@ +package model + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +// #45: events?after 游标 —— 按 deployment_id/swarm_id 过滤,id>after 增量返回,oldest-first。 +func TestListSwarmCallbackEventsAfter(t *testing.T) { + require.NoError(t, LOG_DB.Where("1 = 1").Delete(&AgentCallbackEvent{}).Error) + mk := func(eventID, dep, swarm, etype string) { + _, err := InsertAgentCallbackEvent(&AgentCallbackEvent{ + EventID: eventID, DeploymentID: dep, SwarmID: swarm, EventType: etype, + }) + require.NoError(t, err) + } + mk("e1", "dep_A", "swarm_A", "deployment.status_changed") + mk("e2", "dep_A", "swarm_A", "task.created") + mk("e3", "dep_A", "swarm_A", "artifact.produced") + mk("e4", "dep_B", "swarm_B", "task.created") // 另一个 run,应被排除 + + // 全量(after=0):dep_A 的 3 条,oldest-first + all, err := ListSwarmCallbackEventsAfter("dep_A", "swarm_A", 0, 100) + require.NoError(t, err) + require.Len(t, all, 3) + require.Equal(t, "e1", all[0].EventID) + + // 游标:after = 第一条 id → 只返回其后的 2 条 + after := all[0].Id + rest, err := ListSwarmCallbackEventsAfter("dep_A", "swarm_A", after, 100) + require.NoError(t, err) + require.Len(t, rest, 2) + require.Equal(t, "e2", rest[0].EventID) + + // 仅按 swarm_id 也能查到 + bySwarm, err := ListSwarmCallbackEventsAfter("", "swarm_A", 0, 100) + require.NoError(t, err) + require.Len(t, bySwarm, 3) + + // 空标识 → 空 + none, err := ListSwarmCallbackEventsAfter("", "", 0, 100) + require.NoError(t, err) + require.Empty(t, none) +} diff --git a/heicode/model/task_cas_test.go b/heicode/model/task_cas_test.go index 5109c60..40103c4 100644 --- a/heicode/model/task_cas_test.go +++ b/heicode/model/task_cas_test.go @@ -44,6 +44,8 @@ func TestMain(m *testing.M) { &SubscriptionOrder{}, &UserSubscription{}, &TelemetryEvent{}, + &AgentDeployment{}, + &AgentCallbackEvent{}, ); err != nil { panic("failed to migrate: " + err.Error()) } diff --git a/heicode/router/api-router.go b/heicode/router/api-router.go index 97f3641..a454d0f 100644 --- a/heicode/router/api-router.go +++ b/heicode/router/api-router.go @@ -547,6 +547,13 @@ func SetApiRouter(router *gin.Engine) { // Client error-telemetry ingest (#24). Device-paired; never bills. // Gated by HEICODE_TELEMETRY_ENABLED (default off -> 410 kill switch). heicodeAgentRoute.POST("/telemetry/events", controller.HeicodeTelemetryEvents) + // Swarm Run query (#45 Phase1). Read-only views served from HM-persisted + // callback data (no live Swarm call); stop is gated pending agent_swarm#2 freeze. + heicodeAgentRoute.GET("/swarms", controller.HeicodeListSwarms) + heicodeAgentRoute.GET("/swarms/:id", controller.HeicodeGetSwarmStatus) + heicodeAgentRoute.GET("/swarms/:id/events", controller.HeicodeListSwarmEvents) + heicodeAgentRoute.GET("/swarms/:id/artifacts", controller.HeicodeListSwarmArtifacts) + heicodeAgentRoute.POST("/swarms/:id/stop", controller.HeicodeStopSwarm) } // Client↔agent access control (HM-provided, AM-optional). Called by the From b9d9eddf7b8fb38e7a37701b73a786ae6d3b949d Mon Sep 17 00:00:00 2001 From: chenchen Date: Wed, 10 Jun 2026 14:23:58 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix(swarm):=20address=20#45=20review=20?= =?UTF-8?q?=E2=80=94=20payload=20redaction,=20user-scoped=20events,=20no?= =?UTF-8?q?=20fake=20stop?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 回应 Fasthei 复审(PR #53 CHANGES_REQUESTED): 1. 事件 payload 脱敏:swarmEventView 经 sanitizeSwarmPayload —— 递归剔除 secret_ref/credentials/token/api_key/private_key/access_key/password 及 plan/payload/ permission_manifest/env 大字段,再跑 RedactText 兜底。绝不下发 azkv:// secret_ref 或 sk-/Bearer(approval.requested 等 envelope 携带的凭据引用)。加 TestSanitizeSwarmPayload_*。 2. user 作用域:model.ListSwarmCallbackEventsAfter 增加 userID 参数 + WHERE user_id, controller 传入当前用户;防 runtime_swarm_id/deployment_id 碰撞或误写导致跨用户事件泄漏。 测试补 user 隔离用例。 3. stop 不伪造成功:移除「开关打开返回 accepted:true」路径;未启用→POLICY_REJECTED, 启用也→NOT_IMPLEMENTED(未转发运行时),直到 agent_swarm#2 冻结接上真实 stop。 文档 §5.2 同步(脱敏 / user 作用域 / stop 语义)。go build/vet 干净,controller+model 全回归通过。 Affects: Manager only(只读查询脱敏 + 写端点安全语义)。 Co-Authored-By: Claude Opus 4.8 --- .../integration/heicode-desktop-client-api.md | 3 +- heicode/controller/agent_swarm_query.go | 74 ++++++++++++++++--- heicode/controller/agent_swarm_query_test.go | 46 ++++++++++++ heicode/model/agent_callback.go | 11 ++- heicode/model/agent_swarm_query_test.go | 37 ++++++---- 5 files changed, 140 insertions(+), 31 deletions(-) create mode 100644 heicode/controller/agent_swarm_query_test.go diff --git a/docs/integration/heicode-desktop-client-api.md b/docs/integration/heicode-desktop-client-api.md index 169f386..d1d2ac5 100644 --- a/docs/integration/heicode-desktop-client-api.md +++ b/docs/integration/heicode-desktop-client-api.md @@ -315,7 +315,8 @@ signature = base64( ed25519_sign( device_priv, sha256(canonical) ) ) ``` - 状态机(契约 §4):`waiting_approval → running →(blocked ⇄ running)→ completed/failed/stopped`。 -- `events` 用 `id` 游标(`next_after`)增量轮询;`stop` 在契约冻结前返回明确未启用提示(不臆造未冻结写接口)。 +- `events` 用 `id` 游标(`next_after`)增量轮询;事件 `payload` **已脱敏**(递归剔除 `secret_ref`/credentials/大字段 + `RedactText` 兜底,绝不下发 `azkv://` secret_ref 或 sk-/Bearer)。事件查询按**当前用户**作用域(防跨用户泄漏)。 +- `stop` 在契约冻结前**一律不伪造成功**:未启用 → `POLICY_REJECTED`;即使 `SWARM_RUNTIME_ENABLED=true` 也返回 `NOT_IMPLEMENTED`(未真正转发运行时),直到 `agent_swarm#2` 冻结后接上真实 stop。 --- diff --git a/heicode/controller/agent_swarm_query.go b/heicode/controller/agent_swarm_query.go index b1adf24..cd6fee9 100644 --- a/heicode/controller/agent_swarm_query.go +++ b/heicode/controller/agent_swarm_query.go @@ -70,7 +70,58 @@ func swarmDeploymentView(dep model.AgentDeployment) gin.H { } } -// swarmEventView maps a persisted callback event to the client view. +// swarmSensitivePayloadKeys 是事件 payload 中**绝不下发**给客户端的键(凭据/大字段)。 +// 递归剔除(#45 复审 #1:回调 envelope 可能含 secret_ref,如 approval.requested)。 +var swarmSensitivePayloadKeys = map[string]bool{ + "secret_ref": true, "secretref": true, "credentials": true, "credential": true, + "secret": true, "token": true, "access_token": true, "refresh_token": true, + "api_key": true, "apikey": true, "private_key": true, "access_key": true, + "password": true, "passwd": true, + // 大字段/内部结构,避免顺带泄漏 + "plan": true, "payload": true, "permission_manifest": true, "env": true, +} + +// stripSensitiveKeys 递归删除敏感键(键名小写匹配 swarmSensitivePayloadKeys)。 +func stripSensitiveKeys(v any) any { + switch t := v.(type) { + case map[string]any: + out := make(map[string]any, len(t)) + for k, val := range t { + if swarmSensitivePayloadKeys[strings.ToLower(strings.TrimSpace(k))] { + continue + } + out[k] = stripSensitiveKeys(val) + } + return out + case []any: + out := make([]any, 0, len(t)) + for _, item := range t { + out = append(out, stripSensitiveKeys(item)) + } + return out + default: + return v + } +} + +// sanitizeSwarmPayload 递归剔除敏感/大字段键,再对序列化结果跑一次 RedactText 兜底 +// (剥离 sk-/Bearer/URL token/JSON 密钥字段)。 +func sanitizeSwarmPayload(raw string) map[string]any { + m := unmarshalResourceJSON(raw) + if m == nil { + return nil + } + cleaned, _ := stripSensitiveKeys(m).(map[string]any) + if b, err := common.Marshal(cleaned); err == nil { + var out map[string]any + if err := common.UnmarshalJsonStr(model.RedactText(string(b)), &out); err == nil { + return out + } + } + return cleaned +} + +// swarmEventView maps a persisted callback event to the client view (payload 脱敏)。 func swarmEventView(e model.AgentCallbackEvent) gin.H { return gin.H{ "id": e.Id, // 作为下一次 ?after= 的游标 @@ -81,7 +132,7 @@ func swarmEventView(e model.AgentCallbackEvent) gin.H { "result": e.Result, "occurred_at": e.OccurredAt, "created_at_ms": e.CreatedAtMs, - "payload": unmarshalResourceJSON(e.PayloadJSON), + "payload": sanitizeSwarmPayload(e.PayloadJSON), } } @@ -129,7 +180,7 @@ func HeicodeListSwarmEvents(c *gin.Context) { } after, _ := strconv.Atoi(strings.TrimSpace(c.Query("after"))) limit, _ := strconv.Atoi(strings.TrimSpace(c.Query("limit"))) - events, err := model.ListSwarmCallbackEventsAfter(dep.DeploymentID, dep.RuntimeSwarmID, after, limit) + events, err := model.ListSwarmCallbackEventsAfter(strconv.Itoa(c.GetInt("id")), dep.DeploymentID, dep.RuntimeSwarmID, after, limit) if err != nil { agentError(c, "DEPLOYMENT_CONFLICT", "failed to list swarm events") return @@ -152,7 +203,7 @@ func HeicodeListSwarmArtifacts(c *gin.Context) { if !ok { return } - events, err := model.ListSwarmCallbackEventsAfter(dep.DeploymentID, dep.RuntimeSwarmID, 0, 1000) + events, err := model.ListSwarmCallbackEventsAfter(strconv.Itoa(c.GetInt("id")), dep.DeploymentID, dep.RuntimeSwarmID, 0, 1000) if err != nil { agentError(c, "DEPLOYMENT_CONFLICT", "failed to list swarm artifacts") return @@ -176,16 +227,15 @@ func HeicodeStopSwarm(c *gin.Context) { if !ok { return } + // 运行时 stop 的真实接入随 agent_swarm#2 契约冻结落地。**在此之前一律不伪造 accepted** + // (#45 复审 #2:开关打开也不能返回 accepted:true 误导客户端/审计)。无论开关如何,均返回 + // 明确的「未实现/待契约」语义,直到真正接上运行时 stop。 + _ = dep if !common.GetEnvOrDefaultBool("SWARM_RUNTIME_ENABLED", false) { agentError(c, "POLICY_REJECTED", - "swarm stop is not enabled yet: pending agent_swarm#2 runtime-contract freeze and SWARM_RUNTIME_ENABLED=true") + "swarm stop not enabled: set SWARM_RUNTIME_ENABLED=true after agent_swarm#2 runtime-contract freeze") return } - // 契约冻结后在此调用 Swarm 运行时 stop;当前仅在本地记录意图,避免对未冻结写接口下注。 - common.ApiSuccess(c, gin.H{ - "deployment_id": dep.DeploymentID, - "swarm_id": dep.RuntimeSwarmID, - "accepted": true, - "note": "runtime stop wiring lands with agent_swarm#2 contract freeze", - }) + agentError(c, "NOT_IMPLEMENTED", + "swarm stop runtime wiring is pending agent_swarm#2 contract freeze; not forwarded to runtime") } diff --git a/heicode/controller/agent_swarm_query_test.go b/heicode/controller/agent_swarm_query_test.go new file mode 100644 index 0000000..aed81e9 --- /dev/null +++ b/heicode/controller/agent_swarm_query_test.go @@ -0,0 +1,46 @@ +package controller + +import ( + "encoding/json" + "testing" + + "github.com/stretchr/testify/require" +) + +// #45 复审 #1:事件 payload 必须脱敏 —— 递归剔除 secret_ref / credentials / 大字段, +// 并对结果再跑 RedactText 兜底,绝不把 azkv:// secret_ref 或 sk-/Bearer 下发给客户端。 +func TestSanitizeSwarmPayload_StripsSecrets(t *testing.T) { + raw := `{ + "task_id":"t1", + "approval":{"secret_ref":"azkv://heicode-kv.vault.azure.net/secrets/git-pat","note":"deploy"}, + "credentials":{"access_key":"AKIA123","secret_access_key":"xxx"}, + "headers":{"authorization":"Bearer aZ09tokenVALUE"}, + "api_key":"sk-abcDEF1234567890", + "stack":["plain frame","key=sk-leak0987654321ABCD"], + "ok":true + }` + out := sanitizeSwarmPayload(raw) + b, err := json.Marshal(out) + require.NoError(t, err) + s := string(b) + + // 敏感键被递归剔除 + require.NotContains(t, s, "secret_ref") + require.NotContains(t, s, "azkv://") + require.NotContains(t, s, "git-pat") + require.NotContains(t, s, "credentials") + require.NotContains(t, s, "AKIA123") + require.NotContains(t, s, "api_key") + // RedactText 兜底:残留在普通字段里的 sk-/Bearer 也被打码 + require.NotContains(t, s, "sk-leak0987654321ABCD") + require.NotContains(t, s, "aZ09tokenVALUE") + // 非敏感内容保留 + require.Contains(t, s, "t1") + require.Contains(t, s, "ok") +} + +func TestSanitizeSwarmPayload_EmptyAndPlain(t *testing.T) { + require.Empty(t, sanitizeSwarmPayload("")) + out := sanitizeSwarmPayload(`{"route":"chat","n":3}`) + require.Equal(t, "chat", out["route"]) +} diff --git a/heicode/model/agent_callback.go b/heicode/model/agent_callback.go index 19480be..85cb14a 100644 --- a/heicode/model/agent_callback.go +++ b/heicode/model/agent_callback.go @@ -1,6 +1,9 @@ package model -import "errors" +import ( + "errors" + "strings" +) type AgentCallbackEvent struct { Id int `gorm:"primaryKey" json:"id"` @@ -60,7 +63,7 @@ func InsertAgentCallbackEvent(row *AgentCallbackEvent) (bool, error) { // incremental polling (#45 events?after). Ordered oldest-first so the client can // append; the caller uses the last returned Id as the next `after`. Reading from // HM-persisted callback rows means this needs no live Swarm call. -func ListSwarmCallbackEventsAfter(deploymentID, swarmID string, afterID, limit int) ([]AgentCallbackEvent, error) { +func ListSwarmCallbackEventsAfter(userID, deploymentID, swarmID string, afterID, limit int) ([]AgentCallbackEvent, error) { if DB == nil { return nil, nil } @@ -75,6 +78,10 @@ func ListSwarmCallbackEventsAfter(deploymentID, swarmID string, afterID, limit i default: return nil, nil } + // 防跨用户泄漏:即使 runtime_swarm_id/deployment_id 碰撞或误写,也按 user_id 收口(#45 复审 #3)。 + if strings.TrimSpace(userID) != "" { + q = q.Where("user_id = ?", userID) + } if afterID > 0 { q = q.Where("id > ?", afterID) } diff --git a/heicode/model/agent_swarm_query_test.go b/heicode/model/agent_swarm_query_test.go index be3f4ac..e6f0c03 100644 --- a/heicode/model/agent_swarm_query_test.go +++ b/heicode/model/agent_swarm_query_test.go @@ -6,40 +6,45 @@ import ( "github.com/stretchr/testify/require" ) -// #45: events?after 游标 —— 按 deployment_id/swarm_id 过滤,id>after 增量返回,oldest-first。 +// #45: events?after 游标 + user 作用域 —— 按 deployment_id/swarm_id 过滤,id>after 增量,oldest-first; +// 传入 userID 时按 user_id 收口(防跨用户泄漏,复审 #3)。 func TestListSwarmCallbackEventsAfter(t *testing.T) { require.NoError(t, LOG_DB.Where("1 = 1").Delete(&AgentCallbackEvent{}).Error) - mk := func(eventID, dep, swarm, etype string) { + mk := func(eventID, uid, dep, swarm, etype string) { _, err := InsertAgentCallbackEvent(&AgentCallbackEvent{ - EventID: eventID, DeploymentID: dep, SwarmID: swarm, EventType: etype, + EventID: eventID, UserID: uid, DeploymentID: dep, SwarmID: swarm, EventType: etype, }) require.NoError(t, err) } - mk("e1", "dep_A", "swarm_A", "deployment.status_changed") - mk("e2", "dep_A", "swarm_A", "task.created") - mk("e3", "dep_A", "swarm_A", "artifact.produced") - mk("e4", "dep_B", "swarm_B", "task.created") // 另一个 run,应被排除 + mk("e1", "7", "dep_A", "swarm_A", "deployment.status_changed") + mk("e2", "7", "dep_A", "swarm_A", "task.created") + mk("e3", "7", "dep_A", "swarm_A", "artifact.produced") + mk("e4", "7", "dep_B", "swarm_B", "task.created") // 另一个 run + mk("e5", "9", "dep_A", "swarm_A", "task.created") // 同 dep/swarm 但别的用户 → 不应泄漏给 user 7 - // 全量(after=0):dep_A 的 3 条,oldest-first - all, err := ListSwarmCallbackEventsAfter("dep_A", "swarm_A", 0, 100) + // user 7 + dep_A:3 条(e5 属 user 9,被排除),oldest-first + all, err := ListSwarmCallbackEventsAfter("7", "dep_A", "swarm_A", 0, 100) require.NoError(t, err) require.Len(t, all, 3) require.Equal(t, "e1", all[0].EventID) + for _, e := range all { + require.NotEqual(t, "e5", e.EventID, "不得返回别的用户的事件") + } - // 游标:after = 第一条 id → 只返回其后的 2 条 - after := all[0].Id - rest, err := ListSwarmCallbackEventsAfter("dep_A", "swarm_A", after, 100) + // 游标:after = 第一条 id → 其后 2 条 + rest, err := ListSwarmCallbackEventsAfter("7", "dep_A", "swarm_A", all[0].Id, 100) require.NoError(t, err) require.Len(t, rest, 2) require.Equal(t, "e2", rest[0].EventID) - // 仅按 swarm_id 也能查到 - bySwarm, err := ListSwarmCallbackEventsAfter("", "swarm_A", 0, 100) + // user 9 只看到自己的 e5 + u9, err := ListSwarmCallbackEventsAfter("9", "dep_A", "swarm_A", 0, 100) require.NoError(t, err) - require.Len(t, bySwarm, 3) + require.Len(t, u9, 1) + require.Equal(t, "e5", u9[0].EventID) // 空标识 → 空 - none, err := ListSwarmCallbackEventsAfter("", "", 0, 100) + none, err := ListSwarmCallbackEventsAfter("7", "", "", 0, 100) require.NoError(t, err) require.Empty(t, none) }