Files
dsh_ai1net_server/交付物/平台侧-插件数据面执行器-20260926.mjs
admin c1b5e4d966 chore(工作区): 全量入库 + 补齐 .gitignore(以工作区为准)
- 变更规模:新增 514 / 修改 62 / 重命名 155 / 删除 4(归档重组与文档轮次)
- .gitignore 修:`归档/**/db-cwd归一-备份-*/` —— 原规则写绝对层级(归档/db-cwd归一-…),
  目录搬进 归档/配置与备份/ 后**静默失效**,43 MB 的 DB 备份又变成未跟踪
- .gitignore 补:嵌套 git 内部数据(归档/内嵌git-20261008/、归档/skills-git-旧线-20261007/dotgit-原样移出/)
- .gitignore 补:运行态与部署副本(.workbuddy/collab/、.workbuddy/tools/、.workbuddy/.load-pending、.workbuddy/tmp-*)
- .gitignore 补:备份件(*.bak-*)
- 未跟踪文件从 2190 降到 890(其余为 归档/ 归档件与 .workbuddy/memory/ 知识文件,按口径入库)
2026-10-10 23:13:22 +08:00

316 lines
17 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env node
/**
* 平台侧 · 插件数据面执行器(S3「数据面」· 2026-09-26 第 41 棒 · 插件投放与分库线)
*
* 定位:把「一个包一个库」那四件事在**平台侧**跑起来 —— 声明 / 建库 / 迁移(按模块分步)/
* 对账,外带两个取证件:**回读校验**(负控的靶子)与**单模块回滚演练**(事务内 · 零副作用)。
*
* 🔴 为什么本文件在**平台侧**、⛔ 不在插件包内:
* 「DB-03 §二」红线 1/2 —— 插件**不得直连数据库**、**不得写原生 SQL**。包内只有
* **声明**(`dsh.data.yaml`)与**插件侧常量**(`lib/data-limits.js`);建库与 DDL 只能由平台执行。
* ⇒ 本文件属"平台进程侧的动作",与插件包是两回事(⛔ 别把它拷进包里)。
*
* 🔴 复用既有机制(⛔ 一行都不重造):解析器 `schema.js` / 执行器 `datastore.js` /
* 计划与指纹 `diff.js` —— 全部 import 平台**生产构建产物** `/opt/dshs/lib/db/plugin-data/*.js`,
* 与门户 `POST /api/plugins/business/datastore*` 走的是**同一批函数**。
*
* 用法(在 47 上跑;`--dir` 指向**含 `package.json` + `dsh.data.yaml` 的目录**):
* node p41-datastore-exec.mjs decl --dir <dir>
* node p41-datastore-exec.mjs plan --dir <dir> [--ledger]
* node p41-datastore-exec.mjs create --dir <dir>
* node p41-datastore-exec.mjs migrate --dir <dir> --modules M1,M2,M3,M4,M5
* node p41-datastore-exec.mjs snapshot --dir <dir> # 登记声明快照(回读的基准)
* node p41-datastore-exec.mjs readback --dir <dir> # 回读比对(不一致 ⇒ rc=1)
* node p41-datastore-exec.mjs reconcile # pg_database × 台账
* node p41-datastore-exec.mjs rollback-probe --module M1 # 单模块回滚演练(事务内回滚)
*
* 退出码:0 绿 / 1 红 / 3 未取到(与 S0 四把尺子同一口径)。
* ⛔ 本脚本不 DROP 任何库、不删任何表 —— 库级回滚口径是「改名冻结」(方案 §3.10.5)。
*/
import { createHash } from 'node:crypto'
const LIB = '/opt/dshs/lib/db/plugin-data'
const { parseDeclFromDir } = await import(`${LIB}/schema.js`)
const {
openControlPool,
createPluginDatabase,
readCurrentSchema,
executePlan,
reconcilePluginDatabases,
pluginDbUrlOf,
} = await import(`${LIB}/datastore.js`)
const { diffDecl } = await import(`${LIB}/diff.js`)
const pgMod = await import('/opt/dshs/node_modules/pg/lib/index.js')
const PgClient = (pgMod.default ?? pgMod).Client
const PKG_NAME = '@dsh-local/ai1net'
const ACTOR = 'p41-exec-plugin-dataline'
/** 五个模块段(与包内 `scripts/build.mjs` 门禁 ①d 的 MODULE_SEGS 同表 · ⛔ 同改)。 */
const MODULES = [
{ key: 'M1', seg: 'tenant', label: '多租户' },
{ key: 'M2', seg: 'model', label: '模型管理' },
{ key: 'M3', seg: 'plugin', label: '技能插件管理' },
{ key: 'M4', seg: 'net', label: '覆盖网络界面' },
{ key: 'M5', seg: 'im', label: 'IM 插件侧' },
]
const moduleOf = (tableName) => MODULES.find((m) => tableName.startsWith(m.seg + '_')) ?? null
const argv = process.argv.slice(2)
const CMD = argv[0] ?? ''
const flag = (name) => argv.includes(`--${name}`)
const opt = (name, dflt = '') => {
const i = argv.indexOf(`--${name}`)
return i >= 0 && argv[i + 1] !== undefined && !argv[i + 1].startsWith('--') ? argv[i + 1] : dflt
}
let fails = 0
const ok = (m) => console.log(` ✅ ${m}`)
const bad = (m) => {
fails++
console.log(` ❌ ${m}`)
}
const info = (m) => console.log(` · ${m}`)
function die(code, msg) {
console.log(msg)
process.exit(code)
}
const DB_URL = process.env.DSHS_DB_URL ?? ''
if (DB_URL === '') die(3, '未取到:环境变量 DSHS_DB_URL 为空(平台库连接串只在 drop-in cluster.conf 里,⛔ 不硬编码)')
/** 声明结构摘要 —— 与"当前库现状 / 计划项"**无关**,只覆盖声明本身(故可跨建表前后比对)。 */
function declDigestOf(decl) {
const canon = decl.tables
.slice()
.sort((a, b) => a.name.localeCompare(b.name))
.map((t) => {
const cols = t.columns
.slice()
.sort((a, b) => a.name.localeCompare(b.name))
.map((c) => `${c.name}:${c.type}:${c.notNull ? 1 : 0}:${c.default === undefined ? '-' : JSON.stringify(c.default)}:${c.maxBytes ?? ''}`)
.join(',')
const idx = (t.indexes ?? []).map((i) => `${i.columns.join('+')}${i.unique ? '#u' : ''}`).sort().join(';')
return `${t.name}(${t.scope})[${cols}]{${idx}}`
})
.join('\n')
return createHash('sha256').update(`v${decl.schemaVersion}\n${canon}`).digest('hex')
}
async function loadDecl() {
const dir = opt('dir')
if (dir === '') die(3, '未取到:缺 --dir(须指向含 package.json + dsh.data.yaml 的目录)')
const res = parseDeclFromDir(dir, PKG_NAME)
if (res.level !== 'ok' || res.decl === null) {
console.log(` 声明校验:level=${res.level} origin=${res.origin}`)
for (const f of res.findings ?? []) console.log(` - [${f.code}] ${f.target}: ${f.detail}`)
die(1, '❌ 声明不合法(平台解析器裁决)')
}
ok(`声明解析:level=ok origin=${res.origin} 库名=${res.decl.dbName} schemaVersion=${res.decl.schemaVersion} 表=${res.decl.tables.length}`)
return res.decl
}
const ledgerGet = async (pool, pluginId) => {
const { rows } = await pool.query(
'SELECT plugin_id, db_name, state, schema_version, plan_hash, last_error, created_at, updated_at FROM plugin_datastores WHERE plugin_id = $1',
[pluginId],
)
return rows[0] ?? null
}
/** 台账 upsert —— 逐字照抄 `repo.ts` 的 `upsertPluginDatastore`(含 `plan_hash` 的 COALESCE 语义)。 */
async function ledgerUpsert(pool, input) {
const now = Date.now()
const prev = await ledgerGet(pool, input.pluginId)
const planHash = input.planHash === undefined ? (prev?.plan_hash ?? null) : input.planHash
await pool.query(
`INSERT INTO plugin_datastores (plugin_id, db_name, state, schema_version, plan_hash, last_error, created_at, updated_at)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8)
ON CONFLICT(plugin_id) DO UPDATE SET
db_name = EXCLUDED.db_name, state = EXCLUDED.state, schema_version = EXCLUDED.schema_version,
plan_hash = EXCLUDED.plan_hash, last_error = EXCLUDED.last_error, updated_at = EXCLUDED.updated_at`,
[input.pluginId, input.dbName, input.state, input.schemaVersion ?? null, planHash, input.lastError ?? null, prev?.created_at ?? now, now],
)
}
const audit = (pool, pluginId, action, detail) =>
pool.query('INSERT INTO plugin_data_audit (ts, actor, plugin_id, action, detail) VALUES ($1,$2,$3,$4,$5)', [
Date.now(),
ACTOR,
pluginId,
action,
detail ?? null,
])
/** 取最近一条声明快照的摘要(回读基准)。 */
async function latestSnapshot(pool, pluginId) {
const { rows } = await pool.query(
`SELECT ts, detail FROM plugin_data_audit WHERE plugin_id = $1 AND action = 'decl_snapshot' ORDER BY ts DESC LIMIT 1`,
[pluginId],
)
if (rows.length === 0) return null
try {
return { ts: Number(rows[0].ts), digest: JSON.parse(String(rows[0].detail)).declDigest ?? null }
} catch {
return { ts: Number(rows[0].ts), digest: null }
}
}
const pool = openControlPool(DB_URL)
try {
if (CMD === 'decl') {
const decl = await loadDecl()
const byMod = MODULES.map((m) => [m, decl.tables.filter((t) => moduleOf(t.name)?.key === m.key)])
console.log(`\n── 模块段覆盖(⛔ 不建"无模块段"的表)`)
for (const [m, ts] of byMod) {
const names = ts.map((t) => `${t.name} → ${t.physicalName}`)
if (names.length === 0) bad(`${m.key} ${m.label}(${m.seg}_)零张表`)
else ok(`${m.key} ${m.label}(${m.seg}_)${ts.length} 张:${names.join(' | ')}`)
}
const orphan = decl.tables.filter((t) => moduleOf(t.name) === null)
if (orphan.length > 0) bad(`无模块段的表:${orphan.map((t) => t.name).join(', ')}`)
console.log(`\n── 库名与长度`)
ok(`库名 ${decl.dbName}(长度 ${decl.dbName.length} ≤ 48)`)
info(`库名正则:${/^dshs_pl_[a-z][a-z0-9_]{0,40}$/.test(decl.dbName) ? '匹配 ^dshs_pl_[a-z][a-z0-9_]{0,40}$' : '⛔ 不匹配'}`)
console.log(`\n── 声明结构摘要 declDigest=${declDigestOf(decl)}`)
} else if (CMD === 'plan') {
const decl = await loadDecl()
const current = await readCurrentSchema(DB_URL, decl.dbName)
const plan = diffDecl(decl, current, { includeCreateDatabase: true })
console.log(`\n── 预演(⛔ 本步不动任何库)`)
ok(`库存在=${current.exists} | 计划项=${plan.items.length} | 禁止项=${plan.forbidden.length}`)
info(`planHash=${plan.planHash}`)
for (const it of plan.items.slice(0, 12)) info(`计划:${it.kind} ${it.target}`)
if (plan.items.length > 12) info(`… 余 ${plan.items.length - 12} 项`)
for (const f of plan.forbidden) bad(`禁止项:${f.kind} ${f.target}`)
if (flag('ledger')) {
await ledgerUpsert(pool, { pluginId: PKG_NAME, dbName: decl.dbName, state: current.exists ? 'created' : 'pending', schemaVersion: current.exists ? current.schemaVersion : null, planHash: plan.planHash })
await audit(pool, PKG_NAME, 'datastore_plan', JSON.stringify({ dbName: decl.dbName, planHash: plan.planHash, items: plan.items.length }))
ok('台账已写(state=pending/created + plan_hash;COALESCE 语义)')
}
} else if (CMD === 'create') {
const decl = await loadDecl()
const prev = await ledgerGet(pool, PKG_NAME)
console.log(`\n── 建库(先建库再投放 · 单语句非事务 · 幂等)`)
const r = await createPluginDatabase(pool, decl.dbName)
if (r.outcome === 'created') ok(`CREATE DATABASE ${decl.dbName} OWNER dshs ENCODING UTF8 TEMPLATE template0 ⇒ created`)
else ok(`${decl.dbName} 已存在 ⇒ exists(幂等,⛔ 不重建)`)
await ledgerUpsert(pool, { pluginId: PKG_NAME, dbName: decl.dbName, state: 'created', schemaVersion: prev?.schema_version ?? null })
await audit(pool, PKG_NAME, 'datastore_create', JSON.stringify({ dbName: decl.dbName, outcome: r.outcome }))
ok('台账已写(state=created)+ 审计一条')
} else if (CMD === 'migrate') {
const decl = await loadDecl()
const want = opt('modules', MODULES.map((m) => m.key).join(','))
.split(',')
.map((s) => s.trim())
.filter(Boolean)
console.log(`\n── 迁移(单入口 · 按模块分步 · 每模块独立事务)`)
info(`本棒执行模块:${want.join(' → ')}`)
let badModules = 0
for (const key of want) {
const m = MODULES.find((x) => x.key === key)
if (m === undefined) {
bad(`未知模块 ${key}`)
badModules++
continue
}
const sub = { ...decl, tables: decl.tables.filter((t) => moduleOf(t.name)?.key === key) }
const current = await readCurrentSchema(DB_URL, decl.dbName)
const plan = diffDecl(sub, current, { includeCreateDatabase: false })
const res = await executePlan(pool, DB_URL, decl.dbName, plan, decl.schemaVersion)
const failed = res.executed.filter((e) => !e.ok)
console.log(` ── ${m.key} ${m.label}:表 ${sub.tables.length} | 计划 ${plan.items.length} | 执行 ${res.executed.length} | 失败 ${failed.length} | ${res.elapsedMs}ms`)
for (const e of res.executed) info(`${e.ok ? '✓' : '✗'} ${e.kind} ${e.target}`)
for (const e of failed) bad(`${e.kind} ${e.target}:${e.error}`)
if (failed.length > 0 || res.failed > 0) badModules++
await audit(pool, PKG_NAME, 'datastore_migrate', JSON.stringify({ module: m.key, tables: sub.tables.map((t) => t.physicalName), failed: failed.length }))
}
const total = decl.tables.length
const current2 = await readCurrentSchema(DB_URL, decl.dbName)
console.log(` ── 迁移后现状:exists=${current2.exists} schemaVersion=${current2.schemaVersion} 表=${current2.tables.length}/${total}`)
if (badModules > 0) bad(`${badModules} 个模块有失败项 ⇒ 台账置 error`)
await ledgerUpsert(pool, {
pluginId: PKG_NAME,
dbName: decl.dbName,
state: badModules > 0 ? 'error' : 'ready',
schemaVersion: badModules > 0 ? null : current2.schemaVersion,
lastError: badModules > 0 ? 'migrate_failed' : null,
})
// ⚠️ `planHash` **不传** ⇒ 按 COALESCE 语义保留(⛔ 传 null 等于把 TOCTOU 指纹拆掉)
ok(`台账已写(state=${badModules > 0 ? 'error' : 'ready'},plan_hash 保留)`)
} else if (CMD === 'snapshot') {
const decl = await loadDecl()
const digest = declDigestOf(decl)
await audit(pool, PKG_NAME, 'decl_snapshot', JSON.stringify({ declDigest: digest, schemaVersion: decl.schemaVersion, tables: decl.tables.length }))
console.log(`\n── 声明快照已登记(回读的基准)`)
ok(`declDigest=${digest}`)
ok(`表=${decl.tables.length} | 列=${decl.tables.reduce((n, t) => n + t.columns.length, 0)} | schemaVersion=${decl.schemaVersion}`)
} else if (CMD === 'readback') {
const decl = await loadDecl()
const now = declDigestOf(decl)
const snap = await latestSnapshot(pool, PKG_NAME)
console.log(`\n── 回读校验(yaml 现算 × 已登记快照)`)
if (snap === null) die(3, ' 未取到:审计里没有 decl_snapshot 记录(先跑 snapshot)')
info(`快照 ts=${new Date(snap.ts).toISOString()} digest=${snap.digest}`)
info(`现读 digest=${now}`)
if (snap.digest === now) ok('一致 ⇒ 声明与已登记的快照逐字相同(回读这一步在起作用)')
else bad('不一致 ⇒ 声明被改过却没回读登记(负控即命中此条)')
} else if (CMD === 'reconcile') {
const { rows } = await pool.query('SELECT db_name FROM plugin_datastores ORDER BY db_name ASC')
const r = await reconcilePluginDatabases(pool, rows.map((x) => x.db_name))
console.log(`\n── 对账(pg_database 的 dshs_pl_* × 控制面台账)`)
ok(`pg_database=${r.databases.length} 台账=${r.ledgered.length}`)
info(`库侧:${r.databases.join(', ')}`)
info(`台账:${r.ledgered.join(', ')}`)
if (r.orphans.length === 0 && r.missing.length === 0) ok('两集合一致(差集双空)')
if (r.orphans.length > 0) bad(`orphan(库里有、台账没有 ⇒ 绕过平台建的):${r.orphans.join(', ')}`)
if (r.missing.length > 0) bad(`missing(台账有、库里没有):${r.missing.join(', ')}`)
if (!r.databases.includes('dshs_pl_ai1net')) bad('本棒的库 dshs_pl_ai1net 不在 pg_database 里')
else ok('本棒的库 dshs_pl_ai1net 在 pg_database 里')
} else if (CMD === 'rollback-probe') {
const key = opt('module', 'M1')
const m = MODULES.find((x) => x.key === key)
if (m === undefined) die(1, `❌ 未知模块 ${key}`)
const decl = await loadDecl()
const tables = decl.tables.filter((t) => moduleOf(t.name)?.key === key)
console.log(`\n── 单模块回滚演练(${m.key} ${m.label})—— **事务内**,结束后 ROLLBACK,库内零残留`)
info(`本模块表:${tables.map((t) => t.physicalName).join(', ')}`)
const others = decl.tables.filter((t) => moduleOf(t.name)?.key !== key)
info(`其他模块表 ${others.length} 张(演练必须不碰它们)`)
const client = new PgClient({ connectionString: pluginDbUrlOf(DB_URL, decl.dbName) })
await client.connect()
try {
const before = await client.query(`SELECT tablename FROM pg_tables WHERE schemaname='public' ORDER BY tablename`)
const target = tables[0].physicalName
await client.query('BEGIN')
await client.query(`ALTER TABLE ${target} RENAME TO ${target}__rbprobe`)
const mid = await client.query(`SELECT count(*)::int AS n FROM pg_tables WHERE schemaname='public' AND tablename = $1`, [`${target}__rbprobe`])
ok(`事务内改名生效:${target} ⇒ ${target}__rbprobe(命中 ${mid.rows[0].n} 行 pg_tables)`)
await client.query('ROLLBACK')
const after = await client.query(`SELECT tablename FROM pg_tables WHERE schemaname='public' ORDER BY tablename`)
const same = JSON.stringify(before.rows) === JSON.stringify(after.rows)
if (same) ok(`ROLLBACK 后表集合与演练前逐字相同(${after.rows.length} 张)⇒ 单模块 DDL 可独立回滚且零残留`)
else bad('回滚后表集合与演练前不同 ⇒ 库内可能留残留')
} catch (err) {
await client.query('ROLLBACK').catch(() => undefined)
bad(`演练失败:${String(err.message)}`)
} finally {
await client.end().catch(() => undefined)
}
} else {
die(2, '用法:decl | plan | create | migrate | snapshot | readback | reconcile | rollback-probe')
}
} finally {
await pool.end().catch(() => undefined)
}
if (fails > 0) {
console.log(`\n❌ 失败 ${fails} 项`)
process.exit(1)
}
console.log('\n✅ 通过')
process.exit(0)