Compare commits
14 commits
main
...
rel/10.4.9
| Author | SHA1 | Date | |
|---|---|---|---|
| eec90bdd69 | |||
| 2648f56ddf | |||
| 8a2ebacf02 | |||
| d5147ed608 | |||
| 6a89adc468 | |||
| 88e79729d3 | |||
| 077cbc0b6f | |||
| 5e3b343a0e | |||
| 4014e0f125 | |||
| 73500e7d10 | |||
| 0f0022b1c9 | |||
| d6bcb14f69 | |||
| a963a744cc | |||
|
|
c99308710a |
23 changed files with 1753 additions and 145 deletions
23
CHANGELOG.md
23
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)
|
||||
|
|
|
|||
4
package-lock.json
generated
4
package-lock.json
generated
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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}"
|
||||
|
|
|
|||
349
src/brainy.ts
349
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<T = any> implements BrainyInterface<T> {
|
|||
// 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<T = any> implements BrainyInterface<T> {
|
|||
)
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
*/
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
*/
|
||||
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<void> {
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
* 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<T = any> implements BrainyInterface<T> {
|
|||
private async recoverPendingEmbedsFromLog(): Promise<void> {
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
*/
|
||||
/**
|
||||
* @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<string, unknown>
|
||||
let hasBigint = false
|
||||
for (const k in rec) {
|
||||
if (typeof rec[k] === 'bigint') { hasBigint = true; break }
|
||||
}
|
||||
if (!hasBigint) return metadata
|
||||
const out: Record<string, unknown> = {}
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
// 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<T = any> implements BrainyInterface<T> {
|
|||
}
|
||||
}
|
||||
|
||||
// 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<T = any> implements BrainyInterface<T> {
|
|||
* 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<void> {
|
||||
if (this.isReadOnly) return
|
||||
|
|
@ -12507,7 +12605,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
|||
])
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
|
||||
/**
|
||||
* @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<void> {
|
||||
let stamp: FamilyStamp | null
|
||||
|
|
@ -12546,11 +12652,16 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
|||
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<T = any> implements BrainyInterface<T> {
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* @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<void> {
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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<string[]> {
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
}
|
||||
|
||||
/**
|
||||
* 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<T>, existingResults: Result<T>[]): Promise<Result<T>[]> {
|
||||
if (!params.connected) return existingResults
|
||||
private async resolveConnectedIds(params: FindParams<T>): Promise<string[]> {
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
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<T = any> implements BrainyInterface<T> {
|
|||
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<T>, ids: string[]): Promise<Result<T>[]> {
|
||||
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<T>[] = []
|
||||
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<T = any> implements BrainyInterface<T> {
|
|||
* terminal releases have run.
|
||||
*/
|
||||
async close(): Promise<void> {
|
||||
if (this._pendingEmbedIds.size === 0) await this.writeEmbedLowWater()
|
||||
let closeFailure: unknown = null
|
||||
try {
|
||||
await this.closeDurableSteps()
|
||||
|
|
|
|||
|
|
@ -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 }
|
||||
|
|
|
|||
|
|
@ -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<Array<[number, number]>> {
|
||||
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`
|
||||
)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -411,6 +411,19 @@ export interface MetadataIndexProvider {
|
|||
* @returns The matching id universe as an opaque set.
|
||||
*/
|
||||
getIdSetForFilter?(filter: any): Promise<OpaqueIdSet>
|
||||
/**
|
||||
* @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<string[]>
|
||||
getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>>
|
||||
getSortedIdsForFilter(filter: any, orderBy: string, order?: 'asc' | 'desc', topK?: number): Promise<string[]>
|
||||
getFilterValues(field: string): Promise<string[]>
|
||||
|
|
|
|||
|
|
@ -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<void> | null = null
|
||||
/** The one trailing persist a burst has queued behind the in-flight one. */
|
||||
private countPersistTrailing: Promise<void> | 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-<pid>-<ms>`): 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
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -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<void> {
|
||||
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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<string>()
|
||||
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<string>()
|
||||
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<string>()
|
||||
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<string>()
|
||||
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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
47
src/utils/jsonSafeIndexMetadata.ts
Normal file
47
src/utils/jsonSafeIndexMetadata.ts
Normal file
|
|
@ -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<string, unknown>
|
||||
let hasBigint = false
|
||||
for (const k in rec) {
|
||||
if (typeof rec[k] === 'bigint') { hasBigint = true; break }
|
||||
}
|
||||
if (!hasBigint) return metadata
|
||||
const out: Record<string, unknown> = {}
|
||||
for (const k in rec) {
|
||||
if (typeof rec[k] !== 'bigint') out[k] = rec[k]
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
|
@ -2575,6 +2575,19 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
|||
/** Once-per-field flag for the fallback-degradation announcement. */
|
||||
private static announcedFallbackSorts = new Set<string>()
|
||||
|
||||
/**
|
||||
* 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<string[]> {
|
||||
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,
|
||||
|
|
|
|||
111
tests/integration/counts-persist-single-flight.test.ts
Normal file
111
tests/integration/counts-persist-single-flight.test.ts
Normal file
|
|
@ -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-<pid>-<ms>`). 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)
|
||||
})
|
||||
})
|
||||
|
|
@ -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 = {
|
||||
|
|
|
|||
165
tests/integration/find-connected-order.test.ts
Normal file
165
tests/integration/find-connected-order.test.ts
Normal file
|
|
@ -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<any>
|
||||
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<string>()
|
||||
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()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
|
@ -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<Brainy> => {
|
||||
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()
|
||||
|
|
|
|||
141
tests/integration/pending-embed-low-water.test.ts
Normal file
141
tests/integration/pending-embed-low-water.test.ts
Normal file
|
|
@ -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<Brainy<any>> => {
|
||||
const brain = new Brainy<any>({
|
||||
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()
|
||||
})
|
||||
})
|
||||
89
tests/integration/related-verb-array.test.ts
Normal file
89
tests/integration/related-verb-array.test.ts
Normal file
|
|
@ -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<any>
|
||||
|
||||
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)
|
||||
})
|
||||
})
|
||||
184
tests/integration/transact-edge-delete-bigint-aliasing.test.ts
Normal file
184
tests/integration/transact-edge-delete-bigint-aliasing.test.ts
Normal file
|
|
@ -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<string, unknown>)
|
||||
.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<string, unknown> = { 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<string, unknown> = { 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([[], []])
|
||||
})
|
||||
})
|
||||
|
|
@ -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<void> => {
|
||||
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/
|
||||
)
|
||||
})
|
||||
})
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue