Files
dsh_shenxian/test/relay.test.mjs
T
admin 146c3d25ef feat(overlay): 覆盖网络线序①–⑮ 代码与测试产物入库
覆盖网络线累积产物(此前只在工作区、未入版本库):
- 新增 relay 子系统 src/net/relay/**(wire/duplex/server/client/dialer/switcher/directory/identity/keys/placement/network/addr-override/main/index)
- 新增 src/worker/relay-tunnel.ts、src/web/routes/overlay.ts
- 新增观测/演练脚本 overlay-probe、overlay-failover-drill、overlay-keyring、overlay-holepunch、overlay-jitter、overlay-wan、overlay-relaykey-add、relay-mem-calibrate
- 新增测试 12 个(relay / relay-failover / remote-spawner / instance-port / overlay-{network,auth,bootstrap,identity} / remote-user-fs 等)

验收基线:npm test = 162 pass / 0 fail / 1 skip;--scene all = 12 PASS / 0 SKIP / 0 FAIL;overlay-probe = 12/12
2026-09-17 17:03:27 +08:00

1319 lines
58 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 { createServer as createTcpServer, connect } from 'node:net'
import { test } from 'node:test'
import { MUX, OPS_NETWORK, RelayClient, RelayDialer, RelayServer, chooseNode, decodeMux, encodeJsonFrame, encodeMux, logicalName, parseKeysInline } 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-106)上,**同一条 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} 应在听`)
})