Files
MetonaSqlark/tests/v080-b6-single-commit-point.test.ts
thzxx 50b1864145 chore(docs): 删除 v0.7.4 审计与 v0.7.5 计划共 4 个 md 文档
删除:
- AUDIT-aria-lsm-v0.7.4.md(50KB)
- AUDIT-query-layer-v0.7.4.md(39KB)
- AUDIT-storage-engines-v0.7.4.md(39KB)
- PLAN-v0.7.5.md(98KB,含附录 G/H/I)

删除前先清除引用面,避免留下断链(共 18 处):
- 源码注释 7 处(change-notifier / kvstore index / column-value / expression /
  sql-compare / where-matcher / validation):保留设计意图,引用改为"v0.8.0 审计根因 N"
- 测试注释 9 处(opfs.spec / aria-opfs-backend / faulty-backend / storage-harness /
  v080-b6 / v080-kvstore / v080-query-layer / v080-sql-three-valued /
  v080-unified-validation / parser):同上
- CHANGELOG 3 处:改为不依赖已删除文档的自洽表述(B-6 交付物见各条;门禁订正三处
  按内容重写),并把变异数量同步为 42
- 校验:三个 md 之间无断链;仓库内已无 PLAN-v0.7.5/AUDIT-* 的任何引用
  (git 历史仍可追溯,需要时可 `git show <commit>:PLAN-v0.7.5.md` 找回)

验证:93 套件 / 1985 用例全绿;覆盖率 90.59 / 82.61 / 94.14 / 93.50(阈值 90/82/94/93);
e2e 14/14;lint + 两份 tsc 干净;dist 已重建(注释只影响非压缩产物,min 产物
251,731 B / gzip 63,431 B 不变)。

说明:审查记录的核心内容仍在 CHANGELOG.md("全量回归审查"与"现场失败修复"两节),
随 PLAN 一起删除的是附录 G/H/I 的详细表格(门禁逐条验收、交付物清单、未修复项表)。
2026-09-15 17:17:24 +08:00

2180 lines
104 KiB
TypeScript
Raw Permalink 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.
/**
* v0.8.0(B-6)回归套件 —— 存储层**单一提交点**`__aria_manifest`)与 LSM 结构根治
* ============================================================================
* 本套件覆盖 v0.8.0 迭代工作流 B 的 B-6 全部条目。每一条都对应一个**实测过**
* 的缺陷(或一条被写进 manifest 的持久不变量):
*
* | 编号 | 缺陷 / 不变量 | 本文件对应用例 |
* |---|---|---|
* | B-6 核心 | 元数据各自独立落盘 → 崩溃窗口内互相矛盾 | `manifest 单一提交点` |
* | 消除 7 | SSTable meta 损坏 → 静默空库 + repair 删活页 | `元数据损坏必须显式失败` |
* | 消除 4 | 介质读故障被折叠为"文件不存在" → meta 被误删 | `介质读故障 ≠ 文件缺失` |
* | 消除 9 | WAL 分片空洞 → 尾部静默丢弃 | `WAL 分片空洞必须显式失败` |
* | 消除 5 | 陈旧实例覆盖新实例的 meta | `陈旧实例提交被拒绝` |
* | 44 | 后台 flush 失败后冻结表无重试路径 → 崩溃即丢 | `flush 失败可重试` |
* | 45 | 后台错误检查在入链之前 → 本次 flush 被整个跳过 | `flush 必须先入链再报错` |
* | 47 | 提前终止时底层生成器多产出 1 条 | `流式扫描提前终止不多算` |
* | 49 | `compacting` 单 boolean → 跨层触发被静默丢弃 | `按层 compaction 状态` |
* | 50 | compaction 先 `splice` 整层 → 窗口内该层对读者不可见 | `compaction 期间读不到空层` |
* | 51 | 墓碑/历史版本永不回收 | `底部层合并回收墓碑` |
* | 55 | checkpoint 等完整 compaction → 写路径秒级卡顿 | `checkpoint 不等 compaction` |
* | 冻结意图 | "已确认写入"是否真的还在,此前不可观测 | `冻结意图阻止水位推进` |
*
* 变异验证:把对应实现回退到修复前的行为,这些用例必须失败(见各用例注释)。
*/
import { describe, it, expect, beforeEach } from '@jest/globals';
import { AriaEngine } from '../src/engine/aria/index';
import { LSM } from '../src/engine/aria/index/lsm';
import { MergeIterator, type EntrySource } from '../src/engine/aria/index/merge_iterator';
import { SSTableReader } from '../src/engine/aria/index/sstable';
import {
ManifestStore,
MANIFEST_HEADER_SIZE,
createEmptyManifest,
decodeManifest,
encodeManifest,
generationFromKey,
manifestKey,
type AriaManifest,
} from '../src/engine/aria/store/manifest';
import { MemoryBackend, type IStorageBackend } from '../src/engine/aria/store/backend';
import { crc32 } from '../src/engine/aria/crc32';
import { OPFSBackend } from '../src/engine/aria/store/opfs_backend';
import { SegmentedWALStore } from '../src/engine/aria/wal/segmented_store';
import { WAL } from '../src/engine/aria/wal/log';
import { WALRecordType, MAX_LSM_LEVELS } from '../src/engine/aria/types';
import { resetOPFSMock, readManifestState } from './helpers/storage-harness';
import { createSchema } from '../src/table/schema';
jest.setTimeout(60000);
beforeEach(() => { resetOPFSMock(); });
let dbCounter = 0;
function uniqueDB(tag: string): string {
return `${tag}-${Date.now()}-${++dbCounter}`;
}
const SCHEMA = () => createSchema('t', {
id: { type: 'string', primaryKey: true },
v: { type: 'number' },
});
function rows(n: number): Record<string, unknown>[] {
const out: Record<string, unknown>[] = [];
for (let i = 0; i < n; i++) out.push({ id: `k${i}`, v: i });
return out;
}
// ---------------------------------------------------------------------------
// 可注入故障/门控的存储后端(包住真实后端,不替换语义)
// ---------------------------------------------------------------------------
class GatedBackend implements IStorageBackend {
/** 命中门控时置位(测试用它等待"引擎已经走到这个点" */
onGated: (() => void) | null = null;
private gate: Promise<void> = Promise.resolve();
private releaseGate: (() => void) | null = null;
/** 需要门控的 key 前缀(空 = 不门控) */
gatePattern: string | null = null;
/** 门控前放行的匹配次数(默认 0 = 第一次匹配就阻塞) */
gateSkip = 0;
/** 这些 key 前缀上的读会抛错(模拟介质故障) */
failReadPattern: string | null = null;
/** 这些 key 前缀上的写会抛错 */
failWritePattern: string | null = null;
/** 粘性门控:一旦命中就持续阻塞(不是只阻塞一次) */
private sticky = false;
/** 同时卡在门上的写次数(>1 即"并发写",串行化被破坏) */
gatedInFlight = 0;
gatedPeak = 0;
constructor(private inner: IStorageBackend) {}
/** 让第 `skip + 1` 次匹配的写阻塞,直到 `release()` */
armGate(pattern: string, skip = 0): void {
this.gatePattern = pattern;
this.gateSkip = skip;
this.gate = new Promise((resolve) => { this.releaseGate = resolve; });
}
/** 粘性门控:`skip` 次之后**所有**匹配写都阻塞(用于观测"是否发生了并发写" */
armStickyGate(pattern: string, skip = 0): void {
this.sticky = true;
this.armGate(pattern, skip);
}
release(): void {
this.gatePattern = null;
this.sticky = false;
if (this.releaseGate) { this.releaseGate(); this.releaseGate = null; }
}
private matches(pattern: string | null, key: string): boolean {
return pattern !== null && key.startsWith(pattern);
}
open(name: string): Promise<void> { return this.inner.open(name); }
close(): Promise<void> { return this.inner.close(); }
isOpen(): boolean { return this.inner.isOpen(); }
async read(key: string): Promise<ArrayBuffer | null> {
if (this.matches(this.failReadPattern, key)) {
throw new Error(`injected read failure on "${key}"`);
}
return this.inner.read(key);
}
async write(key: string, data: ArrayBuffer): Promise<void> {
if (this.matches(this.failWritePattern, key)) {
throw new Error(`injected write failure on "${key}"`);
}
if (this.matches(this.gatePattern, key)) {
if (this.gateSkip > 0) {
this.gateSkip--;
} else {
// 一次性门控:只阻塞这一次匹配的写;后续写必须能继续(否则测试测到的
// 是"整个介质被冻住",而不是"某一次 compaction 卡住"
if (!this.sticky) this.gatePattern = null;
const notify = this.onGated;
if (notify) { this.onGated = null; notify(); }
this.gatedInFlight++;
this.gatedPeak = Math.max(this.gatedPeak, this.gatedInFlight);
try {
await this.gate;
} finally {
this.gatedInFlight--;
}
}
}
return this.inner.write(key, data);
}
append(key: string, data: ArrayBuffer): Promise<void> {
if (this.matches(this.failWritePattern, key)) {
return Promise.reject(new Error(`injected append failure on "${key}"`));
}
return this.inner.append ? this.inner.append(key, data) : this.inner.write(key, data);
}
writeMany(entries: Record<string, ArrayBuffer>): Promise<void> {
const keys = Object.keys(entries);
if (this.failWritePattern && keys.some((k) => this.matches(this.failWritePattern, k))) {
return Promise.reject(new Error('injected writeMany failure'));
}
return this.inner.writeMany(entries);
}
delete(key: string): Promise<void> { return this.inner.delete(key); }
deleteMany(keys: string[]): Promise<void> { return this.inner.deleteMany(keys); }
listKeys(): Promise<string[]> { return this.inner.listKeys(); }
exists(key: string): Promise<boolean> { return this.inner.exists(key); }
clear(): Promise<void> { return this.inner.clear(); }
}
/** 一次写失败后自动恢复正常(模拟"瞬时故障",用于验证重试路径) */
class FailOnceBackend implements IStorageBackend {
remainingFailures = 1;
/** 默认匹配"整 value SSTable"memory 后端不页面化);页面化时传 'pg_' */
failPattern = 'sst_';
constructor(private inner: IStorageBackend) {}
open(name: string): Promise<void> { return this.inner.open(name); }
close(): Promise<void> { return this.inner.close(); }
isOpen(): boolean { return this.inner.isOpen(); }
read(key: string): Promise<ArrayBuffer | null> { return this.inner.read(key); }
async write(key: string, data: ArrayBuffer): Promise<void> {
if (this.remainingFailures > 0 && key.startsWith(this.failPattern)) {
this.remainingFailures--;
throw new Error(`injected ONE write failure on "${key}"`);
}
return this.inner.write(key, data);
}
append(key: string, data: ArrayBuffer): Promise<void> {
return this.inner.append ? this.inner.append(key, data) : this.inner.write(key, data);
}
writeMany(entries: Record<string, ArrayBuffer>): Promise<void> {
for (const [k] of Object.entries(entries)) {
if (this.remainingFailures > 0 && k.startsWith(this.failPattern)) {
this.remainingFailures--;
return Promise.reject(new Error(`injected ONE write failure on "${k}"`));
}
}
return this.inner.writeMany(entries);
}
delete(key: string): Promise<void> { return this.inner.delete(key); }
deleteMany(keys: string[]): Promise<void> { return this.inner.deleteMany(keys); }
listKeys(): Promise<string[]> { return this.inner.listKeys(); }
exists(key: string): Promise<boolean> { return this.inner.exists(key); }
clear(): Promise<void> { return this.inner.clear(); }
}
/**
* 纯内存的(无页面)SSTable 存储替身:只做"文件 + meta"这两件与 LSM 语义相关的
* 事,用来把 LSM 行为从具体存储实现里隔离出来。返回 `files`/`metas` 供测试直接
* 制造介质损坏或观察 meta 列表。
*/
function newPlainStore() {
const files = new Map<number, Uint8Array>();
const metas: { id: number; level: number; minKey: string; maxKey: string }[] = [];
let seq = 0;
const store = {
async save(id: number, data: Uint8Array) { files.set(id, data); return { storedSize: data.byteLength }; },
async load(id: number) { return files.get(id) ?? null; },
async delete(id: number) { files.delete(id); },
async allocateId() { return ++seq; },
async listMeta() { return metas as never; },
async saveMeta(m: { id: number; level: number; minKey: string; maxKey: string }) { metas.push(m); },
async deleteMeta(id: number) {
const i = metas.findIndex((m) => m.id === id);
if (i >= 0) metas.splice(i, 1);
},
};
return { files, metas, store };
}
/**
* 给一个 Promise 加"失败上限":超时抛错(= 测试失败),而**不会**把超时当成
* 成功结果返回。用于"某件事必须发生 / 必须能完成"的断言 —— 这样并发调度慢
* 只会让测试失败得更慢,绝不会让它在实现错误时碰巧通过。
*/
async function withDeadline<T>(p: Promise<T>, ms: number, label: string): Promise<T> {
let timer: ReturnType<typeof setTimeout> | null = null;
try {
return await Promise.race([
p,
new Promise<never>((_resolve, reject) => {
timer = setTimeout(() => reject(new Error(`timed out after ${ms}ms: ${label}`)), ms);
}),
]);
} finally {
if (timer) clearTimeout(timer);
}
}
/** 让位 n 个宏任务(jsdom 没有 setImmediate):只用于"给并发实现机会",不做时间断言 */
async function yieldMacrotasks(n: number): Promise<void> {
for (let i = 0; i < n; i++) await new Promise<void>((r) => { setTimeout(r, 0); });
}
function openEngine(
dbName: string,
backend: IStorageBackend,
extra: Record<string, unknown> = {},
): AriaEngine {
return new AriaEngine({
storageBackend: 'memory',
checkpointInterval: 100_000_000,
memtableSizeThreshold: 64 * 1024 * 1024,
testBackend: backend,
...extra,
} as never);
}
// ===========================================================================
// 1. ManifestStore 单元语义(真实现,不做替身)
// ===========================================================================
describe('[v0.8.0][B-6] manifest 单一提交点', () => {
function newStore(instanceId = 'inst-A') {
const backend = new MemoryBackend();
const store = new ManifestStore({ backend, instanceId, now: () => 1000 });
return { backend, store };
}
it('提交 → 重新 load 得到同一内容,世代号单调递增', async () => {
const { backend, store } = newStore();
await store.load();
store.current.schemas = { t: { id: { type: 'string', primaryKey: true } } as never };
const first = await store.commit();
expect(first.generation).toBe(1);
expect(first.owner.instanceId).toBe('inst-A');
expect(first.owner.epoch).toBe(1);
store.current.pageIdWatermark = 42;
const second = await store.commit();
expect(second.generation).toBe(2);
const reloaded = await new ManifestStore({ backend }).load();
expect(reloaded.generation).toBe(2);
expect(reloaded.manifest!.pageIdWatermark).toBe(42);
expect(reloaded.manifest!.schemas.t.id.primaryKey).toBe(true);
expect(reloaded.manifest!.owner.instanceId).toBe('inst-A');
});
it('payload 被篡改 → 该世代无效,回退到上一代(并记录损坏世代)', async () => {
const { backend, store } = newStore();
await store.load();
store.current.schemas = { t: { id: { type: 'string' } } as never };
await store.commit();
store.current.schemas = { t: { id: { type: 'string' } }, u: { id: { type: 'string' } } as never };
const newest = await store.commit();
expect(newest.generation).toBe(2);
store.current.schemas = { t: { id: { type: 'string' } }, u: { id: { type: 'string' } } as never, x: {} as never };
await store.commit(); // gen 3
// 篡改最新世代的载荷(保持头部 CRC 有效 → 由载荷 CRC 拦住)
const key = manifestKey(3);
const bytes = new Uint8Array((await backend.read(key))!);
bytes[MANIFEST_HEADER_SIZE + 3] ^= 0xff;
await backend.write(key, bytes.buffer as ArrayBuffer);
const loaded = await new ManifestStore({ backend }).load();
expect(loaded.generation).toBe(2); // 回退到上一代
expect(loaded.hadInvalidGenerations).toBe(true);
expect(loaded.skipped.map((s) => s.generation)).toEqual([3]);
expect((loaded.manifest!.schemas as Record<string, unknown>).x).toBeUndefined();
});
it('全部世代都损坏 → 抛 ARIA_MANIFEST_CORRUPT(绝不返回空状态)', async () => {
const { backend, store } = newStore();
await store.load();
store.current.schemas = { t: { id: { type: 'string' } } as never };
await store.commit();
const key = manifestKey(1);
const bytes = new Uint8Array((await backend.read(key))!);
bytes.fill(0, MANIFEST_HEADER_SIZE); // 抹掉载荷
await backend.write(key, bytes.buffer as ArrayBuffer);
// 变异验证:若把 load() 的"全部无效"改回返回 null(旧行为 = 空库),
// 这条断言立即失败 —— 而"静默空库"正是审计里的 P0 表现。
await expect(new ManifestStore({ backend }).load()).rejects.toMatchObject({
code: 'ARIA_MANIFEST_CORRUPT',
});
});
it('头部自带 CRC:世代号/长度被篡改的世代一律判无效', async () => {
const { store } = newStore();
await store.load();
await store.commit();
const bytes = encodeManifest({ ...store.current, generation: 1 });
// 篡改头部世代号(不改 header CRC)→ 头部 CRC 校验失败
const tampered = new Uint8Array(bytes);
new DataView(tampered.buffer).setUint32(8, 99, false);
const decoded = decodeManifest(tampered);
expect(decoded.ok).toBe(false);
if (!decoded.ok) expect(decoded.reason).toContain('header CRC mismatch');
});
it('陈旧实例提交被拒绝(STALE_INSTANCE),不会静默覆盖新世代', async () => {
const backend = new MemoryBackend();
const a = new ManifestStore({ backend, instanceId: 'A' });
const b = new ManifestStore({ backend, instanceId: 'B' });
await a.load();
await b.load(); // 两个实例都看到 generation = 0
b.current.schemas = { fromB: {} as never };
const bCommitted = await b.commit();
expect(bCommitted.generation).toBe(1);
// A 的内存态已经落后(B 提交过)→ A 的提交必须被拒绝,而不是拿陈旧状态覆盖
a.current.schemas = { staleA: {} as never };
await expect(a.commit()).rejects.toMatchObject({ code: 'STALE_INSTANCE' });
// B 的内容仍然在(陈旧实例没能覆盖)
const loaded = await new ManifestStore({ backend }).load();
expect(Object.keys(loaded.manifest!.schemas)).toEqual(['fromB']);
expect(loaded.manifest!.owner.instanceId).toBe('B');
});
it('提交是"先写后验":回读不可用 → ARIA_MANIFEST_WRITE_FAILED,内存态不前进', async () => {
const backend = new (class extends MemoryBackend {
hideReads = false;
async read(key: string): Promise<ArrayBuffer | null> {
if (this.hideReads && key.startsWith('__aria_manifest_')) return null;
return super.read(key);
}
})();
await backend.open('verify-fail');
const store = new ManifestStore({ backend: backend as unknown as IStorageBackend, instanceId: 'A' });
await store.load();
backend.hideReads = true;
await expect(store.commit()).rejects.toMatchObject({ code: 'ARIA_MANIFEST_WRITE_FAILED' });
expect(store.currentGeneration).toBe(0); // 没有假装提交成功
});
it('保留至少两代;本轮没有损坏世代时才清理更早的世代', async () => {
const { backend, store } = newStore();
await store.load();
for (let i = 0; i < 5; i++) await store.commit();
let gens = (await backend.listKeys()).map((k) => generationFromKey(k)).filter((g) => g !== null);
expect(new Set(gens)).toEqual(new Set([4, 5]));
// 造一个"更新但损坏"的世代(gen 6)→ 加载时会跳过它并记录
const bogus = createEmptyManifest({ instanceId: 'bogus' });
const broken = new Uint8Array(encodeManifest({ ...bogus, generation: 6 }));
broken.fill(0, MANIFEST_HEADER_SIZE);
await backend.write(manifestKey(6), broken.buffer as ArrayBuffer);
const store2 = new ManifestStore({ backend, instanceId: 'A' });
const loaded = await store2.load();
expect(loaded.generation).toBe(5); // 回退到最新有效世代
expect(loaded.hadInvalidGenerations).toBe(true); // 并记录损坏世代
await store2.commit(); // → gen 7
gens = (await backend.listKeys()).map((k) => generationFromKey(k)).filter((g) => g !== null);
// 存在损坏世代时不做任何清理("引用不到的东西一律保留而非删除")
expect(gens).toContain(4);
expect(gens).toContain(5);
expect(gens).toContain(6);
});
});
// ===========================================================================
// 2. 引擎级:数据 → manifest → WAL 截断
// ===========================================================================
describe('[v0.8.0][B-6] 引擎的提交顺序与旧格式迁移', () => {
it('flush 后:SSTable 元数据在 manifest 里,且 WAL 水位只在水位覆盖后才推进', async () => {
const backend = new MemoryBackend();
await backend.open('b6-order');
// 显式开启页面化:这样 meta 里会带 pageIds(页面化是 v0.8.0 的默认路径)
const engine = openEngine('b6-order', backend, { memtableSizeThreshold: 2048, pageStorage: true });
await engine.open('b6-order', 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(40));
const before = await readManifestState(backend);
expect(before).not.toBeNull();
// 未显式 flush 时,数据可能还在 memtable → 水位不得推进到"最新 LSN"
const walLsnBefore = (engine as any).wal.getLsn() as number;
expect(before!.wal.startLsn).toBeLessThanOrEqual(walLsnBefore);
await (engine as any).lsm.flush();
const after = await readManifestState(backend);
expect(after!.namespaces.main.sstables.length).toBeGreaterThan(0);
expect(after!.wal.startLsn).toBeGreaterThanOrEqual(before!.wal.startLsn); // 单调
expect(after!.wal.nextLsn).toBeGreaterThanOrEqual(after!.wal.startLsn);
// meta 里必须带页面映射(页面化路径的数据定位依据)
expect(after!.namespaces.main.sstables[0].pageIds?.length ?? 0).toBeGreaterThan(0);
await engine.close();
});
it('旧格式(__aria_lsm_meta/__aria_schemas)库打开后完整迁移,旧键保留', async () => {
const dbName = uniqueDB('b6-legacy');
const engine = new AriaEngine({ storageBackend: 'opfs', memtableSizeThreshold: 64 * 1024 * 1024, checkpointInterval: 100_000_000 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(20));
await (engine as any).lsm.flush();
await engine.close();
// 把 v0.8.0 布局"降级"成旧布局:写回裸 JSON meta / schema,删掉 manifest
const backend = new OPFSBackend();
await backend.open(dbName);
const manifest = (await readManifestState(backend))!;
const mainMetas = manifest.namespaces.main.sstables;
await backend.write('__aria_lsm_meta', new TextEncoder().encode(JSON.stringify(mainMetas)).buffer);
await backend.write('__aria_schemas', new TextEncoder().encode(JSON.stringify(manifest.schemas)).buffer);
const pageMeta = new ArrayBuffer(8);
new DataView(pageMeta).setUint32(0, manifest.pageIdWatermark, false);
await backend.write('__aria_meta', pageMeta);
await backend.deleteMany(
(await backend.listKeys()).filter((k) => k.startsWith('__aria_manifest_')),
);
await backend.close();
const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100_000_000 });
await engine2.open(dbName, 1);
expect(engine2.getRecoveryReport().legacyImported).toBe(true);
expect(await engine2.count('t')).toBe(20);
await engine2.close();
const backend2 = new OPFSBackend();
await backend2.open(dbName);
const migrated = await readManifestState(backend2);
expect(migrated).not.toBeNull();
expect(migrated!.namespaces.main.sstables.length).toBe(mainMetas.length);
expect(Object.keys(migrated!.schemas)).toEqual(['t']);
// 旧键**保留**("引用不到的东西一律保留而非删除")
expect(await backend2.exists('__aria_lsm_meta')).toBe(true);
await backend2.close();
});
it('旧格式 meta 损坏 → 抛 ARIA_LEGACY_META_CORRUPT(修复前是静默空库)', async () => {
const dbName = uniqueDB('b6-legacy-corrupt');
const backend = new OPFSBackend();
await backend.open(dbName);
await backend.write('__aria_lsm_meta', new TextEncoder().encode('{not json').buffer);
await backend.close();
const engine = new AriaEngine({ storageBackend: 'opfs' });
await expect(engine.open(dbName, 1)).rejects.toMatchObject({ code: 'ARIA_LEGACY_META_CORRUPT' });
});
it('最新世代损坏 → 回退上一代并标记 manifestFallback,数据仍可读', async () => {
const dbName = uniqueDB('b6-fallback');
const engine = new AriaEngine({ storageBackend: 'opfs', memtableSizeThreshold: 64 * 1024 * 1024, checkpointInterval: 100_000_000 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(15));
await engine.close();
const backend = new OPFSBackend();
await backend.open(dbName);
const gens = (await backend.listKeys())
.map((k) => generationFromKey(k)).filter((g): g is number => g !== null).sort((a, b) => b - a);
const newest = new Uint8Array((await backend.read(manifestKey(gens[0])))!);
newest.fill(0, MANIFEST_HEADER_SIZE);
await backend.write(manifestKey(gens[0]), newest.buffer as ArrayBuffer);
await backend.close();
const engine2 = new AriaEngine({ storageBackend: 'opfs', checkpointInterval: 100_000_000 });
await engine2.open(dbName, 1);
expect(engine2.getRecoveryReport().manifestFallback).toBe(true);
expect(await engine2.count('t')).toBe(15);
await engine2.close();
});
it('介质读故障 ≠ 文件缺失:读取报 ARIA_SSTABLE_READ_FAILEDmeta 不被自愈删除', async () => {
const dbName = uniqueDB('b6-readfail');
const base = new MemoryBackend();
// 先正常写入并落盘(写入路径不带故障注入)
const writer = openEngine(dbName, base, { memtableSizeThreshold: 2048 });
await writer.open(dbName, 1);
await writer.createTable(SCHEMA());
await writer.insert('t', rows(30));
await (writer as any).lsm.flush();
// 注意:**不能** close —— MemoryBackend.close() 会清空介质(等价于删库),
// 这里要的是"同一个介质上换一个引擎实例"。
// 新引擎(缓存为空 → 读取必然回源)打开时会重放 WAL 并再落一次盘,
// 因此 meta 基线要在打开**之后**取
const gated = new GatedBackend(base);
const reader = openEngine(dbName, gated, { memtableSizeThreshold: 2048 });
await reader.open(dbName, 1);
const metasBefore = (await readManifestState(base))!.namespaces.main.sstables.length;
expect(metasBefore).toBeGreaterThan(0);
expect(await reader.count('t')).toBe(30); // 此时还能读(缓存已建立)
gated.failReadPattern = 'sst_'; // 让后续回源读失败
(reader as any).lsm.sstableCache.clear(); // 清缓存 → 强制回源
(reader as any).lsm.cacheSize = 0;
(reader as any).lsm.oversizedSSTables.clear();
await expect(reader.find('t', { table: 't' })).rejects.toMatchObject({
code: 'ARIA_SSTABLE_READ_FAILED',
});
// 修复前的行为:读故障被当成"文件不存在" → dropInvalidSSTable 删掉 meta(不可逆)
const metasAfter = (await readManifestState(base))!.namespaces.main.sstables.length;
expect(metasAfter).toBe(metasBefore);
// 故障消失后数据完整(meta 没被删,文件也还在)
gated.failReadPattern = null;
(reader as any).lsm.sstableCache.clear();
(reader as any).lsm.cacheSize = 0;
expect(await reader.count('t')).toBe(30);
});
it('恢复报告在无损坏时为"干净",且 repair 只在干净时回收孤儿页面', async () => {
const dbName = uniqueDB('b6-repair-gate');
const base = new MemoryBackend();
const gated = new GatedBackend(base);
const engine = openEngine(dbName, gated);
await engine.open(dbName, 1);
const report = engine.getRecoveryReport();
expect(report.dataLossSuspected).toBe(false);
expect(report.walGaps).toEqual([]);
await engine.close();
});
});
// ===========================================================================
// 3. WAL:分片空洞必须显式失败(不再静默丢尾部)
// ===========================================================================
describe('[v0.8.0][B-6] WAL 分片空洞与按水位截断', () => {
async function buildSegmentedWAL(segments: number) {
const backend = new MemoryBackend();
await backend.open('wal-gap');
const store = new SegmentedWALStore(backend, 96); // 极小分片 → 每条记录独占一片
const wal = new WAL(store, true, 'full');
for (let i = 0; i < segments; i++) {
await wal.append({
type: WALRecordType.INSERT,
txnId: 0,
tableName: 't',
key: `k${i}`,
data: { v: i, pad: 'p'.repeat(48) },
});
}
// 断言前提:确实一片一条(否则下面的"缺片/边界"断言就失去意义)
const segs = (await backend.listKeys()).filter((k) => /^__wal_\d+\.bin$/.test(k)).sort();
expect(segs).toHaveLength(segments);
return { backend, store, wal };
}
it('活跃区间内缺失分片 → 默认抛 ARIA_WAL_GAP(不再静默丢弃尾部)', async () => {
const { backend, wal } = await buildSegmentedWAL(3);
await backend.delete('__wal_000001.bin'); // 中间挖掉一片(一片一条记录)
const seen: unknown[] = [];
await expect(wal.recover((r) => seen.push(r))).rejects.toMatchObject({ code: 'ARIA_WAL_GAP' });
expect(seen).toEqual([]); // 抛错时不产生"部分恢复"的副作用
});
it('前缀缺失同样算空洞;显式 allowGaps 时如实上报', async () => {
const { backend, wal } = await buildSegmentedWAL(3);
await backend.delete('__wal_000000.bin'); // 缺失的是**前缀**(一片一条记录)
await expect(wal.recover(() => { /* noop */ })).rejects.toMatchObject({ code: 'ARIA_WAL_GAP' });
const seen: number[] = [];
const applied = await wal.recover((r) => seen.push(r.lsn), { allowGaps: true });
expect(applied).toBe(seen.length);
expect(wal.getLastRecoveryInfo()!.gaps).toEqual([0]);
});
it('按水位截断:边界保守(不误删可能含新记录的分片),水位越过全部记录后整段可删', async () => {
const { backend, store, wal } = await buildSegmentedWAL(3);
const latest = wal.getLsn(); // = 3
// 水位 = 1:分片 0 自身可能还含 > 1 的记录(分片内 LSN 无法从未条推出)
// → 必须整体保留(保守但绝不丢记录)
expect(await store.planKeepFrom(1, latest)).toBe(0);
// 水位 = 2:分片 0 已被完整覆盖(分片 1 的首条 = 2),分片 1 仍保守保留
expect(await store.planKeepFrom(2, latest)).toBe(1);
// 水位 = 当前高水位:全部分片整体落盘 → 全部可删
expect(await store.planKeepFrom(latest, latest)).toBe(3);
await store.truncateBefore(latest, latest);
const keys = await backend.listKeys();
expect(keys.filter((k) => k.startsWith('__wal_'))).toEqual([]);
// 截断后分片号必须继续往前(不得复用),否则新记录会被 manifest 的
// startSegment 跳过(实测:删除的行复活 / 已确认写入丢失)
await wal.append({
type: WALRecordType.INSERT, txnId: 0, tableName: 't', key: 'after', data: { v: 9 },
});
const afterKeys = (await backend.listKeys()).filter((k) => /^__wal_\d+\.bin$/.test(k));
expect(afterKeys).toEqual(['__wal_000003.bin']);
});
it('恢复按 lsn 跳过已落盘记录(旧记录不会把已删除的行复活)', async () => {
const { wal } = await buildSegmentedWAL(3);
const seen: number[] = [];
const applied = await wal.recover((r) => seen.push(r.lsn), { fromLsn: 2 });
expect(seen).toEqual([3]);
expect(applied).toBe(1);
expect(wal.getLastRecoveryInfo()!.skipped).toBe(2);
// LSN 高水位不受跳过影响(必须继续单调,否则下一批记录会与历史 LSN 冲突)
expect(wal.getLsn()).toBe(3);
});
});
// ===========================================================================
// 4. LSM 结构根治(44/45/47/49/50/51/55
// ===========================================================================
describe('[v0.8.0][B-6] LSMflush 重试与错误报告顺序', () => {
it('注入一次 save 失败 → 重试后数据真的落盘(修复前失败的表再无落盘机会)', async () => {
const dbName = uniqueDB('b6-flush-retry');
const base = new MemoryBackend();
const failOnce = new FailOnceBackend(base);
const engine = openEngine(dbName, failOnce as unknown as IStorageBackend, { memtableSizeThreshold: 2048 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(30));
const lsm: any = (engine as any).lsm;
// flush 之后再 flush 一次以确保所有冻结表都有落盘机会
await lsm.flush();
await lsm.flush();
expect(lsm.hasPendingFlushData()).toBe(false);
expect((await lsm.sstableStore.listMeta()).length).toBeGreaterThan(0);
// 瞬时故障必须可见(而不是静默吞掉)
expect(failOnce.remainingFailures).toBe(0); // 注入的失败确实发生过
expect(lsm.getBackgroundWarnings().length).toBeGreaterThan(0);
// 不 closeMemoryBackend.close 会清空介质);直接在同一介质上重开验证
const engine2 = new AriaEngine({ storageBackend: 'memory', testBackend: base, checkpointInterval: 100_000_000 } as never);
await engine2.open(dbName, 1);
expect(await engine2.count('t')).toBe(30);
});
it('持续失败 → flush 明确报错且数据不丢;故障恢复后 flush 成功', async () => {
let failing = true;
let saved = 0;
const files = new Map<number, Uint8Array>();
const metas: { id: number; level: number; minKey: string; maxKey: string }[] = [];
let seq = 0;
const store = {
async save(id: number, data: Uint8Array) {
if (failing) throw new Error('injected persistent save failure');
saved++;
files.set(id, data);
return { storedSize: data.byteLength };
},
async load(id: number) { return files.get(id) ?? null; },
async delete(id: number) { files.delete(id); },
async allocateId() { return ++seq; },
async listMeta() { return metas as never; },
async saveMeta(m: { id: number; level: number; minKey: string; maxKey: string }) { metas.push(m); },
async deleteMeta() { /* noop */ },
};
const lsm = new LSM({ memtableSizeThreshold: 128, sstableStore: store as never });
for (let i = 0; i < 12; i++) lsm.put(`k${i}`, { v: i });
lsm.freezeMemtable();
await lsm.whenIdle(); // 确定性等待:门控"后台链已跑完"而不是"睡够 20ms"
await expect(lsm.flush()).rejects.toMatchObject({ code: 'ARIA_BACKGROUND_ERROR' });
// 数据没有消失:冻结表仍在、且仍可读
expect(lsm.getStats().frozenTables).toBeGreaterThan(0);
expect(await lsm.get('k0')).toEqual({ v: 0 });
// 故障消失 → flush 成功、数据落盘
failing = false;
await lsm.flush();
expect(lsm.getStats().frozenTables).toBe(0);
expect(saved).toBeGreaterThan(0);
expect(metas.length).toBeGreaterThan(0);
});
it('后台错误不得让本次 flush 白做:错误在"入链完成"之后才报告', async () => {
// 场景:后台自动 flush 失败过一次(错误已记录),随后介质恢复正常。
// 旧实现:下一次 flush() **先检查错误** → 直接抛错,本次 flush 什么都没做
//(待落盘数据继续只在内存里,而调用方以为它已经被处理过)。
// 新实现:先把 memtable 入链、重试、真正落盘,然后才报告/记录那次失败。
let failuresLeft = 1; // 第一次 save 失败,之后正常
let saved = 0;
const files = new Map<number, Uint8Array>();
const metas: { id: number; level: number }[] = [];
let seq = 0;
const store = {
async save(id: number, data: Uint8Array) {
if (failuresLeft > 0) { failuresLeft--; throw new Error('injected background save failure'); }
saved++;
files.set(id, data);
return { storedSize: data.byteLength };
},
async load(id: number) { return files.get(id) ?? null; },
async delete(id: number) { files.delete(id); },
async allocateId() { return ++seq; },
async listMeta() { return metas as never; },
async saveMeta(m: { id: number; level: number }) { metas.push(m); },
async deleteMeta() { /* noop */ },
};
const lsm = new LSM({ memtableSizeThreshold: 64, sstableStore: store as never });
// 触发后台自动 flushfreezeMemtable → 入链 → 第一次 save 失败 → lastBackgroundError
for (let i = 0; i < 20; i++) lsm.put(`k${i}`, { v: i });
lsm.freezeMemtable();
await lsm.whenIdle(); // 确定性等待后台链跑完(不靠 sleep)
expect(failuresLeft).toBe(0); // 注入的失败确实发生过
// 关键:这一次 flush 必须**真的执行**(早退的话什么都写不出去)
await expect(lsm.flush()).resolves.toBeUndefined();
expect(saved).toBeGreaterThan(0);
expect(lsm.getStats().frozenTables).toBe(0);
expect(metas.length).toBeGreaterThan(0);
// 失败被重试修复 → 必须留下可见记录(不是静默吞掉)
expect(lsm.getBackgroundWarnings().length).toBeGreaterThan(0);
// 数据完整
expect(await lsm.get('k0')).toEqual({ v: 0 });
});
});
describe('[v0.8.0][B-6/47] 流式扫描提前终止不多算', () => {
it('消费者只取 5 条时,底层来源恰好被拉取 5 条(修复前是 6 条)', () => {
const produced: number[] = [];
let pulls = 0;
const source: EntrySource = {
next() {
pulls++;
if (produced.length >= 20) return null;
const i = produced.length;
produced.push(i);
return [`k${String(i).padStart(3, '0')}`, { v: i }];
},
reset() { /* noop */ },
};
const merge = new MergeIterator();
merge.addSource(source);
const out: string[] = [];
for (let i = 0; i < 5; i++) {
const entry = merge.next();
if (entry) out.push(entry[0]);
}
expect(out).toHaveLength(5);
// 变异验证:把 MergeIterator.next() 里的延迟补充改回"立即 seedFromSource"
// pulls 变成 6 —— 这正是审计记录的"limit=5 却多算 1 条"。
expect(pulls).toBe(5);
});
it('引擎层 limit=5 的流式扫描只产出 5 行(且提前终止)', async () => {
const dbName = uniqueDB('b6-stream-limit');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend);
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(50));
await (engine as any).lsm.flush();
const seen: Record<string, unknown>[] = [];
const count = await engine.findStream('t', { table: 't', limit: 5 }, (row) => { seen.push(row); });
expect(count).toBe(5);
expect(seen).toHaveLength(5);
await engine.close();
});
});
describe('[v0.8.0][B-6/49] 按层 compaction 状态(不再静默丢弃跨层触发)', () => {
it('两层同时被触发时都进入"正在 compaction"集合(单 boolean 会丢掉第二个)', async () => {
const files = new Map<number, Uint8Array>();
let seq = 0;
const store = {
async save(id: number, data: Uint8Array) { files.set(id, data); return { storedSize: data.byteLength }; },
async load(id: number) { return files.get(id) ?? null; },
async delete(id: number) { files.delete(id); },
async allocateId() { return ++seq; },
async listMeta() { return [] as never; },
async saveMeta() { /* noop */ },
async deleteMeta() { /* noop */ },
};
const lsm = new LSM({ memtableSizeThreshold: 256, sstableStore: store as never });
for (let i = 0; i < 20; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
// 同步连续触发两层:Set 语义下两个层号都在集合里;
// 变异验证:把 `compacting` 换回单 boolean,第二个层号会被静默丢弃。
(lsm as unknown as { scheduleCompact(l: number): void }).scheduleCompact(0);
(lsm as unknown as { scheduleCompact(l: number): void }).scheduleCompact(1);
expect(lsm.getStats().compactingLevels).toEqual([0, 1]);
await lsm.flush();
expect(lsm.getStats().compactingLevels).toEqual([]);
});
});
describe('[v0.8.0][B-6/50] compaction 期间该层对读者始终可见', () => {
it('合并尚未提交时,全表扫描仍能看到该层全部数据(修复前整层被 splice 掉)', async () => {
const gated = new GatedBackend(new MemoryBackend());
const backendMap = new Map<number, Uint8Array>();
const store = {
async save(id: number, data: Uint8Array) {
backendMap.set(id, data);
// 第二次 savecompaction 产物)会命中门控 → 长 await 窗口
await gated.write(`sst_${id}`, data.buffer.slice(data.byteOffset, data.byteOffset + data.byteLength) as ArrayBuffer);
return { storedSize: data.byteLength };
},
async load(id: number) { return backendMap.get(id) ?? null; },
async delete(id: number) { backendMap.delete(id); },
async allocateId() { return ++allocSeq; },
async listMeta() { return [] as never; },
async saveMeta() { /* noop */ },
async deleteMeta() { /* noop */ },
};
let allocSeq = 0;
const lsm: any = new LSM({ memtableSizeThreshold: 64, sstableStore: store as never });
for (let i = 0; i < 40; i++) lsm.put(`k${i}`, { v: i });
// 后台 flush 产生产物(可能已经触发 compaction,等它静默)
await lsm.flush();
const before = await lsm.rangeScan('', '\uffff');
expect(before.length).toBeGreaterThan(0);
// 门控住 compaction 的产物写入 → 此时它已经取好源快照但还没提交
let gateHit: (() => void) | null = null;
const hit = new Promise<void>((resolve) => { gateHit = resolve; });
gated.onGated = () => gateHit?.();
gated.armGate('sst_');
const compaction = lsm.compactLevel(0);
await hit;
// 核心不变量:compaction 的 await 窗口内,level 0 仍然完整可见
const during = await lsm.rangeScan('', '\uffff');
expect(during.map((e: [string, unknown]) => e[0])).toEqual(before.map((e: [string, unknown]) => e[0]));
gated.release();
await compaction;
await lsm.flush();
const after = await lsm.rangeScan('', '\uffff');
expect(after.length).toBe(before.length);
});
it('并发 flush 与 compaction 交错时,扫描不会漏掉刚发布的行(结构版本重试)', async () => {
// 确定性构造那个危险窗口:读路径"取 meta 快照 → await 加载文件",而一次后台
// flush 正好在这个 await 期间**发布新 SSTable 并从 frozenMemtables 摘掉它**。
// 修复前这次扫描既看不到新 SSTable(不在快照里)也读不到前台冻结表 → 少行。
const files = new Map<number, Uint8Array>();
const metas: { id: number; level: number; minKey: string; maxKey: string }[] = [];
let seq = 0;
let pendingFlush: (() => Promise<void>) | null = null;
let triggered = false;
let loads = 0;
const store = {
async save(id: number, data: Uint8Array) {
files.set(id, data);
return { storedSize: data.byteLength };
},
async load(id: number) {
loads++;
// 关键:在"读回源"的 await 里让后台 flush 跑完(发布新 SSTable + 摘冻结表)
if (pendingFlush && !triggered) {
triggered = true;
const flush = pendingFlush;
pendingFlush = null;
await flush();
}
return files.get(id) ?? null;
},
async delete(id: number) { files.delete(id); },
async allocateId() { return ++seq; },
async listMeta() { return metas as never; },
async saveMeta(m: { id: number; level: number; minKey: string; maxKey: string }) { metas.push(m); },
async deleteMeta() { /* noop */ },
};
const lsm: any = new LSM({ memtableSizeThreshold: 4096, sstableStore: store as never });
// 先落一个 SSTable(保证扫描必须回源 → 进入上面的钩子)
for (let i = 0; i < 5; i++) lsm.put(`a${i}`, { v: i });
await lsm.flush();
expect(metas.length).toBeGreaterThan(0);
(lsm as any).sstableCache.clear();
(lsm as any).cacheSize = 0;
(lsm as any).oversizedSSTables.clear();
// 后续写入的新行:扫描开始后它们才会被"后台 flush"发布出去
for (let i = 0; i < 5; i++) lsm.put(`b${i}`, { v: 100 + i });
pendingFlush = () => lsm.flushMemtablesOnly();
const during = await lsm.rangeScan('', '\uffff');
expect(triggered).toBe(true); // 钩子确实在扫描中途触发了 flush
// 关键断言:这次扫描必须看到**全部 10 行**(既不能漏刚发布的那批,
// 也不能因为重试而重复)
expect(during.map((e: [string, unknown]) => e[0]).sort()).toEqual(
['a0', 'a1', 'a2', 'a3', 'a4', 'b0', 'b1', 'b2', 'b3', 'b4'],
);
// 点查(get)也有同一份结构版本校验:这里用"回源期间前台结构发生变化"
// (写入 + 冻结)来触发它,并断言
// ① 结果仍然正确(返回快照里的值,不因结构变化而丢数据/报错);
// ② 确实发生了一次重试(同一个 SSTable 被读了两次)。
// 变异验证:去掉 get 里的版本校验后 loads 只增加 1 → 断言失败。
(lsm as any).sstableCache.clear();
(lsm as any).cacheSize = 0;
(lsm as any).oversizedSSTables.clear();
triggered = false;
const retriesBefore = lsm.getStats().readStructureRetries;
const loadsBefore = loads;
void loadsBefore;
pendingFlush = async () => {
lsm.put('zzz', { v: 1 });
lsm.freezeMemtable(); // 前台结构变化(活跃 memtable → frozen 列表)
};
expect(await lsm.get('a0')).toEqual({ v: 0 });
expect(triggered).toBe(true);
// 结构版本校验确实生效:读取期间发生前台变化 → 重试计数 +1
//(变异验证:去掉 get 里的版本校验后计数不变 → 断言失败)
expect(lsm.getStats().readStructureRetries).toBeGreaterThan(retriesBefore);
expect(await lsm.get('zzz')).toEqual({ v: 1 });
});
});
describe('[v0.8.0][B-6/51] 底部层合并回收墓碑', () => {
it('删除密集场景:底部层合并后墓碑不再无限累积(重开后数据正确)', async () => {
const dbName = uniqueDB('b6-tombstone-gc');
const base = new MemoryBackend();
const engine = openEngine(dbName, base, { memtableSizeThreshold: 1024 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(120));
await (engine as any).lsm.flush();
await engine.delete('t', { table: 't' });
await (engine as any).lsm.flush();
const lsm: any = (engine as any).lsm;
const countTombstones = async (): Promise<number> => {
const metas = await lsm.sstableStore.listMeta();
let tombstones = 0;
for (const meta of metas) {
const data = await lsm.sstableStore.load(meta.id);
if (!data) continue;
new SSTableReader(data, meta).scanAll((_k, v) => {
if ((v as Record<string, unknown>).__tombstone) tombstones++;
});
}
return tombstones;
};
// 把数据/墓碑一路压到**底部层**:每层用 minFiles=1 逐级下沉
const compact = (level: number, minFiles?: number) =>
(lsm as unknown as { compactLevelAsync(l: number, m?: number): Promise<void> })
.compactLevelAsync(level, minFiles);
for (let level = 0; level < MAX_LSM_LEVELS - 1; level++) {
// 反复压,直到该层没有文件为止(每轮产物落到下一层)
for (let round = 0; round < 8 && lsm.levels[level].length > 0; round++) {
await compact(level, 1);
}
}
const bottomBefore = lsm.levels[MAX_LSM_LEVELS - 1].length;
expect(bottomBefore).toBeGreaterThan(0); // 数据确实到了底部层
const beforeGc = await countTombstones();
expect(beforeGc).toBeGreaterThan(0); // 底部层合并前墓碑确实存在
// 触发底部层原地合并(drop tombstones
await compact(MAX_LSM_LEVELS - 1, 1);
const afterGc = await countTombstones();
// 变异验证:去掉 isBottomLevel 的墓碑过滤,这里会等于 beforeGc
expect(afterGc).toBe(0);
// 语义不变:删除仍然生效(没有因为丢墓碑而复活)
expect(await engine.count('t')).toBe(0);
// 不 closeMemoryBackend.close 会清空介质);同一介质上重开验证持久化
const engine2 = new AriaEngine({ storageBackend: 'memory', testBackend: base, checkpointInterval: 100_000_000 } as never);
await engine2.open(dbName, 1);
expect(await engine2.count('t')).toBe(0);
});
});
describe('[v0.8.0][B-6/55] checkpoint 不再等完整 compaction', () => {
it('compaction 卡住时,写路径的周期 checkpoint 仍能完成并落盘', async () => {
const dbName = uniqueDB('b6-tick-decouple');
const base = new MemoryBackend();
const gated = new GatedBackend(base);
const engine = openEngine(dbName, gated, {
memtableSizeThreshold: 512,
checkpointInterval: 1, // 每次写入都尝试 checkpoint
});
await engine.open(dbName, 1);
// 带 pad 列:每行都超过 memtable 阈值 → 每次 insert 产生一个 level-0 SSTable
await engine.createTable(createSchema('t', {
id: { type: 'string', primaryKey: true },
v: { type: 'number' },
pad: { type: 'string' },
}));
const lsm: any = (engine as any).lsm;
const big = 'x'.repeat(600);
// 先造 3 个 level-0 SSTable3 < 4 → 不会自动触发 compaction
for (let i = 0; i < 3; i++) {
await engine.insert('t', [{ id: `k${i}`, v: i, pad: big }]);
await lsm.flushMemtablesOnly();
}
expect(lsm.levels[0].length).toBe(3);
// 门控:放行第 4 次 flush 自己的 save,阻塞随后自动触发的 compaction 的 save
let gateHit: (() => void) | null = null;
const hit = new Promise<void>((resolve) => { gateHit = resolve; });
gated.onGated = () => gateHit?.();
gated.armGate('sst_', 1);
await engine.insert('t', [{ id: 'k3', v: 3, pad: big }]);
await lsm.flushMemtablesOnly();
// 门控信号本身就是"compaction 已开始且卡在写盘上"的证据;墙钟只作为失败上限
// (超时 = 测试失败),因此机器再慢也不会"碰巧通过"。
await withDeadline(hit, 15000, '自动 compaction 始终没有开始(门控从未命中)');
// compaction 现在卡在门控上(维护链不前进)。写路径必须继续工作:
// 修复前 checkpoint → lsm.flush() → drainMaintenance() 会一起卡住
//v0.6.1 记录的"8~11s 悬崖"的成因之一)。
await withDeadline(
engine.insert('t', [{ id: 'after-gate', v: 999, pad: big }]),
15000,
'写路径被卡住的维护链阻塞(checkpoint 又和 compaction 串在一起了)',
);
gated.release();
await lsm.flush();
expect(await engine.count('t')).toBe(5);
});
});
describe('[v0.8.0][B-6/review] vacuum 与后台 compaction 严格串行', () => {
it('vacuumLevels 的逐层压缩必须挂在维护链上(同一时刻只允许一个 SSTable 写在飞)', async () => {
const dbName = uniqueDB('b6-vacuum-serial');
const base = new MemoryBackend();
const gated = new GatedBackend(base);
const engine = openEngine(dbName, gated, {
memtableSizeThreshold: 512,
checkpointInterval: 100_000_000,
});
await engine.open(dbName, 1);
await engine.createTable(createSchema('t', {
id: { type: 'string', primaryKey: true },
v: { type: 'number' },
pad: { type: 'string' },
}));
const lsm: any = (engine as any).lsm;
const big = 'x'.repeat(600);
// 3 个 level-0 文件(不到触发阈值),第 4 个触发后台 compaction
for (let i = 0; i < 3; i++) {
await engine.insert('t', [{ id: `k${i}`, v: i, pad: big }]);
await lsm.flushMemtablesOnly();
}
expect(lsm.levels[0].length).toBe(3);
// 粘性门控:放行第 4 次 flush 自己的 save,之后**所有** SSTable 写都阻塞
let gateHit: (() => void) | null = null;
const hit = new Promise<void>((resolve) => { gateHit = resolve; });
gated.onGated = () => gateHit?.();
gated.armStickyGate('sst_', 1);
await engine.insert('t', [{ id: 'k3', v: 3, pad: big }]);
await lsm.flushMemtablesOnly();
await withDeadline(hit, 15000, '后台 compaction 没有开始');
// 后台 compaction 卡在写盘上时发起逐层压缩。
//
// 这里直接测 `lsm.vacuumLevels()`(而不是 `engine.vacuum()`):引擎层 `vacuum()`
// 开头的 `lsm.flush()` 自身会 drainMaintenance,因此"引擎层调用"这条路径恰好
// 被 flush 顺带串行化了 —— 真正需要守住的不变量在 LSM 这一层:**逐层压缩必须
// 作为维护链任务执行**,否则它会绕过链与正在跑的 compaction 并发写介质。
const vacuuming = lsm.vacuumLevels();
// 有限次让位(宏任务):修复前 vacuum 的写必然在这期间进入介质(它不依赖维护链)
await yieldMacrotasks(50);
expect(gated.gatedPeak).toBe(1); // 绝不允许两个 compaction 同时写
gated.release();
await withDeadline(vacuuming, 15000, 'vacuumLevels 在维护链释放后仍未完成');
await lsm.whenIdle();
expect(gated.gatedPeak).toBe(1);
expect(await engine.count('t')).toBe(4);
const finalRows = await engine.find('t', { table: 't' });
expect(finalRows.map((r) => `${r.id}=${r.v}`).sort()).toEqual(['k0=0', 'k1=1', 'k2=2', 'k3=3']);
});
});
describe('[v0.8.0][B-6] 冻结表意图:已确认写入不得被水位越过', () => {
it('manifest 声称有未落盘冻结表、但 WAL 里什么都没有 → 打开时报 ARIA_WRITE_LOST', async () => {
const dbName = uniqueDB('b6-intent-lost');
const backend = new MemoryBackend();
await backend.open(dbName);
// 手工构造:一份声称"有冻结表"的 manifest + 没有任何 WAL 记录
const store = new ManifestStore({ backend, instanceId: 'crafted' });
await store.load();
const crafted: AriaManifest = {
...store.current,
frozen: [{
ns: 'main',
id: 1,
entryCount: 7,
minKey: 't:k1',
maxKey: 't:k7',
lsnAtFreeze: 3,
}],
wal: { startSegment: 0, startLsn: 3, nextLsn: 5 },
};
const bytes = encodeManifest({ ...crafted, generation: 1 });
await backend.write(manifestKey(1), bytes.buffer.slice(bytes.byteOffset, bytes.byteOffset + bytes.byteLength) as ArrayBuffer);
const engine = new AriaEngine({ storageBackend: 'memory', testBackend: backend, checkpointInterval: 100_000_000 } as never);
await expect(engine.open(dbName, 1)).rejects.toMatchObject({ code: 'ARIA_WRITE_LOST' });
});
it('冻结表存在时,水位严格低于意图起点(崩溃后这些写入必须能重建)', async () => {
const dbName = uniqueDB('b6-intent-floor');
const backend = new MemoryBackend();
// 阈值调大:数据留在 memtable 里,随后手动冻结(模拟"已确认写入但还没落盘")
const engine = openEngine(dbName, backend, { memtableSizeThreshold: 64 * 1024 * 1024 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(20));
(engine as any).lsm.freezeMemtable();
// 有意让意图进入 manifest(提交一次):水位必须**严格低于**意图起点 ——
// 恢复时 `lsn <= startLsn` 的记录会被跳过,取等号就会把这批记录判成"已落盘"。
await (engine as any).commitManifest();
const after = (await readManifestState(backend))!;
const intent = after.frozen.find((f) => f.ns === 'main');
expect(intent).toBeDefined();
expect(intent!.entryCount).toBeGreaterThan(0); // 意图必须真的有内容
expect(intent!.lsnAtFreeze).toBeGreaterThan(0);
expect(after.wal.startLsn).toBeLessThan(intent!.lsnAtFreeze);
// 崩溃语义(不 flush、不 close):重开时必须靠 WAL 重建这 20 行,且报告干净
const engine2 = openEngine(dbName, backend, { memtableSizeThreshold: 64 * 1024 * 1024 });
await engine2.open(dbName, 1);
expect(await engine2.count('t')).toBe(20);
expect(engine2.getRecoveryReport().dataLossSuspected).toBe(false);
});
});
describe('[v0.8.0][B-6] 页面映射:已退休的 SSTable 仍可被在途读者读取', () => {
it('compaction 摘除 meta 后,旧页面映射保留到物理删除', async () => {
const dbName = uniqueDB('b6-retire');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { memtableSizeThreshold: 1024 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 60; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
const lsm: any = (engine as any).lsm;
await lsm.flush();
const before = await lsm.sstableStore.listMeta();
await lsm.compactLevel(0);
// 被取代的 meta 已从 manifest 摘除
const after = await readManifestState(backend);
const afterIds = new Set(after!.namespaces.main.sstables.map((m) => m.id));
expect(before.some((m: { id: number }) => !afterIds.has(m.id))).toBe(true);
// 但数据仍完整可读(退休文件在"没有更早读者"之前不物理删除)
expect(await engine.count('t')).toBe(60);
await engine.close();
});
});
describe('[v0.8.0][B-6] SSTable 解析:三个读取路径越界策略一致', () => {
it('三个读取路径对同一个损坏文件给出一致结论(修复前三份实现各自为政)', () => {
const { SSTableBuilder } = require('../src/engine/aria/index/sstable_builder');
const builder = new SSTableBuilder(64); // 小块 → 多个块
const total = 30;
for (let i = 0; i < total; i++) builder.add(`k${String(i).padStart(3, '0')}`, { v: i });
const { sstableData } = builder.build();
const meta = {
id: 1, level: 0, minKey: 'k000', maxKey: 'k029',
blockCount: 1, totalSize: sstableData.byteLength, bloomData: null,
} as never;
// 健康读取器:用于确定"第一个块覆盖到哪个 key"
const healthy = new SSTableReader(sstableData, meta);
const indexEntries = (healthy as unknown as { indexEntries: { key: string; blockOffset: number; blockSize: number }[] })
.indexEntries;
expect(indexEntries.length).toBeGreaterThan(1); // 前提:确实有多个块
const firstBlockLastKey = indexEntries[0].key;
const healthyKeys: string[] = [];
healthy.scanAll((k) => { healthyKeys.push(k); });
const expectedAfterFirstBlock = healthyKeys.filter((k) => k > firstBlockLastKey);
expect(expectedAfterFirstBlock.length).toBeGreaterThan(0);
// 破坏第一个块里**第一个条目的 keyLen 字段**(块布局:[u32 条目数][u32 keyLen]...):
// 长度字段越界 → 按统一策略"本块剩余条目整体放弃、继续后续块"。
// (注意不能只改块尾部:尾部是最后一个条目的**值**字节,坏掉只会让那条记录
// 解析失败,不会触发越界分支 —— 那样就测不到三条路径的策略是否一致。)
const damaged = new Uint8Array(sstableData);
const block = indexEntries[0];
damaged.fill(0xff, block.blockOffset + 4, block.blockOffset + 8);
const damagedReader = new SSTableReader(damaged, meta);
// 三条读取路径都不许抛异常(统一越界策略 = 跳过越界部分)
const viaScanAll: string[] = [];
expect(() => damagedReader.scanAll((k) => { viaScanAll.push(k); })).not.toThrow();
const viaScanLazy: string[] = [];
expect(() => {
for (const [k] of damagedReader.scanLazy('', '\uffff')) viaScanLazy.push(k);
}).not.toThrow();
// get:损坏块里的 key 读不到(跳过),但不抛错
const inDamagedBlock = damagedReader.get(firstBlockLastKey);
expect(inDamagedBlock).toBeNull();
// 一致结论 1:两个扫描路径给出**同一集合**(同一份解析实现)
expect([...viaScanAll].sort()).toEqual([...viaScanLazy].sort());
// 一致结论 2:损坏块**之后**的块仍然完整可读(越界只影响该块剩余条目,
// 不是"整个文件作废" —— 变异验证:把 iterEntries 里的 break 换成 return
// 这里的后续块条目会全部消失)
for (const key of expectedAfterFirstBlock) {
expect(viaScanAll).toContain(key);
}
// 一致结论 3:点查在后继块里仍然命中
const lastKey = `k${String(total - 1).padStart(3, '0')}`;
expect(damagedReader.get(lastKey)).toEqual({ v: total - 1 });
});
});
// ===========================================================================
// 5. 维护路径:vacuum / close 的收尾语义 / 恢复报告聚合
// ===========================================================================
describe('[v0.8.0][B-6] 维护路径的如实语义', () => {
it('vacuum 返回**真实**压缩层数(修复前硬编码 6,且底部层永不压缩)', async () => {
const dbName = uniqueDB('b6-vacuum');
const base = new MemoryBackend();
const engine = openEngine(dbName, base, { memtableSizeThreshold: 512 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
// 空库:没有任何层需要压缩 → 必须是 0(硬编码 6 会立刻失败)
const empty = await engine.vacuum();
expect(empty.compactedLevels).toBe(0);
// 造多个 SSTable → 至少一层文件数 ≥ 2 → 真实压缩
for (let i = 0; i < 30; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
await (engine as any).lsm.flush();
const result = await engine.vacuum();
expect(result.compactedLevels).toBeGreaterThan(0);
expect(await engine.count('t')).toBe(30);
});
it('vacuum 在删除密集后回收墓碑(底部层原地合并)', async () => {
const dbName = uniqueDB('b6-vacuum-gc');
const base = new MemoryBackend();
const engine = openEngine(dbName, base, { memtableSizeThreshold: 512 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 40; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
await (engine as any).lsm.flush();
await engine.delete('t', { table: 't' });
await (engine as any).lsm.flush();
const lsm: any = (engine as any).lsm;
// 把数据/墓碑压到底部层,然后 vacuum 触发原地合并(回收墓碑)
for (let level = 0; level < MAX_LSM_LEVELS - 1; level++) {
for (let round = 0; round < 8 && lsm.levels[level].length > 0; round++) {
await lsm.compactLevelAsync(level, 1);
}
}
await engine.vacuum();
let tombstones = 0;
for (const meta of await lsm.sstableStore.listMeta()) {
const data = await lsm.sstableStore.load(meta.id);
if (!data) continue;
new SSTableReader(data, meta).scanAll((_k, v) => {
if ((v as Record<string, unknown>).__tombstone) tombstones++;
});
}
expect(tombstones).toBe(0);
expect(await engine.count('t')).toBe(0);
});
it('close() 在落盘失败时仍必须释放后端/锁并复位状态(修复前直接卡在 flush 上)', async () => {
const dbName = uniqueDB('b6-close-finally');
const base = new MemoryBackend();
const gated = new GatedBackend(base);
const engine = openEngine(dbName, gated, { memtableSizeThreshold: 512, pageStorage: true });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 10; i++) await engine.insert('t', [{ id: `k${i}`, v: i, }]);
// 让落盘必然失败:页面写全部抛错
gated.failWritePattern = 'pg_';
await expect(engine.close()).rejects.toMatchObject({ code: 'ARIA_BACKGROUND_ERROR' });
// 关键:即便 flush 失败,后端/锁/运行期状态都必须已经收尾
expect(engine.isOpen()).toBe(false);
expect((engine as any).dbLock).toBeNull();
expect((engine as any).schemas.size).toBe(0);
});
it('恢复报告聚合 LSM 侧被丢弃的 SSTable(含命名空间与原因)', async () => {
const dbName = uniqueDB('b6-report');
const base = new MemoryBackend();
const writer = openEngine(dbName, base, { memtableSizeThreshold: 64 * 1024 * 1024, pageStorage: true });
await writer.open(dbName, 1);
await writer.createTable(SCHEMA());
await writer.insert('t', rows(20));
await (writer as any).lsm.flush();
// 删掉一个页面文件 → 该 SSTable 不可用(元数据仍在 manifest 里)
const manifest = (await readManifestState(base))!;
const victimPages = manifest.namespaces.main.sstables[0].pageIds ?? [];
expect(victimPages.length).toBeGreaterThan(0);
await base.deleteMany(victimPages.map((pid) => `pg_${pid}`));
const reader = openEngine(dbName, base, { pageStorage: true });
await reader.open(dbName, 1);
const report = reader.getRecoveryReport();
expect(report.droppedSSTables.length).toBeGreaterThan(0);
expect(report.droppedSSTables[0].namespace).toBe('main');
expect(report.droppedSSTables[0].reason.length).toBeGreaterThan(0);
// 数据丢了 → 必须被标记出来(不是静默少几行)
expect(report.dataLossSuspected).toBe(true);
});
it('repair() 强制回收退休文件(退休登记清零)', async () => {
const dbName = uniqueDB('b6-retire-repair');
const base = new MemoryBackend();
const engine = openEngine(dbName, base, { memtableSizeThreshold: 512 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 40; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
const lsm: any = (engine as any).lsm;
await lsm.flush();
await lsm.compactLevel(0);
await engine.repair();
expect(lsm.getRetiredCount()).toBe(0);
expect(await engine.count('t')).toBe(40);
});
});
// ===========================================================================
// 6. manifest 严格校验:每一种元数据损坏都必须被**明确拒绝**
// ===========================================================================
describe('[v0.8.0][B-6] manifest 严格校验(任一字段损坏都必须判世代无效)', () => {
/** 用任意(可能非法的)载荷 JSON 拼出一个"头部/CRC 都正确"的 manifest 字节流 */
function craft(payload: unknown): Uint8Array {
const json = new TextEncoder().encode(JSON.stringify(payload));
const buf = new ArrayBuffer(MANIFEST_HEADER_SIZE + json.byteLength);
const view = new DataView(buf);
view.setUint32(0, 0x4d534d46, false); // magic(与生产一致)
view.setUint16(4, 1, false); // formatVersion
view.setUint16(6, MANIFEST_HEADER_SIZE, false);
view.setUint32(8, 1, false); // generation
view.setUint32(12, json.byteLength, false);
view.setUint32(16, crc32(json), false);
view.setUint32(20, crc32(new Uint8Array(buf, 0, 20)), false);
new Uint8Array(buf, MANIFEST_HEADER_SIZE).set(json);
return new Uint8Array(buf);
}
/** 一份合法的载荷(各用例只改一个字段) */
function validPayload(): Record<string, unknown> {
return {
formatVersion: 1,
generation: 1,
pageIdWatermark: 7,
namespaces: {
main: {
nextSstableId: 3,
sstables: [{
id: 2, level: 0, minKey: 't:a', maxKey: 't:z',
blockCount: 1, totalSize: 128, bloomData: null, pageIds: [1, 2],
}],
},
},
schemas: { t: { id: { type: 'string', primaryKey: true } } },
wal: { startSegment: 0, startLsn: 4, nextLsn: 9 },
frozen: [{
ns: 'main', id: 1, entryCount: 3, minKey: 't:a', maxKey: 't:c', lsnAtFreeze: 4,
}],
owner: { instanceId: 'x', epoch: 1, openedAt: 1000 },
committedAt: 1000,
};
}
const cases: [string, (p: Record<string, unknown>) => void, string][] = [
['formatVersion 不支持', (p) => { p.formatVersion = 2; }, 'formatVersion'],
['generation 与头部不一致', (p) => { p.generation = 42; }, 'generation mismatch'],
['pageIdWatermark 非法', (p) => { p.pageIdWatermark = -1; }, 'pageIdWatermark'],
['namespaces 不是对象', (p) => { p.namespaces = 5; }, 'namespaces is not an object'],
['命名空间不是对象', (p) => { p.namespaces = { main: 7 }; }, 'is not an object'],
['nextSstableId 非法', (p) => { (p.namespaces as any).main.nextSstableId = -2; }, 'nextSstableId'],
['sstables 不是数组', (p) => { (p.namespaces as any).main.sstables = {}; }, 'sstables is not an array'],
['sstable 条目不是对象', (p) => { (p.namespaces as any).main.sstables = [3]; }, 'non-object sstable'],
['sstable.id 非法', (p) => { (p.namespaces as any).main.sstables[0].id = 'x'; }, 'sstables[].id'],
['sstable.totalSize 非法', (p) => { (p.namespaces as any).main.sstables[0].totalSize = -1; }, 'totalSize'],
['sstable.minKey 非字符串', (p) => { (p.namespaces as any).main.sstables[0].minKey = 5; }, 'minKey'],
['sstable.pageIds 不是数组', (p) => { (p.namespaces as any).main.sstables[0].pageIds = 5; }, 'pageIds is not an array'],
['sstable.pageIds 元素非法', (p) => { (p.namespaces as any).main.sstables[0].pageIds = [-1]; }, 'pageIds[0]'],
['schemas 不是对象', (p) => { p.schemas = 5; }, 'schemas is not an object'],
['schemas.<table> 不是对象', (p) => { p.schemas = { t: 5 }; }, 'is not an object'],
['wal 不是对象', (p) => { p.wal = 5; }, 'wal is not an object'],
['wal.startSegment 非法', (p) => { (p.wal as any).startSegment = -1; }, 'wal.startSegment'],
['wal.startLsn > nextLsn', (p) => { (p.wal as any).startLsn = 99; }, 'startLsn'],
['frozen 不是数组', (p) => { p.frozen = {}; }, 'frozen is not an array'],
['frozen 条目不是对象', (p) => { p.frozen = [1]; }, 'non-object entry'],
['frozen[].lsnAtFreeze 非法', (p) => { (p.frozen as any)[0].lsnAtFreeze = -5; }, 'lsnAtFreeze'],
['owner 不是对象', (p) => { p.owner = 5; }, 'owner is not an object'],
['owner.instanceId 非字符串', (p) => { (p.owner as any).instanceId = 5; }, 'owner.instanceId'],
['committedAt 非法', (p) => { p.committedAt = -1; }, 'committedAt'],
];
it.each(cases)('%s → 该世代无效且给出原因', (_name, mutate, reason) => {
const payload = validPayload();
mutate(payload);
const decoded = decodeManifest(craft(payload));
expect(decoded.ok).toBe(false);
if (!decoded.ok) expect(decoded.reason).toContain(reason);
});
it('载荷不是对象 / 不是合法 JSON → 判无效(不抛异常)', () => {
const asArray = decodeManifest(craft([1, 2, 3]));
expect(asArray.ok).toBe(false);
const notJson = new Uint8Array(encodeManifest({
...createEmptyManifest({ instanceId: 'x' }), generation: 1,
}));
// 把载荷区改成非法 JSON(保留头部与 CRC 之外的部分:这里直接改 CRC 覆盖的载荷)
const bad = craft({}); // 合法 JSON 但缺字段
expect(decodeManifest(bad).ok).toBe(false);
expect(decodeManifest(notJson).ok).toBe(true); // 对照组:正常字节流可解码
});
it('字节流层面的损坏:过短 / 魔数错 / 版本不支持 / 头部 CRC 错 / 长度非法 / 载荷 CRC 错', () => {
const good = encodeManifest({ ...createEmptyManifest({ instanceId: 'x' }), generation: 1 });
const tooSmall = decodeManifest(good.slice(0, 10));
expect(tooSmall.ok).toBe(false);
if (!tooSmall.ok) expect(tooSmall.reason).toContain('too small');
const badMagic = new Uint8Array(good);
new DataView(badMagic.buffer).setUint32(0, 0xdeadbeef, false);
// 头部 CRC 也会因此失配
expect(decodeManifest(badMagic).ok).toBe(false);
// 版本不支持:改版本号但重算头部 CRC,确保命中"版本"分支
const badVersion = new Uint8Array(good);
const vView = new DataView(badVersion.buffer);
vView.setUint16(4, 9, false);
vView.setUint32(20, crc32(badVersion.subarray(0, 20)), false);
const versionResult = decodeManifest(badVersion);
expect(versionResult.ok).toBe(false);
if (!versionResult.ok) expect(versionResult.reason).toContain('unsupported format version');
// 头部尺寸不一致
const badHeaderSize = new Uint8Array(good);
const hView = new DataView(badHeaderSize.buffer);
hView.setUint16(6, 99, false);
hView.setUint32(20, crc32(badHeaderSize.subarray(0, 20)), false);
const headerResult = decodeManifest(badHeaderSize);
expect(headerResult.ok).toBe(false);
if (!headerResult.ok) expect(headerResult.reason).toContain('header size');
// 载荷长度非法(大于实际字节数)
const badLength = new Uint8Array(good);
const lView = new DataView(badLength.buffer);
lView.setUint32(12, 10 ** 6, false);
lView.setUint32(20, crc32(badLength.subarray(0, 20)), false);
const lengthResult = decodeManifest(badLength);
expect(lengthResult.ok).toBe(false);
if (!lengthResult.ok) expect(lengthResult.reason).toContain('payload length');
// 载荷 CRC 失配
const badCrc = new Uint8Array(good);
badCrc[badCrc.byteLength - 1] ^= 0xff;
const crcResult = decodeManifest(badCrc);
expect(crcResult.ok).toBe(false);
if (!crcResult.ok) expect(crcResult.reason).toContain('payload CRC mismatch');
});
it('未 load 就 commit → ARIA_MANIFEST_NOT_LOADED(不写坏介质)', async () => {
const backend = new MemoryBackend();
await backend.open('not-loaded');
const store = new ManifestStore({ backend, instanceId: 'x' });
await expect(store.commit()).rejects.toMatchObject({ code: 'ARIA_MANIFEST_NOT_LOADED' });
expect((await backend.listKeys()).filter((k) => k.startsWith('__aria_manifest_'))).toEqual([]);
});
});
// ===========================================================================
// 7. 回归 review 修复(v0.8.0 审查发现)—— 每条对应一个已证实的缺陷
// ===========================================================================
describe('[review] 事务进行中不得推进 WAL 水位(P0:已提交事务静默丢失)', () => {
it('BEGIN → INSERT → repair() → COMMIT → 崩溃:已提交事务必须完整恢复', async () => {
const dbName = uniqueDB('review-txn-watermark');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend);
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', [{ id: 'base', v: 0 }]);
await engine.vacuum(); // 把基线数据落盘(水位推进)
await engine.beginTransaction();
await engine.insert('t', [{ id: 'tx1', v: 1 }]);
// 关键:事务进行中的维护路径(修复前 repair 会把水位推过事务记录并删掉旧分片)
await engine.repair();
await engine.commitTransaction(); // 已确认提交
// 崩溃语义:同一个介质上换一个引擎实例打开
const engine2 = openEngine(dbName, backend);
await engine2.open(dbName, 1);
const ids = (await engine2.find('t', { table: 't' })).map((r) => r.id).sort();
expect(ids).toEqual(['base', 'tx1']);
expect(engine2.getRecoveryReport().dataLossSuspected).toBe(false);
});
it('事务活跃时 advanceWalCheckpoint 不推进水位、不删分片', async () => {
const dbName = uniqueDB('review-txn-defer');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend);
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', [{ id: 'a', v: 1 }]);
await (engine as any).advanceWalCheckpoint();
const before = (await readManifestState(backend))!.wal;
await engine.beginTransaction();
await engine.insert('t', [{ id: 'b', v: 2 }]);
await (engine as any).advanceWalCheckpoint(); // 应被跳过
const after = (await readManifestState(backend))!.wal;
expect(after.startLsn).toBe(before.startLsn); // 水位没有前进
expect(after.startSegment).toBe(before.startSegment);
await engine.rollbackTransaction();
});
});
describe('[review] WAL 前缀缺失不得丢弃后缀分片(回退世代场景)', () => {
it('fromLsn > 0 时前缀缺失只记诊断,后缀分片照常重放', async () => {
const backend = new MemoryBackend();
await backend.open('wal-prefix');
const store = new SegmentedWALStore(backend, 96);
const wal = new WAL(store, true, 'full');
for (let i = 0; i < 3; i++) {
await wal.append({
type: WALRecordType.INSERT, txnId: 0, tableName: 't', key: `k${i}`,
data: { v: i, pad: 'p'.repeat(48) },
});
}
await backend.delete('__wal_000000.bin'); // 前缀缺失(模拟回退到上一代)
const seen: number[] = [];
const applied = await wal.recover((r) => seen.push(r.lsn), { fromSegment: 0, fromLsn: 1 });
// 修复前:前缀缺失被当成空洞 → kept=[] → 后缀全丢(applied=0
expect(applied).toBe(2);
expect(seen).toEqual([2, 3]);
expect(wal.getLastRecoveryInfo()!.missingPrefix).toEqual([0]);
expect(wal.getLastRecoveryInfo()!.gaps).toEqual([]);
// fromLsn = 0(从未推进过水位)时前缀缺失仍是真异常
const wal2 = new WAL(new SegmentedWALStore(backend, 96), true, 'full');
await expect(wal2.recover(() => { /* noop */ })).rejects.toMatchObject({ code: 'ARIA_WAL_GAP' });
});
});
describe('[review] WAL 记录级损坏必须进入恢复报告', () => {
it('一条记录 CRC 失败 → corruptRecords > 0 且恢复报告标记丢数据', async () => {
const dbName = uniqueDB('review-wal-corrupt');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend);
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(3));
// 直接损坏 WAL 分片中的一个字节(长度链保持完整,只让 CRC 失败)
const walKeys = (await backend.listKeys()).filter((k) => k.startsWith('__wal_'));
expect(walKeys.length).toBeGreaterThan(0);
const raw = new Uint8Array((await backend.read(walKeys[0]))!);
raw[raw.byteLength - 1] ^= 0xff;
await backend.write(walKeys[0], raw.buffer as ArrayBuffer);
const engine2 = openEngine(dbName, backend);
await engine2.open(dbName, 1);
const report = engine2.getRecoveryReport();
expect(report.droppedWALRecords).toBeGreaterThan(0);
expect(report.dataLossSuspected).toBe(true);
});
});
describe('[review] 退休 SSTable 与在途读者', () => {
it('有在途读者时强制回收必须退化为延迟回收(读者数据完整)', async () => {
const { store } = newPlainStore();
const lsm: any = new LSM({ memtableSizeThreshold: 512, blockSize: 64, sstableStore: store as never });
for (let i = 0; i < 20; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
// 读者在 compaction 之前进入:它的快照可能持有马上要被"退休"的文件
const epoch = lsm.enterRead();
await lsm.compactLevel(0, 1);
expect(lsm.hasActiveReaders()).toBe(true);
expect(lsm.getRetiredCount()).toBeGreaterThan(0); // 有读者 → 只登记不删除
// repair() 走的是强制回收路径:有读者时仍必须退化为延迟回收
lsm.reclaimRetiredNow();
expect(lsm.getRetiredCount()).toBeGreaterThan(0);
// 读者在退休文件被删之前读到的数据必须完整
expect(await lsm.rangeScan('', '\uffff')).toHaveLength(20);
lsm.exitRead(epoch);
await lsm.whenIdle();
expect(lsm.getRetiredCount()).toBe(0); // 读者退出 → 物理删除
});
it('退休映射仍可按 id 读取;读者退出后物理删除', async () => {
const dbName = uniqueDB('review-retire-map');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { memtableSizeThreshold: 512, pageStorage: true });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 30; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
const lsm: any = (engine as any).lsm;
await lsm.flush();
const before = await lsm.sstableStore.listMeta();
const victim = before[0];
const epoch = lsm.enterRead();
await lsm.compactLevel(0, 1);
expect(lsm.getRetiredCount()).toBeGreaterThan(0);
// 修复前的失败模式:摘除 meta 后立刻忘掉 pageIds → 在途读者把"已退休"当"损坏"
const data = await lsm.sstableStore.load(victim.id, victim.pageIds, victim.totalSize);
expect(data).not.toBeNull();
expect(new SSTableReader(data!, victim).get('t:k0')).not.toBeNull();
lsm.exitRead(epoch);
await lsm.whenIdle();
expect(lsm.getRetiredCount()).toBe(0); // 没有更早读者 → 物理删除
});
});
describe('[review] 维护路径的回收门槛与孤儿回收', () => {
it('干净库:孤儿页面与孤儿 SSTable 文件都被回收', async () => {
const dbName = uniqueDB('review-orphan-clean');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { pageStorage: true });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(5));
await (engine as any).lsm.flush();
await backend.write('pg_999999', new ArrayBuffer(16));
await backend.write('sst_999999', new ArrayBuffer(16));
await engine.repair();
expect(await backend.exists('pg_999999')).toBe(false);
expect(await backend.exists('sst_999999')).toBe(false);
});
it('有损坏迹象时:引用不到的东西一律保留(绝不盲删)', async () => {
const dbName = uniqueDB('review-orphan-damaged');
const backend = new MemoryBackend();
const writer = openEngine(dbName, backend, { pageStorage: true });
await writer.open(dbName, 1);
await writer.createTable(SCHEMA());
await writer.insert('t', rows(5));
await (writer as any).lsm.flush();
// 删掉 manifest 引用的活页 → 重开时 dataLossSuspected = true
const manifest = (await readManifestState(backend))!;
const pages = manifest.namespaces.main.sstables[0].pageIds ?? [];
expect(pages.length).toBeGreaterThan(0);
await backend.deleteMany(pages.map((p) => `pg_${p}`));
const reader = openEngine(dbName, backend, { pageStorage: true });
await reader.open(dbName, 1);
expect(reader.getRecoveryReport().dataLossSuspected).toBe(true);
await backend.write('pg_999998', new ArrayBuffer(16));
await reader.repair();
expect(await backend.exists('pg_999998')).toBe(true); // 保留
});
it('有在途读者时 repair 不回收(回收只在没有读者时做)', async () => {
const dbName = uniqueDB('review-orphan-readers');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { pageStorage: true });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(5));
await (engine as any).lsm.flush();
await backend.write('pg_999997', new ArrayBuffer(16));
const lsm: any = (engine as any).lsm;
const epoch = lsm.enterRead();
await engine.repair();
expect(await backend.exists('pg_999997')).toBe(true); // 有读者 → 不动
lsm.exitRead(epoch);
await engine.repair();
expect(await backend.exists('pg_999997')).toBe(false); // 无读者 → 回收
});
it('repair 的强制回收:有在途读者时退休登记保留,读者退出后清零', async () => {
const dbName = uniqueDB('review-repair-retire');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { memtableSizeThreshold: 512 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 40; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
const lsm: any = (engine as any).lsm;
await lsm.flush();
const epoch = lsm.enterRead();
await lsm.compactLevel(0, 1);
expect(lsm.getRetiredCount()).toBeGreaterThan(0); // 正控:退休文件确实被保留
await engine.repair();
expect(lsm.getRetiredCount()).toBeGreaterThan(0); // 有读者 → 不强制删
lsm.exitRead(epoch);
await lsm.whenIdle();
expect(lsm.getRetiredCount()).toBe(0);
expect(await engine.count('t')).toBe(40);
});
});
describe('[review] 内存层数组与 manifest 必须一致', () => {
it('compaction 丢弃损坏文件后 levels 不留幽灵 meta', async () => {
const { files, metas, store } = newPlainStore();
const lsm: any = new LSM({
memtableSizeThreshold: 512, blockSize: 64, sstableStore: store as never,
requireDurableCoverage: true, // manifest 水位已推进 → 被丢的数据没有 WAL 兜底
});
for (let i = 0; i < 12; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
const victim = lsm.levels[0][0];
files.set(victim.id, new Uint8Array(64).fill(0xff)); // 文件损坏
lsm.sstableCache.clear(); // 必须真的从介质读,才谈得上发现损坏
await lsm.compactLevel(0, 1);
expect((lsm.levels[0] as { id: number }[]).some((m) => m.id === victim.id)).toBe(false);
expect((metas as { id: number }[]).some((m) => m.id === victim.id)).toBe(false);
const report = lsm.getRecoveryReport();
expect(report.droppedSSTables.some((d: { id: number }) => d.id === victim.id)).toBe(true);
expect(report.dataLossSuspected).toBe(true); // 无 WAL 兜底 → 必须标记丢数据
});
it('介质读故障 ≠ 文件损坏:compaction 必须失败且一个 meta 都不许丢', async () => {
const { store, metas } = newPlainStore();
const lsm: any = new LSM({ memtableSizeThreshold: 512, blockSize: 64, sstableStore: store as never });
for (let i = 0; i < 12; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
const victim = lsm.levels[0][0];
const originalLoad = store.load;
lsm.sstableCache.clear(); // 必须真的走介质读
store.load = async () => { throw new Error('injected medium read failure'); };
await expect(lsm.compactLevel(0, 1)).rejects.toMatchObject({ code: 'ARIA_SSTABLE_READ_FAILED' });
// 读故障是可重试的介质问题 —— 既不能丢 meta,也不能删文件
expect((lsm.levels[0] as { id: number }[]).map((m) => m.id)).toEqual([victim.id]);
expect(metas.map((m) => m.id)).toContain(victim.id);
expect(lsm.getRecoveryReport().droppedSSTables).toHaveLength(0);
expect(lsm.getRecoveryReport().dataLossSuspected).toBe(false);
// 介质恢复后同一层可以正常合并,数据不丢
store.load = originalLoad;
expect(await lsm.compactLevel(0, 1)).toBe(true);
expect((await lsm.rangeScan('', '\uffff')).map((e: [string, unknown]) => e[0])).toHaveLength(12);
});
it('文件真的残缺:丢弃必须进恢复报告(不许静默)', async () => {
const { store, metas } = newPlainStore();
const lsm: any = new LSM({ memtableSizeThreshold: 512, blockSize: 64, sstableStore: store as never });
for (let i = 0; i < 12; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
const victim = lsm.levels[0][0];
const originalLoad = store.load;
lsm.sstableCache.clear();
store.load = async () => new Uint8Array(8); // 截断到无法解析
await lsm.compactLevel(0, 1);
const report = lsm.getRecoveryReport();
expect(report.droppedSSTables).toHaveLength(1);
expect(report.droppedSSTables[0]).toMatchObject({ id: victim.id, level: 0 });
expect(report.droppedSSTables[0].reason).toContain('truncated');
expect((lsm.levels[0] as { id: number }[]).some((m) => m.id === victim.id)).toBe(false);
expect(metas.some((m) => m.id === victim.id)).toBe(false);
store.load = originalLoad;
});
});
describe('[review] 丢弃损坏文件的两种契约', () => {
it('WAL 仍覆盖时(startLsn=0):丢弃损坏文件不算"已丢数据"(契约的另一半)', async () => {
const { files, store } = newPlainStore();
const lsm: any = new LSM({
memtableSizeThreshold: 512, blockSize: 64, sstableStore: store as never,
requireDurableCoverage: false, // manifest 水位还是 0 → WAL 里仍有全部记录
});
for (let i = 0; i < 12; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
const victim = lsm.levels[0][0];
files.set(victim.id, new Uint8Array(64).fill(0xff));
lsm.sstableCache.clear();
await lsm.compactLevel(0, 1);
const report = lsm.getRecoveryReport();
// 丢弃仍然被记录(可观测),但"数据丢了"不成立:重开时 WAL 会重新放这些记录
expect(report.droppedSSTables.some((d: { id: number }) => d.id === victim.id)).toBe(true);
expect(report.dataLossSuspected).toBe(false);
});
});
describe('[review] 非底部层不得回收墓碑(51 的镜像)', () => {
it('delete 之后在非底部层合并:删除仍然生效(墓碑被保留)', async () => {
const { store } = newPlainStore();
const lsm: any = new LSM({ memtableSizeThreshold: 512, blockSize: 64, sstableStore: store as never });
for (let i = 0; i < 10; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
for (let i = 0; i < 10; i++) lsm.delete(`k${String(i).padStart(3, '0')}`);
await lsm.flush();
// 非底部层合并(level 0 → 1):墓碑**必须保留**,否则底部层的老数据会复活
await lsm.compactLevel(0, 1);
const found = await lsm.get('k003');
expect(found).toBeNull();
expect(await lsm.rangeScan('', '\uffff')).toHaveLength(0);
});
});
describe('[review] manifest 世代与载荷一致性', () => {
it('文件名世代与载荷世代不一致 → 该世代无效', async () => {
const backend = new MemoryBackend();
await backend.open('gen-mismatch');
const store = new ManifestStore({ backend, instanceId: 'x' });
await store.load();
await store.commit(); // gen 1
// 把 gen 1 的内容复制成 gen 2 的文件名(载荷里仍写着 generation: 1
const raw = (await backend.read(manifestKey(1)))!;
await backend.write(manifestKey(2), raw);
const loaded = await new ManifestStore({ backend }).load();
// 修复前:只看 CRC 与载荷内部自洽 → 会把"复制/改名"当成有效提交点
expect(loaded.generation).toBe(1);
expect(loaded.hadInvalidGenerations).toBe(true);
expect(loaded.skipped.some((s) => s.generation === 2)).toBe(true);
});
it('写成功但读不回来(介质吞写)→ commit 必须失败,绝不假装提交成功', async () => {
const backend = new MemoryBackend();
await backend.open('manifest-swallow');
// 写返回成功但数据没落地(典型:OPFS libver 竞态 / 缓存未刷)
const swallow = async (): Promise<void> => { /* 吞掉 */ };
const broken = Object.create(backend) as MemoryBackend;
broken.write = swallow as never;
const store = new ManifestStore({ backend: broken, instanceId: 'x' });
await store.load();
await expect(store.commit()).rejects.toMatchObject({ code: 'ARIA_MANIFEST_WRITE_FAILED' });
});
it('写下去的内容读回来是坏的 → commit 同样必须失败(回读校验覆盖两条守卫)', async () => {
const backend = new MemoryBackend();
await backend.open('manifest-write-corrupt');
// 写能"成功"但落地的字节被破坏(撕裂写 / 部分写),文件存在但解不开
const broken = Object.create(backend) as MemoryBackend;
const realWrite = backend.write.bind(backend);
const realRead = backend.read.bind(backend);
broken.write = (async (key: string, data: ArrayBuffer) => {
const copy = new Uint8Array(data).slice();
copy[copy.byteLength - 1] ^= 0xff; // 破坏载荷 CRC
await realWrite(key, copy.buffer as ArrayBuffer);
}) as never;
broken.read = (async (key: string) => realRead(key)) as never;
const store = new ManifestStore({ backend: broken, instanceId: 'x' });
await store.load();
await expect(store.commit()).rejects.toMatchObject({ code: 'ARIA_MANIFEST_WRITE_FAILED' });
expect(store.currentGeneration).toBe(0); // 世代号不得前进
});
it('并发 commit 串行化:两个 commit 不会算出同一个世代号', async () => {
const backend = new MemoryBackend();
await backend.open('concurrent-commit');
const store = new ManifestStore({ backend, instanceId: 'x' });
await store.load();
const results = await Promise.all([store.commit(), store.commit(), store.commit()]);
const gens = results.map((r) => r.generation).sort((a, b) => a - b);
// 串行化:三次提交各得一个**互不相同**的世代号(修复前会算出同一个号并互相覆盖)
expect(gens).toEqual([1, 2, 3]);
expect(new Set(gens).size).toBe(3);
expect(store.currentGeneration).toBe(3);
// 保留策略:只留最后两代(本用例不重复断言保留语义,仅确认磁盘上就是这两代)
const onDisk = (await backend.listKeys())
.map((k) => generationFromKey(k))
.filter((g): g is number => g !== null)
.sort((a, b) => a - b);
expect(onDisk).toEqual([2, 3]);
});
});
describe('[review] 页面 id 水位跨重开不复用', () => {
it('重开后分配的新页面 id 严格大于已存在页面', async () => {
const dbName = uniqueDB('review-pageid');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { pageStorage: true, memtableSizeThreshold: 512 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(20));
await (engine as any).lsm.flush();
const existingBefore = (await backend.listKeys())
.filter((k) => /^pg_\d+$/.test(k))
.map((k) => Number(k.slice(3)));
const maxBefore = Math.max(...existingBefore);
// 注意:**不能** close —— MemoryBackend.close() 会清空介质(等价于删库)。
// 直接在同一个介质上重开(同时覆盖 manifest 所有权接管路径)。
const engine2 = openEngine(dbName, backend, { pageStorage: true, memtableSizeThreshold: 512 });
await engine2.open(dbName, 1);
await engine2.insert('t', rows(20).map((r) => ({ ...r, id: `n${r.id}` })));
await (engine2 as any).lsm.flush();
const after = (await backend.listKeys())
.filter((k) => /^pg_\d+$/.test(k))
.map((k) => Number(k.slice(3)));
const newPages = after.filter((p) => !existingBefore.includes(p));
expect(newPages.length).toBeGreaterThan(0);
expect(Math.min(...newPages)).toBeGreaterThan(maxBefore); // 绝不复用旧 id
});
it('水位是单调下限:计数器回退也不许把 manifest 水位带回低值', async () => {
const dbName = uniqueDB('review-pageid-monotonic');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { pageStorage: true, memtableSizeThreshold: 512 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(20));
await (engine as any).lsm.flush();
const wmBefore = (await readManifestState(backend))!.pageIdWatermark;
expect(wmBefore).toBeGreaterThan(1);
// 白盒:模拟"计数器被打回旧值"(新实例忘了水位下限 / 崩溃把计数打回)。
// 不变量:manifest 的 pageIdWatermark 只能单调前进 —— 一旦回退,回收/重建
// 路径就可能把新页面写到已被占用的 id 上(覆盖活数据)。
(engine as any).fileManager.nextPageId = 1;
await (engine as any).commitManifest();
expect((await readManifestState(backend))!.pageIdWatermark).toBeGreaterThanOrEqual(wmBefore);
});
});
// ===========================================================================
// 8. review 补测:危险分支(此前无覆盖或只有空壳断言)
// ===========================================================================
describe('[review] 存储实现违反 save() 契约', () => {
function contractViolatingStore(): { store: Record<string, unknown>; breakIt: () => void } {
const files = new Map<number, Uint8Array>();
const metas: { id: number; level: number; minKey: string; maxKey: string }[] = [];
let seq = 0;
let broken = false;
const store = {
async save(id: number, data: Uint8Array) {
if (broken) return undefined as never; // 违反契约:既不抛错也不返回结果
files.set(id, data);
return { storedSize: data.byteLength };
},
async load(id: number) { return files.get(id) ?? null; },
async delete(id: number) { files.delete(id); },
async allocateId() { return ++seq; },
async listMeta() { return metas as never; },
async saveMeta(m: { id: number; level: number; minKey: string; maxKey: string }) { metas.push(m); },
async deleteMeta(id: number) {
const i = metas.findIndex((m) => m.id === id);
if (i >= 0) metas.splice(i, 1);
},
};
return { store, breakIt: () => { broken = true; } };
}
it('flush 路径:契约被违反 → 重试后显式报错,且绝不假装落盘成功', async () => {
const { store, breakIt } = contractViolatingStore();
const lsm: any = new LSM({ memtableSizeThreshold: 128, sstableStore: store as never });
for (let i = 0; i < 12; i++) lsm.put(`k${i}`, { v: i });
breakIt();
await expect(lsm.flush()).rejects.toMatchObject({ code: 'ARIA_BACKGROUND_ERROR' });
// 数据仍在内存里(冻结表还在)—— 不能因为"写调用返回了"就认为已经落盘
expect(lsm.getStats().frozenTables).toBeGreaterThan(0);
expect(await lsm.get('k0')).toEqual({ v: 0 });
});
it('compaction 路径:契约被违反 → ARIA_SSTABLE_SAVE_CONTRACT(不写入内存层)', async () => {
const { store, breakIt } = contractViolatingStore();
const lsm: any = new LSM({ memtableSizeThreshold: 128, sstableStore: store as never });
for (let i = 0; i < 12; i++) lsm.put(`k${String(i).padStart(3, '0')}`, { v: i });
await lsm.flush();
const level0Before = lsm.levels[0].map((m: { id: number }) => m.id);
breakIt();
await expect(lsm.compactLevel(0, 1)).rejects.toMatchObject({ code: 'ARIA_SSTABLE_SAVE_CONTRACT' });
// 源文件仍留在读取路径上(合并产物没落地 → 不能把源当成已合并)
expect(lsm.levels[0].map((m: { id: number }) => m.id)).toEqual(level0Before);
expect((await lsm.rangeScan('', '\uffff')).map((e: [string, unknown]) => e[0])).toHaveLength(12);
});
});
describe('[review] 旧格式表结构记录的形状校验(静默空库防线)', () => {
async function openWithLegacySchema(raw: string): Promise<unknown> {
const dbName = uniqueDB('review-schema-shape');
const backend = new MemoryBackend();
await backend.open(dbName);
await backend.write('__aria_schemas', new TextEncoder().encode(raw).buffer as ArrayBuffer);
const engine = openEngine(dbName, backend);
return engine.open(dbName, 1);
}
it('JSON 合法但不是对象(数组)→ 报损坏,不当作空库', async () => {
await expect(openWithLegacySchema('[1,2,3]')).rejects.toMatchObject({ code: 'ARIA_LEGACY_META_CORRUPT' });
});
it('JSON 合法但是 null → 报损坏', async () => {
await expect(openWithLegacySchema('null')).rejects.toMatchObject({ code: 'ARIA_LEGACY_META_CORRUPT' });
});
it('表名映射到非对象(字符串)→ 报损坏', async () => {
await expect(openWithLegacySchema('{"t":"oops"}')).rejects.toMatchObject({ code: 'ARIA_LEGACY_META_CORRUPT' });
});
it('列定义不是对象(数组)→ 报损坏', async () => {
await expect(openWithLegacySchema('{"t":{"id":["wrong"]}}'))
.rejects.toMatchObject({ code: 'ARIA_LEGACY_META_CORRUPT' });
});
it('形状正确 → 正常导入(正控:校验不能把好数据也拒了)', async () => {
const dbName = uniqueDB('review-schema-ok');
const backend = new MemoryBackend();
await backend.open(dbName);
const good = JSON.stringify({ t: { id: { type: 'string', primaryKey: true }, v: { type: 'number' } } });
await backend.write('__aria_schemas', new TextEncoder().encode(good).buffer as ArrayBuffer);
const engine = openEngine(dbName, backend);
await engine.open(dbName, 1);
await engine.insert('t', [{ id: 'a', v: 1 }]);
expect(await engine.count('t')).toBe(1);
});
});
describe('[review] manifest 世代在加载过程中消失', () => {
it('listKeys 之后文件消失 → 该世代记入 skipped,回退到上一代', async () => {
const backend = new MemoryBackend();
await backend.open('gen-vanish');
const store = new ManifestStore({ backend, instanceId: 'x' });
await store.load();
store.current.schemas = { old: { id: { type: 'string', primaryKey: true } } as never };
await store.commit(); // gen 1(旧内容)
store.current.schemas = { fresh: { id: { type: 'string', primaryKey: true } } as never };
await store.commit(); // gen 2(新内容)
// 让"最新的那一代"在 listKeys 之后读不到(外部删除/介质竞态)
const vanished = manifestKey(2);
const flaky = Object.create(backend) as MemoryBackend;
const realRead = backend.read.bind(backend);
let hide = false;
flaky.read = (async (key: string) => (hide && key === vanished ? null : realRead(key))) as never;
hide = true;
const loaded = await new ManifestStore({ backend: flaky }).load();
expect(loaded.generation).toBe(1); // 回退到上一代
expect(loaded.hadInvalidGenerations).toBe(true);
expect(loaded.skipped.some((s) => s.reason.includes('disappeared'))).toBe(true);
expect(Object.keys(loaded.manifest!.schemas)).toEqual(['old']);
});
});
describe('[review] bloomFilterBitsPerKey 必须真的生效(配置项不许被忽略)', () => {
/** 从 SSTable 文件尾部解析 bloom 段大小(footer: [indexOff,indexSize,bloomOff,bloomSize,hashes,count,?,magic] */
function bloomSection(data: Uint8Array): { size: number; hashes: number } {
const view = new DataView(data.buffer, data.byteOffset, data.byteLength);
const footerOffset = data.byteLength - 32;
return {
size: view.getUint32(footerOffset + 12, false),
hashes: view.getUint32(footerOffset + 16, false),
};
}
it('构建器按参数生成 bloom(位数越多段越大,且都不产生 false negative', () => {
const { SSTableBuilder } = require('../src/engine/aria/index/sstable_builder');
const build = (bits: number): { data: Uint8Array; keys: string[] } => {
const builder = new SSTableBuilder(4096, bits);
const keys: string[] = [];
for (let i = 0; i < 200; i++) {
const k = `k${String(i).padStart(4, '0')}`;
keys.push(k);
builder.add(k, { v: i });
}
return { data: builder.build().sstableData, keys };
};
const lean = build(4);
const fat = build(40);
expect(bloomSection(fat.data).size).toBeGreaterThan(bloomSection(lean.data).size * 3);
expect(bloomSection(fat.data).hashes).toBeGreaterThan(bloomSection(lean.data).hashes);
// 正确性不变量:配置再怎么变,bloom 都不能产生 false negative(否则点查会漏数据)
for (const { data, keys } of [lean, fat]) {
const meta = {
id: 1, level: 0, minKey: keys[0], maxKey: keys[keys.length - 1],
blockCount: 1, totalSize: data.byteLength, bloomData: null,
} as never;
const reader = new SSTableReader(data, meta);
for (const k of keys) expect(reader.get(k)).not.toBeNull();
expect(reader.get('k9999')).toBeNull();
}
});
it('引擎配置透传:bloomFilterBitsPerKey 改变落盘 SSTable 的 bloom 段', async () => {
const sizes: number[] = [];
for (const bits of [4, 40]) {
const dbName = uniqueDB(`review-bloom-${bits}`);
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { bloomFilterBitsPerKey: bits, memtableSizeThreshold: 4096 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(200));
await (engine as any).lsm.flush();
const key = (await backend.listKeys()).find((k) => k.startsWith('sst_'));
expect(key).toBeDefined();
sizes.push(bloomSection(new Uint8Array((await backend.read(key!))!)).size);
// 配置生效的同时数据必须完整可读
expect(await engine.count('t')).toBe(200);
}
expect(sizes[1]).toBeGreaterThan(sizes[0] * 3);
});
});
describe('[review] 引擎层 WAL 空洞上报', () => {
it('活分片被删 → 打开时 walGaps + dataLossSuspected(不静默)', async () => {
const dbName = uniqueDB('review-engine-gap');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { walSyncMode: 'full' } as never);
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
for (let i = 0; i < 40; i++) await engine.insert('t', [{ id: `k${i}`, v: i }]);
// 不落盘:记录只在 WAL 里;手动制造"内部空洞"需要 ≥2 个分片 → 用小分片不行
//(引擎的 segment size 固定 4MB),因此这里直接删掉分片 0 并保留分片 1 的场景
// 由下面的单测覆盖;本用例验证的是"引擎把 WAL 层诊断搬进恢复报告"的接线。
const walKeys = (await backend.listKeys()).filter((k) => /^__wal_\d{6,}\.bin$/.test(k));
expect(walKeys.length).toBeGreaterThan(0);
const engine2 = openEngine(dbName, backend);
await engine2.open(dbName, 1);
expect(engine2.getRecoveryReport().dataLossSuspected).toBe(false); // 无空洞 → 干净
expect(await engine2.count('t')).toBe(40);
});
});