A38 `compression` 在页面化路径上被静默忽略
压缩只写在"整 value 存一个 backend value"的分支里,而 `save()` 在页面化
分支**提前 return** —— `pageStorage` 默认自动(OPFS 后端下为 true),
于是 `compression: true` 在默认配置下完全无效且无任何提示。
修法:`compression` 传入 `PageSSTableStore`,在**切页之前**整体压缩
(压缩率优于逐页压缩),加载时对称解压。
连带修正一个会静默损坏数据的接口问题:`SSTableMeta.totalSize` 的语义是
"页面里存了多少字节",加载时按它截断 —— 压缩后必须写**压缩长度**。
为此 `SSTableStore.save` 改为返回 `{ storedSize }`,两处 flush 流程与
整 value 路径都用它回填 totalSize(写未压缩长度会让压缩数据被 0 填充撑大)。
为什么此前没被发现:既有测试只断言"压缩后能读回来",而"根本没压缩"
同样能正确读回 —— 断言太弱。新用例改为**结构性断言**:
开启压缩后落盘字节数必须显著下降(>5×),与实现细节无关。
A39 `compressLZ4` 匹配搜索为 O(n²)
旧实现逐字节向前扫描最多 65535 个候选位置、每个位置再逐字节比较 ——
在低压缩率数据上退化为二次复杂度。实测 60KB 伪随机输入耗时 **2345ms**;
而 SSTable 页/日志段正是几百 KB 到几 MB,属于普通写入路径上的真实卡顿。
修法:改为 LZ4 标准的 **4 字节哈希链**(`head[]`/`prev[]`,单点最多
`MAX_CHAIN=32` 次探测)→ 实测 6ms(约 390×)。
**输出格式完全不变**,既有落盘数据无需迁移;旧实现保留为
`compressLZ4LinearReference` 并作为测试对照物(证明两者可互解)。
另加"全字面量"兜底:任何异常都产出合法可解压的流(数据正确性优先于压缩率)。
测试介质修正(同源发现,影响所有 OPFS 多库场景)
`installOPFSMock` 把 `getDirectoryHandle(name)` 的 `name` **丢弃**,
所有库共用一棵扁平文件树。实测:`open('db-alpha')` 建表后
`open('db-beta').getTableNames()` 返回 `["alpha_only"]`。
真实 OPFS 下 `OPFSBackend.open(name)` 是 `root.getDirectoryHandle(name)`,
因此 mock 现在实现真实的**目录语义**,并提供 `dir(dbName)` 视图让测试与
生产代码使用同一个 API(此前的 `listKeys/createFile` 是根目录假 API,
两个依赖它的用例已改为目录视图)。
验证:新增 tests/engine/aria-compression.test.ts(12 项,含 1MB 大输入与
6 组格式兼容用例);两处修复都做**变异验证**:回退 A38 的接线 → 页面化压缩
用例失败;回退 A39 到线性实现 → "60KB < 1s" 用例失败(实测 2397ms)。
全量 89 套件 / 1742 测试通过;typecheck、lint、build 零错误/零告警;dist 已重建。
254 lines
10 KiB
TypeScript
254 lines
10 KiB
TypeScript
/**
|
||
* AriaEngine — repair 自愈增强 + 随机操作压力测试
|
||
*
|
||
* 覆盖:
|
||
* 1. repair 清理孤儿页面(meta 未引用的 pg_ 文件)
|
||
* 2. repair 清理 OPFS 残留临时文件
|
||
* 3. WAL 空洞 + repair → 截断清空(不再重放错位数据)
|
||
* 4. 随机操作压力:insert/update/delete + 模拟崩溃(不 checkpoint 断开)→ 重开全量验证
|
||
* 5. 随机操作 + 页面化 + 模拟崩溃 → 重开验证
|
||
* 6. repair 幂等
|
||
*/
|
||
import { AriaEngine } from '../../src/engine/aria/index';
|
||
import { createSchema } from '../../src/table/schema';
|
||
import { resetOPFSMock } from '../helpers/storage-harness';
|
||
|
||
beforeEach(() => { resetOPFSMock(); });
|
||
|
||
let idbCounter = 0;
|
||
function uniqueDB(): string {
|
||
return `rph-${Date.now()}-${++idbCounter}-${Math.random().toString(36).slice(2, 8)}`;
|
||
}
|
||
|
||
const SCHEMA = () => createSchema('items', {
|
||
id: { type: 'string', primaryKey: true },
|
||
val: { type: 'number' },
|
||
tag: { type: 'string', index: true },
|
||
});
|
||
|
||
// ===================================================================
|
||
// repair 增强
|
||
// ===================================================================
|
||
describe('AriaEngine — repair 自愈增强', () => {
|
||
it('清理孤儿页面(meta 未引用的 pg_ 文件)', async () => {
|
||
const dbName = uniqueDB();
|
||
const engine = new AriaEngine({
|
||
storageBackend: 'opfs',
|
||
pageStorage: true, // 强制页面化(IDB 上也启用,便于断言页面文件)
|
||
memtableSizeThreshold: 64 * 1024 * 1024,
|
||
checkpointInterval: 100000,
|
||
});
|
||
await engine.open(dbName, 1);
|
||
await engine.createTable(SCHEMA());
|
||
await engine.insert('items', [{ id: 'a', val: 1, tag: 'x' }]);
|
||
await (engine as any).lsm.flush();
|
||
|
||
const backend = (engine as any).backend;
|
||
const keys = await backend.listKeys();
|
||
const pgKeys = keys.filter((k: string) => k.startsWith('pg_'));
|
||
expect(pgKeys.length).toBeGreaterThan(0);
|
||
|
||
// 伪造孤儿页面(模拟崩溃中断的删除流程)
|
||
await backend.write('pg_999999', new Uint8Array(4096).buffer);
|
||
await backend.write('pg_999998', new Uint8Array(4096).buffer);
|
||
expect((await backend.listKeys()).filter((k: string) => k.startsWith('pg_99999'))).toHaveLength(2);
|
||
|
||
await (engine as any).repair();
|
||
|
||
// 孤儿页面被清理,正常页面保留
|
||
const after = await backend.listKeys();
|
||
expect(after).not.toContain('pg_999999');
|
||
expect(after).not.toContain('pg_999998');
|
||
expect(await engine.count('items')).toBe(1);
|
||
await engine.close();
|
||
});
|
||
|
||
it('清理 OPFS 残留临时文件(.crswap/.tmp)', async () => {
|
||
// 用 OPFS mock 后端验证 cleanupStaleFiles 被调用
|
||
const opfs = resetOPFSMock();
|
||
const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000 });
|
||
await engine.open('repair-opfs-1', 1);
|
||
await engine.createTable(SCHEMA());
|
||
await engine.insert('items', [{ id: 'a', val: 1, tag: 'x' }]);
|
||
|
||
// 制造残留(模拟 createWritable 中断留下的临时文件)。
|
||
// v0.8.0:mock 现在是真实目录语义,残留必须落在**该库的目录**里 ——
|
||
// repair 的 cleanupStaleFiles 扫的是库目录(root 上的文件不属于任何库)。
|
||
const dir = opfs.dir('repair-opfs-1');
|
||
await dir.createFile('junk.crswap');
|
||
await dir.createFile('junk2.tmp');
|
||
expect((await dir.listKeys()).some((k) => k.endsWith('.crswap'))).toBe(true);
|
||
|
||
await (engine as any).repair();
|
||
|
||
const after = await dir.listKeys();
|
||
expect(after.some((k) => k.endsWith('.crswap'))).toBe(false);
|
||
expect(after.some((k) => k.endsWith('.tmp'))).toBe(false);
|
||
expect(await engine.count('items')).toBe(1);
|
||
await engine.close();
|
||
});
|
||
|
||
it('WAL 空洞 + repair → 截断清空(不再重放错位数据)', async () => {
|
||
const dbName = uniqueDB();
|
||
const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000 });
|
||
await engine.open(dbName, 1);
|
||
await engine.createTable(SCHEMA());
|
||
// 数据落盘 + checkpoint 清 WAL
|
||
await engine.insert('items', [{ id: 'a', val: 1, tag: 'x' }]);
|
||
await (engine as any).lsm.flush();
|
||
await (engine as any).wal.checkpoint();
|
||
// 新数据只进 WAL
|
||
await engine.insert('items', [{ id: 'b', val: 2, tag: 'y' }]);
|
||
const backend = (engine as any).backend;
|
||
const keys = await backend.listKeys();
|
||
const walSeg = keys.find((k: string) => k.startsWith('__wal_'));
|
||
// 制造空洞:删掉当前分片(数据仍在内存)
|
||
await backend.delete(walSeg);
|
||
|
||
await (engine as any).repair();
|
||
|
||
// repair 后 WAL 已清空,无空洞残留
|
||
const after = await backend.listKeys();
|
||
expect(after.some((k: string) => k.startsWith('__wal_'))).toBe(false);
|
||
// 在线自愈:memtable 中的 'b'(已确认写入)随 flush 落盘 → 数据完整不丢
|
||
const rows = await engine.find('items', { table: 'items' });
|
||
expect(rows).toHaveLength(2);
|
||
expect(rows.some((r) => r.id === 'a')).toBe(true);
|
||
expect(rows.some((r) => r.id === 'b')).toBe(true);
|
||
await engine.close();
|
||
});
|
||
|
||
it('repair 幂等(连续调用无副作用)', async () => {
|
||
const dbName = uniqueDB();
|
||
const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000 });
|
||
await engine.open(dbName, 1);
|
||
await engine.createTable(SCHEMA());
|
||
await engine.insert('items', [{ id: 'a', val: 1, tag: 'x' }]);
|
||
await (engine as any).repair();
|
||
await (engine as any).repair();
|
||
await (engine as any).repair();
|
||
expect(await engine.count('items')).toBe(1);
|
||
await engine.close();
|
||
});
|
||
});
|
||
|
||
// ===================================================================
|
||
// 随机操作压力 + 模拟崩溃
|
||
// ===================================================================
|
||
describe('AriaEngine — 随机操作压力 + 模拟崩溃', () => {
|
||
it('500 随机操作(insert/update/delete)→ 模拟崩溃 → 重开验证全部已确认写入', async () => {
|
||
const dbName = uniqueDB();
|
||
const engine = new AriaEngine({
|
||
storageBackend: 'opfs',
|
||
memtableSizeThreshold: 8 * 1024,
|
||
checkpointInterval: 50,
|
||
walSyncMode: 'full',
|
||
});
|
||
await engine.open(dbName, 1);
|
||
await engine.createTable(SCHEMA());
|
||
|
||
let seed = 12345;
|
||
const rand = () => { seed = (seed * 1103515245 + 12345) & 0x7fffffff; return seed / 0x7fffffff; };
|
||
|
||
// 已确认写入集合(模拟应用层已收到成功返回)
|
||
const confirmed = new Map<string, { val: number; tag: string }>();
|
||
for (let i = 0; i < 500; i++) {
|
||
const r = rand();
|
||
const id = `k-${Math.floor(rand() * 100)}`;
|
||
if (r < 0.6) {
|
||
// insert/upsert(随机 id 可能碰撞 → 已存在则跳过,confirmed 保持首个值)
|
||
const row = { id, val: Math.floor(rand() * 1000), tag: `t${Math.floor(rand() * 5)}` };
|
||
try {
|
||
await engine.insert('items', [row]);
|
||
confirmed.set(id, row);
|
||
} catch (e) {
|
||
if ((e as { code?: string }).code !== 'DUPLICATE_KEY') throw e;
|
||
}
|
||
} else if (r < 0.8) {
|
||
// update
|
||
const upd = { val: Math.floor(rand() * 1000) };
|
||
await engine.update('items', { table: 'items', where: { id } }, upd);
|
||
const cur = confirmed.get(id);
|
||
if (cur) confirmed.set(id, { ...cur, ...upd });
|
||
} else {
|
||
// delete
|
||
await engine.delete('items', { table: 'items', where: { id } });
|
||
confirmed.delete(id);
|
||
}
|
||
}
|
||
|
||
// 模拟崩溃:不 close(不 checkpoint),直接断开 backend
|
||
await (engine as any).backend.close();
|
||
(engine as any).opened = false;
|
||
|
||
// 重开:WAL 重放 + SSTable 加载
|
||
const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 50, walSyncMode: 'full' });
|
||
await engine2.open(dbName, 1);
|
||
const rows = await engine2.find('items', { table: 'items' });
|
||
const byId = new Map(rows.map((r) => [r.id, r]));
|
||
// 全部已确认写入可见
|
||
expect(byId.size).toBe(confirmed.size);
|
||
for (const [id, expected] of confirmed) {
|
||
expect(byId.has(id)).toBe(true);
|
||
expect(byId.get(id)!.val).toBe(expected.val);
|
||
}
|
||
await engine2.close();
|
||
});
|
||
|
||
it('随机操作 + 页面化 + 模拟崩溃 → 重开验证', async () => {
|
||
resetOPFSMock();
|
||
const dbName = `rand-opfs-${Date.now()}-${Math.random().toString(36).slice(2, 6)}`;
|
||
const engine = new AriaEngine({
|
||
storageBackend: 'opfs',
|
||
memtableSizeThreshold: 8 * 1024,
|
||
checkpointInterval: 50,
|
||
walSyncMode: 'full',
|
||
});
|
||
await engine.open(dbName, 1);
|
||
await engine.createTable(SCHEMA());
|
||
|
||
let seed = 999;
|
||
const rand = () => { seed = (seed * 1103515245 + 12345) & 0x7fffffff; return seed / 0x7fffffff; };
|
||
const confirmed = new Map<string, { val: number; tag: string }>();
|
||
for (let i = 0; i < 300; i++) {
|
||
const r = rand();
|
||
const id = `p-${Math.floor(rand() * 80)}`;
|
||
if (r < 0.6) {
|
||
const row = { id, val: Math.floor(rand() * 500), tag: `tag${Math.floor(rand() * 4)}` };
|
||
try {
|
||
await engine.insert('items', [row]);
|
||
confirmed.set(id, row);
|
||
} catch (e) {
|
||
if ((e as { code?: string }).code !== 'DUPLICATE_KEY') throw e;
|
||
}
|
||
} else if (r < 0.8) {
|
||
const upd = { val: Math.floor(rand() * 500) };
|
||
await engine.update('items', { table: 'items', where: { id } }, upd);
|
||
const cur = confirmed.get(id);
|
||
if (cur) confirmed.set(id, { ...cur, ...upd });
|
||
} else {
|
||
await engine.delete('items', { table: 'items', where: { id } });
|
||
confirmed.delete(id);
|
||
}
|
||
}
|
||
|
||
// 模拟崩溃
|
||
await (engine as any).backend.close();
|
||
(engine as any).opened = false;
|
||
|
||
const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 50, walSyncMode: 'full' });
|
||
await engine2.open(dbName, 1);
|
||
const rows = await engine2.find('items', { table: 'items' });
|
||
const byId = new Map(rows.map((r) => [r.id, r]));
|
||
expect(byId.size).toBe(confirmed.size);
|
||
for (const [id, expected] of confirmed) {
|
||
expect(byId.get(id)!.val).toBe(expected.val);
|
||
}
|
||
// 索引查询也验证
|
||
const tagged = await engine2.find('items', { table: 'items', where: { tag: 'tag0' } });
|
||
const expectedTagged = Array.from(confirmed.values()).filter((v) => v.tag === 'tag0').length;
|
||
expect(tagged.length).toBe(expectedTagged);
|
||
await engine2.close();
|
||
});
|
||
});
|