- 变更规模:新增 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/ 知识文件,按口径入库)
316 lines
17 KiB
JavaScript
316 lines
17 KiB
JavaScript
#!/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)
|