把散落在代码里的真实部署值统一收进 config/,代码改为引用配置, 使仓库副本/开源导出不再带出生产域名、IP、内网路径与凭据。 新增 config/:platform.env.example(模板)· load.sh(shell 加载器)· index.cjs(node 加载器)· README.md(键一览与优先级)。 真实值放 config/platform.env —— 已 .gitignore 排除,不入库、不进导出。 TS 侧新增 src/platform-paths.ts 作部署路径的唯一解析处(零副作用): platformDir/stateDir/backupDir/artifactDir/installDir/scriptPath。 config.ts 接入这些字段;内置中继种子由生产 URL 改为空(改由 DSHS_OVERLAY_BOOTSTRAP_SEEDS 提供)。修掉 5 处硬编码绝对路径, src/** 注释中性化 116 行/53 文件。 scripts/** 36 个内部运维脚本:真令牌/PG 口令/隧道目标/主机号/路径 一律改从配置取;web/wake.html 的注册域白名单改为运行时从 location.hostname 推导;test/** 夹具 119 行/13 文件改 RFC 2606/5737 保留值,并把「内置种子必须为空」固化为回归断言。 取证:tsc 0 错;npm test 373/375(唯一失败 lease 属既有); 全仓扫描(大小写不敏感)代码面涉密标识 = 0;已部署 47 并零回归 (/opt/dsh/* 未搬家,/var/lib/dshs/platform 未被误建)。
2012 lines
94 KiB
JavaScript
2012 lines
94 KiB
JavaScript
/**
|
||
* 覆盖网络 · relay R1 单测(**真起服务、真握手、真泵字节**,不 mock 传输层)。
|
||
*
|
||
* ## 为什么必须是"真链路"测试
|
||
* relay 要替掉的是 sshd 反向隧道,而 sshd 那条路的病根正是**静默失败**(`-R` 撞号时没人
|
||
* 检查返回值)。所以本文件的验收标准不是"函数返回了 true",而是:
|
||
* 1. **字节真的过去了**(端到端回环泵,T3);
|
||
* 2. **失败真的被拒了**(错密钥 / 重放 / 越界端口,T4–T6)—— 且**有计数**可查。
|
||
*
|
||
* 运行:`node --test test/relay.test.mjs`(已登记进 `npm run verify`)。
|
||
*
|
||
* @module test/relay
|
||
*/
|
||
|
||
import assert from 'node:assert/strict'
|
||
import { createHmac, randomBytes } from 'node:crypto'
|
||
import { readFile } from 'node:fs/promises'
|
||
import { createServer as createTcpServer, connect } from 'node:net'
|
||
import { test } from 'node:test'
|
||
import { MUX, OPS_NETWORK, RelayClient, RelayDialer, RelayRendezvous, RelayServer, chooseNode, decodeMux, encodeJsonFrame, encodeMux, hostNameIndex, logicalName, parseKeysInline, relayEndpointTarget } from '../lib/net/relay/index.js'
|
||
|
||
const BASE = 45000
|
||
const SPAN = 200
|
||
const PATH = '/dshs-relay'
|
||
|
||
/* ─────────── 小工具 ─────────── */
|
||
|
||
async function listenInRange(server, lo, hi) {
|
||
for (let p = lo; p < hi; p++) {
|
||
const ok = await new Promise((resolve) => {
|
||
const onErr = () => {
|
||
server.off('error', onErr)
|
||
resolve(false)
|
||
}
|
||
server.once('error', onErr)
|
||
server.listen(p, '127.0.0.1', () => {
|
||
server.off('error', onErr)
|
||
resolve(true)
|
||
})
|
||
})
|
||
if (ok) return server.address().port
|
||
}
|
||
throw new Error(`no free port in ${lo}..${hi}`)
|
||
}
|
||
|
||
async function waitFor(cond, ms = 5000) {
|
||
const t0 = Date.now()
|
||
while (Date.now() - t0 < ms) {
|
||
if (cond()) return true
|
||
await new Promise((r) => setTimeout(r, 20))
|
||
}
|
||
return cond()
|
||
}
|
||
|
||
/** 连一个端口、写一段字、读回同样的字(带超时,避免测试悬挂)。 */
|
||
function roundTrip(port, text, ms = 5000) {
|
||
return new Promise((resolve, reject) => {
|
||
const sock = connect(port, '127.0.0.1')
|
||
let got = ''
|
||
const timer = setTimeout(() => {
|
||
sock.destroy()
|
||
reject(new Error(`roundTrip timeout after ${ms}ms (got ${got.length}/${text.length} bytes)`))
|
||
}, ms)
|
||
sock.on('connect', () => sock.write(text))
|
||
sock.on('data', (chunk) => {
|
||
got += chunk.toString('utf8')
|
||
if (got.length >= text.length) {
|
||
clearTimeout(timer)
|
||
sock.end()
|
||
resolve(got)
|
||
}
|
||
})
|
||
sock.on('error', (err) => {
|
||
clearTimeout(timer)
|
||
reject(err)
|
||
})
|
||
})
|
||
}
|
||
|
||
/** 用内建 WebSocket 手工发一帧 HELLO(用于构造"正常 client 做不到"的非法输入)。 */
|
||
async function rawHello(wsUrl, hostId, secret, portsCsv, nonce) {
|
||
const ws = new WebSocket(wsUrl)
|
||
ws.binaryType = 'arraybuffer'
|
||
await new Promise((resolve, reject) => {
|
||
ws.addEventListener('open', resolve, { once: true })
|
||
ws.addEventListener('error', () => reject(new Error('ws open failed')), { once: true })
|
||
})
|
||
const ts = Date.now()
|
||
const mac = createHmac('sha256', Buffer.from(secret, 'hex')).update(`${hostId}|${ts}|${nonce}|${portsCsv}`).digest('hex')
|
||
ws.send(encodeJsonFrame(MUX.HELLO, 0, { v: 1, hostId, ts, nonce, portsCsv, mac }))
|
||
const frame = await new Promise((resolve) => {
|
||
const timer = setTimeout(() => resolve(undefined), 3000)
|
||
ws.addEventListener('message', (ev) => {
|
||
clearTimeout(timer)
|
||
resolve(decodeMux(Buffer.from(ev.data)))
|
||
})
|
||
ws.addEventListener('close', () => {
|
||
clearTimeout(timer)
|
||
resolve(undefined)
|
||
})
|
||
})
|
||
return { ws, frame }
|
||
}
|
||
|
||
/* ─────────── T1 / T2:纯函数 ─────────── */
|
||
|
||
test('T1 mux 帧编解码往返(含 streamId 边界与空负载)', () => {
|
||
for (const [type, id, payload] of [
|
||
[MUX.HELLO, 0, Buffer.from('x')],
|
||
[MUX.DATA, 1, Buffer.alloc(0)],
|
||
[MUX.DATA, 0xffffffff, Buffer.from([0, 1, 2, 3])],
|
||
[MUX.OPEN, 4294967295 - 1, Buffer.alloc(300)],
|
||
]) {
|
||
const raw = encodeMux(type, id, payload)
|
||
assert.equal(raw.length, 5 + payload.length)
|
||
const back = decodeMux(raw)
|
||
assert.equal(back.type, type)
|
||
assert.equal(back.streamId, id)
|
||
assert.deepEqual(Buffer.from(back.payload), payload)
|
||
}
|
||
assert.equal(decodeMux(Buffer.alloc(4)), null, '小于 5 字节必须判为协议错误')
|
||
})
|
||
|
||
test('T2 密钥装载:只接受 64 位 hex,短密钥直接拒', () => {
|
||
const good = randomBytes(32).toString('hex')
|
||
const keys = parseKeysInline(`w-a:${good}`)
|
||
// 序③ 起:`parseKeysInline` 的返回值带上了"属于哪张网"(**成员资格的唯一判据**)。
|
||
// 裸 `hostId:secret` 是 R5 的旧写法 ⇒ 归入运维网 `ops`(现网 relay-keys.json 一字不改照旧可用)。
|
||
assert.deepEqual(keys.get('ops/w-a'), { network: 'ops', secret: good })
|
||
assert.throws(() => parseKeysInline('w-a:deadbeef'), /64 hex/)
|
||
assert.throws(() => parseKeysInline('no-colon'), /malformed/)
|
||
})
|
||
|
||
/* ─────────── T3:端到端(核心) ─────────── */
|
||
|
||
test('T3 端到端:Manager ⇒ relay 回环口 ⇒ relay ⇒ worker 本地端口(字节真过)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((sock) => sock.pipe(sock))
|
||
const workerPort = await listenInRange(echo, BASE, BASE + SPAN)
|
||
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-t', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
|
||
let client
|
||
t.after(async () => {
|
||
client?.stop()
|
||
await server.stop()
|
||
echo.close()
|
||
})
|
||
|
||
client = new RelayClient({ url: `ws://127.0.0.1:${server.boundPort}${PATH}`, hostId: 'w-t', secret, ports: [workerPort], log: () => {} })
|
||
client.start()
|
||
|
||
assert.ok(await waitFor(() => client.status().state === 'up'), '客户端未在 5s 内完成注册')
|
||
assert.ok(await waitFor(() => server.isOnline('w-t')), '服务端未看到该 host 上线')
|
||
const localPort = server.localPortOf('w-t', workerPort)
|
||
assert.ok(typeof localPort === 'number' && localPort > 0, '未分配回环端口')
|
||
assert.notEqual(localPort, workerPort, '回环口不应与 worker 端口同号(会撞本机实例)')
|
||
|
||
const text = 'relay-round-trip-0123456789'
|
||
assert.equal(await roundTrip(localPort, text), text)
|
||
|
||
// 并发流:同一 (host, port) 上两条连接必须能同时工作(多路复用真的在复用)。
|
||
const [a, b] = await Promise.all([roundTrip(localPort, 'AAAA'), roundTrip(localPort, 'BBBBBBBB')])
|
||
assert.equal(a, 'AAAA')
|
||
assert.equal(b, 'BBBBBBBB')
|
||
|
||
const st = server.status()
|
||
assert.ok(st.counters.streamsOpened >= 3, `期望至少 3 条流,实得 ${st.counters.streamsOpened}`)
|
||
assert.equal(st.counters.authFailed, 0)
|
||
assert.equal(st.counters.dropped, 0)
|
||
assert.equal(client.status().denied, 0)
|
||
})
|
||
|
||
/* ─────────── T4:认证失败必须被拒 ─────────── */
|
||
|
||
test('T4 错密钥 ⇒ 拒绝(不 up、有计数、不静默)', async (t) => {
|
||
const right = randomBytes(32).toString('hex')
|
||
const wrong = randomBytes(32).toString('hex')
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-x', right]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
const client = new RelayClient({
|
||
url: `ws://127.0.0.1:${server.boundPort}${PATH}`,
|
||
hostId: 'w-x',
|
||
secret: wrong,
|
||
ports: [BASE + 1],
|
||
reconnectMinMs: 50,
|
||
reconnectMaxMs: 100,
|
||
log: () => {},
|
||
})
|
||
t.after(async () => {
|
||
client.stop()
|
||
await server.stop()
|
||
})
|
||
client.start()
|
||
assert.ok(await waitFor(() => server.status().counters.authFailed > 0, 4000), '服务端没有记录认证失败')
|
||
assert.notEqual(client.status().state, 'up', '错密钥不得进入 up')
|
||
})
|
||
|
||
/* ─────────── T5:重放必须被拒 ─────────── */
|
||
|
||
test('T5 同 nonce 二次注册 ⇒ 重放被拒(第二条连接拿不到 HELLO_ACK)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-r', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
t.after(() => server.stop())
|
||
|
||
const csv = String(BASE + 2)
|
||
const nonce = randomBytes(16).toString('hex')
|
||
const first = await rawHello(`ws://127.0.0.1:${server.boundPort}${PATH}`, 'w-r', secret, csv, nonce)
|
||
t.after(() => first.ws.close())
|
||
assert.ok(first.frame !== undefined && first.frame.type === MUX.HELLO_ACK, '首次注册应成功')
|
||
|
||
const second = await rawHello(`ws://127.0.0.1:${server.boundPort}${PATH}`, 'w-r', secret, csv, nonce)
|
||
t.after(() => second.ws.close())
|
||
// 拒绝可以是"关闭"或"结构化拒绝帧",但**绝不能是 HELLO_ACK**(且不得静默)。
|
||
assert.notEqual(second.frame?.type, MUX.HELLO_ACK, '重放必须被拒')
|
||
assert.equal(second.frame?.type, MUX.HELLO_ERR, '拒绝必须是结构化的 HELLO_ERR')
|
||
assert.equal(JSON.parse(Buffer.from(second.frame.payload).toString('utf8')).reason, 'nonce-replay')
|
||
assert.ok(server.status().counters.authFailed >= 1, '重放应计入 authFailed')
|
||
})
|
||
|
||
/* ─────────── T6:越界端口必须被拒 ─────────── */
|
||
|
||
test('T6 声明实例区间外的端口 ⇒ 拒绝注册(爆炸半径不外扩)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-b', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
t.after(() => server.stop())
|
||
|
||
const outside = BASE + SPAN + 5
|
||
const probe = await rawHello(`ws://127.0.0.1:${server.boundPort}${PATH}`, 'w-b', secret, String(outside), randomBytes(16).toString('hex'))
|
||
t.after(() => probe.ws.close())
|
||
assert.notEqual(probe.frame?.type, MUX.HELLO_ACK, `端口 ${outside} 在区间外,必须拒绝`)
|
||
assert.equal(probe.frame?.type, MUX.HELLO_ERR)
|
||
assert.equal(JSON.parse(Buffer.from(probe.frame.payload).toString('utf8')).reason, 'port-out-of-range')
|
||
assert.ok(server.status().counters.authFailed >= 1)
|
||
assert.equal(server.localPortOf('w-b', outside), undefined, '被拒的端口不得留下回环监听')
|
||
})
|
||
|
||
/* ─────────── T7:未声明端口没有回环口(默认拒绝) ─────────── */
|
||
|
||
test('T7 未注册的端口不存在回环监听(默认拒绝,不是默认放行)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((sock) => sock.pipe(sock))
|
||
const declared = await listenInRange(echo, BASE, BASE + SPAN)
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-d', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
let client
|
||
t.after(async () => {
|
||
client?.stop()
|
||
await server.stop()
|
||
echo.close()
|
||
})
|
||
client = new RelayClient({ url: `ws://127.0.0.1:${server.boundPort}${PATH}`, hostId: 'w-d', secret, ports: [declared], log: () => {} })
|
||
client.start()
|
||
assert.ok(await waitFor(() => client.status().state === 'up'))
|
||
|
||
assert.equal(server.localPortOf('w-d', declared + 1), undefined)
|
||
const endpoints = server.status().endpoints
|
||
assert.equal(endpoints.length, 1, `只应为声明端口开监听,实得 ${endpoints.length} 个`)
|
||
assert.equal(endpoints[0].port, declared)
|
||
})
|
||
|
||
|
||
/* ═══════════════════════════════════════════════════════════════════════════
|
||
* T8–T12:韧性 —— 节点启停 / 网络变化 / 网络中断 / 网络异常 / 时钟漂移
|
||
*
|
||
* 这五条对应传输方案 §12 的场景矩阵。判据不是"有没有重连",而是:
|
||
* **恢复得快不快(计划内 vs 故障)、前提变了会不会立刻纠正(地址/时钟)、静默异常能不能被判死。**
|
||
* ═══════════════════════════════════════════════════════════════════════════ */
|
||
|
||
/** 一个"连不上"的假 WebSocket:注册监听后立刻派发 error(模拟断网期间的 dial 失败)。 */
|
||
function makeFailingWs(counter) {
|
||
return class FailingWs {
|
||
constructor() {
|
||
counter.dials += 1
|
||
this.readyState = 0
|
||
this.binaryType = ''
|
||
this.bufferedAmount = 0
|
||
this.ls = {}
|
||
}
|
||
addEventListener(type, fn) {
|
||
;(this.ls[type] ||= []).push(fn)
|
||
if (type === 'error') setTimeout(() => this.fire('error', {}), 5)
|
||
}
|
||
fire(type, ev) {
|
||
for (const fn of this.ls[type] ?? []) fn(ev)
|
||
}
|
||
send() {}
|
||
close() {
|
||
this.readyState = 3
|
||
}
|
||
}
|
||
}
|
||
|
||
/** 一个"握手成功但随后彻底静默"的假 WebSocket(模拟半开:TCP 没断,但再也不来帧)。 */
|
||
function makeSilentWs(counter, acceptedPort, sessionId = 'silent-sess') {
|
||
return class SilentWs {
|
||
constructor() {
|
||
counter.dials += 1
|
||
this.readyState = 0
|
||
this.binaryType = ''
|
||
this.bufferedAmount = 0
|
||
this.ls = {}
|
||
setTimeout(() => {
|
||
this.readyState = 1
|
||
this.fire('open', {})
|
||
}, 5)
|
||
}
|
||
addEventListener(type, fn) {
|
||
;(this.ls[type] ||= []).push(fn)
|
||
}
|
||
fire(type, ev) {
|
||
for (const fn of this.ls[type] ?? []) fn(ev)
|
||
}
|
||
send(buf) {
|
||
const frame = decodeMux(Buffer.from(buf))
|
||
if (frame !== null && frame.type === MUX.HELLO) {
|
||
// hbSec=1 ⇒ 半开阈值 2.5s ⇒ 静默 3s 左右应被判死
|
||
setTimeout(() => this.fire('message', { data: encodeJsonFrame(MUX.HELLO_ACK, 0, { sessionId, accepted: [acceptedPort], hbSec: 1 }) }), 5)
|
||
}
|
||
// 之后**永不回帧**:PING 不回 PONG、不主动发任何东西
|
||
}
|
||
close(code = 1000) {
|
||
this.readyState = 3
|
||
setTimeout(() => this.fire('close', { code }), 1)
|
||
}
|
||
}
|
||
}
|
||
|
||
test('T8 服务端优雅停机 ⇒ 重启窗口内自动恢复(close 1001 / BYE 语义,不消耗退避)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const instPort = BASE + 10
|
||
// 用一个固定端口:**重启后必须回到同一地址**,才算证明"原地恢复"。
|
||
const probe = createTcpServer()
|
||
const fixedPort = await listenInRange(probe, BASE + 100, BASE + SPAN)
|
||
await new Promise((r) => probe.close(r))
|
||
|
||
const mk = () => new RelayServer({ port: fixedPort, keys: new Map([['w-r8', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
const s1 = mk()
|
||
await s1.start()
|
||
const client = new RelayClient({
|
||
url: `ws://127.0.0.1:${fixedPort}${PATH}`,
|
||
hostId: 'w-r8',
|
||
secret,
|
||
ports: [instPort],
|
||
// 退避下限故意设成 60s:若还走普通退避,下面的恢复断言**必然失败** ⇒ 断言才有区分度。
|
||
reconnectMinMs: 60_000,
|
||
reconnectMaxMs: 60_000,
|
||
gracefulRetryMs: 250,
|
||
// 用"**窗口**"而不是次数:重启耗时不可预测(实测 stop() 自身就要 1.8s),窗口给足才稳。
|
||
gracefulBurstMs: 60_000,
|
||
log: () => {},
|
||
})
|
||
t.after(() => client.stop())
|
||
client.start()
|
||
assert.ok(await waitFor(() => client.status().state === 'up'), '首次注册未完成')
|
||
|
||
await s1.stop() // 优雅停机:BYE + close 1001
|
||
assert.ok(await waitFor(() => client.status().restarts >= 1, 3000), '未识别出对端的优雅下线信号')
|
||
// 关键:#1 必须**立刻**走短间隔重试,而不是把 60s 退避当第一反应。
|
||
assert.ok(
|
||
(client.status().nextRetryMs ?? 99999) <= 1000,
|
||
`graceful 后应处于短间隔重试窗口,实得 nextRetryMs=${client.status().nextRetryMs}`,
|
||
)
|
||
|
||
// systemd Restart= 的典型耗时窗口:400ms 后新进程在同一端口起来了。
|
||
await new Promise((r) => setTimeout(r, 400))
|
||
const s2 = mk()
|
||
await s2.start()
|
||
t.after(async () => {
|
||
await s2.stop()
|
||
})
|
||
|
||
assert.ok(
|
||
await waitFor(() => client.status().state === 'up', 5000),
|
||
`未在重启窗口内恢复(state=${client.status().state} attempts=${client.status().attempts} err=${client.status().lastError})`,
|
||
)
|
||
assert.ok(client.status().reconnects >= 1, '应计一次重连')
|
||
assert.ok(await waitFor(() => s2.isOnline('w-r8'), 2000), '新实例未看到该 host 上线')
|
||
})
|
||
|
||
test('T9 客户端优雅停机 ⇒ 服务端**秒级**标离线(BYE,不等心跳超时)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const port = BASE + 11
|
||
const server = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([['w-r9', secret]]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
idleTimeoutMs: 60_000, // 把心跳超时拉到 60s ⇒ 下面"1.5s 内离线"只可能来自 BYE
|
||
log: () => {},
|
||
})
|
||
await server.start()
|
||
const client = new RelayClient({ url: `ws://127.0.0.1:${server.boundPort}${PATH}`, hostId: 'w-r9', secret, ports: [port], log: () => {} })
|
||
t.after(async () => {
|
||
client.stop()
|
||
await server.stop()
|
||
})
|
||
client.start()
|
||
assert.ok(await waitFor(() => client.status().state === 'up'))
|
||
assert.ok(typeof server.localPortOf('w-r9', port) === 'number')
|
||
|
||
const t0 = Date.now()
|
||
client.stop()
|
||
assert.ok(await waitFor(() => !server.isOnline('w-r9'), 1500), 'BYE 未被及时处理(退化成心跳超时了)')
|
||
assert.ok(Date.now() - t0 < 1500, `离线耗时 ${Date.now() - t0}ms 过长`)
|
||
assert.equal(server.localPortOf('w-r9', port), undefined, '离线后不得再给出回环口(Manager 会打到死地址)')
|
||
})
|
||
|
||
test('T10 时钟漂移 > 认证窗口 ⇒ 用服务端时间戳自愈后注册成功(否则该节点永久失联)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const port = BASE + 12
|
||
const skew = 600_000 // 本机比 relay 快 10 分钟(远超 ±60s 窗口)
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-r10', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
const client = new RelayClient({
|
||
url: `ws://127.0.0.1:${server.boundPort}${PATH}`,
|
||
hostId: 'w-r10',
|
||
secret,
|
||
ports: [port],
|
||
reconnectMinMs: 50,
|
||
reconnectMaxMs: 200,
|
||
nowImpl: () => Date.now() + skew, // 注入漂移时钟
|
||
log: () => {},
|
||
})
|
||
t.after(async () => {
|
||
client.stop()
|
||
await server.stop()
|
||
})
|
||
client.start()
|
||
|
||
assert.ok(await waitFor(() => client.status().state === 'up', 6000), '时钟校正后仍未注册成功 ⇒ 该节点会永久失联')
|
||
assert.ok(Math.abs(client.status().clockSkewMs - skew) < 5_000, `clockSkewMs=${client.status().clockSkewMs} 未反映真实漂移`)
|
||
assert.ok(server.status().counters.authFailed >= 1, '首次握手应被拒(否则本用例没测到东西)')
|
||
})
|
||
|
||
test('T11 网络异常(半开)⇒ 无帧超过阈值即主动重连,不等 OS 的 TCP 超时', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const counter = { dials: 0 }
|
||
const port = BASE + 13
|
||
const client = new RelayClient({
|
||
url: 'ws://127.0.0.1:1/dshs-relay', // 地址不可达无所谓:传输层被下面的假 WS 完全替换
|
||
hostId: 'w-r11',
|
||
secret,
|
||
ports: [port],
|
||
reconnectMinMs: 100,
|
||
reconnectMaxMs: 200,
|
||
webSocketCtor: makeSilentWs(counter, port),
|
||
log: () => {},
|
||
})
|
||
t.after(() => client.stop())
|
||
client.start()
|
||
|
||
assert.ok(await waitFor(() => client.status().state === 'up', 3000), '假 WS 未完成注册')
|
||
// 服务端静默 ⇒ 应在 ~2.5×hb(1s)=2.5s 后被判死并重连
|
||
assert.ok(await waitFor(() => client.status().state === 'backoff', 8000), '半开未被判死(会一直挂着直到 OS 超时)')
|
||
assert.match(String(client.status().lastError), /half-open/, `lastError 应标明半开,实得 ${client.status().lastError}`)
|
||
assert.ok(await waitFor(() => counter.dials >= 2, 3000), '判死后没有重拨')
|
||
})
|
||
|
||
test('T12 网络变化 ⇒ 取消剩余退避、立即重拨(旧退避的前提已失效)', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const counter = { dials: 0 }
|
||
let snap = 0
|
||
const client = new RelayClient({
|
||
url: 'ws://127.0.0.1:1/dshs-relay',
|
||
hostId: 'w-r12',
|
||
secret,
|
||
ports: [BASE + 14],
|
||
// 退避下限 60s:只有"地址变化触发立即重试"这条路才能在几秒内重拨 ⇒ 断言才有区分度
|
||
reconnectMinMs: 60_000,
|
||
reconnectMaxMs: 60_000,
|
||
netWatchMs: 40,
|
||
netSnapshot: () => `snap-${snap++}`,
|
||
webSocketCtor: makeFailingWs(counter),
|
||
log: () => {},
|
||
})
|
||
t.after(() => client.stop())
|
||
client.start()
|
||
|
||
assert.ok(await waitFor(() => client.status().state === 'backoff', 3000), '未进入退避')
|
||
const firstRetry = client.status().nextRetryMs ?? 0
|
||
assert.ok(firstRetry > 30_000, `普通退避应很长,实得 ${firstRetry}ms`)
|
||
|
||
assert.ok(await waitFor(() => client.status().networkChanges >= 1, 2000), '未检测到本机地址变化')
|
||
assert.ok(await waitFor(() => counter.dials >= 2, 2000), '地址变化后未立即重拨')
|
||
})
|
||
|
||
/* ═══════════════════════════════════════════════════════════════════════════
|
||
* T13–T15:容量准入与节点选择 —— 自动(速度 + 负载) / 手动 / 满载排队
|
||
* ═══════════════════════════════════════════════════════════════════════════ */
|
||
|
||
/** 一个可控的假 WebSocket:`gate.full=true` 时回 `at-capacity`,否则回 ACK(模拟"等位成功后放行")。 */
|
||
function makeCapacityWs(counter, acceptedPort, gate, sessionId = 'q-1') {
|
||
return class CapacityWs {
|
||
constructor() {
|
||
counter.dials += 1
|
||
this.readyState = 0
|
||
this.binaryType = ''
|
||
this.bufferedAmount = 0
|
||
this.ls = {}
|
||
setTimeout(() => {
|
||
this.readyState = 1
|
||
this.fire('open', {})
|
||
}, 5)
|
||
}
|
||
addEventListener(type, fn) {
|
||
;(this.ls[type] ||= []).push(fn)
|
||
}
|
||
fire(type, ev) {
|
||
for (const fn of this.ls[type] ?? []) fn(ev)
|
||
}
|
||
send(buf) {
|
||
const frame = decodeMux(Buffer.from(buf))
|
||
if (frame === null || frame.type !== MUX.HELLO) return
|
||
const data = gate.full
|
||
? encodeJsonFrame(MUX.HELLO_ERR, 0, {
|
||
reason: 'at-capacity',
|
||
retryable: true,
|
||
retryAfterMs: 60,
|
||
serverTime: Date.now(),
|
||
capacity: { max: 1, used: 1, free: 0 },
|
||
})
|
||
: encodeJsonFrame(MUX.HELLO_ACK, 0, { sessionId, accepted: [acceptedPort], hbSec: 15, serverTime: Date.now() })
|
||
setTimeout(() => this.fire('message', { data }), 5)
|
||
}
|
||
close(code = 1000) {
|
||
this.readyState = 3
|
||
setTimeout(() => this.fire('close', { code }), 1)
|
||
}
|
||
}
|
||
}
|
||
|
||
test('T13 容量准入:满载时新节点被拒(at-capacity + 排队建议),已在册节点重连优先', async (t) => {
|
||
const s1 = randomBytes(32).toString('hex')
|
||
const s2 = randomBytes(32).toString('hex')
|
||
const server = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([
|
||
['w-a', s1],
|
||
['w-b', s2],
|
||
]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
maxHosts: 1, // 只能有一个节点在线
|
||
log: () => {},
|
||
})
|
||
await server.start()
|
||
t.after(() => server.stop())
|
||
|
||
const c1 = new RelayClient({ url: `ws://127.0.0.1:${server.boundPort}${PATH}`, hostId: 'w-a', secret: s1, ports: [BASE + 40], log: () => {} })
|
||
c1.start()
|
||
t.after(() => c1.stop())
|
||
assert.ok(await waitFor(() => c1.status().state === 'up'), '第一个节点未上线')
|
||
|
||
// 第二个新节点:必须被结构性拒绝,并拿到"多久后再来"
|
||
const probe = await rawHello(`ws://127.0.0.1:${server.boundPort}${PATH}`, 'w-b', s2, String(BASE + 41), randomBytes(16).toString('hex'))
|
||
t.after(() => probe.ws.close())
|
||
assert.equal(probe.frame?.type, MUX.HELLO_ERR, '满载必须回结构化拒绝,而不是静默断开')
|
||
const msg = JSON.parse(Buffer.from(probe.frame.payload).toString('utf8'))
|
||
assert.equal(msg.reason, 'at-capacity')
|
||
assert.equal(msg.retryable, true, '满载属于"过会儿再来",不是永久拒绝')
|
||
assert.ok(msg.retryAfterMs > 0, '必须给出排队建议时长')
|
||
assert.equal(msg.capacity.free, 0)
|
||
assert.ok(!server.isOnline('w-b'), '被拒的节点不得上线')
|
||
assert.equal(server.localPortOf('w-b', BASE + 41), undefined, '被拒的节点不得留下回环监听')
|
||
assert.equal(server.status().capacity.free, 0)
|
||
|
||
// **已在册节点重连优先**:位子本来就是它的,不能被自己的满载规则挡在门外
|
||
const again = await rawHello(`ws://127.0.0.1:${server.boundPort}${PATH}`, 'w-a', s1, String(BASE + 40), randomBytes(16).toString('hex'))
|
||
t.after(() => again.ws.close())
|
||
assert.equal(again.frame?.type, MUX.HELLO_ACK, '已在册节点重连被误拦 ⇒ 会造成"重启后再也连不上"')
|
||
})
|
||
|
||
test('T14 客户端满载排队:进 queued、不消耗退避、位子一空即注册成功', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const counter = { dials: 0 }
|
||
const gate = { full: true }
|
||
const port = BASE + 42
|
||
const client = new RelayClient({
|
||
url: 'ws://127.0.0.1:1/dshs-relay',
|
||
hostId: 'w-q',
|
||
secret,
|
||
ports: [port],
|
||
reconnectMinMs: 60_000,
|
||
reconnectMaxMs: 60_000,
|
||
webSocketCtor: makeCapacityWs(counter, port, gate),
|
||
log: () => {},
|
||
})
|
||
t.after(() => client.stop())
|
||
client.start()
|
||
|
||
assert.ok(await waitFor(() => client.status().state === 'queued', 3000), `未进入 queued,实得 ${client.status().state}`)
|
||
assert.ok(client.status().queueWaits >= 1, '未记录排队次数')
|
||
assert.ok((client.status().nextRetryMs ?? 99999) <= 5000, '排队应按服务端给的 retryAfterMs 回来,而不是 60s 退避')
|
||
assert.equal(client.status().attempts, 0, '排队不是故障,不得累计退避')
|
||
assert.match(String(client.status().lastError), /at-capacity/)
|
||
|
||
gate.full = false // 位子空出来了
|
||
assert.ok(await waitFor(() => client.status().state === 'up', 3000), '空位出现后未注册成功')
|
||
// 排队→成功是**首次**注册(不是重连),所以 `reconnects` 应为 0;这里正是要确认语义没被搞混。
|
||
assert.equal(client.status().reconnects, 0, '排队后首次成功不该被记成"重连"')
|
||
})
|
||
|
||
test('T15 选点判据:速度 + 负载打分、满载硬门、手动优先(纯函数)', () => {
|
||
const near = { id: 'near', rttMs: 20, capacity: { max: 10, used: 2 } }
|
||
const far = { id: 'far', rttMs: 200, capacity: { max: 10, used: 1 } }
|
||
const fullNode = { id: 'full', rttMs: 5, capacity: { max: 4, used: 4 } }
|
||
|
||
const d1 = chooseNode([near, far, fullNode])
|
||
assert.equal(d1.outcome, 'chosen')
|
||
assert.equal(d1.chosen.id, 'near', `速度占优者应胜出,实得 ${d1.chosen?.id}`)
|
||
assert.equal(d1.ranking.find((r) => r.id === 'full').blocked, 'full', '满载必须是硬门')
|
||
|
||
// 只有满载候选 ⇒ 排队(而不是"挑个满的凑合")
|
||
const d2 = chooseNode([fullNode])
|
||
assert.equal(d2.outcome, 'queued')
|
||
assert.equal(d2.chosen, undefined)
|
||
assert.ok(d2.retryAfterMs > 0)
|
||
// 不允许排队 ⇒ 直接拒绝加入
|
||
assert.equal(chooseNode([fullNode], { allowQueue: false }).outcome, 'rejected')
|
||
|
||
// 负载能压过速度:近但 90% 满 vs 远但空
|
||
const busyNear = { id: 'busyNear', rttMs: 20, capacity: { max: 10, used: 9 } }
|
||
const idleFar = { id: 'idleFar', rttMs: 60, capacity: { max: 10, used: 0 } }
|
||
assert.equal(chooseNode([busyNear, idleFar]).chosen.id, 'idleFar', '负载接近满载时应让位给空闲节点')
|
||
|
||
// 手动指定优先,且**不会**被自动算法改掉
|
||
assert.equal(chooseNode([near, far], { manualId: 'far' }).chosen.id, 'far')
|
||
// 手动指定但满载 ⇒ 排队,不静默换节点
|
||
const d5 = chooseNode([near, fullNode], { manualId: 'full' })
|
||
assert.equal(d5.outcome, 'queued')
|
||
assert.equal(d5.chosen, undefined)
|
||
// 手动指定的 id 不存在 ⇒ 明确拒绝(不偷偷选别的)
|
||
assert.equal(chooseNode([near], { manualId: 'ghost' }).outcome, 'rejected')
|
||
|
||
// 近期失败降权
|
||
const flaky = { id: 'flaky', rttMs: 20, capacity: { max: 10, used: 0 }, recentFailures: 4 }
|
||
const steady = { id: 'steady', rttMs: 45, capacity: { max: 10, used: 3 } }
|
||
assert.equal(chooseNode([flaky, steady]).chosen.id, 'steady', '近期反复失败应被降权')
|
||
|
||
// 空候选 ⇒ 明确拒绝,不抛异常
|
||
assert.equal(chooseNode([]).outcome, 'rejected')
|
||
})
|
||
|
||
/* ─────────── T16 / T17:运行期端口增删(R4,替掉 `ssh -O forward/cancel`)─────────── */
|
||
|
||
/** 连一个端口,**期望连不上**(用来证明监听真的被收掉了,而不是"记账删了但口还开着")。 */
|
||
function expectRefused(port, ms = 1500) {
|
||
return new Promise((resolve) => {
|
||
const sock = connect(port, '127.0.0.1')
|
||
let settled = false
|
||
const done = (v) => {
|
||
if (settled) return
|
||
settled = true
|
||
try {
|
||
sock.destroy()
|
||
} catch {
|
||
/* 已断 */
|
||
}
|
||
resolve(v)
|
||
}
|
||
sock.on('error', () => done(true))
|
||
sock.on('connect', () => done(false))
|
||
setTimeout(() => done(false), ms)
|
||
})
|
||
}
|
||
|
||
test('T16 运行期加端口:PORT_ADD ⇒ 真回环口 + 字节真过;越界被拒;PORT_DEL ⇒ 口真的收掉', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((s) => s.pipe(s))
|
||
// ⚠️ 故意**不在注册时声明它** —— 这正是实例端口的形态(运行期才知道)
|
||
const workerPort = await listenInRange(echo, BASE + 50, BASE + SPAN - 1)
|
||
const placeholder = BASE // HELLO 要求端口表非空(`no-ports`),agent 侧对应"agent 自身端口"
|
||
const server = new RelayServer({ port: 0, keys: new Map([['w-p', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server.start()
|
||
const client = new RelayClient({
|
||
url: `ws://127.0.0.1:${server.boundPort}${PATH}`,
|
||
hostId: 'w-p',
|
||
secret,
|
||
ports: [placeholder],
|
||
reconnectMinMs: 50,
|
||
reconnectMaxMs: 200,
|
||
log: () => {},
|
||
})
|
||
client.start()
|
||
t.after(async () => {
|
||
client.stop()
|
||
await server.stop()
|
||
echo.close()
|
||
})
|
||
assert.ok(await waitFor(() => client.status().state === 'up'), '客户端未在 5s 内注册')
|
||
|
||
// ① 默认拒绝:没声明的端口**不存在**回环监听(与 T7 同一条姿态,只是换个入口)
|
||
assert.equal(server.localPortOf('w-p', workerPort), undefined, '未声明的端口不该有回环口')
|
||
|
||
// ② 加端口 ⇒ 拿到真口号、字节真过
|
||
assert.equal(await client.addPort(workerPort), true, 'addPort 应成功')
|
||
assert.ok(await waitFor(() => (server.localPortOf('w-p', workerPort) ?? 0) > 0), '未分配回环口')
|
||
const localPort = server.localPortOf('w-p', workerPort)
|
||
assert.ok(typeof localPort === 'number' && localPort > 0, '未分配回环口')
|
||
assert.equal(await roundTrip(localPort, 'runtime-port-ok'), 'runtime-port-ok')
|
||
assert.deepEqual(client.status().dynamicPorts, [workerPort])
|
||
|
||
// ③ 幂等:重复声明不该报错(重连后重放会重复调它)
|
||
assert.equal(await client.addPort(workerPort), true, '重复 addPort 应幂等成功')
|
||
|
||
// ④ 越界 / 非法端口 ⇒ **明确被拒**(不静默成功)
|
||
assert.equal(await client.addPort(BASE + SPAN + 5), false, '区间外端口必须被拒')
|
||
assert.equal(await client.addPort(0), false, '非法端口必须被拒')
|
||
assert.equal(server.localPortOf('w-p', BASE + SPAN + 5), undefined)
|
||
assert.deepEqual(client.status().dynamicPorts, [workerPort], '被拒的端口不该进记账')
|
||
|
||
// ⑤ 撤端口 ⇒ 记账清掉 + **监听真的收掉**(不然就是"以为撤了、口还开着")
|
||
assert.equal(await client.removePort(workerPort), true)
|
||
assert.ok(await waitFor(() => server.localPortOf('w-p', workerPort) === undefined), '撤销后服务端不该再有该端点')
|
||
assert.equal(await expectRefused(localPort), true, '撤销后旧回环口必须连不上')
|
||
assert.deepEqual(client.status().dynamicPorts, [])
|
||
assert.equal(server.status().counters.protocolErrors, 0)
|
||
})
|
||
|
||
test('T17 断链重连后**重放运行期端口**(否则"本地以为转着、Manager 侧其实没有")', async (t) => {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((s) => s.pipe(s))
|
||
const workerPort = await listenInRange(echo, BASE + 60, BASE + SPAN - 1)
|
||
const PORT = BASE
|
||
const server1 = new RelayServer({ port: 0, keys: new Map([['w-r', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server1.start()
|
||
const client = new RelayClient({
|
||
url: `ws://127.0.0.1:${server1.boundPort}${PATH}`,
|
||
hostId: 'w-r',
|
||
secret,
|
||
ports: [PORT],
|
||
reconnectMinMs: 50,
|
||
reconnectMaxMs: 200,
|
||
gracefulRetryMs: 50,
|
||
log: () => {},
|
||
})
|
||
client.start()
|
||
t.after(async () => {
|
||
client.stop()
|
||
await server1.stop()
|
||
echo.close()
|
||
})
|
||
assert.ok(await waitFor(() => client.status().state === 'up'))
|
||
assert.equal(await client.addPort(workerPort), true)
|
||
assert.ok(await waitFor(() => (server1.localPortOf('w-r', workerPort) ?? 0) > 0), '首轮未开回环口')
|
||
|
||
// relay 侧重启(**同一端口**复用同一个 server 实例做不到 ⇒ 直接 stop 再来一个同端口的)
|
||
const bound = server1.boundPort
|
||
await server1.stop()
|
||
const server2 = new RelayServer({ port: bound, keys: new Map([['w-r', secret]]), instancePortBase: BASE, instancePortSpan: SPAN, log: () => {} })
|
||
await server2.start()
|
||
t.after(async () => {
|
||
await server2.stop()
|
||
})
|
||
|
||
// 重连 + **重放**:不重放的话,服务端这边永远是空的(本地却仍说 dynamicPorts 有值)
|
||
assert.ok(
|
||
await waitFor(() => (server2.localPortOf('w-r', workerPort) ?? 0) > 0, 15000),
|
||
'重连后未重放运行期端口 ⇒ Manager 侧会静默失去这个口',
|
||
)
|
||
const localPort2 = server2.localPortOf('w-r', workerPort)
|
||
assert.equal(await roundTrip(localPort2, 'replayed-ok'), 'replayed-ok')
|
||
assert.equal(server2.status().counters.authFailed, 0)
|
||
})
|
||
|
||
/** 往**拨号流**(`openStream()` 返回的 `Duplex`)里写一段字、读回同样的字。 */
|
||
function dialRoundTrip(duplex, text, ms = 5000) {
|
||
return new Promise((resolve, reject) => {
|
||
let got = ''
|
||
const timer = setTimeout(() => {
|
||
duplex.destroy()
|
||
reject(new Error(`dial roundTrip timeout after ${ms}ms (got ${got.length} chars)`))
|
||
}, ms)
|
||
duplex.on('data', (c) => {
|
||
got += c.toString()
|
||
if (got.length >= text.length) {
|
||
clearTimeout(timer)
|
||
resolve(got)
|
||
}
|
||
})
|
||
duplex.on('error', (e) => {
|
||
clearTimeout(timer)
|
||
reject(e)
|
||
})
|
||
duplex.write(text)
|
||
})
|
||
}
|
||
|
||
test('T18 拨号方(R5):openStream ⇒ 字节真过;离线 / 越权 / 配错的拒绝**都是显式的**', async (t) => {
|
||
const wSecret = randomBytes(32).toString('hex')
|
||
const mSecret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((s) => s.pipe(s))
|
||
const workerPort = await listenInRange(echo, BASE + 70, BASE + SPAN - 1)
|
||
const srv = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([
|
||
['w-d', wSecret],
|
||
['manager', mSecret],
|
||
]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
dialers: new Set(['manager']),
|
||
log: () => {},
|
||
})
|
||
await srv.start()
|
||
const url = `ws://127.0.0.1:${srv.boundPort}${PATH}`
|
||
const worker = new RelayClient({ url, hostId: 'w-d', secret: wSecret, ports: [workerPort], log: () => {} })
|
||
const dialer = new RelayClient({ url, hostId: 'manager', secret: mSecret, ports: [], dialer: true, log: () => {} })
|
||
worker.start()
|
||
dialer.start()
|
||
t.after(async () => {
|
||
worker.stop()
|
||
dialer.stop()
|
||
echo.close()
|
||
await srv.stop()
|
||
})
|
||
assert.ok(await waitFor(() => worker.status().state === 'up'), 'worker 未注册')
|
||
assert.ok(await waitFor(() => dialer.status().state === 'up'), '拨号方未注册')
|
||
assert.deepEqual(srv.status().dialers, ['manager'], '白名单要能在 /status 里看到')
|
||
|
||
// ① 真字节:写进去、读回来 —— `openStream` 返回的就是可 `pipe()` 的 Duplex
|
||
const duplex = await dialer.openStream('w-d', workerPort)
|
||
assert.equal(await dialRoundTrip(duplex, 'dial-ok'), 'dial-ok')
|
||
assert.equal(dialer.status().dialStreams, 1)
|
||
|
||
// ② 关流 ⇒ 两端都清干净(worker 侧计数回落,不留悬挂流)
|
||
duplex.destroy()
|
||
assert.ok(await waitFor(() => dialer.status().dialStreams === 0), '关流后拨号侧应清空')
|
||
assert.ok(
|
||
await waitFor(
|
||
() => (srv.status().endpoints.find((e) => e.hostId === 'w-d' && e.port === workerPort)?.streams ?? -1) === 0,
|
||
),
|
||
'worker 侧该端口的流计数应回落',
|
||
)
|
||
|
||
// ③ 离线目标 ⇒ **显式抛错**(不许返回一个"看起来能用"的对象)
|
||
await assert.rejects(() => dialer.openStream('w-nope', workerPort), /target-offline/, '离线目标必须显式拒绝')
|
||
|
||
// ④ 非拨号方不能拨(客户端侧先拦)
|
||
await assert.rejects(() => worker.openStream('w-d', workerPort), /requires dialer mode/)
|
||
|
||
// ⑤ 拨号方不得声明端口(与服务端 `dialer-must-not-declare-ports` 同口径)
|
||
assert.throws(
|
||
() =>
|
||
new RelayClient({ url: `ws://127.0.0.1:1${PATH}`, hostId: 'manager', secret: mSecret, ports: [1234], dialer: true }),
|
||
/must not declare ports/,
|
||
)
|
||
|
||
// ⑥ R1 的**默认拒绝**姿态没被改掉:端口表为空的普通客户端照样进不来
|
||
const before = srv.status().counters.authFailed
|
||
const bogus = new RelayClient({
|
||
url,
|
||
hostId: 'w-d',
|
||
secret: wSecret,
|
||
ports: [],
|
||
reconnectMinMs: 50,
|
||
reconnectMaxMs: 100,
|
||
log: () => {},
|
||
})
|
||
bogus.start()
|
||
t.after(() => bogus.stop())
|
||
assert.ok(await waitFor(() => srv.status().counters.authFailed > before), '空端口表的普通客户端必须被拒(no-ports)')
|
||
assert.equal(srv.status().counters.protocolErrors, 0, '拒绝要走 HELLO_ERR,不该变成协议错')
|
||
})
|
||
|
||
test('T19 【会合可换机】relay 不开任何回环口(纯流转发)⇒ 业务照样全通', async (t) => {
|
||
const wSecret = randomBytes(32).toString('hex')
|
||
const mSecret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((s) => s.pipe(s))
|
||
const workerPort = await listenInRange(echo, BASE + 80, BASE + SPAN - 1)
|
||
const srv = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([
|
||
['w-d', wSecret],
|
||
['manager', mSecret],
|
||
]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
dialers: new Set(['manager']),
|
||
// 🔴 关键:**一个本地口都不绑** —— 等价于"relay 在另一台机器上,Manager 够不到它的回环"
|
||
exposeLoopback: false,
|
||
log: () => {},
|
||
})
|
||
await srv.start()
|
||
const url = `ws://127.0.0.1:${srv.boundPort}${PATH}`
|
||
const worker = new RelayClient({ url, hostId: 'w-d', secret: wSecret, ports: [workerPort], log: () => {} })
|
||
const dialer = new RelayClient({ url, hostId: 'manager', secret: mSecret, ports: [], dialer: true, log: () => {} })
|
||
worker.start()
|
||
dialer.start()
|
||
t.after(async () => {
|
||
worker.stop()
|
||
dialer.stop()
|
||
echo.close()
|
||
await srv.stop()
|
||
})
|
||
assert.ok(await waitFor(() => worker.status().state === 'up'), 'worker 未注册')
|
||
assert.ok(await waitFor(() => dialer.status().state === 'up'), '拨号方未注册')
|
||
|
||
// 🔴 核心判据一:relay **没有任何回环落点**(`localPort` 拿不到)
|
||
assert.equal(srv.localPortOf('w-d', workerPort), undefined, '纯流转发模式下不该有回环落点')
|
||
|
||
// 🔴 核心判据二:**业务照样通** ⇒ "relay 必须与 Manager 同机"这个前提已经不存在
|
||
const duplex = await dialer.openStream('w-d', workerPort)
|
||
assert.equal(await dialRoundTrip(duplex, 'no-loopback-needed'), 'no-loopback-needed')
|
||
assert.equal(dialer.status().dialStreams, 1)
|
||
|
||
// 端点表仍在册(在线可观测),但 `localPort` 恒 0 ⇒ 靠 `/status` 找落点的那条路**失败关闭**
|
||
const ep = srv.status().endpoints.find((e) => e.hostId === 'w-d' && e.port === workerPort)
|
||
assert.ok(ep !== undefined, '端点条目仍应在册(在线可观测)')
|
||
assert.equal(ep.localPort, 0, '纯流转发模式下 localPort 必须为 0')
|
||
assert.equal(ep.online, true)
|
||
assert.equal(srv.status().counters.protocolErrors, 0)
|
||
duplex.destroy()
|
||
})
|
||
|
||
/* ─────────── T20:序⑤ 判别器计数(观测最小集的"最关键一条") ─────────── */
|
||
|
||
/**
|
||
* 为什么需要这一组计数:443 单 §12 的教训原文是「**静默失效靠判别器定位**」,判别器就是
|
||
* 「relay 到底有没有 `DIAL`」;而在此之前它**只存在于日志行**,脚本无法断言。
|
||
*
|
||
* 三个分支必须**逐条**可断言,且**互斥可加和**:
|
||
* | 分支 | 触发方式(本用例用的就是这些) | 期望计数 |
|
||
* |---|---|---|
|
||
* | 放行 | 拨号方 `openStream` 成功 | `dial` +1 |
|
||
* | 策略拒绝 | 普通 worker 手工发 `DIAL`(不在拨号方白名单) | `dialDenied` +1 |
|
||
* | 请求非法 | 拨号方发 `port=0`(`bad-target`) | `dialFailed` +1 |
|
||
* | 目标不可达 | 拨号方 `openStream('w-nope')` | `dialFailed` +1 |
|
||
*
|
||
* ⚠️ 为什么"白名单拒绝"必须用**裸 ws 手工发帧**:普通 client 根本发不出这一帧 ——
|
||
* ① 客户端侧 `openStream` 有「必须拨号方模式」前置闸门(T18 ④);
|
||
* ② 注册侧两道闸门互斥(`no-ports` / `dialer-must-not-declare-ports`)⇒ 不在白名单的 host
|
||
* 要么带端口注册(假不了拨号方)、要么空端口被拒。⇒ 只有手工帧能构造这个输入。
|
||
*/
|
||
test('T20 序⑤ 判别器计数:DIAL 放行 / 策略拒绝 / 目标不可达三分支都能被 /status 断言', async (t) => {
|
||
const wSecret = randomBytes(32).toString('hex')
|
||
const mSecret = randomBytes(32).toString('hex')
|
||
const xSecret = randomBytes(32).toString('hex')
|
||
const m2Secret = randomBytes(32).toString('hex')
|
||
const echo = createTcpServer((s) => s.pipe(s))
|
||
const workerPort = await listenInRange(echo, BASE + 90, BASE + SPAN - 1)
|
||
const srv = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([
|
||
['w-d', wSecret],
|
||
['w-x', xSecret],
|
||
['manager', mSecret],
|
||
['manager2', m2Secret],
|
||
]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
dialers: new Set(['manager', 'manager2']),
|
||
log: () => {},
|
||
})
|
||
await srv.start()
|
||
const url = `ws://127.0.0.1:${srv.boundPort}${PATH}`
|
||
const worker = new RelayClient({ url, hostId: 'w-d', secret: wSecret, ports: [workerPort], log: () => {} })
|
||
const dialer = new RelayClient({ url, hostId: 'manager', secret: mSecret, ports: [], dialer: true, log: () => {} })
|
||
worker.start()
|
||
dialer.start()
|
||
t.after(async () => {
|
||
worker.stop()
|
||
dialer.stop()
|
||
echo.close()
|
||
await srv.stop()
|
||
})
|
||
assert.ok(await waitFor(() => worker.status().state === 'up'), 'worker 未注册')
|
||
assert.ok(await waitFor(() => dialer.status().state === 'up'), '拨号方未注册')
|
||
|
||
// ⓪ 三个计数器**必须存在** —— 这一条就是"先红":加计数之前它们是 `undefined`,
|
||
// 于是"脚本能不能断言拨号路径"这个问题的答案就是"不能"。
|
||
const c0 = srv.status().counters
|
||
assert.equal(typeof c0.dial, 'number', 'counters.dial 必须存在(否则脚本无法断言拨号路径)')
|
||
assert.equal(typeof c0.dialDenied, 'number', 'counters.dialDenied 必须存在')
|
||
assert.equal(typeof c0.dialFailed, 'number', 'counters.dialFailed 必须存在')
|
||
|
||
// ① 放行 ⇒ dial +1(且与 streamsOpened 同点自增,两条必须同步)
|
||
const duplex = await dialer.openStream('w-d', workerPort)
|
||
assert.equal(await dialRoundTrip(duplex, 'dial-ok'), 'dial-ok')
|
||
assert.equal(srv.status().counters.dial, c0.dial + 1, '成功拨号必须计入 dial')
|
||
duplex.destroy()
|
||
assert.ok(await waitFor(() => dialer.status().dialStreams === 0), '关流后拨号侧应清空')
|
||
|
||
// ② 目标不可达 ⇒ dialFailed +1(`w-nope` 不在册)
|
||
await assert.rejects(() => dialer.openStream('w-nope', workerPort), /target-offline/, '离线目标必须显式拒绝')
|
||
assert.equal(srv.status().counters.dialFailed, c0.dialFailed + 1, '目标不可达必须计入 dialFailed')
|
||
|
||
// ③ 白名单拒绝 ⇒ dialDenied +1(普通 worker 手工发 DIAL;它注册得成、但没资格拨)
|
||
const raw = await rawHello(url, 'w-x', xSecret, String(workerPort), randomBytes(8).toString('hex'))
|
||
t.after(() => raw.ws.close())
|
||
assert.ok(raw.frame !== undefined && raw.frame.type === MUX.HELLO_ACK, 'w-x 应先以普通 worker 身份注册成功')
|
||
const ack = await new Promise((resolve) => {
|
||
const timer = setTimeout(() => resolve(undefined), 3000)
|
||
raw.ws.addEventListener(
|
||
'message',
|
||
(ev) => {
|
||
clearTimeout(timer)
|
||
resolve(decodeMux(Buffer.from(ev.data)))
|
||
},
|
||
{ once: true },
|
||
)
|
||
raw.ws.send(encodeJsonFrame(MUX.DIAL, 7, { target: 'w-d', port: workerPort }))
|
||
})
|
||
assert.ok(ack !== undefined, '拒绝也必须回 DIAL_ACK(不许静默)')
|
||
assert.equal(ack.type, MUX.DIAL_ACK, '拒绝也要走 DIAL_ACK')
|
||
assert.equal(JSON.parse(Buffer.from(ack.payload).toString('utf8')).error, 'not-a-dialer')
|
||
assert.equal(srv.status().counters.dialDenied, c0.dialDenied + 1, '白名单拒绝必须计入 dialDenied')
|
||
assert.equal(srv.status().counters.dial, c0.dial + 1, '被拒的拨号**不算** dial')
|
||
|
||
// ④ 请求非法(`port=0`)⇒ dialFailed 再 +1。
|
||
// ⚠️ 必须用**在白名单里的**拨号方会话来构造:`onDial` 的**第一道门就是白名单**,
|
||
// 拿非拨号方去打非法参数只会拿到 `not-a-dialer`(③ 已证)—— 门是**串行**的,不是并行判的。
|
||
const rawDialer = await rawHello(url, 'manager2', m2Secret, '', randomBytes(8).toString('hex'))
|
||
t.after(() => rawDialer.ws.close())
|
||
assert.ok(
|
||
rawDialer.frame !== undefined && rawDialer.frame.type === MUX.HELLO_ACK,
|
||
'manager2(白名单内的拨号方)应能以空端口表注册',
|
||
)
|
||
const bad = await new Promise((resolve) => {
|
||
const timer = setTimeout(() => resolve(undefined), 3000)
|
||
rawDialer.ws.addEventListener(
|
||
'message',
|
||
(ev) => {
|
||
clearTimeout(timer)
|
||
resolve(decodeMux(Buffer.from(ev.data)))
|
||
},
|
||
{ once: true },
|
||
)
|
||
rawDialer.ws.send(encodeJsonFrame(MUX.DIAL, 8, { target: 'w-d', port: 0 }))
|
||
})
|
||
assert.equal(JSON.parse(Buffer.from(bad.payload).toString('utf8')).error, 'bad-target')
|
||
assert.equal(srv.status().counters.dialFailed, c0.dialFailed + 2, '请求非法必须计入 dialFailed')
|
||
|
||
// ⑤ 三条计数互斥可加和:本次共 4 次 DIAL(1 放行 + 1 策略拒绝 + 2 失败)
|
||
const cf = srv.status().counters
|
||
assert.equal(cf.dial + cf.dialDenied + cf.dialFailed, c0.dial + c0.dialDenied + c0.dialFailed + 4)
|
||
})
|
||
|
||
/**
|
||
* # T23 · 拨号流**严格单向**(序 ⑭ · 数据面缺陷回归)
|
||
*
|
||
* ## 生产现象(用户可见)
|
||
* 任何 `via='relay'` 的 host(今天 = w-2)上,**同一条 keep-alive 连接的第 2 条** agent 请求
|
||
* 必回 `400 clientError`(Fastify `clientError` 兜底)⇒ 用户 `POST /api/dsh/enter` 回 **500**
|
||
* ⇒ **"登录直达工作区"整体不可用**。
|
||
*
|
||
* ## 机制(106 抓包定死,本测试逐帧复现)
|
||
* `<worker 侧连接> > 19000` 的载荷里出现 **`HTTP/1.1 200 OK …`** —— 拨号方把**入向的响应**
|
||
* 当成**出向的字节**又打了回去;worker 把它写进 agent socket,agent 拿响应行当请求行解析
|
||
* ⇒ 非法字节流 ⇒ `clientError 400`。
|
||
*
|
||
* ## 判据
|
||
* 目标端**分开记账**:"请求"(`POST …`)与"非请求字节"(garbage)。一旦入向被回灌,
|
||
* `garbage` 立即非空 ⇒ 断言点名,**不会静默通过**。
|
||
*/
|
||
test('T23 拨号流严格单向:入向的响应 ⛔ 不得被回灌进 agent socket(keep-alive 复用回归)', async (t) => {
|
||
const wSecret = randomBytes(32).toString('hex')
|
||
const mSecret = randomBytes(32).toString('hex')
|
||
/** 目标端"响应"—— 故意带响应行:被回灌时它就是那条最刺眼的证据。 */
|
||
const RESP = 'HTTP/1.1 200 OK\r\ncontent-length: 19\r\n\r\n{"isDirectory":true}'
|
||
const REQ = 'POST /fs/isdir HTTP/1.1\r\nhost: agent\r\ncontent-length: 0\r\n\r\n'
|
||
const requests = []
|
||
const garbage = []
|
||
const target = createTcpServer((s) => {
|
||
s.on('data', (c) => {
|
||
const text = c.toString()
|
||
if (text.startsWith('POST ')) {
|
||
requests.push(text)
|
||
s.write(RESP) // 一条请求回一份响应
|
||
return
|
||
}
|
||
// 到这里就是"不该出现的字节"——回灌的响应正是从这里现形
|
||
garbage.push(text)
|
||
})
|
||
})
|
||
const workerPort = await listenInRange(target, BASE + 90, BASE + SPAN - 1)
|
||
const srv = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([
|
||
['w-one', wSecret],
|
||
['manager', mSecret],
|
||
]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
dialers: new Set(['manager']),
|
||
log: () => {},
|
||
})
|
||
await srv.start()
|
||
const url = `ws://127.0.0.1:${srv.boundPort}${PATH}`
|
||
const worker = new RelayClient({ url, hostId: 'w-one', secret: wSecret, ports: [workerPort], log: () => {} })
|
||
const dialer = new RelayClient({ url, hostId: 'manager', secret: mSecret, ports: [], dialer: true, log: () => {} })
|
||
/** @type {import('node:stream').Duplex | undefined} */
|
||
let duplex
|
||
worker.start()
|
||
dialer.start()
|
||
t.after(async () => {
|
||
// ⚠️ 必须自己收尾:`onDown` 会走 `teardownDialStreams()` → `duplex.destroy(new Error('link down: …'))`,
|
||
// 而 `MuxDuplex` 上没有 `'error'` 监听 ⇒ Node 会把它抛成 **uncaughtException**,
|
||
// 表现为"测试通过但文件红"(node:test 报 async activity after the test ended)。
|
||
// 生产侧 `dialer.ts` 有 `duplex.on('error', …)`,这里补上同形的监听 + 先手动 destroy。
|
||
try {
|
||
duplex.destroy()
|
||
} catch {
|
||
/* 已断 */
|
||
}
|
||
worker.stop()
|
||
dialer.stop()
|
||
await waitFor(() => worker.status().state !== 'up' && dialer.status().state !== 'up', 2_000)
|
||
target.close()
|
||
await srv.stop()
|
||
await new Promise((r) => setTimeout(r, 50))
|
||
})
|
||
assert.ok(await waitFor(() => worker.status().state === 'up'), 'worker 未注册')
|
||
assert.ok(await waitFor(() => dialer.status().state === 'up'), '拨号方未注册')
|
||
|
||
// 与生产同形:**一条复用**的流上连发两条请求(对应 `dialer.ts` 的 `tcp.pipe(duplex).pipe(tcp)`,
|
||
// 这里用同一对 `write`/`data` 手工表达 —— 同一个 streamId,不重拨)。
|
||
const duplex0 = await dialer.openStream('w-one', workerPort)
|
||
duplex = duplex0
|
||
let inBytes = ''
|
||
duplex.on('data', (c) => {
|
||
inBytes += c.toString()
|
||
})
|
||
// 生产侧 `dialer.ts` 同形的收尾监听(见上面 t.after 的说明)
|
||
duplex.on('error', () => {})
|
||
|
||
duplex.write(REQ)
|
||
assert.ok(await waitFor(() => requests.length >= 1), '第 1 条请求没到目标端')
|
||
await new Promise((r) => setTimeout(r, 200))
|
||
duplex.write(REQ)
|
||
assert.ok(await waitFor(() => requests.length >= 2), '第 2 条请求没到目标端(复用同一条流)')
|
||
await new Promise((r) => setTimeout(r, 300))
|
||
|
||
// ① 出向:两条请求**都必须**原样到达
|
||
assert.equal(requests.length, 2, `目标端应收到 2 条请求,实收 ${requests.length} 条`)
|
||
// ② 入向:拨号侧应**恰好**收到两份响应(不重不漏)
|
||
assert.equal(inBytes.split(RESP).length - 1, 2, `拨号侧应收到 2 份响应,实见:${JSON.stringify(inBytes)}`)
|
||
// ③ 🔴 判据:入向 ⛔ 一个字节都不许回到目标端
|
||
assert.deepEqual(
|
||
garbage,
|
||
[],
|
||
`拨号流不是单向的 —— 入向字节被回灌进 agent socket(共 ${garbage.length} 段):${JSON.stringify(garbage)}`,
|
||
)
|
||
// ④ 顺带守一条:拨号流的出向/入向计数都记在 `bytesOut`/`bytesIn`
|
||
assert.ok(dialer.status().dialStreams === 1, '复用期间流的条数应恒为 1')
|
||
})
|
||
|
||
/* ─────────── T24:序⑮ 缺陷 B 回归 ─────────── */
|
||
|
||
/**
|
||
* 该口**是否有人在听**(bind 失败 = 有人在听)。
|
||
* ⚠️ 这才是"实际在听口号数"的正确量法:**零副作用** —— 它不对池口发连接,
|
||
* 所以不会触发被测缺陷、也不会污染 `stray` 计数(缺陷 B 的触发器恰恰是"对未分配口发连接")。
|
||
*/
|
||
function poolPortBusy(port, ms = 1500) {
|
||
return new Promise((resolve) => {
|
||
const s = createTcpServer()
|
||
const timer = setTimeout(() => {
|
||
try {
|
||
s.close()
|
||
} catch {
|
||
/* 已关 */
|
||
}
|
||
resolve(false)
|
||
}, ms)
|
||
s.once('error', () => {
|
||
clearTimeout(timer)
|
||
resolve(true)
|
||
})
|
||
s.listen(port, '127.0.0.1', () => {
|
||
clearTimeout(timer)
|
||
s.close(() => resolve(false))
|
||
})
|
||
})
|
||
}
|
||
|
||
/** 找一段 **n 个连续空闲** 的口,让"池到底绑了哪几个口"可被断言(否则判据 a 无从对账)。 */
|
||
async function findFreeBlock(n, lo) {
|
||
for (let b = lo; b < lo + 300; b++) {
|
||
let ok = true
|
||
for (let i = 0; i < n; i++) {
|
||
if (await poolPortBusy(b + i)) {
|
||
ok = false
|
||
break
|
||
}
|
||
}
|
||
if (ok) return b
|
||
}
|
||
throw new Error(`range ${lo}..${lo + 300} 里找不到 ${n} 个连续空闲口`)
|
||
}
|
||
|
||
/** 对某个口发一条**真实连接**(缺陷 B 的触发器)⇒ `'connected'` | `'error:CODE'` | `'timeout'`。 */
|
||
function strayConnect(port, ms = 1500) {
|
||
return new Promise((resolve) => {
|
||
const s = connect(port, '127.0.0.1')
|
||
const timer = setTimeout(() => {
|
||
s.destroy()
|
||
resolve('timeout')
|
||
}, ms)
|
||
s.on('connect', () => {
|
||
clearTimeout(timer)
|
||
s.destroy()
|
||
resolve('connected')
|
||
})
|
||
s.on('error', (e) => {
|
||
clearTimeout(timer)
|
||
resolve(`error:${e.code}`)
|
||
})
|
||
})
|
||
}
|
||
|
||
test('T24 拨号池:未分配槽位被一条连接命中后 ⛔ 不得自毁(池账=实际在听 · localPortFor 不得发死口号)', async (t) => {
|
||
/**
|
||
* ## 生产现象(用户可见)
|
||
* `POST /api/dsh/enter` 回 **500 `fetch failed: connect ECONNREFUSED 127.0.0.1:25000`**
|
||
* (47 上实测 3 条:`2026-09-17 16:00:31 / 16:00:35 / 16:00:40`,`RemoteSpawner.status` 抛出);
|
||
* 现场形态是 **"池报 64 个口、`ss -lntp` 只见 63"** —— 缺的那个 `25000` 早就是死的了。
|
||
*
|
||
* ## 机制
|
||
* `dialer.ts#onConn` 见到 `slot.key === undefined`(未分配槽位)时只 `slot.server.close()`:
|
||
* 关掉的是**服务器**,而槽位**仍在 `slots` 里、`key` 仍是 `undefined`** ⇒ 这个口永久没人听,
|
||
* 可 `localPortFor()` 的 `slots.find((s) => s.key === undefined)` **下次还会选中它**
|
||
* ⇒ 把它当"可用落点"发出去。**失败被推迟到下一次分配**,现场只有 `ECONNREFUSED`。
|
||
*
|
||
* ## 判据(缺一不可)
|
||
* a. `status().pool` 与实际在听口数**恒等**(⛔ 池不许撒谎);
|
||
* b. `localPortFor()` 返回的落点口**真的能拨通**(**真泵一次字节**,⛔ 只看返回值不算过);
|
||
* c. 该路径**不许静默**(要有计数 + 点名日志)—— 本缺陷最难查之处正是它**零日志**。
|
||
*/
|
||
const wSecret = randomBytes(32).toString('hex')
|
||
const mSecret = randomBytes(32).toString('hex')
|
||
const logs = []
|
||
const POOL = 3
|
||
const poolBase = await findFreeBlock(POOL, 46_800)
|
||
const poolPorts = Array.from({ length: POOL }, (_, i) => poolBase + i)
|
||
const echo = createTcpServer((s) => {
|
||
s.on('error', () => {})
|
||
s.on('data', (c) => s.write(c))
|
||
})
|
||
const echoPort = await listenInRange(echo, BASE + 30, BASE + SPAN - 1)
|
||
const srv = new RelayServer({
|
||
port: 0,
|
||
keys: new Map([
|
||
['w-xb', wSecret],
|
||
['manager', mSecret],
|
||
]),
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
dialers: new Set(['manager']),
|
||
log: () => {},
|
||
})
|
||
await srv.start()
|
||
const url = `ws://127.0.0.1:${srv.boundPort}${PATH}`
|
||
const worker = new RelayClient({ url, hostId: 'w-xb', secret: wSecret, ports: [echoPort], log: () => {} })
|
||
const mClient = new RelayClient({ url, hostId: 'manager', secret: mSecret, ports: [], dialer: true, log: () => {} })
|
||
const dialer = new RelayDialer({
|
||
client: mClient,
|
||
portBase: poolBase,
|
||
portSpan: POOL + 2,
|
||
poolSize: POOL,
|
||
log: (l) => logs.push(l),
|
||
})
|
||
worker.start()
|
||
mClient.start()
|
||
// ⚠️ **必须 await**:口池绑完之前 `localPortFor()` 恒返回 `undefined`(失败关闭)。
|
||
await dialer.start()
|
||
t.after(async () => {
|
||
dialer.close()
|
||
mClient.stop()
|
||
worker.stop()
|
||
await waitFor(() => mClient.status().state !== 'up' && worker.status().state !== 'up', 2_000)
|
||
echo.close()
|
||
await srv.stop()
|
||
await new Promise((r) => setTimeout(r, 50))
|
||
})
|
||
assert.ok(await waitFor(() => srv.isOnline('w-xb', OPS_NETWORK)), `worker 未注册:${logs.slice(-4).join(' | ')}`)
|
||
assert.ok(await waitFor(() => srv.isOnline('manager', OPS_NETWORK)), 'manager 未注册')
|
||
|
||
const liveCount = async () => {
|
||
let n = 0
|
||
for (const p of poolPorts) if (await poolPortBusy(p)) n++
|
||
return n
|
||
}
|
||
assert.equal(dialer.status().pool, POOL, '池就绪后 status().pool')
|
||
assert.equal(await liveCount(), POOL, `池就绪后 ${POOL} 个口都应在听:${logs.join(' | ')}`)
|
||
|
||
// ② 触发缺陷 B:对**未分配**槽位发一条连接(生产上 = 取证探针 / 任何本地扫描)
|
||
assert.equal(await strayConnect(poolBase), 'connected', '池口在听 ⇒ TCP 连接先成功,随后才被丢弃')
|
||
await new Promise((r) => setTimeout(r, 200))
|
||
|
||
// ③ 🔴 判据 a:池账必须诚实(修复前:pool=3 而实际在听=2)
|
||
const live = await liveCount()
|
||
assert.equal(
|
||
dialer.status().pool,
|
||
live,
|
||
`判据 a:池账(${dialer.status().pool}) 与实际在听(${live}) 必须恒等 —— 不等即"池在撒谎"`,
|
||
)
|
||
// ④ 判据 c:该路径不许静默(本缺陷最难查之处)
|
||
assert.equal(dialer.status().stray, 1, '未分配槽位被命中必须**有计数**')
|
||
assert.ok(
|
||
logs.some((l) => l.includes('未分配落点') && l.includes(String(poolBase))),
|
||
`必须有一条点名该口的日志:${logs.join(' | ')}`,
|
||
)
|
||
|
||
// ⑤ 🔴 判据 b:分配落点 ⇒ 必须**真的能拨通**
|
||
//(修复前这里返回的**正是刚被打死的那个口** ⇒ ECONNREFUSED,本断言就是那条生产故障的复现)
|
||
const local = dialer.localPortFor(logicalName(OPS_NETWORK, 'w-xb'), echoPort)
|
||
assert.equal(typeof local, 'number', '同网必须给落点')
|
||
assert.equal(local, poolBase, '落点应命中那个"刚被误连过"的首个槽位 —— 判据 b 才有鉴别力')
|
||
assert.equal(await roundTrip(local, 'PING', 3000), 'PING', '判据 b:落点口必须真的能拨通(修复前 = ECONNREFUSED)')
|
||
|
||
// ⑥ 分配之后再核一次:池账仍诚实、且该口确实在听
|
||
assert.equal(dialer.status().pool, await liveCount(), '分配后池账仍须与实际在听恒等')
|
||
assert.ok(await poolPortBusy(local), `已分配的落点口 ${local} 应在听`)
|
||
})
|
||
|
||
/* ═══════════════════ 序⑲ presence(节点在线态)T25–T31 ═══════════════════ */
|
||
|
||
/**
|
||
* presence 时序口径的**测试档**(生产默认 = grace 10 s / debounce 30 s / batch 1 s / TTL 45 s)。
|
||
* 单测不可能真等 40 s ⇒ 全部注入。⚠️ 口径本身**不改**,只改"等多久"。
|
||
*/
|
||
const PT = { presenceGraceMs: 150, presenceOfflineDebounceMs: 250, presenceBatchMs: 40, presenceTtlMs: 3_000 }
|
||
const sleep = (ms) => new Promise((r) => setTimeout(r, ms))
|
||
|
||
/** 起一套"订阅方 + 若干 worker"的现场(⛔ `sub` 只是普通客户端,不需要拨号方身份)。 */
|
||
async function presenceScene(t, workerIds, tune = {}) {
|
||
const secret = randomBytes(32).toString('hex')
|
||
const keys = new Map(['sub', ...workerIds].map((id) => [id, secret]))
|
||
const logs = []
|
||
const server = new RelayServer({
|
||
port: 0,
|
||
keys,
|
||
instancePortBase: BASE,
|
||
instancePortSpan: SPAN,
|
||
...PT,
|
||
...tune,
|
||
log: (l) => logs.push(l),
|
||
})
|
||
await server.start()
|
||
const url = `ws://127.0.0.1:${server.boundPort}${PATH}`
|
||
const clients = []
|
||
const mk = (hostId) => {
|
||
/**
|
||
* ⚠️ 两条硬约束(都是本轮实测踩到的,⛔ 别再踩):
|
||
* ① 非拨号方客户端**必须至少声明一个端口**(`handleHello` 回 `no-ports` 且 `retryable=false` ⇒ 永不 `up`)。
|
||
* 这里声明 `BASE` 只为满足该约束 —— relay 为它绑的是**动态回环口**(`listen(0)`),声明值只是转发目标,⛔ 不占本机端口。
|
||
* ② **每个 hostId 都必须有 key**(否则 `AUTH DENY unknown-host`)⇒ 按需登记,
|
||
* 这样 `presenceScene(t, [])` 之后仍可 `mk('w-xx')`(⛔ 不用把待造主机在入参里预先枚举一遍)。
|
||
*/
|
||
keys.set(hostId, secret)
|
||
const c = new RelayClient({ url, hostId, secret, ports: [BASE], log: () => {} })
|
||
clients.push(c)
|
||
c.start()
|
||
return c
|
||
}
|
||
t.after(async () => {
|
||
for (const c of clients) c.stop()
|
||
await server.stop()
|
||
})
|
||
const sub = mk('sub')
|
||
const up = (c) => waitFor(() => c.status().state === 'up', 5_000)
|
||
assert.ok(await up(sub), '订阅方未注册成功')
|
||
return { server, sub, mk, up, logs }
|
||
}
|
||
|
||
/** 读一次真 `/status`(顺便驱动 `statusHits` 判别器)。 */
|
||
async function readStatus(server) {
|
||
const res = await fetch(`http://127.0.0.1:${server.boundPort}/status`)
|
||
return res.json()
|
||
}
|
||
|
||
/**
|
||
* 极简**原始** WS 客户端 —— 只为"未知帧号必须被显式拒绝(⛔ 不静默丢弃)"这一条判据。
|
||
*
|
||
* 故意不做认证:它要在**未认证**那一刻发帧 ⇒ 不需要 MAC(也就不能伪造合法注册)。
|
||
* 返回服务端回的 WS close 码;`null` = 没等到 close(= 静默丢弃)。
|
||
*/
|
||
function rawWsProbe(port, muxType, ms = 3_000) {
|
||
return new Promise((resolve) => {
|
||
const sock = connect(port, '127.0.0.1')
|
||
let buf = Buffer.alloc(0)
|
||
let headerEnd = -1
|
||
let sent = false
|
||
let done = false
|
||
let timer
|
||
const finish = (code) => {
|
||
if (done) return
|
||
done = true
|
||
clearTimeout(timer)
|
||
// `end()` 而不是 `destroy()`:给上面那帧 "close 回应" 一个真的发出去的机会(见下)。
|
||
sock.end()
|
||
resolve(code)
|
||
}
|
||
timer = setTimeout(() => finish(null), ms)
|
||
sock.on('error', () => finish(null))
|
||
sock.on('data', (chunk) => {
|
||
buf = Buffer.concat([buf, chunk])
|
||
if (!sent) {
|
||
headerEnd = buf.indexOf('\r\n\r\n')
|
||
if (headerEnd < 0) return
|
||
sent = true
|
||
// 掩码的 binary mux 帧(客户端必须 mask,RFC 6455 §5.1)
|
||
const body = encodeMux(muxType, 0, Buffer.alloc(0))
|
||
const mask = randomBytes(4)
|
||
const masked = Buffer.from(body)
|
||
for (let i = 0; i < masked.length; i++) masked[i] ^= mask[i & 3]
|
||
const head = Buffer.alloc(6)
|
||
head[0] = 0x82
|
||
head[1] = 0x80 | masked.length
|
||
mask.copy(head, 2)
|
||
sock.write(Buffer.concat([head, masked]))
|
||
}
|
||
// 在握手之后的数据里找服务端发的 close 帧(opcode 0x8;服务端→客户端**不掩码**)
|
||
for (let i = headerEnd + 4; i + 3 < buf.length; ) {
|
||
const opcode = buf[i] & 0x0f
|
||
const len = buf[i + 1] & 0x7f
|
||
if (opcode === 0x8) {
|
||
const code = buf.readUInt16BE(i + 2)
|
||
/**
|
||
* 🔴 收到 close **必须回一个 close**(RFC 6455 §5.5.1),然后 `end()` 走优雅 TCP 收尾。
|
||
* 实测(本轮踩到):回都不回就直接 `destroy()` ⇒ 服务端的 ws 会一直等对端 close 帧
|
||
* (默认 30 s)⇒ 这条连接一直挂在 http server 上 ⇒ `server.stop()` 里的
|
||
* `http.close(cb)` **永不回调** ⇒ 整个测试文件在 T26 之后被父级取消
|
||
* (报 `Promise resolution is still pending but the event loop has already resolved`)。
|
||
* 这是**测试夹具**的坑,不是产品缺陷 —— 生产停机有 `closeAllConnections()` 兜底。
|
||
*/
|
||
const mask = randomBytes(4)
|
||
const payload = Buffer.alloc(2)
|
||
payload.writeUInt16BE(code, 0)
|
||
const masked = Buffer.from(payload)
|
||
for (let k = 0; k < masked.length; k++) masked[k] ^= mask[k & 3]
|
||
const head = Buffer.alloc(6)
|
||
head[0] = 0x88
|
||
head[1] = 0x80 | masked.length
|
||
mask.copy(head, 2)
|
||
try {
|
||
sock.write(Buffer.concat([head, masked]))
|
||
} catch {
|
||
/* 对端已经走了 */
|
||
}
|
||
finish(code)
|
||
return
|
||
}
|
||
i += 2 + len
|
||
}
|
||
})
|
||
sock.write(
|
||
`GET ${PATH} HTTP/1.1\r\nHost: 127.0.0.1\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n` +
|
||
`Sec-WebSocket-Key: ${randomBytes(16).toString('base64')}\r\nSec-WebSocket-Version: 13\r\n\r\n`,
|
||
)
|
||
})
|
||
}
|
||
|
||
test('T25 presence 生命周期:注册即在线 · 连接更替不闪烁 · 断连过 grace+debounce 才离线', async (t) => {
|
||
const { server, sub, mk, up } = await presenceScene(t, [])
|
||
const name = logicalName(OPS_NETWORK, 'w-p')
|
||
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
// E4 首帧即全量:`SNAP` **一帧拿全**,⛔ 不是逐 host 拉
|
||
assert.equal(sub.presenceStatus().snapFrames, 1, '首帧必须是 SNAP')
|
||
assert.equal(sub.presenceStatus().pushFrames, 0, '`SNAP` 之前 ⛔ 不得有增量帧')
|
||
|
||
// E1 稳态:没有任何状态变化 ⇒ ⛔ 一个帧都不推
|
||
const idle = sub.presenceStatus().pushFrames
|
||
await sleep(300)
|
||
assert.equal(sub.presenceStatus().pushFrames, idle, '稳态必须 0 帧(变化驱动,无变化不推)')
|
||
|
||
const w = mk('w-p')
|
||
assert.ok(await up(w), 'worker 未注册')
|
||
assert.ok(await waitFor(() => sub.presenceStatus().pushFrames === idle + 1, 3_000), '上线应恰好推 1 帧')
|
||
assert.equal(sub.presenceStatus().lastFrameEntries, 1, '该帧只应带这 1 条')
|
||
assert.equal(sub.presenceOf(name)?.online, true, '上线事件未更新本地镜像')
|
||
assert.equal(sub.presenceOf(name)?.devices, 1)
|
||
|
||
/**
|
||
* E6(聚合口径的**可观测后果**):同 hostId 的第二条连接注册时,服务端会**先顶掉旧会话再接入新会话**
|
||
* (`handleHello` 的 `superseded by new session`)⇒ 此刻"一条连接关了、另一条开了"。
|
||
* 聚合口径要求:**这中间不许产生任何状态事件**(否则每次重连都会让上层看到一次闪烁 = 净退化)。
|
||
*/
|
||
const w2 = mk('w-p')
|
||
assert.ok(await up(w2), '第二条连接未注册')
|
||
await sleep(400)
|
||
assert.equal(sub.presenceStatus().pushFrames, idle + 1, '连接更替(顶旧接新)⛔ 不得产生任何新帧')
|
||
assert.equal(sub.presenceOf(name)?.online, true, '更替期间必须**始终**在线(⛔ 不得闪一下离线)')
|
||
assert.equal(sub.presenceOf(name)?.devices, 1, '旧会话已被顶掉 ⇒ 活连接数仍是 1(devices 必须诚实)')
|
||
|
||
// E5:真正全断 ⇒ 先过 grace 仍在线,再过 debounce 才转离线
|
||
w2.stop()
|
||
await sleep(120)
|
||
const during = sub.presenceOf(name)
|
||
assert.equal(during?.online, true, `≤ grace(${PT.presenceGraceMs}ms) 必须仍在线(防抖动闪烁)`)
|
||
/**
|
||
* 倒计时(`offlineInMs`)**只认 `/status`(或 `SNAP`)上的值,⛔ 不从推送镜像里读**:
|
||
* 它是"生成那一刻"的相对量,而推送是**变化驱动**的(无变化不推)⇒ 镜像里那个值一发出就过期。
|
||
* 想让它实时更新就只能周期性推帧 —— 那正是 E1「稳态 0 帧」要干掉的东西。
|
||
* 换句话说:**事实走推送,带时钟刻度的心跳量走兜底读取**(D5 的兜底就不只是"降级可用",而是分工)。
|
||
*/
|
||
const st = await readStatus(server)
|
||
const mine = st.presence.find((p) => p.name === name)
|
||
assert.equal(typeof mine?.offlineInMs, 'number', 'grace 窗口内 `/status` 必须给出"还有多久转离线"')
|
||
assert.ok(
|
||
mine.offlineInMs > 0 && mine.offlineInMs <= PT.presenceGraceMs + PT.presenceOfflineDebounceMs,
|
||
`倒计时必须落在 (0, grace+debounce] 内,实测 ${mine?.offlineInMs}`,
|
||
)
|
||
assert.equal(sub.presenceStatus().pushFrames, idle + 1, 'grace 窗口内 ⛔ 一个帧都不许推(E2:一次变化 ≤1 帧)')
|
||
assert.ok(await waitFor(() => sub.presenceOf(name)?.online === false, 4_000), '过 grace+debounce 必须转离线')
|
||
assert.equal(sub.presenceStatus().pushFrames, idle + 2, '离线应恰好再推 1 帧')
|
||
})
|
||
|
||
test('T26 线协议:SUB/UNSUB/PRESENCE/SNAP 帧号**末尾追加**且与既有集合不重叠', async (t) => {
|
||
// 帧号:既有的 0x01–0x0f 语义一字未动,新帧全部 > 0x0f
|
||
assert.equal(MUX.SUB, 0x10)
|
||
assert.equal(MUX.UNSUB, 0x11)
|
||
assert.equal(MUX.PRESENCE, 0x12)
|
||
assert.equal(MUX.SNAP, 0x13)
|
||
for (const t2 of [MUX.SUB, MUX.UNSUB, MUX.PRESENCE, MUX.SNAP]) {
|
||
assert.ok(t2 > MUX.DIAL_ACK, `新帧号 ${t2} 必须**追加**在既有分配表末尾(⛔ 不改既有语义)`)
|
||
}
|
||
const codes = Object.values(MUX)
|
||
assert.equal(new Set(codes).size, codes.length, '帧号必须两两不同')
|
||
const back = decodeMux(encodeJsonFrame(MUX.SUB, 0, { all: true }))
|
||
assert.equal(back.type, MUX.SUB, 'SUB 帧编解码往返失败')
|
||
|
||
// 🔴 "未知帧号 ⇒ 显式报错,⛔ 不静默丢弃"(本线头号教训):用一个服务端**不认识**的帧号打它
|
||
const { server } = await presenceScene(t, [])
|
||
const before = server.status().counters.authFailed
|
||
const code = await rawWsProbe(server.boundPort, 0x99)
|
||
assert.equal(code, 1008, `未知帧号必须被**显式拒绝**(期望 close 1008 policy-violation,实得 ${code})`)
|
||
assert.equal(server.status().counters.authFailed, before + 1, '拒绝必须**有计数**(否则等于没记)')
|
||
})
|
||
|
||
test('T27 批合并:同一 1 s 窗口内 N 次状态变化只推 1 帧(E2/E4)', async (t) => {
|
||
const ids = Array.from({ length: 6 }, (_, i) => `w-b${i}`)
|
||
// 窗口取 300ms:6 次回环握手远小于它 ⇒ 判据确定性足够(⛔ 不靠"碰巧合上")
|
||
const { sub, mk, up } = await presenceScene(t, [], { presenceBatchMs: 300 })
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
const base = sub.presenceStatus().pushFrames
|
||
|
||
const ws = ids.map((id) => mk(id))
|
||
await Promise.all(ws.map((c) => up(c)))
|
||
const name0 = logicalName(OPS_NETWORK, ids[0])
|
||
assert.ok(await waitFor(() => sub.presenceOf(name0)?.online === true, 3_000), '上线事件未到达')
|
||
await sleep(900) // 让所有可能的批窗口都过去
|
||
assert.equal(sub.presenceStatus().pushFrames, base + 1, `${ids.length} 台同窗口上线 ⇒ 必须合并成**1 帧**`)
|
||
assert.equal(sub.presenceStatus().lastFrameEntries, ids.length, '这一帧必须**带数组**(6 条),⛔ 不是逐个 host 一条')
|
||
|
||
// 反向:同窗口全部下线 ⇒ 同样只 1 帧(合并方向也要成立,⛔ 不能只测上线)
|
||
for (const c of ws) c.stop()
|
||
await waitFor(() => sub.presenceOf(name0)?.online === false, 5_000)
|
||
await sleep(900)
|
||
assert.equal(sub.presenceStatus().pushFrames, base + 2, `${ids.length} 台同窗口下线 ⇒ 必须合并成 1 帧`)
|
||
assert.equal(sub.presenceStatus().lastFrameEntries, ids.length, '离线帧同样必须带全 6 条')
|
||
})
|
||
|
||
test('T28 订阅可见性**只收窄**:跨网订阅必须显式拒绝 + 计数(⛔ 不静默返空,E8/D6)', async (t) => {
|
||
const { server, sub } = await presenceScene(t, ['w-p'])
|
||
sub.subscribePresence([logicalName('u:5', 'd1')])
|
||
assert.ok(await waitFor(() => sub.presenceStatus().rejected === 1, 3_000), '跨网订阅必须被**显式拒绝**')
|
||
const st = await readStatus(server)
|
||
assert.equal(st.counters.rejected, 1, '服务端必须**有计数**(⛔ 静默返空 = 假绿)')
|
||
assert.equal(sub.presenceOf('u:5/d1'), undefined, '被拒的订阅 ⛔ 不得留下任何条目')
|
||
assert.equal(sub.presenceStatus().state, 'idle', '被拒后状态必须回到 idle(上层据此回退 /status)')
|
||
assert.ok((await readStatus(server)).counters.subs === 0, '被拒的订阅 ⛔ 不得计入 subs')
|
||
})
|
||
|
||
test('T29 TTL 安全网:漏掉 close 事件的"幽灵连接"超 TTL 后被摘掉并收口离线(E5 第三支)', async (t) => {
|
||
const { server, sub } = await presenceScene(t, [], {
|
||
presenceTtlMs: 400,
|
||
presenceGraceMs: 100,
|
||
presenceOfflineDebounceMs: 150,
|
||
})
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
const name = logicalName(OPS_NETWORK, 'w-ghost')
|
||
|
||
/**
|
||
* **故障注入**(⛔ 不是模拟业务,而是模拟**漏掉了 close 事件**这一种故障):
|
||
* 只往 presence 表里记一条"连接",既不建真连接、也就永远不会收到 close ⇒ 这正是 TTL 要兜的事。
|
||
* 不这么做的话,"漏事件 ⇒ 永久假在线"这条路径**根本无法被触发**(正常路径总会 dropSession)。
|
||
*/
|
||
server.presenceTouch({ id: 'ghost-1', hostId: 'w-ghost', network: OPS_NETWORK, ports: new Set() })
|
||
assert.ok(await waitFor(() => sub.presenceOf(name)?.online === true, 3_000), '幽灵连接应先被认定为在线')
|
||
|
||
assert.ok(
|
||
await waitFor(() => sub.presenceOf(name)?.online === false, 5_000),
|
||
'TTL 安全网未能自愈 ⇒ 漏掉 close 的事件会**永久**留在册',
|
||
)
|
||
})
|
||
|
||
test('T30 主路径=订阅 / 兜底=/status:订阅新鲜时 ⛔ 不读兜底;不可用时**必须**回退(D5/E7)', async (t) => {
|
||
const calls = []
|
||
const name = logicalName(OPS_NETWORK, 'w-1')
|
||
const rv = (presence) =>
|
||
new RelayRendezvous({
|
||
dialTargetUrl: 'wss://example.invalid/dshs-relay',
|
||
addressOf: () => '127.0.0.1:19100',
|
||
online: (n) => {
|
||
calls.push(`/status兜底:${n}`)
|
||
return true
|
||
},
|
||
presence,
|
||
})
|
||
|
||
// ① 订阅新鲜 ⇒ **以它为准**,兜底一次都不读
|
||
assert.ok((await rv(() => true).resolve(name)) !== undefined, '订阅说在线 ⇒ 必须解析成功')
|
||
assert.equal(calls.length, 0, '订阅新鲜时 ⛔ 不许读兜底(/status)')
|
||
// ② 订阅说"不在" ⇒ 直接不认识(同样不读兜底)
|
||
assert.equal(await rv(() => false).resolve(name), undefined, '订阅说离线 ⇒ 必须回 undefined')
|
||
assert.equal(calls.length, 0, '订阅能给答案时 ⛔ 不许读兜底')
|
||
// ③ 订阅"不知道"(undefined)⇒ **必须回退**,且照样拿到在线态
|
||
assert.ok((await rv(() => undefined).resolve(name)) !== undefined, '订阅不可用时必须能回退且不瞎')
|
||
assert.equal(calls.length, 1, '回退必须**真的调用**兜底(否则就是"订阅一断就全瞎")')
|
||
|
||
// E7 最终一致:在线态**没有任何**跨节点同步通道(两台 relay 各管各的 ⇒ 天然无脑裂源)
|
||
const secret = randomBytes(32).toString('hex')
|
||
const keys = new Map([['w-e7', secret]])
|
||
const s1 = new RelayServer({ port: 0, keys, instancePortBase: BASE, instancePortSpan: SPAN, ...PT, log: () => {} })
|
||
const s2 = new RelayServer({ port: 0, keys, instancePortBase: BASE, instancePortSpan: SPAN, ...PT, log: () => {} })
|
||
await s1.start()
|
||
await s2.start()
|
||
const c = new RelayClient({
|
||
url: `ws://127.0.0.1:${s1.boundPort}${PATH}`,
|
||
hostId: 'w-e7',
|
||
secret,
|
||
ports: [BASE],
|
||
log: () => {},
|
||
})
|
||
t.after(async () => {
|
||
c.stop()
|
||
await s1.stop()
|
||
await s2.stop()
|
||
})
|
||
c.start()
|
||
assert.ok(await waitFor(() => c.status().state === 'up', 5_000), 'client 未注册')
|
||
assert.equal(s1.status().presence.length, 1, '第一台应有该 host')
|
||
assert.equal(s2.status().presence.length, 0, '⛔ 另一台**不得**知道它(有同步才是缺陷:那是脑裂源)')
|
||
})
|
||
|
||
test('T31 presence 判别器:subs / pushed / rejected / statusHits 都能被断言(⛔ 不许只写日志)', async (t) => {
|
||
const { server, sub, mk, up } = await presenceScene(t, [])
|
||
assert.equal((await readStatus(server)).counters.subs, 0, '没人订阅 ⇒ subs=0')
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
assert.equal((await readStatus(server)).counters.subs, 1, '订阅生效 ⇒ subs=1(gauge)')
|
||
|
||
const pushed0 = (await readStatus(server)).counters.pushed
|
||
const w = mk('w-p')
|
||
assert.ok(await up(w), 'worker 未注册')
|
||
assert.ok(await waitFor(() => sub.presenceStatus().pushFrames >= 1, 3_000), '未收到推送')
|
||
await sleep(400)
|
||
const pushed1 = (await readStatus(server)).counters.pushed
|
||
assert.ok(pushed1 > pushed0, 'pushed 必须随真实推送增长(否则判别器是死的)')
|
||
await sleep(400)
|
||
assert.equal((await readStatus(server)).counters.pushed, pushed1, '稳态下 pushed 必须**停住不走**(E1)')
|
||
|
||
// statusHits:两次读数之差 = 1 ⇒ 期间**没有别人**在读 /status
|
||
const a = (await readStatus(server)).counters.statusHits
|
||
const b = (await readStatus(server)).counters.statusHits
|
||
assert.equal(b - a, 1, `两次读数之差应为 1(期间只有本测试在读),实得 ${b - a}`)
|
||
|
||
/**
|
||
* `snaps`(E4 的机器可读判据):有订阅者却 `snaps = 0` ⇒ 首帧走的不是 `SNAP`。
|
||
* 同理把 `presenceTiming` 钉住 —— 探针(`OBS-13`)拿它与参数表 `PRESENCE_*` 对口径,
|
||
* 对不上就是**口径漂移**(改了默认值却没改表 ⇒ 表在撒谎)。
|
||
*/
|
||
const st = await readStatus(server)
|
||
assert.equal(st.counters.snaps, 1, '本轮只有 1 次订阅 ⇒ 必须恰好 1 帧 SNAP')
|
||
assert.ok(st.counters.snaps <= st.counters.pushed, 'snaps 是 pushed 的子集(⛔ 不得大于)')
|
||
assert.deepEqual(
|
||
{
|
||
graceMs: st.presenceTiming.graceMs,
|
||
offlineDebounceMs: st.presenceTiming.offlineDebounceMs,
|
||
batchMs: st.presenceTiming.batchMs,
|
||
ttlMs: st.presenceTiming.ttlMs,
|
||
subMax: st.presenceTiming.subMax,
|
||
},
|
||
{
|
||
graceMs: PT.presenceGraceMs,
|
||
offlineDebounceMs: PT.presenceOfflineDebounceMs,
|
||
batchMs: PT.presenceBatchMs,
|
||
ttlMs: PT.presenceTtlMs,
|
||
subMax: 0,
|
||
},
|
||
'`presenceTiming` 必须把注入的时序口径如实下发(探针 `OBS-13` 靠它对口径)',
|
||
)
|
||
|
||
sub.stop()
|
||
await sleep(200)
|
||
assert.equal((await readStatus(server)).counters.subs, 0, '连接断了订阅必须随之消失(→ 回到 0)')
|
||
})
|
||
|
||
test('T32 落点不丢:注册即在线那一帧必须带**非 0** 落点,且端口变更会被推送(序⑲ 收口实测踩到的假死)', async (t) => {
|
||
const { sub, mk, up } = await presenceScene(t, [])
|
||
const name = logicalName(OPS_NETWORK, 'w-l')
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
|
||
/**
|
||
* 🔴 这一条对应一个**实测踩到的假死**:`presenceTouch` 曾在 `ensureEndpoint` **之前**调用 ⇒
|
||
* 首帧里 `localPorts[].localPort = 0`;而发布只在"在线态翻转"时发生 ⇒ 那个 0 永远修不回来 ⇒
|
||
* 订阅方(Manager)`addressOf` 查不到落点 ⇒ 实例页**打不开但不报错**。
|
||
* 判据必须卡在"**订阅之后**才上线的主机"上:老主机早就在册,落点已被别的路径补过。
|
||
*/
|
||
const w = mk('w-l')
|
||
assert.ok(await up(w), 'worker 未注册')
|
||
assert.ok(await waitFor(() => sub.presenceOf(name)?.online === true, 3_000), '上线事件未到达')
|
||
const lp0 = sub.presenceOf(name)?.localPorts ?? []
|
||
assert.equal(lp0.length, 1, '首帧必须带该 host 的落点条目')
|
||
assert.equal(lp0[0].port, BASE, '落点条目的 port 必须与声明一致')
|
||
assert.ok(lp0[0].localPort > 0, `落点口号必须**非 0**(实测 ${lp0[0].localPort} ⇒ 0 = addressOf 查不到 ⇒ 页面假死)`)
|
||
|
||
// 运行期加一个端口(`PORT_ADD`)⇒ 新落点也必须**推给订阅方**(否则同样查不到)
|
||
const frames0 = sub.presenceStatus().pushFrames
|
||
const NEW_PORT = BASE + 1
|
||
assert.equal(await w.addPort(NEW_PORT), true, 'PORT_ADD 未被接受')
|
||
assert.ok(await waitFor(() => (sub.presenceOf(name)?.ports ?? []).length === 2, 3_000), '端口变更未推送')
|
||
await sleep(400)
|
||
const lp1 = sub.presenceOf(name)?.localPorts ?? []
|
||
assert.equal(lp1.length, 2, '两个端口都必须在落点表里')
|
||
assert.ok(
|
||
lp1.every((x) => x.localPort > 0),
|
||
`新端口落点同样必须非 0:${JSON.stringify(lp1)}`,
|
||
)
|
||
assert.equal(sub.presenceStatus().pushFrames, frames0 + 1, '端口变更 = 一次状态变化 ⇒ 恰好 1 帧(E2 配额内)')
|
||
})
|
||
|
||
/* ═══════════ 序㉑ P-1 修复(门的判据 = 订阅已建立 ∧ 链路活着)T33–T34 ═══════════ */
|
||
|
||
/**
|
||
* 🔴 **P-1 回归**(在册缺陷,2026-09-17 序 ⑳ 实测)。
|
||
*
|
||
* 病根:门(`presenceFresh()`)原判据 = "最近一次 presence **载荷**距今 ≤ TTL"。而 presence 是
|
||
* **变化驱动**的 —— 稳态下一帧都不推 ⇒ 45 s 后必然过期 ⇒ 门自己重开、`/status` 轮询照旧在跑
|
||
* (真机实测降幅仅 **1.10×**,设计目标 ≥ 10×)。**"没有变化"被读成了"没有数据"**。
|
||
*
|
||
* 本用例把"稳态 + 超过 TTL"这个组合钉死:帧数必须仍是 0(E1 不破),门必须**仍然关着**。
|
||
* ⚠️ 旧实现下必红(载荷年龄 > TTL ⇒ `presenceFresh()` 翻假)—— 这就是"先红后绿"的那条断言。
|
||
*/
|
||
test('T33 P-1:稳态零帧下门不得自己重开(判据 = 订阅已建立 ∧ 链路活着,⛔ 不是载荷年龄)', async (t) => {
|
||
/**
|
||
* ⚠️ `hbSec: 1` 是**夹具前提**,不是产品口径:测试档把 presence TTL 压到 3 s,而生产心跳是 15 s
|
||
* ⇒ 不压心跳的话,relay 的 `presenceDevices`(`now − conns[ts] ≤ ttl`)会在 3 s 后把连接判死
|
||
* ⇒ 自己制造出"离线→在线"的状态变化(**假帧**),把本用例的稳态前提破坏掉。
|
||
* 生产里 `TTL(45 s) > 心跳(15 s)` ⇒ 不存在这个组合(这正是 TTL 因子取 3 的原因)。
|
||
*/
|
||
const { sub } = await presenceScene(t, [], { hbSec: 1 })
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
|
||
/**
|
||
* ⚠️ 先**等静默下来**再取基线:`SNAP` 之后还有一次**由落点落地驱动**的强制推(序⑲ T32 的假死修复)
|
||
* —— 它属于"上线那一件事"的收尾,⛔ 不是稳态帧。不先等它,基线就取在稳态之前。
|
||
*/
|
||
await sleep(1_000)
|
||
const idle = sub.presenceStatus().pushFrames
|
||
// 旧判据看的是"**最近一次载荷**"⇒ 前提要按它的口径算(`max(snap,push)` ⇒ 取**年龄最小**的那个)。
|
||
const ageOf = (s) => Math.min(s.lastSnapAgoMs ?? 0, s.lastPushAgoMs ?? s.lastSnapAgoMs ?? 0)
|
||
const before = sub.presenceStatus()
|
||
|
||
// 稳态(无任何状态变化)⇒ 一帧都不推;等到**超过 TTL**(测试档 3 s)再看门。
|
||
await sleep(PT.presenceTtlMs + 1_200)
|
||
const st = sub.presenceStatus()
|
||
assert.equal(st.pushFrames, idle, '稳态必须仍然是 0 帧(E1:变化驱动,无变化不推)')
|
||
assert.equal(st.state, 'subscribed', '链路上订阅应仍生效(服务端心跳在 ⇒ 半开巡检不会断它)')
|
||
assert.ok(
|
||
ageOf(st) >= PT.presenceTtlMs,
|
||
`前提未成立:最近一次载荷年龄应已超过 TTL,实得 ${ageOf(st)}ms(旧=${ageOf(before)}ms;否则这条用例证不了 P-1)`,
|
||
)
|
||
assert.equal(
|
||
st.fresh,
|
||
true,
|
||
`载荷 ${ageOf(st)}ms 没来、但链路活着(入站静默 ${st.lastInboundAgoMs}ms ≤ 上界 ${st.linkSilentMaxMs}ms)⇒ 门必须保持关闭`,
|
||
)
|
||
assert.equal(sub.presenceFresh(), true, 'P-1 回归:`presenceFresh()` ⛔ 不得因"没有变化"而翻假')
|
||
})
|
||
|
||
/**
|
||
* **反方向**(失败关闭):链路活着 ⛔ 不足以判"新鲜" —— 还必须**订阅真的生效**。
|
||
* 少这一条,"永远返回 true"也能让 T33 绿 ⇒ 过修无法被发现。
|
||
*/
|
||
test('T34 P-1 反向:未订阅 / 已退订 / 被拒 ⇒ 一律不得判"新鲜"(失败关闭,回退 `/status`)', async (t) => {
|
||
const { sub } = await presenceScene(t, [])
|
||
// ① 从没订阅 ⇒ 不新鲜
|
||
assert.equal(sub.presenceFresh(), false, '未订阅 ⇒ 必须回退 /status')
|
||
assert.equal(sub.presenceStatus().fresh, false, '状态视图必须如实报门是开的')
|
||
|
||
// ② 订阅生效 ⇒ 新鲜;退订 ⇒ **立刻**回到不新鲜(链路还活着也不例外)
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '订阅未生效')
|
||
assert.equal(sub.presenceFresh(), true, '订阅刚生效 ⇒ 应判新鲜(否则主路径永远用不上)')
|
||
sub.unsubscribePresence()
|
||
assert.equal(sub.presenceFresh(), false, '退订 ⇒ 必须立刻不新鲜(⛔ 不许靠"链路活着"继续给绿)')
|
||
|
||
// ③ 被**显式拒绝**的订阅(跨网)⇒ 同样不新鲜(⛔ 静默返空与"本网没人"同形,本线头号教训)
|
||
sub.subscribePresence(['u:5/d1'])
|
||
assert.ok(await waitFor(() => sub.presenceStatus().rejected === 1, 3_000), '跨网订阅必须被显式拒绝')
|
||
assert.equal(sub.presenceFresh(), false, '被拒的订阅 ⇒ 必须回退 /status')
|
||
|
||
// ④ 链路断 ⇒ 订阅随之消失 ⇒ 不新鲜(`onPresenceDown` 的唯一职责)
|
||
sub.subscribePresence()
|
||
assert.ok(await waitFor(() => sub.presenceStatus().state === 'subscribed', 3_000), '重新订阅未生效')
|
||
sub.stop()
|
||
await sleep(200)
|
||
assert.equal(sub.presenceFresh(), false, '链路断 ⇒ 必须不新鲜(否则会拿过期镜像当事实)')
|
||
})
|
||
|
||
/* ═══════════ 序㉑ P-2 修复(键口径 ⇒ 抽纯函数)T35–T36 ═══════════ */
|
||
|
||
/**
|
||
* 🔴 **P-2**(在册缺陷):`translateEndpoint` 的键口径 —— `hostVia` / `relayEndpoints` / 拨号池
|
||
* 全按**逻辑名**建键,而调用方 `RemoteSpawner.translateEndpoint(host.hostId, …)` 只给得到
|
||
* **裸 hostId** ⇒ `hostVia.get(hostId)` 恒 `undefined` ⇒ 早退原样透传 ⇒ **闭包整体是死分支**。
|
||
*
|
||
* 本用例钉住**判定本体**(抽成的纯函数);键口径那一半由 T36 钉。
|
||
*/
|
||
test('T35 P-2:`relayEndpointTarget` 四支判定(透传 / 拨号 / 快照 / 失败关闭)⛔ 不误触发拨号池', () => {
|
||
let dialedCalls = 0
|
||
const dial = (p) => () => {
|
||
dialedCalls += 1
|
||
return p
|
||
}
|
||
|
||
// ① 未知 host(不在 `dsh_hosts`)⇒ 原样透传,且**不许**碰拨号池
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: false, via: undefined, dialedPort: dial(25000), snapshotLocalPort: 41000 }),
|
||
{ kind: 'passthrough', why: 'unknown-host' },
|
||
'未知 host 必须保持老行为(单机 / 默认 host 不受影响)',
|
||
)
|
||
assert.equal(dialedCalls, 0, '未知 host ⛔ 不许查拨号池(`localPortFor` 会**按需绑池口**,是有副作用的调用)')
|
||
|
||
// ② `via` 不是 relay ⇒ 原样透传(隧道同号反向转发,不需要翻译);同样不碰池
|
||
for (const via of ['local', 'manager-ssh']) {
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: true, via, dialedPort: dial(25000) }),
|
||
{ kind: 'passthrough', why: 'not-relay' },
|
||
`via=${via} 两侧口号相同 ⇒ 翻译既不需要也不该做`,
|
||
)
|
||
}
|
||
assert.equal(dialedCalls, 0, '非 relay ⛔ 不许查拨号池')
|
||
|
||
// ③ via=relay ⇒ ①拨号落点优先(R5:落点在 Manager 本机 ⇒ relay 换机器也成立)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: true, via: 'relay', dialedPort: dial(25000), snapshotLocalPort: 41000 }),
|
||
{ kind: 'local', port: 25000, via: 'dialed' },
|
||
)
|
||
// ④ 拨号拿不到 ⇒ 回落 relay 快照
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: true, via: 'relay', dialedPort: dial(undefined), snapshotLocalPort: 41000 }),
|
||
{ kind: 'local', port: 41000, via: 'snapshot' },
|
||
)
|
||
// ⑤ 两条都没有 ⇒ **失败关闭**(⛔ 不是原样透传:那会拿 Worker 侧口号拨 Manager 本机)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: true, via: 'relay', dialedPort: dial(undefined) }),
|
||
{ kind: 'unreachable', why: 'no-dialed-port' },
|
||
)
|
||
// ⑥ `0` 是"没有落点"的哨兵值(不是合法口号)⇒ 必须仍判失败关闭
|
||
assert.equal(
|
||
relayEndpointTarget({ known: true, via: 'relay', dialedPort: dial(0), snapshotLocalPort: 0 }).kind,
|
||
'unreachable',
|
||
'落点 0 ⇒ 无落点(与 relay `/status` 同口径)',
|
||
)
|
||
})
|
||
|
||
test('T36 P-2:键口径 —— 裸 `hostId` 必须经 `hostNameIndex` 换到逻辑名(⛔ 闭包不得再拿 hostId 当键)', async () => {
|
||
const idx = hostNameIndex([
|
||
{ id: 'w-1', networkId: 'ops' },
|
||
{ id: 'w-2', networkId: '' }, // 空 ⇒ 归属网取兜底(与 DB 列默认值同口径)
|
||
{ id: 'd1', networkId: 'u:5' },
|
||
])
|
||
assert.equal(idx.get('w-1'), 'ops/w-1')
|
||
assert.equal(idx.get('w-2'), 'ops/w-2', '空 network_id 必须按兜底网补全(否则与 DB 行写的键不一致)')
|
||
assert.equal(idx.get('d1'), 'u:5/d1', '跨网同 hostId 各算一台(P0-3)')
|
||
assert.equal(idx.get('ops/w-2'), undefined, '索引的键是**裸 hostId** ⇒ 拿逻辑名查不到(两侧口径必须显式转换)')
|
||
assert.equal(hostNameIndex([{ id: 'x', networkId: '' }], 'u:9').get('x'), 'u:9/x', '兜底网可注入(⛔ 不写死 ops)')
|
||
|
||
/**
|
||
* 源码级守卫(这类"整个闭包静默失效"的缺陷只靠运行时断言抓不到 —— 无实例时分支根本不执行):
|
||
* ⛔ 闭包不得再出现"拿裸 hostId 当控制面表的键"的写法。
|
||
*/
|
||
const src = await readFile(new URL('../src/web/server.ts', import.meta.url), 'utf8')
|
||
// ⚠️ 只看**代码行**:注释里会引用反例("原实现直接 `hostVia.get(hostId)`…"),拿整文件匹配会自伤。
|
||
const code = src
|
||
.split('\n')
|
||
.filter((l) => !/^\s*(\/\/|\*|\/\*)/.test(l))
|
||
.join('\n')
|
||
assert.equal(code.includes('hostVia.get(hostId)'), false, '⛔ 不得再用裸 hostId 查 `hostVia`(恒 undefined ⇒ 死分支)')
|
||
assert.equal(code.includes('localPortFor(hostId'), false, '⛔ 不得再用裸 hostId 查拨号池(同上)')
|
||
assert.equal(code.includes('relayEndpoints.get(`${hostId}'), false, '⛔ 快照回退键同样必须是逻辑名')
|
||
assert.ok(code.includes('hostNameById.get(hostId)'), '闭包必须经 `hostNameById` 换到逻辑名')
|
||
})
|
||
|
||
/* ═══════════ 序㉒ P-2b 修复(候选链补「订阅推送落点」一级)T37 ═══════════ */
|
||
|
||
/**
|
||
* 🔴 **P-2b**(在册缺陷):`relayEndpointTarget` 只有「拨号落点 → relay `/status` 快照」**两级**,
|
||
* 而地址解析链(`src/web/server.ts#RelayRendezvous.addressOf`)是**三级**:
|
||
* ① 拨号落点 → ② **订阅推送落点**(`presenceLocalPort`)→ ③ relay 快照。
|
||
*
|
||
* 为什么在 P-1 修好之后这条变成**真缺陷**:P-1 把门判据改成「订阅已建立 ∧ 链路活着」之后,
|
||
* 订阅新鲜期**长期成立** ⇒ 快照刷新(`relayEndpoints`)**趋冷**,而拨号池在"该 host 的槽位
|
||
* 分不出来"(跨网被拒 / 池满 / 尚未绑口)时也给不出落点 ⇒ 本判定会落到"两级都没有"
|
||
* ⇒ **判实例不可达(失败关闭)**,尽管**订阅推送里明明有落点**(同一时刻 `addressOf` 能答出来)。
|
||
* ⇒ 修法 = 把订阅推送插成 **②' 级**,与 `addressOf` 的三级链**逐级对齐**。
|
||
*
|
||
* ⚠️ 本用例只钉**优先级与失败语义**;"闭包有没有把这一级传进来"由同文件 `T38` 的源码级守卫钉。
|
||
*/
|
||
test('T37 P-2b:候选链三级(拨号 → 订阅推送 → relay 快照)逐支可判,⛔ 不误触发拨号池', () => {
|
||
let dialedCalls = 0
|
||
const dial = (p) => () => {
|
||
dialedCalls += 1
|
||
return p
|
||
}
|
||
|
||
// ① 三级全有 ⇒ **拨号落点优先**(R5:落点在 Manager 本机 ⇒ relay 换机器也成立)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({
|
||
known: true,
|
||
via: 'relay',
|
||
dialedPort: dial(25000),
|
||
pushedLocalPort: 37057,
|
||
snapshotLocalPort: 41000,
|
||
}),
|
||
{ kind: 'local', port: 25000, via: 'dialed' },
|
||
'拨号落点是第一优先(它与订阅推送、快照三者必须逐支可分辨)',
|
||
)
|
||
|
||
// ② 拨号分不出槽位,但**订阅推送里有落点** ⇒ 取推送(**P-2b 的实体**:修前这一支落到 ③/失败关闭)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({
|
||
known: true,
|
||
via: 'relay',
|
||
dialedPort: dial(undefined),
|
||
pushedLocalPort: 37057,
|
||
snapshotLocalPort: 41000,
|
||
}),
|
||
{ kind: 'local', port: 37057, via: 'pushed' },
|
||
'P-2b:拨号给不出时,订阅推送的落点必须先于 relay 快照被采用(与 `addressOf` 对齐)',
|
||
)
|
||
|
||
// ③ 订阅不新鲜 / 该 host 不在推送范围 ⇒ 回落到 relay 快照(③ 级语义**一行未改**)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({
|
||
known: true,
|
||
via: 'relay',
|
||
dialedPort: dial(undefined),
|
||
pushedLocalPort: undefined,
|
||
snapshotLocalPort: 41000,
|
||
}),
|
||
{ kind: 'local', port: 41000, via: 'snapshot' },
|
||
'订阅给不出 ⇒ 必须仍能落到快照(⛔ 不是失败关闭)',
|
||
)
|
||
|
||
// ④ 三级都没有 ⇒ **失败关闭**(⛔ 绝不原样透传:那会拿 Worker 侧口号拨 Manager 本机)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: true, via: 'relay', dialedPort: dial(undefined) }),
|
||
{ kind: 'unreachable', why: 'no-dialed-port' },
|
||
'三级全无 ⇒ 失败关闭',
|
||
)
|
||
|
||
// ⑤ `0` 是"没有落点"的哨兵(与 relay `/status` 同口径)⇒ 订阅那级也必须按"没有"处理
|
||
assert.deepEqual(
|
||
relayEndpointTarget({
|
||
known: true,
|
||
via: 'relay',
|
||
dialedPort: dial(0),
|
||
pushedLocalPort: 0,
|
||
snapshotLocalPort: 41000,
|
||
}),
|
||
{ kind: 'local', port: 41000, via: 'snapshot' },
|
||
'落点 0 ⇒ 视为无落点,逐级下探(⛔ 不许把 0 当合法口号)',
|
||
)
|
||
|
||
// ⑥ 未知 host / 非 relay ⇒ **原样透传**,且订阅那一级同样不许改写结果(⛔ 不是"新增一条翻译路径")
|
||
const callsBefore6 = dialedCalls
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: false, via: undefined, dialedPort: dial(25000), pushedLocalPort: 37057 }),
|
||
{ kind: 'passthrough', why: 'unknown-host' },
|
||
)
|
||
assert.deepEqual(
|
||
relayEndpointTarget({ known: true, via: 'local', dialedPort: dial(25000), pushedLocalPort: 37057 }),
|
||
{ kind: 'passthrough', why: 'not-relay' },
|
||
)
|
||
assert.equal(
|
||
dialedCalls - callsBefore6,
|
||
0,
|
||
'未知 host / 非 relay ⛔ 不许碰拨号池(`localPortFor` 会**按需绑池口**,是有副作用的调用)',
|
||
)
|
||
})
|
||
|
||
test('T38 P-2b:闭包必须把「订阅推送落点」传进判定(源码级守卫 —— 无实例时该分支不执行)', async () => {
|
||
const src = await readFile(new URL('../src/web/server.ts', import.meta.url), 'utf8')
|
||
const code = src
|
||
.split('\n')
|
||
.filter((l) => !/^\s*(\/\/|\*|\/\*)/.test(l))
|
||
.join('\n')
|
||
assert.ok(
|
||
code.includes('presenceLocalPort(name, ep.port)'),
|
||
'闭包必须把 `presenceLocalPort(name, ep.port)` 交给 `relayEndpointTarget`(否则三级链少一级 = P-2b 原样)',
|
||
)
|
||
assert.ok(
|
||
code.includes('pushedLocalPort:'),
|
||
'判定入参必须显式命名(⛔ 不许靠位置参数 / 事后补丁)',
|
||
)
|
||
})
|