fix: 避免网络心跳触发页面刷新
This commit is contained in:
parent
a757c89276
commit
b5041af7c6
2
API.md
2
API.md
@ -112,7 +112,7 @@ 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,也不使用定时轮询。
|
||||
Admin 首次进入网络管理页通过 HTTP 读取快照,随后使用 `/system/network/events/stream` 接收 `network-state-changed`。API 只在 `reported`、`status`、`events` 对应事务提交且语义状态实际变化后发出事件;MQTT QoS 1 重投、仅推进 `lastHeartbeatAt` 的状态心跳以及仅推进 `currentObservedAt/currentValidUntil/lastObservedAt` 的租约续期继续持久化,但不触发刷新。公网 IP/端口、Keeper/同步/错误/删除状态或 Agent 在线会话变化仍发布事件。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 处理。
|
||||
|
||||
|
||||
@ -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 阻塞。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 控制链路未接线。
|
||||
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 幂等重投、`status` 心跳时间推进和 `reported` 租约时间续期仍写入数据库,但不触发页面刷新。公网 IP/端口、Keeper/同步/错误/删除状态或 Agent 在线会话变化仍发布事件。SSE 心跳复用最近一次真实状态事件游标,尚无状态事件时显式发送空游标,避免 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 默认值。
|
||||
|
||||
|
||||
@ -403,7 +403,10 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
'Reported mapping Keeper intent does not match',
|
||||
);
|
||||
}
|
||||
const mappingBefore = this.reportedMappingStateFingerprint(mapping);
|
||||
const persistedMappingBefore =
|
||||
this.reportedPersistedMappingStateFingerprint(mapping);
|
||||
const refreshMappingBefore =
|
||||
this.reportedRefreshMappingStateFingerprint(mapping);
|
||||
mapping.reportedRevision = String(item.revision);
|
||||
mapping.syncStatus = item.syncStatus;
|
||||
mapping.keeperStatus = item.keeperStatus;
|
||||
@ -432,10 +435,18 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
mapping.isDeleted = true;
|
||||
finalizedDeletion = true;
|
||||
}
|
||||
if (this.reportedMappingStateFingerprint(mapping) !== mappingBefore) {
|
||||
visibleStateChanged = true;
|
||||
if (
|
||||
this.reportedPersistedMappingStateFingerprint(mapping) !==
|
||||
persistedMappingBefore
|
||||
) {
|
||||
await mappingRepository.save(mapping);
|
||||
}
|
||||
if (
|
||||
this.reportedRefreshMappingStateFingerprint(mapping) !==
|
||||
refreshMappingBefore
|
||||
) {
|
||||
visibleStateChanged = true;
|
||||
}
|
||||
}
|
||||
|
||||
if (BigInt(report.appliedRevision) > BigInt(state.appliedRevision)) {
|
||||
@ -520,7 +531,7 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
/**
|
||||
* 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.
|
||||
* @returns True only when semantic Admin status changed; heartbeat time is still persisted.
|
||||
*/
|
||||
private async applyStatus(status: NetworkStatusSnapshot): Promise<boolean> {
|
||||
this.assertAgentId(status.agentId);
|
||||
@ -559,7 +570,8 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
const stateBefore = this.statusStateFingerprint(state);
|
||||
const persistedStateBefore = this.statusPersistedStateFingerprint(state);
|
||||
const refreshStateBefore = this.statusRefreshStateFingerprint(state);
|
||||
state.online = status.online;
|
||||
state.version = status.version || null;
|
||||
state.startedAt = incomingStartedAt
|
||||
@ -570,9 +582,12 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
}
|
||||
state.lastMqttErrorCode = status.errorCode || null;
|
||||
state.lastMqttErrorMessage = status.errorMessage || null;
|
||||
if (this.statusStateFingerprint(state) === stateBefore) return false;
|
||||
if (
|
||||
this.statusPersistedStateFingerprint(state) !== persistedStateBefore
|
||||
) {
|
||||
await repository.save(state);
|
||||
return true;
|
||||
}
|
||||
return this.statusRefreshStateFingerprint(state) !== refreshStateBefore;
|
||||
});
|
||||
}
|
||||
|
||||
@ -624,11 +639,13 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
}
|
||||
|
||||
/**
|
||||
* Serializes only mapping fields rendered by Admin or used for row availability.
|
||||
* Serializes every report-owned mapping field that must be persisted.
|
||||
* @param mapping - Current persisted mapping before or after one report application.
|
||||
* @returns Stable comparison string for suppressing QoS 1 redelivery refreshes.
|
||||
* @returns Stable comparison including lease timestamps for database writes.
|
||||
*/
|
||||
private reportedMappingStateFingerprint(mapping: NetworkPortForward): string {
|
||||
private reportedPersistedMappingStateFingerprint(
|
||||
mapping: NetworkPortForward,
|
||||
): string {
|
||||
return JSON.stringify([
|
||||
mapping.activeKey,
|
||||
mapping.currentObservedAt,
|
||||
@ -647,6 +664,29 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
]);
|
||||
}
|
||||
|
||||
/**
|
||||
* Serializes only semantic mapping changes that justify reloading the Admin page.
|
||||
* @param mapping - Current persisted mapping before or after one report application.
|
||||
* @returns Stable comparison excluding lease-renewal timestamps.
|
||||
*/
|
||||
private reportedRefreshMappingStateFingerprint(
|
||||
mapping: NetworkPortForward,
|
||||
): string {
|
||||
return JSON.stringify([
|
||||
mapping.activeKey,
|
||||
mapping.currentPublicIpv4,
|
||||
mapping.currentPublicPort,
|
||||
mapping.isDeleted,
|
||||
mapping.keeperStatus,
|
||||
mapping.lastErrorCode,
|
||||
mapping.lastErrorMessage,
|
||||
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.
|
||||
@ -663,11 +703,11 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
}
|
||||
|
||||
/**
|
||||
* Serializes status-topic fields shown by the network page.
|
||||
* Serializes every status-topic field that must be persisted.
|
||||
* @param state - Agent singleton before or after one status snapshot.
|
||||
* @returns Stable comparison string for duplicate status suppression.
|
||||
* @returns Stable comparison including heartbeat time for database writes.
|
||||
*/
|
||||
private statusStateFingerprint(state: NetworkAgentState): string {
|
||||
private statusPersistedStateFingerprint(state: NetworkAgentState): string {
|
||||
return JSON.stringify([
|
||||
state.lastHeartbeatAt,
|
||||
state.lastMqttErrorCode,
|
||||
@ -678,6 +718,21 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy {
|
||||
]);
|
||||
}
|
||||
|
||||
/**
|
||||
* Serializes semantic Agent status without the continuously advancing heartbeat.
|
||||
* @param state - Agent singleton before or after one status snapshot.
|
||||
* @returns Stable comparison used to suppress heartbeat-only browser reloads.
|
||||
*/
|
||||
private statusRefreshStateFingerprint(state: NetworkAgentState): string {
|
||||
return JSON.stringify([
|
||||
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()) {
|
||||
|
||||
@ -20,11 +20,13 @@ type MqttHarness = {
|
||||
client: MqttClient & EventEmitter;
|
||||
clientOptions: () => IClientOptions;
|
||||
histories: NetworkEndpointHistory[];
|
||||
mappingSave: jest.Mock;
|
||||
publishCommitted: jest.Mock;
|
||||
mapping: NetworkPortForward;
|
||||
publishCallback: () => (error?: Error) => void;
|
||||
service: NetworkAgentMqttService;
|
||||
state: NetworkAgentState;
|
||||
stateSave: jest.Mock;
|
||||
transactionCalls: () => number;
|
||||
};
|
||||
|
||||
@ -61,14 +63,16 @@ function createHarness(): MqttHarness {
|
||||
targetIpv4: '192.168.31.224',
|
||||
});
|
||||
const histories: NetworkEndpointHistory[] = [];
|
||||
const stateSave = jest.fn(async (value) => Object.assign(state, value));
|
||||
const stateRepository = {
|
||||
findOne: async () => state,
|
||||
save: async (value) => Object.assign(state, value),
|
||||
save: stateSave,
|
||||
} as unknown as Repository<NetworkAgentState>;
|
||||
const mappingSave = jest.fn(async (value) => Object.assign(mapping, value));
|
||||
const mappingRepository = {
|
||||
find: async () => (mapping.isDeleted ? [] : [mapping]),
|
||||
findOne: async ({ where }) => (where.id === mapping.id ? mapping : null),
|
||||
save: async (value) => Object.assign(mapping, value),
|
||||
save: mappingSave,
|
||||
} as unknown as Repository<NetworkPortForward>;
|
||||
const historyRepository = {
|
||||
create: (input) => Object.assign(new NetworkEndpointHistory(), input),
|
||||
@ -143,11 +147,13 @@ function createHarness(): MqttHarness {
|
||||
client,
|
||||
clientOptions: () => options,
|
||||
histories,
|
||||
mappingSave,
|
||||
mapping,
|
||||
publishCallback: () => publishAck,
|
||||
publishCommitted,
|
||||
service,
|
||||
state,
|
||||
stateSave,
|
||||
transactionCalls: () => transactionCallCount,
|
||||
};
|
||||
}
|
||||
@ -482,8 +488,64 @@ describe('NetworkAgentMqttService', () => {
|
||||
});
|
||||
|
||||
await harness.service.consumeMessage(topic, payload);
|
||||
const mappingSavesAfterFirstReport = harness.mappingSave.mock.calls.length;
|
||||
const stateSavesAfterFirstReport = harness.stateSave.mock.calls.length;
|
||||
await harness.service.consumeMessage(topic, payload);
|
||||
|
||||
expect(harness.mappingSave).toHaveBeenCalledTimes(
|
||||
mappingSavesAfterFirstReport,
|
||||
);
|
||||
expect(harness.stateSave).toHaveBeenCalledTimes(stateSavesAfterFirstReport);
|
||||
expect(harness.publishCommitted).toHaveBeenCalledTimes(1);
|
||||
expect(harness.publishCommitted).toHaveBeenCalledWith('reported');
|
||||
});
|
||||
|
||||
/** Proves lease renewal is persisted without turning timestamps into page reloads. */
|
||||
it('suppresses Admin events for reported lease-only timestamp renewal', async () => {
|
||||
const harness = createHarness();
|
||||
const topic = 'kt/network/v1/agents/nas-main/reported';
|
||||
|
||||
await harness.service.consumeMessage(topic, reported(harness, 7));
|
||||
harness.publishCommitted.mockClear();
|
||||
const savesBeforeRenewal = harness.mappingSave.mock.calls.length;
|
||||
await harness.service.consumeMessage(
|
||||
topic,
|
||||
reported(harness, 7, {
|
||||
currentEndpoint: {
|
||||
observedAt: '2026-07-22T01:03:04.000Z',
|
||||
publicIpv4: '8.8.8.8',
|
||||
publicPort: 45000,
|
||||
validUntil: '2026-07-22T01:05:04.000Z',
|
||||
},
|
||||
lastObservedEndpoint: {
|
||||
observedAt: '2026-07-22T01:03:04.000Z',
|
||||
publicIpv4: '8.8.8.8',
|
||||
publicPort: 45000,
|
||||
validUntil: '2026-07-22T01:05:04.000Z',
|
||||
},
|
||||
}),
|
||||
);
|
||||
|
||||
expect(harness.mapping.currentValidUntil?.toISOString()).toBe(
|
||||
'2026-07-22T01:05:04.000Z',
|
||||
);
|
||||
expect(harness.mapping.lastObservedAt?.toISOString()).toBe(
|
||||
'2026-07-22T01:03:04.000Z',
|
||||
);
|
||||
expect(harness.mappingSave).toHaveBeenCalledTimes(savesBeforeRenewal + 1);
|
||||
expect(harness.publishCommitted).not.toHaveBeenCalled();
|
||||
|
||||
await harness.service.consumeMessage(
|
||||
topic,
|
||||
reported(harness, 7, {
|
||||
currentEndpoint: {
|
||||
observedAt: '2026-07-22T01:04:04.000Z',
|
||||
publicIpv4: '8.8.8.8',
|
||||
publicPort: 45001,
|
||||
validUntil: '2026-07-22T01:06:04.000Z',
|
||||
},
|
||||
}),
|
||||
);
|
||||
expect(harness.publishCommitted).toHaveBeenCalledTimes(1);
|
||||
expect(harness.publishCommitted).toHaveBeenCalledWith('reported');
|
||||
});
|
||||
@ -692,6 +754,47 @@ describe('NetworkAgentMqttService', () => {
|
||||
expect(harness.publishCommitted).toHaveBeenCalledWith('events');
|
||||
});
|
||||
|
||||
/** Proves liveness persistence does not become a periodic browser refresh. */
|
||||
it('suppresses Admin events for status heartbeat-only timestamp renewal', async () => {
|
||||
const harness = createHarness();
|
||||
const topic = 'kt/network/v1/agents/nas-main/status';
|
||||
const status = (observedAt: string, online = true) =>
|
||||
Buffer.from(
|
||||
JSON.stringify({
|
||||
agentId: 'nas-main',
|
||||
observedAt,
|
||||
online,
|
||||
schemaVersion: 1,
|
||||
startedAt: '2026-07-22T01:00:00.000Z',
|
||||
version: '0.1.0',
|
||||
}),
|
||||
);
|
||||
|
||||
await harness.service.consumeMessage(
|
||||
topic,
|
||||
status('2026-07-22T01:01:00.000Z'),
|
||||
);
|
||||
harness.publishCommitted.mockClear();
|
||||
const savesBeforeHeartbeat = harness.stateSave.mock.calls.length;
|
||||
await harness.service.consumeMessage(
|
||||
topic,
|
||||
status('2026-07-22T01:02:00.000Z'),
|
||||
);
|
||||
|
||||
expect(harness.state.lastHeartbeatAt?.toISOString()).toBe(
|
||||
'2026-07-22T01:02:00.000Z',
|
||||
);
|
||||
expect(harness.stateSave).toHaveBeenCalledTimes(savesBeforeHeartbeat + 1);
|
||||
expect(harness.publishCommitted).not.toHaveBeenCalled();
|
||||
|
||||
await harness.service.consumeMessage(
|
||||
topic,
|
||||
status('2026-07-22T01:03:00.000Z', false),
|
||||
);
|
||||
expect(harness.publishCommitted).toHaveBeenCalledTimes(1);
|
||||
expect(harness.publishCommitted).toHaveBeenCalledWith('status');
|
||||
});
|
||||
|
||||
it('accepts a same-instance LWT without regressing heartbeat and ignores an old-instance LWT', async () => {
|
||||
const harness = createHarness();
|
||||
const topic = 'kt/network/v1/agents/nas-main/status';
|
||||
|
||||
Loading…
Reference in New Issue
Block a user