/** * 覆盖网络 · 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 抓包定死,本测试逐帧复现) * ` > 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} 应在听`) })