Files
dsh_shenxian/test/relay.test.mjs
T
admin 452924d89c feat(config): 涉密内容外置到配置目录(档案 140)
把散落在代码里的真实部署值统一收进 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 未被误建)。
2026-09-19 15:12:19 +08:00

2012 lines
94 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
* 覆盖网络 · 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:'),
'判定入参必须显式命名(⛔ 不许靠位置参数 / 事后补丁)',
)
})