feat(cluster): 集群化落地 —— Manager/Worker 拆分 + 归属租约 + 跨机验证(T08)

背景:把平台从「单机单进程」改造成「1 组 Manager + N 台 Worker + 共享归属状态」,
硬约束 = 全程兼容单例模式(deployMode 默认 local;生产切换前 47 一行未动)。

主要改动
1) 数据模型 v7(SQLite 与 PG 两方言同步):新增 dsh_hosts 注册表 +
   dsh_instances.{host_id,epoch,heartbeat_at,lease_until};claimInstance 原子抢占
   (UPDATE … WHERE host_id IS NULL OR lease_until < now)+ pinInstanceHost 钉住归属。
2) 租约与 fencing:src/supervisor/lease.ts(acquire/renew/release + stillHolder 判据 +
   ttl > 2×renew 硬校验);心跳里续租,失权即向 worker 下发更高 epoch(self-fencing)。
   ⚠️ release 只清租约(lease_until),**保留 host_id** —— host_id 是「用户数据在哪台」的锚点。
3) Worker agent(src/worker/agent.ts,子命令 dshs worker):实例生命周期 + 文件面 /fs/*
   + 幂等键(operationId)+ 鉴权(timingSafeEqual);Worker 不写控制面数据
   (apiKey/uid 由 Manager 随 launch 投递,R5 收窄)。
4) 远端 Spawner + LeasedSpawner:按 host 路由(**粘性优先**:有历史归属且那台 up 就留在原地,
   否则按容量选最空的)+ 容量准入 + deployMode=cluster 装配(systemd drop-in,可回滚)。
5) bwrap 修正:**所有挂载点的中间目录统一前置 + 去重 + 由外到内**(「就近创建」会在嵌套前缀下
   遮掉已绑挂载点 ⇒ bwrap: Can't chdir);且**只能用 --tmpfs**,用 --perms 会让 47 的
   bwrap 0.4.0 直接拒启动(沙箱全挂)。
6) 跨机隧道 src/worker/tunnel.ts:SSH ControlMaster + 动态 -R 转发;**自愈由 agent 本地
   20s 定时器驱动**(不能只放 /healthz —— 心跳本身经隧道进来,断了就没人触发它)。
7) 文件面按归属路由(RemoteUserFs):实例与文件必须落在同一台机器,否则实例看不到自己的文件。
8) 观测面:dshs doctor / dshs cluster status。

验证(本次均已实跑)
- test/lease.test.mjs:SQLite 10/10 == PG 10/10
- 组件级端到端 5 个:verify-cluster-{agent,lease,fs,migrate,live}.mjs
- 真跨机(47 Manager / 106 Worker,跨云 + 反向隧道)verify-cluster-cross.mjs 九步全绿
- 域名形态访问 verify-cluster-domain.mjs(<user>.域名 → Manager → 远端实例;越权 403)
- 冒烟 scripts/smoke-*:6/8,失败项与改动前基线完全相同(无回归)
- 生产切换与回滚剧本见 dsh-server-docs/交接单/T08-集群化落地-兼容单例模式.md §16
This commit is contained in:
admin committed 2026-09-15 18:47:02 +08:00
1 parent 68c0a320ed
commit c70d5d860e
47 files changed
+5367 -14

No files matched your search

+450
View File
@@ -0,0 +1,450 @@
/**
* Worker agent(T08 S3;设计 §11.2)。
*
* **它是什么**:Worker 上唯一的"被拨入口" —— 一个内部 HTTP 服务,把实例生命周期
* 暴露给 Manager。**它不做归属决策**(谁托管谁是 Manager + PG 的事),只负责
* "在这台机器上把实例起停好",并复用 **`LocalSpawner`**,因此 bwrap/uid/scope
* 隔离、内存配额推导、崩溃退避与熔断、插件探活这些**本地语义全部原样保留**
* (这是本方案相对 k8s 路线最大的成本优势)。
*
* 四条协议纪律(设计 §11.3):
* 1. **单向拨入**:Worker 不反向连 Manager、不写控制面数据(自己可以有库,见下);
* 2. **幂等键**:每个变更请求带 `operationId`,重复请求**回放上次结果**
* (否则 Manager 超时重试会起两个实例);
* 3. **最小接口**:只接受白名单动作,参数受限(folder 由 Manager 解析、patch 有长度上限)
* —— agent 若能被当任意命令执行器,Worker 沦陷 = 全集群沦陷;
* 4. **self-fencing**:`POST /fence {userId, epoch}` —— 本地记录的 epoch 落后于
* Manager 下发的值 ⇒ **主动停掉该实例**(防双写的最后一道防线)。
*
* @module dshs/worker/agent
*/
import { timingSafeEqual } from 'node:crypto'
import Fastify, { type FastifyInstance, type FastifyReply } from 'fastify'
import type { ServerConfig } from '../config.js'
import { LocalUserFs } from '../fs/local-user-fs.js'
import { userRoot } from '../fs/workspace.js'
import { isUserFsErrorCode, UserFsError } from '../fs/user-fs.js'
import { hashUid } from '../isolation.js'
import { LocalSpawner } from '../supervisor/orchestrator.js'
import { SshTunnel } from './tunnel.js'
import type { Instance } from '../supervisor/spawner.js'
/** 绑定的头部名(Manager/agent 双方约定)。 */
export const AGENT_TOKEN_HEADER = 'x-dsh-agent-token'
export interface WorkerAgentOptions {
/** 本机在 `dsh_hosts.id` 里的标识。 */
hostId: string
/** 共享密钥(仅内网 + nft 白名单;本版是 bearer 式比较,HMAC/防重放留待后续)。 */
token: string
/** 监听端口。 */
port: number
/** 绑定地址(默认 `0.0.0.0`,靠 nft 只放行 Manager 网段)。 */
host?: string
/** 返回给 Manager 做代理的地址(同机 1a 用 `127.0.0.1`;跨机时填内网 IP)。 */
instanceHost?: string
/** 日志级别。 */
logLevel?: string
/**
* **反向隧道**(跨机演练):Worker 主动拨 Manager,形如 `[email protected]:32022`。
* 不设则完全关闭(同机/单机形态零影响)。见 `tunnel.ts` 头注释。
*/
tunnelTarget?: string
/** 隧道私钥(默认 `~/.ssh/tunnel_ed25519`)。 */
tunnelIdentity?: string
/** ControlMaster socket(默认 `/tmp/dshs-tunnel-<hostId>.sock`)。 */
tunnelControlPath?: string
}
/** 变更类请求的幂等缓存条数上限(超出后丢最旧的 —— 只是省重试,不是审计)。 */
const OP_CACHE_MAX = 512
/** patch 内容长度上限(防把 agent 当大对象存储)。 */
const MAX_PATCH_BYTES = 256 * 1024
interface OpCache {
order: string[]
results: Map<string, unknown>
}
/**
* 组装 agent。**复用 `LocalSpawner`**(隔离/配额/退避/熔断/探活全部原样保留)。
*
* ⚠️ **边界要读准**(2026-09-15 用户纠正):本 agent **不写控制面数据**(尤其归属/租约 ——
* 双写就是脑裂),所以 `apiKey` 与 `uid` **不由本机查控制面库**,而是 Manager 在
* `POST /launch` 时随请求投递(与 k8s 用 per-user Secret 同一思路),只存内存。
* 但这**不等于"Worker 不许有数据库"**:插件的 per-user 数据(如 `home/.dsh/mcn-plugin.db`)
* 属于**实例业务数据**,由实例自己读写、跟着 home 走;Worker 也可以有自己的运维库。
* 完整判据见设计 §1.3「数据分层」。
*/
export function buildWorkerAgent(
config: ServerConfig,
options: WorkerAgentOptions,
): { app: FastifyInstance; spawner: LocalSpawner; stop: () => Promise<void> } {
/** launch 时投递、仅存内存的凭据与 uid(Worker 不连 DB)。 */
const apiKeys = new Map<string, string>()
const uids = new Map<string, number>()
const spawner = new LocalSpawner(
config,
async (userId: string) => apiKeys.get(userId) ?? null,
async (userId: string) => uids.get(userId) ?? hashUid(userId, config.baseUid),
)
/**
* 文件面(T08 S5):**复用同一个 `LocalUserFs`** —— 用户卷本来就在本机,
* 所以"跨机文件面"= 把这个实现经 HTTP 暴露出去,而不是重新实现一套路径语义。
*/
const userFs = new LocalUserFs((userId: string) => userRoot(config.dataRoot, userId))
/**
* 反向隧道(可选)。静态转发 = **agent 自身端口** + `DSHS_TUNNEL_STATIC_PORTS`(如控制面 PG);
* 实例端口在 launch/stop 时动态加减,并在 `/healthz`(Manager 的心跳)里**对账自愈**。
*/
const tunnelTarget = options.tunnelTarget ?? process.env.DSHS_TUNNEL_TARGET ?? ''
const staticPorts = [
options.port,
...(process.env.DSHS_TUNNEL_STATIC_PORTS ?? '')
.split(',')
.map((v) => Number(v.trim()))
.filter((v) => Number.isInteger(v) && v > 0),
]
const tunnel =
tunnelTarget === ''
? undefined
: new SshTunnel({
target: tunnelTarget,
identity:
options.tunnelIdentity ??
process.env.DSHS_TUNNEL_IDENTITY ??
`${process.env.HOME ?? '/root'}/.ssh/tunnel_ed25519`,
controlPath: options.tunnelControlPath ?? `/tmp/dshs-tunnel-${options.hostId}.sock`,
staticPorts,
})
let tunnelReady = tunnel === undefined
/**
* **隧道自愈**:master 失联就重建(重建会自动补回 staticPorts 的静态转发)。
*
* ⚠️ 为什么不能只在 `/healthz` 里做(2026-09-15 想清楚的一个死角):`/healthz` 是**经隧道**
* 才打得进来的 —— 隧道一断,Manager 的心跳就进不来,自愈**永远不会被触发**(自己把自己锁死)。
* ⇒ 必须由 **agent 本地定时器**驱动(下面 20s 一跳),`/healthz` 里再顺手做一次。
*/
const healTunnel = async (): Promise<void> => {
if (tunnel === undefined) return
if (tunnelReady && (await tunnel.isMasterAlive())) return
tunnelReady = false
try {
await tunnel.ensureMaster()
tunnelReady = true
console.error('[tunnel] master 失联 → 已重建(含静态转发)')
} catch (err) {
console.error('[tunnel] 重建失败,下轮再试:', err instanceof Error ? err.message : err)
}
}
/** 把活着的实例端口补齐、把已消失的撤掉(崩溃退出也走这里收敛,不必逐个挂 exit 钩子)。 */
const reconcileTunnel = async (): Promise<void> => {
if (tunnel === undefined || !tunnelReady) return
const live = new Set((await spawner.listUserInstances()).map((i) => i.port).filter((p): p is number => p !== undefined))
for (const port of live) await tunnel.forward(port)
for (const port of tunnel.ports) {
if (!live.has(port) && !staticPorts.includes(port)) await tunnel.cancel(port)
}
}
let tunnelTimer: NodeJS.Timeout | undefined
if (tunnel !== undefined) {
void tunnel
.ensureMaster()
.then(() => {
tunnelReady = true
})
.catch((err: unknown) => {
console.error('[tunnel] 建立失败(跨机代理将不可用,本机功能不受影响):', err instanceof Error ? err.message : err)
})
tunnelTimer = setInterval(() => {
void healTunnel().then(reconcileTunnel)
}, 20_000)
tunnelTimer.unref?.()
}
const app = Fastify({ logger: { level: options.logLevel ?? 'info' }, bodyLimit: MAX_PATCH_BYTES + 4096 })
const cache: OpCache = { order: [], results: new Map() }
/** agent 侧记住的 epoch(self-fencing 判据)。 */
const epochs = new Map<string, number>()
const remember = (op: string, value: unknown): void => {
if (cache.results.has(op)) return
cache.results.set(op, value)
cache.order.push(op)
while (cache.order.length > OP_CACHE_MAX) {
const oldest = cache.order.shift()
if (oldest !== undefined) cache.results.delete(oldest)
}
}
app.addHook('onRequest', async (request, reply) => {
if (request.url === '/healthz') return // 存活探测不带凭据
const given = request.headers[AGENT_TOKEN_HEADER]
const expected = options.token
const a = Buffer.from(typeof given === 'string' ? given : '')
const b = Buffer.from(expected)
if (a.length !== b.length || !timingSafeEqual(a, b)) {
await reply.code(401).send({ error: 'unauthorized' })
}
})
app.get('/healthz', async () => {
const instances = await spawner.listUserInstances()
// 心跳里顺手自愈 + 对账(主驱动是本地定时器,见 healTunnel 的注释)
await healTunnel()
await reconcileTunnel()
return {
ok: true,
hostId: options.hostId,
instances: instances.length,
tunnel: tunnel === undefined ? null : { ready: tunnelReady, ports: tunnel.ports },
}
})
/** 对账用:**一次拿回整机**(设计 §11.6,替代逐用户查询)。 */
app.get('/instances', async () => ({ instances: await spawner.listUserInstances() }))
app.post('/launch', async (request, reply) => {
const body = request.body as {
userId?: string
folder?: string
patch?: string
epoch?: number
operationId?: string
apiKey?: string | null
uid?: number
}
if (body.userId === undefined || body.operationId === undefined) {
return reply.code(400).send({ error: 'userId and operationId are required' })
}
const cached = cache.results.get(body.operationId)
if (cached !== undefined) return cached // 幂等回放
// 先落凭据/uid/**epoch**(重放路径也安全:同值覆盖)。
// ⚠️ epoch 记录的是「Manager 的意图」,因此必须在**尝试 spawn 之前**落 ——
// "实例本来就在跑"(AlreadyRunningError 分支)时也要记,否则 `/fence` 拿不到
// 我的 epoch,self-fencing 就永远不触发(2026-09-15 T08 S3 实测踩到)。
if (body.apiKey !== undefined && body.apiKey !== null) apiKeys.set(body.userId, body.apiKey)
if (body.uid !== undefined) uids.set(body.userId, body.uid)
if (body.epoch !== undefined) epochs.set(body.userId, body.epoch)
try {
const instance = await spawner.launch(body.userId, body.folder ?? '', body.patch)
// 跨机:把该实例端口经隧道打到 Manager 侧(失败不阻断 —— 本机仍可用)
if (tunnel !== undefined && tunnelReady && instance.port !== undefined) await tunnel.forward(instance.port)
const payload = { instance: { ...instance, launchToken: spawner.launchTokenOf(body.userId) } }
remember(body.operationId, payload)
return payload
} catch (err) {
// 已在跑:**返回现有实例**而不是报错 —— 这让重试天然安全(与 AlreadyRunningError 语义对齐)。
const msg = err instanceof Error ? err.message : String(err)
if (/already has a running/i.test(msg)) {
const currents = await spawner.listUserInstances()
const found = currents.find((i) => i.userId === body.userId)
if (found !== undefined) {
const payload = { instance: { ...found, launchToken: spawner.launchTokenOf(body.userId) }, note: 'already-running' }
remember(body.operationId, payload)
return payload
}
}
return reply.code(500).send({ error: msg })
}
})
app.post('/stop', async (request, reply) => {
const body = request.body as { userId?: string; operationId?: string }
if (body.userId === undefined || body.operationId === undefined) {
return reply.code(400).send({ error: 'userId and operationId are required' })
}
const cached = cache.results.get(body.operationId)
if (cached !== undefined) return cached
const before = await spawner.status(body.userId)
await spawner.stop(body.userId)
if (tunnel !== undefined && tunnelReady && before.main?.port !== undefined) await tunnel.cancel(before.main.port)
epochs.delete(body.userId)
apiKeys.delete(body.userId) // 凭据只该活在实例生命周期内
const payload = { ok: true }
remember(body.operationId, payload)
return payload
})
app.get('/status/:userId', async (request, reply) => {
const { userId } = request.params as { userId: string }
const status = await spawner.status(userId)
return reply.send({ userId, main: status.main ?? null })
})
/** 代理目标(Manager 用)。未运行时返回 `{ running: false }`。 */
app.get('/endpoint/:userId', async (request) => {
const { userId } = request.params as { userId: string }
const endpoint = await spawner.endpointFor(userId)
return endpoint === undefined
? { running: false }
: { running: true, host: options.instanceHost ?? '127.0.0.1', port: endpoint.port }
})
/** self-fencing:我持有的 epoch 落后于 Manager 下发的值 ⇒ **自杀**。 */
app.post('/fence', async (request, reply) => {
const body = request.body as { userId?: string; epoch?: number }
if (body.userId === undefined || body.epoch === undefined) {
return reply.code(400).send({ error: 'userId and epoch are required' })
}
const mine = epochs.get(body.userId)
if (mine === undefined || mine >= body.epoch) return { fenced: false, mine: mine ?? null }
await spawner.stop(body.userId)
epochs.delete(body.userId)
return { fenced: true, mine }
})
app.post('/restart-probe/:userId', async (request) => {
const { userId } = request.params as { userId: string }
return spawner.restartAndProbe(userId)
})
/**
* 活动信号转发(Manager 代理到用户流量时调用)。
* 为什么要转发:idle-reap 是**本地语义**(`LocalSpawner` 的 `lastActive` + TTL/LRU),
* 不转发的话 worker 会以为实例一直没人用、把它回收掉(档案 08)。
*/
app.post('/touch/:userId', async (request) => {
const { userId } = request.params as { userId: string }
spawner.touch(userId)
return { ok: true }
})
app.post('/watchdog/:userId', async (request) => {
const { userId } = request.params as { userId: string }
return { instance: (await spawner.spawnWatchdog(userId)) ?? null }
})
// ── 文件面(T08 S5;供 Manager 的 RemoteUserFs 调用)──────────────────────
// 请求体/响应都是**工作区相对路径 + base64**,与 `UserFs` 的语义一一对应;
// 失败时回 `{error: code}` 并把 `UserFsError.code` 映射成对应 HTTP 状态
// —— 这正是 `user-fs.ts` 里那个 seam 设计的用法(路由按 code 回给前端)。
/** 统一包装:把 `UserFsError` 还原成 wire 形态(其余错误 → 500)。 */
const fsCall = async <T>(reply: FastifyReply, fn: () => Promise<T>): Promise<T | undefined> => {
try {
return await fn()
} catch (err) {
if (err instanceof UserFsError) {
await reply.code(err.status).send({ error: err.code })
return undefined
}
const msg = err instanceof Error ? err.message : String(err)
await reply.code(500).send({ error: 'internal', detail: msg.slice(0, 200) })
return undefined
}
}
app.post('/fs/init', async (request, reply) => {
const body = request.body as { userId?: string; uid?: number }
if (body.userId === undefined) return reply.code(400).send({ error: 'userId is required' })
return (await fsCall(reply, async () => {
await userFs.initUserRoot(body.userId as string, body.uid)
return { ok: true }
})) ?? reply
})
app.post('/fs/list', async (request, reply) => {
const body = request.body as { userId?: string; relPath?: string }
if (body.userId === undefined) return reply.code(400).send({ error: 'userId is required' })
return (await fsCall(reply, () => userFs.listDir(body.userId as string, body.relPath ?? ''))) ?? reply
})
app.post('/fs/mkdir', async (request, reply) => {
const body = request.body as { userId?: string; relPath?: string }
if (body.userId === undefined) return reply.code(400).send({ error: 'userId is required' })
return (await fsCall(reply, async () => {
await userFs.mkdir(body.userId as string, body.relPath ?? '')
return { ok: true }
})) ?? reply
})
app.post('/fs/create', async (request, reply) => {
const body = request.body as { userId?: string; relPath?: string; name?: string; type?: 'file' | 'dir' }
if (body.userId === undefined || body.name === undefined) {
return reply.code(400).send({ error: 'userId and name are required' })
}
// ⚠️ 必须包成对象:`createEntry` 返回的是字符串(净化后的文件名),直接 return 会被
// Fastify 当 text/plain 发出,而调用方(RemoteUserFs)按 JSON 解析 ⇒ 静默 500。
const created = await fsCall(reply, () =>
userFs.createEntry(body.userId as string, body.relPath ?? '', body.name as string, body.type ?? 'file'),
)
return created === undefined ? reply : { name: created }
})
app.post('/fs/upload', async (request, reply) => {
const body = request.body as { userId?: string; relPath?: string; name?: string; dataBase64?: string }
if (body.userId === undefined || body.name === undefined || body.dataBase64 === undefined) {
return reply.code(400).send({ error: 'userId, name and dataBase64 are required' })
}
// 同上:`upload` 返回的是净化后的文件名,必须包成对象。
const uploaded = await fsCall(reply, () =>
userFs.upload(
body.userId as string,
body.relPath ?? '',
body.name as string,
Buffer.from(body.dataBase64 as string, 'base64'),
),
)
return uploaded === undefined ? reply : { name: uploaded }
})
app.post('/fs/read', async (request, reply) => {
const body = request.body as { userId?: string; relPath?: string; maxBytes?: number }
if (body.userId === undefined || body.relPath === undefined) {
return reply.code(400).send({ error: 'userId and relPath are required' })
}
const out = await fsCall(reply, () => userFs.readFile(body.userId as string, body.relPath as string, body.maxBytes))
if (out === undefined) return reply
return { name: out.name, dataBase64: out.data.toString('base64') }
})
app.post('/fs/isdir', async (request, reply) => {
const body = request.body as { userId?: string; relPath?: string }
if (body.userId === undefined) return reply.code(400).send({ error: 'userId is required' })
const out = await fsCall(reply, () => userFs.isDirectory(body.userId as string, body.relPath ?? ''))
return out === undefined ? reply : { isDirectory: out }
})
app.post('/fs/plugins', async (request, reply) => {
const body = request.body as { userId?: string }
if (body.userId === undefined) return reply.code(400).send({ error: 'userId is required' })
return (await fsCall(reply, () => userFs.listInstalledPlugins(body.userId as string))) ?? reply
})
app.post('/fs/handoff', async (request, reply) => {
const body = request.body as { userId?: string; content?: string }
if (body.userId === undefined || body.content === undefined) {
return reply.code(400).send({ error: 'userId and content are required' })
}
return (await fsCall(reply, async () => {
await userFs.writeHandoff(body.userId as string, body.content as string)
return { ok: true }
})) ?? reply
})
/** 本机 dataRoot(Manager 的 RemoteUserFs 用它做 `resolvePath` 的路径数学)。 */
app.get('/fs/root', async () => ({ dataRoot: config.dataRoot }))
app.get('/', async () => ({ agent: 'dshs-worker', hostId: options.hostId }))
return {
app,
spawner,
/**
* 停机:**先收实例、再关 HTTP**。
* 为什么必须收:实例是 worker 自己的子进程,停机不收就变孤儿(占用端口与内存);
* 而且孤儿会继承 stdout ⇒ 调用方的管道永不关闭(2026-09-15 实测:verify 脚本挂死)。
* 归属与租约由 Manager 侧处理(worker 不写控制面数据),所以这里只停进程。
*/
stop: async (): Promise<void> => {
if (tunnelTimer !== undefined) clearInterval(tunnelTimer)
await app.close()
await spawner.teardown()
await tunnel?.close()
},
}
}
+193
View File
@@ -0,0 +1,193 @@
/**
* Worker 侧的**反向隧道管理器**(T08 跨机演练)。
*
* 为什么需要它:Manager 要连 Worker 上的两样东西 —— **agent 端口**与**每个实例的端口**
* (实例只监听 `127.0.0.1`,这是 portGuard 的设计前提)。而 Worker 公网入方向通常被
* **云安全组**挡住(实测:106 的 19100 从 47 与本机都连不上),放通只能在控制台点。
*
* 绕法:**让 Worker 主动拨 Manager**,用 SSH 反向转发把两边的 `127.0.0.1:<port>` 接起来。
* 好处(实测):
* · **两端都不用新开端口** —— 只用已开放的 SSH 端口(47 是 32022);
* · 链路是加密的,且 Manager 侧落在 loopback(`GatewayPorts no` 默认)⇒ 不对外暴露;
* · 实例端口是**动态**的(`findFreePort()`)⇒ 用 **ControlMaster + `ssh -O forward/cancel`**
* 在**同一条长连接**上加/减转发,不必为每个端口重开连接。
*
* ⚠️ 定位:这是**演练级**传输(生产长期方案见设计 §2.3:受控网段白名单或隧道服务)。
* ⚠️ 默认**关闭**:只有设了 `DSHS_TUNNEL_TARGET` 才启用 ⇒ 对同机/单机形态零影响。
*
* @module dshs/worker/tunnel
*/
import { execFile } from 'node:child_process'
import { existsSync, unlinkSync } from 'node:fs'
import { promisify } from 'node:util'
const run = promisify(execFile)
export interface TunnelOptions {
/** 拨入目标,形如 `[email protected]:32022`。 */
target: string
/** 私钥路径(建议专用、且在 Manager 侧用 `restrict,port-forwarding` 限权)。 */
identity: string
/** ControlMaster socket 路径(同一路径复用同一条连接)。 */
controlPath: string
/** 启动时就转发的端口(agent 自身;还可带控制面 PG 等)。 */
staticPorts?: number[]
/** `ssh` 可执行文件路径。 */
sshBin?: string
}
export class SshTunnel {
/** 内部一律用**已补默认值**的具体类型(否则 `sshBin` 会是 `string | undefined`)。 */
private readonly opts: {
target: string
identity: string
controlPath: string
staticPorts: number[]
sshBin: string
}
private readonly forwarded = new Set<number>()
private readonly hostPart: string
private readonly portPart: number | undefined
constructor(options: TunnelOptions) {
// `user@host:port` 里的 port 是 **SSH 端口**(不是转发的端口)—— 47 上用 32022,
// 必须经 `-p` 传,否则会去连 22 而失败。
const [hostPart, portPart] = options.target.split(':')
this.hostPart = hostPart
this.portPart = portPart === undefined ? undefined : Number(portPart)
this.opts = {
target: options.target,
identity: options.identity,
controlPath: options.controlPath,
staticPorts: options.staticPorts ?? [],
sshBin: options.sshBin ?? '/usr/bin/ssh',
}
}
/** 所有 ssh 调用的公共参数(`-p` 只在目标里显式给了端口时才加)。 */
private baseArgs(): string[] {
return this.portPart === undefined ? [] : ['-p', String(this.portPart)]
}
/** 当前已转发的端口(诊断用)。 */
get ports(): number[] {
return [...this.forwarded]
}
/**
* 建立(或复用)ControlMaster 长连接,并把 `staticPorts` 转发上去。
* 幂等:socket 已存在且 master 还活着就直接返回。
*/
async ensureMaster(): Promise<void> {
if (existsSync(this.opts.controlPath)) {
try {
await run(this.opts.sshBin, [...this.baseArgs(), '-S', this.opts.controlPath, '-O', 'check', this.hostPart])
// master 活着 ⇒ 只需补齐静态转发
for (const port of this.opts.staticPorts ?? []) await this.forward(port)
return
} catch {
try {
unlinkSync(this.opts.controlPath) // 僵尸 socket:清掉重建
} catch {
/* 无所谓 */
}
}
}
const args = [
'-M',
'-N',
'-f',
...this.baseArgs(),
'-S',
this.opts.controlPath,
'-i',
this.opts.identity,
'-o',
'BatchMode=yes',
'-o',
'StrictHostKeyChecking=accept-new',
'-o',
'ExitOnForwardFailure=yes',
'-o',
'ServerAliveInterval=15',
'-o',
'ServerAliveCountMax=4',
]
for (const port of this.opts.staticPorts ?? []) args.push('-R', `${port}:127.0.0.1:${port}`)
args.push(this.hostPart)
await run(this.opts.sshBin, args, { timeout: 20_000 })
for (const port of this.opts.staticPorts ?? []) this.forwarded.add(port)
}
/**
* master 是否还活着(`ssh -O check`)。
*
* 为什么需要:**对端 SSH 重启/断链后,转发会全部消失,而本地 `forwarded` 集合并不知情**
* ⇒ 若只看本地状态,会以为"隧道还好",实际 Manager 已经连不上这台 Worker
* (2026-09-15 收口"以跑通为目的"时补)。
*/
async isMasterAlive(): Promise<boolean> {
try {
await run(this.opts.sshBin, [...this.baseArgs(), '-S', this.opts.controlPath, '-O', 'check', this.hostPart], {
timeout: 8_000,
})
return true
} catch {
// check 失败 ⇒ master 不在了;顺手清掉本地记账,避免"以为还转着"
this.forwarded.clear()
return false
}
}
/**
* 动态加一条反向转发(实例起来时调用)。
* 用**同一个端口号**:实例在 Worker 上是 `127.0.0.1:<port>`,反向转发落到 Manager 的
* `127.0.0.1:<port>` ⇒ Manager 侧无需端口映射表,`endpointFor` 直接回 `127.0.0.1`。
*/
async forward(port: number): Promise<boolean> {
if (this.forwarded.has(port)) return true
try {
await run(
this.opts.sshBin,
[...this.baseArgs(), '-S', this.opts.controlPath, '-O', 'forward', '-R', `${port}:127.0.0.1:${port}`, this.hostPart],
{ timeout: 10_000 },
)
this.forwarded.add(port)
return true
} catch {
return false // 失败不抛:实例本身仍在本机可用,只是跨机代理这跳不可用
}
}
/** 撤销一条转发(实例停止/退出时调用)。 */
async cancel(port: number): Promise<void> {
if (!this.forwarded.has(port)) return
try {
await run(
this.opts.sshBin,
[...this.baseArgs(), '-S', this.opts.controlPath, '-O', 'cancel', '-R', `${port}:127.0.0.1:${port}`, this.hostPart],
{ timeout: 10_000 },
)
} catch {
/* 连接已断也一样算撤销 */
}
this.forwarded.delete(port)
}
/** 关闭 master(进程退出时)。 */
async close(): Promise<void> {
try {
await run(this.opts.sshBin, [...this.baseArgs(), '-S', this.opts.controlPath, '-O', 'exit', this.hostPart], {
timeout: 10_000,
})
} catch {
/* 已退出 */
}
this.forwarded.clear()
}
/** 目标 SSH 端口(`root@h:32022` → 32022)。 */
get targetPort(): number | undefined {
return this.portPart
}
}