diff --git a/README.md b/README.md index 1c098ff..359acd3 100644 --- a/README.md +++ b/README.md @@ -123,6 +123,7 @@ node lib/cli.js --port 3080 --db ./dev.local.db | [K8s 踩坑记录](docs/k8s-deploy.md) | 模式 B 实机部署的坑与根因(镜像源 / 存储 / 网络 / 备案) | 模式 B 部署卡住时查 | | [技术蓝图](docs/blueprint.md) | 权威技术设计:拓扑 / 双 DSH / 数据模型 / API / 安全模型 | 想了解原理或参与开发 | | [PoC 实验记录](poc/README.md) | k8s 关键假设的四项实测(RWX / socat / NetworkPolicy / CNPG) | 评估模式 B 可行性时 | +| [代码分层与依赖方向](docs/architecture.md) | 四层划分 + 依赖方向规则 + 三条工程纪律(契约前置 / 模块头注释 = 索引) | **动手改 `src/` 前读一遍**,尤其新增模块或目录时 | > 学习路线:**模式 A** 先读「部署教程」跑通 → 再看「硬隔离」加固,出问题查「常见问题排查」;**模式 B** 直接从「K8s 部署教程」开始,卡住查「K8s 踩坑记录」。 diff --git a/docs/architecture.md b/docs/architecture.md new file mode 100644 index 0000000..e0d966f --- /dev/null +++ b/docs/architecture.md @@ -0,0 +1,101 @@ +# 代码分层与依赖方向 + +> **目的**:让"改一个模块"只需要读 **索引 + 相关模块**,而不是读全仓。 +> **背景与实测数据**:`项目代码_分层范式与迭代风险评估_20260916.md`(含风险分档与"改一处读全仓"的量化评估)。 +> **状态**:2026-09-16 立。当前**只是把规则写下来**——本仓**不需要重构**,需要的是别在扩张期把方向走丢。 + +--- + +## 1. 四层与依赖方向 + +``` +① 入口层 src/cli.ts · src/web/server.ts · src/web/routes/* · src/web/middleware/* +② 领域层 src/domain/* ← 【新增,当前为空】 +③ 能力层 src/supervisor · db · fs · worker · net · nginx +④ 基础层 src/config.ts · crypto.ts · isolation.ts · index.ts +``` + +### 依赖方向:`① → ② → ③ → ④` **单向,⛔ 不得回指** + +| # | 规则 | 判据(怎么算违反) | +|---|---|---| +| **R1** | 上层可依赖下层;**下层不得 import 上层** | 在 `④` 的文件里出现 `from '../web/...'` 之类的**子目录 import** = 违反 | +| **R2** | 同层之间允许,但**新增同层依赖前先问一句"它到底属哪一层"** | 说不清的依赖,通常意味着**该有一条领域层接口** | +| **R3** | `net/*` 是**能力层的一员**(与 `supervisor`/`db`/`fs` 同级),⛔ **不是跨层的特权模块** | `net/*` 不得 import `web/*`;`web` 可以依赖 `net` | +| **R4** | `src/{config,crypto,isolation,index}.ts` **零业务依赖** | 这四个文件出现任何子目录 import = 违反(当前**成立**,见 §3) | + +--- + +## 2. 三条工程纪律 + +**2.1 契约前置** —— 新增模块**先出接口与类型文件**,再写实现。 +> 范例:`src/net/reachability.ts` + `src/net/rendezvous.ts`(覆盖网络 S0)。 +> 做法 = **纯新增文件 + 零行为变化**,验收靠 `test/reachability.test.mjs` 把**现网真实数据写死为判据**(不靠肉眼)。 + +**2.2 模块头注释 = 索引** —— 每个模块目录的入口文件,头注释必须写明三件事: +> **职责 / 依赖谁 / 被谁依赖**。 +> 本仓现有头注释质量已经很高(是"读索引就能定位"的基础),缺的只是把这三项**变成必须写的字段**。 +> 效果:把"读全仓"降级为"读索引 + 读相关模块"。 + +**2.3 动手前的自查** —— 问两句: +> ① 这次改动涉及哪几层?② 有没有从下往上(下层 import 上层)的引用? + +--- + +## 3. 现状实测(2026-09-16) + +**规模**:`src` **58 个 .ts / 14,053 行**。`web` 占 40%(5,593 行)⇒ 入口层事实上承担了业务层职责。 + +**🔴 一处勘误**(早期统计有误,以此处为准): +> 曾有结论称「有 **3 处基础层反向依赖**(`(root)→web` 2 + `(root)→db` 1),趁早清掉」。 +> **实测不成立** —— 它把 `src/cli.ts` 误算进了"基础层"。`cli.ts` 是**入口层①**,它 import `db`/`fs`/`web` 属 **① → ③ 的合法方向**。 +> 证据:`grep -n "from '\./\(web\|db\|fs\|supervisor\|worker\|net\)/" src/*.ts` ⇒ **只命中 `cli.ts`**;`config.ts` / `crypto.ts` / `isolation.ts` / `index.ts` **零 import**。 +> ⇒ **没有可清的反向依赖;R4 当前已经成立。** + +**真正的成本不在反向依赖,在"边界模糊"**: +`web/routes/business-plugins.ts`(743 行)· `skills.ts`(687 行)把业务规则写在路由里,规则的另一半在 `supervisor/orchestrator.ts`(1,259 行)与 `db/repo.ts`(752 行) +⇒ **改一条规则要连读 3 个大文件 ≈ 2,500 行**。这才是该被 §4 第 1 条解决的。 + +--- + +## 4. 未做(按优先级) + +| 优先级 | 事项 | 为什么排这个位置 | +|---|---|---| +| **P0(本文件)** | 把四层规则写下来 + 在 `README.md` 文档表登记 | 零代码风险;**扩张期一旦走丢,回头改的成本高得多** | +| **P1** | 补 `src/domain/*`:把 `web/routes` 的业务规则抽出来 | ⚠️ **行为敏感重构**(不是纯类型);影响面 **>10 文件** ⇒ **动工前先出清单** | +| **P2** | 按 §2.2 给各模块入口补「职责 / 依赖谁 / 被谁依赖」 | 批量改 >10 文件,同样先出清单 | + +> **顺序建议**:P0 ✅ → 覆盖网络 S1–S4 → P1。 +> P1 **不阻塞**覆盖网络;反过来,覆盖网络的新模块(`net/*`)正好是 §2.1「契约前置」的第一个样板。 + +--- + +## 5. 可执行判据(✅ 2026-09-16 补 —— 规则从此**不靠自觉**) + +```bash +npm run check:layering # = node scripts/check-layering.mjs +``` + +- `scripts/check-layering.mjs` 静态扫描 `src/**/*.ts`,按 §1 四层判 **R1 / R3 / R4**;未归类的目录会被顶出来(逼你定层)。 +- **ratchet(棘轮,不是豁免)**:存量违规登记在 `scripts/layering-baseline.json`,键 = 「源→目标」的**边**(不含行号 ⇒ 抗行号漂移)。 + · **新增**违规 ⇒ **退出码 1**(⛔ 不得引入);· 基线内**已消除**项 ⇒ 打印提醒,`--update-baseline` 收紧(**只许减**)。 +- 已挂进 **`npm run verify`**(`build` 之后、单测之前)。 + +### 5.1 首跑实测(2026-09-16 · 60 个 .ts) + +| 层 | ①入口 | ②领域 | ③能力 | ④基础 | 未归类 | +|---|---|---|---|---|---| +| 文件数 | 24 | **0** | 32 | 4 | 0 | + +**现存违规 5 条**,全部是 `③能力 → ①入口` 的反向依赖(= §3 那句「4 组双向依赖」的内核): + +| 源 | 目标 | 修法方向 | +|---|---|---| +| `supervisor/proxy.ts` | `web/auth.ts`(`hashSessionToken` / `parseCookie`) | 纯函数 ⇒ **下沉到基础层** | +| `supervisor/proxy.ts` | `web/middleware/authn.ts`(`requireAuth`) | 会话校验是**能力**不是入口 ⇒ 下沉到能力层 | +| `fs/local-user-fs.ts` · `fs/remote-user-fs.ts` · `fs/workspace.ts` | `web/middleware/fs-guard.ts` | 路径守卫 = **领域规则** ⇒ 归 `src/domain/`(或 `src/fs/`) | + +> ⚠️ **与 §3 勘误不矛盾**:§3 证伪的是「**基础层**(④)有 3 处反向依赖」⇒ 那部分确实 **0 处**;本节这 5 条是**能力层→入口层**(R1),是另一件事。 +> ⚠️ 这 5 条**未整改**(命中 R7:行为敏感重构且影响 >10 文件 ⇒ 先出清单)⇒ 先建基线,让"**不新增**"立刻生效。 +> 📌 **机制已兑现价值**:脚本首跑就抓出两件我肉眼漏掉的事 —— ① 完整的 R1 检查不是 0 违规(有 5 条);② 我自己漏定了 `src/nginx/` 的层归属。 diff --git a/package.json b/package.json index 900b869..f78463c 100644 --- a/package.json +++ b/package.json @@ -22,8 +22,9 @@ "prepare": "npm run build", "dev": "node lib/cli.js", "typecheck": "tsc -p tsconfig.json --noEmit", - "verify": "npm run build && node --test test/db.test.mjs test/local-user-fs.test.mjs test/crash-policy.test.mjs test/locale-pref.test.mjs test/worker-provision.test.mjs && node scripts/verify-inject.cjs lib/supervisor/proxy.js && node scripts/verify-static.mjs && node scripts/verify-platform-admin-section.mjs && node scripts/verify-mem-model.mjs && node scripts/verify-model-landing.mjs && node scripts/verify-dsh-install.mjs && node scripts/verify-models-dict.mjs && node scripts/verify-models-render.cjs && node scripts/verify-dsh-compat.mjs && node scripts/verify-my-skills.mjs && node scripts/verify-portal-entry.mjs", - "test": "npm run build && node --test test/db.test.mjs test/local-user-fs.test.mjs test/crash-policy.test.mjs && node scripts/verify-inject.cjs lib/supervisor/proxy.js", + "check:layering": "node scripts/check-layering.mjs", + "verify": "npm run build && node scripts/check-layering.mjs && node --test test/db.test.mjs test/local-user-fs.test.mjs test/crash-policy.test.mjs test/locale-pref.test.mjs test/worker-provision.test.mjs test/reachability.test.mjs && node scripts/verify-inject.cjs lib/supervisor/proxy.js && node scripts/verify-static.mjs && node scripts/verify-platform-admin-section.mjs && node scripts/verify-mem-model.mjs && node scripts/verify-model-landing.mjs && node scripts/verify-dsh-install.mjs && node scripts/verify-models-dict.mjs && node scripts/verify-models-render.cjs && node scripts/verify-dsh-compat.mjs && node scripts/verify-my-skills.mjs && node scripts/verify-portal-entry.mjs", + "test": "npm run build && node --test test/db.test.mjs test/local-user-fs.test.mjs test/crash-policy.test.mjs test/reachability.test.mjs && node scripts/verify-inject.cjs lib/supervisor/proxy.js", "smoke": "node scripts/smoke.mjs", "smoke:admin": "node scripts/smoke-admin.mjs", "smoke:auth": "node scripts/smoke-auth.mjs", diff --git a/scripts/check-layering.mjs b/scripts/check-layering.mjs new file mode 100644 index 0000000..11284e9 --- /dev/null +++ b/scripts/check-layering.mjs @@ -0,0 +1,128 @@ +#!/usr/bin/env node +/** + * 分层依赖方向检查 —— `docs/architecture.md` 的 R1–R4 **可执行判据**。 + * + * ## 为什么必须有它 + * 规则**只有文字**时靠自觉 ⇒ 扩张期一定会走丢。本脚本把四条规则变成一次静态扫描, + * "分层有没有走形"不再靠人判断(`npm run check:layering`)。 + * + * ## 判据(层号 ① 最上、④ 最下;**目标层号 < 源层号 ⇒ 违规**) + * | 规则 | 内容 | + * |---|---| + * | R1 | 上层可依赖下层,**下层不得 import 上层** | + * | R3 | `net/*` 属能力层,**不得 import `web/*`** | + * | R4 | `config/crypto/isolation/index` **零业务依赖**(任何子目录 import 都算违规) | + * + * ## 基线(ratchet,不是豁免) + * 现存违规登记在 `scripts/layering-baseline.json`(**键 = 源→目标 的边**,不含行号 ⇒ 抗行号漂移)。 + * · **新增**违规 ⇒ 退出码 **1**(规则从今天起生效) + * · 基线里**已消失**的项 ⇒ 打印提醒,可 `--update-baseline` 收紧 + * 这样"立规矩"不会因为存量债而把 `npm run verify` 卡红。 + * + * 用法:`node scripts/check-layering.mjs [--update-baseline] [--json]` + */ +import { readFileSync, readdirSync, statSync, writeFileSync, existsSync } from 'node:fs' +import { join, dirname, relative, sep } from 'node:path' + +const ROOT = process.cwd() +const SRC = join(ROOT, 'src') +const BASELINE = join(ROOT, 'scripts', 'layering-baseline.json') +const args = new Set(process.argv.slice(2)) +const posixify = (p) => p.split(sep).join('/') + +/** 层归属:**只按目录判**,不做语义推断(判据要能一眼看懂)。 */ +function layerOf(file) { + const p = posixify(file) + if (p === 'src/cli.ts' || p.startsWith('src/web/')) return 1 + if (p.startsWith('src/domain/')) return 2 + if (/^src\/(supervisor|db|fs|worker|net|nginx)\//.test(p)) return 3 + if (/^src\/(config|crypto|isolation|index)\.ts$/.test(p)) return 4 + return null // 未归类 ⇒ 新增目录必须先想清"它属于哪一层" +} +const NAME = { 1: '①入口', 2: '②领域', 3: '③能力', 4: '④基础' } + +function walk(dir) { + const out = [] + for (const e of readdirSync(dir)) { + const full = join(dir, e) + if (statSync(full).isDirectory()) out.push(...walk(full)) + else if (e.endsWith('.ts')) out.push(full) + } + return out +} + +/** 所有 `from '<相对>'` / `from "<相对>"`(含 `export … from`)。 */ +function importsOf(text) { + const out = [] + text.split('\n').forEach((line, i) => { + const re = /\bfrom\s+['"]([^'"]+)['"]/g + let m + while ((m = re.exec(line)) !== null) if (m[1].startsWith('.')) out.push({ spec: m[1], line: i + 1 }) + }) + return out +} + +// ⚠️ 统一成**相对仓根**的 posix 路径 —— 层判据是按 `src/...` 写的, +// 直接喂绝对路径会让每个文件都落进"未归类"(首版就踩了这个坑,静默全绿)。 +const files = walk(SRC).map((f) => posixify(relative(ROOT, f))) +const counts = { 1: 0, 2: 0, 3: 0, 4: 0 } +const unclassified = [] +const found = [] // { edge, detail } + +for (const file of files) { + const srcLayer = layerOf(file) + if (srcLayer === null) { + unclassified.push(file) + continue + } + counts[srcLayer] += 1 + for (const { spec, line } of importsOf(readFileSync(join(ROOT, file), 'utf8'))) { + const target = posixify(join(dirname(file), spec)).replace(/\.js$/, '.ts') + const tgtLayer = layerOf(target) + if (tgtLayer === null) continue + const edge = `${file} → ${target}` + const detail = `${file}:${line} → ${target}` + if (tgtLayer < srcLayer) { + const tag = srcLayer === 3 && file.includes('/net/') && target.startsWith('src/web/') ? 'R3' : 'R1' + found.push({ edge, detail: `${tag} ${detail} [${NAME[srcLayer]} → ${NAME[tgtLayer]}]` }) + } else if (srcLayer === 4) { + found.push({ edge, detail: `R4 ${detail} [④基础层必须零业务依赖]` }) + } + } +} + +// ── 基线比对 ──────────────────────────────────────────────────────────────── +const baseline = existsSync(BASELINE) ? new Set(JSON.parse(readFileSync(BASELINE, 'utf8')).edges ?? []) : null +const isNew = (v) => baseline === null || !baseline.has(v.edge) +const added = found.filter(isNew) +const known = found.filter((v) => !isNew(v)) +const fixed = baseline === null ? [] : [...baseline].filter((e) => !found.some((v) => v.edge === e)) + +if (args.has('--update-baseline')) { + writeFileSync(BASELINE, `${JSON.stringify({ note: 'R1–R4 存量违规基线(ratchet):只许减、不许增;新增违规会让 check:layering 失败', edges: [...new Set(found.map((v) => v.edge))].sort() }, null, 2)}\n`) + console.log(`✓ 基线已更新:${found.length} 条边 → ${posixify(BASELINE)}`) + process.exit(0) +} + +if (args.has('--json')) { + console.log(JSON.stringify({ counts, unclassified, added, known, fixed }, null, 2)) + process.exit(added.length > 0 ? 1 : 0) +} + +console.log('=== 分层依赖检查(docs/architecture.md R1–R4)===') +console.log(`扫描:${files.length} 个 .ts|归属 ①${counts[1]} ②${counts[2]} ③${counts[3]} ④${counts[4]}|未归类 ${unclassified.length}`) +if (unclassified.length > 0) console.log(`⚠️ 未归类(新增目录必须先定层):${unclassified.slice(0, 5).join(', ')}${unclassified.length > 5 ? ' …' : ''}`) +console.log(`现存违规 ${found.length} 条(其中基线内 ${known.length} 条)`) + +if (added.length > 0) { + console.log(`\n🔴 新增违规 ${added.length} 条(⛔ 规则已生效,不得引入):`) + for (const v of added) console.log(` ${v.detail}`) +} +if (fixed.length > 0) { + console.log(`\n🟢 基线内已消除 ${fixed.length} 条(可 --update-baseline 收紧):`) + for (const e of fixed.slice(0, 8)) console.log(` ${e}`) +} +if (added.length === 0) { + console.log(baseline === null ? '\n⚠️ 尚无基线文件:首次请先跑 --update-baseline 建立(否则任何存量违规都会被算作"新增")' : '\n✅ 无新增违规') +} +process.exit(added.length > 0 ? 1 : 0) diff --git a/scripts/layering-baseline.json b/scripts/layering-baseline.json new file mode 100644 index 0000000..a94dd00 --- /dev/null +++ b/scripts/layering-baseline.json @@ -0,0 +1,10 @@ +{ + "note": "R1–R4 存量违规基线(ratchet):只许减、不许增;新增违规会让 check:layering 失败", + "edges": [ + "src/fs/local-user-fs.ts → src/web/middleware/fs-guard.ts", + "src/fs/remote-user-fs.ts → src/web/middleware/fs-guard.ts", + "src/fs/workspace.ts → src/web/middleware/fs-guard.ts", + "src/supervisor/proxy.ts → src/web/auth.ts", + "src/supervisor/proxy.ts → src/web/middleware/authn.ts" + ] +} diff --git a/src/config.ts b/src/config.ts index 92ca994..a0e570c 100644 --- a/src/config.ts +++ b/src/config.ts @@ -111,6 +111,13 @@ export interface ServerConfig { * 所有 worker 的 dataRoot 必须是同一个绝对路径(同镜像即可满足,设计 §14.3)。 */ clusterWorkerDataRoot: string + /** + * **worker 侧**会合地址(覆盖网络 S1):worker 主动拨入的 SSH 反向隧道落点。 + * 形如 `ssh://root@47.77.182.89:32022`(也接受不带 scheme 的 `root@host:port`)。空 = 隧道关闭。 + * ⚠️ 与旧变量 `DSHS_TUNNEL_TARGET` **双路径并存**(新变量优先、旧变量兜底)⇒ + * 删掉新 env 即回到旧路径,**零代码回滚**。 + */ + clusterRendezvousUrl: string } /** Untyped overrides collected from argv / env. */ @@ -157,6 +164,7 @@ export interface ConfigOverrides { clusterAgentToken?: string clusterInstanceHost?: string clusterWorkerDataRoot?: string + clusterRendezvousUrl?: string } const DEFAULT_HOST = '127.0.0.1' @@ -352,5 +360,12 @@ export function resolveConfig(overrides: ConfigOverrides = {}): ServerConfig { clusterAgentToken: overrides.clusterAgentToken ?? process.env.DSHS_CLUSTER_AGENT_TOKEN ?? '', clusterInstanceHost: overrides.clusterInstanceHost ?? process.env.DSHS_CLUSTER_INSTANCE_HOST ?? '127.0.0.1', clusterWorkerDataRoot: overrides.clusterWorkerDataRoot ?? process.env.DSHS_CLUSTER_WORKER_DATA_ROOT ?? '', + // 覆盖网络 S1:会合地址出 env。新变量 `DSHS_RENDEZVOUS_URL` 优先,旧变量 `DSHS_TUNNEL_TARGET` + // 兜底(两台机器可分先后改;删掉新 env 即回滚到旧路径,**不需要回滚代码**)。 + clusterRendezvousUrl: + overrides.clusterRendezvousUrl ?? + process.env.DSHS_RENDEZVOUS_URL ?? + process.env.DSHS_TUNNEL_TARGET ?? + '', } } diff --git a/src/net/reachability.ts b/src/net/reachability.ts new file mode 100644 index 0000000..c5a3ffe --- /dev/null +++ b/src/net/reachability.ts @@ -0,0 +1,99 @@ +/** + * 可达性(Reachability)—— **"怎么到这台 worker"的可序列化描述**(覆盖网络 S0)。 + * + * ## 为什么需要它 + * 现状里"怎么到一台 worker"被写死成两件事:`dsh_hosts.endpoint` 存一个 URL, + * `clusterInstanceHost` 存一个**全局**主机名。两者都**表达不出"经谁中转"**,而现网 + * 恰好有两类语义完全不同的 host,字符串却同形: + * + * | hostId | endpoint | 真实语义 | + * |---|---|---| + * | `w-47` | `http://127.0.0.1:19100` | **直连本机**(Manager 与 worker 同机,node 直接监听) | + * | `w-106` | `http://127.0.0.1:19000` | **经 47 上 sshd 的反向隧道落点**(隧道的副作用) | + * + * ⇒ 换会合 / 中继组件时,表里**无法表达**「via(经哪个中继)+ 真实可达地址」, + * 只能改表;越晚改代价越大(会合中继拆分方案 §2 C3)。 + * + * `Reachability` 把这件事显式化:`via` 指向一个 `Rendezvous` 实现,`address` 是真实地址。 + * + * ## 边界(勿破) + * 本模块**只做地址的表征与解析**,不承载任何权威状态 —— 归属 / 租约 / 骨干资格 + * 一律仍只由控制面写(与 `集群化改造方案 §1.3` 数据分层一致)。 + * + * @module dshs/net/reachability + */ + +/** 同机直连:Manager 与 worker 在同一台机器上,不经任何中转。 */ +export const VIA_LOCAL = 'local' + +/** 今天唯一在跑的中转方式 = **Manager 主机上的 sshd 反向隧道**(S4 之后应被 relay 取代)。 */ +export const VIA_MANAGER_SSH = 'manager-ssh' + +export interface Reachability { + /** 哪台 worker。 */ + hostId: string + /** **经谁可达** —— 一个 `Rendezvous` 实现的 id。 */ + via: string + /** agent 的真实地址 `host:port`(**不含 scheme**)。 */ + address: string + /** 传输层。 */ + scheme: 'http' | 'https' +} + +/** 只取址所需的最小形状 —— 避免本模块反向依赖 `supervisor`。 */ +export interface HostAddressable { + hostId?: string + agentUrl?: string + reachability?: Reachability +} + +/** + * 把可达性拼成 agent 基址 —— **全仓唯一的拼接点**,别在别处再拼 `scheme://address`。 + */ +export function agentBaseUrl(reach: Reachability): string { + return `${reach.scheme}://${reach.address}`.replace(/\/+$/, '') +} + +/** + * 取一台 host 的 agent 基址:**可达性优先,回退旧 `agentUrl`**。 + * + * S0 阶段的等价性:现网每个 host 都只有 `agentUrl`(`reachability` 全为 `undefined`) + * ⇒ 本函数返回的就是原先直接用的那个字符串,**行为零变化**。 + * + * ⚠️ 两者皆缺时**抛错**,不返回空串 —— "静默打到空地址"是跨机下最难查的失败。 + */ +export function agentBaseUrlOf(host: HostAddressable): string { + if (host.reachability !== undefined) return agentBaseUrl(host.reachability) + if (host.agentUrl !== undefined && host.agentUrl !== '') return host.agentUrl.replace(/\/+$/, '') + throw new Error(`host "${host.hostId ?? '?'}" 既无 reachability 也无 agentUrl:拒绝静默降级`) +} + +/** + * 从旧的 `endpoint` 字符串解析出 `Reachability`(S2 迁移回填用)。 + * + * 兼容面:`endpoint` 历史上是完整 URL(`http://127.0.0.1:19000`),也容忍裸 + * `host:port` —— 没写 scheme 时按 `http` 处理,与 `RemoteSpawner` 原先"直接把它当 + * fetch 基址"的行为一致(fetch 会补 `http://`)。 + */ +export function parseReachability( + hostId: string, + endpoint: string, + via: string = VIA_MANAGER_SSH, +): Reachability { + const trimmed = endpoint.trim() + const matched = /^(https?):\/\/(.*)$/i.exec(trimmed) + if (matched !== null) { + return { + hostId, + via, + address: matched[2].replace(/\/+$/, ''), + scheme: matched[1].toLowerCase() === 'https' ? 'https' : 'http', + } + } + return { hostId, via, address: trimmed.replace(/\/+$/, ''), scheme: 'http' } +} + +/** `Reachability` → 旧 `endpoint` 字符串(与 `parseReachability` 互逆,回填/回滚用)。 */ +export function toEndpoint(reach: Reachability): string { + return agentBaseUrl(reach) +} diff --git a/src/net/rendezvous.ts b/src/net/rendezvous.ts new file mode 100644 index 0000000..bc1f363 --- /dev/null +++ b/src/net/rendezvous.ts @@ -0,0 +1,112 @@ +/** + * 会合(Rendezvous)—— **"该拨谁、经谁到"的解析器**(覆盖网络 S0)。 + * + * ## 定位(分层口径,勿破) + * | 组件 | 职责 | 可多实例? | 权威状态 | + * |---|---|---|---| + * | **会合 (rendezvous)** | 收 worker 注册、回"该拨谁"、下发中继分配;**不承载数据面流量** | ✅ 无状态可复制 | ❌ 只有位置视图 | + * | **中继 (relay)** | 数据面:worker 拨它 → Manager / 其他节点经它到 worker | ✅ | ❌ | + * | **控制面 (Manager)** | 归属 / 租约 / 骨干资格 / 容量准入 | ❌ 单点 | ✅ 唯一写入者 | + * + * ⛔ **硬约束**:会合与中继**不得**写入 `dsh_instances.host_id` / `epoch` / 骨干资格 —— + * 否则就是双写脑裂。 + * + * ## 现状与目标 + * 今天"会合 + 中继"**不是一个组件,而是 Manager 主机上 sshd 的副作用**: + * 会合点 = `47.77.182.89:32022`,中继落点 = Manager 的 `127.0.0.1`。 + * 本模块的作用是**先把接口抽出来**,让 SSH 隧道退化成"第一个可替换实现" + * —— ⛔ 这一步**不换协议**,只换绑定与寻址(换 WireGuard / TURN 属远期)。 + * + * @module dshs/net/rendezvous + */ + +import { VIA_LOCAL, VIA_MANAGER_SSH, type Reachability } from './reachability.js' + +export interface Rendezvous { + /** 实现 id —— `Reachability.via` 指向它。 */ + readonly id: string + /** 这个会合点**本身**怎么拨(诊断 / 管理面展示用)。 */ + dialTarget(): string + /** + * 解析某台 worker 的可达性。 + * ⚠️ **本实现管不到 ⇒ 回 `undefined`,不抛** —— 由调用方决定回退哪种实现 + * (抛错会让"多实现并存"的过渡期没法跑)。 + */ + resolve(hostId: string): Promise +} + +/** 由调用方提供"hostId → `host:port`"的查表函数(会合实现不直接连 DB)。 */ +export type AddressLookup = (hostId: string) => string | undefined + +/** 同机直连:Manager 能直接连到 worker 的端口,不经任何中转(`w-47` 就是这一类)。 */ +export class LocalRendezvous implements Rendezvous { + readonly id = VIA_LOCAL + + constructor(private readonly addressOf: AddressLookup) {} + + dialTarget(): string { + return '(direct)' + } + + async resolve(hostId: string): Promise { + const address = this.addressOf(hostId) + return address === undefined ? undefined : { hostId, via: this.id, address, scheme: 'http' } + } +} + +/** + * Manager 主机上的 sshd 反向隧道 —— **当前唯一在跑的实现**。 + * + * 语义:worker 主动 `ssh -R :127.0.0.1: root@:`,把端口投到 + * Manager 的 **loopback**(`127.0.0.1:<同号端口>`)⇒ Manager 经 `127.0.0.1:` 到达它。 + * + * ⚠️ 这正是要拆掉的那一层(C1 会合地址硬编码 + C2 中继落点 = Manager loopback)。 + * S4 把"接收 worker 反拨"搬进独立单元 `dshs-relay.service` 后,本类应被 `relay:` + * 实现替换,而**调用方不需要改**(只认 `Rendezvous` 接口)。 + */ +export class ManagerSshRendezvous implements Rendezvous { + readonly id = VIA_MANAGER_SSH + + constructor( + private readonly opts: { + /** 会合点的 SSH 目标,如 `root@47.77.182.89:32022`。 */ + target: string + addressOf: AddressLookup + }, + ) {} + + dialTarget(): string { + return this.opts.target + } + + async resolve(hostId: string): Promise { + const address = this.opts.addressOf(hostId) + return address === undefined ? undefined : { hostId, via: this.id, address, scheme: 'http' } + } +} + +/** + * 按 `via` 选实现的注册表。 + * + * 为什么需要:S2 起 `dsh_hosts` 会带 `via` 列,`hostsProvider` 必须"**先读 via → + * 选对应实现 → 解析成 `Reachability`**";`via` 未设时回退旧 `endpoint` 语义。 + */ +export class RendezvousRegistry { + private readonly impls = new Map() + + constructor(impls: readonly Rendezvous[] = []) { + for (const impl of impls) this.register(impl) + } + + register(impl: Rendezvous): void { + this.impls.set(impl.id, impl) + } + + get(id: string): Rendezvous | undefined { + return this.impls.get(id) + } + + ids(): string[] { + return [...this.impls.keys()] + } +} diff --git a/src/supervisor/remote-spawner.ts b/src/supervisor/remote-spawner.ts index 7f4e5a5..0568698 100644 --- a/src/supervisor/remote-spawner.ts +++ b/src/supervisor/remote-spawner.ts @@ -18,14 +18,26 @@ * @module dshs/supervisor/remote-spawner */ import { randomUUID } from 'node:crypto' +import { agentBaseUrlOf, type Reachability } from '../net/reachability.js' import { AGENT_TOKEN_HEADER } from '../worker/agent.js' import type { Endpoint, Instance, Spawner, UserStatus } from './spawner.js' -/** 一台 worker 的接入信息。 */ +/** + * 一台 worker 的接入信息。 + * + * ⚠️ `agentUrl` 与 `reachability` **至少要有一个** —— 取址一律走 `agentBaseUrlOf()` + * (唯一入口,别在调用点自己拼字符串): + * · `reachability` = S0 引入的**可达性描述**,比 `agentUrl` 多一层语义 —— + * **经谁中转**(`via`)。现网两类 host 的 `endpoint` 字符串同形但语义完全不同 + * (`w-47` 是直连本机、`w-106` 是 Manager 上的隧道落点),只有它能表达。 + * · `agentUrl` = 旧字段,保留向后兼容。 + */ export interface ClusterHost { hostId: string - /** agent 基址。 */ - agentUrl: string + /** agent 基址(旧字段;与 `reachability` 至少给一个)。 */ + agentUrl?: string + /** **可达性**:经谁中转 + 真实地址。给了就优先于 `agentUrl`(S2 起由 `dsh_hosts.via` 驱动)。 */ + reachability?: Reachability /** 与 agent 约定的共享密钥。 */ token: string /** 代理时使用的主机(同机 1a = `127.0.0.1`;跨机填 Worker 内网 IP)。 */ @@ -81,15 +93,17 @@ export class RemoteSpawner implements Spawner { private directoryLoadedAt = 0 constructor(options: RemoteSpawnerOptions) { + // ⚠️ 这里**不再做** `replace(/\/$/,'')` 归一化 —— 归一化统一在 `agentBaseUrlOf()` + // 里做(幂等)。存归一化值会让"可达性与 agentUrl 两套表示"混在一起,S2 之后难拆。 this.defaultHost = { hostId: options.defaultHostId ?? 'local', - agentUrl: options.agentUrl.replace(/\/$/, ''), + agentUrl: options.agentUrl, token: options.token, instanceHost: options.instanceHost ?? '127.0.0.1', } this.hosts.set(this.defaultHost.hostId, this.defaultHost) for (const host of options.hosts ?? []) { - this.hosts.set(host.hostId, { ...host, agentUrl: host.agentUrl.replace(/\/$/, '') }) + this.hosts.set(host.hostId, { ...host }) } this.timeoutMs = options.timeoutMs ?? 10_000 this.doFetch = options.fetchImpl ?? fetch @@ -110,7 +124,7 @@ export class RemoteSpawner implements Spawner { this.directoryLoadedAt = Date.now() try { for (const host of await this.hostsProvider()) { - this.hosts.set(host.hostId, { ...host, agentUrl: host.agentUrl.replace(/\/$/, '') }) + this.hosts.set(host.hostId, { ...host }) } this.hosts.set(this.defaultHost.hostId, this.defaultHost) } catch { @@ -177,7 +191,7 @@ export class RemoteSpawner implements Spawner { for (let attempt = 0; attempt <= RETRY_DELAYS_MS.length; attempt += 1) { if (attempt > 0) await new Promise((r) => setTimeout(r, RETRY_DELAYS_MS[attempt - 1])) try { - const res = await this.doFetch(`${host.agentUrl}${path}`, { + const res = await this.doFetch(`${agentBaseUrlOf(host)}${path}`, { method, headers: { [AGENT_TOKEN_HEADER]: host.token, @@ -314,7 +328,7 @@ export class RemoteSpawner implements Spawner { touch(userId: string): void { void this.hostFor(userId) .then((host) => - this.doFetch(`${host.agentUrl}/touch/${encodeURIComponent(userId)}`, { + this.doFetch(`${agentBaseUrlOf(host)}/touch/${encodeURIComponent(userId)}`, { method: 'POST', headers: { [AGENT_TOKEN_HEADER]: host.token }, }), diff --git a/src/web/server.ts b/src/web/server.ts index 64da356..44bde23 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -17,6 +17,7 @@ import { RemoteUserFs } from '../fs/remote-user-fs.js' import type { UserFs } from '../fs/user-fs.js' import { decrypt, deriveKey } from '../crypto.js' import { hashUid } from '../isolation.js' +import { agentBaseUrlOf } from '../net/reachability.js' import { LocalSpawner } from '../supervisor/orchestrator.js' import { LeasedSpawner } from '../supervisor/leased-spawner.js' import { RemoteSpawner, type ClusterHost } from '../supervisor/remote-spawner.js' @@ -339,7 +340,8 @@ export async function buildServer(config: ServerConfig): Promise { const h = hostDirectory.get(hostId) - return h === undefined ? undefined : { agentUrl: h.agentUrl, token: h.token } + // 取址统一走可达性入口(S0)—— `agentUrl` 已降级为可选旧字段 + return h === undefined ? undefined : { agentUrl: agentBaseUrlOf(h), token: h.token } }, }, )) @@ -354,7 +356,8 @@ export async function buildServer(config: ServerConfig): Promise { const h = hostDirectory.get(hostId) - return h === undefined ? undefined : { agentUrl: h.agentUrl, token: h.token } + // 同上一处:取址只经可达性入口(文件面与实例面共用同一份路由,别各自拼) + return h === undefined ? undefined : { agentUrl: agentBaseUrlOf(h), token: h.token } }, }) // T08 S5:cluster 模式下**所有 worker 的 dataRoot 必须是同一绝对路径**(基线约定, diff --git a/src/worker/agent.ts b/src/worker/agent.ts index 349e956..219a1d7 100644 --- a/src/worker/agent.ts +++ b/src/worker/agent.ts @@ -28,7 +28,7 @@ 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 { normalizeTunnelTarget, SshTunnel } from './tunnel.js' import type { Instance } from '../supervisor/spawner.js' /** 绑定的头部名(Manager/agent 双方约定)。 */ @@ -150,8 +150,13 @@ export function buildWorkerAgent( /** * 反向隧道(可选)。静态转发 = **agent 自身端口** + `DSHS_TUNNEL_STATIC_PORTS`(如控制面 PG); * 实例端口在 launch/stop 时动态加减,并在 `/healthz`(Manager 的心跳)里**对账自愈**。 + * + * S1:会合地址**优先取 config**(`DSHS_RENDEZVOUS_URL` → 兜底 `DSHS_TUNNEL_TARGET`,在 + * `config.ts` 里单点解析);显式 `options.tunnelTarget` 仍是最优先(测试/嵌入用)。 */ - const tunnelTarget = options.tunnelTarget ?? process.env.DSHS_TUNNEL_TARGET ?? '' + const tunnelTarget = normalizeTunnelTarget( + options.tunnelTarget ?? config.clusterRendezvousUrl ?? '', + ) const staticPorts = [ options.port, ...(process.env.DSHS_TUNNEL_STATIC_PORTS ?? '') diff --git a/src/worker/tunnel.ts b/src/worker/tunnel.ts index 0b16baa..aa575c2 100644 --- a/src/worker/tunnel.ts +++ b/src/worker/tunnel.ts @@ -13,7 +13,8 @@ * 在**同一条长连接**上加/减转发,不必为每个端口重开连接。 * * ⚠️ 定位:这是**演练级**传输(生产长期方案见设计 §2.3:受控网段白名单或隧道服务)。 - * ⚠️ 默认**关闭**:只有设了 `DSHS_TUNNEL_TARGET` 才启用 ⇒ 对同机/单机形态零影响。 + * ⚠️ 默认**关闭**:只有设了 `DSHS_RENDEZVOUS_URL` 才启用 ⇒ 对同机/单机形态零影响。 + * 旧变量 `DSHS_TUNNEL_TARGET` 作为**兜底**保留(新变量未设时才用它)。 * * @module dshs/worker/tunnel */ @@ -23,8 +24,28 @@ import { promisify } from 'node:util' const run = promisify(execFile) +/** + * 归一化会合地址(覆盖网络 S1)。 + * + * 为什么要它:会合点从"硬写在 env 里的 `user@host:port`"升格为**带 scheme 的 URL** + * (`ssh://root@47.77.182.89:32022`)—— 以后换传输协议(中继/隧道服务)只改 scheme。 + * 而本类其余代码如下按 `user@host:port` 切分 ⇒ 必须在**入口处**剥掉 scheme: + * 否则 `'ssh://root@h:32022'.split(':')` 会切成三截,把 `ssh` 当成主机名。 + * + * 两种写法都接受(**单点归一,调用方不必判断**): + * · `ssh://root@47.77.182.89:32022` → `root@47.77.182.89:32022` + * · `root@47.77.182.89:32022` → 原样(兼容历史 env `DSHS_TUNNEL_TARGET`) + */ +export function normalizeTunnelTarget(raw: string): string { + const trimmed = raw.trim() + if (trimmed === '') return '' + return trimmed + .replace(/^[a-z][a-z0-9+.-]*:\/\//i, '') // 剥 scheme(ssh:// / dshs+ssh:// …) + .replace(/\/+$/, '') // 去掉可能的尾斜杠 +} + export interface TunnelOptions { - /** 拨入目标,形如 `root@47.77.182.89:32022`。 */ + /** 拨入目标,形如 `root@47.77.182.89:32022`(带 `ssh://` 前缀也接受,见 {@link normalizeTunnelTarget})。 */ target: string /** 私钥路径(建议专用、且在 Manager 侧用 `restrict,port-forwarding` 限权)。 */ identity: string @@ -52,11 +73,17 @@ export class SshTunnel { constructor(options: TunnelOptions) { // `user@host:port` 里的 port 是 **SSH 端口**(不是转发的端口)—— 47 上用 32022, // 必须经 `-p` 传,否则会去连 22 而失败。 - const [hostPart, portPart] = options.target.split(':') + // S1:先归一化(剥 `ssh://` scheme),再按 `user@host:port` 切分。 + const normalized = normalizeTunnelTarget(options.target) + const [hostPart, portPart] = normalized.split(':') + if (portPart !== undefined && !/^\d+$/.test(portPart.trim())) { + // 宁可起不来也不要"静默连到 22 端口":地址写错必须吵。 + throw new Error(`非法的会合地址(端口必须是数字):${options.target}`) + } this.hostPart = hostPart this.portPart = portPart === undefined ? undefined : Number(portPart) this.opts = { - target: options.target, + target: normalized, identity: options.identity, controlPath: options.controlPath, staticPorts: options.staticPorts ?? [], diff --git a/test/reachability.test.mjs b/test/reachability.test.mjs new file mode 100644 index 0000000..c40f11b --- /dev/null +++ b/test/reachability.test.mjs @@ -0,0 +1,126 @@ +/** + * 覆盖网络 S0 · 可达性 / 会合接口单测(纯函数,不连网、不起进程)。 + * 运行:node --test test/reachability.test.mjs(已含在 npm test / npm verify 中) + * + * ## 这个测试的职责 + * S0 的唯一验收标准是「**行为零变化**」。所以本文件的核心断言是 + * **「可达性解析前后的 agent 基址逐字相等」** —— 把现网 `dsh_hosts` 的真实两行 + * 原样写进判据里(**不靠肉眼、不靠人工比对**): + * + * | hostId | dsh_hosts.endpoint(2026-09-16 实测) | 真实语义 | + * |---|---|---| + * | `w-106` | `http://127.0.0.1:19000` | 经 47 上 sshd 的反向隧道落点 | + * | `w-47` | `http://127.0.0.1:19100` | 同机直连(node 自己监听) | + */ +import { test } from 'node:test' +import assert from 'node:assert/strict' + +import { + agentBaseUrl, + agentBaseUrlOf, + parseReachability, + toEndpoint, + VIA_LOCAL, + VIA_MANAGER_SSH, +} from '../lib/net/reachability.js' +import { LocalRendezvous, ManagerSshRendezvous, RendezvousRegistry } from '../lib/net/rendezvous.js' + +/** 现网真实两条(2026-09-16 在 47 上 `SELECT id, endpoint FROM dsh_hosts` 实测)。 */ +const LIVE_HOSTS = [ + { hostId: 'w-106', endpoint: 'http://127.0.0.1:19000', via: VIA_MANAGER_SSH }, + { hostId: 'w-47', endpoint: 'http://127.0.0.1:19100', via: VIA_LOCAL }, +] + +test('S0 等价性:旧 agentUrl 取址与可达性取址逐字相等', () => { + for (const h of LIVE_HOSTS) { + // 旧路径(今天生产用的):直接把 endpoint 当 agent 基址 + const legacy = h.endpoint.replace(/\/$/, '') + // 新路径(S0 起的唯一入口):没给 reachability ⇒ 必须回退到同一个值 + assert.equal(agentBaseUrlOf({ hostId: h.hostId, agentUrl: h.endpoint }), legacy, h.hostId) + } +}) + +test('parseReachability → agentBaseUrl 对现网两行是往返恒等的', () => { + for (const h of LIVE_HOSTS) { + const reach = parseReachability(h.hostId, h.endpoint, h.via) + assert.equal(agentBaseUrl(reach), h.endpoint, `${h.hostId} 往返不一致`) + assert.equal(toEndpoint(reach), h.endpoint, `${h.hostId} toEndpoint 不一致`) + assert.equal(reach.scheme, 'http') + assert.deepEqual( + { hostId: reach.hostId, via: reach.via, address: reach.address }, + { hostId: h.hostId, via: h.via, address: h.endpoint.slice('http://'.length) }, + ) + } +}) + +test('parseReachability:https / 裸 host:port / 尾斜杠 三种兼容面', () => { + assert.deepEqual(parseReachability('h', 'https://a.example:8443', 'x'), { + hostId: 'h', + via: 'x', + address: 'a.example:8443', + scheme: 'https', + }) + // 没写 scheme ⇒ 按 http(与 fetch 的补全行为一致) + assert.deepEqual(parseReachability('h', '10.0.0.5:19000', 'x'), { + hostId: 'h', + via: 'x', + address: '10.0.0.5:19000', + scheme: 'http', + }) + assert.equal(agentBaseUrl(parseReachability('h', 'http://127.0.0.1:19000/', 'x')), 'http://127.0.0.1:19000') +}) + +test('可达性优先于旧 agentUrl', () => { + const host = { + hostId: 'w-106', + agentUrl: 'http://stale.example:1', + reachability: { hostId: 'w-106', via: VIA_MANAGER_SSH, address: '127.0.0.1:19000', scheme: 'http' }, + } + assert.equal(agentBaseUrlOf(host), 'http://127.0.0.1:19000') +}) + +test('两者皆缺 ⇒ 抛错(禁止静默打到空地址)', () => { + assert.throws(() => agentBaseUrlOf({ hostId: 'w-x' }), /既无 reachability 也无 agentUrl/) + assert.throws(() => agentBaseUrlOf({ hostId: 'w-x', agentUrl: '' }), /既无 reachability 也无 agentUrl/) +}) + +test('LocalRendezvous:命中给 local,未命中回 undefined 不抛', async () => { + const table = new Map([['w-47', '127.0.0.1:19100']]) + const rv = new LocalRendezvous((id) => table.get(id)) + assert.equal(rv.id, VIA_LOCAL) + assert.equal(rv.dialTarget(), '(direct)') + assert.deepEqual(await rv.resolve('w-47'), { + hostId: 'w-47', + via: VIA_LOCAL, + address: '127.0.0.1:19100', + scheme: 'http', + }) + assert.equal(await rv.resolve('w-106'), undefined) +}) + +test('ManagerSshRendezvous:解析出的基址必须等于现网 endpoint(S2 迁移判据)', async () => { + const table = new Map([ + ['w-106', '127.0.0.1:19000'], + ['w-47', '127.0.0.1:19100'], + ]) + const rv = new ManagerSshRendezvous({ + target: 'root@47.77.182.89:32022', + addressOf: (id) => table.get(id), + }) + assert.equal(rv.id, VIA_MANAGER_SSH) + assert.equal(rv.dialTarget(), 'root@47.77.182.89:32022') + for (const h of LIVE_HOSTS) { + const resolved = await rv.resolve(h.hostId) + assert.equal(agentBaseUrl(resolved), h.endpoint, `${h.hostId} 迁移后基址变了`) + } +}) + +test('RendezvousRegistry:按 via 取实现', () => { + const local = new LocalRendezvous(() => undefined) + const ssh = new ManagerSshRendezvous({ target: 't', addressOf: () => undefined }) + const reg = new RendezvousRegistry([local, ssh]) + assert.equal(reg.get(VIA_LOCAL), local) + assert.equal(reg.get(VIA_MANAGER_SSH), ssh) + assert.equal(reg.get('relay:backbone-1'), undefined) + assert.deepEqual(reg.ids().sort(), [VIA_LOCAL, VIA_MANAGER_SSH].sort()) +})