Files
MetonaSqlark/tests/engine/kvstore.test.ts
T
thzxx 25bab0f9aa
CI / test (20.x) (push) Successful in 10m55s
CI / test (22.x) (push) Successful in 10m49s
CI / e2e (push) Successful in 9m54s
CI / test (18.x) (push) Successful in 10m58s
CI / test (24.x) (push) Successful in 10m41s
feat: v0.6.1 — AriaEngine 可选自研 KVStore 后端(storageBackend: 'kv')+ KVStore APPEND 日志类型 + 10 个 aria+kv 集成测试 + 文档全量同步
2026-08-10 14:52:38 +08:00

469 lines
17 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.
/**
* KVStore — 自研 KV 引擎单元测试
*
* 覆盖:
* 1. 基本读写(put/get/delete/exists/listKeys
* 2. putMany/deleteMany 原子性(多 key 一次落盘)
* 3. 跨实例持久化(SharedMemory 全局注册表,close 不清数据)
* 4. 崩溃恢复:日志重放(未 checkpoint 数据恢复)
* 5. checkpoint 后恢复(快照 + 水位跳过)
* 6. 快照损坏 → 全量日志重放自愈
* 7. 日志损坏 → 截断至损坏处(丢弃未确认尾部)
* 8. 写入失败原子性(日志失败不更新内存索引)
* 9. clear / repair
* 10. 并发写与 checkpoint 串行(无交错丢数据)
*/
import { KVStore } from '../../src/engine/kvstore/index';
import { SharedMemoryBackend } from '../../src/engine/kvstore/shared_memory_medium';
import { encodeLogRecord, parseLogRecords, KVLogOp } from '../../src/engine/kvstore/log';
import { encodeSnapshot, decodeSnapshot } from '../../src/engine/kvstore/snapshot';
let dbCounter = 0;
function uniqueDB(): string {
return `kv-${Date.now()}-${++dbCounter}-${Math.random().toString(36).slice(2, 8)}`;
}
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('KVStore — 基本读写', () => {
it('put/get/delete/exists/listKeys/size', async () => {
const kv = new KVStore(new SharedMemoryBackend());
await kv.open(uniqueDB());
await kv.put('a', enc('AAA'));
await kv.put('b', enc('BBB'));
expect(dec(await kv.get('a'))).toBe('AAA');
expect(await kv.exists('b')).toBe(true);
expect(await kv.exists('nope')).toBe(false);
expect((await kv.listKeys()).sort()).toEqual(['a', 'b']);
expect(kv.size()).toBe(2);
await kv.delete('a');
expect(await kv.exists('a')).toBe(false);
expect(kv.size()).toBe(1);
await kv.close();
});
it('putMany 多 key 原子写入', async () => {
const kv = new KVStore(new SharedMemoryBackend());
await kv.open(uniqueDB());
await kv.putMany({ a: enc('1'), b: enc('2'), c: enc('3') });
expect(dec(await kv.get('a'))).toBe('1');
expect(dec(await kv.get('c'))).toBe('3');
await kv.deleteMany(['a', 'c']);
expect(await kv.exists('a')).toBe(false);
expect(await kv.exists('c')).toBe(false);
expect(await kv.exists('b')).toBe(true);
await kv.close();
});
});
describe('KVStore — 持久化与崩溃恢复', () => {
it('未 checkpoint 的数据:close 后重开经日志重放恢复', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium);
await kv1.open(dbName);
await kv1.putMany({ a: enc('AAA'), b: enc('BBB') });
await kv1.put('c', enc('CCC'));
await kv1.delete('b');
await kv1.close();
// 模拟"崩溃后重开"(新实例,共享介质)
const kv2 = new KVStore(medium);
await kv2.open(dbName);
expect(dec(await kv2.get('a'))).toBe('AAA');
expect(await kv2.exists('b')).toBe(false);
expect(dec(await kv2.get('c'))).toBe('CCC');
await kv2.close();
});
it('checkpoint 后重开:快照加载 + 日志水位跳过(无重复/无丢失)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0); // 关闭自动 checkpoint
await kv1.open(dbName);
await kv1.putMany({ a: enc('A1'), b: enc('B1') });
await kv1.checkpoint();
// checkpoint 后新写入(进日志)
await kv1.put('c', enc('C1'));
await kv1.put('a', enc('A2'));
await kv1.close();
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
expect(dec(await kv2.get('a'))).toBe('A2'); // 日志重放覆盖快照值
expect(dec(await kv2.get('b'))).toBe('B1');
expect(dec(await kv2.get('c'))).toBe('C1');
await kv2.close();
});
it('多次 checkpoint + 截断日志后重开正确', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0);
await kv1.open(dbName);
for (let i = 0; i < 10; i++) {
await kv1.put(`k${i}`, enc(`v${i}`));
await kv1.checkpoint();
}
await kv1.put('last', enc('L'));
await kv1.close();
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
expect(kv2.size()).toBe(11);
expect(dec(await kv2.get('k9'))).toBe('v9');
expect(dec(await kv2.get('last'))).toBe('L');
await kv2.close();
});
it('快照损坏 → 全量日志重放:日志数据保留,快照唯一副本丢失(自愈不崩)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0);
await kv1.open(dbName);
await kv1.putMany({ a: enc('AAA'), b: enc('BBB') });
await kv1.checkpoint(); // 快照成为 a/b 的唯一副本(日志已截断)
await kv1.put('c', enc('CCC')); // 日志中(seq=2
await kv1.close();
// 篡改快照(close 后介质不可用,重新打开访问"磁盘")
const disk = new SharedMemoryBackend();
await disk.open(dbName);
const snap = await disk.read('__kv_snapshot');
expect(snap).not.toBeNull();
const corrupted = new Uint8Array(snap as ArrayBuffer);
corrupted[20] ^= 0xff; // 破坏 entry 长度字段 → CRC/边界校验失败
await disk.write('__kv_snapshot', corrupted.buffer as ArrayBuffer);
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
// 快照损坏 → 全量日志重放:c(日志中)保留;a/b(快照唯一副本)丢失
expect(await kv2.exists('a')).toBe(false);
expect(await kv2.exists('b')).toBe(false);
expect(dec(await kv2.get('c'))).toBe('CCC');
await kv2.close();
});
it('回归:快照完好但 meta.seq 超前 → 日志不被 metaSeq 跳过(水位只信任快照)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0);
await kv1.open(dbName);
await kv1.put('a', enc('AAA'));
await kv1.checkpoint(); // 快照 seq=1meta seq=1,日志截断
await kv1.put('b', enc('BBB')); // 日志 seq=2
await kv1.close();
// 手工把 meta.seq 改大(模拟异常状态:meta 比日志内容超前)
const disk = new SharedMemoryBackend();
await disk.open(dbName);
await disk.write('__kv_meta', new TextEncoder().encode(JSON.stringify({ seq: 999 })).buffer);
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
// 快照完好(seq=1)→ 水位=1 → 日志 seq2 应用;meta.seq=999 不参与水位判断
expect(dec(await kv2.get('a'))).toBe('AAA');
expect(dec(await kv2.get('b'))).toBe('BBB');
await kv2.close();
});
it('日志损坏尾部 → 截断至损坏处(丢弃未确认尾部,已确认数据保留)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0);
await kv1.open(dbName);
await kv1.put('confirmed', enc('KEEP'));
await kv1.put('confirmed2', enc('KEEP2'));
await kv1.close();
// 追加一条损坏记录(模拟写入中断:半写记录)
const bad = encodeLogRecord(999, { damaged: enc('X') });
const corrupted = new Uint8Array(bad);
corrupted[20] ^= 0xff; // 破坏 CRC
// close 后介质实例不可用,重新打开访问"磁盘"
const disk = new SharedMemoryBackend();
await disk.open(dbName);
const log = await disk.read('__kv_log');
const combined = new Uint8Array((log as ArrayBuffer).byteLength + corrupted.byteLength);
combined.set(new Uint8Array(log as ArrayBuffer), 0);
combined.set(corrupted, (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('confirmed'))).toBe('KEEP');
expect(dec(await kv2.get('confirmed2'))).toBe('KEEP2');
expect(await kv2.exists('damaged')).toBe(false);
// 损坏日志已被截断(重开后再写正常)
await kv2.put('after', enc('OK'));
await kv2.close();
const kv3 = new KVStore(medium, 0);
await kv3.open(dbName);
expect(dec(await kv3.get('after'))).toBe('OK');
await kv3.close();
});
it('写入失败 → 内存索引不更新(原子性)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv = new KVStore(medium, 0);
await kv.open(dbName);
// 注入日志追加失败
const origAppend = medium.append!.bind(medium);
medium.append = async () => { throw new Error('disk full'); };
await expect(kv.put('x', enc('X'))).rejects.toMatchObject({ code: 'KV_LOG_ERROR' });
expect(await kv.exists('x')).toBe(false);
medium.append = origAppend;
await kv.put('y', enc('Y'));
expect(dec(await kv.get('y'))).toBe('Y');
await kv.close();
});
});
describe('KVStore — 维护与并发', () => {
it('clear 清空全部(保留库),重开为空', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0);
await kv1.open(dbName);
await kv1.putMany({ a: enc('1'), b: enc('2') });
await kv1.checkpoint();
await kv1.put('c', enc('3'));
await kv1.clear();
expect(kv1.size()).toBe(0);
await kv1.close();
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
expect(kv2.size()).toBe(0);
await kv2.close();
});
it('repair 清理损坏快照与日志(损坏快照数据无法恢复,日志数据保留)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv1 = new KVStore(medium, 0);
await kv1.open(dbName);
await kv1.put('a', enc('A'));
await kv1.checkpoint(); // 快照成为 a 的唯一副本(日志已截断)
await kv1.put('b', enc('B')); // 日志中(seq=2
await kv1.close();
// 快照损坏 + 日志尾部损坏
const disk = new SharedMemoryBackend();
await disk.open(dbName);
const snap = new Uint8Array(await disk.read('__kv_snapshot') as ArrayBuffer);
snap[10] ^= 0xff;
await disk.write('__kv_snapshot', snap.buffer as ArrayBuffer);
const bad = encodeLogRecord(500, { junk: enc('J') });
bad[15] ^= 0xff;
const log = new Uint8Array(await disk.read('__kv_log') as ArrayBuffer);
const combined = new Uint8Array(log.byteLength + bad.byteLength);
combined.set(log, 0); combined.set(bad, log.byteLength);
await disk.write('__kv_log', combined.buffer as ArrayBuffer);
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
const discarded = await kv2.repair();
expect(discarded).toBeGreaterThan(0);
// 快照损坏(a 的唯一副本丢失)+ 损坏日志尾部被截断 → 日志数据 b 完整
expect(await kv2.exists('a')).toBe(false);
expect(dec(await kv2.get('b'))).toBe('B');
// repair 后介质上的损坏快照已清理
expect(await medium.exists('__kv_snapshot')).toBe(false);
await kv2.close();
});
it('并发写 + checkpoint 串行(无交错丢数据)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv = new KVStore(medium, 0);
await kv.open(dbName);
// 并发发起大量写入 + 中途 checkpoint
const writes = [] as Promise<void>[];
for (let i = 0; i < 200; i++) {
writes.push(kv.put(`k${i}`, enc(`v${i}`)));
if (i === 100) writes.push(kv.checkpoint());
}
await Promise.all(writes);
expect(kv.size()).toBe(200);
await kv.close();
// 重开验证全部恢复
const kv2 = new KVStore(medium, 0);
await kv2.open(dbName);
expect(kv2.size()).toBe(200);
expect(dec(await kv2.get('k150'))).toBe('v150');
await kv2.close();
});
it('自动 checkpoint 阈值触发(日志不无限增长)', async () => {
const dbName = uniqueDB();
const medium = new SharedMemoryBackend();
const kv = new KVStore(medium, 256); // 小阈值
await kv.open(dbName);
for (let i = 0; i < 50; i++) {
await kv.put(`k${i}`, enc(`value-${i}-`.repeat(10)));
}
// 自动 checkpoint 后日志应被截断
const log = await medium.read('__kv_log');
expect((log as ArrayBuffer).byteLength).toBeLessThan(4096);
await kv.close();
const kv2 = new KVStore(medium, 256);
await kv2.open(dbName);
expect(kv2.size()).toBe(50);
await kv2.close();
});
});
describe('KVStore — 编解码单元', () => {
it('encodeLogRecord / parseLogRecords 往返', () => {
const rec = encodeLogRecord(1, { a: enc('A'), b: enc('BB') }, ['del']);
const records: { seq: number; entries: { op: KVLogOp; key: string; value: ArrayBuffer }[] }[] = [];
parseLogRecords(rec, (r) => records.push(r));
expect(records).toHaveLength(1);
expect(records[0].seq).toBe(1);
expect(records[0].entries).toHaveLength(3);
expect(dec(records[0].entries[0].value)).toBe('A');
expect(records[0].entries[2].op).toBe(KVLogOp.DELETE);
});
it('多记录日志顺序解析', () => {
const r1 = encodeLogRecord(1, { a: enc('1') });
const r2 = encodeLogRecord(2, { b: enc('2') });
const combined = new Uint8Array(r1.byteLength + r2.byteLength);
combined.set(r1, 0); combined.set(r2, r1.byteLength);
const seqs: number[] = [];
parseLogRecords(combined, (r) => seqs.push(r.seq));
expect(seqs).toEqual([1, 2]);
});
it('encodeSnapshot / decodeSnapshot 往返', () => {
const entries = new Map<string, ArrayBuffer>([['a', enc('A')], ['b', enc('BB')]]);
const bytes = encodeSnapshot(42, entries);
const snap = decodeSnapshot(bytes);
expect(snap).not.toBeNull();
expect(snap!.seq).toBe(42);
expect(dec(snap!.entries.get('a'))).toBe('A');
expect(dec(snap!.entries.get('b'))).toBe('BB');
});
it('损坏快照 decode 返回 null', () => {
const entries = new Map<string, ArrayBuffer>([['a', enc('A')]]);
const bytes = encodeSnapshot(1, entries);
bytes[10] ^= 0xff;
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();
});
});