From a757c89276baadb3b658b071531ee6e9d7b903b6 Mon Sep 17 00:00:00 2001 From: sunlei Date: Thu, 23 Jul 2026 09:30:19 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=A2=9E=E5=8A=A0=E7=BD=91=E7=BB=9C?= =?UTF-8?q?=E7=AE=A1=E7=90=86=E7=8A=B6=E6=80=81=E4=BA=8B=E4=BB=B6=E6=B5=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- API.md | 29 +- README.md | 4 +- .../admin-platform-config.module.ts | 2 + .../network-agent-mqtt.service.ts | 385 +++++++++++------- ...network-management-event-stream.service.ts | 122 ++++++ .../network-management.controller.ts | 24 +- .../network-management.types.ts | 7 + .../network-agent-mqtt.service.spec.ts | 28 ++ ...rk-management-event-stream.service.spec.ts | 92 +++++ .../network-management.controller.spec.ts | 73 ++++ 10 files changed, 602 insertions(+), 164 deletions(-) create mode 100644 src/modules/admin/platform-config/network-management/network-management-event-stream.service.ts create mode 100644 test/admin/network-management/network-management-event-stream.service.spec.ts diff --git a/API.md b/API.md index cfb44fb..63c292d 100644 --- a/API.md +++ b/API.md @@ -90,18 +90,19 @@ Admin、Component、Dict、MinIO、Blog 管理、WordPress 管理和 QQBot 管 ## System 网络端口转发管理 -| 方法 | 路径 | 认证 | 说明 | -| -------- | -------------------------------------------------------------- | ------- | ----------------------------------------- | -| `GET` | `/system/network/port-forward/list` | `super` | 分页查询期望、同步、Keeper 与端点租约状态 | -| `POST` | `/system/network/port-forward` | `super` | 新增 TCP/UDP 期望记录 | -| `PUT` | `/system/network/port-forward/:id` | `super` | 修改名称、协议和端口期望 | -| `DELETE` | `/system/network/port-forward/:id` | `super` | 写入 absent tombstone,等待 Agent 确认 | -| `POST` | `/system/network/port-forward/:id/retry` | `super` | 提升 revision 并重试协调 | -| `POST` | `/system/network/port-forward/:id/keeper/enable` | `super` | 启用同源端口 UDP Keeper 并立即探测 | -| `POST` | `/system/network/port-forward/:id/keeper/disable` | `super` | 停用 UDP Keeper 并撤下当前端点 | -| `POST` | `/system/network/port-forward/:id/probe` | `super` | 为已启用 Keeper 生成新 probeRequestId | -| `GET` | `/system/network/port-forward/:id/endpoint-history` | `super` | 查询端点状态变化历史 | -| `GET` | `/system/network/agent/status` | `super` | 查询 Agent 在线与 revision 收敛状态 | +| 方法 | 路径 | 认证 | 说明 | +| -------- | --------------------------------------------------- | ------- | ----------------------------------------- | +| `GET` | `/system/network/port-forward/list` | `super` | 分页查询期望、同步、Keeper 与端点租约状态 | +| `POST` | `/system/network/port-forward` | `super` | 新增 TCP/UDP 期望记录 | +| `PUT` | `/system/network/port-forward/:id` | `super` | 修改名称、协议和端口期望 | +| `DELETE` | `/system/network/port-forward/:id` | `super` | 写入 absent tombstone,等待 Agent 确认 | +| `POST` | `/system/network/port-forward/:id/retry` | `super` | 提升 revision 并重试协调 | +| `POST` | `/system/network/port-forward/:id/keeper/enable` | `super` | 启用同源端口 UDP Keeper 并立即探测 | +| `POST` | `/system/network/port-forward/:id/keeper/disable` | `super` | 停用 UDP Keeper 并撤下当前端点 | +| `POST` | `/system/network/port-forward/:id/probe` | `super` | 为已启用 Keeper 生成新 probeRequestId | +| `GET` | `/system/network/port-forward/:id/endpoint-history` | `super` | 查询端点状态变化历史 | +| `GET` | `/system/network/agent/status` | `super` | 查询 Agent 在线与 revision 收敛状态 | +| `GET` | `/system/network/events/stream` | `super` | SSE 推送已提交的 MQTT 状态变化 | 新增和修改请求只接受名称、备注、`tcp|udp`、外部端口和内部端口;目标 NAS IPv4 固定来自 `NETWORK_AGENT_TARGET_IPV4`,请求体中的未知字段会返回 400。Snowflake ID 与 revision 在 HTTP JSON 中保留为字符串。所有动态响应设置 `Cache-Control: no-store`。 @@ -111,6 +112,8 @@ API 数据库是唯一事实源。每次合法期望变更在同一事务中锁 Wire contract 采用 `kt-network-agent/internal/contract` schema-v1:desired mapping 使用 `state=present|absent`;reported 使用 `appliedRevision`、`desiredDigest`、helper 状态和逐条 router/route/Keeper 证据;endpoint event 使用唯一 `eventId`。API 不发布数据库备注,也不在 MQTT、HTTP 或数据库中接收/保存路由器密码和 token。事件按 `eventId` 幂等追加;删除只有在 Agent 明确回报 absent、synced、router/route 均不存在、Keeper 期望关闭且实际 disabled、current endpoint 为空,并且 helper 已确认且 `helperAppliedRevision === appliedRevision` 后才完成,随后生成新 revision 移除 tombstone。 +Admin 首次进入网络管理页通过 HTTP 读取快照,随后使用 `/system/network/events/stream` 接收 `network-state-changed`。API 只在 `reported`、`status`、`events` 对应事务提交且持久化可见状态实际变化后发出事件;MQTT QoS 1 重投和相同状态不会触发刷新。SSE 心跳只维持连接并复用最近一次真实状态事件 ID,尚无状态事件时显式发送空 ID,避免 Nest 自动生成游标;浏览器重连通过 `Last-Event-ID` 或 `lastEventId` 重放有限窗口,游标失效时收到一次 `snapshot-required` 并重新读取 HTTP 快照。前端不直接订阅 MQTT,也不使用定时轮询。 + TCP 记录支持 API CRUD,但不提供 STUN/Keeper;当前已验证的 Agent 切片尚未启用 TCP 路由器写入,会明确回报 `tcp_router_write_gated`,不能把 pending/failed TCP 记录描述为已生效转发。UDP 只有 `externalPort === internalPort` 时可启用 Keeper。当前端点只有在 `currentValidUntil` 未过期时才返回为可用值;租约过期不会删除最近观测或 `network_endpoint_history`。API 不直接访问小米路由器、不修改 NAS 路由;真实路由器、raw UDP 与回程规则只由固定 NAS Agent/helper 处理。 ## 环境变量分组 @@ -126,7 +129,7 @@ TCP 记录支持 API CRUD,但不提供 STUN/Keeper;当前已验证的 Agent | NapCat | `NAPCAT_WEBUI_BASE_URL`、`NAPCAT_WEBUI_TOKEN`、`QQBOT_NAPCAT_*` | | MQTT | `MQTT_URL`、`MQTT_USERNAME`、`MQTT_PASSWORD`、`MQTT_CLIENT_ID` | | Env Dashboard | `ENV_DASHBOARD_CACHE_TTL_MS`、`ENV_DASHBOARD_SIGNAL_TIMEOUT_MS`、`ENV_DASHBOARD_EVENT_BUS`、`ENV_DASHBOARD_MQTT_*`、`ENV_DASHBOARD_SSE_*`、`ENV_DASHBOARD_JENKINS_*`、`ENV_DASHBOARD_K8S_*`、`ENV_DASHBOARD_TENCENT_*`、`ENV_DASHBOARD_CADDY_*`、`ENV_DASHBOARD_R4SE_*` | -| Network | `NETWORK_AGENT_ID`、`NETWORK_AGENT_TARGET_IPV4`、`NETWORK_AGENT_MQTT_URL`、`NETWORK_AGENT_MQTT_CLIENT_ID`、`NETWORK_AGENT_MQTT_USERNAME`、`NETWORK_AGENT_MQTT_PASSWORD`、`NETWORK_AGENT_MQTT_RETRY_MS` | +| Network | `NETWORK_AGENT_ID`、`NETWORK_AGENT_TARGET_IPV4`、`NETWORK_AGENT_MQTT_URL`、`NETWORK_AGENT_MQTT_CLIENT_ID`、`NETWORK_AGENT_MQTT_USERNAME`、`NETWORK_AGENT_MQTT_PASSWORD`、`NETWORK_AGENT_MQTT_RETRY_MS`、`NETWORK_MANAGEMENT_SSE_HEARTBEAT_MS`、`NETWORK_MANAGEMENT_SSE_REPLAY_LIMIT` | | BangDream | `BANGDREAM_TSUGU_MAIN_SERVER`、`BANGDREAM_TSUGU_DISPLAYED_SERVERS`、`BANGDREAM_TSUGU_CACHE_ROOT` | | FF14 Market | `FF14_XIVAPI_BASE_URL`、`FF14_UNIVERSALIS_BASE_URL`、`FF14_DEFAULT_WORLD` | | FFLogs | `FFLOGS_GRAPHQL_URL`、`FFLOGS_TOKEN_URL`、`FFLOGS_CLIENT_ID`、`FFLOGS_CLIENT_SECRET` | diff --git a/README.md b/README.md index 6986f5e..34eac21 100644 --- a/README.md +++ b/README.md @@ -65,7 +65,7 @@ ci/ Jenkins Agent/Docker 辅助文件 | Logging/Loki | `LOG_LEVEL`、`LOG_APP_NAME`、`LOKI_URL`、`LOKI_QUERY_HOST`、`LOKI_*` | | QQBot/NapCat | `QQBOT_ENABLED`、`QQBOT_ACCOUNT_SECRET_KEY`、`QQBOT_REVERSE_WS_*`、`QQBOT_SEND_*`、`QQBOT_PLUGIN_QUEUE_REDIS_*`、`QQBOT_PLUGIN_TASK_QUEUE_REDIS_*`、`QQBOT_PLUGIN_QUEUE_WAIT_TIMEOUT_MS`、`QQBOT_COMMAND_MIN_COOLDOWN_MS`、`QQBOT_RULE_MIN_COOLDOWN_MS`、`QQBOT_REPEATER_*`、`NAPCAT_*`、`QQBOT_NAPCAT_*`、`MQTT_*` | | Environment Dashboard | `ENV_DASHBOARD_CACHE_TTL_MS`、`ENV_DASHBOARD_SIGNAL_TIMEOUT_MS`、`ENV_DASHBOARD_EVENT_BUS`、`ENV_DASHBOARD_MQTT_*`、`ENV_DASHBOARD_SSE_*`、`ENV_DASHBOARD_JENKINS_*`、`ENV_DASHBOARD_K8S_*`、`ENV_DASHBOARD_TENCENT_*`、`ENV_DASHBOARD_CADDY_*`、`ENV_DASHBOARD_R4SE_*` | -| Network Management | `NETWORK_AGENT_ID`、`NETWORK_AGENT_TARGET_IPV4`、`NETWORK_AGENT_MQTT_URL`、`NETWORK_AGENT_MQTT_CLIENT_ID`、`NETWORK_AGENT_MQTT_USERNAME`、`NETWORK_AGENT_MQTT_PASSWORD`、`NETWORK_AGENT_MQTT_RETRY_MS` | +| Network Management | `NETWORK_AGENT_ID`、`NETWORK_AGENT_TARGET_IPV4`、`NETWORK_AGENT_MQTT_URL`、`NETWORK_AGENT_MQTT_CLIENT_ID`、`NETWORK_AGENT_MQTT_USERNAME`、`NETWORK_AGENT_MQTT_PASSWORD`、`NETWORK_AGENT_MQTT_RETRY_MS`、`NETWORK_MANAGEMENT_SSE_HEARTBEAT_MS`、`NETWORK_MANAGEMENT_SSE_REPLAY_LIMIT` | | BangDream | `BANGDREAM_TSUGU_MAIN_SERVER`、`BANGDREAM_TSUGU_DISPLAYED_SERVERS`、`BANGDREAM_TSUGU_CACHE_ROOT` | | FF14 Market | `FF14_XIVAPI_BASE_URL`、`FF14_UNIVERSALIS_BASE_URL`、`FF14_MARKET_CACHE_TTL_MS` | | FFLogs | `FFLOGS_BASE_URL`、`FFLOGS_GRAPHQL_URL`、`FFLOGS_TOKEN_URL`、`FFLOGS_CLIENT_ID`、`FFLOGS_CLIENT_SECRET` | @@ -80,7 +80,7 @@ QQBot 插件定时任务由 manifest 的 `tasks` 声明,平台持久化到 `qq Admin 环境总览面板使用 `ENV_DASHBOARD_*` 只读配置聚合 local-dev、NAS 线上、腾讯云和 r4se 状态。`ENV_DASHBOARD_ADMIN_LOCAL_URL` / `ENV_DASHBOARD_ADMIN_PUBLIC_URL` 只用于展示 Admin 本机与线上入口证据。HTTP 快照提供当前拓扑,后端 local/MQTT 事件总线通过 SSE 推送增量事件给 Admin;前端不直连 MQTT,也不轮询刷新。Jenkins、K8s、Tencent Cloud、Caddy、WireGuard、Mihomo/OpenClash 未配置时会显示 `unwired` 证据,不能渲染成健康假象;第一版不暴露重启、部署、迁移、容器重建、插件启停或代理切换等写操作。 -System 网络管理以 MySQL 中的 TCP/UDP 端口转发期望状态为唯一事实源。`super` 通过统一 CRUD 和 UDP Keeper 动作修改期望状态;API 在事务内单调提升 revision,提交后使用固定 `kt/network/v1/agents/{agentId}` MQTT topic、QoS 1 retained 完整快照通知 NAS `kt-network-agent`,自身不登录路由器、不接收路由器密码,也不执行 raw socket。Agent 失联或 MQTT 暂不可用时合法请求仍保存为 pending,恢复后按 revision 自动收敛;消费端发生瞬时数据库错误或 SUBACK 失败时主动重连并依赖 broker 重投,非法负载则确认后丢弃,避免 poison message 阻塞。TCP 目前仅保存 CRUD 期望,真实路由器写入仍受设备协议证据门禁并回报 `tcp_router_write_gated`;只有外部端口等于内部端口的 UDP 记录允许启停 Keeper 和立即探测。当前公网端点受 `currentValidUntil` 租约约束,过期后列表隐藏当前值但保留最近观测与历史。生产发布同时把完整 `NETWORK_AGENT_*` 连接配置作为 Jenkins 私有 env 和 `/health/runtime` 必需项,任一项缺失时拒绝发布或报告运行态阻断,避免页面可见但 MQTT 控制链路未接线。 +System 网络管理以 MySQL 中的 TCP/UDP 端口转发期望状态为唯一事实源。`super` 通过统一 CRUD 和 UDP Keeper 动作修改期望状态;API 在事务内单调提升 revision,提交后使用固定 `kt/network/v1/agents/{agentId}` MQTT topic、QoS 1 retained 完整快照通知 NAS `kt-network-agent`,自身不登录路由器、不接收路由器密码,也不执行 raw socket。Agent 失联或 MQTT 暂不可用时合法请求仍保存为 pending,恢复后按 revision 自动收敛;消费端发生瞬时数据库错误或 SUBACK 失败时主动重连并依赖 broker 重投,非法负载则确认后丢弃,避免 poison message 阻塞。API 仅在入站 MQTT 事务提交且可见状态实际变化后通过 `/system/network/events/stream` 向 Admin 发布 SSE;QoS 1 幂等重投与心跳均不触发页面刷新。心跳复用最近一次真实状态事件游标,尚无状态事件时显式发送空游标,避免 Nest 自动生成的 ID 污染重放位置;有限重放缺口只要求一次 HTTP 快照。TCP 目前仅保存 CRUD 期望,真实路由器写入仍受设备协议证据门禁并回报 `tcp_router_write_gated`;只有外部端口等于内部端口的 UDP 记录允许启停 Keeper 和立即探测。当前公网端点受 `currentValidUntil` 租约约束,过期后列表隐藏当前值但保留最近观测与历史。生产发布同时把完整 `NETWORK_AGENT_*` 连接配置作为 Jenkins 私有 env 和 `/health/runtime` 必需项,任一项缺失时拒绝发布或报告运行态阻断,避免页面可见但 MQTT 控制链路未接线。 NapCat Runtime/Protocol Profile 已完成本地 API/Admin 实施,线上发布和账号闭环按 `docs/plans/2026-06-18-qqbot-napcat-runtime-protocol-profile-implementation-plan.md` 的 Task 10 执行。当前实现覆盖运行态/协议/会话行为/历史登录事件兼容表/风险模式表,真实物理设备风格 hostname/MAC,NapCat/OneBot 配置 hash,KT `zh_CN.UTF-8` 中国桌面派生镜像资产,只读 `/qqbot/napcat/runtime/detail` 证据接口,watchdog 离线巡检告警,以及 Admin 账号页“运行态”抽屉;不绕过 QQ/Tencent 验证码、不修改 QQ/NTQQ 签名协议、不启用 privileged/host network,也不做账号级每小时/每日累计发送预算。NapCat Chinese Desktop Runtime v20 使用 KT `NapCatQQ` fork 源码构建出的 `NapCat.Shell` artifact,并在 QQ `KickedOffLine` 后标记 native login service stale;API 在源 Docker 容器在线但 WebUI 明确 QQ 离线时会同容器调用 `RestartNapCat` 重启 NapCat worker,重建 QQCore login service 后再推进 quick/password/qrcode,不做 Docker 重建、补 env 或设备身份迁移,且同一个更新登录 session 只消费一次 worker restart 预算;v14 起还会对 QQ/NapCat/Xvfb 长期进程的 `/proc//mountinfo` 做 PID 级遮蔽,防止 `overlay`、`/vol1/docker`、`docker-init`、`/docker/containers`、`napcat-instances` 等宿主路径泄露;v15 修复扫码成功时 `QQLoginInfo` 晚于登录态写入造成的 QQ 号回读空窗;v16 在 native reset 缺少 `offline()` 时改用 `destroy()` 硬重置半登录服务,并让镜像 verify 等待 mountinfo guard 收敛;v17/v18 增加 WebUI 鉴权的 `/api/Debug/RuntimeViewProbe` 同进程诊断并修正 native maps 截断导致的 hook 证据假阴性;v19 保留 WebUI `RestartNapCat` 重启 worker 时的 `-q ` 快速登录参数,避免重启后退回无账号扫码;v20 保护 API 预写的 `/app/napcat/config`,避免上游首次解包 `NapCat.Shell/*` 覆盖 `bypass.*=true` 与 `o3HookMode=0`。镜像必须先用 `scripts/napcat-desktop-cn-stage-build.mjs` staged build context,生产 `QQBOT_NAPCAT_IMAGE` 应指向验证过的 `kt-napcat-desktop-cn:desktop-cn-v20` digest。`k8s/prod/api.yaml` 保留 `desktop-cn-v20` 稳定默认值;Jenkins `QQBOT_NAPCAT_IMAGE_OVERRIDE` 和 `QQBOT_NAPCAT_DESKTOP_PROFILE_VERSION_OVERRIDE` 仅在填写时通过 `kubectl set env` 推广已验证运行时镜像/profile,空值会继续使用 manifest/default env。回滚时重新运行 Jenkins 并填入上一版 digest/profile,或清空两个 override 后重新部署 manifest 默认值。 diff --git a/src/modules/admin/platform-config/admin-platform-config.module.ts b/src/modules/admin/platform-config/admin-platform-config.module.ts index 1b5896c..82acd6b 100644 --- a/src/modules/admin/platform-config/admin-platform-config.module.ts +++ b/src/modules/admin/platform-config/admin-platform-config.module.ts @@ -12,6 +12,7 @@ import { NetworkAgentMqttService } from '@/modules/admin/platform-config/network import { NetworkAgentState } from '@/modules/admin/platform-config/network-management/network-agent-state.entity'; import { NetworkEndpointHistory } from '@/modules/admin/platform-config/network-management/network-endpoint-history.entity'; import { NetworkManagementController } from '@/modules/admin/platform-config/network-management/network-management.controller'; +import { NetworkManagementEventStreamService } from '@/modules/admin/platform-config/network-management/network-management-event-stream.service'; import { NetworkPortForward } from '@/modules/admin/platform-config/network-management/network-management.entity'; import { NetworkManagementService } from '@/modules/admin/platform-config/network-management/network-management.service'; import { SystemLogController } from '@/modules/admin/platform-config/system-log/system-log.controller'; @@ -81,6 +82,7 @@ export const ADMIN_PLATFORM_CONFIG_PROVIDERS = [ EnvironmentEventMaterializer, EnvironmentEventStreamService, NetworkManagementService, + NetworkManagementEventStreamService, NetworkAgentMqttService, ]; diff --git a/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts b/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts index 0c020f7..774dc36 100644 --- a/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts +++ b/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts @@ -13,6 +13,7 @@ import type { IClientOptions, MqttClient } from 'mqtt'; import { KtDateTime } from '@/common'; import { NetworkAgentState } from './network-agent-state.entity'; import { NetworkEndpointHistory } from './network-endpoint-history.entity'; +import { NetworkManagementEventStreamService } from './network-management-event-stream.service'; import { NetworkPortForward } from './network-management.entity'; import { buildDesiredSnapshot, @@ -24,6 +25,7 @@ import { parseStatusSnapshot, type NetworkEndpointEvent, type NetworkReportedSnapshot, + type NetworkStateChangeSource, type NetworkStatusSnapshot, } from './network-management.types'; @@ -56,11 +58,13 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { * Creates the dedicated network Agent MQTT bridge. * @param configService - Runtime broker, Agent, and client identity settings. * @param dataSource - Transaction boundary for publish acknowledgements and inbound state. + * @param eventStream - SSE fan-out notified only after accepted inbound commits. * @param clientFactory - Optional deterministic MQTT client factory used by tests. */ constructor( private readonly configService: ConfigService, private readonly dataSource: DataSource, + private readonly eventStream: NetworkManagementEventStreamService, @Optional() @Inject(NETWORK_MQTT_CLIENT_FACTORY) private readonly clientFactory?: NetworkMqttClientFactory, @@ -170,19 +174,21 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { throw new NetworkMessageValidationError('Invalid network MQTT JSON'); } + let changed = false; + let source: NetworkStateChangeSource; if (topic === this.topic('reported')) { - await this.applyReported(parseReportedSnapshot(parsed)); - return; + source = 'reported'; + changed = await this.applyReported(parseReportedSnapshot(parsed)); + } else if (topic === this.topic('status')) { + source = 'status'; + changed = await this.applyStatus(parseStatusSnapshot(parsed)); + } else if (topic === this.topic('events')) { + source = 'events'; + changed = await this.appendEndpointEvent(parseEndpointEvent(parsed)); + } else { + throw new NetworkMessageValidationError('Unexpected network MQTT topic'); } - if (topic === this.topic('status')) { - await this.applyStatus(parseStatusSnapshot(parsed)); - return; - } - if (topic === this.topic('events')) { - await this.appendEndpointEvent(parseEndpointEvent(parsed)); - return; - } - throw new NetworkMessageValidationError('Unexpected network MQTT topic'); + if (changed) this.eventStream.publishCommitted(source); } /** Restores exact subscriptions and one retained desired republish after each connection. */ @@ -306,145 +312,160 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { }); } - /** Applies a full reported snapshot transactionally without creating desired rows. */ - private async applyReported(report: NetworkReportedSnapshot): Promise { + /** + * Applies a full reported snapshot transactionally without creating desired rows. + * @param report - Strict Agent report parsed from the retained topic. + * @returns True only when persisted Admin-visible state changed. + */ + private async applyReported( + report: NetworkReportedSnapshot, + ): Promise { this.assertAgentId(report.agentId); - const desiredChanged = await this.dataSource.transaction( - async (manager) => { - const stateRepository = manager.getRepository(NetworkAgentState); - const mappingRepository = manager.getRepository(NetworkPortForward); - const state = await stateRepository.findOne({ - lock: { mode: 'pessimistic_write' }, - where: { agentId: report.agentId }, + const result = await this.dataSource.transaction(async (manager) => { + const stateRepository = manager.getRepository(NetworkAgentState); + const mappingRepository = manager.getRepository(NetworkPortForward); + const state = await stateRepository.findOne({ + lock: { mode: 'pessimistic_write' }, + where: { agentId: report.agentId }, + }); + if ( + !state || + BigInt(report.appliedRevision) > BigInt(state.desiredRevision) + ) { + throw new NetworkMessageValidationError('Invalid reported revision'); + } + if (BigInt(report.appliedRevision) < BigInt(state.appliedRevision)) { + return { desiredChanged: false, visibleStateChanged: false }; + } + const stateBefore = this.reportedAgentStateFingerprint(state); + let visibleStateChanged = false; + const isCurrentRevision = + BigInt(report.appliedRevision) === BigInt(state.desiredRevision); + const desiredMappings = await mappingRepository.find({ + where: { isDeleted: false }, + }); + if (isCurrentRevision) { + const currentSnapshot = buildDesiredSnapshot(state, desiredMappings); + if (desiredSnapshotDigest(currentSnapshot) !== report.desiredDigest) { + throw new NetworkMessageValidationError( + 'Reported desired digest does not match current revision', + ); + } + } + const mappingById = new Map( + desiredMappings.map((mapping) => [mapping.id, mapping]), + ); + const finalizedIds = new Set(); + let finalizedDeletion = false; + for (const item of report.mappings) { + const mapping = mappingById.get(item.id); + if (mapping) continue; + if (isCurrentRevision) { + throw new NetworkMessageValidationError( + 'Unknown mapping in current reported snapshot', + ); + } + const historical = await mappingRepository.findOne({ + where: { id: item.id }, }); if ( - !state || - BigInt(report.appliedRevision) > BigInt(state.desiredRevision) + historical?.isDeleted && + item.desiredState === 'absent' && + item.syncStatus === 'synced' ) { - throw new NetworkMessageValidationError('Invalid reported revision'); + finalizedIds.add(item.id); + } else { + throw new NetworkMessageValidationError('Unknown reported mapping'); } - if (BigInt(report.appliedRevision) < BigInt(state.appliedRevision)) { - return; - } - const isCurrentRevision = - BigInt(report.appliedRevision) === BigInt(state.desiredRevision); - const desiredMappings = await mappingRepository.find({ - where: { isDeleted: false }, - }); - if (isCurrentRevision) { - const currentSnapshot = buildDesiredSnapshot(state, desiredMappings); - if (desiredSnapshotDigest(currentSnapshot) !== report.desiredDigest) { - throw new NetworkMessageValidationError( - 'Reported desired digest does not match current revision', - ); - } - } - const mappingById = new Map( - desiredMappings.map((mapping) => [mapping.id, mapping]), - ); - const finalizedIds = new Set(); - let finalizedDeletion = false; - for (const item of report.mappings) { - const mapping = mappingById.get(item.id); - if (mapping) continue; - if (isCurrentRevision) { - throw new NetworkMessageValidationError( - 'Unknown mapping in current reported snapshot', - ); - } - const historical = await mappingRepository.findOne({ - where: { id: item.id }, - }); - if ( - historical?.isDeleted && - item.desiredState === 'absent' && - item.syncStatus === 'synced' - ) { - finalizedIds.add(item.id); - } else { - throw new NetworkMessageValidationError('Unknown reported mapping'); - } - } - if (isCurrentRevision) { - const reportedIds = new Set(report.mappings.map((item) => item.id)); - if (desiredMappings.some((mapping) => !reportedIds.has(mapping.id))) { - throw new NetworkMessageValidationError( - 'Incomplete current reported snapshot', - ); - } - } - - for (const item of report.mappings) { - if (finalizedIds.has(item.id)) continue; - const mapping = mappingById.get(item.id) as NetworkPortForward; - if (BigInt(item.revision) < BigInt(mapping.desiredRevision)) { - continue; - } - if (item.desiredState !== mapping.desiredPresence) { - throw new NetworkMessageValidationError( - 'Reported mapping desired state does not match', - ); - } - if (item.keeperDesiredEnabled !== mapping.keeperDesiredEnabled) { - throw new NetworkMessageValidationError( - 'Reported mapping Keeper intent does not match', - ); - } - mapping.reportedRevision = String(item.revision); - mapping.syncStatus = item.syncStatus; - mapping.keeperStatus = item.keeperStatus; - mapping.lastErrorCode = item.errorCode || null; - mapping.lastErrorMessage = item.errorMessage || null; - this.applyReportedEndpoints( - mapping, - item.currentEndpoint, - item.lastObservedEndpoint, - report.reportedAt, + } + if (isCurrentRevision) { + const reportedIds = new Set(report.mappings.map((item) => item.id)); + if (desiredMappings.some((mapping) => !reportedIds.has(mapping.id))) { + throw new NetworkMessageValidationError( + 'Incomplete current reported snapshot', ); + } + } - if ( - mapping.desiredPresence === 'absent' && - item.desiredState === 'absent' && - item.syncStatus === 'synced' && - item.routerPresent === false && - item.routePresent === false && - report.helperStatus === 'confirmed' && - report.helperAppliedRevision === report.appliedRevision && - item.keeperDesiredEnabled === false && - item.keeperStatus === 'disabled' && - !item.currentEndpoint - ) { - mapping.activeKey = null; - mapping.isDeleted = true; - finalizedDeletion = true; - } + for (const item of report.mappings) { + if (finalizedIds.has(item.id)) continue; + const mapping = mappingById.get(item.id) as NetworkPortForward; + if (BigInt(item.revision) < BigInt(mapping.desiredRevision)) { + continue; + } + if (item.desiredState !== mapping.desiredPresence) { + throw new NetworkMessageValidationError( + 'Reported mapping desired state does not match', + ); + } + if (item.keeperDesiredEnabled !== mapping.keeperDesiredEnabled) { + throw new NetworkMessageValidationError( + 'Reported mapping Keeper intent does not match', + ); + } + const mappingBefore = this.reportedMappingStateFingerprint(mapping); + mapping.reportedRevision = String(item.revision); + mapping.syncStatus = item.syncStatus; + mapping.keeperStatus = item.keeperStatus; + mapping.lastErrorCode = item.errorCode || null; + mapping.lastErrorMessage = item.errorMessage || null; + this.applyReportedEndpoints( + mapping, + item.currentEndpoint, + item.lastObservedEndpoint, + report.reportedAt, + ); + + if ( + mapping.desiredPresence === 'absent' && + item.desiredState === 'absent' && + item.syncStatus === 'synced' && + item.routerPresent === false && + item.routePresent === false && + report.helperStatus === 'confirmed' && + report.helperAppliedRevision === report.appliedRevision && + item.keeperDesiredEnabled === false && + item.keeperStatus === 'disabled' && + !item.currentEndpoint + ) { + mapping.activeKey = null; + mapping.isDeleted = true; + finalizedDeletion = true; + } + if (this.reportedMappingStateFingerprint(mapping) !== mappingBefore) { + visibleStateChanged = true; await mappingRepository.save(mapping); } + } - if (BigInt(report.appliedRevision) > BigInt(state.appliedRevision)) { - state.appliedRevision = String(report.appliedRevision); - } - const failedMapping = report.mappings.find( - (item) => - item.syncStatus === 'conflict' || item.syncStatus === 'failed', - ); - state.lastReconcileErrorCode = failedMapping - ? failedMapping.errorCode || `sync_${failedMapping.syncStatus}` - : report.helperStatus === 'failed' - ? 'route_helper_failed' - : null; - state.lastReconcileErrorMessage = failedMapping?.errorMessage || null; - if (finalizedDeletion) { - state.desiredRevision = ( - BigInt(state.desiredRevision) + 1n - ).toString(); - state.desiredIssuedAt = new KtDateTime(); - } + if (BigInt(report.appliedRevision) > BigInt(state.appliedRevision)) { + state.appliedRevision = String(report.appliedRevision); + } + const failedMapping = report.mappings.find( + (item) => + item.syncStatus === 'conflict' || item.syncStatus === 'failed', + ); + state.lastReconcileErrorCode = failedMapping + ? failedMapping.errorCode || `sync_${failedMapping.syncStatus}` + : report.helperStatus === 'failed' + ? 'route_helper_failed' + : null; + state.lastReconcileErrorMessage = failedMapping?.errorMessage || null; + if (finalizedDeletion) { + state.desiredRevision = (BigInt(state.desiredRevision) + 1n).toString(); + state.desiredIssuedAt = new KtDateTime(); + } + if (this.reportedAgentStateFingerprint(state) !== stateBefore) { + visibleStateChanged = true; await stateRepository.save(state); - return finalizedDeletion; - }, - ); - if (desiredChanged) this.requestDesiredPublish(); + } + return { + desiredChanged: finalizedDeletion, + visibleStateChanged, + }; + }); + if (result.desiredChanged) this.requestDesiredPublish(); + return result.visibleStateChanged; } /** @@ -496,10 +517,14 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { } } - /** Applies retained Agent online status without changing mapping sync semantics. */ - private async applyStatus(status: NetworkStatusSnapshot): Promise { + /** + * Applies retained Agent online status without changing mapping sync semantics. + * @param status - Strict retained status or LWT snapshot. + * @returns True only when persisted Agent status changed. + */ + private async applyStatus(status: NetworkStatusSnapshot): Promise { this.assertAgentId(status.agentId); - await this.dataSource.transaction(async (manager) => { + return await this.dataSource.transaction(async (manager) => { const repository = manager.getRepository(NetworkAgentState); const state = await repository.findOne({ lock: { mode: 'pessimistic_write' }, @@ -520,7 +545,7 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { currentStartedAt && incomingStartedAt.getTime() < currentStartedAt.getTime() ) { - return; + return false; } const isSameSessionWill = status.online === false && @@ -532,8 +557,9 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { observedAt.getTime() < new Date(state.lastHeartbeatAt).getTime() && !isSameSessionWill ) { - return; + return false; } + const stateBefore = this.statusStateFingerprint(state); state.online = status.online; state.version = status.version || null; state.startedAt = incomingStartedAt @@ -544,16 +570,22 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { } state.lastMqttErrorCode = status.errorCode || null; state.lastMqttErrorMessage = status.errorMessage || null; + if (this.statusStateFingerprint(state) === stateBefore) return false; await repository.save(state); + return true; }); } - /** Appends one endpoint change event exactly once by event ID. */ + /** + * Appends one endpoint change event exactly once by event ID. + * @param event - Strict endpoint transition parsed from the events topic. + * @returns True only when a new history row committed. + */ private async appendEndpointEvent( event: NetworkEndpointEvent, - ): Promise { + ): Promise { this.assertAgentId(event.agentId); - await this.dataSource.transaction(async (manager) => { + return await this.dataSource.transaction(async (manager) => { const state = await manager.getRepository(NetworkAgentState).findOne({ where: { agentId: event.agentId }, }); @@ -568,7 +600,7 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { } const repository = manager.getRepository(NetworkEndpointHistory); if (await repository.findOne({ where: { eventId: event.eventId } })) { - return; + return false; } const history = repository.create({ eventId: event.eventId, @@ -583,12 +615,69 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { }); try { await repository.save(history); + return true; } catch (error) { if (!this.isDuplicateKeyError(error)) throw error; + return false; } }); } + /** + * Serializes only mapping fields rendered by Admin or used for row availability. + * @param mapping - Current persisted mapping before or after one report application. + * @returns Stable comparison string for suppressing QoS 1 redelivery refreshes. + */ + private reportedMappingStateFingerprint(mapping: NetworkPortForward): string { + return JSON.stringify([ + mapping.activeKey, + mapping.currentObservedAt, + mapping.currentPublicIpv4, + mapping.currentPublicPort, + mapping.currentValidUntil, + mapping.isDeleted, + mapping.keeperStatus, + mapping.lastErrorCode, + mapping.lastErrorMessage, + mapping.lastObservedAt, + mapping.lastObservedIpv4, + mapping.lastObservedPort, + mapping.reportedRevision, + mapping.syncStatus, + ]); + } + + /** + * Serializes report-owned Agent fields shown by the network page. + * @param state - Agent singleton before or after one reported snapshot. + * @returns Stable comparison string for committed report changes. + */ + private reportedAgentStateFingerprint(state: NetworkAgentState): string { + return JSON.stringify([ + state.appliedRevision, + state.desiredIssuedAt, + state.desiredRevision, + state.lastReconcileErrorCode, + state.lastReconcileErrorMessage, + ]); + } + + /** + * Serializes status-topic fields shown by the network page. + * @param state - Agent singleton before or after one status snapshot. + * @returns Stable comparison string for duplicate status suppression. + */ + private statusStateFingerprint(state: NetworkAgentState): string { + return JSON.stringify([ + state.lastHeartbeatAt, + state.lastMqttErrorCode, + state.lastMqttErrorMessage, + state.online, + state.startedAt, + state.version, + ]); + } + /** Validates that an inbound message belongs to the one configured Agent. */ private assertAgentId(agentId: string): void { if (agentId !== this.agentId()) { diff --git a/src/modules/admin/platform-config/network-management/network-management-event-stream.service.ts b/src/modules/admin/platform-config/network-management/network-management-event-stream.service.ts new file mode 100644 index 0000000..9b53234 --- /dev/null +++ b/src/modules/admin/platform-config/network-management/network-management-event-stream.service.ts @@ -0,0 +1,122 @@ +import { Injectable, Optional } from '@nestjs/common'; +import { merge, Observable, of, Subject, timer } from 'rxjs'; +import { map } from 'rxjs/operators'; +import type { + NetworkStateChangeEvent, + NetworkStateChangeSource, +} from './network-management.types'; + +export interface NetworkManagementEventStreamOptions { + heartbeatMs?: number; + replayLimit?: number; +} + +export interface NetworkManagementStreamEvent { + data: NetworkStateChangeEvent | { message: string; observedAt: string }; + id: string; + type: 'heartbeat' | 'network-state-changed' | 'snapshot-required'; +} + +@Injectable() +export class NetworkManagementEventStreamService { + private readonly replay: NetworkManagementStreamEvent[] = []; + private readonly streamSubject = new Subject(); + private readonly heartbeatMs: number; + private readonly replayLimit: number; + private eventSequence = 0; + + /** + * Creates the network-management SSE fan-out without exposing MQTT to browsers. + * @param options - Optional heartbeat and replay bounds used by runtime or tests. + */ + constructor(@Optional() options: NetworkManagementEventStreamOptions = {}) { + this.heartbeatMs = + options.heartbeatMs || + Number(process.env.NETWORK_MANAGEMENT_SSE_HEARTBEAT_MS) || + 25_000; + this.replayLimit = + options.replayLimit || + Number(process.env.NETWORK_MANAGEMENT_SSE_REPLAY_LIMIT) || + 100; + } + + /** + * Opens a replayable SSE stream plus heartbeats pinned to the committed cursor. + * @param lastEventId - Browser replay cursor from the previous committed state event. + * @returns Observable containing replay, live changes, or keepalive messages. + */ + stream(lastEventId?: string): Observable { + const replayEvents = this.getReplayEvents(lastEventId); + const heartbeat$ = timer(this.heartbeatMs, this.heartbeatMs).pipe( + map(() => this.createHeartbeatEvent()), + ); + return merge( + ...replayEvents.map((event) => of(event)), + this.streamSubject, + heartbeat$, + ); + } + + /** + * Publishes one browser-safe notice after an inbound MQTT transaction commits. + * @param source - Accepted Agent topic category that changed persisted state. + * @returns Exact stream event stored for reconnect replay. + */ + publishCommitted( + source: NetworkStateChangeSource, + ): NetworkManagementStreamEvent { + const observedAt = new Date().toISOString(); + const eventId = `network-${Date.now()}-${++this.eventSequence}`; + const event: NetworkManagementStreamEvent = { + data: { eventId, observedAt, source }, + id: eventId, + type: 'network-state-changed', + }; + this.replay.push(event); + if (this.replay.length > this.replayLimit) { + this.replay.splice(0, this.replay.length - this.replayLimit); + } + this.streamSubject.next(event); + return event; + } + + /** + * Resolves the bounded replay window without treating a first connection as stale. + * @param lastEventId - Last committed event applied by the current Admin page. + * @returns Subsequent events or one snapshot instruction when continuity is lost. + */ + private getReplayEvents( + lastEventId?: string, + ): NetworkManagementStreamEvent[] { + if (!lastEventId) return []; + const index = this.replay.findIndex((event) => event.id === lastEventId); + if (index === -1) return [this.createSnapshotRequiredEvent()]; + return this.replay.slice(index + 1); + } + + /** Creates a keepalive pinned to the latest committed replay cursor. */ + private createHeartbeatEvent(): NetworkManagementStreamEvent { + return { + data: { message: 'alive', observedAt: new Date().toISOString() }, + id: this.currentReplayCursor(), + type: 'heartbeat', + }; + } + + /** Creates a snapshot instruction aligned to the latest committed replay cursor. */ + private createSnapshotRequiredEvent(): NetworkManagementStreamEvent { + return { + data: { + message: 'snapshot-required', + observedAt: new Date().toISOString(), + }, + id: this.currentReplayCursor(), + type: 'snapshot-required', + }; + } + + /** Returns the latest real state-event ID or an explicit empty initial cursor. */ + private currentReplayCursor(): string { + return this.replay.at(-1)?.id || ''; + } +} diff --git a/src/modules/admin/platform-config/network-management/network-management.controller.ts b/src/modules/admin/platform-config/network-management/network-management.controller.ts index f8c91d1..52641a6 100644 --- a/src/modules/admin/platform-config/network-management/network-management.controller.ts +++ b/src/modules/admin/platform-config/network-management/network-management.controller.ts @@ -3,6 +3,7 @@ import { Controller, Delete, Get, + Headers, HttpCode, HttpStatus, Param, @@ -10,6 +11,7 @@ import { Put, Query, Res, + Sse, UseGuards, UsePipes, ValidationPipe, @@ -26,6 +28,7 @@ import { NetworkPortForwardUpdateDto, } from './network-management.dto'; import { NetworkManagementService } from './network-management.service'; +import { NetworkManagementEventStreamService } from './network-management-event-stream.service'; @ApiTags('Admin - 网络端口转发') @Controller('system/network') @@ -41,8 +44,27 @@ export class NetworkManagementController { /** * Creates the super-admin-only network desired-state controller. * @param service - Persisted port-forward and Agent state service. + * @param eventStream - Committed MQTT change stream exposed to Admin through SSE. */ - constructor(private readonly service: NetworkManagementService) {} + constructor( + private readonly service: NetworkManagementService, + private readonly eventStream: NetworkManagementEventStreamService, + ) {} + + /** + * Subscribes Admin to committed network-state changes without exposing MQTT credentials. + * @param lastEventIdHeader - Native EventSource replay cursor. + * @param lastEventIdQuery - Query fallback retained across keep-alive deactivation. + * @returns SSE observable containing typed state changes and cursor-pinned heartbeats. + */ + @Sse('events/stream') + @ApiOperation({ summary: '订阅网络管理状态变化' }) + stream( + @Headers('last-event-id') lastEventIdHeader?: string, + @Query('lastEventId') lastEventIdQuery?: string, + ) { + return this.eventStream.stream(lastEventIdHeader || lastEventIdQuery); + } /** Lists persisted mappings without returning expired endpoint leases. */ @Get('port-forward/list') diff --git a/src/modules/admin/platform-config/network-management/network-management.types.ts b/src/modules/admin/platform-config/network-management/network-management.types.ts index 5eade1a..c47ca92 100644 --- a/src/modules/admin/platform-config/network-management/network-management.types.ts +++ b/src/modules/admin/platform-config/network-management/network-management.types.ts @@ -28,6 +28,13 @@ export type EndpointEventType = | 'published' | 'restored' | 'withdrawn'; +export type NetworkStateChangeSource = 'events' | 'reported' | 'status'; + +export type NetworkStateChangeEvent = { + eventId: string; + observedAt: string; + source: NetworkStateChangeSource; +}; export type NetworkDesiredMapping = { externalPort: number; diff --git a/test/admin/network-management/network-agent-mqtt.service.spec.ts b/test/admin/network-management/network-agent-mqtt.service.spec.ts index 391d9ed..0cb08a1 100644 --- a/test/admin/network-management/network-agent-mqtt.service.spec.ts +++ b/test/admin/network-management/network-agent-mqtt.service.spec.ts @@ -9,6 +9,7 @@ import { } from '../../../src/modules/admin/platform-config/network-management/network-agent-mqtt.service'; import { NetworkAgentState } from '../../../src/modules/admin/platform-config/network-management/network-agent-state.entity'; import { NetworkEndpointHistory } from '../../../src/modules/admin/platform-config/network-management/network-endpoint-history.entity'; +import type { NetworkManagementEventStreamService } from '../../../src/modules/admin/platform-config/network-management/network-management-event-stream.service'; import { NetworkPortForward } from '../../../src/modules/admin/platform-config/network-management/network-management.entity'; import { buildDesiredSnapshot, @@ -19,6 +20,7 @@ type MqttHarness = { client: MqttClient & EventEmitter; clientOptions: () => IClientOptions; histories: NetworkEndpointHistory[]; + publishCommitted: jest.Mock; mapping: NetworkPortForward; publishCallback: () => (error?: Error) => void; service: NetworkAgentMqttService; @@ -127,9 +129,14 @@ function createHarness(): MqttHarness { options = clientOptions; return client; }; + const publishCommitted = jest.fn(); + const eventStream = { + publishCommitted, + } as unknown as NetworkManagementEventStreamService; const service = new NetworkAgentMqttService( configService, dataSource, + eventStream, factory, ); return { @@ -138,6 +145,7 @@ function createHarness(): MqttHarness { histories, mapping, publishCallback: () => publishAck, + publishCommitted, service, state, transactionCalls: () => transactionCallCount, @@ -464,6 +472,22 @@ describe('NetworkAgentMqttService', () => { ); }); + it('publishes one Admin event only after a reported transaction changes visible state', async () => { + const harness = createHarness(); + const topic = 'kt/network/v1/agents/nas-main/reported'; + const payload = reported(harness, 7); + harness.publishCommitted.mockImplementation(() => { + expect(harness.mapping.syncStatus).toBe('synced'); + expect(harness.state.appliedRevision).toBe('7'); + }); + + await harness.service.consumeMessage(topic, payload); + await harness.service.consumeMessage(topic, payload); + + expect(harness.publishCommitted).toHaveBeenCalledTimes(1); + expect(harness.publishCommitted).toHaveBeenCalledWith('reported'); + }); + it('does not let an out-of-order same-revision withdrawal erase a newer lease', async () => { const harness = createHarness(); const topic = 'kt/network/v1/agents/nas-main/reported'; @@ -664,6 +688,8 @@ describe('NetworkAgentMqttService', () => { await harness.service.consumeMessage(topic, payload); await harness.service.consumeMessage(topic, payload); expect(harness.histories).toHaveLength(1); + expect(harness.publishCommitted).toHaveBeenCalledTimes(1); + expect(harness.publishCommitted).toHaveBeenCalledWith('events'); }); it('accepts a same-instance LWT without regressing heartbeat and ignores an old-instance LWT', async () => { @@ -697,5 +723,7 @@ describe('NetworkAgentMqttService', () => { expect(harness.state.startedAt?.toISOString()).toBe( '2026-07-22T02:00:00.000Z', ); + expect(harness.publishCommitted).toHaveBeenCalledTimes(1); + expect(harness.publishCommitted).toHaveBeenCalledWith('status'); }); }); diff --git a/test/admin/network-management/network-management-event-stream.service.spec.ts b/test/admin/network-management/network-management-event-stream.service.spec.ts new file mode 100644 index 0000000..f2cf7a7 --- /dev/null +++ b/test/admin/network-management/network-management-event-stream.service.spec.ts @@ -0,0 +1,92 @@ +import { filter, firstValueFrom, take } from 'rxjs'; +import { NetworkManagementEventStreamService } from '../../../src/modules/admin/platform-config/network-management/network-management-event-stream.service'; + +describe('NetworkManagementEventStreamService', () => { + it('fans out committed MQTT changes and replays only events after the cursor', async () => { + const service = new NetworkManagementEventStreamService({ + heartbeatMs: 60_000, + replayLimit: 5, + }); + const liveEvent = firstValueFrom( + service.stream().pipe( + filter((event) => event.type === 'network-state-changed'), + take(1), + ), + ); + + const first = service.publishCommitted('reported'); + await expect(liveEvent).resolves.toEqual(first); + const second = service.publishCommitted('status'); + + await expect( + firstValueFrom(service.stream(first.id).pipe(take(1))), + ).resolves.toEqual(second); + }); + + it('requests one snapshot for an unknown replay cursor', async () => { + const service = new NetworkManagementEventStreamService({ + heartbeatMs: 60_000, + replayLimit: 1, + }); + const latest = service.publishCommitted('events'); + + const event = await firstValueFrom( + service.stream('missing-event').pipe(take(1)), + ); + + expect(event).toMatchObject({ type: 'snapshot-required' }); + expect(event.id).toBe(latest.id); + }); + + it('keeps heartbeat on the latest committed replay cursor', async () => { + jest.useFakeTimers(); + try { + const service = new NetworkManagementEventStreamService({ + heartbeatMs: 1_000, + replayLimit: 1, + }); + const latest = service.publishCommitted('status'); + const heartbeat = firstValueFrom( + service.stream(latest.id).pipe( + filter((event) => event.type === 'heartbeat'), + take(1), + ), + ); + + jest.advanceTimersByTime(1_000); + await Promise.resolve(); + + const event = await heartbeat; + expect(event).toMatchObject({ type: 'heartbeat' }); + expect(event.id).toBe(latest.id); + } finally { + jest.useRealTimers(); + } + }); + + it('uses an explicit empty heartbeat cursor before any state event', async () => { + jest.useFakeTimers(); + try { + const service = new NetworkManagementEventStreamService({ + heartbeatMs: 1_000, + replayLimit: 1, + }); + const heartbeat = firstValueFrom( + service.stream().pipe( + filter((event) => event.type === 'heartbeat'), + take(1), + ), + ); + + jest.advanceTimersByTime(1_000); + await Promise.resolve(); + + await expect(heartbeat).resolves.toMatchObject({ + id: '', + type: 'heartbeat', + }); + } finally { + jest.useRealTimers(); + } + }); +}); diff --git a/test/admin/network-management/network-management.controller.spec.ts b/test/admin/network-management/network-management.controller.spec.ts index d235908..81b63d3 100644 --- a/test/admin/network-management/network-management.controller.spec.ts +++ b/test/admin/network-management/network-management.controller.spec.ts @@ -1,9 +1,14 @@ import type { INestApplication } from '@nestjs/common'; import { Test } from '@nestjs/testing'; +import { type Observable, of } from 'rxjs'; import * as request from 'supertest'; import { AdminSuperGuard } from '../../../src/modules/admin/identity/auth/admin-super.guard'; import { JwtAuthGuard } from '../../../src/modules/admin/identity/auth/jwt-auth.guard'; import { NetworkManagementController } from '../../../src/modules/admin/platform-config/network-management/network-management.controller'; +import { + NetworkManagementEventStreamService, + type NetworkManagementStreamEvent, +} from '../../../src/modules/admin/platform-config/network-management/network-management-event-stream.service'; import { NetworkManagementService } from '../../../src/modules/admin/platform-config/network-management/network-management.service'; describe('NetworkManagementController', () => { @@ -21,6 +26,22 @@ describe('NetworkManagementController', () => { retry: jest.fn(), update: jest.fn(), }; + const eventStream = { + stream: jest.fn< + Observable, + [lastEventId?: string] + >(() => + of({ + data: { + eventId: 'network-event-1', + observedAt: '2026-07-23T00:00:00.000Z', + source: 'reported', + }, + id: 'network-event-1', + type: 'network-state-changed', + }), + ), + }; beforeAll(async () => { const moduleRef = await Test.createTestingModule({ @@ -28,6 +49,10 @@ describe('NetworkManagementController', () => { providers: [ AdminSuperGuard, { provide: NetworkManagementService, useValue: service }, + { + provide: NetworkManagementEventStreamService, + useValue: eventStream, + }, ], }) .overrideGuard(JwtAuthGuard) @@ -160,4 +185,52 @@ describe('NetworkManagementController', () => { lastErrorMessage: 'router conflict', }); }); + + it('streams committed MQTT updates through a real SSE HTTP request', async () => { + await request(apiUrl) + .get('/system/network/events/stream?lastEventId=query-event') + .set('Last-Event-ID', 'header-event') + .buffer(true) + .parse((response, callback) => { + response.once('data', () => callback(null, 'ok')); + }) + .expect('content-type', /text\/event-stream/) + .expect(200); + + expect(eventStream.stream).toHaveBeenCalledWith('header-event'); + expect(service.list).not.toHaveBeenCalled(); + expect(service.agentStatus).not.toHaveBeenCalled(); + }); + + it('keeps an empty heartbeat cursor instead of accepting a Nest-generated ID', async () => { + let parsed = false; + let serialized = ''; + eventStream.stream.mockReturnValueOnce( + of({ + data: { + message: 'alive', + observedAt: '2026-07-23T00:00:00.000Z', + }, + id: '', + type: 'heartbeat', + }), + ); + + await request(apiUrl) + .get('/system/network/events/stream') + .buffer(true) + .parse((response, callback) => { + response.on('data', (chunk) => { + serialized += chunk.toString('utf8'); + if (!parsed && serialized.includes('event: heartbeat')) { + parsed = true; + callback(null, 'ok'); + } + }); + }) + .expect(200); + + expect(serialized).toContain('event: heartbeat\nid: \n'); + expect(serialized).not.toContain('id: 1\n'); + }); });