Files
MetonaSqlark/src/core.ts
T
thzxx cbe407eb49
CI / test (22.x) (push) Successful in 17m24s
CI / e2e (push) Successful in 10m39s
CI / test (18.x) (push) Successful in 19m5s
CI / test (20.x) (push) Successful in 18m4s
CI / test (24.x) (push) Successful in 23m19s
fix: v0.7.1 API 修复与防御统一 — static create / 未 open 防护 / 原型污染 / lint 清零
- MetonaSqlark.create 静态工厂(README 示例在 ESM/Node 下此前 TypeError),
  独立 create 函数委托静态实现
- close 未初始化防御(engine undefined 不再崩溃)
- AriaEngine hasTable/getTableNames/getTableSchema 统一 ensureOpen
- 事务回滚失败不掩盖原始错误
- __proto__ 列名防护:Executor 列映射 Object.create(null) + schema 校验拒绝
- lint 清零(移除 5 处未使用导入)

测试 1147 → 1155(73 套件);行覆盖率 89.8%;版本 0.7.1
2026-08-13 11:14:42 +08:00

615 lines
20 KiB
TypeScript
Raw 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.
/**
* metona-sqlark Core — 数据库主类
* @module core
*
* 管理数据库生命周期、引擎调度、表操作、SQL 查询、事务和插件。
*/
import type { IStorageEngine } from './engine/interface';
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.7.1: 静态工厂(与 connect/disconnect 同一入口风格)。
* 此前 create 仅存在于 api 对象 / window 挂载 —— README/站点示例的
* `MetonaSqlark.create({...})` 在 ESM/Node 下是 undefinedTypeError)。
*/
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;
this.emit(msg.table ?? '', { type: 'external', table: msg.table ?? '' });
// Hybrid 引擎:从磁盘重载内存,保证读到其他标签页的最新数据
if (this.engine instanceof HybridEngine) {
(this.engine as HybridEngine).reloadMemoryFromDisk().catch(() => {
// 重载失败不影响主流程(下次读可能短暂过期)
});
}
};
}
}
// ---- 初始化 ----
/** 初始化数据库(创建引擎、打开连接) */
async init(): Promise<void> {
// 创建引擎
this.engine = this.createEngine();
// 打开连接
await this.engine.open(this.name, this.version);
// 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);
// 注册插件
if (this.config.plugins) {
for (const plugin of this.config.plugins) {
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,
(tableName) => this.broadcastChange(tableName),
(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.3.2: 写语句广播表变更(多标签页同步)
const table = this.writeStatementTable(stmt);
if (table) this.broadcastChange(table);
}
} 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 ----
/**
* 流式查询:逐行回调,不一次性物化全部结果(大表友好)。
* 支持简单 SELECTWHERE/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);
* });
* ```
*/
async queryStream<T extends Record<string, unknown> = Record<string, unknown>>(
sql: string,
onRow: (row: T) => void,
): Promise<number> {
this.ensureReady();
const stmt = parseAll(sql)[0];
if (!stmt || stmt.type !== 'SELECT') {
throw new DatabaseError('queryStream only supports SELECT statements', 'NOT_SUPPORTED');
}
const select = stmt as import('./query/ast').SelectStatement;
// 不可流式场景:JOIN / GROUP BY / HAVING / DISTINCT / 聚合 / UNION / 关联子查询 / ORDER BY
const aggregate = select.columns.some((c) => /^(COUNT|SUM|AVG|MIN|MAX)\(/i.test(c));
const streamable = !select.joins && !select.groupBy && !select.having && !select.distinct
&& !aggregate && !(select.orderBy && select.orderBy.length > 0)
&& !(select.where && select.where['$exists'] !== undefined);
if (streamable && typeof this.engine.findStream === 'function') {
// 用户回调为 async(返回 Promise)时引擎同步扫描无法 await → 回退物化
const isAsync = (onRow as { constructor?: { name?: string } }).constructor?.name === 'AsyncFunction';
if (!isAsync) {
const where = this.normalizeWhereForStream(select);
const plainCols = select.columns.filter((c) => !/\s+AS\s+\w+$/i.test(c));
return this.engine.findStream(select.from, {
table: select.from,
columns: plainCols.length > 0 && plainCols[0] !== '*' ? plainCols : ['*'],
where: where && Object.keys(where).length > 0 ? where : undefined,
limit: select.limit,
offset: select.offset,
}, onRow as (row: Record<string, unknown>) => void);
}
}
// 回退:物化后逐行回调
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;
}
/** 流式查询用:剥离主表别名前缀(复用 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: 在线备份 — 导出全库一致性快照。
* Aria 引擎走引擎级 backup()(MVCC 一致性视图);其余引擎回退 exportAll()。
*/
async backup(): Promise<Record<string, Record<string, unknown>[]>> {
this.ensureReady();
if (typeof this.engine.backup === 'function') {
return this.engine.backup();
}
return this.exportAll();
}
// ---- 发布订阅 ----
private listeners: Map<string, Set<(data: unknown) => void>> = new Map();
/** 订阅表变更 */
subscribe(tableName: string, callback: (event: { type: string; row?: unknown; table?: string }) => 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);
}
/** 触发变更事件 */
emit(tableName: string, event: { type: string; row?: unknown; table?: string }): void {
const key = `change:${tableName}`;
this.listeners.get(key)?.forEach((cb) => cb(event));
}
// ---- 多标签页同步(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': {
const rows: Record<string, unknown>[] = (stmt.values ?? []).map((vals: unknown[]) => {
const row: Record<string, unknown> = {};
const cols = stmt.columns ?? [];
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;
}
/** 获取底层引擎 */
getEngine(): IStorageEngine {
return 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);
}
}
}