fix(P0): SSTable rangeScan 尾块漏读 — 索引块记录块尾 key 但二分按块首语义定位,endKey 落在块尾时下一块被排除导致二级索引查询丢数据(5万行丢106~771条);endBlockIdx 多扫一块 + 4 个 sstable 边界回归 + 重建二级索引大数据量回归(含崩溃恢复)。附带:checkpoint 同步 flush 二级索引 LSM(P1)、kv 后端页面化 + SharedMemoryBackend chunk 化、生产负载验证测试
This commit is contained in:
@@ -0,0 +1,293 @@
|
||||
/**
|
||||
* AriaEngine — 生产负载验证(正确性优先,规模可完成)
|
||||
*
|
||||
* 覆盖 README 宣称的核心能力在生产负载下的正确性:
|
||||
* 1. 5 万行写入(多级 Compaction)→ 完整查询 → 崩溃恢复
|
||||
* 2. 高频更新/删除(Compaction 回收墓碑)→ 重启后一致
|
||||
* 3. 大 value(100KB×50)→ 编码/压缩/恢复
|
||||
* 4. 混合操作 + 崩溃 → 已确认写入零丢失
|
||||
* 5. 大量删除(90%)+ Compaction → 重启无残留
|
||||
* 6. kv 后端 5 万行(页面化路径)
|
||||
*
|
||||
* 注:10 万级 kv 后端性能专项见 CHANGELOG v0.6.1 待办(KVStore 日志增长优化)。
|
||||
*/
|
||||
import { AriaEngine } from '../../src/engine/aria/index';
|
||||
import { createSchema } from '../../src/table/schema';
|
||||
import { SharedMemoryBackend } from '../../src/engine/kvstore/shared_memory_medium';
|
||||
import { installOPFSMock } from '../helpers/opfs-mock';
|
||||
|
||||
let counter = 0;
|
||||
function uniqueDB(): string {
|
||||
return `pload-${Date.now()}-${++counter}-${Math.random().toString(36).slice(2, 6)}`;
|
||||
}
|
||||
|
||||
const SCHEMA = () => createSchema('big', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
val: { type: 'number' },
|
||||
tag: { type: 'string', index: true },
|
||||
name: { type: 'string' },
|
||||
});
|
||||
|
||||
beforeEach(() => {
|
||||
SharedMemoryBackend.clearRegistry();
|
||||
installOPFSMock(new Map());
|
||||
});
|
||||
|
||||
describe('AriaEngine — 生产负载验证', () => {
|
||||
it('5 万行写入(多级 Compaction + 页面化)→ 完整查询 → 崩溃恢复(OPFS 后端)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 256 * 1024,
|
||||
checkpointInterval: 2000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine.open(dbName, 1);
|
||||
await engine.createTable(SCHEMA());
|
||||
|
||||
const TOTAL = 50000;
|
||||
for (let batch = 0; batch < TOTAL / 1000; batch++) {
|
||||
const rows = [] as Record<string, unknown>[];
|
||||
for (let i = 0; i < 1000; i++) {
|
||||
const idx = batch * 1000 + i;
|
||||
rows.push({ id: `k${idx}`, val: idx, tag: `t${idx % 10}`, name: `User${idx}` });
|
||||
}
|
||||
await engine.insert('big', rows);
|
||||
}
|
||||
expect(await engine.count('big')).toBe(TOTAL);
|
||||
|
||||
// 多级 compaction(VACUUM 语义)
|
||||
await (engine as any).lsm.flush();
|
||||
for (let level = 0; level < 4; level++) {
|
||||
await (engine as any).lsm.compactLevel(level);
|
||||
}
|
||||
const stats = (engine as any).lsm.getStats();
|
||||
expect((stats.levelCounts as number[]).reduce((a: number, b: number) => a + b, 0)).toBeGreaterThanOrEqual(1);
|
||||
|
||||
// 完整查询
|
||||
expect((await engine.find('big', { table: 'big' }))).toHaveLength(TOTAL);
|
||||
|
||||
// 崩溃恢复
|
||||
await (engine as any).backend.close();
|
||||
(engine as any).opened = false;
|
||||
const engine2 = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 256 * 1024,
|
||||
checkpointInterval: 2000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine2.open(dbName, 1);
|
||||
expect(await engine2.count('big')).toBe(TOTAL);
|
||||
for (const id of ['k0', 'k25000', 'k49999']) {
|
||||
expect(await engine2.find('big', { table: 'big', where: { id } })).toHaveLength(1);
|
||||
}
|
||||
// 索引(重启重建)
|
||||
expect(await engine2.find('big', { table: 'big', where: { tag: 't5' } })).toHaveLength(5000);
|
||||
await engine2.close();
|
||||
}, 180000);
|
||||
|
||||
it('高频更新/删除(Compaction 回收墓碑)→ 重启后一致', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 128 * 1024,
|
||||
checkpointInterval: 1000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine.open(dbName, 1);
|
||||
await engine.createTable(SCHEMA());
|
||||
|
||||
const rows = [] as Record<string, unknown>[];
|
||||
for (let i = 0; i < 10000; i++) rows.push({ id: `k${i}`, val: i, tag: `t${i % 5}` });
|
||||
await engine.insert('big', rows);
|
||||
|
||||
let seed = 7;
|
||||
const rand = () => { seed = (seed * 1103515245 + 12345) & 0x7fffffff; return seed / 0x7fffffff; };
|
||||
for (let i = 0; i < 2000; i++) {
|
||||
const id = `k${Math.floor(rand() * 10000)}`;
|
||||
if (rand() < 0.5) {
|
||||
await engine.update('big', { table: 'big', where: { id } }, { val: Math.floor(rand() * 1e9) });
|
||||
} else {
|
||||
await engine.delete('big', { table: 'big', where: { id } });
|
||||
}
|
||||
}
|
||||
await (engine as any).lsm.flush();
|
||||
await (engine as any).lsm.compactLevel(0);
|
||||
|
||||
const count = await engine.count('big');
|
||||
expect(count).toBeGreaterThan(0);
|
||||
await engine.close();
|
||||
|
||||
const engine2 = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 128 * 1024,
|
||||
checkpointInterval: 1000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine2.open(dbName, 1);
|
||||
expect(await engine2.count('big')).toBe(count);
|
||||
await engine2.close();
|
||||
}, 120000);
|
||||
|
||||
it('大 value(100KB × 50)压缩写入/恢复完整', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
compression: true,
|
||||
memtableSizeThreshold: 512 * 1024,
|
||||
checkpointInterval: 5000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine.open(dbName, 1);
|
||||
await engine.createTable(createSchema('docs', {
|
||||
id: { type: 'string', primaryKey: true },
|
||||
body: { type: 'string' },
|
||||
}));
|
||||
const chunk = '这是大段生产数据内容。'.repeat(5000); // ~100KB
|
||||
for (let i = 0; i < 50; i++) {
|
||||
await engine.insert('docs', [{ id: `d${i}`, body: chunk }]);
|
||||
}
|
||||
await engine.close();
|
||||
|
||||
const engine2 = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
compression: true,
|
||||
memtableSizeThreshold: 512 * 1024,
|
||||
checkpointInterval: 5000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine2.open(dbName, 1);
|
||||
expect(await engine2.count('docs')).toBe(50);
|
||||
const one = await engine2.find('docs', { table: 'docs', where: { id: 'd25' } });
|
||||
expect((one[0].body as string).length).toBe(chunk.length);
|
||||
await engine2.close();
|
||||
}, 120000);
|
||||
|
||||
it('混合操作 + 崩溃:已确认写入零丢失(20000 操作)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 128 * 1024,
|
||||
checkpointInterval: 1000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine.open(dbName, 1);
|
||||
await engine.createTable(SCHEMA());
|
||||
|
||||
let seed = 123;
|
||||
const rand = () => { seed = (seed * 1103515245 + 12345) & 0x7fffffff; return seed / 0x7fffffff; };
|
||||
const confirmed = new Map<string, { val: number; tag: string }>();
|
||||
for (let i = 0; i < 20000; i++) {
|
||||
const r = rand();
|
||||
const id = `k${Math.floor(rand() * 10000)}`;
|
||||
if (r < 0.5) {
|
||||
const row = { id, val: Math.floor(rand() * 1e9), tag: `t${Math.floor(rand() * 5)}` };
|
||||
try {
|
||||
await engine.insert('big', [row]);
|
||||
confirmed.set(id, row);
|
||||
} catch (e) {
|
||||
if ((e as { code?: string }).code !== 'DUPLICATE_KEY') throw e;
|
||||
}
|
||||
} else if (r < 0.8) {
|
||||
const val = Math.floor(rand() * 1e9);
|
||||
await engine.update('big', { table: 'big', where: { id } }, { val });
|
||||
if (confirmed.has(id)) confirmed.set(id, { ...confirmed.get(id)!, val });
|
||||
} else {
|
||||
await engine.delete('big', { table: 'big', where: { id } });
|
||||
confirmed.delete(id);
|
||||
}
|
||||
}
|
||||
|
||||
await (engine as any).backend.close();
|
||||
(engine as any).opened = false;
|
||||
|
||||
const engine2 = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 128 * 1024,
|
||||
checkpointInterval: 1000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine2.open(dbName, 1);
|
||||
expect(await engine2.count('big')).toBe(confirmed.size);
|
||||
let sampled = 0;
|
||||
for (const [id, expected] of confirmed) {
|
||||
if (sampled++ > 1000) break;
|
||||
const rows = await engine2.find('big', { table: 'big', where: { id } });
|
||||
expect(rows).toHaveLength(1);
|
||||
expect(rows[0].val).toBe(expected.val);
|
||||
}
|
||||
await engine2.close();
|
||||
}, 120000);
|
||||
|
||||
it('大量删除(90% 行)+ Compaction → 重启无残留(墓碑清理)', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 128 * 1024,
|
||||
checkpointInterval: 1000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine.open(dbName, 1);
|
||||
await engine.createTable(SCHEMA());
|
||||
const rows = [] as Record<string, unknown>[];
|
||||
for (let i = 0; i < 10000; i++) rows.push({ id: `k${i}`, val: i, tag: `t${i % 5}` });
|
||||
await engine.insert('big', rows);
|
||||
|
||||
const toDelete = [] as string[];
|
||||
for (let i = 1000; i < 10000; i++) toDelete.push(`k${i}`);
|
||||
await engine.delete('big', { table: 'big', where: { id: { $in: toDelete } } });
|
||||
expect(await engine.count('big')).toBe(1000);
|
||||
|
||||
await (engine as any).lsm.flush();
|
||||
for (let l = 0; l < 4; l++) await (engine as any).lsm.compactLevel(l);
|
||||
await engine.close();
|
||||
|
||||
const engine2 = new AriaEngine({
|
||||
storageBackend: 'opfs',
|
||||
memtableSizeThreshold: 128 * 1024,
|
||||
checkpointInterval: 1000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine2.open(dbName, 1);
|
||||
expect(await engine2.count('big')).toBe(1000);
|
||||
const all = await engine2.find('big', { table: 'big' });
|
||||
expect(all.every((r) => Number(String(r.id).slice(1)) < 1000)).toBe(true);
|
||||
await engine2.close();
|
||||
}, 120000);
|
||||
|
||||
it('kv 后端 5 万行(页面化路径):写入 → 崩溃 → 恢复完整', async () => {
|
||||
const dbName = uniqueDB();
|
||||
const engine = new AriaEngine({
|
||||
storageBackend: 'kv',
|
||||
memtableSizeThreshold: 512 * 1024,
|
||||
checkpointInterval: 5000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine.open(dbName, 1);
|
||||
await engine.createTable(SCHEMA());
|
||||
|
||||
const TOTAL = 50000;
|
||||
for (let batch = 0; batch < TOTAL / 1000; batch++) {
|
||||
const rows = [] as Record<string, unknown>[];
|
||||
for (let i = 0; i < 1000; i++) {
|
||||
const idx = batch * 1000 + i;
|
||||
rows.push({ id: `k${idx}`, val: idx, tag: `t${idx % 10}`, name: `User${idx}` });
|
||||
}
|
||||
await engine.insert('big', rows);
|
||||
}
|
||||
expect(await engine.count('big')).toBe(TOTAL);
|
||||
await (engine as any).backend.close();
|
||||
(engine as any).opened = false;
|
||||
|
||||
const engine2 = new AriaEngine({
|
||||
storageBackend: 'kv',
|
||||
memtableSizeThreshold: 512 * 1024,
|
||||
checkpointInterval: 5000,
|
||||
walSyncMode: 'full',
|
||||
});
|
||||
await engine2.open(dbName, 1);
|
||||
expect(await engine2.count('big')).toBe(TOTAL);
|
||||
expect(await engine2.find('big', { table: 'big', where: { tag: 't3' } })).toHaveLength(5000);
|
||||
await engine2.close();
|
||||
}, 180000);
|
||||
});
|
||||
Reference in New Issue
Block a user