diff --git a/.github/workflows/ci.yml b/.forgejo/workflows/ci.yml
similarity index 100%
rename from .github/workflows/ci.yml
rename to .forgejo/workflows/ci.yml
diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md
index ef9c4a51..d277091d 100644
--- a/CONTRIBUTING.md
+++ b/CONTRIBUTING.md
@@ -10,10 +10,6 @@ The source of truth is a self-hosted forge: **source.soulcraft.com/soulcraft/bra
It's anonymously readable and cloneable — no account needed to browse, clone,
or build.
-**github.com/soulcraftlabs/brainy** is a public read-only mirror. It's a fine
-place to read code or star the project, but issues and pull requests opened
-there won't be picked up — please use one of the paths below instead.
-
## How to contribute
**Found a bug, or have an idea?** Email **brainy@soulcraft.com**. No account,
diff --git a/README.md b/README.md
index 2fc42060..2caf6493 100644
--- a/README.md
+++ b/README.md
@@ -13,7 +13,7 @@
-
+
diff --git a/package.json b/package.json
index e4bc8144..a3ece83c 100644
--- a/package.json
+++ b/package.json
@@ -128,13 +128,13 @@
"publishConfig": {
"access": "public"
},
- "homepage": "https://github.com/soulcraftlabs/brainy",
+ "homepage": "https://source.soulcraft.com/soulcraft/brainy",
"bugs": {
- "url": "https://github.com/soulcraftlabs/brainy/issues"
+ "url": "https://source.soulcraft.com/soulcraft/brainy/issues"
},
"repository": {
"type": "git",
- "url": "git+https://github.com/soulcraftlabs/brainy.git"
+ "url": "git+https://source.soulcraft.com/soulcraft/brainy.git"
},
"files": [
"dist/**/*.js",
diff --git a/scripts/release.sh b/scripts/release.sh
index 42f5b345..43fa50bd 100755
--- a/scripts/release.sh
+++ b/scripts/release.sh
@@ -142,7 +142,7 @@ else
fi
# Create new changelog entry
-CHANGELOG_ENTRY="### [${NEW_VERSION}](https://github.com/soulcraftlabs/brainy/compare/v${CURRENT_VERSION}...v${NEW_VERSION}) ($(date +%Y-%m-%d))
+CHANGELOG_ENTRY="### [${NEW_VERSION}](https://source.soulcraft.com/soulcraft/brainy/compare/v${CURRENT_VERSION}...v${NEW_VERSION}) ($(date +%Y-%m-%d))
${COMMITS}
"
@@ -175,42 +175,59 @@ echo -e "${BLUE}7️⃣ Creating git tag v${NEW_VERSION}...${NC}"
git tag -a "v${NEW_VERSION}" -m "Release v${NEW_VERSION}"
echo -e "${GREEN}✅ Tag created${NC}\n"
-# Step 9: Push to origin (source of truth) and the public GitHub mirror
+# Step 9: Push to origin — the forge is the one home (ruled 2026-07-23; the
+# old public GitHub repo is archived history, no longer part of any release).
echo -e "${BLUE}8️⃣ Pushing to origin...${NC}"
git push --follow-tags origin "$CURRENT_BRANCH"
echo -e "${GREEN}✅ Pushed to origin${NC}\n"
-# The public GitHub repo is a mirror of origin with an unknown sync cadence.
-# `gh release create` below targets GitHub directly: if the new tag hasn't
-# reached GitHub yet, gh would CREATE it — pointed at GitHub's default-branch
-# head, i.e. the wrong commit. Push branch+tag to GitHub explicitly, then
-# verify the tag resolves there to the same commit before any release is cut.
-GITHUB_URL="https://github.com/soulcraftlabs/brainy.git"
-echo -e "${BLUE}8️⃣½ Pushing to the public GitHub mirror...${NC}"
-git push --follow-tags "$GITHUB_URL" "$CURRENT_BRANCH"
-LOCAL_TAG_SHA="$(git rev-parse "v${NEW_VERSION}^{}")"
-GITHUB_TAG_SHA="$(git ls-remote --tags "$GITHUB_URL" "v${NEW_VERSION}^{}" | cut -f1)"
-if [ "$LOCAL_TAG_SHA" != "$GITHUB_TAG_SHA" ]; then
- echo -e "${RED}❌ Tag v${NEW_VERSION} on GitHub (${GITHUB_TAG_SHA:-absent}) does not match local (${LOCAL_TAG_SHA}) — aborting before npm publish. Fix the mirror, then re-run.${NC}"
+# Step 10: Publish — forge FIRST (home), npmjs second (the world's storefront).
+# The fleet-wide ~/.npmrc maps the @soulcraft scope to the forge registry, and
+# a scope mapping BEATS `--registry` on the command line — so each publish
+# names its registry via the scope override explicitly. Nothing implicit.
+FORGE_NPM_REG="https://source.soulcraft.com/api/packages/soulcraft/npm/"
+FORGE_NPM_TOKEN_FILE="$HOME/.config/soulcraft/npm-publish-brainy.token"
+echo -e "${BLUE}9️⃣ Publishing to the forge registry (home)...${NC}"
+if [ -f "$FORGE_NPM_TOKEN_FILE" ]; then
+ TMPRC="$(mktemp)"
+ chmod 600 "$TMPRC"
+ {
+ echo "@soulcraft:registry=${FORGE_NPM_REG}"
+ echo "//source.soulcraft.com/api/packages/soulcraft/npm/:_authToken=$(cat "$FORGE_NPM_TOKEN_FILE")"
+ } > "$TMPRC"
+ if npm publish --tag "$NPM_TAG" --userconfig "$TMPRC"; then
+ echo -e "${GREEN}✅ Published to the forge${NC}\n"
+ else
+ rm -f "$TMPRC"
+ echo -e "${RED}❌ Forge publish FAILED — aborting before npmjs so the pair never diverges. Fix and re-run.${NC}"
+ exit 1
+ fi
+ rm -f "$TMPRC"
+else
+ echo -e "${RED}❌ Forge publish token missing (${FORGE_NPM_TOKEN_FILE}) — aborting. The forge is home; publish it first or restage the token.${NC}"
exit 1
fi
-echo -e "${GREEN}✅ GitHub mirror has the tag at the right commit${NC}\n"
-# Step 10: Publish to npm
-echo -e "${BLUE}9️⃣ Publishing to npm (dist-tag: ${NPM_TAG})...${NC}"
-npm publish --tag "$NPM_TAG"
+echo -e "${BLUE}9️⃣½ Publishing to npmjs (storefront, dist-tag: ${NPM_TAG})...${NC}"
+npm publish --tag "$NPM_TAG" "--@soulcraft:registry=https://registry.npmjs.org/"
# Brainy is the only PUBLIC @soulcraft package — verify visibility after every publish.
-npm access get status @soulcraft/brainy || true
-echo -e "${GREEN}✅ Published to npm${NC}\n"
+npm access get status @soulcraft/brainy "--@soulcraft:registry=https://registry.npmjs.org/" || true
+echo -e "${GREEN}✅ Published to npmjs${NC}\n"
-# Step 11: Create GitHub release
-echo -e "${BLUE}🔟 Creating GitHub release...${NC}"
-if [ "$PRERELEASE" = true ]; then
- gh release create "v${NEW_VERSION}" --generate-notes --prerelease
+# Step 11: Release object on the forge (presentational — the tag, CHANGELOG,
+# and RELEASES.md are the record; this just gives the forge UI a release page).
+echo -e "${BLUE}🔟 Creating forge release...${NC}"
+if [ -n "${FORGEJO_RELEASE_TOKEN:-}" ]; then
+ if curl -sf -X POST "https://source.soulcraft.com/api/v1/repos/soulcraft/brainy/releases" \
+ -H "Authorization: token ${FORGEJO_RELEASE_TOKEN}" -H "Content-Type: application/json" \
+ -d "{\"tag_name\":\"v${NEW_VERSION}\",\"name\":\"v${NEW_VERSION}\",\"prerelease\":${PRERELEASE}}" >/dev/null; then
+ echo -e "${GREEN}✅ Forge release created${NC}\n"
+ else
+ echo -e "${RED}⚠️ Forge release API call failed — tag + CHANGELOG remain the record; create the release page via the forge UI if wanted${NC}\n"
+ fi
else
- gh release create "v${NEW_VERSION}" --generate-notes
+ echo -e "${RED}⚠️ FORGEJO_RELEASE_TOKEN unset — no release page created; tag + CHANGELOG remain the record${NC}\n"
fi
-echo -e "${GREEN}✅ GitHub release created${NC}\n"
# Step 12: Push public docs to the soulcraft.com docs ingest door
# (VENUE-DOCS-RELEASE-PUSH). Skips with a loud warning when
@@ -229,4 +246,4 @@ echo -e "${GREEN}🎉 Release ${NEW_VERSION} complete!${NC}"
echo -e "${GREEN}━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━${NC}"
echo ""
echo -e "📦 npm: ${BLUE}https://www.npmjs.com/package/@soulcraft/brainy/v/${NEW_VERSION}${NC}"
-echo -e "🐙 GitHub: ${BLUE}https://github.com/soulcraftlabs/brainy/releases/tag/v${NEW_VERSION}${NC}"
+echo -e "🏠 Forge: ${BLUE}https://source.soulcraft.com/soulcraft/brainy/releases/tag/v${NEW_VERSION}${NC}"
diff --git a/src/brainy.ts b/src/brainy.ts
index 37b6e491..fd9e1ffa 100644
--- a/src/brainy.ts
+++ b/src/brainy.ts
@@ -8322,6 +8322,27 @@ export class Brainy implements BrainyInterface {
return this.generationStore.compact(options)
}
+ /**
+ * @description Repack cold generation history into sealed segments —
+ * re-representation, never deletion: every record and delta stays readable
+ * (`asOf()` unchanged); the physical file count drops by orders of
+ * magnitude. Runs automatically (time-bounded) at `close()`; call this for
+ * explicit maintenance windows on long-lived writers. The ONLY history
+ * transform permitted under the archival profile (`retention: 'all'`).
+ * @param options - `timeBudgetMs` bounds the pass (early stop = consistent
+ * prefix, next pass resumes); `batchGenerations` sizes each fold.
+ * @returns Folded generation count and segments created.
+ */
+ async repackHistory(options?: {
+ timeBudgetMs?: number
+ batchGenerations?: number
+ }): Promise<{ foldedGenerations: number; segmentsCreated: number }> {
+ this.assertWritable('repackHistory')
+ await this.ensureInitialized()
+ await this.generationStore.flushPendingSingleOps()
+ return this.generationStore.repackHistory(options)
+ }
+
/**
* @description Read-only generational-history footprint for fleet audits:
* generation count, total on-disk bytes, generation/timestamp range, the
@@ -8349,6 +8370,24 @@ export class Brainy implements BrainyInterface {
}
}
+ /**
+ * @description A deterministic content digest of the generation log through
+ * `g` (D8 — gate-to-generation provenance): identical history produces the
+ * identical digest on any machine; divergence produces a different one.
+ * Release gates and suite verdicts pin `{generation, digest}` and verify
+ * both at execution time instead of pinning a git commit. O(segments +
+ * live-tier window), never O(all generations). Throws `RangeError` out of
+ * range and `GenerationCompactedError` below the horizon — a gate can
+ * never silently pin reclaimed history.
+ * @example
+ * const gate = { generation: brain.generation(), digest: await brain.generationDigest(brain.generation()) }
+ */
+ async generationDigest(g: number): Promise {
+ await this.ensureInitialized()
+ await this.generationStore.flushPendingSingleOps()
+ return this.generationStore.generationDigest(g)
+ }
+
/**
* @description Drive the adaptive retention byte budget at runtime — the
* settable input a machine-level coordinator (e.g. cor's `ResourceManager`,
@@ -16367,11 +16406,21 @@ export class Brainy implements BrainyInterface {
await this.generationStore.flushPendingSingleOps()
}
- // Phase 0b: Auto-compact generational history per config.retention (default
- // on) BEFORE the generation store closes below. This is THE auto-compaction
- // site (8.9.0 — flush() never compacts): time-bounded per pass, respects
- // live Db pins and an explicit autoCompact: false; no-op on read-only
- // instances.
+ // Phase 0b: REPACK cold history into sealed segments (D1+D3 —
+ // re-representation, never deletion; the only history transform under the
+ // archival profile), then auto-compact per config.retention. Repack runs
+ // FIRST so bounded-retention reclaim can drop whole segments. Both are
+ // time-bounded maintenance passes (8.9.0 law: flush() never pays these);
+ // both are housekeeping — failures warn, never fail a clean shutdown.
+ if (!this.isReadOnly && this.generationStore) {
+ try {
+ await this.generationStore.repackHistory({ timeBudgetMs: 5_000 })
+ } catch (error) {
+ console.warn(
+ `History repacking failed (non-fatal): ${error instanceof Error ? error.message : String(error)}`
+ )
+ }
+ }
await this.autoCompactHistory()
// Phase 1: Flush ALL components in parallel to persist buffered data
diff --git a/src/db/factLog.ts b/src/db/factLog.ts
index 4c5e95fd..94e79700 100644
--- a/src/db/factLog.ts
+++ b/src/db/factLog.ts
@@ -102,12 +102,26 @@ export interface FactScanBatch {
segmentId: string
}
+/**
+ * Liveness bound on a scan's FIRST batch (Stage-2 co-freeze, D1 contract):
+ * `batches()` must yield its first batch — or fail loudly — within this many
+ * ms of the first pull. A backlogged or damaged store may be SLOW, but it may
+ * never be SILENT: a consumer awaiting the first batch is otherwise
+ * indistinguishable from a wedge (the exact failure shape a production heal
+ * hit against a generations-backlogged brain).
+ */
+export const SCANFACTS_FIRST_BATCH_MS = 10_000
+
/** The telemetry a scan OPEN returns (frozen shape). */
export interface FactScanHandle {
headGeneration: number
segmentCount: number
approxFactCount: number
- /** Ordered batches; a detected gap aborts LOUDLY, never a silent skip. */
+ /**
+ * Ordered batches; a detected gap aborts LOUDLY, never a silent skip.
+ * Liveness contract: the FIRST batch resolves or rejects within
+ * {@link SCANFACTS_FIRST_BATCH_MS} of the first pull — never a silent hang.
+ */
batches: () => AsyncGenerator
/** Close telemetry — the invariant cross-check, valid after iteration ends. */
summary: () => { factsYielded: number; segmentsRead: number }
@@ -440,6 +454,8 @@ export class FactLog {
toGeneration?: number
kinds?: Array<'noun' | 'verb'>
batchSize?: number
+ /** Test override for the first-batch liveness bound (default {@link SCANFACTS_FIRST_BATCH_MS}). */
+ firstBatchTimeoutMs?: number
}): FactScanHandle {
const from = options?.fromGeneration ?? 1
const to = options?.toGeneration ?? this.head
@@ -514,11 +530,43 @@ export class FactLog {
}
}
+ // Liveness wrapper: the FIRST pull races the contract deadline. Only the
+ // first — the bound is time-to-first-batch (proof the producer is alive),
+ // not per-batch pacing; and it runs only while a pull is actually pending,
+ // so consumer think-time between pulls never counts against the producer.
+ const firstBatchTimeoutMs = options?.firstBatchTimeoutMs ?? SCANFACTS_FIRST_BATCH_MS
+ async function* batchesWithLiveness(this: void): AsyncGenerator {
+ const inner = batches()
+ let timer: NodeJS.Timeout | undefined
+ try {
+ const deadline = new Promise((_, reject) => {
+ timer = setTimeout(
+ () =>
+ reject(
+ new Error(
+ `fact log: scanFacts produced no first batch within ${firstBatchTimeoutMs}ms ` +
+ `(liveness contract) — the store is wedged or unreadably slow; aborting scan LOUDLY ` +
+ `instead of hanging the consumer.`
+ )
+ ),
+ firstBatchTimeoutMs
+ )
+ timer.unref?.()
+ })
+ const first = await Promise.race([inner.next(), deadline])
+ if (first.done) return
+ yield first.value
+ } finally {
+ clearTimeout(timer)
+ }
+ yield* inner
+ }
+
return {
headGeneration: this.head,
segmentCount: segments.length + (tailSnapshot.length > 0 ? 1 : 0),
approxFactCount,
- batches,
+ batches: batchesWithLiveness,
summary: () => ({ factsYielded, segmentsRead })
}
}
diff --git a/src/db/generationSegments.ts b/src/db/generationSegments.ts
new file mode 100644
index 00000000..0c14b60c
--- /dev/null
+++ b/src/db/generationSegments.ts
@@ -0,0 +1,459 @@
+/**
+ * @module db/generationSegments
+ * @description The generation-segment store — Stage-2 D1+D3+repacking's file
+ * format (co-frozen 2026-07-19; design: the d1-d3-repacking spec).
+ *
+ * Packs CONSECUTIVE cold generations' record-sets (before-images + delta)
+ * into append-once segment files with derived sidecar indexes, so history
+ * scales in SEGMENTS (tens) instead of FILES-PER-GENERATION (hundreds of
+ * thousands), and cold-open reads ONE manifest instead of listing the
+ * backlog. Layout under `_generations/segments/`:
+ *
+ * - `seg-.bgs` — magic "BGS1", then one frame per
+ * generation: `u32 payloadLen | u32 crc32c | msgpack payload`. Payload is
+ * POSITIONAL: `[generation, timestamp, delta, records[], flags]` with
+ * records `[kindByte, id, record]`. `flags` reserves encoding evolution
+ * (bit 0 = compressed payload — v1 always 0; a future writer upgrade,
+ * never a format break). Sealed segments are IMMUTABLE — the fact log's
+ * own law, generalized.
+ * - `seg-.idx` — DERIVED sidecar (msgpack): per-generation frame
+ * offsets (point reads = one ranged read, never a listing) + per-id
+ * generation postings (per-id chain rebuilds read only what they need).
+ * Corrupt/missing → rebuilt from its segment in one sequential read,
+ * loudly.
+ * - `manifest.json` — the segment catalogue + `compactedBelow` (D3's
+ * horizon marker). Cold-open reads THIS; the packed backlog is never
+ * listed.
+ *
+ * D3 semantics carried here: bounded-retention reclaim drops WHOLE segments
+ * at boundaries (O(1) per segment, no rewrite); under the archival profile
+ * (`retention: 'all'`) nothing here is ever dropped — folding is the only
+ * transform (re-representation, never deletion).
+ */
+
+import { encode as msgpackEncode, decode as msgpackDecode } from '@msgpack/msgpack'
+import { crc32c } from '../utils/crc32c.js'
+import type { FactLogStorage } from './factLog.js'
+import { prodLog } from '../utils/logger.js'
+
+/** Directory for segment files + manifest, under the generations prefix. */
+export const SEGMENTS_PREFIX = '_generations/segments'
+
+/** Target sealed-segment size (co-freeze proposal; tunable on evidence). */
+export const SEGMENT_TARGET_BYTES = 64 * 1024 * 1024
+
+const MAGIC = new TextEncoder().encode('BGS1')
+const FRAME_PREFIX_BYTES = 8 // u32 payloadLen + u32 crc32c
+const MANIFEST_PATH = `${SEGMENTS_PREFIX}/manifest.json`
+
+/** One generation's fold input — exactly what the live tier holds for it. */
+export interface FoldGeneration {
+ generation: number
+ timestamp: number
+ /** The tx.json delta object, carried verbatim. */
+ delta: unknown
+ /** The before-image record-set (empty for record-less generations). */
+ records: Array<{ kind: 'noun' | 'verb'; id: string; record: unknown }>
+}
+
+/** Manifest entry for one sealed segment. */
+export interface SegmentMeta {
+ file: string
+ firstGeneration: number
+ lastGeneration: number
+ frames: number
+ bytes: number
+ /** crc32c of the full segment byte stream — the digest chain's link. */
+ checksum: number
+}
+
+interface SegmentManifest {
+ version: 1
+ compactedBelow: number
+ segments: SegmentMeta[]
+}
+
+interface SidecarIndex {
+ version: 1
+ /** [generation, frameOffset, frameLen] ascending by generation. */
+ generations: Array<[number, number, number]>
+ /** `${kindByte}:${id}` → ascending generations holding a record for it. */
+ ids: Record
+}
+
+const segmentFileName = (firstGeneration: number): string =>
+ `seg-${String(firstGeneration).padStart(20, '0')}.bgs`
+const sidecarFileName = (firstGeneration: number): string =>
+ `seg-${String(firstGeneration).padStart(20, '0')}.idx`
+
+/**
+ * The generation-segment store. Owns the packed tier ONLY — the live
+ * per-generation tier and the routing between tiers belong to
+ * `GenerationStore`. All mutating entry points here are called under the
+ * generation store's commit mutex.
+ */
+export class GenerationSegmentStore {
+ private readonly storage: FactLogStorage
+ private manifest: SegmentManifest = { version: 1, compactedBelow: 0, segments: [] }
+ /** Sidecar cache — segments are immutable, so entries never invalidate. */
+ private readonly sidecars = new Map()
+
+ constructor(storage: FactLogStorage) {
+ this.storage = storage
+ }
+
+ /** Load the manifest (ONE read — never a directory listing). */
+ async open(): Promise {
+ const raw = (await this.storage.readRawObject(MANIFEST_PATH)) as SegmentManifest | null
+ if (raw) {
+ if (raw.version !== 1) {
+ throw new Error(
+ `[GenerationSegments] manifest version ${String(raw.version)} is newer than this ` +
+ `engine understands — refusing to serve partial history. Upgrade the engine.`
+ )
+ }
+ this.manifest = raw
+ }
+ }
+
+ /** The packed tier's catalogue (ascending, immutable snapshot). */
+ segments(): readonly SegmentMeta[] {
+ return this.manifest.segments
+ }
+
+ /** D3's horizon marker: generations below this were reclaimed (bounded profiles only). */
+ compactedBelow(): number {
+ return this.manifest.compactedBelow
+ }
+
+ /** The covering sealed segment for `gen`, or null if it lives outside the packed tier. */
+ private coveringSegment(gen: number): SegmentMeta | null {
+ // Manifest is ascending and ranges never overlap — binary search.
+ const segs = this.manifest.segments
+ let lo = 0
+ let hi = segs.length - 1
+ while (lo <= hi) {
+ const mid = (lo + hi) >> 1
+ const s = segs[mid]
+ if (gen < s.firstGeneration) hi = mid - 1
+ else if (gen > s.lastGeneration) lo = mid + 1
+ else return s
+ }
+ return null
+ }
+
+ /** True when `gen` is packed (readable from this tier). */
+ hasGeneration(gen: number): boolean {
+ return this.coveringSegment(gen) !== null
+ }
+
+ /**
+ * Fold consecutive generations into ONE new sealed segment + sidecar and
+ * append it to the manifest atomically. Caller guarantees: `gens` is
+ * ascending, contiguous with the packed tier (first = last packed + 1 when
+ * segments exist), and already durable in the live tier. Crash between the
+ * segment write and the caller's live-tier delete leaves a DUPLICATE
+ * representation — resolved live-tier-wins by the reader; never a gap.
+ */
+ async fold(gens: FoldGeneration[]): Promise {
+ if (gens.length === 0) {
+ throw new Error('[GenerationSegments] fold() requires at least one generation')
+ }
+ for (let i = 1; i < gens.length; i++) {
+ if (gens[i].generation <= gens[i - 1].generation) {
+ throw new Error('[GenerationSegments] fold() input must be strictly ascending')
+ }
+ }
+ const last = this.manifest.segments[this.manifest.segments.length - 1]
+ if (last && gens[0].generation <= last.lastGeneration) {
+ throw new Error(
+ `[GenerationSegments] fold() overlaps the packed tier: ${gens[0].generation} ≤ ` +
+ `sealed ${last.lastGeneration} — segments are immutable, never rewritten`
+ )
+ }
+
+ const first = gens[0].generation
+ const file = segmentFileName(first)
+ const sidecar: SidecarIndex = { version: 1, generations: [], ids: {} }
+
+ // Encode all frames, tracking offsets for the sidecar.
+ const parts: Uint8Array[] = [MAGIC]
+ let offset = MAGIC.length
+ for (const g of gens) {
+ const payload = msgpackEncode([
+ g.generation,
+ g.timestamp,
+ g.delta,
+ g.records.map((r) => [r.kind === 'noun' ? 0 : 1, r.id, r.record]),
+ 0 // flags: v1 = uncompressed
+ ])
+ const frame = new Uint8Array(FRAME_PREFIX_BYTES + payload.length)
+ const view = new DataView(frame.buffer)
+ view.setUint32(0, payload.length, true)
+ view.setUint32(4, crc32c(payload), true)
+ frame.set(payload, FRAME_PREFIX_BYTES)
+ sidecar.generations.push([g.generation, offset, frame.length])
+ for (const r of g.records) {
+ const key = `${r.kind === 'noun' ? 0 : 1}:${r.id}`
+ ;(sidecar.ids[key] ??= []).push(g.generation)
+ }
+ parts.push(frame)
+ offset += frame.length
+ }
+ const total = parts.reduce((n, p) => n + p.length, 0)
+ const bytes = new Uint8Array(total)
+ let at = 0
+ for (const p of parts) {
+ bytes.set(p, at)
+ at += p.length
+ }
+
+ const meta: SegmentMeta = {
+ file,
+ firstGeneration: first,
+ lastGeneration: gens[gens.length - 1].generation,
+ frames: gens.length,
+ bytes: total,
+ checksum: crc32c(bytes)
+ }
+
+ // Durability order: segment + sidecar fsync'd BEFORE the manifest names
+ // them (a crash before the manifest = invisible orphan files, harmless);
+ // manifest last, atomically.
+ const segPath = `${SEGMENTS_PREFIX}/${file}`
+ const idxPath = `${SEGMENTS_PREFIX}/${sidecarFileName(first)}`
+ await this.storage.writeRawBytes(segPath, bytes)
+ await this.storage.writeRawBytes(idxPath, msgpackEncode(sidecar))
+ await this.storage.syncRawObjects([segPath, idxPath])
+ const next: SegmentManifest = {
+ ...this.manifest,
+ segments: [...this.manifest.segments, meta]
+ }
+ await this.storage.writeRawObject(MANIFEST_PATH, next)
+ await this.storage.syncRawObjects([MANIFEST_PATH])
+ this.manifest = next
+ this.sidecars.set(file, sidecar)
+ return meta
+ }
+
+ /** Load (or rebuild, loudly) a segment's sidecar. */
+ private async sidecarFor(meta: SegmentMeta): Promise {
+ const cached = this.sidecars.get(meta.file)
+ if (cached) return cached
+ const idxPath = `${SEGMENTS_PREFIX}/${sidecarFileName(meta.firstGeneration)}`
+ const raw = await this.storage.readRawBytes(idxPath)
+ if (raw) {
+ try {
+ const idx = msgpackDecode(raw) as SidecarIndex
+ if (idx.version === 1) {
+ this.sidecars.set(meta.file, idx)
+ return idx
+ }
+ } catch {
+ // fall through to rebuild
+ }
+ }
+ // Sidecars are DERIVED: rebuild from the segment, loudly — never serve
+ // wrong offsets silently.
+ prodLog.warn(
+ `[GenerationSegments] sidecar for ${meta.file} missing or unreadable — rebuilding from the segment`
+ )
+ const rebuilt = await this.rebuildSidecar(meta)
+ await this.storage.writeRawBytes(idxPath, msgpackEncode(rebuilt))
+ this.sidecars.set(meta.file, rebuilt)
+ return rebuilt
+ }
+
+ /** One sequential read of the segment → a fresh sidecar. Verifies every frame CRC. */
+ private async rebuildSidecar(meta: SegmentMeta): Promise {
+ const frames = await this.readAllFrames(meta)
+ const idx: SidecarIndex = { version: 1, generations: [], ids: {} }
+ for (const f of frames) {
+ idx.generations.push([f.generation, f.offset, f.frameLen])
+ for (const r of f.records) {
+ const key = `${r.kind === 'noun' ? 0 : 1}:${r.id}`
+ ;(idx.ids[key] ??= []).push(f.generation)
+ }
+ }
+ return idx
+ }
+
+ private decodeFrame(
+ payload: Uint8Array
+ ): { generation: number; timestamp: number; delta: unknown; records: FoldGeneration['records'] } {
+ const [generation, timestamp, delta, rawRecords] = msgpackDecode(payload) as [
+ number,
+ number,
+ unknown,
+ Array<[number, string, unknown]>,
+ number
+ ]
+ return {
+ generation,
+ timestamp,
+ delta,
+ records: rawRecords.map(([kindByte, id, record]) => ({
+ kind: kindByte === 0 ? ('noun' as const) : ('verb' as const),
+ id,
+ record
+ }))
+ }
+ }
+
+ private async readAllFrames(meta: SegmentMeta): Promise<
+ Array & { offset: number; frameLen: number }>
+ > {
+ const bytes = await this.storage.readRawBytes(`${SEGMENTS_PREFIX}/${meta.file}`)
+ if (!bytes) {
+ throw new Error(
+ `[GenerationSegments] sealed segment ${meta.file} is MISSING — packed history is damaged; ` +
+ `refusing to continue silently`
+ )
+ }
+ const out: Array & { offset: number; frameLen: number }> = []
+ let at = MAGIC.length
+ const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength)
+ while (at + FRAME_PREFIX_BYTES <= bytes.length) {
+ const payloadLen = view.getUint32(at, true)
+ const crc = view.getUint32(at + 4, true)
+ const payload = bytes.subarray(at + FRAME_PREFIX_BYTES, at + FRAME_PREFIX_BYTES + payloadLen)
+ if (payload.length !== payloadLen || crc32c(payload) !== crc) {
+ throw new Error(
+ `[GenerationSegments] frame CRC mismatch in ${meta.file} at offset ${at} — ` +
+ `packed history is damaged; refusing to serve it`
+ )
+ }
+ out.push({ ...this.decodeFrame(payload), offset: at, frameLen: FRAME_PREFIX_BYTES + payloadLen })
+ at += FRAME_PREFIX_BYTES + payloadLen
+ }
+ return out
+ }
+
+ /** Read one packed generation's frame via its sidecar offset (one ranged read). */
+ private async readFrame(
+ gen: number
+ ): Promise | null> {
+ const meta = this.coveringSegment(gen)
+ if (!meta) return null
+ const idx = await this.sidecarFor(meta)
+ // generations ascending → binary search.
+ const gens = idx.generations
+ let lo = 0
+ let hi = gens.length - 1
+ while (lo <= hi) {
+ const mid = (lo + hi) >> 1
+ if (gens[mid][0] < gen) lo = mid + 1
+ else if (gens[mid][0] > gen) hi = mid - 1
+ else {
+ const [, offset, frameLen] = gens[mid]
+ const bytes = await this.storage.readRawBytes(`${SEGMENTS_PREFIX}/${meta.file}`)
+ if (!bytes) {
+ throw new Error(`[GenerationSegments] sealed segment ${meta.file} is MISSING`)
+ }
+ const frame = bytes.subarray(offset, offset + frameLen)
+ const view = new DataView(frame.buffer, frame.byteOffset, frame.byteLength)
+ const payloadLen = view.getUint32(0, true)
+ const crc = view.getUint32(4, true)
+ const payload = frame.subarray(FRAME_PREFIX_BYTES, FRAME_PREFIX_BYTES + payloadLen)
+ if (payload.length !== payloadLen || crc32c(payload) !== crc) {
+ throw new Error(
+ `[GenerationSegments] frame CRC mismatch for generation ${gen} in ${meta.file} — ` +
+ `packed history is damaged; refusing to serve it`
+ )
+ }
+ return this.decodeFrame(payload)
+ }
+ }
+ // In the covering range but not present: the packed tier is dense by
+ // construction (fold packs every generation it is handed, including
+ // record-less ones) — absence inside a sealed range is damage.
+ throw new Error(
+ `[GenerationSegments] generation ${gen} is inside sealed segment ${meta.file}'s declared ` +
+ `range but has no frame — packed history is damaged`
+ )
+ }
+
+ /** The packed tier's delta for `gen` (null = not packed). */
+ async readDelta(gen: number): Promise<{ delta: unknown; timestamp: number } | null> {
+ const frame = await this.readFrame(gen)
+ return frame ? { delta: frame.delta, timestamp: frame.timestamp } : null
+ }
+
+ /** The packed tier's full record-set for `gen` (null = not packed). */
+ async readRecords(gen: number): Promise {
+ const frame = await this.readFrame(gen)
+ return frame ? frame.records : null
+ }
+
+ /** One packed before-image (null = not packed OR no record for the id in that generation). */
+ async readRecord(gen: number, kind: 'noun' | 'verb', id: string): Promise {
+ const frame = await this.readFrame(gen)
+ if (!frame) return null
+ const hit = frame.records.find((r) => r.kind === kind && r.id === id)
+ return hit ? hit.record : null
+ }
+
+ /**
+ * D3 reclaim: drop WHOLE segments whose lastGeneration < `belowGeneration`
+ * and bump `compactedBelow`. Partial segments are never dropped — the
+ * boundary waits. NEVER called under the archival profile (the caller
+ * enforces retention semantics; this method only executes boundary drops).
+ */
+ async dropSegmentsBelow(belowGeneration: number): Promise<{ dropped: number; compactedBelow: number }> {
+ const keep: SegmentMeta[] = []
+ const drop: SegmentMeta[] = []
+ for (const s of this.manifest.segments) {
+ ;(s.lastGeneration < belowGeneration ? drop : keep).push(s)
+ }
+ if (drop.length === 0) {
+ return { dropped: 0, compactedBelow: this.manifest.compactedBelow }
+ }
+ const compactedBelow = Math.max(
+ this.manifest.compactedBelow,
+ drop[drop.length - 1].lastGeneration + 1
+ )
+ // Manifest first (the drop is authoritative once named), then bytes —
+ // a crash between leaves orphan segment files invisible to the manifest,
+ // harmless and re-collectable.
+ const next: SegmentManifest = { ...this.manifest, compactedBelow, segments: keep }
+ await this.storage.writeRawObject(MANIFEST_PATH, next)
+ await this.storage.syncRawObjects([MANIFEST_PATH])
+ this.manifest = next
+ for (const s of drop) {
+ await this.storage.deleteRawObject(`${SEGMENTS_PREFIX}/${s.file}`)
+ await this.storage.deleteRawObject(`${SEGMENTS_PREFIX}/${sidecarFileName(s.firstGeneration)}`)
+ this.sidecars.delete(s.file)
+ }
+ return { dropped: drop.length, compactedBelow }
+ }
+
+ /**
+ * D8 rider — the packed portion of `generationDigest(g)`: a deterministic
+ * crc32c chain over sealed-segment checksums fully below `g`, plus the
+ * frame CRC of `g`'s own frame when `g` is mid-segment. O(segments), not
+ * O(generations); identical history ⇒ identical digest on any machine.
+ * The live-tier portion is composed by the caller.
+ */
+ async digestThroughPacked(g: number): Promise {
+ let digest = 0
+ let covered = false
+ for (const s of this.manifest.segments) {
+ if (s.lastGeneration <= g) {
+ digest = crc32c(new TextEncoder().encode(`${digest}:${s.checksum}`))
+ if (s.lastGeneration === g) covered = true
+ } else if (s.firstGeneration <= g) {
+ // g is mid-segment: chain the partial prefix via g's frame CRC.
+ const frame = await this.readFrame(g)
+ if (frame === null) return null
+ const idx = await this.sidecarFor(s)
+ const upTo = idx.generations.filter(([gen]) => gen <= g)
+ for (const [gen, offset, frameLen] of upTo) {
+ digest = crc32c(new TextEncoder().encode(`${digest}:${gen}:${offset}:${frameLen}`))
+ }
+ covered = true
+ break
+ }
+ }
+ return covered || this.manifest.segments.length > 0 ? digest : null
+ }
+}
diff --git a/src/db/generationStore.ts b/src/db/generationStore.ts
index 4e3738d9..aede17a4 100644
--- a/src/db/generationStore.ts
+++ b/src/db/generationStore.ts
@@ -46,6 +46,8 @@ import type {
TxLogEntry
} from './types.js'
import { FactLog, storageSupportsFactLog, type CommitFact, type FactOp } from './factLog.js'
+import { GenerationSegmentStore, type FoldGeneration } from './generationSegments.js'
+import { crc32c } from '../utils/crc32c.js'
/**
* The byte-identical before-images of every id a commit touches, read UNDER
@@ -266,6 +268,21 @@ export class GenerationStore {
*/
private historyBytesTotal: number | null = null
+ /**
+ * The packed tier (D1+D3): sealed segments holding folded cold
+ * generations. Null until {@link open} wires it (and on storage adapters
+ * without raw-byte primitives — the live tier then carries everything,
+ * exactly as before the packed tier existed).
+ */
+ private segments: GenerationSegmentStore | null = null
+
+ /**
+ * Live-tier window: generations newer than `committed - REPACK_LIVE_WINDOW`
+ * are never folded — the hot tail stays in the per-generation layout the
+ * write path owns. Matches the resident chain window's scale.
+ */
+ static readonly REPACK_LIVE_WINDOW = 1024
+
/**
* Model-B per-write group-commit — the in-memory PENDING tier.
*
@@ -433,6 +450,33 @@ export class GenerationStore {
this.factLog = null
}
+ // PACKED TIER (D1+D3): same capability gate as the fact log. Opening
+ // reads ONE manifest — never a listing of the packed backlog — and seeds
+ // committedRanges with the sealed ranges so packed generations resolve
+ // exactly like live ones.
+ if (storageSupportsFactLog(this.storage)) {
+ this.segments = new GenerationSegmentStore(this.storage)
+ await this.segments.open()
+ const packedRanges = this.segments
+ .segments()
+ .map((s): [number, number] => [s.firstGeneration, Math.min(s.lastGeneration, this.committed)])
+ .filter(([lo, hi]) => lo <= hi)
+ if (packedRanges.length > 0) {
+ // Merge packed (older) + live (newer) interval sets — both ascending;
+ // coalesce adjacency so range arithmetic stays interval-exact.
+ const merged: Array<[number, number]> = []
+ for (const r of [...packedRanges, ...this.committedRanges].sort((a, b) => a[0] - b[0])) {
+ const last = merged[merged.length - 1]
+ if (last && r[0] <= last[1] + 1) last[1] = Math.max(last[1], r[1])
+ else merged.push([r[0], r[1]])
+ }
+ this.committedRanges = merged
+ }
+ this.horizonGen = Math.max(this.horizonGen, this.segments.compactedBelow() - 1)
+ } else {
+ this.segments = null
+ }
+
// Hook single-op write batches so generation() is always meaningful.
// Suppressed while a transact batch executes (the batch is ONE generation).
if (!options?.readOnly) {
@@ -500,6 +544,51 @@ export class GenerationStore {
* deltas (cache-bounded reads).
* @returns Counts, bytes, generation range, and the compaction horizon.
*/
+ /**
+ * @description D8 (gate-to-generation provenance): a deterministic content
+ * digest of the generation log THROUGH `g` — identical history ⇒ identical
+ * digest on any machine; any divergence (different records, different
+ * order, reclaimed range) ⇒ different digest. Composed from the packed
+ * tier's sealed-segment checksum chain (O(segments)) plus the live tier's
+ * per-generation delta digests (O(live window at most)). Release gates pin
+ * {generation, digest} and verify both at execution time.
+ * @param g - The generation to digest through (≤ committed).
+ * @returns A hex digest string, stable across reopen and repacking states
+ * ONLY for fully-packed prefixes — repacking changes representation, so
+ * the composed digest is defined over CONTENT: live-tier gens hash their
+ * delta + record ids, packed gens hash via frame CRCs. A gate should pin
+ * after a repack pass for long-term stability, or re-pin on repack.
+ */
+ async generationDigest(g: number): Promise {
+ if (!Number.isInteger(g) || g < 1 || g > this.committed) {
+ throw new RangeError(
+ `generationDigest(): generation ${g} is out of range [1, ${this.committed}]`
+ )
+ }
+ if (g <= this.horizonGen) {
+ throw new GenerationCompactedError(g, this.horizonGen)
+ }
+ let digest = 0
+ const enc = new TextEncoder()
+ if (this.segments) {
+ const packed = await this.segments.digestThroughPacked(g)
+ if (packed !== null) digest = packed
+ }
+ // Live-tier composition: every committed gen ≤ g not covered by a sealed
+ // segment hashes its delta content in ascending order.
+ for (const gen of this.committedGensAsc()) {
+ if (gen > g) break
+ if (this.segments?.hasGeneration(gen)) continue
+ const delta = await this.getDelta(gen)
+ digest = crc32c(
+ enc.encode(
+ `${digest}:${gen}:${delta.timestamp}:${[...delta.nouns].sort().join(',')}:${[...delta.verbs].sort().join(',')}`
+ )
+ )
+ }
+ return digest.toString(16).padStart(8, '0')
+ }
+
async historyStats(): Promise<{
generations: number
bytes: number
@@ -538,14 +627,17 @@ export class GenerationStore {
try {
paths = await this.storage.listRawObjects(`${GENERATIONS_PREFIX}/${gen}/prev`)
} catch {
- return []
+ paths = []
}
const records: GenerationRecord[] = []
for (const p of paths) {
const record = (await this.storage.readRawObject(p)) as GenerationRecord | null
if (record) records.push(record)
}
- return records
+ if (records.length > 0) return records
+ // Two-tier: folded generations serve their record-set from the segment.
+ const packed = await this.segments?.readRecords(gen)
+ return packed ? (packed.map((r) => r.record) as GenerationRecord[]) : []
}
/**
@@ -1783,9 +1875,15 @@ export class GenerationStore {
if (pending) {
return (kind === 'noun' ? pending.nouns : pending.verbs).get(id) ?? null
}
- return (await this.storage.readRawObject(
+ const live = (await this.storage.readRawObject(
`${GENERATIONS_PREFIX}/${gen}/prev/${id}.json`
)) as GenerationRecord | null
+ if (live) return live
+ // Two-tier: the packed tier serves folded generations (live-tier-wins).
+ if (this.segments?.hasGeneration(gen)) {
+ return (await this.segments.readRecord(gen, kind, id)) as GenerationRecord | null
+ }
+ return null
}
/**
@@ -2132,6 +2230,21 @@ export class GenerationStore {
`${GENERATIONS_PREFIX}/${gen}/tx.json`
)) as GenerationDelta | null
if (delta === null) {
+ // Two-tier read (D1+D3): not in the live tier → the packed tier.
+ // Live-tier-wins ordering (a crash mid-fold leaves a duplicate, never
+ // a gap), so the segment lookup runs only after the live miss.
+ const packed = await this.segments?.readDelta(gen)
+ if (packed) {
+ const d = packed.delta as GenerationDelta
+ const entry = {
+ nouns: new Set(d.nouns),
+ verbs: new Set(d.verbs),
+ timestamp: packed.timestamp,
+ bytes: d.bytes ?? 0
+ }
+ this.setDelta(gen, entry)
+ return entry
+ }
throw new Error(
`Generation delta missing: ${GENERATIONS_PREFIX}/${gen}/tx.json ` +
`(store corrupted or records removed outside compactHistory())`
@@ -2213,6 +2326,94 @@ export class GenerationStore {
* @param options - Retention caps (see {@link CompactHistoryOptions}).
* @returns Count of removed record-sets and the new horizon.
*/
+ /**
+ * @description The REPACKER (D1+D3+repacking): fold cold live-tier
+ * generations into sealed segments — re-representation, never deletion.
+ * Every record and delta stays readable (asOf/chains unchanged); the
+ * per-generation directories are deleted only AFTER their segment is
+ * durable (crash between = duplicate representation, resolved
+ * live-tier-wins by every reader; never a gap). This is the transform that
+ * takes a 70k-file history to tens of segment files, and the ONLY history
+ * transform permitted under the archival profile.
+ *
+ * Folds oldest-first, contiguous from the packed boundary, in batches, and
+ * stops at the live window ({@link GenerationStore.REPACK_LIVE_WINDOW})
+ * or when `timeBudgetMs` is spent — an early stop is a consistent prefix;
+ * the next pass resumes.
+ */
+ async repackHistory(options?: { timeBudgetMs?: number; batchGenerations?: number }): Promise<{
+ foldedGenerations: number
+ segmentsCreated: number
+ }> {
+ if (!this.segments) return { foldedGenerations: 0, segmentsCreated: 0 }
+ const segments = this.segments
+ return this.withMutex(async () => {
+ const deadline =
+ options?.timeBudgetMs !== undefined ? Date.now() + options.timeBudgetMs : undefined
+ const batchSize = options?.batchGenerations ?? 512
+ const coldCeiling = this.committed - GenerationStore.REPACK_LIVE_WINDOW
+ const packedThrough =
+ segments.segments().length > 0
+ ? segments.segments()[segments.segments().length - 1].lastGeneration
+ : 0
+
+ // Cold, unpacked, committed generations — ascending, contiguous scan.
+ const eligible: number[] = []
+ for (const gen of this.committedGensAsc()) {
+ if (gen > coldCeiling) break
+ if (gen <= packedThrough) continue // already packed (dup fold barred)
+ if (this.pendingBuffer.has(gen)) continue // un-flushed = live by definition
+ eligible.push(gen)
+ }
+
+ let folded = 0
+ let segmentsCreated = 0
+ for (let i = 0; i < eligible.length; i += batchSize) {
+ if (deadline !== undefined && Date.now() >= deadline) break
+ const batch = eligible.slice(i, i + batchSize)
+ const foldInput: FoldGeneration[] = []
+ for (const gen of batch) {
+ const delta = (await this.storage.readRawObject(
+ `${GENERATIONS_PREFIX}/${gen}/tx.json`
+ )) as GenerationDelta | null
+ if (delta === null) {
+ // Already folded by a prior crashed pass whose dirs were removed,
+ // or damage — getDelta's two-tier read decides which, loudly,
+ // when someone asks. Skip; never fold a generation we cannot read.
+ continue
+ }
+ const records: FoldGeneration['records'] = []
+ for (const [kind, ids] of [
+ ['noun', delta.nouns] as const,
+ ['verb', delta.verbs] as const
+ ]) {
+ for (const id of ids) {
+ const record = await this.storage.readRawObject(
+ `${GENERATIONS_PREFIX}/${gen}/prev/${id}.json`
+ )
+ if (record) records.push({ kind, id, record })
+ }
+ }
+ foldInput.push({ generation: gen, timestamp: delta.timestamp, delta, records })
+ }
+ if (foldInput.length === 0) continue
+ await segments.fold(foldInput)
+ segmentsCreated++
+ // Segment + manifest durable → the live copies retire.
+ for (const g of foldInput) {
+ await this.storage.removeRawPrefix(`${GENERATIONS_PREFIX}/${g.generation}`)
+ }
+ folded += foldInput.length
+ }
+ if (folded > 0) {
+ prodLog.info(
+ `[GenerationStore] repacked ${folded} cold generation(s) into ${segmentsCreated} segment(s) — history preserved, file count reduced`
+ )
+ }
+ return { foldedGenerations: folded, segmentsCreated }
+ })
+ }
+
async compact(options?: CompactHistoryOptions): Promise {
return this.withMutex(async () => {
const minPinned = this.minPinnedGeneration()
@@ -2304,6 +2505,16 @@ export class GenerationStore {
// Reclaimed generations leave the per-id chains stale → rebuild on next read.
this.invalidateChains()
this.horizonGen = Math.max(this.horizonGen, highestRemoved)
+ // Packed-tier reclaim (D3): a packed generation's bytes live in a
+ // sealed segment — removeRawPrefix above was a no-op for it. Drop
+ // WHOLE segments now fully below the horizon; a partially-reclaimed
+ // segment keeps its bytes until the boundary passes it (the frozen
+ // partial-segments-wait rule; logical reclamation above still holds —
+ // the generations left committedRanges and asOf below the horizon
+ // throws regardless).
+ if (this.segments) {
+ await this.segments.dropSegmentsBelow(this.horizonGen + 1)
+ }
const manifest: GenerationManifest = {
version: 1,
generation: this.committed,
diff --git a/src/index.ts b/src/index.ts
index 00adc191..bfab39da 100644
--- a/src/index.ts
+++ b/src/index.ts
@@ -216,6 +216,7 @@ export type {
CommitFact,
FactOp,
FactScanBatch,
+ SCANFACTS_FIRST_BATCH_MS,
FactScanHandle
} from './db/factLog.js'
// The generalized family stamp — which source generation a projection
diff --git a/src/storage/adapters/memoryStorage.ts b/src/storage/adapters/memoryStorage.ts
index bab9d4d9..1b1f412e 100644
--- a/src/storage/adapters/memoryStorage.ts
+++ b/src/storage/adapters/memoryStorage.ts
@@ -133,6 +133,11 @@ export class MemoryStorage extends BaseStorage {
*/
protected async deleteObjectFromPath(path: string): Promise {
this.objectStore.delete(path)
+ // Filesystem parity: on disk, objects and raw BYTE files are both just
+ // files — unlink removes whichever exists. Without this, deleteRawObject
+ // on a raw-bytes path (fact-log/generation segments) silently no-ops on
+ // memory storage: the delete "succeeds" and the bytes remain.
+ this.rawBytesStore.delete(path)
}
/**
diff --git a/tests/integration/history-repacking.test.ts b/tests/integration/history-repacking.test.ts
new file mode 100644
index 00000000..2bcee038
--- /dev/null
+++ b/tests/integration/history-repacking.test.ts
@@ -0,0 +1,186 @@
+/**
+ * @module tests/integration/history-repacking
+ * @description The D1+D3 two-tier history lifecycle end-to-end on a real
+ * brain. Laws: (1) repacking is RE-REPRESENTATION — after folding, every
+ * asOf() read below the fold boundary answers exactly as before, across a
+ * cold reopen; (2) folded per-generation directories are physically gone
+ * (the file-count cure is real, not cosmetic); (3) repack + reclaim compose:
+ * bounded retention after repacking drops whole segments and asOf below the
+ * horizon throws GenerationCompactedError; (4) repackHistory is explicit
+ * API and time-bounded (spent budget = consistent no-op).
+ *
+ * Uses a tiny REPACK_LIVE_WINDOW override so a small history has a cold
+ * tier at all (the production window is 1024).
+ */
+import { describe, it, expect, afterEach } from 'vitest'
+import * as fs from 'node:fs'
+import * as path from 'node:path'
+import * as os from 'node:os'
+import { Brainy } from '../../src/brainy.js'
+import { NounType } from '../../src/types/graphTypes.js'
+import { GenerationStore } from '../../src/db/generationStore.js'
+import { GenerationCompactedError } from '../../src/db/errors.js'
+import { SEGMENTS_PREFIX } from '../../src/db/generationSegments.js'
+
+const stub = async (text: string): Promise => {
+ const h = text.split('').reduce((a, c) => a + c.charCodeAt(0), 0)
+ return new Array(384).fill(0).map((_, i) => Math.sin(h + i))
+}
+
+const openBrain = async (dir: string): Promise => {
+ const brain = new Brainy({
+ requireSubtype: false,
+ storage: { type: 'filesystem', path: dir },
+ embeddingFunction: stub
+ })
+ await brain.init()
+ return brain
+}
+
+describe('history repacking — the two-tier lifecycle', () => {
+ const dirs: string[] = []
+ const tempDir = (): string => {
+ const d = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-repack-'))
+ dirs.push(d)
+ return d
+ }
+ const originalWindow = GenerationStore.REPACK_LIVE_WINDOW
+
+ afterEach(() => {
+ ;(GenerationStore as any).REPACK_LIVE_WINDOW = originalWindow
+ for (const d of dirs.splice(0)) {
+ try {
+ fs.rmSync(d, { recursive: true, force: true })
+ } catch {
+ /* best effort */
+ }
+ }
+ })
+
+ it('repack preserves every historical read across cold reopen; folded dirs are gone', async () => {
+ ;(GenerationStore as any).REPACK_LIVE_WINDOW = 3
+ const dir = tempDir()
+ const brain = await openBrain(dir)
+
+ const id = await brain.add({
+ data: 'versioned-entity',
+ type: NounType.Document,
+ metadata: { v: 0 }
+ })
+ for (let v = 1; v <= 10; v++) await brain.update({ id, metadata: { v } })
+ await brain.flush()
+
+ // Ground truth BEFORE repacking: capture asOf views for early generations.
+ const before: Record = {}
+ for (const g of [2, 4, 6]) {
+ const db = await brain.asOf(g)
+ before[g] = (await db.get(id))?.metadata?.v as number
+ await db.release()
+ }
+
+ const result = await brain.repackHistory()
+ expect(result.foldedGenerations).toBeGreaterThan(0)
+ expect(result.segmentsCreated).toBeGreaterThan(0)
+
+ // The folded per-generation directories are PHYSICALLY gone…
+ const genDirs = fs
+ .readdirSync(path.join(dir, '_generations'), { withFileTypes: true })
+ .filter((e) => e.isDirectory() && /^\d+$/.test(e.name)).length
+ expect(genDirs).toBeLessThanOrEqual(4) // live window (3) + at most the newest
+ // …and the segment tier exists (the filesystem adapter stores objects
+ // gzipped, so the manifest may live at either spelling).
+ const segDir = path.join(dir, SEGMENTS_PREFIX)
+ expect(
+ fs.existsSync(path.join(segDir, 'manifest.json')) ||
+ fs.existsSync(path.join(segDir, 'manifest.json.gz'))
+ ).toBe(true)
+ expect(fs.readdirSync(segDir).some((f) => f.endsWith('.bgs'))).toBe(true)
+
+ // Same asOf answers from the packed tier, same process…
+ for (const g of [2, 4, 6]) {
+ const db = await brain.asOf(g)
+ expect((await db.get(id))?.metadata?.v).toBe(before[g])
+ await db.release()
+ }
+ await brain.close()
+
+ // …and across a COLD REOPEN (manifest discovery, no live dirs to list).
+ const reopened = await openBrain(dir)
+ for (const g of [2, 4, 6]) {
+ const db = await reopened.asOf(g)
+ expect((await db.get(id))?.metadata?.v).toBe(before[g])
+ await db.release()
+ }
+ expect((await reopened.get(id))?.metadata?.v).toBe(10) // live state untouched
+ await reopened.close()
+ })
+
+ it('repack + bounded reclaim compose: whole segments drop, horizon is loud', async () => {
+ ;(GenerationStore as any).REPACK_LIVE_WINDOW = 2
+ const dir = tempDir()
+ const brain = await openBrain(dir)
+ const id = await brain.add({ data: 'reclaim-probe', type: NounType.Document, metadata: { v: 0 } })
+ for (let v = 1; v <= 8; v++) await brain.update({ id, metadata: { v } })
+ await brain.flush()
+ await brain.repackHistory()
+
+ // Reclaim down to the 3 newest generations — packed segments below the
+ // horizon drop whole; asOf below throws loudly.
+ const res = await brain.compactHistory({ maxGenerations: 3 })
+ expect(res.removedGenerations).toBeGreaterThan(0)
+ await expect(brain.asOf(1)).rejects.toBeInstanceOf(GenerationCompactedError)
+ expect((await brain.get(id))?.metadata?.v).toBe(8)
+ await brain.close()
+ })
+
+ it('generationDigest: reopen-stable, divergence-sensitive, loud below the horizon', async () => {
+ ;(GenerationStore as any).REPACK_LIVE_WINDOW = 2
+ const dir = tempDir()
+ const brain = await openBrain(dir)
+ const id = await brain.add({ data: 'digest-probe', type: NounType.Document, metadata: { v: 0 } })
+ for (let v = 1; v <= 6; v++) await brain.update({ id, metadata: { v } })
+ await brain.flush()
+ await brain.repackHistory()
+
+ const gen = brain.generation()
+ const atHead = await brain.generationDigest(gen)
+ const atMid = await brain.generationDigest(3)
+ expect(atHead).toMatch(/^[0-9a-f]{8}$/)
+ expect(atMid).not.toBe(atHead) // more history ⇒ different digest
+ await brain.close()
+
+ // Reopen-stable: same history, same digests (packed prefix stability).
+ const reopened = await openBrain(dir)
+ expect(await reopened.generationDigest(gen)).toBe(atHead)
+ expect(await reopened.generationDigest(3)).toBe(atMid)
+
+ // New history diverges the head digest.
+ await reopened.update({ id, metadata: { v: 7 } })
+ await reopened.flush()
+ expect(await reopened.generationDigest(reopened.generation())).not.toBe(atHead)
+
+ // Below the horizon: LOUD, never a silent pin of reclaimed history.
+ await reopened.compactHistory({ maxGenerations: 2 })
+ await expect(reopened.generationDigest(1)).rejects.toBeInstanceOf(GenerationCompactedError)
+ await reopened.close()
+ })
+
+ it('a spent time budget is a consistent no-op; the next pass resumes', async () => {
+ ;(GenerationStore as any).REPACK_LIVE_WINDOW = 2
+ const dir = tempDir()
+ const brain = await openBrain(dir)
+ const id = await brain.add({ data: 'budget-probe', type: NounType.Document, metadata: { v: 0 } })
+ for (let v = 1; v <= 6; v++) await brain.update({ id, metadata: { v } })
+ await brain.flush()
+
+ const bounded = await brain.repackHistory({ timeBudgetMs: 0 })
+ expect(bounded).toEqual({ foldedGenerations: 0, segmentsCreated: 0 })
+
+ const resumed = await brain.repackHistory()
+ expect(resumed.foldedGenerations).toBeGreaterThan(0)
+ const db = await brain.asOf(3)
+ expect((await db.get(id))?.metadata?.v).toBeDefined()
+ await db.release()
+ await brain.close()
+ })
+})
diff --git a/tests/unit/db/fact-log.test.ts b/tests/unit/db/fact-log.test.ts
index abce2dc9..f1c226cc 100644
--- a/tests/unit/db/fact-log.test.ts
+++ b/tests/unit/db/fact-log.test.ts
@@ -186,4 +186,52 @@ describe('fact log — round-trip, framing, reconcile, rotation, scan', () => {
await log.sync()
expect(log.segmentPaths()).toEqual([]) // only a tail exists — nothing sealed
})
+
+ describe('scanFacts liveness contract (Stage-2 D1)', () => {
+ it('a wedged store fails LOUDLY within the first-batch bound — never a silent hang', async () => {
+ // Force a sealed segment (tiny rotateBytes) so the scan must READ from
+ // storage, then wedge that read: the exact production shape (a
+ // backlogged brain whose segment read never returned).
+ const mem: any = new MemoryStorage()
+ await mem.init()
+ const wedgeable = new FactLog(mem, { rotateBytes: 1 })
+ await wedgeable.open(0)
+ await wedgeable.append(fact(1))
+ await wedgeable.append(fact(2)) // second append rotates → seg 1 sealed
+ await wedgeable.sync()
+
+ const realRead = mem.readRawBytes.bind(mem)
+ mem.readRawBytes = (p: string) =>
+ p.includes('facts/seg-') ? new Promise(() => {}) : realRead(p) // hangs forever
+
+ const scan = wedgeable.scanFacts({ firstBatchTimeoutMs: 200 })
+ const started = Date.now()
+ await expect(scan.batches().next()).rejects.toThrow(/no first batch within 200ms/)
+ expect(Date.now() - started).toBeLessThan(5_000) // bound held, not a hang
+ })
+
+ it('a healthy scan is unaffected — first batch well inside the bound, all facts delivered', async () => {
+ for (let g = 1; g <= 5; g++) await log.append(fact(g))
+ await log.sync()
+ const scan = log.scanFacts({ batchSize: 2 })
+ const all: CommitFact[] = []
+ for await (const b of scan.batches()) all.push(...b.facts)
+ expect(all.map((f) => f.generation)).toEqual([1, 2, 3, 4, 5])
+ expect(scan.summary().factsYielded).toBe(5)
+ })
+
+ it('consumer think-time between pulls never counts against the producer', async () => {
+ for (let g = 1; g <= 4; g++) await log.append(fact(g))
+ await log.sync()
+ // Bound tighter than the consumer's pause: only the FIRST pull is
+ // raced, so a slow consumer after batch 1 must not trip the deadline.
+ const gen = log.scanFacts({ batchSize: 2, firstBatchTimeoutMs: 150 }).batches()
+ const first = await gen.next()
+ expect(first.done).toBe(false)
+ await new Promise((r) => setTimeout(r, 400)) // dawdle past the bound
+ const second = await gen.next()
+ expect(second.done).toBe(false)
+ expect((await gen.next()).done).toBe(true)
+ })
+ })
})
diff --git a/tests/unit/db/generation-segments.test.ts b/tests/unit/db/generation-segments.test.ts
new file mode 100644
index 00000000..27ab85cb
--- /dev/null
+++ b/tests/unit/db/generation-segments.test.ts
@@ -0,0 +1,150 @@
+/**
+ * @module tests/unit/db/generation-segments
+ * @description The generation-segment store (Stage-2 D1+D3 file format).
+ * Laws: (1) fold → read round-trips deltas and records byte-faithfully via
+ * sidecar point-reads; (2) the manifest is the ONLY discovery path — reopen
+ * reads one file, never a listing; (3) a lost/corrupt sidecar rebuilds from
+ * its segment loudly, a damaged SEGMENT fails loudly (never silent wrong
+ * data); (4) D3 reclaim drops whole segments only and bumps compactedBelow;
+ * (5) the packed digest is deterministic across reopen; (6) immutability —
+ * fold refuses overlap with sealed ranges.
+ */
+import { describe, it, expect, beforeEach } from 'vitest'
+import { MemoryStorage } from '../../../src/storage/adapters/memoryStorage.js'
+import {
+ GenerationSegmentStore,
+ SEGMENTS_PREFIX,
+ type FoldGeneration
+} from '../../../src/db/generationSegments.js'
+
+const UUID = (n: number): string => `00000000-0000-4000-8000-${String(n).padStart(12, '0')}`
+
+const gen = (g: number, recordCount = 2): FoldGeneration => ({
+ generation: g,
+ timestamp: 1_700_000_000_000 + g,
+ delta: { generation: g, nouns: [UUID(g)], verbs: [], bytes: 123 + g },
+ records: Array.from({ length: recordCount }, (_, i) => ({
+ kind: (i % 2 === 0 ? 'noun' : 'verb') as 'noun' | 'verb',
+ id: UUID(g * 100 + i),
+ record: { metadata: { noun: 'document', v: g }, vector: { v: [g, i] } }
+ }))
+})
+
+describe('db/GenerationSegmentStore — the D1+D3 packed tier', () => {
+ let storage: MemoryStorage
+ let store: GenerationSegmentStore
+
+ beforeEach(async () => {
+ storage = new MemoryStorage()
+ await storage.init()
+ store = new GenerationSegmentStore(storage as any)
+ await store.open()
+ })
+
+ it('fold → read round-trips deltas and records via sidecar point-reads', async () => {
+ const meta = await store.fold([gen(1), gen(2), gen(3)])
+ expect(meta).toMatchObject({ firstGeneration: 1, lastGeneration: 3, frames: 3 })
+ expect(meta.checksum).toBeGreaterThan(0)
+
+ expect(store.hasGeneration(2)).toBe(true)
+ expect(store.hasGeneration(4)).toBe(false)
+
+ const d2 = await store.readDelta(2)
+ expect(d2?.delta).toEqual({ generation: 2, nouns: [UUID(2)], verbs: [], bytes: 125 })
+ expect(d2?.timestamp).toBe(1_700_000_000_002)
+
+ const records = await store.readRecords(3)
+ expect(records).toHaveLength(2)
+ expect(records![0]).toEqual({
+ kind: 'noun',
+ id: UUID(300),
+ record: { metadata: { noun: 'document', v: 3 }, vector: { v: [3, 0] } }
+ })
+ // Point read by id, both kinds.
+ expect(await store.readRecord(3, 'verb', UUID(301))).toEqual({
+ metadata: { noun: 'document', v: 3 },
+ vector: { v: [3, 1] }
+ })
+ expect(await store.readRecord(3, 'noun', UUID(999))).toBeNull()
+ })
+
+ it('reopen discovers everything from the manifest alone — no listing', async () => {
+ await store.fold([gen(1), gen(2)])
+ await store.fold([gen(3), gen(4)])
+
+ const reopened = new GenerationSegmentStore(storage as any)
+ await reopened.open()
+ expect(reopened.segments()).toHaveLength(2)
+ expect(reopened.hasGeneration(4)).toBe(true)
+ expect((await reopened.readDelta(1))?.timestamp).toBe(1_700_000_000_001)
+ })
+
+ it('a lost sidecar rebuilds from its segment; a damaged segment fails LOUDLY', async () => {
+ const meta = await store.fold([gen(1), gen(2)])
+ const idxPath = `${SEGMENTS_PREFIX}/seg-${String(1).padStart(20, '0')}.idx`
+ await storage.deleteRawObject(idxPath)
+
+ const reopened = new GenerationSegmentStore(storage as any)
+ await reopened.open()
+ // Rebuild path: still serves correct data.
+ expect((await reopened.readRecords(2))!).toHaveLength(2)
+
+ // Now damage the SEGMENT itself: flip a payload byte → CRC mismatch, loud.
+ const segPath = `${SEGMENTS_PREFIX}/${meta.file}`
+ const bytes = (await storage.readRawBytes(segPath))!
+ bytes[bytes.length - 3] ^= 0xff
+ await storage.writeRawBytes(segPath, bytes)
+ const damaged = new GenerationSegmentStore(storage as any)
+ await damaged.open()
+ ;(damaged as any).sidecars.clear()
+ await storage.deleteRawObject(idxPath) // force the sequential rebuild over damaged bytes
+ await expect(damaged.readRecords(2)).rejects.toThrow(/CRC mismatch|damaged/)
+ })
+
+ it('D3 reclaim drops whole segments only and bumps compactedBelow', async () => {
+ await store.fold([gen(1), gen(2)])
+ await store.fold([gen(3), gen(4)])
+ await store.fold([gen(5), gen(6)])
+
+ // Horizon mid-segment-2 (below 4): only segment 1 is FULLY below → drops.
+ const r1 = await store.dropSegmentsBelow(4)
+ expect(r1).toEqual({ dropped: 1, compactedBelow: 3 })
+ expect(store.hasGeneration(1)).toBe(false)
+ expect(store.hasGeneration(3)).toBe(true) // partial segment survives whole
+
+ // Bytes actually gone.
+ expect(await storage.readRawBytes(`${SEGMENTS_PREFIX}/seg-${String(1).padStart(20, '0')}.bgs`)).toBeNull()
+
+ // Horizon past everything: the rest drop; compactedBelow is durable.
+ const r2 = await store.dropSegmentsBelow(7)
+ expect(r2.dropped).toBe(2)
+ const reopened = new GenerationSegmentStore(storage as any)
+ await reopened.open()
+ expect(reopened.compactedBelow()).toBe(7)
+ expect(reopened.segments()).toHaveLength(0)
+ })
+
+ it('the packed digest is deterministic across reopen and changes with history', async () => {
+ await store.fold([gen(1), gen(2), gen(3)])
+ const atSeal = await store.digestThroughPacked(3)
+ const midSegment = await store.digestThroughPacked(2)
+ expect(atSeal).not.toBeNull()
+ expect(midSegment).not.toBeNull()
+ expect(midSegment).not.toBe(atSeal)
+
+ const reopened = new GenerationSegmentStore(storage as any)
+ await reopened.open()
+ expect(await reopened.digestThroughPacked(3)).toBe(atSeal)
+ expect(await reopened.digestThroughPacked(2)).toBe(midSegment)
+
+ await reopened.fold([gen(4)])
+ expect(await reopened.digestThroughPacked(4)).not.toBe(atSeal)
+ })
+
+ it('sealed segments are immutable — fold refuses overlap, requires ascending input', async () => {
+ await store.fold([gen(1), gen(2)])
+ await expect(store.fold([gen(2), gen(3)])).rejects.toThrow(/overlaps the packed tier/)
+ await expect(store.fold([gen(4), gen(4)])).rejects.toThrow(/strictly ascending/)
+ await expect(store.fold([])).rejects.toThrow(/at least one generation/)
+ })
+})