Compare commits
2 commits
1fb5109351
...
e9ede5774b
| Author | SHA1 | Date | |
|---|---|---|---|
| e9ede5774b | |||
| 293b131fa2 |
2 changed files with 21 additions and 462 deletions
|
|
@ -40,10 +40,7 @@
|
||||||
* The manifest (`_generations/facts/manifest.json`, JSON — forensics stay
|
* The manifest (`_generations/facts/manifest.json`, JSON — forensics stay
|
||||||
* terminal-readable) is the single source of truth for the segment SET;
|
* terminal-readable) is the single source of truth for the segment SET;
|
||||||
* rotation flips it atomically (write-new → fsync → rename) BEFORE the new
|
* rotation flips it atomically (write-new → fsync → rename) BEFORE the new
|
||||||
* tail's first byte exists, so no segment file is ever unaccounted for. Its
|
* tail's first byte exists, so no segment file is ever unaccounted for.
|
||||||
* per-segment `firstGeneration`/`lastGeneration` are LOAD-BEARING at open: a
|
|
||||||
* recovery pass looking for facts above a bound reads only the segments those
|
|
||||||
* bounds cannot rule out (the prune law — see `segmentsHoldingFactsAbove`).
|
|
||||||
*
|
*
|
||||||
* ## Mixed-version logs (the v2 live-write cutover)
|
* ## Mixed-version logs (the v2 live-write cutover)
|
||||||
*
|
*
|
||||||
|
|
@ -692,74 +689,6 @@ function parseSegment(
|
||||||
return { facts, validBytes: offset, formatVersion: FACT_LOG_FORMAT_V1 }
|
return { facts, validBytes: offset, formatVersion: FACT_LOG_FORMAT_V1 }
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* THE PRUNE LAW — which segment files a pass looking for facts ABOVE
|
|
||||||
* `committedGeneration` actually has to read, and how many the manifest's own
|
|
||||||
* recorded bounds took off the table.
|
|
||||||
*
|
|
||||||
* A sealed segment's `lastGeneration` is written at SEAL time and never
|
|
||||||
* mutated upward afterwards ({@link FactLog.rotate}, unchanged since the log
|
|
||||||
* was introduced): the tail's bytes are fsynced FIRST (`await this.sync()` —
|
|
||||||
* "sealed segments are always fully durable"), the entry is then built from
|
|
||||||
* the content that fsync covered, and only then does the manifest flip —
|
|
||||||
* atomically (tmp+rename) and fsynced — which in the SAME write re-points
|
|
||||||
* `tailSegment` at a new file, so the sealed file is never appended to again.
|
|
||||||
* A crash anywhere in that order is safe in the pruning direction: crash
|
|
||||||
* before the manifest write and the segment is still the TAIL (read whole);
|
|
||||||
* crash after it and the entry describes bytes that were already durable. The
|
|
||||||
* only later mutation of a sealed segment is `open()`'s straddle truncation,
|
|
||||||
* which REMOVES facts and re-derives the entry from the actual bytes — so a
|
|
||||||
* recorded bound can drift DOWN with its file, never up.
|
|
||||||
*
|
|
||||||
* Therefore: `lastGeneration = L` proves the file holds no fact above L, and
|
|
||||||
* a pass above `committedGeneration >= L` can skip it whole — no read, no
|
|
||||||
* CRC decode, no msgpack. What the manifest cannot PROVE is never pruned: an
|
|
||||||
* entry with no numeric `lastGeneration` (a legacy or hand-repaired manifest)
|
|
||||||
* is read, and the unsealed tail is always read.
|
|
||||||
*
|
|
||||||
* This is the difference between an open that costs O(whole fact log) and one
|
|
||||||
* that costs O(the facts that could matter). MEASURED in production: a 16k-row
|
|
||||||
* brain at generation ~478,819 paid 34-37s of segment reads and CRC decoding
|
|
||||||
* in `generation-store-open-fold` on EVERY open — to answer a question whose
|
|
||||||
* answer, after a clean close, is always "nothing".
|
|
||||||
*/
|
|
||||||
function segmentsHoldingFactsAbove(
|
|
||||||
stored: FactsManifest,
|
|
||||||
committedGeneration: number
|
|
||||||
): { files: string[]; pruned: number } {
|
|
||||||
const files: string[] = []
|
|
||||||
let pruned = 0
|
|
||||||
for (const entry of stored.segments) {
|
|
||||||
const last = (entry as Partial<SegmentEntry>).lastGeneration
|
|
||||||
if (typeof last === 'number' && Number.isFinite(last) && last <= committedGeneration) {
|
|
||||||
pruned++
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
files.push(entry.file)
|
|
||||||
}
|
|
||||||
if (stored.tailSegment) files.push(stored.tailSegment)
|
|
||||||
return { files, pruned }
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Say what the open actually read. One line, and only when the log holds more
|
|
||||||
* than one segment (a single-segment log has nothing to prune and nothing to
|
|
||||||
* report) — the operator's receipt that the open is paying for the tail, not
|
|
||||||
* for the whole history.
|
|
||||||
*/
|
|
||||||
function narrateAboveScan(
|
|
||||||
pass: string,
|
|
||||||
committedGeneration: number,
|
|
||||||
read: number,
|
|
||||||
pruned: number
|
|
||||||
): void {
|
|
||||||
if (read + pruned <= 1) return
|
|
||||||
prodLog.narrate(
|
|
||||||
`[FactLog] ${pass} above generation ${committedGeneration}: ${read} segment(s) read, ` +
|
|
||||||
`${pruned} pruned of ${read + pruned} (sealed at or below the bound)`
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The generation fact log. One instance per open store; every method assumes
|
* The generation fact log. One instance per open store; every method assumes
|
||||||
* the single-writer discipline the generation store already enforces (calls
|
* the single-writer discipline the generation store already enforces (calls
|
||||||
|
|
@ -825,6 +754,22 @@ export class FactLog {
|
||||||
return this.manifest.brainId !== undefined || this.tailVersion === FACT_LOG_FORMAT_V2
|
return this.manifest.brainId !== undefined || this.tailVersion === FACT_LOG_FORMAT_V2
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Open the log and reconcile it to committed truth: read the manifest,
|
||||||
|
* establish the tail's intact content (torn-tail scan), then TRUNCATE any
|
||||||
|
* fact with `generation > committedGeneration` — those never committed (a
|
||||||
|
* crash between fact-append and the commit point). After open, the log is
|
||||||
|
* exactly the committed prefix.
|
||||||
|
*/
|
||||||
|
/**
|
||||||
|
* Read (without truncating) every intact fact ABOVE a generation — the
|
||||||
|
* log-authority recovery surface: after a crash, facts beyond the
|
||||||
|
* manifest watermark that survived with valid CRCs are ACKED writes in
|
||||||
|
* durable-at-ack mode, and the owner REPLAYS them instead of letting
|
||||||
|
* open() truncate them. Must be called BEFORE open() (it reads the raw
|
||||||
|
* segments directly; the torn tail's invalid suffix is ignored exactly
|
||||||
|
* like open() would).
|
||||||
|
*/
|
||||||
/**
|
/**
|
||||||
* STREAMING twin of {@link FactLog.peekFactsAbove} for the recovery fold:
|
* STREAMING twin of {@link FactLog.peekFactsAbove} for the recovery fold:
|
||||||
* yields facts above the bound one SEGMENT at a time, ascending, without
|
* yields facts above the bound one SEGMENT at a time, ascending, without
|
||||||
|
|
@ -834,18 +779,13 @@ export class FactLog {
|
||||||
* Works manifest-direct (safe before {@link FactLog.open}). Ordering is
|
* Works manifest-direct (safe before {@link FactLog.open}). Ordering is
|
||||||
* structural (segments rotate in order; appends are ordered within one) and
|
* structural (segments rotate in order; appends are ordered within one) and
|
||||||
* ASSERTED — a violation aborts loudly, never a silent misordered replay.
|
* ASSERTED — a violation aborts loudly, never a silent misordered replay.
|
||||||
*
|
|
||||||
* Reads only the segments that CAN hold a fact above the bound — see
|
|
||||||
* {@link segmentsHoldingFactsAbove}. A bounded fold above a high checkpoint
|
|
||||||
* therefore reads its own tail, not the whole history it already proved
|
|
||||||
* durable.
|
|
||||||
*/
|
*/
|
||||||
async *streamFactsAbove(committedGeneration: number): AsyncGenerator<CommitFact[], void> {
|
async *streamFactsAbove(committedGeneration: number): AsyncGenerator<CommitFact[], void> {
|
||||||
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
||||||
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return
|
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return
|
||||||
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return
|
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return
|
||||||
const { files, pruned } = segmentsHoldingFactsAbove(stored, committedGeneration)
|
const files = [...stored.segments.map((s) => s.file)]
|
||||||
narrateAboveScan('recovery fold', committedGeneration, files.length, pruned)
|
if (stored.tailSegment) files.push(stored.tailSegment)
|
||||||
let lastGen = committedGeneration
|
let lastGen = committedGeneration
|
||||||
for (const file of files) {
|
for (const file of files) {
|
||||||
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
|
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
|
||||||
|
|
@ -867,27 +807,13 @@ export class FactLog {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Read (without truncating) every intact fact ABOVE a generation — the
|
|
||||||
* log-authority recovery surface: after a crash, facts beyond the
|
|
||||||
* manifest watermark that survived with valid CRCs are ACKED writes in
|
|
||||||
* durable-at-ack mode, and the owner REPLAYS them instead of letting
|
|
||||||
* open() truncate them. Must be called BEFORE open() (it reads the raw
|
|
||||||
* segments directly; the torn tail's invalid suffix is ignored exactly
|
|
||||||
* like open() would).
|
|
||||||
*
|
|
||||||
* Reads only the segments that CAN hold such a fact — see
|
|
||||||
* {@link segmentsHoldingFactsAbove}. This runs on EVERY log-authority open,
|
|
||||||
* including the clean one where the answer is always empty, so the segments
|
|
||||||
* the manifest already proves irrelevant are never opened at all.
|
|
||||||
*/
|
|
||||||
async peekFactsAbove(committedGeneration: number): Promise<CommitFact[]> {
|
async peekFactsAbove(committedGeneration: number): Promise<CommitFact[]> {
|
||||||
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
||||||
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return []
|
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return []
|
||||||
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return []
|
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return []
|
||||||
const out: CommitFact[] = []
|
const out: CommitFact[] = []
|
||||||
const { files, pruned } = segmentsHoldingFactsAbove(stored, committedGeneration)
|
const files = [...stored.segments.map((s) => s.file)]
|
||||||
narrateAboveScan('above-manifest peek', committedGeneration, files.length, pruned)
|
if (stored.tailSegment) files.push(stored.tailSegment)
|
||||||
for (const file of files) {
|
for (const file of files) {
|
||||||
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
|
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
|
||||||
if (bytes === null) continue
|
if (bytes === null) continue
|
||||||
|
|
@ -900,13 +826,6 @@ export class FactLog {
|
||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Open the log and reconcile it to committed truth: read the manifest,
|
|
||||||
* establish the tail's intact content (torn-tail scan), then TRUNCATE any
|
|
||||||
* fact with `generation > committedGeneration` — those never committed (a
|
|
||||||
* crash between fact-append and the commit point). After open, the log is
|
|
||||||
* exactly the committed prefix.
|
|
||||||
*/
|
|
||||||
async open(committedGeneration: number): Promise<void> {
|
async open(committedGeneration: number): Promise<void> {
|
||||||
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
||||||
if (stored && typeof stored === 'object' && Array.isArray(stored.segments)) {
|
if (stored && typeof stored === 'object' && Array.isArray(stored.segments)) {
|
||||||
|
|
|
||||||
|
|
@ -1,360 +0,0 @@
|
||||||
/**
|
|
||||||
* @module tests/integration/factlog-open-prune
|
|
||||||
* @description THE OPEN READS THE TAIL, NOT THE HISTORY.
|
|
||||||
*
|
|
||||||
* Every log-authority open asks the fact log one question — "is there a fact
|
|
||||||
* above the committed pointer?" — and until this lane existed it answered by
|
|
||||||
* reading and CRC-decoding EVERY segment file the manifest names. MEASURED in
|
|
||||||
* production on a 16k-row brain at generation ~478,819: 34-37 seconds inside
|
|
||||||
* the `generation-store-open-fold` phase, on every open, including the clean
|
|
||||||
* one where the answer is always "nothing".
|
|
||||||
*
|
|
||||||
* The manifest already records each sealed segment's `lastGeneration`, written
|
|
||||||
* at seal time AFTER the segment's bytes are fsynced and into a manifest that
|
|
||||||
* is itself written atomically and fsynced — and a sealed file is never
|
|
||||||
* appended to again (the same manifest flip re-points `tailSegment`). So an
|
|
||||||
* entry recording `lastGeneration ≤ committed` PROVES its file holds nothing
|
|
||||||
* above the bound, and the open can skip it whole.
|
|
||||||
*
|
|
||||||
* Pinned here, from the log's own counters (the narration line), never a clock:
|
|
||||||
*
|
|
||||||
* 1. A clean close and reopen on a log with ≥4 sealed segments reads
|
|
||||||
* EXACTLY the tail (1 of 6), prunes the rest, and finds nothing.
|
|
||||||
* 2. A real SIGKILLed process that sealed segments holding facts ABOVE the
|
|
||||||
* committed pointer: the reopen READS those sealed segments and recovers
|
|
||||||
* byte-identically to an unpruned open (differential — the same store,
|
|
||||||
* with the provable field stripped from its manifest, takes the full-scan
|
|
||||||
* path and must agree fact for fact, before and after `open()`).
|
|
||||||
* 3. A manifest entry with no `lastGeneration` (legacy, or hand-repaired) is
|
|
||||||
* READ. Never prune what the manifest cannot prove.
|
|
||||||
*/
|
|
||||||
import { describe, it, expect, afterEach } from 'vitest'
|
|
||||||
import * as fs from 'node:fs'
|
|
||||||
import * as os from 'node:os'
|
|
||||||
import * as path from 'node:path'
|
|
||||||
import { spawn } from 'node:child_process'
|
|
||||||
import {
|
|
||||||
FactLog,
|
|
||||||
FACTS_MANIFEST_PATH,
|
|
||||||
type CommitFact,
|
|
||||||
type FactLogStorage
|
|
||||||
} from '../../src/db/factLog.js'
|
|
||||||
import { FileSystemStorage } from '../../src/storage/adapters/fileSystemStorage.js'
|
|
||||||
|
|
||||||
const REPO_ROOT = process.cwd()
|
|
||||||
const TSX = path.join(REPO_ROOT, 'node_modules', '.bin', 'tsx')
|
|
||||||
/** ~1KB frames against a 4KB rotation threshold: ~5 facts per segment. */
|
|
||||||
const ROTATE_BYTES = 4096
|
|
||||||
|
|
||||||
const tmpDirs: string[] = []
|
|
||||||
function makeTempDir(): string {
|
|
||||||
const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-factlog-prune-'))
|
|
||||||
tmpDirs.push(dir)
|
|
||||||
return dir
|
|
||||||
}
|
|
||||||
|
|
||||||
afterEach(() => {
|
|
||||||
for (const dir of tmpDirs.splice(0)) {
|
|
||||||
try {
|
|
||||||
fs.rmSync(dir, { recursive: true, force: true })
|
|
||||||
} catch {
|
|
||||||
/* best effort */
|
|
||||||
}
|
|
||||||
try {
|
|
||||||
fs.rmSync(`${dir}.ready.json`, { force: true })
|
|
||||||
} catch {
|
|
||||||
/* best effort */
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})
|
|
||||||
|
|
||||||
const UUID = (n: number): string => `00000000-0000-4000-8000-${String(n).padStart(12, '0')}`
|
|
||||||
|
|
||||||
/** One ~1KB fact — the padding is what makes rotation cheap to provoke. */
|
|
||||||
function fact(generation: number): CommitFact {
|
|
||||||
return {
|
|
||||||
generation,
|
|
||||||
timestamp: 1_700_000_000_000 + generation,
|
|
||||||
ops: [
|
|
||||||
{
|
|
||||||
kind: 'noun',
|
|
||||||
id: UUID(generation),
|
|
||||||
record: {
|
|
||||||
metadata: { noun: 'document', pad: 'x'.repeat(900), g: generation },
|
|
||||||
vector: null
|
|
||||||
}
|
|
||||||
}
|
|
||||||
]
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A deterministic int minter so the log writes the V2 format production
|
|
||||||
* writes (the prune is a manifest-level decision and never touches segment
|
|
||||||
* bytes — but the pins should run against the bytes the fleet actually has).
|
|
||||||
*/
|
|
||||||
function makeMinter(): (kind: 'noun' | 'verb', id: string) => bigint {
|
|
||||||
const ints = new Map<string, bigint>()
|
|
||||||
return (kind, id) => {
|
|
||||||
const key = `${kind}:${id}`
|
|
||||||
let minted = ints.get(key)
|
|
||||||
if (minted === undefined) {
|
|
||||||
minted = BigInt(ints.size + 1)
|
|
||||||
ints.set(key, minted)
|
|
||||||
}
|
|
||||||
return minted
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Open a fact log over a store directory (a fresh adapter each time — this is
|
|
||||||
* what a reopen actually does). */
|
|
||||||
async function openStore(dir: string): Promise<{ storage: any; log: FactLog }> {
|
|
||||||
const storage: any = new FileSystemStorage(dir)
|
|
||||||
await storage.init()
|
|
||||||
const log = new FactLog(storage as FactLogStorage, { rotateBytes: ROTATE_BYTES })
|
|
||||||
log.setIntMinter(makeMinter())
|
|
||||||
return { storage, log }
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Build a log of `count` facts (rotating every ~5), left durable, not closed. */
|
|
||||||
async function buildLog(dir: string, count: number): Promise<number> {
|
|
||||||
const { log } = await openStore(dir)
|
|
||||||
await log.open(0)
|
|
||||||
for (let g = 1; g <= count; g++) await log.append(fact(g))
|
|
||||||
await log.sync()
|
|
||||||
return log.headGeneration()
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Capture the narration channel (`prodLog.narrate` → console.warn). */
|
|
||||||
async function captureNarration<T>(
|
|
||||||
fn: () => Promise<T>
|
|
||||||
): Promise<{ result: T; lines: string[] }> {
|
|
||||||
const lines: string[] = []
|
|
||||||
const original = console.warn
|
|
||||||
console.warn = ((...args: unknown[]) => {
|
|
||||||
lines.push(args.map((a) => String(a)).join(' '))
|
|
||||||
}) as typeof console.warn
|
|
||||||
try {
|
|
||||||
return { result: await fn(), lines }
|
|
||||||
} finally {
|
|
||||||
console.warn = original
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The counters the open narrated — the pin's only source of truth for what
|
|
||||||
* was read (a wall-clock assertion could pass on a warm page cache). */
|
|
||||||
function scanCounts(lines: string[]): { read: number; pruned: number; total: number } {
|
|
||||||
const line = lines.find((l) => l.includes('[FactLog] above-manifest peek above generation'))
|
|
||||||
if (!line) {
|
|
||||||
throw new Error(`no peek narration in:\n${lines.join('\n')}`)
|
|
||||||
}
|
|
||||||
const match = /(\d+) segment\(s\) read, (\d+) pruned of (\d+)/.exec(line)
|
|
||||||
if (!match) throw new Error(`unparsable peek narration: ${line}`)
|
|
||||||
return { read: Number(match[1]), pruned: Number(match[2]), total: Number(match[3]) }
|
|
||||||
}
|
|
||||||
|
|
||||||
interface SegmentEntryOnDisk {
|
|
||||||
file: string
|
|
||||||
firstGeneration: number
|
|
||||||
lastGeneration?: number
|
|
||||||
facts: number
|
|
||||||
bytes: number
|
|
||||||
}
|
|
||||||
|
|
||||||
async function readManifest(dir: string): Promise<{
|
|
||||||
segments: SegmentEntryOnDisk[]
|
|
||||||
tailSegment: string | null
|
|
||||||
}> {
|
|
||||||
const storage: any = new FileSystemStorage(dir)
|
|
||||||
await storage.init()
|
|
||||||
return (await storage.readRawObject(FACTS_MANIFEST_PATH)) as any
|
|
||||||
}
|
|
||||||
|
|
||||||
async function rewriteManifest(
|
|
||||||
dir: string,
|
|
||||||
mutate: (manifest: any) => void
|
|
||||||
): Promise<void> {
|
|
||||||
const storage: any = new FileSystemStorage(dir)
|
|
||||||
await storage.init()
|
|
||||||
const manifest = await storage.readRawObject(FACTS_MANIFEST_PATH)
|
|
||||||
mutate(manifest)
|
|
||||||
await storage.writeRawObject(FACTS_MANIFEST_PATH, manifest)
|
|
||||||
await storage.syncRawObjects([FACTS_MANIFEST_PATH])
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Every fact the log holds, in order — the recovered state, read back. */
|
|
||||||
async function allFacts(log: FactLog): Promise<CommitFact[]> {
|
|
||||||
const out: CommitFact[] = []
|
|
||||||
const handle = log.scanFacts()
|
|
||||||
for await (const batch of handle.batches()) out.push(...batch.facts)
|
|
||||||
return out
|
|
||||||
}
|
|
||||||
|
|
||||||
describe('fact log — the open reads only the segments that can hold facts above the bound', () => {
|
|
||||||
it('a clean close + reopen over ≥4 sealed segments reads exactly the tail and finds nothing', async () => {
|
|
||||||
const dir = makeTempDir()
|
|
||||||
const head = await buildLog(dir, 30)
|
|
||||||
|
|
||||||
const manifest = await readManifest(dir)
|
|
||||||
expect(manifest.segments.length).toBeGreaterThanOrEqual(4) // the fixture is real
|
|
||||||
expect(manifest.tailSegment).not.toBeNull()
|
|
||||||
|
|
||||||
// The reopen: a clean close means committed === the log's head.
|
|
||||||
const { log } = await openStore(dir)
|
|
||||||
const { result: orphans, lines } = await captureNarration(() => log.peekFactsAbove(head))
|
|
||||||
|
|
||||||
expect(orphans).toEqual([]) // the fold finds nothing, as it always does after a clean close
|
|
||||||
const counts = scanCounts(lines)
|
|
||||||
expect(counts.read).toBe(1) // EXACTLY the tail
|
|
||||||
expect(counts.total).toBe(manifest.segments.length + 1)
|
|
||||||
expect(counts.pruned).toBe(manifest.segments.length)
|
|
||||||
|
|
||||||
// And the reconciling open still lands on the same committed prefix.
|
|
||||||
await log.open(head)
|
|
||||||
expect(log.headGeneration()).toBe(head)
|
|
||||||
expect((await allFacts(log)).map((f) => f.generation)).toEqual(
|
|
||||||
Array.from({ length: head }, (_, i) => i + 1)
|
|
||||||
)
|
|
||||||
})
|
|
||||||
|
|
||||||
it('a manifest entry with no lastGeneration is READ — never prune what you cannot prove', async () => {
|
|
||||||
const dir = makeTempDir()
|
|
||||||
const head = await buildLog(dir, 30)
|
|
||||||
const before = await readManifest(dir)
|
|
||||||
expect(before.segments.length).toBeGreaterThanOrEqual(4)
|
|
||||||
|
|
||||||
// A legacy/hand-repaired entry: the field the prune needs is simply absent.
|
|
||||||
await rewriteManifest(dir, (m) => {
|
|
||||||
delete m.segments[0].lastGeneration
|
|
||||||
})
|
|
||||||
|
|
||||||
const { log } = await openStore(dir)
|
|
||||||
const { result: orphans, lines } = await captureNarration(() => log.peekFactsAbove(head))
|
|
||||||
|
|
||||||
expect(orphans).toEqual([]) // still nothing above the bound — it was READ to find out
|
|
||||||
const counts = scanCounts(lines)
|
|
||||||
expect(counts.read).toBe(2) // the unprovable entry + the tail
|
|
||||||
expect(counts.pruned).toBe(before.segments.length - 1)
|
|
||||||
expect(counts.total).toBe(before.segments.length + 1)
|
|
||||||
})
|
|
||||||
|
|
||||||
it(
|
|
||||||
'a SIGKILLed writer that sealed segments above the committed pointer recovers identically to an unpruned open',
|
|
||||||
async () => {
|
|
||||||
const dir = makeTempDir()
|
|
||||||
const readyPath = `${dir}.ready.json`
|
|
||||||
// A real process death: the child fsyncs its segments, records what it
|
|
||||||
// reached, and SIGKILLs ITSELF — no close, no unwind, no chance to tidy.
|
|
||||||
const script = `
|
|
||||||
import * as fs from 'node:fs'
|
|
||||||
import { FactLog } from ${JSON.stringify(path.join(REPO_ROOT, 'src', 'db', 'factLog.ts'))}
|
|
||||||
import { FileSystemStorage } from ${JSON.stringify(path.join(REPO_ROOT, 'src', 'storage', 'adapters', 'fileSystemStorage.ts'))}
|
|
||||||
const UUID = (n) => '00000000-0000-4000-8000-' + String(n).padStart(12, '0')
|
|
||||||
const fact = (g) => ({
|
|
||||||
generation: g,
|
|
||||||
timestamp: 1700000000000 + g,
|
|
||||||
ops: [{ kind: 'noun', id: UUID(g), record: { metadata: { noun: 'document', pad: 'x'.repeat(900), g }, vector: null } }]
|
|
||||||
})
|
|
||||||
const ints = new Map()
|
|
||||||
const storage = new FileSystemStorage(${JSON.stringify(dir)})
|
|
||||||
await storage.init()
|
|
||||||
const log = new FactLog(storage, { rotateBytes: ${ROTATE_BYTES} })
|
|
||||||
log.setIntMinter((kind, id) => {
|
|
||||||
const key = kind + ':' + id
|
|
||||||
if (!ints.has(key)) ints.set(key, BigInt(ints.size + 1))
|
|
||||||
return ints.get(key)
|
|
||||||
})
|
|
||||||
await log.open(0)
|
|
||||||
for (let g = 1; g <= 30; g++) await log.append(fact(g))
|
|
||||||
await log.sync()
|
|
||||||
fs.writeFileSync(${JSON.stringify(readyPath)}, JSON.stringify({ head: log.headGeneration() }))
|
|
||||||
process.kill(process.pid, 'SIGKILL')
|
|
||||||
`
|
|
||||||
const scriptPath = path.join(dir, 'crash-writer.mts')
|
|
||||||
fs.writeFileSync(scriptPath, script)
|
|
||||||
const child = spawn(TSX, [scriptPath], { cwd: REPO_ROOT, stdio: ['ignore', 'pipe', 'pipe'] })
|
|
||||||
let output = ''
|
|
||||||
child.stdout.on('data', (d) => { output += String(d) })
|
|
||||||
child.stderr.on('data', (d) => { output += String(d) })
|
|
||||||
const exit = await new Promise<{ code: number | null; signal: string | null }>((resolve) =>
|
|
||||||
child.on('exit', (code, signal) => resolve({ code, signal }))
|
|
||||||
)
|
|
||||||
if (!fs.existsSync(readyPath)) {
|
|
||||||
throw new Error(`the crash writer never reached its kill point:\n${output}`)
|
|
||||||
}
|
|
||||||
// Death, not a shutdown: no close(), no unwind, no orderly exit code.
|
|
||||||
expect(exit.signal ?? `code ${exit.code}`).not.toBe('code 0')
|
|
||||||
const head = JSON.parse(fs.readFileSync(readyPath, 'utf8')).head as number
|
|
||||||
expect(head).toBe(30)
|
|
||||||
|
|
||||||
// The committed pointer the survivor comes back on: mid-log, so sealed
|
|
||||||
// segments hold facts ABOVE it — the exact shape the prune must not skip.
|
|
||||||
const committed = 12
|
|
||||||
const manifest = await readManifest(dir)
|
|
||||||
const straddling = manifest.segments.filter(
|
|
||||||
(s) => s.firstGeneration <= committed && (s.lastGeneration ?? 0) > committed
|
|
||||||
)
|
|
||||||
const entirelyAbove = manifest.segments.filter((s) => s.firstGeneration > committed)
|
|
||||||
expect(straddling.length).toBeGreaterThanOrEqual(1)
|
|
||||||
expect(entirelyAbove.length).toBeGreaterThanOrEqual(1)
|
|
||||||
|
|
||||||
// THE DIFFERENTIAL. The unpruned answer, through the SAME code on the
|
|
||||||
// SAME bytes: a peek above generation 0 can prune nothing (no sealed
|
|
||||||
// segment ends at or below 0), so it reads every segment file and
|
|
||||||
// decodes every frame — exactly what this open used to do — and its
|
|
||||||
// facts above the pointer are what the fold is entitled to replay.
|
|
||||||
const { log } = await openStore(dir)
|
|
||||||
const { result: fullScan, lines: fullLines } = await captureNarration(() =>
|
|
||||||
log.peekFactsAbove(0)
|
|
||||||
)
|
|
||||||
expect(scanCounts(fullLines)).toEqual({
|
|
||||||
read: manifest.segments.length + 1,
|
|
||||||
pruned: 0,
|
|
||||||
total: manifest.segments.length + 1
|
|
||||||
})
|
|
||||||
const unprunedAnswer = fullScan.filter((f) => f.generation > committed)
|
|
||||||
|
|
||||||
const { result: prunedAnswer, lines } = await captureNarration(() =>
|
|
||||||
log.peekFactsAbove(committed)
|
|
||||||
)
|
|
||||||
|
|
||||||
// The sealed segments above the bound were READ, not skipped.
|
|
||||||
const counts = scanCounts(lines)
|
|
||||||
expect(counts.read).toBe(straddling.length + entirelyAbove.length + 1)
|
|
||||||
expect(counts.pruned).toBe(manifest.segments.length - straddling.length - entirelyAbove.length)
|
|
||||||
expect(counts.pruned).toBeGreaterThan(0) // the prune did engage, and was still right
|
|
||||||
expect(prunedAnswer.map((f) => f.generation)).toEqual(
|
|
||||||
Array.from({ length: head - committed }, (_, i) => committed + 1 + i)
|
|
||||||
)
|
|
||||||
// Facts that live in a SEALED segment (not the tail) came back.
|
|
||||||
expect(prunedAnswer.some((f) => f.generation <= (straddling[0].lastGeneration ?? 0))).toBe(
|
|
||||||
true
|
|
||||||
)
|
|
||||||
// Fact for fact, the pruned answer IS the unpruned answer — so whatever
|
|
||||||
// the recovery replays, it replays identically.
|
|
||||||
expect(prunedAnswer).toEqual(unprunedAnswer)
|
|
||||||
|
|
||||||
// The fold's streaming twin (the unclean-open path) agrees too.
|
|
||||||
const streamed: CommitFact[] = []
|
|
||||||
for await (const batch of log.streamFactsAbove(committed)) streamed.push(...batch)
|
|
||||||
expect(streamed).toEqual(unprunedAnswer)
|
|
||||||
|
|
||||||
// And the reconciling open rolls back exactly as it always did: the two
|
|
||||||
// never-committed sealed segments dropped, the straddling one cut, the
|
|
||||||
// tail truncated — the log left as the committed prefix.
|
|
||||||
await log.open(committed)
|
|
||||||
expect(log.headGeneration()).toBe(committed)
|
|
||||||
expect((await allFacts(log)).map((f) => f.generation)).toEqual(
|
|
||||||
Array.from({ length: committed }, (_, i) => i + 1)
|
|
||||||
)
|
|
||||||
const after = await readManifest(dir)
|
|
||||||
expect(after.segments.map((s) => s.file)).toEqual(
|
|
||||||
manifest.segments
|
|
||||||
.filter((s) => s.firstGeneration <= committed)
|
|
||||||
.map((s) => s.file)
|
|
||||||
)
|
|
||||||
expect(after.segments[after.segments.length - 1].lastGeneration).toBe(committed)
|
|
||||||
},
|
|
||||||
120_000
|
|
||||||
)
|
|
||||||
})
|
|
||||||
Reference in a new issue