Compare commits
No commits in common. "v10.1.0" and "v10.0.0" have entirely different histories.
19 changed files with 38 additions and 1192 deletions
|
|
@ -2,15 +2,6 @@
|
||||||
|
|
||||||
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.
|
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.1.0](https://source.soulcraft.com/soulcraft/brainy/compare/v10.0.0...v10.1.0) (2026-08-13)
|
|
||||||
|
|
||||||
- docs(releases): the 10.1.0 consumer entry — bounded recovery, restore founding, the two write-path cures (7d3c8696)
|
|
||||||
- fix(restore): a restore is an unclean event — the swap runs quiesced and the snapshot's durability stamps never survive it (9ca80667)
|
|
||||||
- feat(recovery): the fold-checkpoint bound — crash folds (checkpoint, head], never the whole log twice (ff43de1a)
|
|
||||||
- fix(log): pad-frame construction is total; the at-ack sync-failure compensation splits by phase — a production adoption's two write-path defects, cured at their roots (cbe34d11)
|
|
||||||
- feat(query): the sparse-store cut — where on a never-carried field serves operator truth, never a refusal (7b67db4d)
|
|
||||||
|
|
||||||
|
|
||||||
### [10.0.0](https://source.soulcraft.com/soulcraft/brainy/compare/v9.0.0...v10.0.0) (2026-08-12)
|
### [10.0.0](https://source.soulcraft.com/soulcraft/brainy/compare/v9.0.0...v10.0.0) (2026-08-12)
|
||||||
|
|
||||||
- fix(adoption): the baseline backfill cures hydration-law drift — existing brains reach the crash-safe default with zero operator steps (25f0dd96)
|
- fix(adoption): the baseline backfill cures hydration-law drift — existing brains reach the crash-safe default with zero operator steps (25f0dd96)
|
||||||
|
|
|
||||||
39
RELEASES.md
39
RELEASES.md
|
|
@ -31,45 +31,6 @@ is sometimes cited as a 7.x removal — those methods never existed on 7.x; the
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## v10.1.0 — 2026-08-13 (the bounded-recovery and write-path-cure release)
|
|
||||||
|
|
||||||
The theme: **crash recovery is bounded, restores are durably founded, and two
|
|
||||||
production-reported write-path defects are cured at their roots.** Ships together
|
|
||||||
with the matching native accelerator version; adopt as a pair.
|
|
||||||
|
|
||||||
- **Bounded crash recovery (the fold-checkpoint bound).** Recovery after an unclean
|
|
||||||
shutdown now replays only the log segment above a durably-stamped checkpoint
|
|
||||||
instead of the whole log. The checkpoint advances only after a canonical-sync
|
|
||||||
barrier makes every touched record durable (deletes included), so the bound can
|
|
||||||
lag but can never overstate durability. Existing stores converge automatically at
|
|
||||||
their first recovery — zero operator steps; recovery cost stops scaling with
|
|
||||||
store age.
|
|
||||||
- **Restores are unclean events, by construction.** `restore()` now runs its swap
|
|
||||||
fully quiesced (no background flush can race the directory replacement — a
|
|
||||||
consumer-reported `ENOTEMPTY` crash class is dead), and a snapshot's durability
|
|
||||||
stamps never survive the restore: the reopen folds the restored log, re-syncs
|
|
||||||
what it re-applied, and stamps fresh. Restored state is durably founded at
|
|
||||||
restore time instead of inheriting assertions about bytes the disk never synced.
|
|
||||||
- **Write-path cures from a production report.** (1) Log pad-frame construction is
|
|
||||||
total — a size-class boundary hole could previously kill a sync with "pad frame
|
|
||||||
not constructible". (2) The at-ack sync-failure compensation now splits by phase:
|
|
||||||
the generation counter can never re-mint a number the log may already carry, so
|
|
||||||
the non-monotonic append refusal loop reported by a downstream deployment cannot
|
|
||||||
recur. Both pinned with the reporter's exact shapes.
|
|
||||||
- **Operator-truthful sparse queries.** `where` on a field no store row has ever
|
|
||||||
carried now serves the honest answer (`eq`/`in`/range → empty; `ne`/`exists:false`
|
|
||||||
→ all rows; `exists:true` → empty) with a throttled did-you-mean warning, instead
|
|
||||||
of refusing. `orderBy` on unknown fields and ambiguous spellings keep their typed
|
|
||||||
refusals.
|
|
||||||
- **Cross-package error identity.** `UnresolvableFieldError` thrown across package
|
|
||||||
boundaries is re-normalized so `instanceof` checks in consuming applications
|
|
||||||
match regardless of duplicated dependency trees.
|
|
||||||
- Release tooling: publishes now push the tag before the branch (the publish
|
|
||||||
workflow can no longer queue behind a redundant CI run) and verify registry
|
|
||||||
byte-identity with a propagation-tolerant raw-registry probe.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## v10.0.0 — 2026-08-10 (the write-path and lifecycle release)
|
## v10.0.0 — 2026-08-10 (the write-path and lifecycle release)
|
||||||
|
|
||||||
The theme: **writes ack fast and honestly, startup adopts instead of rebuilding, and
|
The theme: **writes ack fast and honestly, startup adopts instead of rebuilding, and
|
||||||
|
|
|
||||||
4
package-lock.json
generated
4
package-lock.json
generated
|
|
@ -1,12 +1,12 @@
|
||||||
{
|
{
|
||||||
"name": "@soulcraft/brainy",
|
"name": "@soulcraft/brainy",
|
||||||
"version": "10.1.0",
|
"version": "10.0.0",
|
||||||
"lockfileVersion": 3,
|
"lockfileVersion": 3,
|
||||||
"requires": true,
|
"requires": true,
|
||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"name": "@soulcraft/brainy",
|
"name": "@soulcraft/brainy",
|
||||||
"version": "10.1.0",
|
"version": "10.0.0",
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@msgpack/msgpack": "^3.1.2",
|
"@msgpack/msgpack": "^3.1.2",
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
{
|
{
|
||||||
"name": "@soulcraft/brainy",
|
"name": "@soulcraft/brainy",
|
||||||
"version": "10.1.0",
|
"version": "10.0.0",
|
||||||
"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.",
|
"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",
|
"main": "dist/index.js",
|
||||||
"module": "dist/index.js",
|
"module": "dist/index.js",
|
||||||
|
|
|
||||||
|
|
@ -177,15 +177,8 @@ echo -e "${GREEN}✅ Tag created${NC}\n"
|
||||||
|
|
||||||
# Step 9: Push to origin — The Source is the one home (ruled 2026-07-23; the
|
# Step 9: Push to origin — The Source is the one home (ruled 2026-07-23; the
|
||||||
# old public GitHub repo is archived history, no longer part of any release).
|
# old public GitHub repo is archived history, no longer part of any release).
|
||||||
# TAG FIRST, branch second — deliberately two pushes: the runner is
|
echo -e "${BLUE}8️⃣ Pushing to origin...${NC}"
|
||||||
# sequential, and a combined push can queue the release commit's ci.yml run
|
git push --follow-tags origin "$CURRENT_BRANCH"
|
||||||
# AHEAD of the tag's publish-source run (observed on 10.0.0: the publish sat
|
|
||||||
# ~37 minutes behind a redundant CI run of the very commit the local gates
|
|
||||||
# had just proven). Pushing the tag alone queues the publish immediately;
|
|
||||||
# the branch push (and its ci.yml run) follows behind it, harmlessly.
|
|
||||||
echo -e "${BLUE}8️⃣ Pushing to origin (tag first — the publish must never queue behind CI)...${NC}"
|
|
||||||
git push origin "v${NEW_VERSION}"
|
|
||||||
git push origin "$CURRENT_BRANCH"
|
|
||||||
echo -e "${GREEN}✅ Pushed to origin${NC}\n"
|
echo -e "${GREEN}✅ Pushed to origin${NC}\n"
|
||||||
|
|
||||||
# Step 10: The home publish (The Source, source.soulcraft.com) is CI's job
|
# Step 10: The home publish (The Source, source.soulcraft.com) is CI's job
|
||||||
|
|
@ -234,32 +227,14 @@ npm publish "$SOURCE_TARBALL" --tag "$NPM_TAG" "--@soulcraft:registry=https://re
|
||||||
rm -rf "$STOREFRONT_TMP"
|
rm -rf "$STOREFRONT_TMP"
|
||||||
# Brainy is the only PUBLIC @soulcraft package — verify visibility after every publish.
|
# Brainy is the only PUBLIC @soulcraft package — verify visibility after every publish.
|
||||||
npm access get status @soulcraft/brainy "--@soulcraft:registry=https://registry.npmjs.org/" || true
|
npm access get status @soulcraft/brainy "--@soulcraft:registry=https://registry.npmjs.org/" || true
|
||||||
# Verify the pair is byte-identical by registry-reported shasum — divergence
|
# Verify the pair is byte-identical by registry-reported shasum — divergence here
|
||||||
# here means the storefront leg must be treated as failed, loudly. RETRIED
|
# means the storefront leg must be treated as failed, loudly.
|
||||||
# with raw curl: npmjs metadata propagates with a lag measured in minutes,
|
|
||||||
# and a one-shot npm-view probe fired a false DIVERGENCE on 10.0.0 while a
|
|
||||||
# raw curl of the registry document already confirmed byte-identity. The
|
|
||||||
# probe now reads the registry JSON directly (no npm cache in the path) and
|
|
||||||
# gives propagation up to 5 minutes before calling the pair divergent.
|
|
||||||
NPMJS_VERIFY_ATTEMPTS=20
|
|
||||||
NPMJS_VERIFY_INTERVAL_S=15 # 20 × 15s = 5 minutes of propagation grace
|
|
||||||
SOURCE_SHA=$(npm view "@soulcraft/brainy@${NEW_VERSION}" dist.shasum "--@soulcraft:registry=${SOURCE_NPM_REG}" 2>/dev/null || echo "source-unavailable")
|
SOURCE_SHA=$(npm view "@soulcraft/brainy@${NEW_VERSION}" dist.shasum "--@soulcraft:registry=${SOURCE_NPM_REG}" 2>/dev/null || echo "source-unavailable")
|
||||||
PAIR_IDENTICAL=false
|
NPMJS_SHA=$(npm view "@soulcraft/brainy@${NEW_VERSION}" dist.shasum "--@soulcraft:registry=https://registry.npmjs.org/" 2>/dev/null || echo "npmjs-unavailable")
|
||||||
for ((attempt = 1; attempt <= NPMJS_VERIFY_ATTEMPTS; attempt++)); do
|
if [ "$SOURCE_SHA" = "$NPMJS_SHA" ]; then
|
||||||
NPMJS_SHA=$(curl -fsSL "https://registry.npmjs.org/@soulcraft%2Fbrainy" 2>/dev/null \
|
|
||||||
| node -e "let d='';process.stdin.on('data',c=>d+=c).on('end',()=>{try{const v=JSON.parse(d).versions[process.argv[1]];console.log(v?v.dist.shasum:'')}catch{console.log('')}})" "${NEW_VERSION}" \
|
|
||||||
|| echo "")
|
|
||||||
if [ -n "$NPMJS_SHA" ] && [ "$SOURCE_SHA" = "$NPMJS_SHA" ]; then
|
|
||||||
PAIR_IDENTICAL=true
|
|
||||||
break
|
|
||||||
fi
|
|
||||||
echo -e "${YELLOW} … npmjs metadata not settled (attempt ${attempt}/${NPMJS_VERIFY_ATTEMPTS}: '${NPMJS_SHA:-absent}' vs '${SOURCE_SHA}'); retrying in ${NPMJS_VERIFY_INTERVAL_S}s${NC}"
|
|
||||||
sleep "$NPMJS_VERIFY_INTERVAL_S"
|
|
||||||
done
|
|
||||||
if [ "$PAIR_IDENTICAL" = true ]; then
|
|
||||||
echo -e "${GREEN}✅ Published to npmjs — byte-identical pair (shasum ${NPMJS_SHA})${NC}\n"
|
echo -e "${GREEN}✅ Published to npmjs — byte-identical pair (shasum ${NPMJS_SHA})${NC}\n"
|
||||||
else
|
else
|
||||||
echo -e "${RED}❌ REGISTRY DIVERGENCE: The Source shasum ${SOURCE_SHA} != npmjs shasum ${NPMJS_SHA} after ${NPMJS_VERIFY_ATTEMPTS} attempts — investigate before announcing${NC}\n"
|
echo -e "${RED}❌ REGISTRY DIVERGENCE: The Source shasum ${SOURCE_SHA} != npmjs shasum ${NPMJS_SHA} — investigate before announcing${NC}\n"
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -197,7 +197,6 @@ import { GenerationConflictError, StoreInconsistentError } from './db/errors.js'
|
||||||
import { BrainyError, GraphIndexNotReadyError, MetadataIndexNotReadyError, MigrationInProgressError, VectorIndexNotReadyError } from './errors/brainyError.js'
|
import { BrainyError, GraphIndexNotReadyError, MetadataIndexNotReadyError, MigrationInProgressError, VectorIndexNotReadyError } from './errors/brainyError.js'
|
||||||
import { assessIndexReadiness } from './utils/indexReadiness.js'
|
import { assessIndexReadiness } from './utils/indexReadiness.js'
|
||||||
import { reconstructNounWrapper } from './db/factLog.js'
|
import { reconstructNounWrapper } from './db/factLog.js'
|
||||||
import { asBrainyFieldRefusal } from './db/fieldAddressing.js'
|
|
||||||
import {
|
import {
|
||||||
readLogAuthority,
|
readLogAuthority,
|
||||||
runLogCompletenessOracle,
|
runLogCompletenessOracle,
|
||||||
|
|
@ -4108,7 +4107,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
|
|
||||||
const probeServes = async (): Promise<boolean> => {
|
const probeServes = async (): Promise<boolean> => {
|
||||||
try {
|
try {
|
||||||
const ids = await this.filterIdsBelted({ [p.field]: p.value })
|
const ids = await this.metadataIndex.getIdsForFilter({ [p.field]: p.value })
|
||||||
return ids.includes(p.id)
|
return ids.includes(p.id)
|
||||||
} catch {
|
} catch {
|
||||||
// FIELD_NOT_INDEXED for a field a persisted entity actually holds is
|
// FIELD_NOT_INDEXED for a field a persisted entity actually holds is
|
||||||
|
|
@ -6448,7 +6447,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
// 'visibility' key would address the USER's metadata bag under the
|
// 'visibility' key would address the USER's metadata bag under the
|
||||||
// field-addressing law and silently hide nothing (VFS/system entities
|
// field-addressing law and silently hide nothing (VFS/system entities
|
||||||
// would leak into every default read).
|
// would leak into every default read).
|
||||||
const ids = await this.filterIdsBelted({
|
const ids = await this.metadataIndex.getIdsForFilter({
|
||||||
'system.visibility': excluded.length === 1 ? excluded[0] : { oneOf: excluded }
|
'system.visibility': excluded.length === 1 ? excluded[0] : { oneOf: excluded }
|
||||||
})
|
})
|
||||||
return new Set(ids)
|
return new Set(ids)
|
||||||
|
|
@ -6656,7 +6655,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
// offset stays 0 because the visibility filter + slice happen here. The JS
|
// offset stays 0 because the visibility filter + slice happen here. The JS
|
||||||
// index ignores the bound and returns all matches (behaviour unchanged).
|
// index ignores the bound and returns all matches (behaviour unchanged).
|
||||||
const pageEnd = (params.offset || 0) + (params.limit || 10) + hiddenIds.size
|
const pageEnd = (params.offset || 0) + (params.limit || 10) + hiddenIds.size
|
||||||
filteredIds = await this.filterIdsBelted(filter, { limit: pageEnd, offset: 0 })
|
filteredIds = await this.metadataIndex.getIdsForFilter(filter, { limit: pageEnd, offset: 0 })
|
||||||
}
|
}
|
||||||
|
|
||||||
// Visibility hard filter — drop hidden ids BEFORE pagination so limit is exact.
|
// Visibility hard filter — drop hidden ids BEFORE pagination so limit is exact.
|
||||||
|
|
@ -6728,7 +6727,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
// filter returns nothing from getIdsForFilter, so the unfiltered case below uses
|
// filter returns nothing from getIdsForFilter, so the unfiltered case below uses
|
||||||
// getNouns instead (it returns all nouns, including their visibility).
|
// getNouns instead (it returns all nouns, including their visibility).
|
||||||
if (Object.keys(filter).length > 0) {
|
if (Object.keys(filter).length > 0) {
|
||||||
let filteredIds = await this.filterIdsBelted(filter)
|
let filteredIds = await this.metadataIndex.getIdsForFilter(filter)
|
||||||
// Visibility hard filter — drop hidden ids BEFORE pagination.
|
// Visibility hard filter — drop hidden ids BEFORE pagination.
|
||||||
if (hiddenIds.size > 0) filteredIds = filteredIds.filter((id) => !hiddenIds.has(id))
|
if (hiddenIds.size > 0) filteredIds = filteredIds.filter((id) => !hiddenIds.has(id))
|
||||||
const pageIds = filteredIds.slice(offset, offset + limit)
|
const pageIds = filteredIds.slice(offset, offset + limit)
|
||||||
|
|
@ -6779,7 +6778,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
|
|
||||||
if (params.where || params.type || params.subtype || params.service || params.excludeVFS) {
|
if (params.where || params.type || params.subtype || params.service || params.excludeVFS) {
|
||||||
preResolvedFilter = this.buildMetadataFilter(params)
|
preResolvedFilter = this.buildMetadataFilter(params)
|
||||||
preResolvedMetadataIds = await this.filterIdsBelted(preResolvedFilter)
|
preResolvedMetadataIds = await this.metadataIndex.getIdsForFilter(preResolvedFilter)
|
||||||
|
|
||||||
// Visibility hard filter — restrict the HNSW candidate set to non-hidden ids.
|
// Visibility hard filter — restrict the HNSW candidate set to non-hidden ids.
|
||||||
if (hiddenIds.size > 0) {
|
if (hiddenIds.size > 0) {
|
||||||
|
|
@ -8239,23 +8238,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
async adoptLogAuthority(): Promise<OracleReport> {
|
async adoptLogAuthority(): Promise<OracleReport> {
|
||||||
await this.ensureInitialized()
|
await this.ensureInitialized()
|
||||||
this.assertWritable('adoptLogAuthority')
|
this.assertWritable('adoptLogAuthority')
|
||||||
// Fold-checkpoint chain, phase 1: a FRESH brain (no committed
|
|
||||||
// generations) arms the chain now so the backfill's re-commits below
|
|
||||||
// feed the canonical-sync accumulator — its first stamp is then total.
|
|
||||||
// A non-fresh flip skips (the store refuses the arm); its chain starts
|
|
||||||
// at the first recovery fold instead. Disarmed on any failure below.
|
|
||||||
this.generationStore.beginFoldCheckpointBootstrap()
|
|
||||||
try {
|
|
||||||
return await this.adoptLogAuthorityInner()
|
|
||||||
} catch (err) {
|
|
||||||
this.generationStore.abandonFoldCheckpointBootstrap()
|
|
||||||
throw err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The adoption body — see {@link Brainy.adoptLogAuthority} (which owns the
|
|
||||||
* fold-checkpoint bootstrap arm/disarm around it). */
|
|
||||||
private async adoptLogAuthorityInner(): Promise<OracleReport> {
|
|
||||||
let report = await this.verifyLogAuthority()
|
let report = await this.verifyLogAuthority()
|
||||||
|
|
||||||
// BASELINE BACKFILL: curable divergences are rows whose CANONICAL truth
|
// BASELINE BACKFILL: curable divergences are rows whose CANONICAL truth
|
||||||
|
|
@ -8351,9 +8333,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
report
|
report
|
||||||
)
|
)
|
||||||
this.generationStore.setLogDurability('at-ack')
|
this.generationStore.setLogDurability('at-ack')
|
||||||
// Fold-checkpoint chain, phase 2: the flip is recorded — open the stamp
|
|
||||||
// gate so the next flush/close barrier writes the first checkpoint.
|
|
||||||
this.generationStore.completeFoldCheckpointBootstrap()
|
|
||||||
return report
|
return report
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -9229,13 +9208,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
}
|
}
|
||||||
|
|
||||||
const floorGeneration = this.generationStore.generation()
|
const floorGeneration = this.generationStore.generation()
|
||||||
// The swap runs inside the generation store's exclusive section: pending
|
await this.storage.restoreFromDirectory(path)
|
||||||
// flush timers are disarmed and buffers discarded BEFORE any directory is
|
|
||||||
// removed, so a background flush can never write into `_system/` mid-swap
|
|
||||||
// (the ENOTEMPTY race a checkpoint stamp once hit).
|
|
||||||
await this.generationStore.runStateReplacement(() =>
|
|
||||||
this.storage.restoreFromDirectory(path)
|
|
||||||
)
|
|
||||||
await this.generationStore.reopenAfterRestore(floorGeneration)
|
await this.generationStore.reopenAfterRestore(floorGeneration)
|
||||||
|
|
||||||
// If the entity-id mapper is a NATIVE provider with a `rebuild()`, reload it
|
// If the entity-id mapper is a NATIVE provider with a `rebuild()`, reload it
|
||||||
|
|
@ -11575,27 +11548,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
* console.log(`Lazy rebuild completed: ${status.lazyRebuildCompleted}`)
|
* console.log(`Lazy rebuild completed: ${status.lazyRebuildCompleted}`)
|
||||||
* ```
|
* ```
|
||||||
*/
|
*/
|
||||||
|
|
||||||
/**
|
|
||||||
* The provider-seam belt for filter reads: whatever manager serves
|
|
||||||
* getIdsForFilter (the JS twin or a native replacement), a field refusal
|
|
||||||
* crossing this seam is normalized to BRAINY'S UnresolvableFieldError —
|
|
||||||
* one class identity for consumers, never a foreign twin that fails
|
|
||||||
* instanceof. All other errors pass through untouched.
|
|
||||||
*/
|
|
||||||
private async filterIdsBelted(
|
|
||||||
filter: unknown,
|
|
||||||
opts?: { limit?: number; offset?: number }
|
|
||||||
): Promise<string[]> {
|
|
||||||
try {
|
|
||||||
return await this.metadataIndex.getIdsForFilter(filter, opts)
|
|
||||||
} catch (err) {
|
|
||||||
const normalized = asBrainyFieldRefusal(err)
|
|
||||||
if (normalized) throw normalized
|
|
||||||
throw err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
async getIndexStatus(): Promise<{
|
async getIndexStatus(): Promise<{
|
||||||
initialized: boolean
|
initialized: boolean
|
||||||
lazyRebuildCompleted: boolean
|
lazyRebuildCompleted: boolean
|
||||||
|
|
@ -12193,7 +12145,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const filteredIds = await this.filterIdsBelted(filter)
|
const filteredIds = await this.metadataIndex.getIdsForFilter(filter)
|
||||||
return filteredIds.length
|
return filteredIds.length
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -12265,7 +12217,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const filteredIds = await this.filterIdsBelted(filterObj)
|
const filteredIds = await this.metadataIndex.getIdsForFilter(filterObj)
|
||||||
|
|
||||||
// Stream filtered entities in batches for memory efficiency
|
// Stream filtered entities in batches for memory efficiency
|
||||||
const batchSize = 100
|
const batchSize = 100
|
||||||
|
|
|
||||||
|
|
@ -1204,30 +1204,6 @@ function buildPadFrame(totalBytes: number): Uint8Array {
|
||||||
fillerLength += diff
|
fillerLength += diff
|
||||||
if (fillerLength < 0) break
|
if (fillerLength < 0) break
|
||||||
}
|
}
|
||||||
if (!converged) {
|
|
||||||
// Class-boundary holes: a single bin filler steps its header by one
|
|
||||||
// byte at each msgpack size class (bin8→bin16→bin32), leaving exactly
|
|
||||||
// one unreachable payload size per boundary (the 291-byte production
|
|
||||||
// case). Bridge with a trailing fixint (+1 byte) beside the bin —
|
|
||||||
// {bin(n)} ∪ {bin(n) + fixint} covers every size ≥ minimum.
|
|
||||||
let bridged = Math.max(0, targetPayload - payload.length - 2)
|
|
||||||
for (let i = 0; i < 8; i++) {
|
|
||||||
const candidate = attempt([
|
|
||||||
LOG_RECORD_TYPES.PAD,
|
|
||||||
LOG_RECORD_VERSION,
|
|
||||||
new Uint8Array(bridged),
|
|
||||||
0
|
|
||||||
])
|
|
||||||
const diff = targetPayload - candidate.length
|
|
||||||
if (diff === 0) {
|
|
||||||
payload = candidate
|
|
||||||
converged = true
|
|
||||||
break
|
|
||||||
}
|
|
||||||
bridged += diff
|
|
||||||
if (bridged < 0) break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (!converged) {
|
if (!converged) {
|
||||||
throw new Error(`fact log v2: a pad frame of ${totalBytes} bytes is not constructible`)
|
throw new Error(`fact log v2: a pad frame of ${totalBytes} bytes is not constructible`)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -261,28 +261,6 @@ export function buildUnresolvableMessage(
|
||||||
* the fix ships inside the error. Thrown by the query layer with index
|
* the fix ships inside the error. Thrown by the query layer with index
|
||||||
* knowledge, never by the pure parser.
|
* knowledge, never by the pure parser.
|
||||||
*/
|
*/
|
||||||
/**
|
|
||||||
* Cross-package identity normalizer (the seam belt): the native accelerator
|
|
||||||
* throws ITS OWN UnresolvableFieldError class, which fails `instanceof`
|
|
||||||
* against this package's export — consumers were forced to match by name.
|
|
||||||
* Every provider-boundary catch routes suspected field-refusals through
|
|
||||||
* here: a foreign refusal (matched by name, duck fields tolerated) is
|
|
||||||
* rethrown as THIS package's class, so exactly one identity ever reaches
|
|
||||||
* consumers. Anything else returns null (caller rethrows the original).
|
|
||||||
*/
|
|
||||||
export function asBrainyFieldRefusal(err: unknown): UnresolvableFieldError | null {
|
|
||||||
if (err instanceof UnresolvableFieldError) return err
|
|
||||||
const e = err as { name?: string; message?: string; raw?: string; kind?: string } | null
|
|
||||||
if (e && e.name === 'UnresolvableFieldError') {
|
|
||||||
return new UnresolvableFieldError(
|
|
||||||
e.raw ?? 'unknown-field',
|
|
||||||
(e.kind as FieldAddressKind) ?? 'entity',
|
|
||||||
e.message
|
|
||||||
)
|
|
||||||
}
|
|
||||||
return null
|
|
||||||
}
|
|
||||||
|
|
||||||
export class UnresolvableFieldError extends Error {
|
export class UnresolvableFieldError extends Error {
|
||||||
public readonly raw: string
|
public readonly raw: string
|
||||||
public readonly kind: FieldAddressKind
|
public readonly kind: FieldAddressKind
|
||||||
|
|
|
||||||
|
|
@ -78,21 +78,10 @@ export const MANIFEST_PATH = '_system/manifest.json'
|
||||||
/**
|
/**
|
||||||
* The clean-shutdown marker (log-authority recovery gate): written+fsynced at
|
* The clean-shutdown marker (log-authority recovery gate): written+fsynced at
|
||||||
* a clean close carrying the committed generation; CONSUMED at every open.
|
* a clean close carrying the committed generation; CONSUMED at every open.
|
||||||
* Absent or generation-mismatched at open = unclean shutdown = the replay
|
* Absent or generation-mismatched at open = unclean shutdown = the whole-log
|
||||||
* fold, bounded below by the fold checkpoint when one is stored (whole-log
|
* replay fold. Its absence is always safe (costs one replay, loses nothing).
|
||||||
* without one). Its absence is always safe (costs one fold, loses nothing).
|
|
||||||
*/
|
*/
|
||||||
export const CLEAN_SHUTDOWN_PATH = '_system/clean-shutdown.json'
|
export const CLEAN_SHUTDOWN_PATH = '_system/clean-shutdown.json'
|
||||||
/**
|
|
||||||
* The fold checkpoint (log-authority recovery BOUND): `{ generation: G }`
|
|
||||||
* asserts that every entity whose latest fact is ≤ G has durable canonical
|
|
||||||
* bytes — so an unclean open folds only `(G, head]` instead of the whole log.
|
|
||||||
* Stamped strictly AFTER a canonical-sync barrier over every live entity
|
|
||||||
* touched since the last stamp (stamp-after-data); absent or torn = fold from
|
|
||||||
* 0 (always safe, just bigger). The chain of stamps starts only at a provable
|
|
||||||
* point: an empty brain, or the end of a whole-log fold.
|
|
||||||
*/
|
|
||||||
export const FOLD_CHECKPOINT_PATH = '_system/fold-checkpoint.json'
|
|
||||||
/** Storage-root-relative prefix of the per-generation record directories. */
|
/** Storage-root-relative prefix of the per-generation record directories. */
|
||||||
export const GENERATIONS_PREFIX = '_generations'
|
export const GENERATIONS_PREFIX = '_generations'
|
||||||
|
|
||||||
|
|
@ -230,37 +219,6 @@ export class GenerationStore {
|
||||||
/** Compaction horizon — record-sets ≤ this are reclaimed. */
|
/** Compaction horizon — record-sets ≤ this are reclaimed. */
|
||||||
private horizonGen = 0
|
private horizonGen = 0
|
||||||
|
|
||||||
/**
|
|
||||||
* Fold-checkpoint accumulator: every entity whose CANONICAL live bytes were
|
|
||||||
* (re)written since the last stamped checkpoint. Drained by
|
|
||||||
* {@link advanceFoldCheckpointUnlocked} — synced first, stamped after; on a
|
|
||||||
* failed barrier the drained ids merge back so the checkpoint can never
|
|
||||||
* advance past unsynced bytes. Fed only while the chain is valid (see
|
|
||||||
* {@link foldCheckpointChainValid}) so tree-authority brains never grow it.
|
|
||||||
*/
|
|
||||||
private checkpointDirtyNouns = new Set<string>()
|
|
||||||
/** @see checkpointDirtyNouns — the verb half of the accumulator. */
|
|
||||||
private checkpointDirtyVerbs = new Set<string>()
|
|
||||||
/**
|
|
||||||
* Whether the checkpoint chain is PROVABLY sound for this brain: true when
|
|
||||||
* a stored checkpoint exists (induction), the brain opened empty (vacuous),
|
|
||||||
* or a whole-log fold just re-applied every fact (base case). While false,
|
|
||||||
* checkpoints are never stamped and the fold bound stays 0 — the honest
|
|
||||||
* 10.0 contract, upgraded at the brain's first recovery fold.
|
|
||||||
*/
|
|
||||||
private foldCheckpointChainValid = false
|
|
||||||
/** Last stamped fold-checkpoint generation (0 = none / fold from origin). */
|
|
||||||
private foldCheckpoint = 0
|
|
||||||
/**
|
|
||||||
* Whether this brain's stored authority is the log — set from the stored
|
|
||||||
* artifact at open, or by {@link completeFoldCheckpointBootstrap} when an
|
|
||||||
* in-session adoption flips it. Checkpoints are only ever STAMPED under log
|
|
||||||
* authority (the artifact bounds the log fold, which only log-authority
|
|
||||||
* recovery runs); the dirty accumulator may fill slightly earlier, during
|
|
||||||
* an adoption in flight (see {@link beginFoldCheckpointBootstrap}).
|
|
||||||
*/
|
|
||||||
private authorityIsLog = false
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Committed generations whose record dirs exist, stored as a SORTED, DISJOINT,
|
* Committed generations whose record dirs exist, stored as a SORTED, DISJOINT,
|
||||||
* ascending list of INCLUSIVE `[start, end]` intervals (a run-length set).
|
* ascending list of INCLUSIVE `[start, end]` intervals (a run-length set).
|
||||||
|
|
@ -594,7 +552,6 @@ export class GenerationStore {
|
||||||
// drift machinery at open — same as group-commit recovery.
|
// drift machinery at open — same as group-commit recovery.
|
||||||
const authority = await readLogAuthority(this.storage)
|
const authority = await readLogAuthority(this.storage)
|
||||||
if (authority.authority === 'log') {
|
if (authority.authority === 'log') {
|
||||||
this.authorityIsLog = true
|
|
||||||
// TWO REPLAY TIERS, gated by the clean-shutdown marker:
|
// TWO REPLAY TIERS, gated by the clean-shutdown marker:
|
||||||
//
|
//
|
||||||
// (1) ABOVE-MANIFEST (always): an intact fact above the manifest is
|
// (1) ABOVE-MANIFEST (always): an intact fact above the manifest is
|
||||||
|
|
@ -617,22 +574,8 @@ export class GenerationStore {
|
||||||
const cleanShutdown = await this.readCleanShutdownMarker()
|
const cleanShutdown = await this.readCleanShutdownMarker()
|
||||||
const orphans = await this.factLog.peekFactsAbove(this.committed)
|
const orphans = await this.factLog.peekFactsAbove(this.committed)
|
||||||
const uncleanOpen = cleanShutdown === null || cleanShutdown !== this.committed
|
const uncleanOpen = cleanShutdown === null || cleanShutdown !== this.committed
|
||||||
// FOLD-CHECKPOINT BOUND: a stored checkpoint G proves every entity
|
|
||||||
// whose latest fact is ≤ G has durable canonical bytes (each stamp
|
|
||||||
// followed a canonical-sync barrier), so the unclean fold only needs
|
|
||||||
// (G, head] — entities untouched since G are already safe, entities
|
|
||||||
// touched after G get their latest after-image re-applied. Absent or
|
|
||||||
// invalid checkpoint = fold from 0 (the 10.0 whole-log contract).
|
|
||||||
const checkpoint = await this.readFoldCheckpoint()
|
|
||||||
const foldBound = checkpoint ?? 0
|
|
||||||
// Chain validity: induction (a stored stamp), vacuous truth (an empty
|
|
||||||
// brain has no bytes to assert), or — set below — the base case (a
|
|
||||||
// whole-log fold re-applies and re-syncs every entity in the log).
|
|
||||||
this.foldCheckpointChainValid = checkpoint !== null || this.committed === 0
|
|
||||||
this.foldCheckpoint = foldBound
|
|
||||||
if (uncleanOpen) this.foldCheckpointChainValid = true
|
|
||||||
const factsToReplay = uncleanOpen
|
const factsToReplay = uncleanOpen
|
||||||
? await this.factLog.peekFactsAbove(foldBound)
|
? await this.factLog.peekFactsAbove(0)
|
||||||
: orphans
|
: orphans
|
||||||
if (factsToReplay.length > 0) {
|
if (factsToReplay.length > 0) {
|
||||||
let replayed = 0
|
let replayed = 0
|
||||||
|
|
@ -644,7 +587,6 @@ export class GenerationStore {
|
||||||
: { metadata: op.record.metadata, vector: op.record.vector }
|
: { metadata: op.record.metadata, vector: op.record.vector }
|
||||||
if (op.kind === 'verb') await this.storage.writeVerbRaw(op.id, image)
|
if (op.kind === 'verb') await this.storage.writeVerbRaw(op.id, image)
|
||||||
else await this.storage.writeNounRaw(op.id, image)
|
else await this.storage.writeNounRaw(op.id, image)
|
||||||
this.noteCheckpointDirty(op.kind, op.id)
|
|
||||||
}
|
}
|
||||||
replayed++
|
replayed++
|
||||||
if (fact.generation > this.committed) {
|
if (fact.generation > this.committed) {
|
||||||
|
|
@ -670,21 +612,10 @@ export class GenerationStore {
|
||||||
await this.storage.syncRawObjects([MANIFEST_PATH])
|
await this.storage.syncRawObjects([MANIFEST_PATH])
|
||||||
prodLog.warn(
|
prodLog.warn(
|
||||||
`[GenerationStore] log-authority recovery replayed ${replayed} fact(s) into ` +
|
`[GenerationStore] log-authority recovery replayed ${replayed} fact(s) into ` +
|
||||||
`canonical (${
|
`canonical (${uncleanOpen ? 'WHOLE-LOG fold — unclean shutdown' : 'above-manifest'}; ` +
|
||||||
uncleanOpen
|
`committed at ${this.committed}) — an acked write is never lost`
|
||||||
? foldBound > 0
|
|
||||||
? `BOUNDED fold above checkpoint ${foldBound} — unclean shutdown`
|
|
||||||
: 'WHOLE-LOG fold — unclean shutdown'
|
|
||||||
: 'above-manifest'
|
|
||||||
}; committed at ${this.committed}) — an acked write is never lost`
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
// A recovery fold re-applied (and the barrier below re-syncs) every
|
|
||||||
// entity in (bound, head] — stamp the checkpoint at the new committed
|
|
||||||
// watermark so the NEXT crash folds only its own tail. This is also
|
|
||||||
// the chain's base case: the first whole-log fold of a pre-checkpoint
|
|
||||||
// brain covers every entity in the log, so its stamp is total.
|
|
||||||
if (uncleanOpen) await this.advanceFoldCheckpointUnlocked()
|
|
||||||
// The marker is consumed: any session that can write invalidates it
|
// The marker is consumed: any session that can write invalidates it
|
||||||
// at first commit (see the commit paths); a clean close re-writes it.
|
// at first commit (see the commit paths); a clean close re-writes it.
|
||||||
await this.clearCleanShutdownMarker()
|
await this.clearCleanShutdownMarker()
|
||||||
|
|
@ -741,12 +672,6 @@ export class GenerationStore {
|
||||||
await this.flushPendingSingleOps()
|
await this.flushPendingSingleOps()
|
||||||
this.storage.setGenerationBumpHook(undefined)
|
this.storage.setGenerationBumpHook(undefined)
|
||||||
await this.persistCounterNow()
|
await this.persistCounterNow()
|
||||||
// Fold-checkpoint barrier BEFORE the clean-shutdown marker: entities that
|
|
||||||
// reached the accumulator outside the pending tier (transact commits,
|
|
||||||
// aborted-write restores) get their canonical bytes synced and the stamp
|
|
||||||
// advanced, so the marker below never vouches for bytes the checkpoint
|
|
||||||
// chain hasn't proven durable.
|
|
||||||
await this.advanceFoldCheckpoint()
|
|
||||||
// Clean-shutdown marker (log-authority recovery gate): everything above
|
// Clean-shutdown marker (log-authority recovery gate): everything above
|
||||||
// is durable; stamp the committed generation so the next open can adopt
|
// is durable; stamp the committed generation so the next open can adopt
|
||||||
// instead of folding the log. Written LAST — a crash before this line is
|
// instead of folding the log. Written LAST — a crash before this line is
|
||||||
|
|
@ -780,131 +705,6 @@ export class GenerationStore {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Read the fold checkpoint's generation, or `null` when absent, torn, or
|
|
||||||
* implausible (> committed) — every invalid shape degrades to the safe
|
|
||||||
* whole-log fold, never to a bound that could skip an acked write.
|
|
||||||
*/
|
|
||||||
private async readFoldCheckpoint(): Promise<number | null> {
|
|
||||||
try {
|
|
||||||
const raw = (await this.storage.readRawObject(FOLD_CHECKPOINT_PATH)) as {
|
|
||||||
generation?: number
|
|
||||||
} | null
|
|
||||||
const gen = raw?.generation
|
|
||||||
if (!Number.isSafeInteger(gen) || (gen as number) < 0) return null
|
|
||||||
if ((gen as number) > this.committed) {
|
|
||||||
prodLog.warn(
|
|
||||||
`[GenerationStore] fold checkpoint ${gen} is ahead of the manifest ` +
|
|
||||||
`(${this.committed}) — ignoring it; recovery folds the whole log`
|
|
||||||
)
|
|
||||||
return null
|
|
||||||
}
|
|
||||||
return gen as number
|
|
||||||
} catch {
|
|
||||||
return null
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Record that an entity's canonical live bytes were (re)written and are not
|
|
||||||
* yet covered by a checkpoint stamp. Gated on chain validity so brains
|
|
||||||
* without a sound chain (tree authority, or log authority before its first
|
|
||||||
* recovery fold) never accumulate — they keep the fold-from-0 contract.
|
|
||||||
*/
|
|
||||||
private noteCheckpointDirty(kind: 'noun' | 'verb', id: string): void {
|
|
||||||
if (!this.foldCheckpointChainValid) return
|
|
||||||
if (kind === 'verb') this.checkpointDirtyVerbs.add(id)
|
|
||||||
else this.checkpointDirtyNouns.add(id)
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* The canonical-sync barrier + checkpoint stamp (must run under the commit
|
|
||||||
* mutex or in single-threaded open). Drains the dirty accumulator, makes
|
|
||||||
* those entities' canonical bytes durable via the adapter barrier, and only
|
|
||||||
* THEN stamps `_system/fold-checkpoint.json` at the committed watermark —
|
|
||||||
* stamp-after-data, always. On any failure the drained ids merge back and
|
|
||||||
* the stored checkpoint stays where it was: the bound can lag (a bigger
|
|
||||||
* fold later) but can never overstate durability (a lost write, outlawed).
|
|
||||||
*/
|
|
||||||
private async advanceFoldCheckpointUnlocked(): Promise<void> {
|
|
||||||
if (!this.foldCheckpointChainValid || !this.authorityIsLog || !this.factLog) return
|
|
||||||
const nouns = [...this.checkpointDirtyNouns]
|
|
||||||
const verbs = [...this.checkpointDirtyVerbs]
|
|
||||||
const target = this.committed
|
|
||||||
if (nouns.length === 0 && verbs.length === 0 && target === this.foldCheckpoint) return
|
|
||||||
this.checkpointDirtyNouns = new Set()
|
|
||||||
this.checkpointDirtyVerbs = new Set()
|
|
||||||
try {
|
|
||||||
if (nouns.length > 0 || verbs.length > 0) {
|
|
||||||
await this.storage.syncEntityCanonical?.(nouns, verbs)
|
|
||||||
}
|
|
||||||
await this.storage.writeRawObject(FOLD_CHECKPOINT_PATH, { generation: target })
|
|
||||||
await this.storage.syncRawObjects([FOLD_CHECKPOINT_PATH])
|
|
||||||
this.foldCheckpoint = target
|
|
||||||
} catch (err) {
|
|
||||||
for (const id of nouns) this.checkpointDirtyNouns.add(id)
|
|
||||||
for (const id of verbs) this.checkpointDirtyVerbs.add(id)
|
|
||||||
prodLog.warn(
|
|
||||||
`[GenerationStore] fold-checkpoint barrier failed at generation ${target} ` +
|
|
||||||
`(${(err as Error).message}) — checkpoint stays at ${this.foldCheckpoint}; ` +
|
|
||||||
`recovery would fold from there (bigger, never lossy). Will retry next flush.`
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @description Public, mutex-serialized fold-checkpoint advance — called by
|
|
||||||
* `close()` after the final flush so entities touched by paths that do not
|
|
||||||
* ride the pending tier (e.g. `transact()`) are covered before the
|
|
||||||
* clean-shutdown marker is written.
|
|
||||||
*/
|
|
||||||
async advanceFoldCheckpoint(): Promise<void> {
|
|
||||||
return this.withMutex(() => this.advanceFoldCheckpointUnlocked())
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @description Adoption-time chain bootstrap, phase 1 — called by
|
|
||||||
* `adoptLogAuthority()` BEFORE its oracle/backfill passes. Only a FRESH
|
|
||||||
* brain (committed === 0) may bootstrap here: with no committed
|
|
||||||
* generations the chain's assertion is vacuously true, and arming it now
|
|
||||||
* means the baseline backfill's own re-commits feed the dirty accumulator,
|
|
||||||
* so the first stamp after the flip covers them. A non-fresh flip skips
|
|
||||||
* this (returns false) — its chain starts at the brain's first recovery
|
|
||||||
* fold instead, because only a whole-log fold can prove coverage of
|
|
||||||
* entities written before the log existed.
|
|
||||||
*/
|
|
||||||
beginFoldCheckpointBootstrap(): boolean {
|
|
||||||
if (this.committed !== 0 || this.foldCheckpointChainValid) {
|
|
||||||
return this.foldCheckpointChainValid
|
|
||||||
}
|
|
||||||
this.foldCheckpointChainValid = true
|
|
||||||
this.foldCheckpoint = 0
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @description Adoption-time chain bootstrap, phase 2 — called after
|
|
||||||
* `flipToLogAuthority` records the flip. Opens the stamp gate; the next
|
|
||||||
* flush/close barrier writes the first checkpoint.
|
|
||||||
*/
|
|
||||||
completeFoldCheckpointBootstrap(): void {
|
|
||||||
this.authorityIsLog = true
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @description Adoption-time chain bootstrap, abort — called when an
|
|
||||||
* adoption attempt throws or refuses after phase 1. Disarms the chain and
|
|
||||||
* drops the accumulator so a tree-authority brain never accumulates or
|
|
||||||
* stamps. (If the chain was valid BEFORE the attempt — a stored checkpoint
|
|
||||||
* exists — it stays valid; only a phase-1 arm is undone.)
|
|
||||||
*/
|
|
||||||
abandonFoldCheckpointBootstrap(): void {
|
|
||||||
if (this.authorityIsLog) return
|
|
||||||
this.foldCheckpointChainValid = false
|
|
||||||
this.checkpointDirtyNouns = new Set()
|
|
||||||
this.checkpointDirtyVerbs = new Set()
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @description TEST-ONLY: install (or clear, with `undefined`) a fault
|
* @description TEST-ONLY: install (or clear, with `undefined`) a fault
|
||||||
* injector that is invoked at each {@link CommitFaultPhase} of the commit
|
* injector that is invoked at each {@link CommitFaultPhase} of the commit
|
||||||
|
|
@ -1438,13 +1238,6 @@ export class GenerationStore {
|
||||||
this.historyBytesTotal += delta.bytes ?? 0
|
this.historyBytesTotal += delta.bytes ?? 0
|
||||||
}
|
}
|
||||||
this.extendChains(gen, nouns, verbs)
|
this.extendChains(gen, nouns, verbs)
|
||||||
// Fold-checkpoint accounting: the write barrier above already synced
|
|
||||||
// this batch's canonical footprint on adapters that have one, but the
|
|
||||||
// accumulator entry is the belt — an adapter without a write barrier
|
|
||||||
// still gets these ids covered by the next checkpoint barrier, and a
|
|
||||||
// redundant fsync of already-durable bytes is cheap and idempotent.
|
|
||||||
for (const id of nouns) this.noteCheckpointDirty('noun', id)
|
|
||||||
for (const id of verbs) this.noteCheckpointDirty('verb', id)
|
|
||||||
const logEntry: TxLogEntry = { generation: gen, timestamp, ...(args.meta && { meta: args.meta }) }
|
const logEntry: TxLogEntry = { generation: gen, timestamp, ...(args.meta && { meta: args.meta }) }
|
||||||
await this.storage.appendTxLogLine(JSON.stringify(logEntry))
|
await this.storage.appendTxLogLine(JSON.stringify(logEntry))
|
||||||
|
|
||||||
|
|
@ -1456,13 +1249,6 @@ export class GenerationStore {
|
||||||
if (crashSimulated) {
|
if (crashSimulated) {
|
||||||
throw err
|
throw err
|
||||||
}
|
}
|
||||||
// Fold-checkpoint accounting: an abort's rollback restores are raw
|
|
||||||
// canonical writes that never reach the transaction write barrier
|
|
||||||
// (flushWriteBarrier only runs on the commit path) — feed them so the
|
|
||||||
// next checkpoint barrier syncs the restored bytes before any stamp
|
|
||||||
// vouches for them.
|
|
||||||
for (const id of nouns) this.noteCheckpointDirty('noun', id)
|
|
||||||
for (const id of verbs) this.noteCheckpointDirty('verb', id)
|
|
||||||
// The trapdoor for a batch: if rollback FAILED to fully apply, canonical
|
// The trapdoor for a batch: if rollback FAILED to fully apply, canonical
|
||||||
// storage may be inconsistent. A batch is never adopted forward (its
|
// storage may be inconsistent. A batch is never adopted forward (its
|
||||||
// other ops were rolled back — partial commit would break atomicity), so
|
// other ops were rolled back — partial commit would break atomicity), so
|
||||||
|
|
@ -1690,13 +1476,6 @@ export class GenerationStore {
|
||||||
await args.execute()
|
await args.execute()
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
this.inTransact = false
|
this.inTransact = false
|
||||||
// Fold-checkpoint accounting: execute() ran, so canonical bytes for
|
|
||||||
// the touched ids changed — whether they now hold the new images, a
|
|
||||||
// restored rollback, or (the trapdoor) something indeterminate, the
|
|
||||||
// next checkpoint stamp must not assert their durability without a
|
|
||||||
// barrier over whatever is actually there.
|
|
||||||
for (const id of nouns) this.noteCheckpointDirty('noun', id)
|
|
||||||
for (const id of verbs) this.noteCheckpointDirty('verb', id)
|
|
||||||
// A failed rollback (TransactionRollbackError) may have left canonical
|
// A failed rollback (TransactionRollbackError) may have left canonical
|
||||||
// storage inconsistent — the trapdoor. Reconcile against the
|
// storage inconsistent — the trapdoor. Reconcile against the
|
||||||
// before-images to decide the honest response (David's ruling:
|
// before-images to decide the honest response (David's ruling:
|
||||||
|
|
@ -1751,10 +1530,6 @@ export class GenerationStore {
|
||||||
throw err
|
throw err
|
||||||
}
|
}
|
||||||
this.inTransact = false
|
this.inTransact = false
|
||||||
// Fold-checkpoint accounting: the live canonical write is applied — it
|
|
||||||
// must ride the next canonical-sync barrier before any stamp covers it.
|
|
||||||
for (const id of nouns) this.noteCheckpointDirty('noun', id)
|
|
||||||
for (const id of verbs) this.noteCheckpointDirty('verb', id)
|
|
||||||
// Test-only crash simulation (direct call — a throw propagates with no
|
// Test-only crash simulation (direct call — a throw propagates with no
|
||||||
// cleanup, exactly like a process death; recovery-on-open restores the
|
// cleanup, exactly like a process death; recovery-on-open restores the
|
||||||
// contract). A crash here must cost only the never-returned ack: the
|
// contract). A crash here must cost only the never-returned ack: the
|
||||||
|
|
@ -1780,21 +1555,6 @@ export class GenerationStore {
|
||||||
// the log's group-commit (many concurrent writers share ONE sync) —
|
// the log's group-commit (many concurrent writers share ONE sync) —
|
||||||
// an acked write's fact survives power loss, by contract.
|
// an acked write's fact survives power loss, by contract.
|
||||||
if (this.factLog) {
|
if (this.factLog) {
|
||||||
// TWO PHASES, TWO DISTINCT COMPENSATIONS (a production adoption
|
|
||||||
// proved the difference the hard way): rewinding the counter after
|
|
||||||
// a SUCCESSFUL append re-mints the same generation and every later
|
|
||||||
// append refuses non-monotonic — the write path wedges in a refusal
|
|
||||||
// loop. The counter may only rewind when the log provably does NOT
|
|
||||||
// carry the generation.
|
|
||||||
const unbuffer = (): void => {
|
|
||||||
this.pendingBuffer.delete(gen)
|
|
||||||
const idx = this.pendingGens.lastIndexOf(gen)
|
|
||||||
if (idx !== -1) this.pendingGens.splice(idx, 1)
|
|
||||||
this.invalidateChains()
|
|
||||||
}
|
|
||||||
// Phase 1 — APPEND. Failure = the log never took the fact: full
|
|
||||||
// compensation (un-buffer + counter rewind); a rejected write must
|
|
||||||
// not commit, and the next mint may safely reuse the number.
|
|
||||||
try {
|
try {
|
||||||
await this.factLog.append(
|
await this.factLog.append(
|
||||||
await this.buildCommitFact({
|
await this.buildCommitFact({
|
||||||
|
|
@ -1805,37 +1565,24 @@ export class GenerationStore {
|
||||||
...(args.records && args.records.length > 0 ? { records: args.records } : {})
|
...(args.records && args.records.length > 0 ? { records: args.records } : {})
|
||||||
})
|
})
|
||||||
)
|
)
|
||||||
|
if (this.logDurability === 'at-ack') {
|
||||||
|
await this.factLog.ensureSynced()
|
||||||
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
unbuffer()
|
// A rejected write must NOT commit: the generation was buffered
|
||||||
|
// before the append, so un-buffer it and return the counter
|
||||||
|
// reservation — otherwise the next flush would durably commit a
|
||||||
|
// generation with NO fact, a silent log gap a later replay would
|
||||||
|
// turn into loss. Canonical bytes from execute() remain as an
|
||||||
|
// uncommitted orphan — identical to a crash at this point; never
|
||||||
|
// a torn committed state.
|
||||||
|
this.pendingBuffer.delete(gen)
|
||||||
|
const idx = this.pendingGens.lastIndexOf(gen)
|
||||||
|
if (idx !== -1) this.pendingGens.splice(idx, 1)
|
||||||
|
this.invalidateChains()
|
||||||
if (this.counter === gen) this.counter = gen - 1
|
if (this.counter === gen) this.counter = gen - 1
|
||||||
throw err
|
throw err
|
||||||
}
|
}
|
||||||
// Phase 2 — the at-ack covering sync. Failure here means the fact
|
|
||||||
// IS in the log (append succeeded) but durability was not promised:
|
|
||||||
// try to remove it (dropAbove); only a SUCCESSFUL drop earns the
|
|
||||||
// counter rewind. If the drop itself fails (e.g. the fact was
|
|
||||||
// sealed by a racing rotation), the generation stays consumed and
|
|
||||||
// buffered — monotonicity holds, the flush path retries durability,
|
|
||||||
// and the caller still gets the loud failure.
|
|
||||||
if (this.logDurability === 'at-ack') {
|
|
||||||
try {
|
|
||||||
await this.factLog.ensureSynced()
|
|
||||||
} catch (err) {
|
|
||||||
try {
|
|
||||||
await this.factLog.dropAbove(gen - 1)
|
|
||||||
unbuffer()
|
|
||||||
if (this.counter === gen) this.counter = gen - 1
|
|
||||||
} catch (dropErr) {
|
|
||||||
prodLog.warn(
|
|
||||||
`[GenerationStore] at-ack sync failed for generation ${gen} and the ` +
|
|
||||||
`appended fact could not be dropped (${(dropErr as Error).message}) — ` +
|
|
||||||
`the generation stays consumed and buffered; the flush path retries ` +
|
|
||||||
`durability. Never re-minting a number the log may carry.`
|
|
||||||
)
|
|
||||||
}
|
|
||||||
throw err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
// Test-only crash simulation. A crash here must cost the buffered
|
// Test-only crash simulation. A crash here must cost the buffered
|
||||||
// history + the appended fact in 'deferred' mode (open() truncates it
|
// history + the appended fact in 'deferred' mode (open() truncates it
|
||||||
|
|
@ -2026,16 +1773,6 @@ export class GenerationStore {
|
||||||
for (const entry of logEntries) {
|
for (const entry of logEntries) {
|
||||||
await this.storage.appendTxLogLine(JSON.stringify(entry))
|
await this.storage.appendTxLogLine(JSON.stringify(entry))
|
||||||
}
|
}
|
||||||
|
|
||||||
// Fold-checkpoint barrier: the window's LIVE canonical bytes (the acked
|
|
||||||
// writes themselves — the staging sync above covered only their history
|
|
||||||
// copies) become durable here, and only then does the checkpoint stamp
|
|
||||||
// advance to the new committed watermark. This is what keeps crash
|
|
||||||
// recovery's log fold bounded to (checkpoint, head] instead of the
|
|
||||||
// whole log. A failure inside is absorbed by the barrier (it warns,
|
|
||||||
// retains the accumulator, and leaves the old bound standing) — history
|
|
||||||
// durability above already succeeded, so the flush itself is good.
|
|
||||||
await this.advanceFoldCheckpointUnlocked()
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -3127,28 +2864,6 @@ export class GenerationStore {
|
||||||
* are never reissued.
|
* are never reissued.
|
||||||
* @param floorGeneration - The counter value before the restore.
|
* @param floorGeneration - The counter value before the restore.
|
||||||
*/
|
*/
|
||||||
/**
|
|
||||||
* @description Run a wholesale state replacement (the restore swap)
|
|
||||||
* EXCLUSIVELY: under the commit mutex, with the pending flush timer
|
|
||||||
* disarmed and the pending tier + fold-checkpoint accumulator discarded
|
|
||||||
* FIRST — so no background flush can write into `_system/` while the
|
|
||||||
* replacement is removing and swapping directories. Observed without this:
|
|
||||||
* a checkpoint stamp raced restore's directory removal and the swap died
|
|
||||||
* ENOTEMPTY mid-flight. The discarded in-memory state describes the store
|
|
||||||
* being replaced — `reopenAfterRestore` (which the caller runs next)
|
|
||||||
* rebuilds everything from the restored bytes.
|
|
||||||
*/
|
|
||||||
async runStateReplacement(replace: () => Promise<void>): Promise<void> {
|
|
||||||
return this.withMutex(async () => {
|
|
||||||
this.clearPendingFlushTimer()
|
|
||||||
this.pendingGens = []
|
|
||||||
this.pendingBuffer.clear()
|
|
||||||
this.checkpointDirtyNouns = new Set()
|
|
||||||
this.checkpointDirtyVerbs = new Set()
|
|
||||||
await replace()
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
async reopenAfterRestore(floorGeneration: number): Promise<void> {
|
async reopenAfterRestore(floorGeneration: number): Promise<void> {
|
||||||
await this.withMutex(async () => {
|
await this.withMutex(async () => {
|
||||||
this.deltaCache.clear()
|
this.deltaCache.clear()
|
||||||
|
|
@ -3161,27 +2876,6 @@ export class GenerationStore {
|
||||||
this.clearPendingFlushTimer()
|
this.clearPendingFlushTimer()
|
||||||
this.pendingGens = []
|
this.pendingGens = []
|
||||||
this.pendingBuffer.clear()
|
this.pendingBuffer.clear()
|
||||||
// The fold-checkpoint accumulator described the replaced state too.
|
|
||||||
this.checkpointDirtyNouns = new Set()
|
|
||||||
this.checkpointDirtyVerbs = new Set()
|
|
||||||
this.foldCheckpointChainValid = false
|
|
||||||
this.foldCheckpoint = 0
|
|
||||||
// A RESTORE IS AN UNCLEAN EVENT, by construction: the snapshot's files
|
|
||||||
// were just bulk-copied WITHOUT per-file fsync, so a power cut here can
|
|
||||||
// tear them — yet the snapshot may CARRY the source brain's
|
|
||||||
// clean-shutdown marker and fold checkpoint, which would together
|
|
||||||
// suppress exactly the recovery fold that cures such a tear. Delete
|
|
||||||
// both BEFORE reopening: the open below then treats the store as
|
|
||||||
// uncleanly shut, folds the restored log into canonical, barrier-syncs
|
|
||||||
// what it re-applied, and stamps a FRESH checkpoint — the restored
|
|
||||||
// state becomes durably founded at restore time instead of inheriting
|
|
||||||
// the source brain's assertions about bytes this disk never synced.
|
|
||||||
try {
|
|
||||||
await this.storage.deleteRawObject(CLEAN_SHUTDOWN_PATH)
|
|
||||||
} catch { /* absent is fine — same outcome */ }
|
|
||||||
try {
|
|
||||||
await this.storage.deleteRawObject(FOLD_CHECKPOINT_PATH)
|
|
||||||
} catch { /* absent is fine — fold from 0 */ }
|
|
||||||
this.opened = false
|
this.opened = false
|
||||||
// open() re-reads counter/manifest and re-registers the bump hook.
|
// open() re-reads counter/manifest and re-registers the bump hook.
|
||||||
await this.open()
|
await this.open()
|
||||||
|
|
@ -3239,8 +2933,6 @@ export class GenerationStore {
|
||||||
private async rollBackUncommittedGeneration(gen: number): Promise<void> {
|
private async rollBackUncommittedGeneration(gen: number): Promise<void> {
|
||||||
const dir = `${GENERATIONS_PREFIX}/${gen}`
|
const dir = `${GENERATIONS_PREFIX}/${gen}`
|
||||||
const prevPaths = await this.storage.listRawObjects(`${dir}/prev`)
|
const prevPaths = await this.storage.listRawObjects(`${dir}/prev`)
|
||||||
const restoredNouns: string[] = []
|
|
||||||
const restoredVerbs: string[] = []
|
|
||||||
for (const recordPath of prevPaths) {
|
for (const recordPath of prevPaths) {
|
||||||
const id = recordIdFromPath(recordPath)
|
const id = recordIdFromPath(recordPath)
|
||||||
if (id === null) continue
|
if (id === null) continue
|
||||||
|
|
@ -3249,20 +2941,10 @@ export class GenerationStore {
|
||||||
const image = { metadata: record.metadata, vector: record.vector }
|
const image = { metadata: record.metadata, vector: record.vector }
|
||||||
if (record.kind === 'verb') {
|
if (record.kind === 'verb') {
|
||||||
await this.storage.writeVerbRaw(id, image)
|
await this.storage.writeVerbRaw(id, image)
|
||||||
restoredVerbs.push(id)
|
|
||||||
} else {
|
} else {
|
||||||
await this.storage.writeNounRaw(id, image)
|
await this.storage.writeNounRaw(id, image)
|
||||||
restoredNouns.push(id)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Make the restores durable IMMEDIATELY (this runs at open, before the
|
|
||||||
// fold-checkpoint chain state is even read): a restored before-image
|
|
||||||
// replaces bytes a stored checkpoint may already vouch for, so it must
|
|
||||||
// reach disk with the same certainty — otherwise a power cut could let
|
|
||||||
// the rolled-back write's bytes resurrect past a bounded fold.
|
|
||||||
if (restoredNouns.length > 0 || restoredVerbs.length > 0) {
|
|
||||||
await this.storage.syncEntityCanonical?.(restoredNouns, restoredVerbs)
|
|
||||||
}
|
|
||||||
await this.storage.removeRawPrefix(dir)
|
await this.storage.removeRawPrefix(dir)
|
||||||
prodLog.warn(
|
prodLog.warn(
|
||||||
`[GenerationStore] rolled back uncommitted generation ${gen} ` +
|
`[GenerationStore] rolled back uncommitted generation ${gen} ` +
|
||||||
|
|
|
||||||
|
|
@ -462,18 +462,6 @@ export interface GenerationStorage {
|
||||||
/** @see beginWriteBarrier — fsync every canonical write since begin. */
|
/** @see beginWriteBarrier — fsync every canonical write since begin. */
|
||||||
flushWriteBarrier?(): Promise<void>
|
flushWriteBarrier?(): Promise<void>
|
||||||
|
|
||||||
/**
|
|
||||||
* OPTIONAL fold-checkpoint durability barrier: make the listed entities'
|
|
||||||
* CANONICAL live objects durable — fsync each present metadata/vector file
|
|
||||||
* AND the parent directory entry of each absent one (so a delete is as
|
|
||||||
* durable as a write). The generation store may only advance the fold
|
|
||||||
* checkpoint (`_system/fold-checkpoint.json`) after this resolves; the
|
|
||||||
* checkpoint bounds crash recovery's log fold to `(checkpoint, head]`.
|
|
||||||
* Adapters whose writes are durable per-call may leave this undefined —
|
|
||||||
* the store then treats canonical durability as immediate.
|
|
||||||
*/
|
|
||||||
syncEntityCanonical?(nouns: string[], verbs: string[]): Promise<void>
|
|
||||||
|
|
||||||
/** Read an entity's raw stored metadata+vector objects. */
|
/** Read an entity's raw stored metadata+vector objects. */
|
||||||
readNounRaw(id: string): Promise<{ metadata: any | null; vector: any | null }>
|
readNounRaw(id: string): Promise<{ metadata: any | null; vector: any | null }>
|
||||||
/** Restore an entity's raw stored objects (`null` part ⇒ delete that file). */
|
/** Restore an entity's raw stored objects (`null` part ⇒ delete that file). */
|
||||||
|
|
|
||||||
|
|
@ -799,7 +799,6 @@ export class FileSystemStorage extends BaseStorage {
|
||||||
|
|
||||||
for (const objectPath of paths) {
|
for (const objectPath of paths) {
|
||||||
const fullPath = path.join(this.rootDir, objectPath)
|
const fullPath = path.join(this.rootDir, objectPath)
|
||||||
let synced = false
|
|
||||||
for (const candidate of [`${fullPath}.gz`, fullPath]) {
|
for (const candidate of [`${fullPath}.gz`, fullPath]) {
|
||||||
let handle: any
|
let handle: any
|
||||||
try {
|
try {
|
||||||
|
|
@ -814,14 +813,8 @@ export class FileSystemStorage extends BaseStorage {
|
||||||
await handle.close()
|
await handle.close()
|
||||||
}
|
}
|
||||||
parentDirs.add(path.dirname(fullPath))
|
parentDirs.add(path.dirname(fullPath))
|
||||||
synced = true
|
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
// An absent path is a state too: fsync the parent directory so a
|
|
||||||
// completed unlink is durable (a delete must survive power loss as
|
|
||||||
// surely as a write — otherwise a bounded log fold could let a
|
|
||||||
// tombstoned record resurrect from a lost directory update).
|
|
||||||
if (!synced) parentDirs.add(path.dirname(fullPath))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
for (const dir of parentDirs) {
|
for (const dir of parentDirs) {
|
||||||
|
|
|
||||||
|
|
@ -1390,29 +1390,6 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
||||||
void paths
|
void paths
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Fold-checkpoint durability barrier: make the listed entities' canonical
|
|
||||||
* live objects durable. Maps each id to its canonical metadata + vector
|
|
||||||
* paths and delegates to {@link BaseStorage.syncRawObjects}, whose
|
|
||||||
* filesystem override fsyncs present files (and their rename directory
|
|
||||||
* entries) and the parent directory of absent ones — so deletes are as
|
|
||||||
* durable as writes. The generation store advances the fold checkpoint
|
|
||||||
* only after this resolves (stamp-after-data).
|
|
||||||
*
|
|
||||||
* @param nouns - Entity ids whose canonical objects must be durable.
|
|
||||||
* @param verbs - Relationship ids whose canonical objects must be durable.
|
|
||||||
*/
|
|
||||||
public async syncEntityCanonical(nouns: string[], verbs: string[]): Promise<void> {
|
|
||||||
const paths: string[] = []
|
|
||||||
for (const id of nouns) {
|
|
||||||
paths.push(getNounMetadataPath(id), getNounVectorPath(id))
|
|
||||||
}
|
|
||||||
for (const id of verbs) {
|
|
||||||
paths.push(getVerbMetadataPath(id), getVerbVectorPath(id))
|
|
||||||
}
|
|
||||||
if (paths.length > 0) await this.syncRawObjects(paths)
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Read an entity's raw stored objects — the exact bytes at its canonical
|
* Read an entity's raw stored objects — the exact bytes at its canonical
|
||||||
* metadata + vector paths (write-cache coherent). Used by the generation
|
* metadata + vector paths (write-cache coherent). Used by the generation
|
||||||
|
|
|
||||||
|
|
@ -1908,35 +1908,6 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
||||||
* index (early-stop at `offset+limit`); the JS index returns ALL matches and lets
|
* index (early-stop at `offset+limit`); the JS index returns ALL matches and lets
|
||||||
* the caller window them, so `_opts` is intentionally ignored here.
|
* the caller window them, so `_opts` is intentionally ignored here.
|
||||||
*/
|
*/
|
||||||
/** Once-per-field throttle for the sparse-store did-you-mean WARN. */
|
|
||||||
private readonly warnedNeverCarried = new Set<string>()
|
|
||||||
|
|
||||||
/**
|
|
||||||
* THE SPARSE-STORE CUT (ruled 2026-08-12): a WHERE filter naming a field
|
|
||||||
* no row carries is SERVED OPERATOR-TRUTHFULLY (eq/range/contains → [];
|
|
||||||
* ne/exists:false → all rows; exists:true → []) — the JS evaluator below
|
|
||||||
* already computes exactly these truths via complements — with the
|
|
||||||
* did-you-mean demoted to this throttled WARN. A fresh store's first
|
|
||||||
* filtered read is a correct empty answer, never a refusal. orderBy and
|
|
||||||
* ambiguous addresses KEEP their hard refusals (no truthful order
|
|
||||||
* exists; ambiguity is a contract error — absence is data).
|
|
||||||
*/
|
|
||||||
/** Is this field known to the index at all (any row ever carried it)? */
|
|
||||||
private fieldRegistryHas(field: string): boolean {
|
|
||||||
return this.fieldStats.has(field)
|
|
||||||
}
|
|
||||||
|
|
||||||
private warnNeverCarriedOnce(field: string): void {
|
|
||||||
if (this.warnedNeverCarried.has(field)) return
|
|
||||||
this.warnedNeverCarried.add(field)
|
|
||||||
prodLog.warn(
|
|
||||||
`[MetadataIndex] filter names field '${field}' which no row carries — ` +
|
|
||||||
`serving the operator-truthful answer (empty for positive matches; ` +
|
|
||||||
`the complement for ne/exists:false). If this is a typo, check the ` +
|
|
||||||
`field name; refusals remain on orderBy.`
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
async getIdsForFilter(filter: any, _opts?: { limit?: number; offset?: number }): Promise<string[]> {
|
async getIdsForFilter(filter: any, _opts?: { limit?: number; offset?: number }): Promise<string[]> {
|
||||||
if (!filter || Object.keys(filter).length === 0) {
|
if (!filter || Object.keys(filter).length === 0) {
|
||||||
return []
|
return []
|
||||||
|
|
@ -2013,17 +1984,6 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
||||||
const address = parseFieldAddress(rawField, 'entity')
|
const address = parseFieldAddress(rawField, 'entity')
|
||||||
const field = address.scope === 'system' ? `system.${address.field}` : address.field
|
const field = address.scope === 'system' ? `system.${address.field}` : address.field
|
||||||
|
|
||||||
// Sparse-store cut: a user field no row carries serves operator-truth
|
|
||||||
// below (the evaluators' complements are already correct) — announce
|
|
||||||
// it once so a typo is findable without breaking a fresh store.
|
|
||||||
if (
|
|
||||||
address.scope !== 'system' &&
|
|
||||||
!(this.columnStore && this.columnStore.hasField(field)) &&
|
|
||||||
!this.fieldRegistryHas(field)
|
|
||||||
) {
|
|
||||||
this.warnNeverCarriedOnce(field)
|
|
||||||
}
|
|
||||||
|
|
||||||
let fieldResults: string[] = []
|
let fieldResults: string[] = []
|
||||||
|
|
||||||
try {
|
try {
|
||||||
|
|
@ -2062,18 +2022,7 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
||||||
// complement as a bitmap difference over the int-id universe rather
|
// complement as a bitmap difference over the int-id universe rather
|
||||||
// than materializing the whole corpus as UUID strings to filter it.
|
// than materializing the whole corpus as UUID strings to filter it.
|
||||||
const excludeInts: number[] = []
|
const excludeInts: number[] = []
|
||||||
// Sparse-store truth: a never-carried field has NOTHING to
|
for (const uuid of await this.getIds(field, operand)) {
|
||||||
// exclude — the complement of nothing is EVERYTHING. getIds
|
|
||||||
// throws FIELD_NOT_INDEXED there; the clause-level catch
|
|
||||||
// would wrongly zero this NEGATIVE operator, so absorb it
|
|
||||||
// here as the empty exclude set (the ruled operator-truth).
|
|
||||||
let neMatches: string[] = []
|
|
||||||
try {
|
|
||||||
neMatches = await this.getIds(field, operand)
|
|
||||||
} catch {
|
|
||||||
neMatches = []
|
|
||||||
}
|
|
||||||
for (const uuid of neMatches) {
|
|
||||||
const intId = this.idMapper.getInt(uuid)
|
const intId = this.idMapper.getInt(uuid)
|
||||||
if (intId !== undefined) excludeInts.push(intId)
|
if (intId !== undefined) excludeInts.push(intId)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,82 +0,0 @@
|
||||||
/**
|
|
||||||
* @module tests/conformance/sparse-store-cut
|
|
||||||
* @description THE SPARSE-STORE CUT (ruled 2026-08-12) — the shared
|
|
||||||
* conformance rows both engines run: a WHERE filter naming a field NO row
|
|
||||||
* carries is SERVED OPERATOR-TRUTHFULLY, never refused:
|
|
||||||
* eq / in / range / contains → [] (nothing carries it → nothing matches)
|
|
||||||
* ne / exists:false → ALL rows (equally true — blanket-empty here
|
|
||||||
* would be the outlawed silent wrong)
|
|
||||||
* exists:true → []
|
|
||||||
* orderBy on an unresolvable field KEEPS the hard refusal (no truthful
|
|
||||||
* order exists). The did-you-mean demotes to a throttled WARN on the serve.
|
|
||||||
* A fresh tenant's first filtered read is a correct empty answer — the
|
|
||||||
* 341-red first-adopter class, closed.
|
|
||||||
*/
|
|
||||||
import { describe, it, expect, afterEach } from 'vitest'
|
|
||||||
import { Brainy, UnresolvableFieldError } from '../../src/index.js'
|
|
||||||
import { NounType } from '../../src/types/graphTypes.js'
|
|
||||||
|
|
||||||
const brains: Brainy[] = []
|
|
||||||
afterEach(async () => {
|
|
||||||
for (const b of brains.splice(0)) await b.close().catch(() => {})
|
|
||||||
})
|
|
||||||
|
|
||||||
async function corpus(): Promise<{ brain: Brainy; ids: string[] }> {
|
|
||||||
const b = new Brainy({ storage: { type: 'memory' }, requireSubtype: false })
|
|
||||||
await b.init()
|
|
||||||
brains.push(b)
|
|
||||||
const ids: string[] = []
|
|
||||||
for (let i = 0; i < 4; i++) {
|
|
||||||
ids.push(
|
|
||||||
await b.add({ data: `row ${i}`, type: NounType.Document, metadata: { carried: i } })
|
|
||||||
)
|
|
||||||
}
|
|
||||||
return { brain: b, ids }
|
|
||||||
}
|
|
||||||
|
|
||||||
describe('sparse-store cut — operator-truthful serve on never-carried fields', () => {
|
|
||||||
it('positive matches serve EMPTY: eq, in, range, contains', async () => {
|
|
||||||
const { brain } = await corpus()
|
|
||||||
expect(await brain.find({ where: { ghost: 'x' }, limit: 10 })).toEqual([])
|
|
||||||
expect(await brain.find({ where: { ghost: { in: ['a', 'b'] } }, limit: 10 })).toEqual([])
|
|
||||||
expect(await brain.find({ where: { ghost: { gt: 5 } }, limit: 10 })).toEqual([])
|
|
||||||
expect(await brain.find({ where: { ghost: { exists: true } }, limit: 10 })).toEqual([])
|
|
||||||
})
|
|
||||||
|
|
||||||
it('negative matches serve ALL rows: ne and exists:false (the truth, not blanket-empty)', async () => {
|
|
||||||
const { brain, ids } = await corpus()
|
|
||||||
const ne = await brain.find({ where: { ghost: { ne: 'x' } }, limit: 10 })
|
|
||||||
expect(ne.map((r) => r.id).sort()).toEqual([...ids].sort())
|
|
||||||
const absent = await brain.find({ where: { ghost: { exists: false } }, limit: 10 })
|
|
||||||
expect(absent.map((r) => r.id).sort()).toEqual([...ids].sort())
|
|
||||||
})
|
|
||||||
|
|
||||||
it('the fresh-tenant day-one shape: an EMPTY store answers its first filtered read with [], never a refusal', async () => {
|
|
||||||
const b = new Brainy({ storage: { type: 'memory' }, requireSubtype: false })
|
|
||||||
await b.init()
|
|
||||||
brains.push(b)
|
|
||||||
expect(await b.find({ where: { status: 'open' }, limit: 50 })).toEqual([])
|
|
||||||
expect(await b.find({ where: { date: { gte: '2026-01-01' } }, limit: 50 })).toEqual([])
|
|
||||||
})
|
|
||||||
|
|
||||||
it('orderBy on an unresolvable field KEEPS the typed refusal', async () => {
|
|
||||||
const { brain } = await corpus()
|
|
||||||
await expect(
|
|
||||||
brain.find({ where: { carried: { gte: 0 } }, orderBy: 'system.notAScalar', limit: 10 })
|
|
||||||
).rejects.toThrow(UnresolvableFieldError)
|
|
||||||
})
|
|
||||||
|
|
||||||
it('compound filters: the never-carried clause composes truthfully with carried clauses', async () => {
|
|
||||||
const { brain, ids } = await corpus()
|
|
||||||
// carried>=2 AND ghost ne 'x' → the carried>=2 rows (ne-clause = all).
|
|
||||||
const both = await brain.find({
|
|
||||||
where: { carried: { gte: 2 }, ghost: { ne: 'x' } },
|
|
||||||
limit: 10
|
|
||||||
})
|
|
||||||
expect(both.map((r) => r.id).sort()).toEqual([ids[2], ids[3]].sort())
|
|
||||||
// carried>=2 AND ghost eq 'x' → [] (eq-clause empties the intersection).
|
|
||||||
expect(
|
|
||||||
await brain.find({ where: { carried: { gte: 2 }, ghost: 'x' }, limit: 10 })
|
|
||||||
).toEqual([])
|
|
||||||
})
|
|
||||||
})
|
|
||||||
|
|
@ -1,235 +0,0 @@
|
||||||
/**
|
|
||||||
* @module tests/integration/fold-checkpoint-bound
|
|
||||||
* @description The fold-checkpoint bound (crash recovery's log fold, bounded):
|
|
||||||
* `_system/fold-checkpoint.json` at generation G asserts every entity whose
|
|
||||||
* latest fact is ≤ G has DURABLE canonical bytes — each stamp strictly follows
|
|
||||||
* a canonical-sync barrier over every live entity touched since the last one
|
|
||||||
* (stamp-after-data). An unclean open then folds only `(G, head]` instead of
|
|
||||||
* the whole log. These pins prove the four load-bearing properties:
|
|
||||||
*
|
|
||||||
* 1. The stamp exists and tracks the committed watermark (flush + close).
|
|
||||||
* 2. The fold is genuinely BOUNDED — facts ≤ G are skipped — while facts in
|
|
||||||
* `(G, head]` are re-applied even BELOW the manifest.
|
|
||||||
* 3. A failed barrier NEVER advances the stamp (the bound can lag, growing
|
|
||||||
* a later fold — it can never overstate durability, losing a write).
|
|
||||||
* 4. A pre-checkpoint brain (the 10.0 shape) bootstraps its chain at its
|
|
||||||
* first whole-log fold; a tree-authority brain never stamps at all.
|
|
||||||
*/
|
|
||||||
import { describe, it, expect, afterEach, vi } from 'vitest'
|
|
||||||
import * as fs from 'node:fs'
|
|
||||||
import * as zlib from 'node:zlib'
|
|
||||||
import { join } from 'node:path'
|
|
||||||
import { Brainy } from '../../src/brainy.js'
|
|
||||||
import { NounType } from '../../src/types/graphTypes.js'
|
|
||||||
import {
|
|
||||||
abandonAsCrashed,
|
|
||||||
dropCanonicalNoun,
|
|
||||||
makeTempDir,
|
|
||||||
openBrain,
|
|
||||||
storeOf
|
|
||||||
} from '../helpers/durabilityKillMatrix.js'
|
|
||||||
|
|
||||||
const CHECKPOINT = join('_system', 'fold-checkpoint.json')
|
|
||||||
|
|
||||||
/** Read the fold-checkpoint artifact's generation from disk, or null. */
|
|
||||||
function readCheckpoint(dir: string): number | null {
|
|
||||||
for (const candidate of [join(dir, `${CHECKPOINT}.gz`), join(dir, CHECKPOINT)]) {
|
|
||||||
if (!fs.existsSync(candidate)) continue
|
|
||||||
const raw = fs.readFileSync(candidate)
|
|
||||||
const text = candidate.endsWith('.gz') ? zlib.gunzipSync(raw).toString('utf8') : raw.toString('utf8')
|
|
||||||
const parsed = JSON.parse(text) as { generation?: number }
|
|
||||||
return Number.isSafeInteger(parsed.generation) ? (parsed.generation as number) : null
|
|
||||||
}
|
|
||||||
return null
|
|
||||||
}
|
|
||||||
|
|
||||||
function removeArtifact(dir: string, rel: string): void {
|
|
||||||
for (const candidate of [join(dir, `${rel}.gz`), join(dir, rel)]) {
|
|
||||||
fs.rmSync(candidate, { force: true })
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
function committedOf(brain: Brainy): number {
|
|
||||||
return (storeOf(brain) as unknown as { committed: number }).committed
|
|
||||||
}
|
|
||||||
|
|
||||||
describe('fold-checkpoint bound — crash recovery folds (checkpoint, head], never less durability than stamped', () => {
|
|
||||||
const dirs: string[] = []
|
|
||||||
const liveBrains: Brainy[] = []
|
|
||||||
afterEach(async () => {
|
|
||||||
vi.restoreAllMocks()
|
|
||||||
for (const b of liveBrains.splice(0)) await b.close().catch(() => {})
|
|
||||||
for (const d of dirs.splice(0)) fs.rmSync(d, { recursive: true, force: true })
|
|
||||||
})
|
|
||||||
function trackDir(): string {
|
|
||||||
const dir = makeTempDir()
|
|
||||||
dirs.push(dir)
|
|
||||||
return dir
|
|
||||||
}
|
|
||||||
|
|
||||||
it('a fresh adopt brain stamps at flush and again at close — the stamp tracks the committed watermark', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
expect(brain.logAuthority().authority).toBe('log')
|
|
||||||
|
|
||||||
await brain.add({ data: 'first', type: NounType.Document, metadata: { n: 1 } })
|
|
||||||
await brain.add({ data: 'second', type: NounType.Document, metadata: { n: 2 } })
|
|
||||||
await brain.flush()
|
|
||||||
const afterFlush = readCheckpoint(dir)
|
|
||||||
expect(afterFlush).toBe(committedOf(brain))
|
|
||||||
expect(afterFlush!).toBeGreaterThan(0)
|
|
||||||
|
|
||||||
await brain.add({ data: 'third', type: NounType.Document, metadata: { n: 3 } })
|
|
||||||
const closingCommit = liveBrains.pop()!
|
|
||||||
await closingCommit.close()
|
|
||||||
// Close flushes, so the stamp advanced with it — and the clean-shutdown
|
|
||||||
// marker it writes afterward never vouches for bytes the stamp has not.
|
|
||||||
expect(readCheckpoint(dir)).toBeGreaterThanOrEqual(afterFlush!)
|
|
||||||
}, 120000)
|
|
||||||
|
|
||||||
it('BOUNDED fold: facts ≤ checkpoint are skipped, facts in (checkpoint, head] are re-applied even below the manifest; a failed barrier retains the old bound', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
|
|
||||||
// Window 1 — flushed and stamped: the checkpoint's covered past.
|
|
||||||
const idA = await brain.add({ data: 'covered by the stamp', type: NounType.Document, metadata: { w: 1 } })
|
|
||||||
await brain.flush()
|
|
||||||
const checkpoint1 = readCheckpoint(dir)
|
|
||||||
expect(checkpoint1).toBe(committedOf(brain))
|
|
||||||
|
|
||||||
// Window 2 — committed BELOW a new manifest but with the checkpoint stamp
|
|
||||||
// FAILING: the barrier throws once, so the manifest advances while the
|
|
||||||
// stamp stays at checkpoint1 (pin 3: a failed barrier never advances it).
|
|
||||||
const storage = (brain as unknown as {
|
|
||||||
storage: { syncEntityCanonical(n: string[], v: string[]): Promise<void> }
|
|
||||||
}).storage
|
|
||||||
const realBarrier = storage.syncEntityCanonical.bind(storage)
|
|
||||||
let failedOnce = false
|
|
||||||
vi.spyOn(storage, 'syncEntityCanonical').mockImplementation(async (n: string[], v: string[]) => {
|
|
||||||
if (!failedOnce) {
|
|
||||||
failedOnce = true
|
|
||||||
throw new Error('injected barrier failure (device hiccup)')
|
|
||||||
}
|
|
||||||
return realBarrier(n, v)
|
|
||||||
})
|
|
||||||
const idB = await brain.add({ data: 'below manifest, above checkpoint', type: NounType.Document, metadata: { w: 2 } })
|
|
||||||
await brain.flush()
|
|
||||||
expect(failedOnce).toBe(true)
|
|
||||||
expect(readCheckpoint(dir)).toBe(checkpoint1) // stamp did NOT advance
|
|
||||||
expect(committedOf(brain)).toBeGreaterThan(checkpoint1!) // manifest DID
|
|
||||||
|
|
||||||
// Crash. Vaporize BOTH canonical records: idB's fact lives in
|
|
||||||
// (checkpoint, manifest] — the bounded fold MUST restore it; idA's fact
|
|
||||||
// is ≤ checkpoint — the fold must SKIP it (its loss here is synthetic:
|
|
||||||
// the stamp's barrier fsynced it, a power cut cannot take it, and the
|
|
||||||
// skip is exactly what makes the fold bounded instead of whole-log).
|
|
||||||
await abandonAsCrashed(liveBrains.pop()!)
|
|
||||||
dropCanonicalNoun(dir, idA)
|
|
||||||
dropCanonicalNoun(dir, idB)
|
|
||||||
|
|
||||||
const reopened = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(reopened)
|
|
||||||
const restoredB = await reopened.get(idB)
|
|
||||||
expect(restoredB, 'a fact above the checkpoint is re-applied even below the manifest').not.toBeNull()
|
|
||||||
const skippedA = await reopened.get(idA)
|
|
||||||
expect(skippedA, 'a fact at-or-below the checkpoint is outside the fold — the bound is real').toBeNull()
|
|
||||||
// And recovery re-stamped at its new committed watermark.
|
|
||||||
expect(readCheckpoint(dir)).toBe(committedOf(reopened))
|
|
||||||
}, 120000)
|
|
||||||
|
|
||||||
it('a pre-checkpoint brain (the 10.0 shape) folds the WHOLE log once, then its chain is established', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
const idA = await brain.add({ data: 'ten-point-oh resident', type: NounType.Document, metadata: { era: '10.0' } })
|
|
||||||
await brain.flush()
|
|
||||||
await liveBrains.pop()!.close()
|
|
||||||
|
|
||||||
// Rewind the brain to the 10.0 shape: no checkpoint artifact, and an
|
|
||||||
// unclean shutdown (marker gone) — exactly what an existing fleet brain
|
|
||||||
// looks like at its first crash under 10.1.
|
|
||||||
removeArtifact(dir, CHECKPOINT)
|
|
||||||
removeArtifact(dir, join('_system', 'clean-shutdown.json'))
|
|
||||||
dropCanonicalNoun(dir, idA)
|
|
||||||
|
|
||||||
const reopened = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(reopened)
|
|
||||||
expect(await reopened.get(idA), 'no checkpoint ⇒ whole-log fold ⇒ every acked write restored').not.toBeNull()
|
|
||||||
const stamped = readCheckpoint(dir)
|
|
||||||
expect(stamped, 'the first whole-log fold is the chain’s base case — it stamps').toBe(committedOf(reopened))
|
|
||||||
}, 120000)
|
|
||||||
|
|
||||||
it('a tree-authority brain never stamps a checkpoint', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'defer' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
expect(brain.logAuthority().authority).not.toBe('log')
|
|
||||||
await brain.add({ data: 'tree resident', type: NounType.Document, metadata: { n: 1 } })
|
|
||||||
await brain.flush()
|
|
||||||
await liveBrains.pop()!.close()
|
|
||||||
expect(readCheckpoint(dir)).toBeNull()
|
|
||||||
}, 120000)
|
|
||||||
|
|
||||||
it('restore is an UNCLEAN event: the snapshot’s stamps do not survive — the reopen fold re-founds and re-stamps the restored state', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
const idA = await brain.add({ data: 'survives the restore', type: NounType.Document, metadata: { n: 1 } })
|
|
||||||
await brain.flush()
|
|
||||||
|
|
||||||
const snapDir = join(trackDir(), 'snap')
|
|
||||||
const db = brain.now()
|
|
||||||
await (db as unknown as { persist(p: string): Promise<void> }).persist(snapDir)
|
|
||||||
await (db as unknown as { release(): Promise<void> }).release()
|
|
||||||
|
|
||||||
// Advance the live brain past the snapshot: a later write, a later flush,
|
|
||||||
// a later checkpoint stamp — none of which may survive the restore.
|
|
||||||
const idB = await brain.add({ data: 'must not survive', type: NounType.Document, metadata: { n: 2 } })
|
|
||||||
await brain.flush()
|
|
||||||
const stampBeforeRestore = readCheckpoint(dir)
|
|
||||||
expect(stampBeforeRestore).toBe(committedOf(brain))
|
|
||||||
|
|
||||||
// Unflushed traffic in flight at restore time — the quiesced swap discards
|
|
||||||
// it under the mutex instead of letting its flush timer race the swap
|
|
||||||
// (the ENOTEMPTY class).
|
|
||||||
await brain.add({ data: 'in-flight at restore', type: NounType.Document, metadata: { n: 3 } })
|
|
||||||
await brain.restore(snapDir, { confirm: true })
|
|
||||||
|
|
||||||
expect(await brain.get(idA), 'snapshot state restored').not.toBeNull()
|
|
||||||
expect(await brain.get(idB), 'post-snapshot state replaced').toBeNull()
|
|
||||||
// The stamp on disk is the REOPEN FOLD's fresh assertion about the
|
|
||||||
// restored (and now barrier-synced) bytes — at the restored watermark,
|
|
||||||
// strictly below the pre-restore stamp that must not survive.
|
|
||||||
const stampAfterRestore = readCheckpoint(dir)
|
|
||||||
expect(stampAfterRestore).toBe(committedOf(brain))
|
|
||||||
expect(stampAfterRestore!).toBeLessThan(stampBeforeRestore!)
|
|
||||||
}, 120000)
|
|
||||||
|
|
||||||
it('a delete rides the barrier: the tombstoned id is in the synced set and the stamp advances past it', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
const id = await brain.add({ data: 'short-lived', type: NounType.Document, metadata: { n: 1 } })
|
|
||||||
await brain.flush()
|
|
||||||
|
|
||||||
const storage = (brain as unknown as {
|
|
||||||
storage: { syncEntityCanonical(n: string[], v: string[]): Promise<void> }
|
|
||||||
}).storage
|
|
||||||
const seen: string[][] = []
|
|
||||||
const realBarrier = storage.syncEntityCanonical.bind(storage)
|
|
||||||
vi.spyOn(storage, 'syncEntityCanonical').mockImplementation(async (n: string[], v: string[]) => {
|
|
||||||
seen.push([...n])
|
|
||||||
return realBarrier(n, v)
|
|
||||||
})
|
|
||||||
|
|
||||||
await brain.remove(id)
|
|
||||||
await brain.flush()
|
|
||||||
expect(
|
|
||||||
seen.some((nouns) => nouns.includes(id)),
|
|
||||||
'the deleted id must reach the canonical barrier (absence is durable state too)'
|
|
||||||
).toBe(true)
|
|
||||||
expect(readCheckpoint(dir)).toBe(committedOf(brain))
|
|
||||||
}, 120000)
|
|
||||||
})
|
|
||||||
|
|
@ -1,78 +0,0 @@
|
||||||
/**
|
|
||||||
* @module tests/integration/sync-fail-compensation
|
|
||||||
* @description The non-monotonic refusal-loop cure (a production adoption's
|
|
||||||
* second defect): when the at-ack covering SYNC fails AFTER a successful
|
|
||||||
* append, the counter must NOT rewind unless the appended fact is provably
|
|
||||||
* removed — rewinding while the log carries the generation re-mints the
|
|
||||||
* same number and every later append refuses non-monotonic, wedging the
|
|
||||||
* write path in a refusal loop ("writes REFUSED until it drains").
|
|
||||||
*/
|
|
||||||
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/index.js'
|
|
||||||
import { NounType } from '../../src/types/graphTypes.js'
|
|
||||||
|
|
||||||
const dirs: string[] = []
|
|
||||||
const brains: Brainy[] = []
|
|
||||||
afterEach(async () => {
|
|
||||||
vi.restoreAllMocks()
|
|
||||||
for (const b of brains.splice(0)) await b.close().catch(() => {})
|
|
||||||
for (const d of dirs.splice(0)) rmSync(d, { recursive: true, force: true })
|
|
||||||
})
|
|
||||||
|
|
||||||
describe('at-ack sync-failure compensation', () => {
|
|
||||||
it('a one-shot sync failure never wedges the write path: the next write mints a FRESH generation and succeeds', async () => {
|
|
||||||
const dir = mkdtempSync(join(tmpdir(), 'brainy-syncfail-'))
|
|
||||||
dirs.push(dir)
|
|
||||||
const brain = new Brainy({ storage: { type: 'filesystem', path: dir }, requireSubtype: false })
|
|
||||||
await brain.init() // adopt-default: log authority, at-ack
|
|
||||||
brains.push(brain)
|
|
||||||
expect(brain.logAuthority().authority).toBe('log')
|
|
||||||
await brain.add({ data: 'baseline', type: NounType.Document, metadata: { n: 0 } })
|
|
||||||
|
|
||||||
// Fail exactly ONE covering sync (after its append lands).
|
|
||||||
// Target ensureSynced (the ACK path's covering sync) — mocking sync()
|
|
||||||
// itself gets eaten by background flushes before the victim write.
|
|
||||||
const factLog = (brain as unknown as {
|
|
||||||
generationStore: { getFactLog(): { ensureSynced(): Promise<void> } }
|
|
||||||
}).generationStore.getFactLog()
|
|
||||||
const realEnsure = factLog.ensureSynced.bind(factLog)
|
|
||||||
let failed = false
|
|
||||||
vi.spyOn(factLog, 'ensureSynced').mockImplementation(async () => {
|
|
||||||
if (!failed) {
|
|
||||||
failed = true
|
|
||||||
throw new Error('injected sync failure (device hiccup)')
|
|
||||||
}
|
|
||||||
return realEnsure()
|
|
||||||
})
|
|
||||||
|
|
||||||
// The write whose sync fails: LOUD failure to the caller — never silent.
|
|
||||||
await expect(
|
|
||||||
brain.add({ data: 'sync victim', type: NounType.Document, metadata: { n: 1 } })
|
|
||||||
).rejects.toThrow(/sync failure/)
|
|
||||||
|
|
||||||
// THE PIN: the very next write mints a fresh generation and SUCCEEDS —
|
|
||||||
// no non-monotonic refusal, no refusal loop, regardless of whether the
|
|
||||||
// failed write's fact was dropped or retained (both are legal outcomes;
|
|
||||||
// an equal-generation re-mint is not).
|
|
||||||
const survivor = await brain.add({ data: 'after the storm', type: NounType.Document, metadata: { n: 2 } })
|
|
||||||
expect((await brain.get(survivor))!.data).toContain('after the storm')
|
|
||||||
await brain.flush()
|
|
||||||
expect(Number.isSafeInteger(brain.generation())).toBe(true)
|
|
||||||
|
|
||||||
// And the log scans clean end-to-end (no torn ordering).
|
|
||||||
const scan = brain.scanFacts()
|
|
||||||
let last = 0
|
|
||||||
if (scan) {
|
|
||||||
for await (const batch of (scan as { batches(): AsyncIterable<{ facts: Array<{ generation: number }> }> }).batches()) {
|
|
||||||
for (const f of batch.facts) {
|
|
||||||
expect(f.generation, 'strictly ascending').toBeGreaterThan(last)
|
|
||||||
last = f.generation
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
expect(last).toBeGreaterThan(0)
|
|
||||||
}, 120000)
|
|
||||||
})
|
|
||||||
|
|
@ -1,149 +0,0 @@
|
||||||
/**
|
|
||||||
* @module tests/integration/write-flow-production-shape
|
|
||||||
* @description The production-shaped WRITE-FLOW gate leg. A downstream
|
|
||||||
* deployment's release gate went all-green on snapshots and rehearsal reads
|
|
||||||
* while two write-path defects (pad-frame constructibility, a counter rewind
|
|
||||||
* after a successful append) waited in ordinary WRITE flows — deferred
|
|
||||||
* embedding retries plus background history-flush concurrency wearing the
|
|
||||||
* stacks. This leg runs that exact shape, permanently:
|
|
||||||
*
|
|
||||||
* - concurrent mixed writes (adds, deferred-embed adds, updates, removes)
|
|
||||||
* - racing explicit flushes (the history tier's group commit, mid-traffic)
|
|
||||||
* - then the three laws: every ack is readable truth, the fact log is
|
|
||||||
* STRICTLY ascending end-to-end, and no write is ever refused.
|
|
||||||
*
|
|
||||||
* Part two crashes the brain mid-traffic (no close — RAM discarded) and
|
|
||||||
* requires every acked write back after reopen: the at-ack contract under
|
|
||||||
* the same production shape, not under a synthetic single write.
|
|
||||||
*/
|
|
||||||
import { describe, it, expect, afterEach } from 'vitest'
|
|
||||||
import * as fs from 'node:fs'
|
|
||||||
import { Brainy } from '../../src/brainy.js'
|
|
||||||
import { NounType } from '../../src/types/graphTypes.js'
|
|
||||||
import {
|
|
||||||
abandonAsCrashed,
|
|
||||||
factGenerations,
|
|
||||||
makeTempDir,
|
|
||||||
openBrain
|
|
||||||
} from '../helpers/durabilityKillMatrix.js'
|
|
||||||
|
|
||||||
describe('write-flow production shape — the pair gate leg from a consumer-reported miss', () => {
|
|
||||||
const dirs: string[] = []
|
|
||||||
const liveBrains: Brainy[] = []
|
|
||||||
afterEach(async () => {
|
|
||||||
for (const b of liveBrains.splice(0)) await b.close().catch(() => {})
|
|
||||||
for (const d of dirs.splice(0)) fs.rmSync(d, { recursive: true, force: true })
|
|
||||||
})
|
|
||||||
function trackDir(): string {
|
|
||||||
const dir = makeTempDir()
|
|
||||||
dirs.push(dir)
|
|
||||||
return dir
|
|
||||||
}
|
|
||||||
|
|
||||||
async function runTrafficWave(
|
|
||||||
brain: Brainy,
|
|
||||||
wave: number,
|
|
||||||
perWave: number
|
|
||||||
): Promise<{ kept: string[]; removed: string[] }> {
|
|
||||||
const kept: string[] = []
|
|
||||||
const removed: string[] = []
|
|
||||||
const work: Promise<unknown>[] = []
|
|
||||||
for (let i = 0; i < perWave; i++) {
|
|
||||||
const n = wave * perWave + i
|
|
||||||
if (i % 4 === 0) {
|
|
||||||
// Deferred-embed add — the retry-marker flow that wore the defect.
|
|
||||||
work.push(
|
|
||||||
brain
|
|
||||||
.add({ data: `deferred payload ${n}`, type: NounType.Document, metadata: { n, defer: true }, deferEmbedding: true })
|
|
||||||
.then((id) => void kept.push(id))
|
|
||||||
)
|
|
||||||
} else if (i % 4 === 1) {
|
|
||||||
// Add, then update it in the same wave (two generations, same id).
|
|
||||||
work.push(
|
|
||||||
brain.add({ data: `versioned payload ${n}`, type: NounType.Document, metadata: { n, v: 1 } }).then(async (id) => {
|
|
||||||
kept.push(id)
|
|
||||||
await brain.update({ id, metadata: { n, v: 2 } })
|
|
||||||
})
|
|
||||||
)
|
|
||||||
} else if (i % 4 === 2) {
|
|
||||||
// Add, then remove — a durable tombstone is an ack too.
|
|
||||||
work.push(
|
|
||||||
brain.add({ data: `ephemeral payload ${n}`, type: NounType.Document, metadata: { n } }).then(async (id) => {
|
|
||||||
await brain.remove(id)
|
|
||||||
removed.push(id)
|
|
||||||
})
|
|
||||||
)
|
|
||||||
} else {
|
|
||||||
work.push(
|
|
||||||
brain.add({ data: `plain payload ${n}`, type: NounType.Document, metadata: { n } }).then((id) => void kept.push(id))
|
|
||||||
)
|
|
||||||
}
|
|
||||||
// Race the history tier's group commit against live traffic.
|
|
||||||
if (i % 5 === 3) work.push(brain.flush())
|
|
||||||
}
|
|
||||||
// NO REFUSALS: every promise must resolve — a single rejection here is
|
|
||||||
// the refusal-loop costume this leg exists to catch.
|
|
||||||
await Promise.all(work)
|
|
||||||
return { kept, removed }
|
|
||||||
}
|
|
||||||
|
|
||||||
it('three waves of mixed traffic with racing flushes: every ack is truth, the log is strictly ascending, nothing refused', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
expect(brain.logAuthority().authority).toBe('log')
|
|
||||||
|
|
||||||
const kept: string[] = []
|
|
||||||
const removed: string[] = []
|
|
||||||
for (let wave = 0; wave < 3; wave++) {
|
|
||||||
const result = await runTrafficWave(brain, wave, 20)
|
|
||||||
kept.push(...result.kept)
|
|
||||||
removed.push(...result.removed)
|
|
||||||
}
|
|
||||||
await brain.flush()
|
|
||||||
|
|
||||||
for (const id of kept) {
|
|
||||||
expect(await brain.get(id), `acked write ${id} must be readable truth`).not.toBeNull()
|
|
||||||
}
|
|
||||||
for (const id of removed) {
|
|
||||||
expect(await brain.get(id), `acked remove ${id} must hold`).toBeNull()
|
|
||||||
}
|
|
||||||
|
|
||||||
const gens = await factGenerations(brain)
|
|
||||||
expect(gens.length).toBeGreaterThan(0)
|
|
||||||
for (let i = 1; i < gens.length; i++) {
|
|
||||||
expect(gens[i], 'fact log strictly ascending end-to-end').toBeGreaterThan(gens[i - 1])
|
|
||||||
}
|
|
||||||
|
|
||||||
// Clean reopen: the same truth survives a restart.
|
|
||||||
await liveBrains.pop()!.close()
|
|
||||||
const reopened = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(reopened)
|
|
||||||
for (const id of kept.slice(0, 10)) {
|
|
||||||
expect(await reopened.get(id)).not.toBeNull()
|
|
||||||
}
|
|
||||||
}, 240000)
|
|
||||||
|
|
||||||
it('crash mid-traffic: every acked write survives the reopen (the at-ack law under the production shape)', async () => {
|
|
||||||
const dir = trackDir()
|
|
||||||
const brain = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(brain)
|
|
||||||
|
|
||||||
const { kept, removed } = await runTrafficWave(brain, 0, 24)
|
|
||||||
// No close, no flush — the process "dies" holding its RAM.
|
|
||||||
await abandonAsCrashed(liveBrains.pop()!)
|
|
||||||
|
|
||||||
const reopened = await openBrain(dir, { logAuthority: 'adopt' })
|
|
||||||
liveBrains.push(reopened)
|
|
||||||
for (const id of kept) {
|
|
||||||
expect(await reopened.get(id), `acked write ${id} must survive the crash`).not.toBeNull()
|
|
||||||
}
|
|
||||||
for (const id of removed) {
|
|
||||||
expect(await reopened.get(id), `acked remove ${id} must survive the crash`).toBeNull()
|
|
||||||
}
|
|
||||||
const gens = await factGenerations(reopened)
|
|
||||||
for (let i = 1; i < gens.length; i++) {
|
|
||||||
expect(gens[i], 'fact log strictly ascending after recovery').toBeGreaterThan(gens[i - 1])
|
|
||||||
}
|
|
||||||
}, 240000)
|
|
||||||
})
|
|
||||||
|
|
@ -1,29 +0,0 @@
|
||||||
/**
|
|
||||||
* @module tests/unit/db/pad-frame-total
|
|
||||||
* @description Pad-frame construction is TOTAL: every size from the minimum
|
|
||||||
* through 4096+257 is constructible byte-exact (a production adoption found
|
|
||||||
* the msgpack class-boundary hole at 291 bytes — sync died whole, and the
|
|
||||||
* failure cascaded into a counter rewind after a successful append). Every
|
|
||||||
* constructed pad decodes as skip-by-definition filler.
|
|
||||||
*/
|
|
||||||
import { describe, it, expect } from 'vitest'
|
|
||||||
import { encodePadFrame, minPadFrameBytes, decodeGroupV2 } from '../../../src/db/factLogFormat.js'
|
|
||||||
|
|
||||||
describe('pad frames are constructible at EVERY size', () => {
|
|
||||||
it('exact construction from the minimum through a full sector + boundary spill', () => {
|
|
||||||
const min = minPadFrameBytes()
|
|
||||||
for (let size = min; size <= 4096 + 257; size++) {
|
|
||||||
const frame = encodePadFrame(size)
|
|
||||||
expect(frame.length, `size ${size}`).toBe(size)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
|
|
||||||
it('the production case (291) and its class-boundary siblings decode as invisible filler', () => {
|
|
||||||
for (const size of [291, minPadFrameBytes(), 300, 511, 512, 513, 4096]) {
|
|
||||||
const frame = encodePadFrame(size)
|
|
||||||
const group = decodeGroupV2(frame)
|
|
||||||
expect(group.facts, `size ${size} is reader-invisible`).toEqual([])
|
|
||||||
expect(group.validBytes).toBe(size)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
})
|
|
||||||
|
|
@ -37,9 +37,6 @@ const MANUAL_ONLY = new Set<string>([
|
||||||
// (byte + fold digests) — runs in the explicit conformance gate stage,
|
// (byte + fold digests) — runs in the explicit conformance gate stage,
|
||||||
// same invocation family as the other conformance suites.
|
// same invocation family as the other conformance suites.
|
||||||
'tests/conformance/golden-log-fold.test.ts',
|
'tests/conformance/golden-log-fold.test.ts',
|
||||||
// The sparse-store cut's shared operator rows (both engines run these):
|
|
||||||
// explicit conformance-gate invocation, like its siblings.
|
|
||||||
'tests/conformance/sparse-store-cut.test.ts',
|
|
||||||
'tests/api/performance-benchmarks.test.ts',
|
'tests/api/performance-benchmarks.test.ts',
|
||||||
'tests/critical-neural-validation.test.ts',
|
'tests/critical-neural-validation.test.ts',
|
||||||
'tests/critical-performance-benchmark.test.ts',
|
'tests/critical-performance-benchmark.test.ts',
|
||||||
|
|
|
||||||
Reference in a new issue