方法:四个对抗性子代理分头审查(数据正确性 / 文档宣称 vs 实现 / 公共 API 契约 / 测试质量),每条结论要求可复现证据;逐条复核 + 探针确认 + 变异验证(40 项全部 被对应用例拦住)。 P0:事务活跃期间 repair()/close()/周期 checkpoint 推进 WAL 水位 → 已 COMMIT 的 事务整批消失且恢复报告"干净"。根因 hasPendingFlushData()/computeDurableLsn() 不看 txnSnapshot;守卫此前只在 CheckpointManager 两个回调里。修复:守卫下沉到 computeDurableLsn() 与 advanceWalCheckpoint() 入口(唯一实现)。 P1: - WAL 前缀缺失丢弃整段活分片(回退上一代 manifest 时 kept 为空)→ 前缀缺失单独 记录,后缀照常重放;仅 fromLsn === 0 时才算真异常 - 孤儿回收门槛只看引擎层 dataLossSuspected,漏掉 LSM 层被丢的 SSTable → 统一 describeRecoveryDamage() 聚合判定(损坏时绝不删"引用不到"的文件) - vacuum() 逐层压缩绕过维护链 → vacuumLevels() 每层作为维护链任务执行 - reclaimRetiredNow() 无视在途读者(读者把"已退休"读成"文件损坏")→ 有读者时 退化为延迟回收 P2:WAL 记录级 CRC 损坏不计数不上报;旧格式表结构记录形状损坏静默当空库; bloomFilterBitsPerKey 配置被接受却完全不生效(构建器写死默认值,实现缺陷); 幽灵 meta;介质读故障等于文件损坏的语义无用例;manifest 回读校验两条守卫无用例; 文件名≠载荷世代判定无用例;pageIdWatermark 单调性无用例;分片号两条真实不变量 无用例。 覆盖率口径(第二处漏洞):interface.ts 混着三个运行时函数(cloneRow 等)却被 描述为"纯类型、不纳入统计" → 实现搬到 src/engine/row_clone.ts;搬完门禁真的 失败(functions 93.84% < 94%),补测退化路径后通过。 测试质量:3 条空壳用例改值级断言;1 条"全损坏"用例实际只走缓存 → 拆成两条真 用例;5 秒墙钟 race 改门控 + 失败上限;setTimeout 改 whenIdle();<= 收紧为 <。 变异脚本加固:正控(干净基线必须全绿)、编译失败/0 用例单独归类、300s 超时、 逐字节 sha256 恢复校验、O_EXCL 进程锁、锚点唯一性;变异 22 → 40 项。 文档两轮订正(16 + 11 条不成立宣称):MVCC 快照隔离、backup 一致性快照、 "空洞检测截断"、体积(251,109 B / gzip 63,145 B)、测试与覆盖率数字、 "5 种存储引擎"、Tree-shakable、错误码表补 16 个码、恢复报告字段、已知限制 (回退单向 / 多实例依赖 Web Locks / manifest 体积 / 尾部 WAL 分片不可识别)。 验证:常规套件 92 套件 / 1980 用例全绿;覆盖率 90.59 / 82.59 / 94.14 / 93.50 (阈值 90/82/94/93);e2e 14/14(真实 Chromium + OPFS + CDP 崩溃); 重型套件 4 套件 / 27 用例;变异 40/40;lint + 两份 tsc 干净;dist 已重建。
2180 lines
104 KiB
TypeScript
2180 lines
104 KiB
TypeScript
/**
|
||
* 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;
|
||
/** 粘性门控:一旦命中就持续阻塞(不是只阻塞一次) */
|
||
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_FAILED,meta 不被自愈删除', 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] LSM:flush 重试与错误报告顺序', () => {
|
||
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);
|
||
// 不 close(MemoryBackend.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 });
|
||
// 触发后台自动 flush(freezeMemtable → 入链 → 第一次 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);
|
||
// 第二次 save(compaction 产物)会命中门控 → 长 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);
|
||
|
||
// 不 close(MemoryBackend.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 SSTable(3 < 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);
|
||
});
|
||
});
|