Merge pull request #64 from xmindlab-heicode/fix/reland-swarm-stop-sse
fix: re-land swarm stop + events SSE into main (#61 误合到错误 base)
This commit is contained in:
@@ -313,17 +313,18 @@ signature = base64( ed25519_sign( device_priv, sha256(canonical) ) )
|
||||
|
||||
---
|
||||
|
||||
## 5.2 Swarm 运行查询 🟡(#45 Phase1 · 只读预览)
|
||||
## 5.2 Swarm 运行查询 🟢(#45/#46 · 契约 v1 已冻结)
|
||||
|
||||
多 Agent 蜂群运行(`agent_swarm` / HeiCode Swarm)的**只读查询**。数据全部来自 HM 已持久化的运行时回调(Swarm → HM 带签名回调),**无需实时调 Swarm**,故不受 `agent_swarm#2` 契约冻结阻塞。
|
||||
多 Agent 蜂群运行(`agent_swarm` / HeiCode Swarm)的查询。读类接口数据全部来自 HM 已持久化的运行时回调(Swarm → HM 带签名回调),**无需实时调 Swarm**;契约已冻结为 `runtime-contract v1`/`event-schema v1`(`agent_swarm#14`/`#15`)。
|
||||
|
||||
| 方法 | 路径 | 说明 | 状态 |
|
||||
|---|---|---|---|
|
||||
| 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/events/stream?after=` | **事件 SSE 实时流**(与 events 同源;进入即回放 `after` 之后历史) | 🟢 |
|
||||
| GET | `/api/heicode/swarms/:id/artifacts` | 从已存事件派生的产物 | 🟢(派生) |
|
||||
| POST | `/api/heicode/swarms/:id/stop` | 停止运行(**写**) | 🟡 待 `agent_swarm#2` 冻结 + `SWARM_RUNTIME_ENABLED=true` |
|
||||
| POST | `/api/heicode/swarms/:id/stop` | 停止运行(**写**,真实调运行时) | 🟢(开关:`SWARM_RUNTIME_ENABLED=true` + base_url + service token;否则 `POLICY_REJECTED`) |
|
||||
|
||||
```json
|
||||
// GET /api/heicode/swarms/:id
|
||||
@@ -350,6 +351,8 @@ signature = base64( ed25519_sign( device_priv, sha256(canonical) ) )
|
||||
- **事件**:`sequence` = agent_swarm event-schema v1 的 **per-swarm 严格递增序号**(每 swarm 从 1、无空洞),客户端用它去重/排序;`id`/`next_after` 是 HM 不透明分页游标(单调,兼容 `sequence` 尚未全量上线)。事件 `payload` **已脱敏**(递归剔除 `secret_ref`/`credential_ref`/`signing_secret_ref`/credentials/大字段 + `RedactText` 兜底)。查询按**当前用户**作用域。
|
||||
- **新增事件类型**(已注册):`swarm.completed/failed/stopped`、`approval.approved/rejected`、`handoff.created`(event-schema v1)。
|
||||
- **artifact**:从 `artifact.created` 派生扁平 `{uri,checksum,task_id,size_bytes?,created_at}`(**无 secret_ref**;size 未知则省略,不伪造)。
|
||||
- **SSE 实时流**(`/events/stream`):`text/event-stream`,帧 `event: message` 的 `data` 即单条事件视图(同 `events` 的 item);命中终态(`swarm.completed/failed/stopped`)发 `event: done`(`{next_after,reason}`)后收流,客户端断开或超时(30min)亦收流。与轮询 `events?after` 等价,二选一;断线后用最后的 `next_after` 重新订阅即可续传。
|
||||
- **stop**:`SWARM_RUNTIME_ENABLED=true` 且配齐 base_url + service token 时真实调用运行时(`X-Idempotency-Key` 幂等),返回 `{accepted:true,runtime_status,...}`;**终态不在响应里**——由运行时异步回推 `swarm.stopped` 事件,客户端经 `events`/`stream` 收最终态。未启用则 `POLICY_REJECTED`(不伪造受理)。
|
||||
- `stop`(写):仍 gated,待 stop 真实接入 PR(复用运行时客户端 + `SWARM_RUNTIME_SERVICE_TOKEN`,契约 §3 已冻结)落地。
|
||||
|
||||
---
|
||||
|
||||
@@ -13,7 +13,12 @@
|
||||
- **HM(Heicode Manager)当前不实现 swarm runtime。** HM 的职责边界是:模型网关(`/v1/*`)+ 资源/权限/计费/审计 + **模板 Agent 部署编排**(经 AM 启动常驻 agent、客户端直连)。多 agent 蜂群编排**不在 HM 端**。
|
||||
- **旧的「HM 内部 sub/蜂群任务编排」模型已作废。** 那套(sub 任务、display_status、HM 侧蜂群 runtime 对接草案)随产品转向「模板 Agent + 客户端直连」一并下线,相关代码删除清单见 [`heicode-hm-legacy-teardown.md`](./heicode-hm-legacy-teardown.md),当前模型见 [`heicode-hm-template-agent-model.md`](./heicode-hm-template-agent-model.md)。
|
||||
- **新版蜂群能力在 `agent_swarm`(HeiCode Swarm)仓,不在本仓。** 当前模型:主控 Agent(Master Agent,`orchestrator/master_agent.py`)把需求**分解**为子任务 → **派发**给不同领域的专家 Agent **并行执行** → 重叠领域**协作/移交** → 主控**评审/重做**循环(受 `MAX_REVIEW_CYCLES` 约束)→ **汇总交付**;Orchestrator(FastAPI) + Redis 权威状态 + WebSocket Agent 协议 + Prometheus 指标。旧文档描述的「HM 主导编排蜂群 / 仅 `/tasks`」已作废。HM 侧最多提供资源/计费/鉴权支撑面,runtime 与编排由 Swarm 承载。
|
||||
- **该仓尚未作为 Heicode 主链路正式 Runtime Backend 接入。** 按 `agent_swarm` README 与 `agent_swarm#2`:编排器已可运行(仓内),但 Manager↔Runtime 生命周期契约、HMAC 签名回调 envelope、稳定 `deployment_id`/`workflow_id`/`trace_id`、统一 usage/审计接入等仍为 🟡 待接入,需各 Team 评审冻结后联调。
|
||||
- **契约已冻结为 v1(2026-06-10)。** Swarm 侧已冻结 `runtime-contract v1`(`agent_swarm#14`:stop = `POST /api/agent/swarm/deployments/{deployment_id}/stop` + `Bearer SWARM_RUNTIME_SERVICE_TOKEN` + `X-Idempotency-Key`,受理后异步回推 `swarm.stopped`;`deployment_id↔swarm_id↔manager_deployment_id` 三映射;状态 §4.1 `blocked→degraded`)与 `event-schema v1`(`agent_swarm#15`:per-swarm 严格递增 `sequence`;新增 6 类事件;`secret_ref`/`credential_ref`/`signing_secret_ref` 按设计以 `azkv://` 引用透传)。
|
||||
- **HM 侧已据冻结契约落地读侧 + stop + SSE(不再 deferred 的部分):**
|
||||
- **只读查询(#45 / PR #59)🟢**:list/status/`events?after`/artifacts 全部基于 HM 已持久化的回调事件(`agent_callback.go` → `AgentCallbackEvent`),按 `user_id` 收口、payload 递归脱敏(剔除 secret_ref/凭据/大字段 + `RedactText` 兜底),**无需调用 Swarm 运行时**。`controller/agent_swarm_query.go`。
|
||||
- **stop(#45)🟢(受开关约束)**:真实调用运行时冻结路径,复用 `agent_runtime_client` 的配置/URL/信封解析;`SWARM_RUNTIME_ENABLED=true` 且配齐 `SWARM_RUNTIME_BASE_URL` + `SWARM_RUNTIME_SERVICE_TOKEN` 时才发起,否则 `POLICY_REJECTED`,**绝不伪造 accepted**。终态不抢写,由 `swarm.stopped` 回调写回。
|
||||
- **事件 SSE 实时流(#46 读侧)🟢**:`GET /api/heicode/swarms/:id/events/stream?after=`,与 `events?after` 同源(HM 持久化事件),命中终态(`swarm.completed/failed/stopped`)/客户端断开/超时结束,**不依赖运行时**。
|
||||
- **仍 deferred / 阻塞的部分:** ① **per-user 计费令牌注入(#60)**——swarm create 目前未像模板 Agent 那样现签隐藏 `sk-` 注入运行时(`OPENAI_API_KEY`+`OPENAI_API_BASE=HM/v1`),是 HM 唯一计费缺口,**阻塞于 Swarm 对 `agent_swarm#16` 注入口径 A/B 的确认**;② 事件 **append 写入 / 主动拉取**仍走回调被动落库,未做 HM→Swarm 主动拉取(无对应需求)。
|
||||
|
||||
---
|
||||
|
||||
@@ -22,10 +27,10 @@
|
||||
| 项 | 归属仓 / 负责人 | 说明 |
|
||||
|---|---|---|
|
||||
| Swarm runtime / 多 agent 编排 | **`agent_swarm`**(产品名 HeiCode Swarm,@Songhaoz666) | 执行面、Master-Agent 编排、回调、Swarm Runtime |
|
||||
| Manager ↔ Swarm 契约 | Swarm 侧**已起草正式契约** `agent_swarm/docs/integration/runtime-contract.md`(沿用 `heicode-am-contract` 的鉴权/回调/env/路径覆盖约定),**待 Manager Runtime Team 评审冻结**;冻结后在本目录 `docs/integration/` 另立 HM 侧对接契约 | Swarm 已实现生命周期接口 create/status/tasks/logs/events/metrics/workflow/diagnostics/stop/approvals(均带 `deployment_id`,三组路径别名 `/api/swarms`、`/api/agent/swarm/deployments`、`/api/agnet/deployments`);契约冻结与主链路接入跟踪在 **`agent_swarm#2`**(`#1` 执行面缺口已关闭) |
|
||||
| Manager ↔ Swarm 契约 | Swarm 侧契约**已冻结为 v1**:`runtime-contract v1`(`agent_swarm#14`)+ `event-schema v1`(`agent_swarm#15`),见 `agent_swarm/docs/integration/`;HM 侧据此对接,本文为 HM 侧锚点 | Swarm 已实现生命周期接口 create/status/tasks/logs/events/metrics/workflow/diagnostics/stop/approvals(均带 `deployment_id`,三组路径别名 `/api/swarms`、`/api/agent/swarm/deployments`、`/api/agnet/deployments`);HM 接入跟踪在 #45/#46/#60 与 PR #59 |
|
||||
| Agent 运行时(单 agent,已落地) | **`agent_management`(AM)**(@azgy) | 模板 Agent 启动/状态/停止/删除,契约见 [`heicode-am-contract.md`](./heicode-am-contract.md) |
|
||||
|
||||
> **HM 侧后续若要支撑蜂群**:不恢复旧文档,按 Swarm 的 `runtime-contract.md` 冻结版在 `docs/integration/` 新立 HM 对接契约;本文作为「蜂群在 HM 侧当前为 deferred」的唯一锚点。HM 侧对接前置依赖见 issue #45(Phase1 只读查询,阻塞于 `agent_swarm#2` 契约冻结)/ #46(Phase2 SSE)。
|
||||
> **HM 侧对接状态(2026-06-10)**:契约已按 `runtime-contract v1`/`event-schema v1` 冻结(`agent_swarm#14`/`#15`)。HM 已落地 #45(只读查询 + stop,PR #59 + 本批)与 #46 读侧(events SSE);剩余阻塞仅 #60(per-user 计费令牌注入,待 `agent_swarm#16` A/B 确认)。本文仍是「蜂群在 HM 侧」的唯一锚点——不恢复旧文档。
|
||||
|
||||
---
|
||||
|
||||
@@ -56,4 +61,4 @@
|
||||
|
||||
## 4. 给后续开发者的一句话
|
||||
|
||||
要找「蜂群在 HM 侧怎么对接」——**当前答案是「HM 不实现 swarm runtime;Swarm 侧已起草 `agent_swarm/docs/integration/runtime-contract.md`,待 Manager Runtime Team 评审冻结(`agent_swarm#2`)后,HM 再据冻结版做只读查询接入(#45/#46)」**;旧设计(HM 主导编排 / 仅 `/tasks` / `HeiCode-Swarm` 仓名)已废,别从 git 历史里捞旧文档当依据,按本文与 `agent_swarm` 仓的最新结论走。
|
||||
要找「蜂群在 HM 侧怎么对接」——**当前答案是「HM 不实现 swarm runtime;但契约已冻结(`runtime-contract v1`/`event-schema v1`),HM 已据冻结版落地只读查询 + stop + 事件 SSE(#45/#46/PR #59),代码在 `controller/agent_swarm_query.go`;唯一剩余计费缺口是 per-user `sk-` 注入(#60),阻塞于 `agent_swarm#16`」**;旧设计(HM 主导编排 / 仅 `/tasks` / `HeiCode-Swarm` 仓名)已废,别从 git 历史里捞旧文档当依据,按本文与 `agent_swarm` 仓的最新结论走。
|
||||
|
||||
@@ -1,9 +1,15 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/heicode/manager/common"
|
||||
@@ -298,23 +304,201 @@ func firstNonNil(vals ...any) any {
|
||||
}
|
||||
|
||||
// HeicodeStopSwarm: POST /api/heicode/swarms/:id/stop — 停止 swarm 运行(写操作)。
|
||||
// 这是唯一需要真正调用 Swarm 运行时的操作。在 agent_swarm#2 契约冻结 + SWARM_RUNTIME_ENABLED=true
|
||||
// 前默认关闭并明确提示(不臆造未冻结的写接口)。冻结后在此接入运行时 stop(草案路径
|
||||
// /api/agent/swarm/deployments/{id}/stop)。
|
||||
// 这是唯一需要真正调用 Swarm 运行时的操作。契约已随 agent_swarm runtime-contract v1 冻结
|
||||
// (#45 复审 + agent_swarm#14):POST /api/agent/swarm/deployments/{deployment_id}/stop,
|
||||
// Bearer SWARM_RUNTIME_SERVICE_TOKEN,X-Idempotency-Key 幂等;运行时受理后异步停机并回推
|
||||
// swarm.stopped 事件。开关 SWARM_RUNTIME_ENABLED=false(默认)或缺 base_url/token 时一律拒绝,
|
||||
// **绝不伪造 accepted**(#45 复审 #2)。
|
||||
func HeicodeStopSwarm(c *gin.Context) {
|
||||
dep, ok := findUserSwarmDeployment(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
// 运行时 stop 的真实接入随 agent_swarm#2 契约冻结落地。**在此之前一律不伪造 accepted**
|
||||
// (#45 复审 #2:开关打开也不能返回 accepted:true 误导客户端/审计)。无论开关如何,均返回
|
||||
// 明确的「未实现/待契约」语义,直到真正接上运行时 stop。
|
||||
_ = dep
|
||||
if !common.GetEnvOrDefaultBool("SWARM_RUNTIME_ENABLED", false) {
|
||||
cfg := agentRuntimeClientConfigForMode(agentRuntimeModeSwarm)
|
||||
// 开关 + base_url + 服务令牌缺一不可:任一缺失都不调用、不伪造受理。
|
||||
if !cfg.Enabled || strings.TrimSpace(cfg.BaseURL) == "" || strings.TrimSpace(cfg.Token) == "" {
|
||||
agentError(c, "POLICY_REJECTED",
|
||||
"swarm stop not enabled: set SWARM_RUNTIME_ENABLED=true after agent_swarm#2 runtime-contract freeze")
|
||||
"swarm stop not enabled: requires SWARM_RUNTIME_ENABLED=true with SWARM_RUNTIME_BASE_URL and SWARM_RUNTIME_SERVICE_TOKEN")
|
||||
return
|
||||
}
|
||||
agentError(c, "NOT_IMPLEMENTED",
|
||||
"swarm stop runtime wiring is pending agent_swarm#2 contract freeze; not forwarded to runtime")
|
||||
var body struct {
|
||||
Reason string `json:"reason"`
|
||||
}
|
||||
_ = c.ShouldBindJSON(&body)
|
||||
result, err := callSwarmRuntimeStop(c.Request.Context(), cfg, dep, body.Reason)
|
||||
if err != nil {
|
||||
common.SysLog("HeicodeStopSwarm: " + err.Error())
|
||||
agentError(c, "RUNTIME_UNAVAILABLE", "failed to stop swarm run at runtime: "+err.Error())
|
||||
return
|
||||
}
|
||||
// 运行时已受理停机;真正终态由运行时异步回推 swarm.stopped(经 /api/agent/callbacks)写回,
|
||||
// HM 此处不抢先改写 dep.Status 以免与回调竞态。只落审计。
|
||||
recordAgentAuditEvent(agentEvent{
|
||||
EventID: "evt_" + common.GetUUID()[:12],
|
||||
Event: "swarm.stop_requested",
|
||||
SchemaVersion: 1,
|
||||
UserID: dep.UserID,
|
||||
ChannelID: dep.ChannelID,
|
||||
BindingScope: dep.BindingScope,
|
||||
DeploymentID: dep.DeploymentID,
|
||||
CorrelationID: dep.CorrelationID,
|
||||
OccurredAt: agentNow(),
|
||||
}, "agent_swarm_query", dep.DeploymentID, agentRequestID(c), "ok")
|
||||
common.ApiSuccess(c, gin.H{
|
||||
"deployment_id": dep.DeploymentID,
|
||||
"swarm_id": dep.RuntimeSwarmID,
|
||||
"accepted": true,
|
||||
"runtime_status": firstNonEmpty(result.RuntimeStatus, "stopping"),
|
||||
"note": "stop accepted by runtime; final state arrives via swarm.stopped callback",
|
||||
})
|
||||
}
|
||||
|
||||
// swarmStopResult 是运行时 stop 调用的归一化结果。
|
||||
type swarmStopResult struct {
|
||||
RuntimeStatus string
|
||||
HTTPStatus int
|
||||
}
|
||||
|
||||
// callSwarmRuntimeStop 向 Swarm 运行时发起真实 stop(runtime-contract v1)。复用 agent_runtime_client
|
||||
// 的配置/URL/信封解析 helper,但按 model.AgentDeployment 直接构造请求(无需 agentDeploymentRecord)。
|
||||
// runtime id 取 runtime_deployment_id,缺失时回退 runtime_swarm_id(契约三映射)。
|
||||
func callSwarmRuntimeStop(ctx context.Context, cfg agentRuntimeConfig, dep model.AgentDeployment, reason string) (swarmStopResult, error) {
|
||||
runtimeID := firstNonEmpty(strings.TrimSpace(dep.RuntimeDeploymentID), strings.TrimSpace(dep.RuntimeSwarmID))
|
||||
if runtimeID == "" {
|
||||
return swarmStopResult{}, errors.New("swarm run has no runtime deployment id yet")
|
||||
}
|
||||
endpoint, err := agentRuntimeURL(cfg.BaseURL, agentRuntimeStopPath(cfg, runtimeID))
|
||||
if err != nil {
|
||||
return swarmStopResult{}, err
|
||||
}
|
||||
payload, err := common.Marshal(gin.H{
|
||||
"reason": firstNonEmpty(strings.TrimSpace(reason), "Heicode Manager requested stop"),
|
||||
"manager_deployment_id": dep.DeploymentID,
|
||||
})
|
||||
if err != nil {
|
||||
return swarmStopResult{}, err
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(payload))
|
||||
if err != nil {
|
||||
return swarmStopResult{}, err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Authorization", "Bearer "+strings.TrimSpace(cfg.Token))
|
||||
// 幂等键固定按 manager deployment_id 派生:同一 stop 重试不会在运行时侧重复执行(契约要求)。
|
||||
req.Header.Set("X-Idempotency-Key", "manager-stop-"+dep.DeploymentID)
|
||||
if cid := strings.TrimSpace(dep.CorrelationID); cid != "" {
|
||||
req.Header.Set("X-Correlation-ID", cid)
|
||||
}
|
||||
client := &http.Client{Timeout: cfg.Timeout}
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
return swarmStopResult{}, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
respBody, readErr := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
if readErr != nil {
|
||||
return swarmStopResult{HTTPStatus: resp.StatusCode}, readErr
|
||||
}
|
||||
if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices {
|
||||
return swarmStopResult{HTTPStatus: resp.StatusCode}, fmt.Errorf("runtime stop returned HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(respBody)))
|
||||
}
|
||||
var envelope map[string]any
|
||||
if len(respBody) > 0 {
|
||||
if err := common.Unmarshal(respBody, &envelope); err != nil {
|
||||
return swarmStopResult{HTTPStatus: resp.StatusCode}, err
|
||||
}
|
||||
}
|
||||
if message := agentRuntimeEnvelopeError(envelope); message != "" {
|
||||
return swarmStopResult{HTTPStatus: resp.StatusCode}, errors.New(message)
|
||||
}
|
||||
data := extractAgentRuntimeData(envelope)
|
||||
return swarmStopResult{
|
||||
RuntimeStatus: stringFromMap(data, "runtime_status", "status"),
|
||||
HTTPStatus: resp.StatusCode,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// swarmTerminalEventTypes 是 agent_swarm event-schema v1(#15)的运行终态事件。命中即可结束 SSE 流。
|
||||
var swarmTerminalEventTypes = map[string]bool{
|
||||
"swarm.completed": true,
|
||||
"swarm.failed": true,
|
||||
"swarm.stopped": true,
|
||||
}
|
||||
|
||||
func swarmEventIsTerminal(eventType string) bool {
|
||||
return swarmTerminalEventTypes[strings.ToLower(strings.TrimSpace(eventType))]
|
||||
}
|
||||
|
||||
// HeicodeStreamSwarmEvents: GET /api/heicode/swarms/:id/events/stream?after= — SSE 实时事件流。
|
||||
//
|
||||
// 与 §events?after 同源:**全部来自 HM 已持久化的回调事件**(model.ListSwarmCallbackEventsAfter,
|
||||
// 按 user_id 收口防跨用户泄漏),服务端轮询新事件后以 SSE 帧推送,**无需调用 Swarm 运行时**——
|
||||
// 因此与 stop 不同,不受 SWARM_RUNTIME_ENABLED 限制。事件 payload 经 swarmEventView 脱敏。
|
||||
// 进入即回放 after 之后的历史(免去客户端先调 events?after 再订阅);命中终态事件
|
||||
// (swarm.completed/failed/stopped)、客户端断开或超过 maxLifetime 即结束。
|
||||
func HeicodeStreamSwarmEvents(c *gin.Context) {
|
||||
dep, ok := findUserSwarmDeployment(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
userID := strconv.Itoa(c.GetInt("id"))
|
||||
after, _ := strconv.Atoi(strings.TrimSpace(c.Query("after")))
|
||||
|
||||
c.Writer.Header().Set("Content-Type", "text/event-stream")
|
||||
c.Writer.Header().Set("Cache-Control", "no-cache")
|
||||
c.Writer.Header().Set("Connection", "keep-alive")
|
||||
c.Writer.Header().Set("X-Accel-Buffering", "no") // 关掉反代缓冲,逐帧下发
|
||||
|
||||
const (
|
||||
pollInterval = 1500 * time.Millisecond
|
||||
maxLifetime = 30 * time.Minute
|
||||
)
|
||||
deadline := time.Now().Add(maxLifetime)
|
||||
ticker := time.NewTicker(pollInterval)
|
||||
defer ticker.Stop()
|
||||
ctx := c.Request.Context()
|
||||
|
||||
// flush 推送 after 之后的新事件,返回 false 表示已命中终态(应结束流)或读失败。
|
||||
flush := func() bool {
|
||||
events, err := model.ListSwarmCallbackEventsAfter(userID, dep.DeploymentID, dep.RuntimeSwarmID, after, 200)
|
||||
if err != nil {
|
||||
c.SSEvent("error", gin.H{"message": "failed to read swarm events"})
|
||||
c.Writer.Flush()
|
||||
return false
|
||||
}
|
||||
terminal := false
|
||||
for _, e := range events {
|
||||
c.SSEvent("message", swarmEventView(e))
|
||||
if e.Id > after {
|
||||
after = e.Id
|
||||
}
|
||||
if swarmEventIsTerminal(e.EventType) {
|
||||
terminal = true
|
||||
}
|
||||
}
|
||||
c.Writer.Flush()
|
||||
return !terminal
|
||||
}
|
||||
if !flush() {
|
||||
c.SSEvent("done", gin.H{"next_after": after, "reason": "terminal"})
|
||||
c.Writer.Flush()
|
||||
return
|
||||
}
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done(): // 客户端断开
|
||||
return
|
||||
case <-ticker.C:
|
||||
if time.Now().After(deadline) {
|
||||
c.SSEvent("done", gin.H{"next_after": after, "reason": "timeout"})
|
||||
c.Writer.Flush()
|
||||
return
|
||||
}
|
||||
if !flush() {
|
||||
c.SSEvent("done", gin.H{"next_after": after, "reason": "terminal"})
|
||||
c.Writer.Flush()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,9 +1,13 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/heicode/manager/model"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -80,6 +84,86 @@ func TestSwarmArtifactView(t *testing.T) {
|
||||
require.Equal(t, "2026-06-10T01:00:00Z", v2["created_at"])
|
||||
}
|
||||
|
||||
// #45/#15: SSE 终态判定 —— 仅 swarm.completed/failed/stopped 收流,其余继续。
|
||||
func TestSwarmEventIsTerminal(t *testing.T) {
|
||||
for _, et := range []string{"swarm.completed", "swarm.failed", "swarm.stopped", "SWARM.STOPPED"} {
|
||||
require.True(t, swarmEventIsTerminal(et), "%s 应为终态", et)
|
||||
}
|
||||
for _, et := range []string{"swarm.started", "handoff.created", "approval.approved", "", "running"} {
|
||||
require.False(t, swarmEventIsTerminal(et), "%s 不应为终态", et)
|
||||
}
|
||||
}
|
||||
|
||||
// #45: stop 真实接入 —— 命中冻结契约路径/鉴权/幂等头,并解析 {success,data} 信封。
|
||||
func TestCallSwarmRuntimeStop(t *testing.T) {
|
||||
var gotPath, gotAuth, gotIdem, gotCorr string
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
gotPath = r.URL.Path
|
||||
gotAuth = r.Header.Get("Authorization")
|
||||
gotIdem = r.Header.Get("X-Idempotency-Key")
|
||||
gotCorr = r.Header.Get("X-Correlation-ID")
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_, _ = w.Write([]byte(`{"success":true,"data":{"runtime_status":"stopping"}}`))
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
cfg := agentRuntimeConfig{
|
||||
Enabled: true,
|
||||
BaseURL: srv.URL,
|
||||
Token: "svc-token",
|
||||
StopPath: "/api/agent/swarm/deployments/{deployment_id}/stop",
|
||||
Timeout: 5 * time.Second,
|
||||
}
|
||||
dep := model.AgentDeployment{
|
||||
DeploymentID: "dep-1",
|
||||
RuntimeDeploymentID: "rt-9",
|
||||
RuntimeSwarmID: "sw-7",
|
||||
CorrelationID: "cor-3",
|
||||
}
|
||||
res, err := callSwarmRuntimeStop(context.Background(), cfg, dep, "user requested")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "stopping", res.RuntimeStatus)
|
||||
require.Equal(t, "/api/agent/swarm/deployments/rt-9/stop", gotPath) // runtime_deployment_id 优先
|
||||
require.Equal(t, "Bearer svc-token", gotAuth)
|
||||
require.Equal(t, "manager-stop-dep-1", gotIdem) // 幂等键按 manager deployment_id 派生
|
||||
require.Equal(t, "cor-3", gotCorr)
|
||||
}
|
||||
|
||||
// runtime_deployment_id 缺失时回退 runtime_swarm_id(契约三映射);无任何 runtime id 则报错不发请求。
|
||||
func TestCallSwarmRuntimeStop_RuntimeIDFallbackAndMissing(t *testing.T) {
|
||||
var gotPath string
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
gotPath = r.URL.Path
|
||||
_, _ = w.Write([]byte(`{"success":true,"data":{}}`))
|
||||
}))
|
||||
defer srv.Close()
|
||||
cfg := agentRuntimeConfig{
|
||||
Enabled: true, BaseURL: srv.URL, Token: "t",
|
||||
StopPath: "/api/agent/swarm/deployments/{deployment_id}/stop", Timeout: 5 * time.Second,
|
||||
}
|
||||
_, err := callSwarmRuntimeStop(context.Background(), cfg, model.AgentDeployment{DeploymentID: "d", RuntimeSwarmID: "sw-only"}, "")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "/api/agent/swarm/deployments/sw-only/stop", gotPath)
|
||||
|
||||
_, err = callSwarmRuntimeStop(context.Background(), cfg, model.AgentDeployment{DeploymentID: "d"}, "")
|
||||
require.Error(t, err) // 无 runtime id → 不发请求
|
||||
}
|
||||
|
||||
// 运行时非 2xx → 返回错误(handler 据此回 RUNTIME_UNAVAILABLE,不伪造受理)。
|
||||
func TestCallSwarmRuntimeStop_RuntimeError(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusBadGateway)
|
||||
_, _ = w.Write([]byte(`upstream down`))
|
||||
}))
|
||||
defer srv.Close()
|
||||
cfg := agentRuntimeConfig{
|
||||
Enabled: true, BaseURL: srv.URL, Token: "t",
|
||||
StopPath: "/api/agent/swarm/deployments/{deployment_id}/stop", Timeout: 5 * time.Second,
|
||||
}
|
||||
_, err := callSwarmRuntimeStop(context.Background(), cfg, model.AgentDeployment{DeploymentID: "d", RuntimeDeploymentID: "rt"}, "")
|
||||
require.Error(t, err)
|
||||
}
|
||||
|
||||
// #15: 6 类新事件已注册(类别 + 必填字段两张表)。
|
||||
func TestSwarmFrozenEventTypesRegistered(t *testing.T) {
|
||||
for _, et := range []string{
|
||||
|
||||
@@ -552,10 +552,12 @@ func SetApiRouter(router *gin.Engine) {
|
||||
// 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.
|
||||
// callback data (no live Swarm call); stop calls the runtime (runtime-contract v1,
|
||||
// gated by SWARM_RUNTIME_ENABLED + base_url + service token).
|
||||
heicodeAgentRoute.GET("/swarms", controller.HeicodeListSwarms)
|
||||
heicodeAgentRoute.GET("/swarms/:id", controller.HeicodeGetSwarmStatus)
|
||||
heicodeAgentRoute.GET("/swarms/:id/events", controller.HeicodeListSwarmEvents)
|
||||
heicodeAgentRoute.GET("/swarms/:id/events/stream", controller.HeicodeStreamSwarmEvents)
|
||||
heicodeAgentRoute.GET("/swarms/:id/artifacts", controller.HeicodeListSwarmArtifacts)
|
||||
heicodeAgentRoute.POST("/swarms/:id/stop", controller.HeicodeStopSwarm)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user