diff --git a/CHANGELOG.md b/CHANGELOG.md index a54d609e..16fb5786 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,29 @@ All notable changes to this project will be documented in this file. See [standard-version](https://github.com/conventional-changelog/standard-version) for commit guidelines. +### [10.4.9](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.6...v10.4.9) (2026-09-02) + +- Merge branch 'fix/pending-embed-low-water' into rel/10.4.9-candidate (2648f56d) +- fix(open): pending-embed recovery keeps the crash-recovery contract — foreground, bounded by the mark (8a2ebacf) +- Merge branches 'fix/connected-find-order', 'fix/pending-embed-low-water' and 'fix/related-verb-array' into rel/10.4.9-candidate (d5147ed6) +- fix(graph): the verb fast paths honour every requested type, source, and target (6a89adc4) +- perf(open): pending-embed recovery is bounded by a low-water mark and runs behind the doors (88e79729) +- fix(find): connected finds are graph-first — neighbours, then the filter over those ids, then the page (077cbc0b) +- fix(storage): counts persistence is single-flight, coalesced, and never races its own temp file (5e3b343a) + + +### [10.4.6](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.5...v10.4.6) (2026-08-31) + +- fix(transact): metadata-index ops take their JSON-safe view at the crossing, not at construction (73500e7d) + + +### [10.4.5](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.4...v10.4.5) (2026-08-31) + +- build(release): the docs-push step retires — this engine documents itself in its own repository (d6bcb14f) +- fix(generations): a sealed segment may only declare the generations it holds (a963a744) +- fix(recovery): a torn generation-log tail is a terminal verdict, never a wait (c9930871) + + ### [10.4.4](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.3...v10.4.4) (2026-08-28) - fix(vfs): the old-root sweep narrates only when it has something to say (d49148e1) diff --git a/package-lock.json b/package-lock.json index c4f66561..fc530baa 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@soulcraftlabs/brainy", - "version": "10.4.4", + "version": "10.4.9", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@soulcraftlabs/brainy", - "version": "10.4.4", + "version": "10.4.9", "license": "MIT", "dependencies": { "@msgpack/msgpack": "^3.1.2", diff --git a/package.json b/package.json index 06ce0253..f07bb94c 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@soulcraftlabs/brainy", - "version": "10.4.4", + "version": "10.4.9", "brainyContract": 1, "description": "Universal Knowledge Protocol™ - World's first Triple Intelligence database unifying vector, graph, and document search in one API. Stage 3 CANONICAL: 42 nouns × 127 verbs covering 96-97% of all human knowledge.", "main": "dist/index.js", diff --git a/scripts/release.sh b/scripts/release.sh index 08293e3a..1a4fe575 100755 --- a/scripts/release.sh +++ b/scripts/release.sh @@ -248,17 +248,12 @@ else echo -e "${RED}⚠️ FORGEJO_RELEASE_TOKEN unset — no release page created; tag + CHANGELOG remain the record${NC}\n" fi -# Step 12: Push public docs to the soulcraft.com docs ingest door -# (VENUE-DOCS-RELEASE-PUSH). Skips with a loud warning when -# DOCS_INGEST_SECRET is unset; fails loudly (without undoing the publish — -# that already happened) when a push errors, so the docs site never -# silently trails npm. -echo -e "${BLUE}1️⃣2️⃣ Pushing public docs to soulcraft.com/docs...${NC}" -if node scripts/push-docs.js; then - echo -e "${GREEN}✅ Docs push step done${NC}\n" -else - echo -e "${RED}❌ Docs push FAILED — soulcraft.com/docs trails npm until re-run or interim sync${NC}\n" -fi +# Step 12 RETIRED (2026-08-31, CORTEX-SITE-BRAINY-RENAME round 12, David-ruled): +# soulcraft.com/docs carries the paid product's documentation only. This +# engine's documentation home is THIS repository — README and docs/ — and the +# site serves 301s for the slugs this rail used to push. The push script stays +# in the tree for history; the rail no longer calls it. +echo -e "${BLUE}Docs step: this engine documents itself in its own repo (site push retired 2026-08-31)${NC}" echo -e "${GREEN}━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━${NC}" echo -e "${GREEN}🎉 Release ${NEW_VERSION} complete!${NC}" diff --git a/src/brainy.ts b/src/brainy.ts index 70d46973..e08ca411 100644 --- a/src/brainy.ts +++ b/src/brainy.ts @@ -15,6 +15,7 @@ import { JsHnswVectorIndex } from './hnsw/hnswIndex.js' import { createStorage, resolveFilesystemRoot } from './storage/storageFactory.js' import type { StorageOptions } from './storage/storageFactory.js' import { rebuildCounts } from './utils/rebuildCounts.js' +import { jsonSafeIndexMetadata } from './utils/jsonSafeIndexMetadata.js' import type { MetadataWriteBuffer } from './utils/metadataWriteBuffer.js' import { BaseStorage } from './storage/baseStorage.js' import { @@ -1819,6 +1820,11 @@ export class Brainy implements BrainyInterface { // a deferred write's ack and its background embed DELAYED a vector; // this is where it lands. if (!this.isReadOnly) { + // Foreground, as the crash-recovery contract pins it: a reopened brain + // has its markers re-armed when open() returns. The low-water mark + // bounds this to the log's tail on any brain that has ever drained — + // milliseconds — so the foreground cost is the unmarked first open + // only, once per upgraded brain. try { await step( 'bridge-pending-embed-sidecars', @@ -1827,7 +1833,7 @@ export class Brainy implements BrainyInterface { ) await step( 'recover-pending-embeds', - 'folding the generation log\'s deferred-embed markers back into the pending set', + 'folding the generation log\'s deferred-embed markers (from the low-water mark) into the pending set', () => this.recoverPendingEmbedsFromLog() ) if (this._pendingEmbedIds.size > 0) { @@ -2407,6 +2413,17 @@ export class Brainy implements BrainyInterface { */ private static readonly PENDING_EMBED_PREFIX = '_system/pending_embeds/' + /** + * Storage-root-relative path of the ADVISORY pending-embed low-water mark: + * `{ generation, writtenAt }`, written whenever the pending set drains to + * empty (and at clean close when empty). Every marker in facts at or below + * `generation` is consumed, so recovery scans from `generation + 1`. The + * mark is advisory and monotone-safe: stale-low costs a longer scan, never + * a lost marker; it is never required for correctness. + */ + private static readonly PENDING_EMBED_LOWWATER_PATH = '_system/pending_embeds_lowwater.json' + + /** * @description Mark a deferred embed pending (MT5): the id joins the * in-memory fast-path set and the returned `embed.pending` record is @@ -2434,6 +2451,40 @@ export class Brainy implements BrainyInterface { */ private clearPendingEmbed(id: string): void { this._pendingEmbedIds.delete(id) + if (this._pendingEmbedIds.size === 0) this.maybeWriteEmbedLowWater() + } + + /** + * @description Advance the advisory low-water mark: called at drain-to-empty + * (and at clean close when empty), it records the fact log's CURRENT head — + * with the set empty, every marker at or below the head has been consumed, + * so the next open's recovery fold scans only what comes after. Fire-and- + * forget at the drain (close() awaits the core); loud on failure: a missed + * write costs the next open a longer scan, never a marker. No-op without a + * fact log (no durable markers exist there) and on read-only opens. + */ + private maybeWriteEmbedLowWater(): void { + void this.writeEmbedLowWater() + } + + /** The awaitable core of {@link maybeWriteEmbedLowWater} — close() awaits it. */ + private async writeEmbedLowWater(): Promise { + if (this.isReadOnly) return + const log = this.generationStore ? this.generationStore.getFactLog() : null + if (!log) return + const generation = log.headGeneration() + if (!(generation > 0)) return + try { + await this.storage.writeRawObject(Brainy.PENDING_EMBED_LOWWATER_PATH, { + generation, + writtenAt: Date.now() + }) + } catch (err) { + prodLog.warn( + `[Brainy] pending-embed low-water write failed at generation ${generation}: ` + + `${(err as Error).message} — the next open scans from the previous mark` + ) + } } /** @@ -2444,9 +2495,14 @@ export class Brainy implements BrainyInterface { * survives the fold is exactly the set of acknowledged deferred writes * whose vectors have not landed. * - * BOUND (honest): no durable low-water mark exists for the earliest - * unconsumed pending, so the fold scans the log's committed facts from - * generation 1 — a sequential read of the log at open, O(log bytes). + * BOUND: the scan starts at the advisory low-water mark + * ({@link Brainy.PENDING_EMBED_LOWWATER_PATH}) — the log head at which the + * pending set last drained to empty — so a settled brain reads only the + * facts since then, not its whole history. Without a mark (first open + * after upgrade) it scans from generation 1, once; a stale-low mark costs + * a longer scan, never a marker. The fold stays on the open's foreground — + * the crash-recovery contract pins that a reopened brain has its markers + * re-armed when open() returns — and the mark is what makes that cheap. * It is SKIPPED WHOLESALE when the log has never had a v2 tail * ({@link FactLog.hasV2History} — v1 facts cannot carry marker records), * so pre-cutover brains pay nothing; on a mixed log the scan still reads @@ -2459,7 +2515,18 @@ export class Brainy implements BrainyInterface { private async recoverPendingEmbedsFromLog(): Promise { const log = this.generationStore.getFactLog() if (!log || !log.hasV2History()) return - const scan = log.scanFacts({ fromGeneration: 1 }) + let fromGeneration = 1 + try { + const mark = (await this.storage.readRawObject(Brainy.PENDING_EMBED_LOWWATER_PATH)) as { + generation?: number + } | null + if (mark && typeof mark.generation === 'number' && mark.generation > 0) { + fromGeneration = mark.generation + 1 + } + } catch { + // No mark (or unreadable): scan from 1 — correctness over cost. + } + const scan = log.scanFacts({ fromGeneration }) for await (const batch of scan.batches()) { for (const fact of batch.facts) { for (const record of fact.records ?? []) { @@ -4203,32 +4270,19 @@ export class Brainy implements BrainyInterface { */ /** * @description A JSON-safe view of a record bound for the metadata-index - * crossing. The seam's metadata is JSON-safe BY CONTRACT (a native provider - * serializes it; u64 ints as Number corrupt above 2^53) — but - * {@link resolveVerbEndpointInts} MIRRORS the resolved endpoint ints onto - * the verb object itself as BigInt (`verb.sourceInt`/`targetInt`), so a - * verb object reused as index metadata carried BigInts into - * JSON.stringify, which throws, aborting the whole transaction (found by - * the first joint pair gate). Endpoint ints ride their OWN op params on the - * graph legs — the metadata crossing drops every BigInt-valued top-level - * key instead of guessing at a lossy numeric encoding. + * crossing — delegates to the shared {@link jsonSafeIndexMetadata} leaf, + * which the metadata-index transaction operations ALSO apply at execute + * and rollback time. This plan-time wrap alone proved insufficient: it + * returns the same reference when the record is clean, and `transact()`'s + * delete legs share that reference with a graph-retraction op whose + * execute-time endpoint resolution mirrors BigInt ints onto it (the full + * aliasing story lives on the leaf module's doc). * @param metadata - The candidate index-metadata record. * @returns The same object when already JSON-safe, else a shallow copy * without the BigInt-valued keys. */ private static jsonSafeIndexMetadata(metadata: unknown): unknown { - if (metadata === null || typeof metadata !== 'object') return metadata - const rec = metadata as Record - let hasBigint = false - for (const k in rec) { - if (typeof rec[k] === 'bigint') { hasBigint = true; break } - } - if (!hasBigint) return metadata - const out: Record = {} - for (const k in rec) { - if (typeof rec[k] !== 'bigint') out[k] = rec[k] - } - return out + return jsonSafeIndexMetadata(metadata) } private metadataIndexRetractionOp( @@ -7482,7 +7536,37 @@ export class Brainy implements BrainyInterface { // JS path — there the materialized `candidateIds` restricts the walk instead. let preResolvedAllowedIds: OpaqueIdSet | undefined - if (params.where || params.type || params.subtype || params.service || params.excludeVFS) { + // Graph-first law (10.4.8, BRAINY-PROD-LATENCY-TRIAD rounds 44/45): with + // `connected` present the NEIGHBOUR SET is the candidate universe. It is + // resolved first from the adjacency (O(neighbours)), the metadata filter + // is evaluated over those ids only, and paging happens LAST. The earlier + // order materialized the whole-store filtered id list, paged it, hydrated + // the page, and only then intersected with the neighbours — O(store) per + // call, and a neighbour outside the first page was silently dropped. + let graphFirstIds: string[] | null = null + if (hasGraphCriteria) { + graphFirstIds = await this.resolveConnectedIds(params) + if (hiddenIds.size > 0) { + graphFirstIds = graphFirstIds.filter((id) => !hiddenIds.has(id)) + } + if ( + graphFirstIds.length > 0 && + (params.where || params.type || params.subtype || params.service || params.excludeVFS) + ) { + preResolvedFilter = this.buildMetadataFilter(params) + graphFirstIds = await this.filterIdsWithinBelted(preResolvedFilter, graphFirstIds) + } + if (graphFirstIds.length === 0) { + return [] + } + if (!hasVectorSearchCriteria) { + return await this.pageConnectedIds(params, graphFirstIds) + } + // The vector leg walks ONLY the neighbours (its candidate walk). The + // filter is already applied above, so no opaque universe is produced — + // it would describe the whole store, not the neighbour set. + preResolvedMetadataIds = graphFirstIds + } else if (params.where || params.type || params.subtype || params.service || params.excludeVFS) { preResolvedFilter = this.buildMetadataFilter(params) preResolvedMetadataIds = await this.filterIdsBelted(preResolvedFilter) @@ -7671,9 +7755,11 @@ export class Brainy implements BrainyInterface { } } - // Graph search component with O(1) traversal - if (params.connected) { - results = await this.executeGraphSearch(params, results) + // The text leg of a hybrid find has no candidate door, so its hits are + // held to the neighbour set here; the vector leg walked only the neighbours. + if (graphFirstIds !== null && results.length > 0) { + const neighbourSet = new Set(graphFirstIds) + results = results.filter((r) => neighbourSet.has(r.id)) } // Apply fusion scoring if requested @@ -12497,6 +12583,18 @@ export class Brainy implements BrainyInterface { * healed by `repairIndex()`, whose unconditional recount rebuilds the * rollups from a canonical walk and re-stamps. Best-effort: a stamp-write * fault warns loudly but never fails the flush that carried real data. + * + * THE SOURCE IS `committedGeneration()`, NEVER `generation()`. The latter is + * the ALLOCATED counter — a number a write in flight has claimed and may + * never commit. Stamping it made the stamp's generation label a claim about + * counts it was not taken at, and every crash inside a write window then + * produced a spurious verdict at the next open: either `sourceGeneration N + * is ahead of the log head N-1` (the allocated generation died with the + * process) or `rollup invariant 'nounCount': stamped X, observed Y` (the + * recovery fold folded facts the stamp's counts predate). MEASURED on the + * crash-consistency lane before this line changed: 4 of 11 SIGKILL cycles on + * a coherent store raised one of those two verdicts, each of them naming + * `repairIndex()` — a whole-store recount — as the cure for nothing. */ private async stampEntityTree(): Promise { if (this.isReadOnly) return @@ -12507,7 +12605,7 @@ export class Brainy implements BrainyInterface { ]) await writeFamilyStamp(this.storage, ENTITY_TREE_STAMP_PATH, { family: 'entity-tree', - sourceGeneration: this.generationStore.generation(), + sourceGeneration: this.generationStore.committedGeneration(), members: { mode: 'rollup', invariants: { nounCount, verbCount } } }) } catch (error) { @@ -12520,16 +12618,24 @@ export class Brainy implements BrainyInterface { /** * @description Open-time coherence check for the entity tree's family stamp: - * compare `sourceGeneration` against the log head and the stamped rollup - * invariants against the live counters. Verdicts: + * compare `sourceGeneration` against the store's COMMITTED generation and + * the stamped rollup invariants against the live counters. Verdicts: * - `coherent` / `absent` (legacy store; first flush stamps) → silent. * - `behind` → benign for the tree (it is written BY the commit; only the * stamp is stale — a crash landed between commit and flush). Refreshed at * the next flush. + * - `torn` → a TORN GENERATION-LOG TAIL, handled by + * {@link demoteTornEntityTreeStamp}: terminal, never a wait. * - `incoherent` → LOUD: the tree or its counters diverged from what was * stamped — `repairIndex()` recounts from canonical and re-stamps. * Never blocks open; a fault reading the stamp is surfaced as unverifiable, * never conflated with absence. + * + * THE COMPARISON IS AGAINST `committedGeneration()`, matching what + * {@link stampEntityTree} writes and what every other open-time watermark in + * this class already reasons about (the fact-scan capability, the metadata / + * graph / HNSW watermark verdicts). Comparing against the allocated counter + * was the one place that disagreed, and disagreeing was the whole defect. */ private async verifyEntityTreeStamp(): Promise { let stamp: FamilyStamp | null @@ -12546,11 +12652,16 @@ export class Brainy implements BrainyInterface { this.storage.getNounCount(), this.storage.getVerbCount() ]) - const verdict = verifyFamilyStamp(stamp, this.generationStore.generation(), { + const verdict = verifyFamilyStamp(stamp, this.generationStore.committedGeneration(), { nounCount, verbCount }) - if (verdict.state === 'incoherent') { + if (verdict.state === 'torn') { + await this.demoteTornEntityTreeStamp(stamp as FamilyStamp, verdict.stampSource, verdict.head, { + nounCount, + verbCount + }) + } else if (verdict.state === 'incoherent') { prodLog.warn( `[Brainy] entity-tree stamp INCOHERENT at open: ${verdict.failures.join('; ')}. ` + `The canonical tree or its counters diverged from the stamped state — run ` + @@ -12564,6 +12675,92 @@ export class Brainy implements BrainyInterface { } } + /** + * @description THE TERMINAL VERDICT for a torn generation-log tail. + * + * A stamp whose `sourceGeneration` sits ABOVE the store's committed + * watermark witnesses a generation that is not in the log: the stamp's fsync + * outlived the tail's. By the time this runs, log-authority recovery has + * already folded every intact fact above the manifest and advanced the + * watermark to cover them — so if the stamp is STILL ahead, the generation + * it names is not merely late, it is GONE. There is nothing to wait for. + * + * That is the whole point of this method. A field report of this class + * (single-process store, abrupt termination mid-fold) described a reopen + * that narrated the tear and then held 100% CPU with zero log growth for + * eight minutes before an operator wiped the directory. A recovery that + * cannot say what it is waiting for has no business spinning; the honest + * answer here is a verdict, taken now, at O(1) cost. + * + * WHAT THE VERDICT DOES — the stamped surface is UNUSABLE, so it is + * discarded rather than believed: the stamped counts describe a generation + * that never became durable, and comparing them against live counters can + * only produce noise. The tree itself is not in question (it IS canonical — + * every commit writes it, and the fold re-applied every after-image the log + * still holds), so the demotion is a re-derivation of this family's verified + * surface at the generation the store can actually show: + * + * - WRITER open → re-stamp at `committedGeneration()` from the live + * counters — exactly what the next flush would write, taken now so the + * tear cannot re-narrate on every subsequent open. Both count sets are + * logged so an operator can see whether anything really moved. + * - READER open → a reader cannot re-stamp. Narrate the same terminal + * verdict with the named cure and carry on serving; a read-only inspector + * is never locked out of a store, and never left waiting either. + * + * BOUNDEDNESS: straight-line code. No loop, no retry, no await on any + * external progress signal — the two counter reads and one stamp write are + * the entire cost, and none of them scales with the store. + */ + private async demoteTornEntityTreeStamp( + stamp: FamilyStamp, + stampSource: number, + head: number, + observed: { nounCount: number; verbCount: number } + ): Promise { + const stamped = stamp.members.mode === 'rollup' ? stamp.members.invariants : {} + const detail = + `[Brainy] TORN GENERATION-LOG TAIL at open: ${ENTITY_TREE_STAMP_PATH} witnesses source ` + + `generation ${stampSource} (stamped ${stamp.committedAt}), but the store's committed ` + + `generation is ${head} after crash recovery — the stamp's fsync outlived the log tail's, ` + + `and generation ${stampSource} is not in the log to arrive. Stamped rollups ` + + `${JSON.stringify(stamped)}; observed ${JSON.stringify(observed)}.` + + if (this.isReadOnly) { + prodLog.warn( + `${detail} This open is READ-ONLY, so the stamp cannot be re-derived: the entity-tree ` + + `family stays UNVERIFIED for this session (reads are unaffected — the canonical tree ` + + `is the truth this stamp only describes). Cure: open the store with a writer, or run ` + + `brain.repairIndex() there, to recount from canonical and re-stamp.` + ) + return + } + + const startedAt = Date.now() + try { + await writeFamilyStamp(this.storage, ENTITY_TREE_STAMP_PATH, { + family: 'entity-tree', + sourceGeneration: head, + members: { + mode: 'rollup', + invariants: { nounCount: observed.nounCount, verbCount: observed.verbCount } + } + }) + prodLog.warn( + `${detail} DEMOTED: the unusable stamp was re-derived at committed generation ${head} ` + + `from the live counters in ${Date.now() - startedAt}ms — terminal, not a wait. If the ` + + `observed counts above look wrong for your data, run brain.repairIndex() to recount ` + + `from canonical.` + ) + } catch (error) { + prodLog.warn( + `${detail} The demotion's re-stamp FAILED (${(error as Error).message}) — the tear will ` + + `narrate again at the next open, which is the honest outcome; the store still serves ` + + `from canonical. Cure: run brain.repairIndex() to recount from canonical and re-stamp.` + ) + } + } + /** * Ask the writer process serving this data directory to flush its in-memory * indexes to disk, so a read-only inspector can observe fresh state. @@ -12677,6 +12874,29 @@ export class Brainy implements BrainyInterface { } } + /** + * The id-scoped twin of {@link filterIdsBelted}: evaluate `filter` over `ids` + * only, through the provider's own evaluation so the answer can never drift + * from `getIdsForFilter`'s. A provider without the door is served by its + * whole-store answer intersected here (the reference index implements the + * door itself). Same belt: field refusals cross as `BrainyFieldRefusal`. + */ + private async filterIdsWithinBelted(filter: unknown, ids: readonly string[]): Promise { + this.ensureIndexesLoaded(['metadata']) + const mip = this.metadataIndex as unknown as MetadataIndexProvider + try { + if (typeof mip.filterIdsWithin === 'function') { + return await mip.filterIdsWithin(filter, ids) + } + const matched = new Set(await this.metadataIndex.getIdsForFilter(filter)) + return ids.filter((id) => matched.has(id)) + } catch (err) { + const normalized = asBrainyFieldRefusal(err) + if (normalized) throw normalized + throw err + } + } + async getIndexStatus(): Promise<{ initialized: boolean /** `true` once open()'s index-build-if-needed step has run. Named for API @@ -15660,16 +15880,16 @@ export class Brainy implements BrainyInterface { } /** - * Execute graph search component. + * Resolve `params.connected` to the neighbour id set — the graph-first + * find's candidate universe (deterministic traversal order, anchors excluded). * * Honors the full `GraphConstraints` contract: multi-hop `depth` (breadth-first via - * `neighbors()`), `via`/`type` verb-type filtering, and `direction`. Previously this read - * only `from`/`to`/`direction` and did a single 1-hop `getNeighbors()`, so `depth` and `via` - * were silently ignored — `find({ connected: { from, depth: 3 } })` returned only the - * immediate neighbour at every depth. + * `neighbors()`), `via`/`type` verb-type filtering, and `direction`. An empty set + * is re-verified against the adjacency before it is believed — a not-serving + * adjacency throws rather than answering `[]` as truth. */ - private async executeGraphSearch(params: FindParams, existingResults: Result[]): Promise[]> { - if (!params.connected) return existingResults + private async resolveConnectedIds(params: FindParams): Promise { + if (!params.connected) return [] const { from, to, depth, direction = 'both' } = params.connected const via = params.connected.via ?? params.connected.type @@ -15723,8 +15943,8 @@ export class Brainy implements BrainyInterface { if (anchorInt === undefined) return new Set() // unmapped → no relations const verbTypeIndex = TypeUtils.getVerbIndex(via as VerbType) - // No limit: match the JS BFS exactly — overall result limiting happens - // downstream against existingResults. + // No limit: match the JS BFS exactly — the page is cut downstream, + // after the metadata filter, by pageConnectedIds / the candidate walk. const reachedInts = await provider.findConnectedSubtype( anchorInt, verbTypeIndex, subtypeArr[0], effectiveDepth, null ) @@ -15809,22 +16029,44 @@ export class Brainy implements BrainyInterface { await this.verifyGraphAdjacencyLive() } - // Filter existing results to only connected entities - if (existingResults.length > 0) { - return existingResults.filter(r => connectedIds.has(r.id)) - } + return [...connectedIds] + } - // Batch-load connected entities for fast cloud-storage performance + /** + * Page and hydrate an already-filtered neighbour set — the pure graph (and + * graph + metadata) find's tail. `orderBy` sorts the WHOLE set by field value + * before the page is cut (never the page after), null values last on `asc` + * and first on `desc`; without `orderBy` the traversal order stands. + */ + private async pageConnectedIds(params: FindParams, ids: string[]): Promise[]> { + const limit = params.limit || 10 + const offset = params.offset || 0 + let ordered = ids + if (params.orderBy) { + const field = params.orderBy + const asc = (params.order || 'asc') === 'asc' + const valued = await Promise.all( + ids.map(async (id) => ({ id, value: await this.metadataIndex.getFieldValueForEntity(id, field) })) + ) + valued.sort((a, b) => { + if (a.value == null && b.value == null) return 0 + if (a.value == null) return asc ? 1 : -1 + if (b.value == null) return asc ? -1 : 1 + if (a.value === b.value) return 0 + const comparison = a.value < b.value ? -1 : 1 + return asc ? comparison : -comparison + }) + ordered = valued.map((v) => v.id) + } + const pageIds = ordered.slice(offset, offset + limit) + const entitiesMap = await this.batchGet(pageIds) const results: Result[] = [] - const ids = [...connectedIds] - const entitiesMap = await this.batchGet(ids) - for (const id of ids) { + for (const id of pageIds) { const entity = entitiesMap.get(id) if (entity) { results.push(this.createResult(id, 1.0, entity)) } } - return results } @@ -19344,6 +19586,7 @@ export class Brainy implements BrainyInterface { * terminal releases have run. */ async close(): Promise { + if (this._pendingEmbedIds.size === 0) await this.writeEmbedLowWater() let closeFailure: unknown = null try { await this.closeDurableSteps() diff --git a/src/db/familyStamp.ts b/src/db/familyStamp.ts index 98342884..2f01e935 100644 --- a/src/db/familyStamp.ts +++ b/src/db/familyStamp.ts @@ -12,9 +12,11 @@ * the verified surface is a small set of rollup invariants (entity/ * relationship counts) plus `sourceGeneration`. * - * `sourceGeneration` is the generation of the source-of-truth log this - * projection reflects — open-time coherence becomes a COMPARISON (stamp vs - * log head), not a walk: + * `sourceGeneration` is the COMMITTED generation of the source-of-truth log + * this projection reflects — never the allocated counter, which names a + * generation that may never commit (see {@link StampVerdict.torn}) — so + * open-time coherence becomes a COMPARISON (stamp vs committed head), not a + * walk: * * - equal + invariants hold → coherent, serve. * - behind → the projection missed the tail (crash between commit and stamp); @@ -24,6 +26,9 @@ * - invariants FAIL at equal generation → genuine incoherence: loud, and the * repair ritual (`repairIndex()`, whose recount rebuilds the rollups from a * canonical walk) heals it. + * - AHEAD → a torn generation-log tail: the stamp's fsync outlived the log + * tail's. TERMINAL, never a wait — the generation the stamp names does not + * exist to arrive. * * Stamps are JSON on purpose — every incident gets debugged by reading a * stamp in a terminal. @@ -70,6 +75,12 @@ export type StampVerdict = | { state: 'coherent' } | { state: 'absent' } // legacy store — first stamp writes at the next flush | { state: 'behind'; stampSource: number; head: number } + /** + * TORN GENERATION-LOG TAIL: the stamp witnesses a source generation the + * store's committed watermark can no longer show. TERMINAL — there is no + * generation to wait for, so the open demotes (or refuses) and never spins. + */ + | { state: 'torn'; stampSource: number; head: number } | { state: 'incoherent'; failures: string[] } | { state: 'unverifiable'; reason: string } // a FAULT reading the stamp — never conflated with absence @@ -118,12 +129,15 @@ export function verifyFamilyStamp( ): StampVerdict { if (stamp === null) return { state: 'absent' } if (stamp.sourceGeneration > head) { - // A stamp AHEAD of the log claims state that never committed — the - // projection was stamped against truth that a crash rolled back. - return { - state: 'incoherent', - failures: [`sourceGeneration ${stamp.sourceGeneration} is ahead of the log head ${head}`] - } + // A stamp AHEAD of committed truth witnesses a generation the store can no + // longer show: the stamp's fsync survived a crash that the log tail did + // not. This is the TORN GENERATION-LOG TAIL — its own class, never folded + // in with `incoherent` (a count that drifted at a generation both sides + // agree on), because the two have opposite cures: incoherence is recounted, + // a tear is DEMOTED. It is also terminal by construction — there is no + // generation the open can wait for, because the one the stamp names is + // gone. + return { state: 'torn', stampSource: stamp.sourceGeneration, head } } if (stamp.sourceGeneration < head) { return { state: 'behind', stampSource: stamp.sourceGeneration, head } diff --git a/src/db/generationSegments.ts b/src/db/generationSegments.ts index 0c14b60c..91451281 100644 --- a/src/db/generationSegments.ts +++ b/src/db/generationSegments.ts @@ -147,6 +147,60 @@ export class GenerationSegmentStore { return this.coveringSegment(gen) !== null } + /** + * @description True when `meta` declares more generations than it holds + * frames — a segment sealed by a writer that folded across a hole. The + * manifest records `frames` at fold time, so this is an O(1) comparison + * against the declared span and needs no I/O. + */ + private isSparse(meta: SegmentMeta): boolean { + return meta.lastGeneration - meta.firstGeneration + 1 !== meta.frames + } + + /** + * @description The generations this tier ACTUALLY holds, as coalesced + * ascending intervals — not what the segments declare. + * + * Dense segments (every one a current writer produces) contribute their + * declared range with no I/O. A SPARSE segment — one sealed before the + * density law was enforced, whose declared range spans generations it has + * no frame for — has its real generation list read from its sidecar and + * contributed instead, with the discrepancy narrated once. + * + * This is what keeps a store that already carries the damage from wedging. + * `open()` seeds `committedRanges` from these intervals, so a hole is never + * re-admitted as a committed generation, and the auto-compaction pass that + * used to fail on every run with "packed history is damaged" simply never + * asks for the missing frame. + * + * @returns Ascending, non-overlapping `[first, last]` intervals. + */ + async actualRanges(): Promise> { + const out: Array<[number, number]> = [] + for (const meta of this.manifest.segments) { + if (!this.isSparse(meta)) { + out.push([meta.firstGeneration, meta.lastGeneration]) + continue + } + const missing = meta.lastGeneration - meta.firstGeneration + 1 - meta.frames + prodLog.warn( + `[GenerationSegments] sealed segment ${meta.file} declares generations ` + + `${meta.firstGeneration}..${meta.lastGeneration} but holds only ${meta.frames} ` + + `frame(s) — ${missing} generation(s) in that span were never folded into it. ` + + `Serving the frames it actually holds; the declared span is not treated as ` + + `committed history. (Written by a pre-density-law writer that folded across a ` + + `gap; the segment itself is intact and no record is lost.)` + ) + const idx = await this.sidecarFor(meta) + for (const [gen] of idx.generations) { + const last = out[out.length - 1] + if (last !== undefined && gen === last[1] + 1) last[1] = gen + else out.push([gen, gen]) + } + } + return out + } + /** * Fold consecutive generations into ONE new sealed segment + sidecar and * append it to the manifest atomically. Caller guarantees: `gens` is @@ -164,6 +218,38 @@ export class GenerationSegmentStore { throw new Error('[GenerationSegments] fold() input must be strictly ascending') } } + // THE DENSITY LAW, MADE MECHANICAL. + // + // A sealed segment declares a CONTIGUOUS range [firstGeneration, + // lastGeneration] and every reader treats that range as containment: + // `coveringSegment` is an interval test, `hasGeneration` returns true for + // anything inside it, and `open()` seeds committedRanges from it. So a + // segment folded from a SPARSE input silently claims generations it does + // not hold, and the first read of one of those holes throws + // "inside sealed segment ... but has no frame — packed history is damaged". + // + // That is exactly how the damage was produced. `repackHistory` skipped + // generations mid-batch — ones absent from committedRanges, ones still in + // the pending buffer, ones whose tx.json would not read — and handed the + // survivors here, where the range was computed from the first and last of + // them. Worse, the mis-declared range was then merged back into + // committedRanges at the next open, which is what turned a quiet hole into + // a repeating auto-compaction failure on every subsequent run. + // + // Callers now split at discontinuities; this refusal is what keeps any + // future caller from reintroducing the class. A refusal here loses + // nothing — the generations stay in the live tier, readable, and the next + // pass folds them correctly. + for (let i = 1; i < gens.length; i++) { + if (gens[i].generation !== gens[i - 1].generation + 1) { + throw new Error( + `[GenerationSegments] fold() input is not contiguous: ${gens[i - 1].generation} → ` + + `${gens[i].generation} skips ${gens[i].generation - gens[i - 1].generation - 1} ` + + `generation(s). A sealed segment declares a dense range, so folding a sparse ` + + `batch would claim generations it does not hold. Split the batch at the gap.` + ) + } + } const last = this.manifest.segments[this.manifest.segments.length - 1] if (last && gens[0].generation <= last.lastGeneration) { throw new Error( @@ -364,12 +450,37 @@ export class GenerationSegmentStore { 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. + // Inside the covering range but with no frame. Two very different causes, + // and conflating them is what made this class wedge every maintenance pass + // on the affected stores. + // + // (1) A SPARSE SEGMENT — the manifest's own `frames` count is smaller than + // the span it declares. That segment was sealed by a writer that + // folded across a hole (the class this file's density law now bars). + // The segment is INTACT and nothing is lost; it simply never held this + // generation. Answering "not packed" is the honest answer, and it lets + // the caller's two-tier read decide what a genuinely absent generation + // means, instead of every compaction pass dying on a repeating throw. + // `actualRanges()` keeps such holes out of committedRanges at open, so + // in a healed store nobody asks this question in the first place. + // + // (2) A DENSE SEGMENT missing a frame it says it has — the manifest and + // the sidecar disagree about a segment that claims to be complete. + // That IS damage, and it stays loud. + if (this.isSparse(meta)) { + prodLog.warn( + `[GenerationSegments] generation ${gen} falls inside sealed segment ${meta.file}'s ` + + `declared range ${meta.firstGeneration}..${meta.lastGeneration}, but that segment ` + + `holds ${meta.frames} frame(s) for a ${meta.lastGeneration - meta.firstGeneration + 1}` + + `-generation span — it was sealed across a gap and never held this generation. ` + + `Reporting it as unpacked rather than as damage; no record is lost.` + ) + return null + } throw new Error( `[GenerationSegments] generation ${gen} is inside sealed segment ${meta.file}'s declared ` + - `range but has no frame — packed history is damaged` + `range but has no frame, and that segment declares a complete ${meta.frames}-frame ` + + `span — the manifest and the sidecar disagree; packed history is damaged` ) } diff --git a/src/db/generationStore.ts b/src/db/generationStore.ts index fd052c31..da21dc61 100644 --- a/src/db/generationStore.ts +++ b/src/db/generationStore.ts @@ -96,6 +96,35 @@ export const FOLD_CHECKPOINT_PATH = '_system/fold-checkpoint.json' /** Storage-root-relative prefix of the per-generation record directories. */ export const GENERATIONS_PREFIX = '_generations' +/** + * @description Split an ascending list of fold candidates into maximal + * CONTIGUOUS runs — `[7,8,9,12,13]` becomes `[[7,8,9],[12,13]]`. + * + * A sealed segment declares one dense range `[firstGeneration, + * lastGeneration]`, and every reader treats that range as containment. So a + * batch with a hole in it must never become one segment: it would claim a + * generation it does not hold, and the first read of that hole reports the + * packed history as damaged. One run, one segment — the ranges then describe + * exactly what the segments contain. + * + * @param gens - Fold candidates, strictly ascending by generation. + * @returns One array per contiguous run, in ascending order. Empty in, empty out. + */ +export function contiguousRuns(gens: FoldGeneration[]): FoldGeneration[][] { + const runs: FoldGeneration[][] = [] + let run: FoldGeneration[] = [] + for (const g of gens) { + const prev = run[run.length - 1] + if (prev !== undefined && g.generation !== prev.generation + 1) { + runs.push(run) + run = [] + } + run.push(g) + } + if (run.length > 0) runs.push(run) + return runs +} + /** * @description Phases of the {@link GenerationStore.commitTransaction} commit * protocol at which a test-only fault injector can simulate a process crash. @@ -784,9 +813,15 @@ export class GenerationStore { 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)]) + // ACTUAL ranges, not declared ones. A segment sealed by a pre-density-law + // writer can declare a span wider than the frames it holds; seeding + // committedRanges from the declared span re-admits those holes as + // committed generations, and every later maintenance pass then asks for a + // frame that was never written. `actualRanges()` reads the real + // generation list from the sidecar for exactly those segments (and does + // no I/O for the dense ones, which is all of them on a healthy store). + const packedRanges = (await this.segments.actualRanges()) + .map((r): [number, number] => [r[0], Math.min(r[1], this.committed)]) .filter(([lo, hi]) => lo <= hi) if (packedRanges.length > 0) { // Merge packed (older) + live (newer) interval sets — both ascending; @@ -3121,13 +3156,26 @@ export class GenerationStore { 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}`) + // SPLIT AT DISCONTINUITIES. `eligible` is NOT contiguous — three + // filters above punch holes in it: a generation missing from + // committedRanges never appears, one still in the pending buffer is + // skipped, and one whose tx.json will not read is skipped. A sealed + // segment declares a DENSE range, so folding across such a hole makes + // the segment claim a generation it does not hold; the next open + // merges that mis-declared range into committedRanges, and every + // subsequent auto-compaction pass then asks for the missing frame and + // fails with "packed history is damaged". Fold each contiguous RUN as + // its own segment instead — same bytes, honest ranges. + for (const run of contiguousRuns(foldInput)) { + if (deadline !== undefined && Date.now() >= deadline) break + await segments.fold(run) + segmentsCreated++ + // Segment + manifest durable → the live copies retire. + for (const g of run) { + await this.storage.removeRawPrefix(`${GENERATIONS_PREFIX}/${g.generation}`) + } + folded += run.length } - folded += foldInput.length } if (folded > 0) { prodLog.info( diff --git a/src/plugin.ts b/src/plugin.ts index b1aef8e0..15b14b4e 100644 --- a/src/plugin.ts +++ b/src/plugin.ts @@ -411,6 +411,19 @@ export interface MetadataIndexProvider { * @returns The matching id universe as an opaque set. */ getIdSetForFilter?(filter: any): Promise + /** + * @description OPTIONAL: evaluate `filter` over `ids` ONLY and return the + * survivors in the caller's order — the door a graph-first + * `find({ connected, where })` walks. The neighbour set is the universe there, + * so the filter must cost O(|ids|) membership checks, never a whole-store + * materialization. A native index answers from its roaring filter result + * (membership by entity int); the reference index answers from its own + * `getIdsForFilter`, so the two doors can never disagree. Absent → Brainy + * intersects `getIdsForFilter`'s answer with `ids` itself (correct, O(store)). + * @param filter - The same filter shape accepted by `getIdsForFilter`. + * @param ids - The candidate ids (canonical). The answer is a subsequence. + */ + filterIdsWithin?(filter: any, ids: readonly string[]): Promise getIdsForTextQuery(query: string): Promise> getSortedIdsForFilter(filter: any, orderBy: string, order?: 'asc' | 'desc', topK?: number): Promise getFilterValues(field: string): Promise diff --git a/src/storage/adapters/baseStorageAdapter.ts b/src/storage/adapters/baseStorageAdapter.ts index cabe2e30..a90adb93 100644 --- a/src/storage/adapters/baseStorageAdapter.ts +++ b/src/storage/adapters/baseStorageAdapter.ts @@ -1089,6 +1089,10 @@ export abstract class BaseStorageAdapter implements StorageAdapter { // Counts changed since the last persist? Drives the write-through flush. protected pendingCountPersist = false + /** The one persist running right now, if any (single-flight law — see flushCounts). */ + private countPersistInFlight: Promise | null = null + /** The one trailing persist a burst has queued behind the in-flight one. */ + private countPersistTrailing: Promise | null = null /** * Get total noun count - O(1) operation @@ -1341,15 +1345,46 @@ export abstract class BaseStorageAdapter implements StorageAdapter { return } - try { - // Persist to storage (implemented by subclass) - await this.persistCounts() - this.pendingCountPersist = false - } catch (error) { - console.error('CRITICAL: Failed to flush counts to storage:', error) - // Keep pending flag set so we retry on next operation - throw error + // SINGLE-FLIGHT, COALESCED. Counts are write-through on every change, so + // a burst of writes used to launch one persist per change, all in flight + // together. Two of them inside the same millisecond shared the atomic + // writer's temp path (`.tmp--`): both wrote it, the first rename + // consumed it, the second rename found nothing — ENOENT, ~1,500 times a + // day on a busy production brain, with a full ledger write per change + // behind it. Now exactly one persist runs at a time; requests that arrive + // while it runs collapse into ONE trailing persist that carries the final + // state. A burst of N changes costs at most two writes and never races + // itself. + if (this.countPersistInFlight) { + // The in-flight write may have already serialised a stale snapshot — + // ask for one more pass after it, and let every caller in this burst + // await that same pass. + if (!this.countPersistTrailing) { + this.countPersistTrailing = this.countPersistInFlight + .catch(() => undefined) + .then(() => { + this.countPersistTrailing = null + return this.flushCounts() + }) + } + return this.countPersistTrailing } + + this.countPersistInFlight = (async () => { + try { + // Persist to storage (implemented by subclass) + this.pendingCountPersist = false + await this.persistCounts() + } catch (error) { + // Keep the flag set so the next operation retries. + this.pendingCountPersist = true + console.error('CRITICAL: Failed to flush counts to storage:', error) + throw error + } finally { + this.countPersistInFlight = null + } + })() + return this.countPersistInFlight } /** diff --git a/src/storage/adapters/fileSystemStorage.ts b/src/storage/adapters/fileSystemStorage.ts index 5ec1d88e..87b6406f 100644 --- a/src/storage/adapters/fileSystemStorage.ts +++ b/src/storage/adapters/fileSystemStorage.ts @@ -2400,8 +2400,15 @@ export class FileSystemStorage extends BaseStorage { * Atomic write via temp-file-then-rename so concurrent readers never see a * half-written lock JSON. Reused by writer-lock writes + heartbeat. */ + /** Monotonic per-process sequence so two atomic writes never share a temp path. */ + private static atomicWriteSeq = 0 + private async writeFileAtomic(filePath: string, contents: string): Promise { - const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}` + // pid + timestamp alone collided: two writers of the same target inside + // one millisecond shared this path, and the loser's rename found the + // winner had already moved it (ENOENT). The sequence makes every call's + // temp path its own. + const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}-${++FileSystemStorage.atomicWriteSeq}` await fs.promises.writeFile(tmp, contents) await fs.promises.rename(tmp, filePath) } diff --git a/src/storage/baseStorage.ts b/src/storage/baseStorage.ts index d8bcb780..a1cc2e35 100644 --- a/src/storage/baseStorage.ts +++ b/src/storage/baseStorage.ts @@ -2942,19 +2942,33 @@ export abstract class BaseStorage extends BaseStorageAdapter { !options.filter.service && !options.filter.metadata ) { - const sourceId = Array.isArray(options.filter.sourceId) - ? options.filter.sourceId[0] - : options.filter.sourceId + const sourceIds = Array.isArray(options.filter.sourceId) + ? options.filter.sourceId + : [options.filter.sourceId] - const verbType = Array.isArray(options.filter.verbType) - ? options.filter.verbType[0] - : options.filter.verbType + // EVERY requested verb type is honoured — an array used to collapse to + // its first element here, silently dropping the rest of the ask. + const verbTypes = new Set( + Array.isArray(options.filter.verbType) + ? options.filter.verbType + : [options.filter.verbType] + ) - // Get verbs by source, then filter by type (O(1) graph lookup + O(n) type filter), - // then apply the subtype / visibility metadata filters on the candidate set. - const verbsBySource = await this.getVerbsBySource_internal(sourceId) + // Get verbs by source (union over every requested source), filter by the + // requested type SET (O(1) graph lookup + O(n) type filter), then apply + // the subtype / visibility metadata filters on the candidate set. + const bySource: HNSWVerbWithMetadata[] = [] + const seenVerbIds = new Set() + for (const oneSource of sourceIds) { + for (const v of await this.getVerbsBySource_internal(oneSource)) { + if (!seenVerbIds.has(v.id)) { + seenVerbIds.add(v.id) + bySource.push(v) + } + } + } const filteredVerbs = this.applyVerbMetadataFilters( - verbsBySource.filter(v => v.verb === verbType), + bySource.filter(v => verbTypes.has(v.verb)), options.filter ) @@ -2985,16 +2999,22 @@ export abstract class BaseStorage extends BaseStorageAdapter { !options.filter.service && !options.filter.metadata ) { - const sourceId = Array.isArray(options.filter.sourceId) - ? options.filter.sourceId[0] - : options.filter.sourceId - - // Get verbs by source directly (hydrated with metadata), then apply the - // subtype / visibility metadata filters on the O(degree) candidate set. - const verbsBySource = this.applyVerbMetadataFilters( - await this.getVerbsBySource_internal(sourceId), - options.filter - ) + // EVERY requested source is honoured — an array used to collapse to + // its first element here, silently dropping the rest of the ask. + const onlySourceIds = Array.isArray(options.filter.sourceId) + ? options.filter.sourceId + : [options.filter.sourceId] + const sourceUnion: HNSWVerbWithMetadata[] = [] + const seenSourceVerbIds = new Set() + for (const oneSource of onlySourceIds) { + for (const v of await this.getVerbsBySource_internal(oneSource)) { + if (!seenSourceVerbIds.has(v.id)) { + seenSourceVerbIds.add(v.id) + sourceUnion.push(v) + } + } + } + const verbsBySource = this.applyVerbMetadataFilters(sourceUnion, options.filter) // Apply pagination const paginatedVerbs = verbsBySource.slice(offset, offset + limit) @@ -3023,16 +3043,22 @@ export abstract class BaseStorage extends BaseStorageAdapter { !options.filter.service && !options.filter.metadata ) { - const targetId = Array.isArray(options.filter.targetId) - ? options.filter.targetId[0] - : options.filter.targetId - - // Get verbs by target directly (hydrated with metadata), then apply the - // subtype / visibility metadata filters on the O(degree) candidate set. - const verbsByTarget = this.applyVerbMetadataFilters( - await this.getVerbsByTarget_internal(targetId), - options.filter - ) + // EVERY requested target is honoured — an array used to collapse to + // its first element here, silently dropping the rest of the ask. + const onlyTargetIds = Array.isArray(options.filter.targetId) + ? options.filter.targetId + : [options.filter.targetId] + const targetUnion: HNSWVerbWithMetadata[] = [] + const seenTargetVerbIds = new Set() + for (const oneTarget of onlyTargetIds) { + for (const v of await this.getVerbsByTarget_internal(oneTarget)) { + if (!seenTargetVerbIds.has(v.id)) { + seenTargetVerbIds.add(v.id) + targetUnion.push(v) + } + } + } + const verbsByTarget = this.applyVerbMetadataFilters(targetUnion, options.filter) // Apply pagination const paginatedVerbs = verbsByTarget.slice(offset, offset + limit) @@ -3061,16 +3087,25 @@ export abstract class BaseStorage extends BaseStorageAdapter { !options.filter.service && !options.filter.metadata ) { - const verbType = Array.isArray(options.filter.verbType) - ? options.filter.verbType[0] - : options.filter.verbType + // EVERY requested verb type is honoured — an array used to collapse to + // its first element here, silently dropping the rest of the ask. + const verbTypes = Array.isArray(options.filter.verbType) + ? options.filter.verbType + : [options.filter.verbType] - // Get verbs by type directly (hydrated with metadata), then apply the - // subtype / visibility metadata filters on the candidate set. - const verbsByType = this.applyVerbMetadataFilters( - await this.getVerbsByType_internal(verbType), - options.filter - ) + // Get verbs by each requested type (hydrated with metadata), deduped by + // id, then apply the subtype / visibility metadata filters on the set. + const byType: HNSWVerbWithMetadata[] = [] + const seenTypeVerbIds = new Set() + for (const oneType of verbTypes) { + for (const v of await this.getVerbsByType_internal(oneType)) { + if (!seenTypeVerbIds.has(v.id)) { + seenTypeVerbIds.add(v.id) + byType.push(v) + } + } + } + const verbsByType = this.applyVerbMetadataFilters(byType, options.filter) // Apply pagination const paginatedVerbs = verbsByType.slice(offset, offset + limit) diff --git a/src/transaction/operations/IndexOperations.ts b/src/transaction/operations/IndexOperations.ts index 1bbbca88..0142dc54 100644 --- a/src/transaction/operations/IndexOperations.ts +++ b/src/transaction/operations/IndexOperations.ts @@ -14,6 +14,7 @@ import type { MetadataIndexManager } from '../../utils/metadataIndex.js' import type { GraphVerb } from '../../coreTypes.js' import type { Operation, RollbackAction } from '../types.js' import { isZeroNormVector } from '../../utils/distance.js' +import { jsonSafeIndexMetadata } from '../../utils/jsonSafeIndexMetadata.js' import { prodLog } from '../../utils/logger.js' /** @@ -390,13 +391,21 @@ export class AddToMetadataIndexOperation implements Operation { // rollback so add + undo reference the same watermark. const generation = this.generationFn?.() - // Add to metadata index (skipFlush=true for transaction atomicity) - await this.index.addToIndex(this.id, this.entity, true, false, generation) + // The JSON-safe view is taken HERE, per crossing, never at construction: + // the entity reference this op holds can be mutated between plan and + // execute (a graph op's execute-time endpoint-int resolution mirrors + // BigInts onto a shared verb object) — see jsonSafeIndexMetadata's + // module doc. + await this.index.addToIndex( + this.id, jsonSafeIndexMetadata(this.entity), true, false, generation + ) // Return rollback action return async () => { // Remove from metadata index - await this.index.removeFromIndex(this.id, this.entity, generation) + await this.index.removeFromIndex( + this.id, jsonSafeIndexMetadata(this.entity), generation + ) } } } @@ -432,13 +441,21 @@ export class RemoveFromMetadataIndexOperation implements Operation { // Resolve the removal generation once; reuse it for the rollback re-add. const generation = this.generationFn?.() - // Remove from metadata index - await this.index.removeFromIndex(this.id, this.entity, generation) + // Sanitized per crossing, never at construction — transact()'s delete + // legs hand this op the SAME verb object the graph-retraction op's + // execute-time endpoint resolution mutates (BigInt sourceInt/targetInt), + // so a plan-time view aliases the pollution. See jsonSafeIndexMetadata's + // module doc. + await this.index.removeFromIndex( + this.id, jsonSafeIndexMetadata(this.entity), generation + ) // Return rollback action return async () => { // Re-add with original metadata (skipFlush=true) - await this.index.addToIndex(this.id, this.entity, true, false, generation) + await this.index.addToIndex( + this.id, jsonSafeIndexMetadata(this.entity), true, false, generation + ) } } } diff --git a/src/utils/jsonSafeIndexMetadata.ts b/src/utils/jsonSafeIndexMetadata.ts new file mode 100644 index 00000000..d3b1be5f --- /dev/null +++ b/src/utils/jsonSafeIndexMetadata.ts @@ -0,0 +1,47 @@ +/** + * @module utils/jsonSafeIndexMetadata + * @description The metadata-index crossing's JSON-safety law, as a leaf + * function both the coordinator and the transaction operations share. + * + * The seam's metadata is JSON-safe BY CONTRACT (a native provider serializes + * it; u64 ints as Number corrupt above 2^53) — but `resolveVerbEndpointInts` + * MIRRORS the resolved endpoint ints onto the verb object itself as BigInt + * (`verb.sourceInt`/`targetInt`), so a verb object reused as index metadata + * carries BigInts into JSON.stringify, which throws, aborting the whole + * transaction. Endpoint ints ride their OWN op params on the graph legs — the + * metadata crossing drops every BigInt-valued top-level key instead of + * guessing at a lossy numeric encoding. + * + * WHY THIS IS A LEAF MODULE, ENFORCED AT THE CROSSING: sanitizing only at + * operation-construction time is not enough. `transact()`'s delete legs pass + * the SAME verb object to both the graph-retraction op (whose endpoint-int + * thunk deliberately resolves at EXECUTE time, for same-batch forward refs) + * and the metadata-retraction op. At plan time the verb is still clean, so a + * plan-time sanitize returns the same reference — then the graph op executes + * first, mirrors the BigInt ints onto the shared object, and the metadata op + * crosses the seam with them (found by the first fleet adoption of the native + * pair: every transact-wrapped edge delete aborted). The crossing itself is + * the only place ordering cannot bypass. + */ + +/** + * A JSON-safe view of a record bound for the metadata-index crossing. + * + * @param metadata - The candidate index-metadata record. + * @returns The same object when already JSON-safe, else a shallow copy + * without the BigInt-valued keys. + */ +export function jsonSafeIndexMetadata(metadata: unknown): unknown { + if (metadata === null || typeof metadata !== 'object') return metadata + const rec = metadata as Record + let hasBigint = false + for (const k in rec) { + if (typeof rec[k] === 'bigint') { hasBigint = true; break } + } + if (!hasBigint) return metadata + const out: Record = {} + for (const k in rec) { + if (typeof rec[k] !== 'bigint') out[k] = rec[k] + } + return out +} diff --git a/src/utils/metadataIndex.ts b/src/utils/metadataIndex.ts index 3e0e3d17..0fd312e2 100644 --- a/src/utils/metadataIndex.ts +++ b/src/utils/metadataIndex.ts @@ -2575,6 +2575,19 @@ export class MetadataIndexManager implements MetadataIndexProvider { /** Once-per-field flag for the fallback-degradation announcement. */ private static announcedFallbackSorts = new Set() + /** + * Evaluate `filter` over `ids` only — the graph-first find's door (the + * neighbour set filtered by id, never the store filtered and then + * intersected). This index answers from its own `getIdsForFilter`, so the + * two doors cannot disagree; the cost is that of the filter over this + * in-memory index, and the answer keeps the caller's order. + */ + async filterIdsWithin(filter: any, ids: readonly string[]): Promise { + if (ids.length === 0) return [] + const matched = new Set(await this.getIdsForFilter(filter)) + return ids.filter((id) => matched.has(id)) + } + async getSortedIdsForFilter( filter: any, orderBy: string, diff --git a/tests/integration/counts-persist-single-flight.test.ts b/tests/integration/counts-persist-single-flight.test.ts new file mode 100644 index 00000000..5acbdcc3 --- /dev/null +++ b/tests/integration/counts-persist-single-flight.test.ts @@ -0,0 +1,111 @@ +/** + * @module tests/integration/counts-persist-single-flight + * @description Regression for a production race in FileSystemStorage's + * counts ledger: `persistCounts()` was write-through on every count change + * with no serialization, and the atomic writer named its temp file with + * millisecond granularity (`.tmp--`). Two persists inside one + * millisecond shared the temp path — both wrote it, the first rename + * consumed it, the second rename found nothing: ENOENT, ~1,500 times a day + * on a busy production brain, with a full ledger write per change behind it. + * + * Under pin: persists are single-flight and coalesced — one in flight, at + * most one trailing pass carrying the burst's final state — and every atomic + * write owns a unique temp path. A burst of N count changes costs at most + * two ledger writes, never errors, and leaves a ledger equal to memory. + */ +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest' +import * as fs from 'node:fs' +import * as os from 'node:os' +import * as path from 'node:path' +import { Brainy } from '../../src/brainy.js' +import { NounType } from '../../src/types/graphTypes.js' + +describe('counts persistence is single-flight, coalesced, and never races its own temp file', () => { + let dir: string + let brain: any + + beforeEach(async () => { + process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true' + dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-counts-race-')) + brain = new Brainy({ + requireSubtype: false, + storage: { type: 'filesystem', path: dir }, + dimensions: 384, + silent: true + }) + await brain.init() + }) + + afterEach(async () => { + vi.restoreAllMocks() + await brain.close() + fs.rmSync(dir, { recursive: true, force: true }) + }) + + it('a burst of concurrent count changes → at most two ledger writes, zero errors, ledger == memory', async () => { + const storage = brain.storage + const countsPath: string = storage.countsFilePath + expect(countsPath, 'the filesystem adapter persists a counts ledger').toBeTruthy() + + // Let init's own persists settle so the burst is measured alone. + await storage.flushCounts?.() + + const renameSpy = vi.spyOn(fs.promises, 'rename') + const errorSpy = vi.spyOn(console, 'error') + + // Twenty-five concurrent count changes — the shape of a write burst; each + // used to launch its own persist. + const BURST = 25 + await Promise.all( + Array.from({ length: BURST }, () => storage.scheduleCountPersist()) + ) + + const ledgerRenames = renameSpy.mock.calls.filter(([, to]) => String(to) === countsPath) + expect(ledgerRenames.length, 'single-flight + one trailing pass').toBeLessThanOrEqual(2) + expect(ledgerRenames.length, 'the burst was persisted at all').toBeGreaterThanOrEqual(1) + + const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts')) + expect(persistErrors).toEqual([]) + + const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8')) + expect(ledger.totalNounCount).toBe(storage.totalNounCount) + expect(ledger.totalVerbCount).toBe(storage.totalVerbCount) + }) + + it('real writes in parallel: the ledger lands complete and no persist error is logged', async () => { + const storage = brain.storage + const countsPath: string = storage.countsFilePath + const errorSpy = vi.spyOn(console, 'error') + + await Promise.all( + Array.from({ length: 12 }, (_, i) => + brain.add({ data: `burst row ${i}`, type: NounType.Thing }) + ) + ) + await storage.flushCounts?.() + + const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts')) + expect(persistErrors).toEqual([]) + const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8')) + expect(ledger.totalNounCount).toBe(storage.totalNounCount) + expect(await brain.getNounCount()).toBe(ledger.totalNounCount) + }) + + it('every atomic write owns its own temp path — two writes in one millisecond never collide', async () => { + const storage = brain.storage + const tmpNames: string[] = [] + vi.spyOn(fs.promises, 'writeFile').mockImplementation(async (p: any) => { + tmpNames.push(String(p)) + }) + vi.spyOn(fs.promises, 'rename').mockImplementation(async () => undefined) + const target = path.join(dir, 'probe.json') + await Promise.all([ + storage.writeFileAtomic(target, '{"a":1}'), + storage.writeFileAtomic(target, '{"a":2}'), + storage.writeFileAtomic(target, '{"a":3}') + ]) + const probeTmps = tmpNames.filter((n) => n.startsWith(`${target}.tmp-`)) + expect(probeTmps.length).toBe(3) + expect(new Set(probeTmps).size, 'no two writes shared a temp path').toBe(3) + }) +}) diff --git a/tests/integration/entity-tree-stamp.test.ts b/tests/integration/entity-tree-stamp.test.ts index deefc5e6..23cc0a15 100644 --- a/tests/integration/entity-tree-stamp.test.ts +++ b/tests/integration/entity-tree-stamp.test.ts @@ -57,7 +57,11 @@ describe('entity-tree family stamp', () => { const invariants = (stamp.members as any).invariants expect(invariants.nounCount).toBe(await brain.storage.getNounCount()) expect(invariants.verbCount).toBe(await brain.storage.getVerbCount()) - expect(stamp.sourceGeneration).toBe(brain.generation()) + // THE SOURCE IS COMMITTED TRUTH, never the allocated counter. Stamping the + // counter labelled the stamp with a generation a write in flight had merely + // claimed, so every crash inside a write window produced a spurious verdict + // at the next open (see the torn-tail pins below). + expect(stamp.sourceGeneration).toBe(brain.generationStore.committedGeneration()) expect(stamp.generation).toBeGreaterThanOrEqual(1) }) @@ -112,6 +116,96 @@ describe('entity-tree family stamp', () => { expect(stillIncoherent).toEqual([]) }) + /** + * Rewrite the on-disk stamp so its `sourceGeneration` sits ABOVE the store's + * committed watermark — the durable shape a torn generation-log tail leaves + * behind (the stamp's fsync outlived the tail's). Fabricated rather than + * crash-produced so the pin is deterministic; the seeded-SIGKILL lane + * (`scripts/crash-consistency.mjs` in the engine repo) produces the same + * shape from a real abrupt termination. + */ + const fabricateTear = (ahead: number): FamilyStamp => { + const file = path.join(dir, `${ENTITY_TREE_STAMP_PATH}.gz`) + const zlib = require('node:zlib') + const raw = JSON.parse(zlib.gunzipSync(fs.readFileSync(file)).toString('utf-8')) as FamilyStamp + const torn: FamilyStamp = { ...raw, sourceGeneration: raw.sourceGeneration + ahead } + fs.writeFileSync(file, zlib.gzipSync(JSON.stringify(torn))) + return torn + } + + it('a torn generation-log tail is a TERMINAL VERDICT at open: narrated, demoted, never a wait', async () => { + for (let i = 0; i < 3; i++) + await brain.add({ data: `torn${i}`, type: 'document', metadata: { i } }) + await brain.close() + const torn = fabricateTear(5) + + const warn = vi.spyOn(prodLog, 'warn') + const startedAt = Date.now() + brain = await open() + const openMs = Date.now() - startedAt + + const tearLines = warn.mock.calls.filter((c) => String(c[0]).includes('TORN GENERATION-LOG TAIL')) + expect(tearLines.length).toBe(1) + const said = String(tearLines[0][0]) + // Narrated PRECISELY: both generations, the file, and the named cure. + expect(said).toContain(`source generation ${torn.sourceGeneration}`) + expect(said).toContain(`committed generation ${brain.generationStore.committedGeneration()}`) + expect(said).toContain(ENTITY_TREE_STAMP_PATH) + expect(said).toContain('DEMOTED') + expect(said).toMatch(/repairIndex\(\)/) + // Terminal, not a wait: the demotion is O(1) straight-line work, so a tear + // cannot turn an open into the 8-minute spin this class was reported as. + expect(openMs).toBeLessThan(30_000) + + // The store SERVES — a tear in a stamp never locks an owner out of the + // canonical tree the stamp merely describes. + expect((await brain.find({ type: 'document', limit: 100 })).length).toBe(3) + + // The demotion CONVERGED: the stamp now names committed truth, and the + // next open is quiet. A verdict that re-narrates every open is a wait + // wearing a different hat. + const restamped = (await readFamilyStamp(brain.storage, ENTITY_TREE_STAMP_PATH)) as FamilyStamp + expect(restamped.sourceGeneration).toBe(brain.generationStore.committedGeneration()) + await brain.close() + const warn2 = vi.spyOn(prodLog, 'warn') + brain = await open() + expect(warn2.mock.calls.filter((c) => String(c[0]).includes('TORN'))).toEqual([]) + }) + + it('a READ-ONLY open on a torn tail refuses to guess: terminal verdict + named cure, no re-stamp', async () => { + await brain.add({ data: 'ro', type: 'document', metadata: {} }) + await brain.close() + const torn = fabricateTear(3) + + const warn = vi.spyOn(prodLog, 'warn') + const reader: any = await Brainy.openReadOnly({ + requireSubtype: false, + storage: { type: 'filesystem', path: dir }, + silent: true, + dimensions: 384 + }) + const tearLines = warn.mock.calls.filter((c) => String(c[0]).includes('TORN GENERATION-LOG TAIL')) + expect(tearLines.length).toBe(1) + const said = String(tearLines[0][0]) + expect(said).toContain('READ-ONLY') + expect(said).toContain('UNVERIFIED') + expect(said).toMatch(/repairIndex\(\)/) + await reader.close() + + // A reader never rewrites the store: read the bytes back off disk (not + // through a writer open, which would demote them) — the torn stamp is + // exactly as it was found. + const onDisk = JSON.parse( + require('node:zlib') + .gunzipSync(fs.readFileSync(path.join(dir, `${ENTITY_TREE_STAMP_PATH}.gz`))) + .toString('utf-8') + ) as FamilyStamp + expect(onDisk.sourceGeneration).toBe(torn.sourceGeneration) + expect(onDisk.generation).toBe(torn.generation) + + brain = await open() + }) + it('the one verifier handles both member modes', () => { const rollup: FamilyStamp = { family: 'x', @@ -127,7 +221,13 @@ describe('entity-tree family stamp', () => { stampSource: 5, head: 9 }) - expect(verifyFamilyStamp(rollup, 3, { nounCount: 10 }).state).toBe('incoherent') // ahead of head + // AHEAD is its own class — a torn generation-log tail, never folded in + // with `incoherent`: the two have opposite cures (recount vs demote). + expect(verifyFamilyStamp(rollup, 3, { nounCount: 10 })).toEqual({ + state: 'torn', + stampSource: 5, + head: 3 + }) expect(verifyFamilyStamp(null, 5, {})).toEqual({ state: 'absent' }) const enumerated: FamilyStamp = { diff --git a/tests/integration/find-connected-order.test.ts b/tests/integration/find-connected-order.test.ts new file mode 100644 index 00000000..b04e7f99 --- /dev/null +++ b/tests/integration/find-connected-order.test.ts @@ -0,0 +1,165 @@ +/** + * @module tests/integration/find-connected-order + * @description The graph-first law for `find({ connected })` (10.4.8). + * + * With `connected` present the neighbour set is the candidate universe: it is + * resolved from the adjacency first, the metadata filter is evaluated over + * those ids only, and the page is cut last. The earlier order materialized the + * whole-store filtered id list, paged it, hydrated the page, and only then + * intersected with the neighbours — so a neighbour outside the first page of + * the filtered STORE was silently dropped, and every call paid O(store). + * + * These pins hold both halves. The answer: every matching neighbour is + * reachable by paging, a non-neighbour never appears, a negation (`missing`) + * is evaluated over the neighbours, `orderBy` sorts the whole neighbour set + * before the page is cut, and the vector leg walks the neighbours only. The + * cost shape: the metadata index is asked about the neighbour ids only, and + * hydration is one page — never the store. + */ +import { describe, it, expect, beforeAll, afterAll, vi } from 'vitest' +import { Brainy } from '../../src/brainy' +import { NounType, VerbType } from '../../src/types/graphTypes' +import { v5 } from '../../src/universal/uuid' +import { generateTestVector } from '../helpers/test-factory' + +/** Matching rows that are NOT neighbours — added FIRST, so the whole-store filtered list leads with them. */ +const NOISE = 120 +/** Matching rows that ARE neighbours of the anchor. */ +const NEIGHBOURS = 30 +/** Neighbours carrying `retracted: true` — excluded by the `missing` negation. */ +const RETRACTED = 4 + +describe('find({ connected }) is graph-first: neighbours → filter → page', () => { + let brain: Brainy + const anchor = 'anchor' + const sharedVector = generateTestVector() + const neighbourIds = new Set(Array.from({ length: NEIGHBOURS }, (_, i) => v5(`nb-${i}`))) + + beforeAll(async () => { + brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } }) + await brain.init() + await brain.add({ + id: anchor, + data: 'the anchor', + type: NounType.Person, + metadata: { kind: 'anchor' }, + vector: generateTestVector() + }) + for (let i = 0; i < NOISE; i++) { + await brain.add({ + id: `noise-${i}`, + data: `noise ${i}`, + type: NounType.Person, + metadata: { kind: 'note', rank: 1000 + i }, + vector: sharedVector + }) + } + for (let i = 0; i < NEIGHBOURS; i++) { + await brain.add({ + id: `nb-${i}`, + data: `neighbour ${i}`, + type: NounType.Person, + metadata: { kind: 'note', rank: i + 1, ...(i < RETRACTED ? { retracted: true } : {}) }, + vector: sharedVector + }) + await brain.relate({ from: anchor, to: `nb-${i}`, type: VerbType.Knows }) + } + }) + + afterAll(async () => { + brain = null as any + }) + + it('returns the matching neighbours page by page — none dropped, never a non-neighbour', async () => { + const seen = new Set() + for (let offset = 0; offset <= NEIGHBOURS; offset += 10) { + const page = await brain.find({ + connected: { from: anchor, direction: 'out' }, + where: { kind: 'note' }, + limit: 10, + offset + }) + expect(page).toHaveLength(offset < NEIGHBOURS ? 10 : 0) + for (const r of page) { + expect(neighbourIds.has(r.entity.id)).toBe(true) + expect(seen.has(r.entity.id)).toBe(false) + seen.add(r.entity.id) + } + } + expect(seen.size).toBe(NEIGHBOURS) + }) + + it('evaluates a negation (`missing`) over the neighbour set, not the store', async () => { + const results = await brain.find({ + connected: { from: anchor, direction: 'out' }, + where: { kind: 'note', retracted: { missing: true } }, + limit: 100 + }) + expect(results).toHaveLength(NEIGHBOURS - RETRACTED) + for (const r of results) { + expect(neighbourIds.has(r.entity.id)).toBe(true) + expect(r.entity.metadata.retracted).toBeUndefined() + } + }) + + it('asks the metadata index about the neighbour ids only, and hydrates one page', async () => { + const index = (brain as any).metadataIndex + const within = vi.spyOn(index, 'filterIdsWithin') + const hydrate = vi.spyOn(brain as any, 'batchGet') + try { + const results = await brain.find({ + connected: { from: anchor, direction: 'out' }, + where: { kind: 'note' }, + limit: 10 + }) + expect(results).toHaveLength(10) + expect(within).toHaveBeenCalledTimes(1) + const askedIds = within.mock.calls[0][1] as string[] + expect(askedIds).toHaveLength(NEIGHBOURS) + for (const id of askedIds) expect(neighbourIds.has(id)).toBe(true) + expect(hydrate).toHaveBeenCalledTimes(1) + expect(hydrate.mock.calls[0][0]).toHaveLength(10) + } finally { + within.mockRestore() + hydrate.mockRestore() + } + }) + + it('orders the WHOLE neighbour set before cutting the page', async () => { + const results = await brain.find({ + connected: { from: anchor, direction: 'out' }, + where: { kind: 'note' }, + orderBy: 'rank', + order: 'desc', + limit: 5 + }) + expect(results.map((r) => r.entity.metadata.rank)).toEqual([30, 29, 28, 27, 26]) + }) + + it('walks the vector leg over the neighbours only', async () => { + const results = await brain.find({ + vector: sharedVector, + connected: { from: anchor, direction: 'out' }, + where: { kind: 'note' }, + limit: 5 + }) + expect(results).toHaveLength(5) + for (const r of results) expect(neighbourIds.has(r.entity.id)).toBe(true) + }) + + it('an anchor without neighbours answers [] before the filter is asked', async () => { + const index = (brain as any).metadataIndex + const within = vi.spyOn(index, 'filterIdsWithin') + try { + const results = await brain.find({ + connected: { from: 'noise-0', direction: 'out' }, + where: { kind: 'note' }, + limit: 10 + }) + expect(results).toEqual([]) + expect(within).not.toHaveBeenCalled() + } finally { + within.mockRestore() + } + }) +}) diff --git a/tests/integration/history-repacking.test.ts b/tests/integration/history-repacking.test.ts index 2bcee038..bb07268d 100644 --- a/tests/integration/history-repacking.test.ts +++ b/tests/integration/history-repacking.test.ts @@ -16,6 +16,7 @@ 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 * as zlib from 'node:zlib' import { Brainy } from '../../src/brainy.js' import { NounType } from '../../src/types/graphTypes.js' import { GenerationStore } from '../../src/db/generationStore.js' @@ -57,6 +58,107 @@ describe('history repacking — the two-tier lifecycle', () => { } }) + /** + * THE HOLE, END TO END — the shape a real store carries. + * + * A forensic fixture was measured with generation directories 1..2503 + * present except for exactly one: 1416. Its fact-log segment already showed + * the tell — `seg-...1410.bfl` declaring firstGeneration 1410, lastGeneration + * 1940 (531 generations) while recording only 530 facts. + * + * Before the fix, repacking such a store folded ACROSS that hole: the batch + * skipped 1416 (no readable delta) and the sealed segment declared a range + * spanning it anyway. The next open merged that declared range back into + * committedRanges, re-admitting 1416 as committed history, and every + * subsequent auto-compaction pass then asked the packed tier for a frame + * that was never written — producing, on EVERY run, the non-fatal narration + * + * Auto-compaction of generational history failed (non-fatal): generation + * N is inside sealed segment seg-....bgs's declared range but has no frame + * — packed history is damaged + * + * This pin removes a generation directory to make the same hole, then + * requires repack + reopen + compaction to complete cleanly. + */ + it('a missing generation directory does not poison the packed tier', async () => { + const dir = tempDir() + // `retention: 'all'` throughout: close() otherwise auto-compacts the + // history away, and this pin needs the cold generations still on disk so + // there is something to punch a hole in. The live window stays at its + // production default for the build phase, so nothing folds yet. + const archival = async (): Promise => { + const b = new Brainy({ + requireSubtype: false, + storage: { type: 'filesystem', path: dir }, + embeddingFunction: stub, + retention: 'all' + }) + await b.init() + return b + } + const brain = await archival() + + const id = await brain.add({ + data: 'holed-entity', + type: NounType.Document, + metadata: { v: 0 } + }) + // One flush per update: single-op writes coalesce inside a flush window, + // so a history deep enough to have a middle needs the windows separated. + for (let v = 1; v <= 12; v++) { + await brain.update({ id, metadata: { v } }) + await brain.flush() + } + await brain.close() + + // Punch the hole: delete ONE generation directory in the middle of the + // cold range, exactly as the real store presents it. + const genRoot = path.join(dir, '_generations') + const numeric = fs + .readdirSync(genRoot, { withFileTypes: true }) + .filter((e) => e.isDirectory() && /^\d+$/.test(e.name)) + .map((e) => Number(e.name)) + .sort((a, b) => a - b) + expect(numeric.length).toBeGreaterThan(6) + const victim = numeric[Math.floor(numeric.length / 2)] + fs.rmSync(path.join(genRoot, String(victim)), { recursive: true, force: true }) + + // Now shrink the live window and reopen. close() repacks automatically + // (brainy.ts phase 0b), so this is the production sequence exactly: a + // store with a hole in its history gets folded by ordinary housekeeping, + // with nobody asking for it. + ;(GenerationStore as any).REPACK_LIVE_WINDOW = 3 + const reopened = await archival() + const result = await reopened.repackHistory() + expect(result.foldedGenerations).toBeGreaterThan(0) + + const segDir = path.join(dir, SEGMENTS_PREFIX) + const manifestPath = ['manifest.json', 'manifest.json.gz'] + .map((f) => path.join(segDir, f)) + .find((p) => fs.existsSync(p))! + const raw = manifestPath.endsWith('.gz') + ? zlib.gunzipSync(fs.readFileSync(manifestPath)).toString('utf8') + : fs.readFileSync(manifestPath, 'utf8') + const manifest = JSON.parse(raw) as { + segments: Array<{ firstGeneration: number; lastGeneration: number; frames: number }> + } + + // THE LAW: every sealed segment declares exactly as many generations as it + // holds frames, and none of them spans the victim. + for (const s of manifest.segments) { + expect(s.lastGeneration - s.firstGeneration + 1).toBe(s.frames) + expect(victim >= s.firstGeneration && victim <= s.lastGeneration).toBe(false) + } + + await reopened.close() + + // And the pass that used to fail on every run now completes: reopen (which + // re-seeds committedRanges from the packed tier) then compact history. + const third = await openBrain(dir) + await expect(third.compactHistory({ maxGenerations: 2 })).resolves.toBeDefined() + await third.close() + }) + it('repack preserves every historical read across cold reopen; folded dirs are gone', async () => { ;(GenerationStore as any).REPACK_LIVE_WINDOW = 3 const dir = tempDir() diff --git a/tests/integration/pending-embed-low-water.test.ts b/tests/integration/pending-embed-low-water.test.ts new file mode 100644 index 00000000..f966d0a1 --- /dev/null +++ b/tests/integration/pending-embed-low-water.test.ts @@ -0,0 +1,141 @@ +/** + * @module tests/integration/pending-embed-low-water + * @description The pending-embed recovery fold is bounded and background (10.4.9). + * + * The fold used to scan the generation log from generation 1 at EVERY open, + * on the open's foreground — O(whole history) per open on long-lived brains. + * Now: an advisory low-water mark (`_system/pending_embeds_lowwater.json`) + * records the committed generation whenever the pending set drains to empty, + * recovery scans from `mark + 1` on the open's foreground — the crash-recovery + * contract keeps markers re-armed when open() returns. The mark is advisory: stale-low costs a longer scan, never a + * marker — a pending embed enqueued before a crash is still recovered. + */ +import { describe, it, expect, afterEach, vi } from 'vitest' +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { Brainy } from '../../src/brainy' +import { NounType } from '../../src/types/graphTypes' + +const LOWWATER_PATH = '_system/pending_embeds_lowwater.json' + +describe('pending-embed recovery: bounded by the low-water mark', () => { + const roots: string[] = [] + const dir = (): string => { + const d = mkdtempSync(join(tmpdir(), 'brainy-lowwater-')) + roots.push(d) + return d + } + const open = async (root: string): Promise> => { + const brain = new Brainy({ + requireSubtype: false, + storage: { type: 'filesystem', path: root } + }) + await brain.init() + return brain + } + + afterEach(() => { + for (const d of roots.splice(0)) rmSync(d, { recursive: true, force: true }) + }) + + it('drain-to-empty writes the mark, and the next open scans from mark + 1', async () => { + const root = dir() + const brain = await open(root) + // Hold the worker so the pending state is observable, then release it. + const realKick = (brain as any).kickEmbedWorker.bind(brain) + ;(brain as any).kickEmbedWorker = () => {} + await brain.add({ + id: 'row-1', + data: 'the first deferred row', + type: NounType.Thing, + deferEmbedding: true + }) + expect(brain.pendingEmbedCount()).toBeGreaterThan(0) + ;(brain as any).kickEmbedWorker = realKick + await brain.awaitPendingEmbeds() + // The drain wrote the advisory mark (fire-and-forget: settle the microtask). + await new Promise((r) => setTimeout(r, 50)) + const mark = (await (brain as any).storage.readRawObject(LOWWATER_PATH)) as { + generation: number + } | null + expect(mark).not.toBeNull() + expect(mark!.generation).toBeGreaterThan(0) + await brain.close() + + const brain2 = await open(root) + const log = (brain2 as any).generationStore.getFactLog() + const scanSpy = vi.spyOn(log, 'scanFacts') + try { + await (brain2 as any).recoverPendingEmbedsFromLog() + expect(scanSpy).toHaveBeenCalledTimes(1) + const opts = scanSpy.mock.calls[0][0] as { fromGeneration?: number } + expect(opts.fromGeneration).toBeGreaterThanOrEqual(mark!.generation + 1) + } finally { + scanSpy.mockRestore() + await brain2.close() + } + }) + + it('a pending embed enqueued after the mark survives an unclean stop', async () => { + const root = dir() + const brain = await open(root) + await brain.add({ id: 'settled', data: 'lands before the mark', type: NounType.Thing }) + await brain.awaitPendingEmbeds() + await new Promise((r) => setTimeout(r, 50)) + + // A deferred write whose embed never lands: block the worker, then drop + // the instance without close() — the unclean-stop shape. + ;(brain as any).kickEmbedWorker = () => {} + await brain.add({ + id: 'orphan', + data: 'enqueued then abandoned', + type: NounType.Thing, + deferEmbedding: true + }) + expect(brain.pendingEmbedCount()).toBeGreaterThan(0) + // No close(): simulate the crash by releasing only the writer lock so the + // next open can proceed. + await (brain as any).storage.releaseWriterLock() + + const brain2 = await open(root) + expect(brain2.pendingEmbedCount()).toBeGreaterThan(0) + await brain2.awaitPendingEmbeds() + expect(brain2.pendingEmbedCount()).toBe(0) + await brain2.close() + // Reap the crashed instance: its fence is gone, so close() fails loudly — + // swallow that here; the point is clearing its watchers and registry entry. + await brain.close().catch(() => undefined) + }) + + it('a reopened brain has its pending set settled when open() returns', async () => { + const root = dir() + const brain = await open(root) + await brain.add({ id: 'a-row', data: 'some data', type: NounType.Thing }) + await brain.awaitPendingEmbeds() + await brain.close() + + const brain2 = await open(root) + // The crash-recovery contract: markers are re-armed by open itself — + // no latch, no background race. (Here the drain landed, so zero.) + expect(brain2.pendingEmbedCount()).toBe(0) + await brain2.close() + }) + + it('a clean close with an empty set writes the mark even if no drain happened', async () => { + const root = dir() + const brain = await open(root) + await brain.add({ id: 'r1', data: 'row one', type: NounType.Thing }) + await brain.awaitPendingEmbeds() + await brain.close() + // Read the mark back through the storage door (the adapter owns the + // on-disk encoding), on a fresh instance. + const brain2 = await open(root) + const mark = (await (brain2 as any).storage.readRawObject(LOWWATER_PATH)) as { + generation: number + } | null + expect(mark).not.toBeNull() + expect(mark!.generation).toBeGreaterThan(0) + await brain2.close() + }) +}) diff --git a/tests/integration/related-verb-array.test.ts b/tests/integration/related-verb-array.test.ts new file mode 100644 index 00000000..36a49850 --- /dev/null +++ b/tests/integration/related-verb-array.test.ts @@ -0,0 +1,89 @@ +/** + * @module tests/integration/related-verb-array + * @description related() honours EVERY verb type in an array (10.4.9). + * + * The storage fast paths for `sourceId + verbType` and `verbType` collapsed a + * verb-type ARRAY to its first element — `related({ from, type: [a, b] })` + * silently returned only `a` edges, whichever order the array came in. The + * same quiet-loss class as the graph-first paging defect, one seam over. + * These pins seed a store where the SECOND requested type's edge must come + * back, on every path the collapse lived in. + */ +import { describe, it, expect, beforeAll, afterAll } from 'vitest' +import { Brainy } from '../../src/brainy' +import { NounType, VerbType } from '../../src/types/graphTypes' +import { v5 } from '../../src/universal/uuid' + +describe('related() with a verb-type array returns every requested type', () => { + let brain: Brainy + + beforeAll(async () => { + brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } }) + await brain.init() + for (const id of ['a', 'b', 'c', 'd']) { + await brain.add({ id, data: `node ${id}`, type: NounType.Person }) + } + await brain.relate({ from: 'a', to: 'b', type: VerbType.Supports }) + await brain.relate({ from: 'a', to: 'c', type: VerbType.RelatedTo }) + await brain.relate({ from: 'a', to: 'd', type: VerbType.Knows }) + await brain.relate({ from: 'b', to: 'c', type: VerbType.RelatedTo }) + }) + + afterAll(async () => { + brain = null as any + }) + + it('from + type array: the second type\'s edge comes back, both orders', async () => { + for (const types of [ + [VerbType.Supports, VerbType.RelatedTo], + [VerbType.RelatedTo, VerbType.Supports] + ]) { + const edges = await brain.related({ from: 'a', type: types }) + const targets = new Set(edges.map((e) => e.to)) + expect(targets.has(v5('b')), `types [${types}] missing Supports edge`).toBe(true) + expect(targets.has(v5('c')), `types [${types}] missing RelatedTo edge`).toBe(true) + expect(targets.has(v5('d'))).toBe(false) + expect(edges).toHaveLength(2) + } + }) + + it('a single-element array behaves exactly like the scalar', async () => { + const scalar = await brain.related({ from: 'a', type: VerbType.Supports }) + const array = await brain.related({ from: 'a', type: [VerbType.Supports] }) + expect(array.map((e) => e.id).sort()).toEqual(scalar.map((e) => e.id).sort()) + expect(array).toHaveLength(1) + }) + + it('no duplicate edges when types overlap the same edge set', async () => { + const edges = await brain.related({ + from: 'a', + type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows] + }) + const ids = edges.map((e) => e.id) + expect(new Set(ids).size).toBe(ids.length) + expect(edges).toHaveLength(3) + }) + + it('type-only asks (no anchor) honour the whole array too', async () => { + const edges = await brain.related({ type: [VerbType.Supports, VerbType.Knows] }) + const verbs = new Set(edges.map((e) => e.type)) + expect(verbs.has(VerbType.Supports)).toBe(true) + expect(verbs.has(VerbType.Knows)).toBe(true) + expect(edges).toHaveLength(2) + }) + + it('to + type array: the target side honours every type too', async () => { + const edges = await brain.related({ to: 'c', type: [VerbType.RelatedTo, VerbType.Supports] }) + const froms = new Set(edges.map((e) => e.from)) + expect(froms.has(v5('a'))).toBe(true) + expect(froms.has(v5('b'))).toBe(true) + expect(edges).toHaveLength(2) + }) + + it('pagination stays consistent across the union', async () => { + const page1 = await brain.related({ from: 'a', type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows], limit: 2 }) + const page2 = await brain.related({ from: 'a', type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows], limit: 2, offset: 2 }) + const all = [...page1, ...page2].map((e) => e.id) + expect(new Set(all).size).toBe(3) + }) +}) diff --git a/tests/integration/transact-edge-delete-bigint-aliasing.test.ts b/tests/integration/transact-edge-delete-bigint-aliasing.test.ts new file mode 100644 index 00000000..b3902571 --- /dev/null +++ b/tests/integration/transact-edge-delete-bigint-aliasing.test.ts @@ -0,0 +1,184 @@ +/** + * @module tests/integration/transact-edge-delete-bigint-aliasing + * @description Regression for a fleet-adoption blocker: ANY edge delete + * inside `transact()` — a direct unrelate or a noun-remove's cascade — + * aborted with the metadata seam's BigInt JSON-guard error on a strict + * (native) metadata provider. + * + * The aliasing chain: `planTxUnrelate`/the remove-cascade pass the SAME verb + * object to the graph-retraction op and the metadata-retraction op. The + * metadata leg's JSON-safe wrap ran at PLAN time, when the verb was still + * clean — so it returned the same reference. At EXECUTE time the graph op + * runs first and `resolveVerbEndpointInts` mirrors BigInt + * `sourceInt`/`targetInt` onto the shared object (deliberately deferred for + * same-batch forward refs — see transact-forward-ref-graph.test.ts); the + * metadata op then crossed the seam with the polluted object. Direct + * `unrelate()` resolves ints at BUILD time, before its sanitize, which is why + * only the transact() shapes ever hit it. + * + * Fix under pin: the JSON-safe view is taken AT THE CROSSING — inside the + * metadata-index operations' execute/rollback — so no plan-vs-execute + * ordering can bypass it. The JS baseline index tolerates BigInts (it would + * mask the bug), so these pins SPY on the seam and assert what actually + * crossed, exactly as a strict native provider would judge it. + */ +import { describe, it, expect, beforeEach, afterEach } from 'vitest' +import * as fs from 'node:fs' +import * as os from 'node:os' +import * as path from 'node:path' +import { Brainy } from '../../src/brainy.js' +import { NounType, VerbType } from '../../src/types/graphTypes.js' +import { + AddToMetadataIndexOperation, + RemoveFromMetadataIndexOperation +} from '../../src/transaction/operations/index.js' + +let seq = 0 +const freshId = (): string => + `00000000-0000-4000-8000-${(++seq).toString(16).padStart(12, '0')}` + +/** Top-level BigInt-valued keys of a candidate seam crossing (the guard's law). */ +const bigintKeys = (metadata: unknown): string[] => { + if (metadata === null || typeof metadata !== 'object') return [] + return Object.entries(metadata as Record) + .filter(([, v]) => typeof v === 'bigint') + .map(([k]) => k) +} + +describe('transact() edge deletes never carry BigInt across the metadata seam', () => { + let dir: string + let brain: any + let crossings: Array<{ door: string; id: string; keys: string[] }> + + beforeEach(async () => { + process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true' + dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-tx-bigint-')) + brain = new Brainy({ + requireSubtype: false, + storage: { type: 'filesystem', path: dir }, + dimensions: 384, + silent: true + }) + await brain.init() + + // Spy on the seam the way a strict native provider judges it: record the + // BigInt-valued top-level keys of every metadata argument that crosses. + // The JS baseline index tolerates BigInts, so without this the baseline + // run would green a shape the native pair aborts on. + crossings = [] + const index = brain.metadataIndex + for (const door of ['addToIndex', 'removeFromIndex'] as const) { + const real = index[door].bind(index) + index[door] = (id: string, metadata: unknown, ...rest: unknown[]) => { + crossings.push({ door, id, keys: bigintKeys(metadata) }) + return real(id, metadata, ...rest) + } + } + }) + + afterEach(async () => { + await brain.close() + fs.rmSync(dir, { recursive: true, force: true }) + }) + + it('CASE 1 (the fleet repro): relate, then transact([{op: unrelate}])', async () => { + const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing }) + const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing }) + const verbId = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo }) + + crossings.length = 0 + await brain.transact([{ op: 'unrelate', id: verbId }]) + + const polluted = crossings.filter((c) => c.keys.length > 0) + expect(polluted).toEqual([]) + expect(await brain.storage.getVerb(verbId)).toBeFalsy() + }) + + it('CASE 2 (the cascade shape): transact([{op: remove}]) cascading edge deletes', async () => { + const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing }) + const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing }) + const c = await brain.add({ id: freshId(), data: 'c', type: NounType.Thing }) + const ab = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo }) + const ca = await brain.relate({ from: c, to: a, type: VerbType.RelatedTo }) + + crossings.length = 0 + await brain.transact([{ op: 'remove', id: a }]) + + const polluted = crossings.filter((c2) => c2.keys.length > 0) + expect(polluted).toEqual([]) + expect(await brain.get(a)).toBeFalsy() + expect(await brain.storage.getVerb(ab)).toBeFalsy() + expect(await brain.storage.getVerb(ca)).toBeFalsy() + }) + + it('CASE 3 (one batch, both legs): adds + relate + unrelate of a pre-existing edge', async () => { + const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing }) + const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing }) + const old = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo }) + + const x = freshId() + crossings.length = 0 + await brain.transact([ + { op: 'add', id: x, data: 'x', type: NounType.Thing }, + { op: 'relate', from: a, to: x, type: VerbType.RelatedTo }, + { op: 'unrelate', id: old } + ]) + + const polluted = crossings.filter((c) => c.keys.length > 0) + expect(polluted).toEqual([]) + expect(await brain.storage.getVerb(old)).toBeFalsy() + const edges = await brain.related({ from: a }) + expect(edges.length).toBe(1) + expect(edges[0].id).not.toBe(old) + }) +}) + +describe('the metadata-index operations sanitize at the crossing, not at construction', () => { + /** A strict seam: refuses BigInts exactly as the native provider does. */ + const strictIndex = () => { + const seen: Array<{ door: string; keys: string[] }> = [] + const judge = (door: string, metadata: unknown) => { + const keys = bigintKeys(metadata) + seen.push({ door, keys }) + if (keys.length > 0) { + throw new Error( + `${door}: the metadata object violates the provider seam's JSON ` + + `contract — BigInt at ${keys.join(', ')}.` + ) + } + } + return { + seen, + addToIndex: async (_id: string, metadata: unknown) => judge('addToIndex', metadata), + removeFromIndex: async (_id: string, metadata: unknown) => judge('removeFromIndex', metadata) + } + } + + it('RemoveFromMetadataIndexOperation: entity mutated AFTER construction still crosses clean', async () => { + const index = strictIndex() + const verb: Record = { id: 'v1', sourceId: 'a', targetId: 'b' } + const op = new RemoveFromMetadataIndexOperation(index as any, 'v1', verb, () => 7n) + + // The graph leg's execute-time endpoint resolution, simulated: the shared + // object is polluted between plan and execute. + verb.sourceInt = 800_000n + verb.targetInt = 800_001n + + const rollback = await op.execute() + await rollback() + expect(index.seen.map((s) => s.keys)).toEqual([[], []]) + }) + + it('AddToMetadataIndexOperation: same law on the add leg and its rollback', async () => { + const index = strictIndex() + const verb: Record = { id: 'v2', sourceId: 'a', targetId: 'b' } + const op = new AddToMetadataIndexOperation(index as any, 'v2', verb, () => 7n) + + verb.sourceInt = 800_000n + verb.targetInt = 800_001n + + const rollback = await op.execute() + await rollback() + expect(index.seen.map((s) => s.keys)).toEqual([[], []]) + }) +}) diff --git a/tests/unit/db/generation-segments.test.ts b/tests/unit/db/generation-segments.test.ts index 27ab85cb..f16e67b3 100644 --- a/tests/unit/db/generation-segments.test.ts +++ b/tests/unit/db/generation-segments.test.ts @@ -147,4 +147,119 @@ describe('db/GenerationSegmentStore — the D1+D3 packed tier', () => { await expect(store.fold([gen(4), gen(4)])).rejects.toThrow(/strictly ascending/) await expect(store.fold([])).rejects.toThrow(/at least one generation/) }) + + // ========================================================================== + // THE DENSITY LAW + // ========================================================================== + // + // A sealed segment declares a CONTIGUOUS range and every reader treats that + // range as containment. Folding a sparse batch therefore makes the segment + // claim generations it does not hold — and because `open()` merges declared + // ranges back into committedRanges, the hole is re-admitted as committed + // history and every later maintenance pass fails asking for a frame that was + // never written. That is the "generation N is inside sealed segment + // seg-....bgs's declared range but has no frame — packed history is damaged" + // narration seen on every run of the affected stores. + + it('fold REFUSES a batch with a hole — a dense range may not be declared over sparse input', async () => { + await expect(store.fold([gen(1), gen(2), gen(4)])).rejects.toThrow( + /not contiguous: 2 → 4 skips 1 generation/ + ) + // The refusal loses nothing: no segment was sealed, so the generations + // stay in the live tier and the next pass folds them correctly. + expect(store.segments()).toHaveLength(0) + expect(store.hasGeneration(1)).toBe(false) + }) + + it('a wider gap names how many generations it would have swallowed', async () => { + await expect(store.fold([gen(10), gen(20)])).rejects.toThrow( + /not contiguous: 10 → 20 skips 9 generation\(s\)/ + ) + }) + + it('two contiguous runs folded separately declare honest ranges', async () => { + // What the caller now does instead of folding across the gap. + const a = await store.fold([gen(1), gen(2), gen(3)]) + const b = await store.fold([gen(7), gen(8)]) + expect(a).toMatchObject({ firstGeneration: 1, lastGeneration: 3, frames: 3 }) + expect(b).toMatchObject({ firstGeneration: 7, lastGeneration: 8, frames: 2 }) + // The gap is honestly outside the packed tier. + for (const g of [4, 5, 6]) expect(store.hasGeneration(g)).toBe(false) + for (const g of [1, 2, 3, 7, 8]) expect(store.hasGeneration(g)).toBe(true) + expect(await store.actualRanges()).toEqual([ + [1, 3], + [7, 8] + ]) + }) + + it('actualRanges() is exact and I/O-free for dense segments', async () => { + await store.fold([gen(1), gen(2)]) + await store.fold([gen(3), gen(4)]) + // Adjacent dense segments each contribute their declared range. + expect(await store.actualRanges()).toEqual([ + [1, 2], + [3, 4] + ]) + }) + + // ---- pre-existing damage: a store sealed by the old writer ---------------- + + /** + * Seal a SPARSE segment the way the pre-fix writer did: write the bytes and + * sidecar for a contiguous run, then rewrite the manifest so the segment + * declares a wider range than the frames it holds. This reproduces on disk + * exactly what the affected stores carry, without needing the old code. + */ + const sealSparseSegment = async (): Promise => { + await store.fold([gen(1), gen(2), gen(3)]) + const manifest = (await storage.readRawObject(`${SEGMENTS_PREFIX}/manifest.json`)) as any + // Declare 1..5 while holding frames for 1..3 — generations 4 and 5 become + // holes inside a sealed range. + manifest.segments[0].lastGeneration = 5 + await storage.writeRawObject(`${SEGMENTS_PREFIX}/manifest.json`, manifest) + } + + it('a pre-existing sparse segment reports its holes as UNPACKED, not as damage', async () => { + await sealSparseSegment() + const reopened = new GenerationSegmentStore(storage as any) + await reopened.open() + + // The frames it really holds still serve, byte-faithfully. + expect((await reopened.readDelta(2))?.timestamp).toBe(1_700_000_000_002) + expect(await reopened.readRecords(3)).toHaveLength(2) + + // The holes answer "not packed" instead of throwing. This is the fix for + // the wedge: the old reader threw here on EVERY maintenance pass. + expect(await reopened.readDelta(4)).toBeNull() + expect(await reopened.readRecords(5)).toBeNull() + }) + + it('actualRanges() excludes the holes so they are never re-admitted as committed', async () => { + await sealSparseSegment() + const reopened = new GenerationSegmentStore(storage as any) + await reopened.open() + // Declared 1..5; actually holds 1..3. The store seeds committedRanges from + // THIS, so generations 4 and 5 never become committed history again. + expect(await reopened.actualRanges()).toEqual([[1, 3]]) + }) + + it('a DENSE segment missing a frame is still loud damage', async () => { + // The other side of the branch: when the manifest claims a complete span, + // a missing frame means the manifest and sidecar disagree — real damage, + // and it must not be quietly downgraded to "unpacked". + await store.fold([gen(1), gen(2), gen(3)]) + const idxPath = `${SEGMENTS_PREFIX}/seg-${String(1).padStart(20, '0')}.idx` + const raw = (await storage.readRawBytes(idxPath))! + const { decode, encode } = await import('@msgpack/msgpack') + const idx = decode(raw) as any + // Drop generation 2's entry while the manifest still declares 3 frames. + idx.generations = idx.generations.filter(([g]: [number]) => g !== 2) + await storage.writeRawBytes(idxPath, encode(idx)) + + const reopened = new GenerationSegmentStore(storage as any) + await reopened.open() + await expect(reopened.readDelta(2)).rejects.toThrow( + /manifest and the sidecar disagree; packed history is damaged/ + ) + }) })