#!/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 * node p41-datastore-exec.mjs plan --dir [--ledger] * node p41-datastore-exec.mjs create --dir * node p41-datastore-exec.mjs migrate --dir --modules M1,M2,M3,M4,M5 * node p41-datastore-exec.mjs snapshot --dir # 登记声明快照(回读的基准) * node p41-datastore-exec.mjs readback --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)