/** * AriaEngine — SegmentedWALStore 分片 WAL 存储测试 * * 覆盖: * 1. append → readAll 往返(分片切换、字节顺序) * 2. 空洞检测:序号不连续 → 空洞后的分片丢弃 * 3. 旧格式 __wal_N + __wal_count 兼容读取 * 4. truncate / exists * 5. 后端无 append 接口时回退 read+write(正确性一致) * 6. AriaEngine 集成:高频写入跨重启数据完整(分片 + WAL 重放) */ import { SegmentedWALStore } from '../../src/engine/aria/wal/segmented_store'; import { MemoryBackend } from '../../src/engine/aria/store/backend'; 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 `seg-${Date.now()}-${++idbCounter}-${Math.random().toString(36).slice(2, 8)}`; } const enc = (s: string) => new Uint8Array(new TextEncoder().encode(s)); /** 构造"首 4 字节为大端 LSN"的记录(分片的水位判定依赖它) */ function rec(lsn: number, pad = ''): Uint8Array { const filler = enc(pad.padEnd(4, '.')); const out = new Uint8Array(4 + filler.byteLength); new DataView(out.buffer).setUint32(0, lsn, false); out.set(filler, 4); return out; } describe('AriaEngine — SegmentedWALStore 单元', () => { it('append → readAll 往返,多批次字节顺序一致', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-1'); const store = new SegmentedWALStore(backend, 1024); expect(await store.exists()).toBe(false); await store.append(enc('AAA')); await store.append(enc('BBB')); await store.append(enc('CCC')); expect(await store.exists()).toBe(true); const all = await store.readAll(); expect(new TextDecoder().decode(all)).toBe('AAABBBCCC'); // 后端是分片文件(单个 key) const keys = await backend.listKeys(); expect(keys).toEqual(['__wal_000000.bin']); await backend.close(); }); it('超过分片阈值自动切换分片', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-2'); const store = new SegmentedWALStore(backend, 16); // 每分片 16 字节 await store.append(enc('AAAAAAAAAA')); // 10B → 分片 0 await store.append(enc('BBBBBBBBBB')); // 10B → 分片 1(0 已 10B+10B > 16) await store.append(enc('CC')); // 2B → 分片 1(10+2 <= 16) const all = await store.readAll(); expect(new TextDecoder().decode(all)).toBe('AAAAAAAAAABBBBBBBBBBCC'); const keys = await backend.listKeys().then((ks) => ks.sort()); expect(keys).toEqual(['__wal_000000.bin', '__wal_000001.bin']); await backend.close(); }); it('尾部分片缺失(无法与更后分片比较)→ 已读到的分片仍然返回', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-3'); const store = new SegmentedWALStore(backend, 16); await store.append(enc('SEG0-SEG0-S')); // 10B → 分片 0 await store.append(enc('SEG1-SEG1-S')); // 10B → 分片 1(0 已 10B+10B > 16) await store.append(enc('S2')); // 2B → 分片 1(10+2 <= 16) // 模拟 truncate 部分完成:删除分片 1(空洞) await backend.delete('__wal_000001.bin'); const all = await store.readAll(); // 介质上只剩分片 0:没有更后的分片作为参照,缺失的"尾部"无法被识别成空洞 //(这是介质信息本身的限制 —— 尾部丢失只能靠 manifest.nextLsn 之外的证据发现) expect(new TextDecoder().decode(all)).toBe('SEG0-SEG0-S'); await backend.close(); }); it('前缀缺失:水位未推进(fromLsn=0)→ 全部丢弃;水位已推进 → 后缀照常读取', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-4'); const store = new SegmentedWALStore(backend, 3); // 2 字节记录 → 每条一个分片 await store.append(enc('X0')); await store.append(enc('X1')); await backend.delete('__wal_000000.bin'); // 从未推进过水位:分片本应从 0 连续存在 → 前缀缺失是异常,保守丢弃 const all = await store.readAll(); expect(all.byteLength).toBe(0); // 水位已推进(例如回退到上一代 manifest,startSegment/startLsn 指向更早的位置): // 前缀那一段本就被水位跳过,**不能**因此丢弃后面的活分片 const fromLater = await store.readAllFrom(0, 1); expect(new TextDecoder().decode(fromLater.data)).toBe('X1'); expect(fromLater.missingPrefix).toEqual([0]); expect(fromLater.gaps).toEqual([]); await backend.close(); }); it('内部空洞:空洞之后的分片整体丢弃,并把空洞号上报', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-4b'); const store = new SegmentedWALStore(backend, 8); // 8 字节记录 → 每条一个分片 await store.append(rec(1, 'aaaa')); // 分片 0 await store.append(rec(2, 'bbbb')); // 分片 1 await store.append(rec(3, 'cccc')); // 分片 2 expect((await store.readAllFrom(0, 0)).segments).toEqual([0, 1, 2]); await backend.delete('__wal_000001.bin'); // 内部空洞 const result = await store.readAllFrom(0, 0); expect(result.segments).toEqual([0]); // 只读了空洞之前的分片 expect(result.gaps).toEqual([1]); // 显式上报(不再静默) expect(result.data.byteLength).toBe(8); // 只保留空洞之前的分片 expect(new DataView(result.data.buffer, result.data.byteOffset).getUint32(0, false)).toBe(1); await backend.close(); }); it('整体清空后分片号绝不回退(可以复用最后用过的号,但绝不能回到 0)', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-4c'); const store = new SegmentedWALStore(backend, 8); // 8 字节记录 → 每条一个分片 await store.append(rec(1, 'aaaa')); // 分片 0 await store.append(rec(2, 'bbbb')); // 分片 1 await store.append(rec(3, 'cccc')); // 分片 2 expect((await store.readAllFrom(0, 0)).segments).toEqual([0, 1, 2]); await store.truncate(); expect((await backend.listKeys()).length).toBe(0); // 介质上确实没有分片了 await store.append(rec(4, 'dddd')); const walKeys = (await backend.listKeys()).filter((k) => k.startsWith('__wal_')); expect(walKeys).toHaveLength(1); const newSeq = Number(walKeys[0].match(/^__wal_(\d{6,})\.bin$/)![1]); // 关键不变量:不回退到 0(回退 + manifest.startSegment>0 = 新记录被恢复过滤掉) expect(newSeq).toBeGreaterThanOrEqual(2); expect(walKeys).not.toContain('__wal_000000.bin'); // 清空后从新分片号读取仍能读回新记录(该契约不能被破坏) const after = await store.readAllFrom(newSeq, 0); expect(after.data.byteLength).toBe(8); expect(new DataView(after.data.buffer, after.data.byteOffset).getUint32(0, false)).toBe(4); await backend.close(); }); it('写入永不落到 manifest 水位下限之下(startSegment=K → 新分片号 >= K)', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-4e'); const store = new SegmentedWALStore(backend, 8); await store.append(rec(1, 'aaaa')); // 分片 0(首条 LSN=1) await store.append(rec(2, 'bbbb')); // 分片 1(首条 LSN=2) await store.append(rec(3, 'cccc')); // 分片 2(首条 LSN=3) // 模拟 checkpoint:先算水位下限并**提交进 manifest**,再删除前缀分片 const keepFrom = await store.planKeepFrom(3, 3); // 全部记录都已落盘 → 全部可删 expect(keepFrom).toBe(3); await store.truncateBefore(3, 3); expect((await backend.listKeys()).filter((k) => k.startsWith('__wal_'))).toEqual([]); // 之后的新写入必须落在 >= keepFrom 的分片里(否则重开时被 `seq >= startSegment` 过滤 → 静默丢失) await store.append(rec(4, 'dddd')); const seqs = (await backend.listKeys()) .filter((k) => k.startsWith('__wal_')) .map((k) => Number(k.match(/^__wal_(\d{6,})\.bin$/)![1])); expect(seqs.length).toBeGreaterThan(0); expect(Math.min(...seqs)).toBeGreaterThanOrEqual(keepFrom); // 而且从 manifest 的水位下限开始读,这条新记录必须读得到(双向验证) const reread = await store.readAllFrom(keepFrom, 3); expect(reread.data.byteLength).toBe(8); expect(new DataView(reread.data.buffer, reread.data.byteOffset).getUint32(0, false)).toBe(4); await backend.close(); }); it('活跃区间已被清理(fromSegment > 0,介质上无分片):新写入从 fromSegment 继续', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-4d'); const store = new SegmentedWALStore(backend, 8); await store.append(enc('AAAA')); await store.append(enc('BBBB')); // 分片 0、1 // 介质上分片被全部清理,manifest 说 startSegment = 5 const later = await store.readAllFrom(5, 100); // 一个分片都没有:无法判断缺了哪些号 → 不谎报空洞,也不丢已有数据 expect(later.data.byteLength).toBe(0); expect(later.segments).toEqual([]); expect(later.gaps).toEqual([]); await store.append(enc('EEEE')); const keys = await backend.listKeys(); expect(keys.some((k) => k === '__wal_000005.bin')).toBe(true); // 从 fromSegment 继续 await backend.close(); }); it('旧格式 __wal_N + __wal_count 兼容读取(迁移前数据)', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-5'); // 构造旧格式数据 await backend.write('__wal_0', enc('LEGACY0').buffer); await backend.write('__wal_1', enc('LEGACY1').buffer); await backend.write('__wal_count', enc('2').buffer); const store = new SegmentedWALStore(backend, 1024); expect(await store.exists()).toBe(true); const all = await store.readAll(); expect(new TextDecoder().decode(all)).toBe('LEGACY0LEGACY1'); // 迁移中追加新格式:新记录写入分片文件(不与旧键冲突) await store.append(enc('NEW')); const keys = await backend.listKeys().then((ks) => ks.sort()); expect(keys).toContain('__wal_000000.bin'); expect(keys).toContain('__wal_0'); // truncate 清空新旧全部 await store.truncate(); expect(await store.exists()).toBe(false); expect((await backend.listKeys()).length).toBe(0); await backend.close(); }); it('truncate 清空分片并重置状态', async () => { const backend = new MemoryBackend(); await backend.open('seg-unit-6'); const store = new SegmentedWALStore(backend, 16); await store.append(enc('AAAAAA')); await store.append(enc('BBBBBB')); await store.truncate(); expect(await store.exists()).toBe(false); expect((await backend.listKeys()).length).toBe(0); // 截断后继续追加从分片 0 重新开始 await store.append(enc('CCC')); const all = await store.readAll(); expect(new TextDecoder().decode(all)).toBe('CCC'); await backend.close(); }); it('后端无 append 接口 → 回退 read+write 语义一致', async () => { // MemoryBackend 没有 append 接口 → 走回退路径 const backend = new MemoryBackend(); await backend.open('seg-unit-7'); const store = new SegmentedWALStore(backend, 1024); await store.append(enc('PART1')); await store.append(enc('PART2')); const all = await store.readAll(); expect(new TextDecoder().decode(all)).toBe('PART1PART2'); await backend.close(); }); }); // =================================================================== // AriaEngine 集成 — 分片 WAL 跨重启 // =================================================================== describe('AriaEngine — 分片 WAL 集成', () => { it('高频写入(多条 WAL 记录)→ close → reopen 数据完整', async () => { const dbName = uniqueDB(); const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000, walSyncMode: 'full', }); await engine.open(dbName, 1); await engine.createTable(createSchema('logs', { id: { type: 'string', primaryKey: true }, msg: { type: 'string' }, })); const total = 300; for (let i = 0; i < total; i++) { await engine.insert('logs', [{ id: `log-${i}`, msg: `message ${i}` }]); } // 不 flush 不 checkpoint,全部留在 WAL → 模拟崩溃后重放 await (engine as any).backend.close(); (engine as any).opened = false; const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000, walSyncMode: 'full', }); await engine2.open(dbName, 1); const rows = await engine2.find('logs', { table: 'logs' }); expect(rows).toHaveLength(total); expect(rows.some((r) => r.id === 'log-299')).toBe(true); await engine2.close(); }); it('分片结构落盘验证(backend 中是分片文件而非旧单记录键)', async () => { const dbName = uniqueDB(); const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000, walSyncMode: 'full', }); await engine.open(dbName, 1); await engine.createTable(createSchema('t', { id: { type: 'string', primaryKey: true }, })); await engine.insert('t', [{ id: '1' }, { id: '2' }, { id: '3' }]); const backend = (engine as any).backend; const keys = await backend.listKeys(); const walKeys = keys.filter((k: string) => k.startsWith('__wal_')); // 新格式:__wal_000000.bin;无旧格式单记录键与 count 键 expect(walKeys).toContain('__wal_000000.bin'); expect(walKeys.some((k: string) => /^__wal_\d+$/.test(k) && !k.endsWith('.bin'))).toBe(false); expect(walKeys).not.toContain('__wal_count'); await engine.close(); }); it('WAL 分片部分残留(空洞)→ 重开恢复已落盘数据,不丢已 checkpoint 数据', async () => { const dbName = uniqueDB(); const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000, walSyncMode: 'full', }); await engine.open(dbName, 1); await engine.createTable(createSchema('t', { id: { type: 'string', primaryKey: true }, })); // 第一批:flush 落盘 + checkpoint 清 WAL await engine.insert('t', [{ id: 'a' }, { id: 'b' }]); await (engine as any).lsm.flush(); await (engine as any).wal.checkpoint(); // 第二批:只进 WAL(不 flush) await engine.insert('t', [{ id: 'c' }, { id: 'd' }]); // 模拟 WAL 分片文件被外部删掉(空洞场景) const backend = (engine as any).backend; await backend.delete('__wal_000000.bin'); await engine.close(); const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000 }); await engine2.open(dbName, 1); // 已落盘的 a/b 必须保留;c/d 在 WAL 中且文件被删 → 恢复不到(可接受的保守丢弃) const rows = await engine2.find('t', { table: 't' }); expect(rows.some((r) => r.id === 'a')).toBe(true); expect(rows.some((r) => r.id === 'b')).toBe(true); await engine2.close(); }); it('旧格式 WAL 升级迁移:旧键重放后 checkpoint 清空', async () => { const dbName = uniqueDB(); // 手工构造旧格式 WAL 库:直接写旧键(模拟 v0.4.4 库崩溃现场) const engine = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000, walSyncMode: 'full', }); await engine.open(dbName, 1); await engine.createTable(createSchema('t', { id: { type: 'string', primaryKey: true }, })); // 写入一批数据但只保留在 WAL(新格式)→ close 前手工转成旧格式 await engine.insert('t', [{ id: 'x' }]); const backend = (engine as any).backend; const raw = await backend.read('__wal_000000.bin'); expect(raw).not.toBeNull(); // 清掉新格式,写成旧格式(单记录键) await backend.delete('__wal_000000.bin'); await backend.write('__wal_0', raw); await backend.write('__wal_count', new TextEncoder().encode('1').buffer); // 同时清掉内存里的 WAL 状态,模拟"重启" await (engine as any).backend.close(); (engine as any).opened = false; const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100000 }); await engine2.open(dbName, 1); const rows = await engine2.find('t', { table: 't' }); expect(rows).toHaveLength(1); expect(rows[0].id).toBe('x'); // 恢复后 checkpoint 清空旧键 const keys = await (engine2 as any).backend.listKeys(); expect(keys.some((k: string) => k.startsWith('__wal_'))).toBe(false); await engine2.close(); }); });