feat: 增加网络管理状态事件流
This commit is contained in:
parent
01b18533fe
commit
a757c89276
29
API.md
29
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` |
|
||||
|
||||
@ -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/<pid>/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 <uin>` 快速登录参数,避免重启后退回无账号扫码;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 默认值。
|
||||
|
||||
|
||||
@ -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,
|
||||
];
|
||||
|
||||
|
||||
@ -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<void> {
|
||||
/**
|
||||
* 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<boolean> {
|
||||
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<string>();
|
||||
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<string>();
|
||||
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<void> {
|
||||
/**
|
||||
* 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<boolean> {
|
||||
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<void> {
|
||||
): Promise<boolean> {
|
||||
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()) {
|
||||
|
||||
@ -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<NetworkManagementStreamEvent>();
|
||||
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<NetworkManagementStreamEvent> {
|
||||
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 || '';
|
||||
}
|
||||
}
|
||||
@ -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')
|
||||
|
||||
@ -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;
|
||||
|
||||
@ -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');
|
||||
});
|
||||
});
|
||||
|
||||
@ -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();
|
||||
}
|
||||
});
|
||||
});
|
||||
@ -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<NetworkManagementStreamEvent>,
|
||||
[lastEventId?: string]
|
||||
>(() =>
|
||||
of<NetworkManagementStreamEvent>({
|
||||
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');
|
||||
});
|
||||
});
|
||||
|
||||
Loading…
Reference in New Issue
Block a user