feat: v0.6.1 — AriaEngine 可选自研 KVStore 后端(storageBackend: 'kv')+ KVStore APPEND 日志类型 + 10 个 aria+kv 集成测试 + 文档全量同步
This commit is contained in:
@@ -0,0 +1,291 @@
|
||||
/**
|
||||
* KVStoreBackend + AriaEngine(kv 后端) 集成测试
|
||||
*
|
||||
* 覆盖:
|
||||
* 1. KVStoreBackend 单元:读写/追加/批量/删除/clear/close
|
||||
* 2. AriaEngine + kv 后端:CRUD 持久化(跨实例恢复)
|
||||
* 3. WAL 分片追加(backend.append 走 KVStore APPEND)
|
||||
* 4. 崩溃恢复:kv 后端下 WAL 重放
|
||||
* 5. 二级索引跨重启
|
||||
* 6. 页面化存储(SSTable → KVStore key)
|
||||
* 7. writeMany 原子性(aria 已不依赖,验证接口)
|
||||
*/
|
||||
import { KVStoreBackend } from '../../src/engine/aria/store/kvstore_backend';
|
||||
import { SharedMemoryBackend } from '../../src/engine/kvstore/shared_memory_medium';
|
||||
import { AriaEngine } from '../../src/engine/aria/index';
|
||||
import { createSchema } from '../../src/table/schema';
|
||||
|
||||
let counter = 0;
|
||||
function uniqueDB(): string {
|
||||
return `kvbe-${Date.now()}-${++counter}-${Math.random().toString(36).slice(2, 6)}`;
|
||||
}
|
||||
|
||||
const enc = (s: string) => new TextEncoder().encode(s).buffer as ArrayBuffer;
|
||||
const dec = (b: ArrayBuffer | null) => (b ? new TextDecoder().decode(b) : null);
|
||||
|
||||
beforeEach(() => {
|
||||
SharedMemoryBackend.clearRegistry();
|
||||
});
|
||||
|
||||
describe('KVStoreBackend — 单元', () => {
|
||||
it('读写/追加/批量/删除/clear 全接口', async () => {
|
||||
const backend = new KVStoreBackend(new SharedMemoryBackend());
|
||||
await backend.open(uniqueDB());
|
||||
await backend.write('k1', enc('V1'));
|
||||
expect(dec(await backend.read('k1'))).toBe('V1');
|
||||
expect(await backend.exists('k1')).toBe(true);
|
||||
|
||||
// 追加
|
||||
await backend.append('wal', enc('A'));
|
||||
await backend.append('wal', enc('B'));
|
||||
expect(dec(await backend.read('wal'))).toBe('AB');
|
||||
|
||||
// 批量原子
|
||||
await backend.writeMany({ a: enc('1'), b: enc('2') });
|
||||
expect(dec(await backend.read('b'))).toBe('2');
|
||||
await backend.deleteMany(['a']);
|
||||
expect(await backend.exists('a')).toBe(false);
|
||||
|
||||
// listKeys / delete / clear
|
||||
const keys = await backend.listKeys();
|
||||
expect(keys).toContain('k1');
|
||||
expect(keys).toContain('wal');
|
||||
await backend.delete('k1');
|
||||
expect(await backend.exists('k1')).toBe(false);
|
||||
await backend.clear();
|
||||
expect(await backend.listKeys()).toEqual([]);
|
||||
await backend.close();
|
||||
expect(backend.isOpen()).toBe(false);
|
||||
});
|
||||
|
||||
it('跨实例持久化(SharedMemory 语义)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const b1 = new KVStoreBackend(new SharedMemoryBackend());
|
||||
await b1.open(dbName);
|
||||
await b1.write('k', enc('PERSIST'));
|
||||
await b1.append('wal', enc('LOG'));
|
||||
await b1.close();
|
||||
|
||||
const b2 = new KVStoreBackend(new SharedMemoryBackend());
|
||||
await b2.open(dbName);
|
||||
expect(dec(await b2.read('k'))).toBe('PERSIST');
|
||||
expect(dec(await b2.read('wal'))).toBe('LOG');
|
||||
await b2.close();
|
||||
});
|
||||
});
|
||||
|
||||
describe('AriaEngine + kv 后端(storageBackend: kv)', () => {
|
||||
function createEngine(): AriaEngine {
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'kv',
|
||||
checkpointInterval: 100000,
|
||||
walSyncMode: 'full',
|
||||
memtableSizeThreshold: 64 * 1024 * 1024,
|
||||
});
|
||||
return engine;
|
||||
}
|
||||
|
||||
it('CRUD 持久化:写 → close → 重开数据完整(KVStore 快照/日志)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const e1 = createEngine();
|
||||
await e1.open(dbName, 1);
|
||||
await e1.createTable(createSchema('users', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
name: { type: 'string', required: true },
|
||||
age: { type: 'number', default: 0 },
|
||||
}));
|
||||
await e1.insert('users', [
|
||||
{ id: '1', name: 'Alice', age: 30 },
|
||||
{ id: '2', name: 'Bob', age: 25 },
|
||||
]);
|
||||
await e1.close();
|
||||
|
||||
const e2 = createEngine();
|
||||
await e2.open(dbName, 1);
|
||||
expect(await e2.count('users')).toBe(2);
|
||||
const rows = await e2.find('users', { table: 'users', where: { name: 'Alice' } });
|
||||
expect(rows[0].age).toBe(30);
|
||||
await e2.close();
|
||||
});
|
||||
|
||||
it('WAL 分片经 KVStore APPEND 追加:崩溃恢复重放', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const e1 = createEngine();
|
||||
await e1.open(dbName, 1);
|
||||
await e1.createTable(createSchema('logs', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
msg: { type: 'string' },
|
||||
}));
|
||||
// 不 flush 不 checkpoint:数据全在 WAL(KVStore APPEND 记录)
|
||||
for (let i = 0; i < 50; i++) {
|
||||
await e1.insert('logs', [{ id: `l-${i}`, msg: `msg-${i}` }]);
|
||||
}
|
||||
// 模拟崩溃:直接断开(不 close)
|
||||
await (e1 as any).backend.close();
|
||||
(e1 as any).opened = false;
|
||||
|
||||
const e2 = createEngine();
|
||||
await e2.open(dbName, 1);
|
||||
expect(await e2.count('logs')).toBe(50);
|
||||
await e2.close();
|
||||
});
|
||||
|
||||
it('checkpoint 后崩溃:快照 + WAL 混合恢复', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const e1 = createEngine();
|
||||
await e1.open(dbName, 1);
|
||||
await e1.createTable(createSchema('t', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
v: { type: 'number' },
|
||||
}));
|
||||
await e1.insert('t', [{ id: 'a', v: 1 }]);
|
||||
// 完整 checkpoint:flush LSM(a 落盘为 SSTable)+ 截断 WAL
|
||||
await (e1 as any).checkpointManager.checkpoint();
|
||||
await e1.insert('t', [{ id: 'b', v: 2 }]); // b 只在 WAL
|
||||
await (e1 as any).backend.close();
|
||||
(e1 as any).opened = false;
|
||||
|
||||
const e2 = createEngine();
|
||||
await e2.open(dbName, 1);
|
||||
expect(await e2.count('t')).toBe(2);
|
||||
await e2.close();
|
||||
});
|
||||
|
||||
it('二级索引跨重启恢复', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const e1 = createEngine();
|
||||
await e1.open(dbName, 1);
|
||||
await e1.createTable(createSchema('items', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
tag: { type: 'string', index: true },
|
||||
}));
|
||||
await e1.insert('items', [
|
||||
{ id: '1', tag: 'a' }, { id: '2', tag: 'b' }, { id: '3', tag: 'a' },
|
||||
]);
|
||||
await e1.close();
|
||||
|
||||
const e2 = createEngine();
|
||||
await e2.open(dbName, 1);
|
||||
const rows = await e2.find('items', { table: 'items', where: { tag: 'a' } });
|
||||
expect(rows).toHaveLength(2);
|
||||
await e2.close();
|
||||
});
|
||||
|
||||
it('页面化存储:SSTable 页面经 KVStore 读写', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const e1 = createEngine();
|
||||
await e1.open(dbName, 1);
|
||||
await e1.createTable(createSchema('docs', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
body: { type: 'string' },
|
||||
}));
|
||||
const rows = [] as Record<string, unknown>[];
|
||||
for (let i = 0; i < 30; i++) {
|
||||
rows.push({ id: `d-${i}`, body: '页面化内容'.repeat(40) });
|
||||
}
|
||||
await e1.insert('docs', rows);
|
||||
await (e1 as any).lsm.flush();
|
||||
// kv 后端:SSTable 整 value 存储(KVStore 日志型天然适合大块,无需拆 4KB 页面)
|
||||
const kv = (e1 as any).backend.getKV();
|
||||
const keys = await kv.listKeys();
|
||||
expect(keys.some((k: string) => k.startsWith('sst_'))).toBe(true);
|
||||
await e1.close();
|
||||
|
||||
const e2 = createEngine();
|
||||
await e2.open(dbName, 1);
|
||||
expect(await e2.count('docs')).toBe(30);
|
||||
await e2.close();
|
||||
});
|
||||
|
||||
it('大数据量 + 崩溃模拟:KV 后端零丢失', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const e1 = createEngine();
|
||||
await e1.open(dbName, 1);
|
||||
await e1.createTable(createSchema('big', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
v: { type: 'number' },
|
||||
}));
|
||||
for (let batch = 0; batch < 10; batch++) {
|
||||
const rows = [] as Record<string, unknown>[];
|
||||
for (let i = 0; i < 100; i++) rows.push({ id: `k${batch}-${i}`, v: i });
|
||||
await e1.insert('big', rows);
|
||||
if ((batch + 1) % 3 === 0) await (e1 as any).kvCheckpoint?.();
|
||||
}
|
||||
await (e1 as any).backend.close();
|
||||
(e1 as any).opened = false;
|
||||
|
||||
const e2 = createEngine();
|
||||
await e2.open(dbName, 1);
|
||||
expect(await e2.count('big')).toBe(1000);
|
||||
await e2.close();
|
||||
}, 60000);
|
||||
});
|
||||
|
||||
// ===================================================================
|
||||
// 高层 API:MetonaSqlark + diskEngine 'kv'
|
||||
// ===================================================================
|
||||
describe('MetonaSqlark + diskEngine: kv(高层 API)', () => {
|
||||
it('创建 → 建表 → CRUD → SQL → 持久化重开', async () => {
|
||||
const { MetonaSqlark } = require('../../src/core');
|
||||
const dbName = uniqueDB();
|
||||
const db = new MetonaSqlark({
|
||||
name: dbName,
|
||||
mode: 'aria',
|
||||
diskEngine: 'kv',
|
||||
aria: { walSyncMode: 'full', checkpointInterval: 100000, memtableSizeThreshold: 64 * 1024 * 1024 },
|
||||
});
|
||||
await db.init();
|
||||
|
||||
await db.defineTable('users', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
name: { type: 'string', required: true },
|
||||
age: { type: 'number', default: 0 },
|
||||
});
|
||||
await db.query("INSERT INTO users VALUES ('1', 'Alice', 30)");
|
||||
await db.query("INSERT INTO users VALUES ('2', 'Bob', 25)");
|
||||
const rows = await db.query("SELECT * FROM users WHERE age > 26") as Record<string, unknown>[];
|
||||
expect(rows).toHaveLength(1);
|
||||
expect(rows[0].name).toBe('Alice');
|
||||
await db.close();
|
||||
|
||||
// 重开:数据完整
|
||||
const db2 = new MetonaSqlark({
|
||||
name: dbName,
|
||||
mode: 'aria',
|
||||
diskEngine: 'kv',
|
||||
aria: { walSyncMode: 'full', checkpointInterval: 100000, memtableSizeThreshold: 64 * 1024 * 1024 },
|
||||
});
|
||||
await db2.init();
|
||||
expect(await db2.table('users').count()).toBe(2);
|
||||
await db2.close();
|
||||
});
|
||||
|
||||
it('事务回滚 + 外键级联在 kv 后端工作', async () => {
|
||||
const { MetonaSqlark } = require('../../src/core');
|
||||
const db = new MetonaSqlark({
|
||||
name: uniqueDB(), mode: 'aria', diskEngine: 'kv',
|
||||
aria: { walSyncMode: 'full', checkpointInterval: 100000 },
|
||||
});
|
||||
await db.init();
|
||||
|
||||
await db.defineTable('users', { id: { type: 'string', primaryKey: true } });
|
||||
await db.defineTable('orders', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
user_id: { type: 'string', references: 'users.id', onDelete: 'CASCADE' },
|
||||
});
|
||||
|
||||
// 事务回滚
|
||||
await expect(db.transaction(async (trx) => {
|
||||
await trx.table('users').insert({ id: '1' });
|
||||
throw new Error('boom');
|
||||
})).rejects.toThrow('boom');
|
||||
expect(await db.table('users').count()).toBe(0);
|
||||
|
||||
// 级联删除
|
||||
await db.table('users').insert({ id: '1' });
|
||||
await db.table('orders').insert({ id: 'o1', user_id: '1' });
|
||||
await db.table('users').delete().where({ id: '1' }).execute();
|
||||
expect(await db.table('orders').count()).toBe(0);
|
||||
await db.close();
|
||||
});
|
||||
});
|
||||
@@ -366,3 +366,103 @@ describe('KVStore — 编解码单元', () => {
|
||||
expect(decodeSnapshot(bytes)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
// ===================================================================
|
||||
// APPEND 追加写入(v0.6.1)
|
||||
// ===================================================================
|
||||
describe('KVStore — APPEND 追加写入', () => {
|
||||
it('appendValue 拼接 + 恢复完整(跨实例)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const medium = new SharedMemoryBackend();
|
||||
const kv = new KVStore(medium, 0);
|
||||
await kv.open(dbName);
|
||||
await kv.appendValue('log', enc('AAA'));
|
||||
await kv.appendValue('log', enc('BBB'));
|
||||
await kv.appendValue('log', enc('CCC'));
|
||||
expect(dec(await kv.get('log'))).toBe('AAABBBCCC');
|
||||
await kv.close();
|
||||
|
||||
const kv2 = new KVStore(medium, 0);
|
||||
await kv2.open(dbName);
|
||||
expect(dec(await kv2.get('log'))).toBe('AAABBBCCC');
|
||||
await kv2.close();
|
||||
});
|
||||
|
||||
it('append 与 put/delete 混合 + checkpoint 后恢复', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const medium = new SharedMemoryBackend();
|
||||
const kv = new KVStore(medium, 0);
|
||||
await kv.open(dbName);
|
||||
await kv.appendValue('log', enc('A'));
|
||||
await kv.put('meta', enc('M'));
|
||||
await kv.appendValue('log', enc('B'));
|
||||
await kv.checkpoint();
|
||||
await kv.appendValue('log', enc('C'));
|
||||
await kv.put('meta2', enc('M2'));
|
||||
await kv.close();
|
||||
|
||||
const kv2 = new KVStore(medium, 0);
|
||||
await kv2.open(dbName);
|
||||
expect(dec(await kv2.get('log'))).toBe('ABC');
|
||||
expect(dec(await kv2.get('meta'))).toBe('M');
|
||||
expect(dec(await kv2.get('meta2'))).toBe('M2');
|
||||
await kv2.close();
|
||||
});
|
||||
|
||||
it('append 覆盖语义:先 put 后 append 拼接', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const medium = new SharedMemoryBackend();
|
||||
const kv = new KVStore(medium, 0);
|
||||
await kv.open(dbName);
|
||||
await kv.put('f', enc('HEAD-'));
|
||||
await kv.appendValue('f', enc('BODY'));
|
||||
expect(dec(await kv.get('f'))).toBe('HEAD-BODY');
|
||||
await kv.close();
|
||||
|
||||
const kv2 = new KVStore(medium, 0);
|
||||
await kv2.open(dbName);
|
||||
expect(dec(await kv2.get('f'))).toBe('HEAD-BODY');
|
||||
await kv2.close();
|
||||
});
|
||||
|
||||
it('append 后 delete → 删除生效(索引与恢复一致)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const medium = new SharedMemoryBackend();
|
||||
const kv = new KVStore(medium, 0);
|
||||
await kv.open(dbName);
|
||||
await kv.appendValue('log', enc('X'));
|
||||
await kv.delete('log');
|
||||
expect(await kv.exists('log')).toBe(false);
|
||||
await kv.close();
|
||||
|
||||
const kv2 = new KVStore(medium, 0);
|
||||
await kv2.open(dbName);
|
||||
expect(await kv2.exists('log')).toBe(false);
|
||||
await kv2.close();
|
||||
});
|
||||
|
||||
it('日志损坏截断后已确认的 append 数据保留', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const medium = new SharedMemoryBackend();
|
||||
const kv = new KVStore(medium, 0);
|
||||
await kv.open(dbName);
|
||||
await kv.appendValue('log', enc('KEEP'));
|
||||
await kv.close();
|
||||
|
||||
// 追加损坏记录
|
||||
const disk = new SharedMemoryBackend();
|
||||
await disk.open(dbName);
|
||||
const log = await disk.read('__kv_log');
|
||||
const bad = new Uint8Array(encodeLogRecord(999, { x: enc('X') }));
|
||||
bad[15] ^= 0xff;
|
||||
const combined = new Uint8Array((log as ArrayBuffer).byteLength + bad.byteLength);
|
||||
combined.set(new Uint8Array(log as ArrayBuffer), 0);
|
||||
combined.set(bad, (log as ArrayBuffer).byteLength);
|
||||
await disk.write('__kv_log', combined.buffer as ArrayBuffer);
|
||||
|
||||
const kv2 = new KVStore(medium, 0);
|
||||
await kv2.open(dbName);
|
||||
expect(dec(await kv2.get('log'))).toBe('KEEP');
|
||||
await kv2.close();
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user