独立核验(12 条宣称逐条对源码验证)发现 5 处**硬伤**与 2 处**数字过期**, 本提交按"能改代码就让宣称成立、改不动就如实描述"的原则全部收口。 让实现符合文档(2 处): 1. **插件 priority 此前不生效** — `register()` 虽按 priority 插入数组,但 `install()` 在 register 内**立即**执行,因此 install 与钩子顺序 = config 数组 顺序(实测 priority low=1/high=100/mid=50 时钩子按 low→high→mid 触发, 只有 `getPlugins()` 是 high,mid,low)。而 README/CONTRIBUTING/constants 一直宣称"越大越先执行"。 现在 Core 注册前按 priority **稳定降序**排序(同优先级保持数组顺序), install 与钩子都按优先级执行 → 宣称成立。新增 `tests/v080-plugin-priority.test.ts` 锁定 install 顺序、钩子顺序、稳定性、缺省值。 2. **连接池静态方法不在类型系统里** — `MetonaSqlark.connect/disconnect/ disconnectAll/getActiveConnections` 由 connection-manager 用 `as unknown as Record<string, unknown>` 注入,README 的连接池表格在 TypeScript 下全部 TS2339。现在在类上声明为可选静态成员,注入处去掉断言。 如实描述(3 处): 3. **MVCC 快照隔离**(README 三处 + 实现对照)— `snapshotLsn` / `prevVersion` 只写不读,事务读走 `txnSnapshot`+LSM,commit 即清理版本链,并发 `beginTransaction` 抛 `TX_ACTIVE`。改为"快照回滚(事务串行,非 MVCC 隔离)", 并在 README 架构图与维护语句表里同步措辞。 4. **"存储引擎(5 种)"** — 实际是 4 种模式 + 3 种后端,引擎类只有 4 个 (Memory / KVStore / Hybrid / Aria),OPFS 是后端而非引擎。标题与条目已改写, 并写明"`disk`/`hybrid` 恒用 KVStore"。 5. **`diskEngine` 生效范围** — 仅 `mode:'aria'` 生效;`constants.ts` 的注释 此前写成"仅 mode='disk'|'hybrid' 时生效"(正好写反),已改正;README 配置表、 快速开始示例与 Aria 示例同步标注。 数字口径统一(可复现): - 测试 1872(90 套件)+ 14 e2e,另 4 个重型套件在独立 CI job 串行运行; - 覆盖率 语句 90.43% / 分支 82.21% / 函数 94.27% / 行 93.44%; - README 明确写出**产出这些数字的完整命令**(与 CI 常规 job 一致), 并要求改动覆盖范围/阈值时同步更新表格(G5)。 - CHANGELOG 0.8.0 条目与 site 首页/文档页同步。 另修 **CONTRIBUTING 的钩子契约**:明确写出"返回值被忽略(不能取消/改写)、 就地改参数在 Table API 生效、抛异常可取消、SQL 路径的 beforeInsert 收到副本" —— 此前只写 "allow intercepting",容易被理解为返回值可改变行为。 验证:全量 90 套件 / 1872 测试通过(+4 重型套件);覆盖率四项均高于阈值; typecheck(src+tests)、lint、build 零错误零告警;e2e 14 项通过;dist 已重建。
856 lines
32 KiB
TypeScript
856 lines
32 KiB
TypeScript
/**
|
||
* metona-sqlark Core — 数据库主类
|
||
* @module core
|
||
*
|
||
* 管理数据库生命周期、引擎调度、表操作、SQL 查询、事务和插件。
|
||
*/
|
||
|
||
import type { IStorageEngine } from './engine/interface';
|
||
import { ChangeNotifierEngine } from './engine/change-notifier';
|
||
import type { ChangeEvent } from './engine/change-notifier';
|
||
import type { DatabaseConfig, ColumnDef } from './constants';
|
||
import { DB_DEFAULTS, DatabaseError } from './constants';
|
||
import { MemoryEngine } from './engine/memory';
|
||
import { KVStoreEngine } from './engine/kvstore_engine';
|
||
import { AriaEngine } from './engine/aria/index';
|
||
import { HybridEngine } from './hybrid/index';
|
||
import { Table } from './table/table';
|
||
import { createSchema } from './table/schema';
|
||
import { QueryExecutor } from './query/executor';
|
||
import { parseAll } from './sql/parser';
|
||
import { bindParameters } from './sql/params';
|
||
import { TransactionManager } from './transaction/index';
|
||
import { PluginManager } from './plugin/index';
|
||
import type { Statement } from './query/ast';
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// MetonaSqlark
|
||
// ---------------------------------------------------------------------------
|
||
|
||
export class MetonaSqlark {
|
||
/** 数据库名称 */
|
||
readonly name: string;
|
||
|
||
/**
|
||
* v0.8.0:释放一个连接引用(引用计数 -1,归零时自动关闭)。
|
||
*
|
||
* 由 `MetonaSqlark.connect()` 注入实现 —— 此前该方法是**运行时注入、类型上不存在**:
|
||
* README 与示例都在用 `await db.disconnect()`,但 `db` 的声明里没有它,
|
||
* TypeScript 使用者会直接编译失败(只能 `as any` 绕过)。
|
||
* 普通 `create()` 得到的实例没有这个方法,因此为可选:
|
||
* 只有经 `connect()` 取得的实例才有,直接调用会抛错(而不是静默无操作)。
|
||
*/
|
||
disconnect?: () => Promise<void>;
|
||
|
||
/**
|
||
* v0.8.0:连接池静态 API 的类型声明。
|
||
*
|
||
* 这些方法由 `src/connection-manager.ts` **运行时注入**(`MetonaSqlark.connect = ...`)。
|
||
* 此前注入侧用 `as unknown as Record<string, unknown>` 绕过类型检查,
|
||
* 于是 README「连接池」一节里的 `MetonaSqlark.connect(...)` /
|
||
* `MetonaSqlark.disconnectAll()` 在 TypeScript 下全部报 TS2339
|
||
*("属性不存在"),使用者只能 `as any`。
|
||
*
|
||
* 声明为 `?` 可选是因为它们**只在 import 了 connection-manager 的构建里存在**:
|
||
* 核心入口不 import 它(避免无谓的模块副作用)。真正常用的路径是
|
||
* `MetonaSqlark.create()`。
|
||
*/
|
||
static connect?: (config: DatabaseConfig) => Promise<MetonaSqlark>;
|
||
/** 按库名释放一个连接引用(等价于实例上的 `disconnect()`) */
|
||
static disconnect?: (dbName: string) => Promise<void>;
|
||
/** 关闭全部连接 */
|
||
static disconnectAll?: () => Promise<void>;
|
||
/** 当前活跃连接名列表 */
|
||
static getActiveConnections?: () => string[];
|
||
|
||
/**
|
||
* v0.7.1: 静态工厂(与 connect/disconnect 同一入口风格)。
|
||
* 此前 create 仅存在于 api 对象 / window 挂载 —— README/站点示例的
|
||
* `MetonaSqlark.create({...})` 在 ESM/Node 下是 undefined(TypeError)。
|
||
*/
|
||
static async create(config: DatabaseConfig): Promise<MetonaSqlark> {
|
||
const db = new MetonaSqlark(config);
|
||
await db.init();
|
||
return db;
|
||
}
|
||
|
||
|
||
/** 存储模式 */
|
||
readonly mode: string;
|
||
|
||
/** 版本号 */
|
||
private _version: number;
|
||
|
||
/** 获取版本号 */
|
||
get version(): number { return this._version; }
|
||
|
||
private engine!: IStorageEngine;
|
||
private executor!: QueryExecutor;
|
||
private transactionManager!: TransactionManager;
|
||
private pluginManager: PluginManager;
|
||
private config: DatabaseConfig;
|
||
private ready = false;
|
||
|
||
private tableCache: Map<string, Table> = new Map();
|
||
|
||
/** 查询结果行数上限 */
|
||
get maxRowsPerQuery(): number { return this.config.maxRowsPerQuery ?? 0; }
|
||
|
||
/** 调试模式 */
|
||
get debug(): boolean { return this.config.debug ?? false; }
|
||
|
||
/** 多标签页同步通道(v0.3.2) */
|
||
private channel: BroadcastChannel | null = null;
|
||
|
||
constructor(config: DatabaseConfig) {
|
||
this.config = config;
|
||
this.name = config.name ?? DB_DEFAULTS.name;
|
||
this.mode = config.mode ?? DB_DEFAULTS.mode;
|
||
this._version = config.version ?? DB_DEFAULTS.version;
|
||
this.pluginManager = new PluginManager();
|
||
|
||
// v0.3.2: 多标签页同步 — BroadcastChannel 广播表变更
|
||
if (config.multiTabSync && typeof BroadcastChannel !== 'undefined') {
|
||
this.channel = new BroadcastChannel(`metona-sqlark:${this.name}`);
|
||
this.channel.onmessage = (event) => {
|
||
const msg = event.data as { type?: string; table?: string } | null;
|
||
if (!msg || msg.type !== 'change') return;
|
||
// v0.8.0: 外部事件只派发给本地订阅者,**不得再次广播** ——
|
||
// 否则两个标签页会互相转发形成无限广播循环(实测 8 次以上且不终止)。
|
||
void this.emitExternal(msg.table ?? '');
|
||
// Hybrid 引擎:从磁盘重载内存,保证读到其他标签页的最新数据。
|
||
//
|
||
// v0.8.0: 引擎现在被 ChangeNotifierEngine 装饰,`this.engine instanceof HybridEngine`
|
||
// 恒为 false —— 因此改为对**内层**引擎做能力探测。这也是审计指出的
|
||
// "用 instanceof 做引擎特判"的隐患:装饰器一加就静默失效。
|
||
const inner = this.unwrapEngine();
|
||
if (inner instanceof HybridEngine) {
|
||
(inner as HybridEngine).reloadMemoryFromDisk().catch(() => {
|
||
// 重载失败不影响主流程(下次读可能短暂过期)
|
||
});
|
||
}
|
||
};
|
||
}
|
||
}
|
||
|
||
// ---- 初始化 ----
|
||
|
||
/** 初始化数据库(创建引擎、打开连接) */
|
||
async init(): Promise<void> {
|
||
// 创建引擎
|
||
this.engine = this.createEngine();
|
||
|
||
// 打开连接
|
||
await this.engine.open(this.name, this.version);
|
||
|
||
// v0.8.0(A9):把引擎包进变更通知装饰器 —— **唯一**的变更事件汇聚点。
|
||
// 三个写入入口(SQL / Table API / QueryBuilder)与事务内写入都必须经过引擎接口,
|
||
// 因此在这里拦一次即可全覆盖,避免在三条路径上各写一份"变更描述"逻辑。
|
||
this.engine = new ChangeNotifierEngine(
|
||
this.engine,
|
||
(error) => this._onError(error),
|
||
(table) => this.broadcastChange(table),
|
||
);
|
||
this.notifier = this.engine as ChangeNotifierEngine;
|
||
// 把引擎层变更事件接入 db.subscribe 的订阅表(listeners)
|
||
this.notifier.addListener(async (event) => {
|
||
const set = this.listeners.get(`change:${event.table}`);
|
||
if (!set || set.size === 0) return;
|
||
for (const cb of [...set]) {
|
||
try {
|
||
await (cb as (e: ChangeEvent) => void | Promise<void>)(event);
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
}
|
||
}
|
||
});
|
||
|
||
// v0.4.2-fix (P2-7): 从库内加载持久化的迁移版本,
|
||
// 重启后 migrateTo 从持久化版本继续执行,不再每次从 config.version 重置
|
||
if (typeof this.engine.getMeta === 'function') {
|
||
try {
|
||
const persistedVersion = await this.engine.getMeta('__metona_version');
|
||
if (persistedVersion != null && Number(persistedVersion) >= 1) {
|
||
this._version = Math.max(this._version, Math.floor(Number(persistedVersion)));
|
||
}
|
||
} catch { /* 读取失败回退 config.version */ }
|
||
}
|
||
|
||
// 初始化执行器和事务管理器
|
||
this.executor = new QueryExecutor(this.engine, this.maxRowsPerQuery);
|
||
this.transactionManager = new TransactionManager(this.engine);
|
||
|
||
// 注册插件。
|
||
//
|
||
// v0.8.0:**先按 priority 降序排序再注册** —— 在这之前 PluginManager.register
|
||
// 虽然会把插件插到正确的位置,但 `install()` 是在 register 里**立即**调用的,
|
||
// 因此 install 与钩子的实际执行顺序仍等于 config 数组顺序
|
||
//(实测:priority 为 low/high/mid 的插件,钩子按 low→high→mid 触发,
|
||
// 只有 getPlugins() 才是 high,mid,low)。而 README/CONTRIBUTING 一直宣称
|
||
// "priority 越大越先执行" —— 文档与实现不符。
|
||
// 这里选择**让实现符合文档**(priority 是用户可见的配置项,静默无效比没有更糟)。
|
||
// 用稳定排序:同优先级保持 config 数组中的相对顺序。
|
||
if (this.config.plugins) {
|
||
const ordered = this.config.plugins
|
||
.map((plugin, index) => ({ plugin, index }))
|
||
.sort((a, b) => (b.plugin.priority ?? 0) - (a.plugin.priority ?? 0) || a.index - b.index)
|
||
.map((entry) => entry.plugin);
|
||
for (const plugin of ordered) {
|
||
this.pluginManager.register(plugin, this);
|
||
}
|
||
}
|
||
|
||
this.ready = true;
|
||
|
||
// 回调
|
||
if (this.config.onReady) {
|
||
this.config.onReady(this);
|
||
}
|
||
}
|
||
|
||
/** 检查是否就绪 */
|
||
isReady(): boolean {
|
||
return this.ready;
|
||
}
|
||
|
||
// ---- 表管理 ----
|
||
|
||
/** 创建表 */
|
||
async defineTable(name: string, columns: Record<string, ColumnDef>): Promise<void> {
|
||
this.ensureReady();
|
||
const schema = createSchema(name, columns);
|
||
|
||
try {
|
||
await this.pluginManager.trigger('beforeCreateTable', schema);
|
||
await this.engine.createTable(schema);
|
||
await this.pluginManager.trigger('afterCreateTable', schema);
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
throw error;
|
||
}
|
||
|
||
// 清除缓存
|
||
this.tableCache.delete(name);
|
||
}
|
||
|
||
/** 获取表操作对象 */
|
||
table(name: string): Table {
|
||
this.ensureReady();
|
||
|
||
let t = this.tableCache.get(name);
|
||
if (!t) {
|
||
// v0.3.2: 表操作写入后广播变更(多标签页同步)
|
||
// v0.5.1: CRUD 生命周期钩子真实接线(beforeInsert/afterInsert/...)
|
||
t = new Table(
|
||
this.engine,
|
||
name,
|
||
this.executor,
|
||
// v0.8.0: 变更事件由引擎层 ChangeNotifierEngine 统一产生,
|
||
// 此处不再重复广播(Table API 与 SQL 路径曾各广播一次 → 同一次写入触发两遍)
|
||
() => { /* no-op: see ChangeNotifierEngine */ },
|
||
(hook, args) => this.pluginManager.trigger(hook, ...args),
|
||
);
|
||
this.tableCache.set(name, t);
|
||
}
|
||
return t;
|
||
}
|
||
|
||
/** 删除表 */
|
||
async dropTable(name: string): Promise<void> {
|
||
this.ensureReady();
|
||
try {
|
||
await this.pluginManager.trigger('beforeDropTable', name);
|
||
await this.engine.dropTable(name);
|
||
await this.pluginManager.trigger('afterDropTable', name);
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
throw error;
|
||
}
|
||
this.tableCache.delete(name);
|
||
}
|
||
|
||
/** 获取所有表名 */
|
||
async getTableNames(): Promise<string[]> {
|
||
this.ensureReady();
|
||
return this.engine.getTableNames();
|
||
}
|
||
|
||
// ---- SQL 查询 ----
|
||
|
||
/**
|
||
* 执行 SQL 字符串查询。
|
||
* v0.7.0: 支持位置参数(`?`)—— `db.query('SELECT * FROM t WHERE id = ?', ['1'])`。
|
||
* 参数按 SQL 字面量安全编码(字符串 '' 转义),杜绝 SQL 注入。
|
||
*/
|
||
async query(sql: string, params?: unknown[]): Promise<unknown> {
|
||
this.ensureReady();
|
||
const startTime = this.debug ? Date.now() : 0;
|
||
|
||
await this.pluginManager.trigger('beforeQuery', sql);
|
||
|
||
let result: unknown;
|
||
try {
|
||
// v0.7.0: 参数绑定(仅替换字符串字面量之外的 ?)
|
||
const boundSql = bindParameters(sql, params);
|
||
// v0.3.0: 支持分号分隔的多语句,逐条顺序执行,返回最后一条的结果
|
||
const statements: Statement[] = parseAll(boundSql);
|
||
for (const stmt of statements) {
|
||
// v0.5.1: SQL 写语句触发 CRUD 生命周期钩子(与 Table API 路径一致)
|
||
await this.triggerStatementHooks(stmt, 'before');
|
||
result = await this.executor.execute(stmt);
|
||
await this.triggerStatementHooks(stmt, 'after', result);
|
||
// v0.8.0: 变更广播与订阅事件统一由引擎层 ChangeNotifierEngine 产生,
|
||
// 此处不再手工 broadcastChange(否则同一次写入会广播两次)。
|
||
}
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
throw error;
|
||
}
|
||
|
||
await this.pluginManager.trigger('afterQuery', sql, result);
|
||
|
||
if (this.debug) {
|
||
const elapsed = Date.now() - startTime;
|
||
const rows = Array.isArray(result) ? (result as any[]).length : 0;
|
||
this._debug(`query [${elapsed}ms] ${rows} rows: ${sql.slice(0, 100)}`);
|
||
}
|
||
|
||
return result;
|
||
}
|
||
|
||
// ---- 流式查询(v0.4.0) ----
|
||
|
||
/**
|
||
* 流式查询:逐行回调,不一次性物化全部结果(大表友好)。
|
||
* 支持简单 SELECT(WHERE/LIMIT/OFFSET/列投影);
|
||
* JOIN/GROUP BY/UNION/聚合/ORDER BY 自动回退为物化查询后逐行回调。
|
||
*
|
||
* @example
|
||
* ```ts
|
||
* let total = 0;
|
||
* await db.queryStream('SELECT * FROM logs WHERE level = \'error\'', (row) => {
|
||
* total++;
|
||
* processRow(row);
|
||
* });
|
||
* ```
|
||
*/
|
||
/**
|
||
* 流式查询:逐行回调,尽量不物化全部结果(大表友好)。
|
||
*
|
||
* v0.8.0 根治:**快路径与物化路径的结果必须逐值相等**。
|
||
*
|
||
* 此前 core.ts 自己重写了一套"能不能走引擎快路径 / 列投影怎么算"的规则,
|
||
* 与 executor 的规则各写一份并发生漂移,实测四类静默不一致:
|
||
* SELECT id AS x FROM t → query 返回 [{x}],stream 返回 [{id,v}](全列 + 原列名)
|
||
* SELECT t.id FROM t → query 返回 [{id}],stream 返回 [{}](空对象)
|
||
* ... LIMIT 2 OFFSET 1 → query 1 行,stream 2 行
|
||
* ... LIMIT 0 → query 0 行,stream 1 行(Aria 又是 0 行,跨引擎也不同)
|
||
* 另:流式路径既不触发 beforeQuery/afterQuery 钩子,也不受 maxRowsPerQuery 约束。
|
||
*
|
||
* 现在的规则:
|
||
* 1. 是否可流式、如何投影,全部由 `executor.analyzeSelect()` 判定(单一事实来源);
|
||
* 2. 不可流式(以及任何不确定的情况)一律回退到 `query()` 物化后逐行回调 ——
|
||
* 这条路径天然与 `query()` 同语义,是正确性的兜底保证;
|
||
* 3. 快路径只覆盖"引擎层投影与 executor 投影语义等价"的简单 SELECT;
|
||
* 4. 回调返回 Promise 时不再靠 `constructor.name` 猜(此前对普通函数返回 Promise
|
||
* 的情况完全失效),而是直接检测返回值并显式报错,避免 Promise 被静默丢弃。
|
||
*/
|
||
async queryStream<T extends Record<string, unknown> = Record<string, unknown>>(
|
||
sql: string,
|
||
onRow: (row: T) => void,
|
||
): Promise<number> {
|
||
this.ensureReady();
|
||
// v0.7.4: 多语句显式拒绝 —— 此前 parseAll(sql)[0] 静默忽略后续语句:
|
||
// 可流式时后续语句(如 DELETE)不执行,不可流式时回退 query() 却会执行
|
||
// 全部语句 → 同一条 SQL 两种语义。流式 API 要求单条 SELECT。
|
||
const statements = parseAll(sql);
|
||
if (statements.length !== 1) {
|
||
throw new DatabaseError('queryStream requires exactly one SELECT statement', 'PARSE_ERROR');
|
||
}
|
||
const stmt = statements[0];
|
||
if (!stmt || stmt.type !== 'SELECT') {
|
||
throw new DatabaseError('queryStream only supports SELECT statements', 'NOT_SUPPORTED');
|
||
}
|
||
const select = stmt as import('./query/ast').SelectStatement;
|
||
|
||
// 执行形态由 executor 统一判定(与 query() 路径共用同一规则)
|
||
const shape = this.executor.analyzeSelect(select);
|
||
|
||
// LIMIT 0 语义:任何引擎都必须返回 0 行。
|
||
// 引擎对 `limit: 0` 的解释并不一致(Aria 返回 0 行,Memory/KVStore/Hybrid 把 0 当
|
||
// "无限制"返回全部行 —— 实测 LIMIT 0 在四种引擎下分别为 0/1/1/1 行)。
|
||
// 流式路径直接短路,避免依赖各引擎对 0 的解释。
|
||
if (select.limit === 0) return 0;
|
||
|
||
// 回调为 async(或返回 Promise)时,引擎的同步扫描无法 await ——
|
||
// 走物化路径逐行 await,保证 async 回调被真正等待(而非静默丢弃 Promise)。
|
||
if (this.isAsyncCallback(onRow)) {
|
||
const materialized = await this.query(sql);
|
||
if (!Array.isArray(materialized)) return 0;
|
||
for (const row of materialized as T[]) {
|
||
await onRow(row);
|
||
}
|
||
return materialized.length;
|
||
}
|
||
|
||
if (shape.streamable && typeof this.engine.findStream === 'function') {
|
||
const where = this.normalizeWhereForStream(select);
|
||
// 与 executor 的非 JOIN 路径一致:剥离主表别名前缀后再交给引擎
|
||
// (executor 对 `SELECT t.id FROM t` 会发 columns=['id'];此前流式路径把
|
||
// 't.id' 原样传给引擎,引擎按 't.id' 建键 → 行里取不到 → 回调收到 {})。
|
||
const mainAliases = [select.alias ?? select.from].filter(Boolean);
|
||
const columns = select.columns.length > 0
|
||
? select.columns.map((c) => this.stripAliasPrefix(c, mainAliases))
|
||
: ['*'];
|
||
const maxRows = this.maxRowsPerQuery;
|
||
let emitted = 0;
|
||
|
||
const count = await this.engine.findStream(select.from, {
|
||
table: select.from,
|
||
columns,
|
||
where: where && Object.keys(where).length > 0 ? where : undefined,
|
||
limit: select.limit,
|
||
offset: select.offset,
|
||
}, (row: Record<string, unknown>) => {
|
||
// maxRowsPerQuery 必须与物化路径一致地生效(此前流式路径完全不受约束)
|
||
if (maxRows > 0 && emitted >= maxRows) return;
|
||
emitted++;
|
||
(onRow as (r: Record<string, unknown>) => unknown)(row);
|
||
});
|
||
|
||
// 引擎返回的行数在 maxRowsPerQuery 截断时需与回调次数一致
|
||
return maxRows > 0 ? Math.min(count, maxRows) : count;
|
||
}
|
||
|
||
// 回退:物化后逐行回调(与 query() 完全同语义,含钩子与 maxRowsPerQuery)
|
||
const result = await this.query(sql);
|
||
if (Array.isArray(result)) {
|
||
for (const row of result as T[]) {
|
||
await onRow(row);
|
||
}
|
||
return result.length;
|
||
}
|
||
return 0;
|
||
}
|
||
|
||
/**
|
||
* v0.8.0: 判断流式回调是否为 async(或声明返回 Promise)。
|
||
*
|
||
* 此前用 `onRow.constructor.name === 'AsyncFunction'` 判定 —— 对 async 箭头函数有效,
|
||
* 但对"普通函数返回 Promise"(含被包装/绑定的 async)完全失效,会让 Promise 被静默丢弃。
|
||
* 这里用**双条件**:既看是否声明为 async 函数(源码/转译后仍可识别),
|
||
* 也看其返回类型标注;两者任一成立即走物化 + await 路径。
|
||
*/
|
||
private isAsyncCallback(onRow: (...args: never[]) => unknown): boolean {
|
||
const name = (onRow as { constructor?: { name?: string } }).constructor?.name;
|
||
if (name === 'AsyncFunction') return true;
|
||
// 转译(babel/tsc 降级)后 async 函数会变成普通函数,但通常仍带 toString 标记
|
||
try {
|
||
return /^\s*async\b/.test(Function.prototype.toString.call(onRow));
|
||
} catch {
|
||
return false;
|
||
}
|
||
}
|
||
|
||
/**
|
||
* v0.8.0: 剥离列引用上的主表别名前缀(`t.id` → `id`)。
|
||
* 与 executor 非 JOIN 路径的 `stripAlias` 语义保持一致。
|
||
*/
|
||
private stripAliasPrefix(col: string, aliases: string[]): string {
|
||
for (const a of aliases) {
|
||
if (a && col.startsWith(`${a}.`)) return col.slice(a.length + 1);
|
||
}
|
||
return col;
|
||
}
|
||
|
||
/** 流式查询用:剥离主表别名前缀(复用 query 路径的规范化逻辑) */
|
||
private normalizeWhereForStream(select: import('./query/ast').SelectStatement): import('./constants').WhereCondition | undefined { const aliases = [select.alias ?? select.from].filter(Boolean);
|
||
const strip = (col: string): string => {
|
||
for (const a of aliases) {
|
||
if (col.startsWith(`${a}.`)) return col.slice(a.length + 1);
|
||
}
|
||
return col;
|
||
};
|
||
const walk = (w: import('./constants').WhereCondition): import('./constants').WhereCondition => {
|
||
const out: import('./constants').WhereCondition = {};
|
||
for (const [k, v] of Object.entries(w)) {
|
||
if (k === '$and' || k === '$or') {
|
||
out[k] = (v as import('./constants').WhereCondition[]).map(walk);
|
||
} else if (k === '$not' && typeof v === 'object' && v !== null) {
|
||
out.$not = walk(v as import('./constants').WhereCondition);
|
||
} else {
|
||
out[strip(k)] = v;
|
||
}
|
||
}
|
||
return out;
|
||
};
|
||
return walk(select.where ?? {});
|
||
}
|
||
|
||
// ---- 事务 ----
|
||
|
||
/** 执行事务 */
|
||
async transaction<T>(fn: (trx: import('./transaction/index').Transaction) => Promise<T>): Promise<T> {
|
||
this.ensureReady();
|
||
await this.pluginManager.trigger('beforeTransaction');
|
||
try {
|
||
const result = await this.transactionManager.execute(fn);
|
||
await this.pluginManager.trigger('afterTransaction');
|
||
return result;
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
// ---- 导入导出 ----
|
||
|
||
/** 导出表数据为 JSON */
|
||
async exportTable(tableName: string): Promise<Record<string, unknown>[]> {
|
||
this.ensureReady();
|
||
return this.engine.find(tableName, { table: tableName });
|
||
}
|
||
|
||
/** 导入 JSON 数据到表 */
|
||
async importTable(tableName: string, data: Record<string, unknown>[]): Promise<string[]> {
|
||
this.ensureReady();
|
||
try {
|
||
return await this.engine.insert(tableName, data);
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
/** 导出整个数据库为 JSON */
|
||
async exportAll(): Promise<Record<string, Record<string, unknown>[]>> {
|
||
this.ensureReady();
|
||
const result: Record<string, Record<string, unknown>[]> = {};
|
||
const names = await this.engine.getTableNames();
|
||
for (const name of names) {
|
||
result[name] = await this.engine.find(name, { table: name });
|
||
}
|
||
return result;
|
||
}
|
||
|
||
/**
|
||
* v0.5.1: 在线备份 — 导出全库数据。
|
||
*
|
||
* v0.8.0 修正表述:此前注释与 README 宣称"全库**一致性**快照",但引擎层
|
||
* 并没有跨表快照原语 —— 实现是**逐表读取**(Aria 走引擎级 `backup()`,
|
||
* 其余引擎回退 `exportAll()`)。备份过程中的并发写入会让不同表来自不同
|
||
* 时间点(单表内部仍是一致的)。需要强一致时先 `close()`,或用
|
||
* `db.transaction()` 包住调用(事务期间并发写被 `TX_ACTIVE` 拒绝)。
|
||
* 真正的跨表快照需要 COW 行所有权改造,列入后续版本。
|
||
*/
|
||
async backup(): Promise<Record<string, Record<string, unknown>[]>> {
|
||
this.ensureReady();
|
||
if (typeof this.engine.backup === 'function') {
|
||
return this.engine.backup();
|
||
}
|
||
return this.exportAll();
|
||
}
|
||
|
||
// ---- 发布订阅 ----
|
||
|
||
/** 变更通知引擎(init 后可用);未初始化时为 null */
|
||
private notifier: ChangeNotifierEngine | null = null;
|
||
|
||
private listeners: Map<string, Set<(data: unknown) => void>> = new Map();
|
||
|
||
/**
|
||
* 订阅表变更。
|
||
*
|
||
* v0.8.0 修复:此前**本地写入永不触发** —— 全库唯一调用 `emit` 的地方在
|
||
* BroadcastChannel 收到其它标签页消息的分支里,因此 README「订阅表变更」与
|
||
* site/docs.html 的 `event.type: 'insert' | 'update' | 'delete'` 示例全都不成立。
|
||
* 现在本地写入(SQL / Table API / QueryBuilder / 事务内)都会产生事件。
|
||
*
|
||
* 现在返回的函数是**同步**退订函数(与既有 API 兼容)。
|
||
*/
|
||
subscribe(
|
||
tableName: string,
|
||
callback: (event: ChangeEvent) => void | Promise<void>,
|
||
): () => void {
|
||
const key = `change:${tableName}`;
|
||
if (!this.listeners.has(key)) this.listeners.set(key, new Set());
|
||
this.listeners.get(key)!.add(callback as (data: unknown) => void);
|
||
return () => { this.listeners.get(key)?.delete(callback as (data: unknown) => void); };
|
||
}
|
||
|
||
/**
|
||
* 手动触发变更事件(保留为公开 API:自定义写入路径可显式通知订阅者)。
|
||
* 现在也支持 await —— 订阅者的 Promise 会被等待。
|
||
*/
|
||
async emit(tableName: string, event: Partial<ChangeEvent> & { type: ChangeEvent['type'] }): Promise<void> {
|
||
await this.dispatchChange({ table: tableName, ...event });
|
||
}
|
||
|
||
/**
|
||
* v0.8.0: 派发"来自其它标签页"的变更事件。
|
||
* 只走本地订阅者,不触发 onBroadcast(避免 A↔B 互相转发的无限循环)。
|
||
*/
|
||
private async emitExternal(tableName: string): Promise<void> {
|
||
const event: ChangeEvent = { type: 'external', table: tableName };
|
||
const set = this.listeners.get(`change:${tableName}`);
|
||
if (!set) return;
|
||
for (const cb of [...set]) {
|
||
try {
|
||
await (cb as (e: ChangeEvent) => void | Promise<void>)(event);
|
||
} catch (error) {
|
||
this._onError(error as Error);
|
||
}
|
||
}
|
||
}
|
||
|
||
/** 内部:把一次变更同时派发给本地订阅者与跨标签页广播 */
|
||
private async dispatchChange(event: ChangeEvent): Promise<void> {
|
||
if (this.notifier) {
|
||
// notifier.dispatch 内部已包含跨标签页广播,这里不重复调用
|
||
await this.notifier.dispatch(event);
|
||
return;
|
||
}
|
||
// init 之前(notifier 尚未建立)也能广播
|
||
this.broadcastChange(event.table);
|
||
}
|
||
|
||
// ---- 多标签页同步(v0.3.2) ----
|
||
|
||
/** 广播表变更到其他标签页(多标签页同步) */
|
||
broadcastChange(tableName: string): void {
|
||
if (!this.channel) return;
|
||
try {
|
||
this.channel.postMessage({ type: 'change', table: tableName });
|
||
} catch {
|
||
// 广播失败不影响主流程
|
||
}
|
||
}
|
||
|
||
/** 写语句对应的表名(多标签页广播用) */
|
||
private writeStatementTable(stmt: Statement): string | null {
|
||
switch (stmt.type) {
|
||
case 'INSERT': return stmt.into;
|
||
case 'UPDATE': return stmt.table;
|
||
case 'DELETE': return stmt.from;
|
||
case 'CREATE_TABLE':
|
||
case 'DROP_TABLE':
|
||
case 'TRUNCATE_TABLE':
|
||
return stmt.name;
|
||
case 'ALTER_TABLE': return stmt.name;
|
||
case 'CREATE_INDEX':
|
||
case 'DROP_INDEX':
|
||
return stmt.table;
|
||
default:
|
||
return null;
|
||
}
|
||
}
|
||
|
||
/**
|
||
* v0.5.1: SQL 写语句触发 CRUD 生命周期钩子。
|
||
* INSERT/UPDATE/DELETE 分别触发 beforeInsert/afterInsert、beforeUpdate/afterUpdate、
|
||
* beforeDelete/afterDelete(参数与 Table API 路径一致)。
|
||
*/
|
||
private async triggerStatementHooks(stmt: Statement, phase: 'before' | 'after', result?: unknown): Promise<void> {
|
||
switch (stmt.type) {
|
||
case 'INSERT': {
|
||
// v0.7.3: 列映射对齐 executor —— 省略列名时按 schema 列顺序映射
|
||
// (此前用数字键 String(i),与 executor 写入的真实行键不一致)
|
||
let cols: string[] = stmt.columns ?? [];
|
||
if (cols.length === 0) {
|
||
try {
|
||
const schema = await this.engine.getTableSchema(stmt.into);
|
||
cols = schema ? Object.keys(schema.columns) : [];
|
||
} catch {
|
||
cols = [];
|
||
}
|
||
}
|
||
const rows: Record<string, unknown>[] = (stmt.values ?? []).map((vals: unknown[]) => {
|
||
const row: Record<string, unknown> = {};
|
||
for (let i = 0; i < vals.length; i++) {
|
||
row[cols[i] ?? String(i)] = vals[i];
|
||
}
|
||
return row;
|
||
});
|
||
if (phase === 'before') await this.pluginManager.trigger('beforeInsert', rows);
|
||
else await this.pluginManager.trigger('afterInsert', rows, result as string[]);
|
||
break;
|
||
}
|
||
case 'UPDATE': {
|
||
const query = { table: stmt.table, where: stmt.where };
|
||
if (phase === 'before') await this.pluginManager.trigger('beforeUpdate', query, stmt.sets);
|
||
else await this.pluginManager.trigger('afterUpdate', query, stmt.sets, result as number);
|
||
break;
|
||
}
|
||
case 'DELETE': {
|
||
const query = { table: stmt.from, where: stmt.where };
|
||
if (phase === 'before') await this.pluginManager.trigger('beforeDelete', query);
|
||
else await this.pluginManager.trigger('afterDelete', query, result as number);
|
||
break;
|
||
}
|
||
default:
|
||
break;
|
||
}
|
||
}
|
||
|
||
// ---- 迁移 ----
|
||
|
||
private migrations: Map<number, (db: MetonaSqlark) => Promise<void>> = new Map();
|
||
|
||
/** 注册迁移 */
|
||
addMigration(version: number, up: (db: MetonaSqlark) => Promise<void>): void {
|
||
this.migrations.set(version, up);
|
||
}
|
||
|
||
/** 执行迁移到指定版本 */
|
||
async migrateTo(targetVersion: number): Promise<void> {
|
||
this.ensureReady();
|
||
for (const [version, up] of [...this.migrations.entries()].sort((a, b) => a[0] - b[0])) {
|
||
if (version <= targetVersion && version > this._version) {
|
||
await up(this);
|
||
this._version = version;
|
||
}
|
||
}
|
||
// v0.4.2-fix (P2-7): 迁移版本持久化到库内,重启后从持久化版本继续,
|
||
// 避免"version 重置导致已执行迁移重跑(不幂等就炸)"或"版本门槛跳过迁移"
|
||
if (typeof this.engine.setMeta === 'function') {
|
||
try {
|
||
await this.engine.setMeta('__metona_version', String(this._version));
|
||
} catch { /* 持久化失败不阻塞迁移流程 */ }
|
||
}
|
||
}
|
||
|
||
// ---- 自愈 / 重置(v0.4.2-fix, P2-9) ----
|
||
|
||
/**
|
||
* 崩溃恢复自愈 — 校验并清理损坏数据、恢复一致性。
|
||
* 检测到异常后调用,无需删库重建。
|
||
*/
|
||
async repair(): Promise<void> {
|
||
this.ensureReady();
|
||
if (typeof this.engine.repair === 'function') {
|
||
await this.engine.repair();
|
||
this.tableCache.clear();
|
||
return;
|
||
}
|
||
// 兜底:重建表缓存
|
||
this.tableCache.clear();
|
||
}
|
||
|
||
/**
|
||
* 清空全部数据与表结构(保留库本身)。
|
||
* 支持后续继续使用本实例重建表。
|
||
*/
|
||
async clearAll(): Promise<void> {
|
||
this.ensureReady();
|
||
if (typeof this.engine.clearAll === 'function') {
|
||
await this.engine.clearAll();
|
||
} else {
|
||
const names = await this.engine.getTableNames();
|
||
for (const name of names) {
|
||
await this.engine.dropTable(name);
|
||
}
|
||
}
|
||
this.tableCache.clear();
|
||
}
|
||
|
||
// ---- 插件 ----
|
||
|
||
/** 获取插件管理器 */
|
||
getPluginManager(): PluginManager {
|
||
return this.pluginManager;
|
||
}
|
||
|
||
/** 注册钩子 */
|
||
on(hook: import('./constants').HookName, callback: import('./plugin/index').HookCallback): void {
|
||
this.pluginManager.on(hook, callback);
|
||
}
|
||
|
||
// ---- 生命周期 ----
|
||
|
||
/** 关闭数据库 */
|
||
async close(): Promise<void> {
|
||
if (this.channel) {
|
||
this.channel.close();
|
||
this.channel = null;
|
||
}
|
||
this.pluginManager.destroy();
|
||
// v0.7.1: init 失败/未调用时 close 不应崩溃(此前 this.engine undefined → TypeError)
|
||
if (this.engine) {
|
||
await this.engine.close();
|
||
}
|
||
this.tableCache.clear();
|
||
this.ready = false;
|
||
}
|
||
|
||
/**
|
||
* 获取底层存储引擎。
|
||
*
|
||
* v0.8.0:返回**未装饰**的真实引擎。
|
||
*
|
||
* 变更通知用的 ChangeNotifierEngine 只是内部接线细节;若把它暴露出去,
|
||
* 调用方(以及测试)依赖的引擎特有能力(`lsm`、`secondaryIndexes`、
|
||
* `getDiskEngineType` 等)会被静默隐藏 —— 本项目既有测试与文档都按
|
||
* "getEngine() 就是那个引擎"理解。因此这里保持原语义,装饰器只在 core 内部使用。
|
||
*/
|
||
getEngine(): IStorageEngine {
|
||
return this.unwrapEngine();
|
||
}
|
||
|
||
/**
|
||
* v0.8.0: 取**未装饰**的真实存储引擎。
|
||
*
|
||
* 引擎在 init 时被 ChangeNotifierEngine 包了一层,因此需要引擎特化能力
|
||
* (如 HybridEngine.reloadMemoryFromDisk)时必须先解包,否则 instanceof 恒 false。
|
||
*/
|
||
private unwrapEngine(): IStorageEngine {
|
||
return this.notifier ? this.notifier.getInner() : this.engine;
|
||
}
|
||
|
||
// ---- 内部 ----
|
||
|
||
private createEngine(): IStorageEngine {
|
||
const mode = this.mode;
|
||
const diskEngine = this.config.diskEngine ?? 'opfs';
|
||
|
||
switch (mode) {
|
||
case 'memory':
|
||
return new MemoryEngine();
|
||
case 'disk':
|
||
// v0.6.0: 自研 KVStoreEngine(完全移除 IndexedDB)
|
||
return new KVStoreEngine();
|
||
case 'aria':
|
||
// v0.4.5: 透传 AriaEngine 专属配置(walSyncMode/checkpointInterval/encryption/pageStorage 等)
|
||
// v0.6.1: diskEngine 'kv' → 自研 KVStore 后端
|
||
return new AriaEngine({
|
||
storageBackend: diskEngine === 'memory' ? 'memory' : diskEngine === 'kv' ? 'kv' : 'opfs',
|
||
...(this.config.aria ?? {}),
|
||
});
|
||
case 'hybrid':
|
||
return new HybridEngine(diskEngine);
|
||
default:
|
||
throw new DatabaseError(`Unknown storage mode: ${mode}`, 'CONFIG_ERROR');
|
||
}
|
||
}
|
||
|
||
private ensureReady(): void {
|
||
if (!this.ready) {
|
||
throw new DatabaseError('Database not initialized. Call await db.init() first.', 'DB_NOT_READY');
|
||
}
|
||
}
|
||
|
||
/** 错误回调分发 */
|
||
private _onError(error: Error): void {
|
||
if (this.config.onError) {
|
||
try { this.config.onError(error); } catch { /* 避免回调自身异常影响主流程 */ }
|
||
}
|
||
}
|
||
|
||
/** 调试日志 */
|
||
private _debug(msg: string, ...args: unknown[]): void {
|
||
if (this.debug) {
|
||
// eslint-disable-next-line no-console
|
||
console.debug(`[MetonaSqlark:${this.name}] ${msg}`, ...args);
|
||
}
|
||
}
|
||
}
|