Files
MetonaSqlark/tests/v080-b6-single-commit-point.test.ts
T
thzxx c5694b1d23 feat(B-6): 存储层单一提交点(__aria_manifest)+ LSM 结构根治
按 PLAN-v0.7.5.md §B-6 的**完整规格**实施(此前只落地了"降级选项"里的五处止血):
B-6 要求的是 `__aria_manifest` 单一提交点 + LSM 单项改造。完整记录见方案附录 H。

一、单一提交点
  - 新增 `src/engine/aria/store/manifest.ts`:`__aria_manifest_<generation>`
    (magic + formatVersion + generation + 头部 CRC + 载荷 CRC;先写后验;保留两代)。
    载荷 = 页面水位 + 各命名空间 SSTable 元数据 + 表结构 + WAL 起始位置 + 待落盘冻结表意图。
  - 顺序固定:**数据落盘 → manifest 提交 → 才允许截断 WAL / 删除旧文件 / 删除旧 SSTable**。
  - 恢复只认最后一份 CRC 通过的世代;全部世代无效 → `ARIA_MANIFEST_CORRUPT`
    (修复前:裸 JSON meta 解析失败 → `[]` → 静默空库,随后 repair 还会删光活页)。
  - 旧格式(__aria_lsm_meta/__aria_schemas/__aria_meta)首次打开自动迁移,旧键保留;
    迁移遇到损坏 → `ARIA_LEGACY_META_CORRUPT`。
  - 陈旧实例保护(STALE_INSTANCE):认领时一次跨过 MANIFEST_TAKEOVER_STRIDE 个世代,
    杜绝"旧实例在途提交落在同一世代号上"(实测第二个实例 open 直接失败)。

二、LSM
  - 44 冻结表成为一等状态:失败保留 + 可重试(修复前失败即永久失去落盘机会)。
  - 45 `flush()` 先入链再报告后台错误(修复前一次后台失败会让之后每次 flush 直接抛错、
       数据永远等不到落盘);被重试修复的失败进 `getBackgroundWarnings()`(可见但不误报失败)。
  - 47 `MergeIterator` 胜出来源的补充推迟到下一次 `next()`:提前终止不再多算一条。
  - 49 `compacting` 由单 boolean 改为按层集合(跨层触发不再被静默丢弃)。
  - 50 compaction 不再"先 splice 整层再合并"(窗口内该层对读者可见);
       被取代的 SSTable 进"退休表" + 读者 epoch,等更早读者退出才物理删除。
  - 51 底部层原地合并回收墓碑(删除密集场景空间不再无界增长);"整层只剩墓碑" 有专门分支
       (修复前会读 `merged[0][0]` 抛 TypeError,compaction 永久失败)。
  - 55 flush 与 compaction 拆成两条链,checkpoint 只落 memtable;删除引擎层全部
       `prefetch*`/`drainChain` 依赖,改为"快照 + 结构版本乐观重试"
       (版本号同时覆盖 levels 与前台 memtable/frozen 的变化)。
  - 读路径自洽:介质读故障抛 `ARIA_SSTABLE_READ_FAILED`,不再折叠成"文件不存在"误删元数据。

三、WAL
  - LSN 全库单调(manifest 记高水位);按水位删除旧分片(`planKeepFrom` → 提交 → 再删除)。
  - **分片号只增不减**:修复前全量截断后重置为 0,会与 manifest 记录的 startSegment 错位,
    实测造成两个方向的损坏(删掉的行复活 / 已确认写入丢失,见随机压力套件)。
  - 分片空洞(含前缀缺失)显式报 `ARIA_WAL_GAP`,不再静默丢弃尾部。

四、其它
  - `sstable.ts` 三份解析循环合并为 `iterEntries()`,越界策略统一。
  - `vacuum()` 返回真实压缩层数(修复前硬编码 6 且底部层永不压缩)。
  - `close()` 加 try/finally(落盘失败也必须释放后端/锁并复位状态)。
  - `getRecoveryReport()`:{droppedSSTables, dataLossSuspected, walGaps, legacyImported,
    manifestFallback} —— "自愈了什么、有没有真丢数据"成为可读返回值。

五、验证
  - 新增 `tests/v080-b6-single-commit-point.test.ts`(63 项,含 manifest 严格校验表驱动 25 例)。
  - 新增 `scripts/mutation-b6.py`:22 项变异验证(把每个修复回退到修复前行为,对应用例必须失败),
    全部被拦住 —— 这批用例不是陪跑。
  - 常规套件 1935 通过 / 91 套件;覆盖率 90.34 / 82.16 / 94.06 / 93.23(阈值 90/82/94/93);
    e2e 14/14;重型套件 4 套件 27 项全绿。
2026-09-15 10:29:03 +08:00

1405 lines
66 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.
/**
* v0.8.0(B-6)回归套件 —— 存储层**单一提交点**`__aria_manifest`)与 LSM 结构根治
* ============================================================================
* 本套件覆盖 PLAN-v0.7.5.md 工作流 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;
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; });
}
release(): void {
this.gatePattern = null;
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 卡住"
this.gatePattern = null;
const notify = this.onGated;
if (notify) { this.onGated = null; notify(); }
await this.gate;
}
}
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(); }
}
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 inner = new MemoryBackend();
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); // 没有假装提交成功
void inner;
});
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 new Promise((r) => setTimeout(r, 20));
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 new Promise((r) => setTimeout(r, 30));
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();
const gatedOutcome = await Promise.race([
hit.then(() => 'gated'),
new Promise((r) => setTimeout(() => r('no-compaction'), 5000)),
]);
expect(gatedOutcome).toBe('gated'); // compaction 确实开始了且被卡住
// compaction 现在卡在门控上(维护链不前进)。写路径必须继续工作:
// 修复前 checkpoint → lsm.flush() → drainMaintenance() 会一起卡住
//v0.6.1 记录的"8~11s 悬崖"的成因之一)。
const outcome = await Promise.race([
engine.insert('t', [{ id: 'after-gate', v: 999, pad: big }]).then(() => 'ok', () => 'error'),
new Promise((r) => setTimeout(() => r('BLOCKED'), 5000)),
]);
expect(outcome).toBe('ok');
gated.release();
await lsm.flush();
expect(await engine.count('t')).toBe(5);
});
});
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('冻结表存在时,水位不会推进越过它(manifest 如实记录意图)', async () => {
const dbName = uniqueDB('b6-intent-floor');
const backend = new MemoryBackend();
const engine = openEngine(dbName, backend, { memtableSizeThreshold: 128 });
await engine.open(dbName, 1);
await engine.createTable(SCHEMA());
await engine.insert('t', rows(20));
// 强行再冻结一份(不 flush),让"待落盘意图"存在
(engine as any).lsm.freezeMemtable();
const manifest = (await readManifestState(backend))!;
// 有意让意图进入 manifest(提交一次),断言水位 <= 最早意图的 lsnAtFreeze
await (engine as any).commitManifest();
const after = (await readManifestState(backend))!;
void manifest;
expect(after.frozen.length).toBeGreaterThan(0);
const minIntentLsn = Math.min(...after.frozen.map((f) => f.lsnAtFreeze));
expect(after.wal.startLsn).toBeLessThanOrEqual(minIntentLsn);
await engine.close();
});
});
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([]);
});
});