Compare commits
21 commits
fix/pendin
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| bc82d36294 | |||
| 97b5ea2d5d | |||
| aa457d7159 | |||
| adcb883e67 | |||
| 85b1fa5c1a | |||
| 8752f11f4d | |||
| 61bc5f423b | |||
| 3835a0e702 | |||
| 08758c254f | |||
| 3dadbec8f2 | |||
| 4f1e27c9a0 | |||
| 297a3d7657 | |||
| 7ab670b525 | |||
| f097cbf6f2 | |||
| e64e2bc175 | |||
| 655aa13ea7 | |||
| 39c71ecdac | |||
| 0759c03a82 | |||
| 9a888c37e9 | |||
|
|
298cb6daca | ||
| b8475cc86a |
33 changed files with 1187 additions and 4159 deletions
|
|
@ -5,6 +5,10 @@ name: CI
|
|||
# sequential, so tag-triggered matrix jobs (~22 min) would queue AHEAD of the
|
||||
# tag's publish-source run and starve every release (observed on 8.10.3 and
|
||||
# 9.0.0: the publish sat behind the tag's own redundant CI).
|
||||
concurrency:
|
||||
group: ci-${{ github.ref }}
|
||||
cancel-in-progress: true
|
||||
|
||||
on:
|
||||
push:
|
||||
branches: ['**']
|
||||
|
|
|
|||
|
|
@ -1,148 +0,0 @@
|
|||
name: Delta Gate
|
||||
|
||||
# On-demand candidate-vs-control gate on the capped functional CI lane
|
||||
# (label: gate-functional). That lane is Bun-only host-mode — there is no
|
||||
# Node.js runtime available to it, so this workflow deliberately avoids every
|
||||
# JS-based action (checkout/setup-node/setup-bun/upload-artifact all require
|
||||
# one) and does everything with plain git + bun in shell steps instead.
|
||||
#
|
||||
# Verdict lines a caller should grep for in the run log:
|
||||
# COLLECTED patch=<n> control=<n> — collection-truncation guard inputs
|
||||
# NEW-RED-COUNT:<n> — failures on candidate absent from control
|
||||
# DELTA-GATE: CLEAN | NEW REDS | INVALID | STOPPED-BY-REGISTRY-TRIPWIRE
|
||||
#
|
||||
# The lane's own housekeeping stops the runner and drops a marker file when
|
||||
# host pressure (I/O, registry latency, disk budget) trips — never ours to
|
||||
# interpret as a red or a green. The final step checks for that marker before
|
||||
# it says anything about pass/fail.
|
||||
|
||||
on:
|
||||
workflow_dispatch:
|
||||
inputs:
|
||||
candidate:
|
||||
description: 'Candidate ref (branch or sha) to gate'
|
||||
required: true
|
||||
type: string
|
||||
control:
|
||||
description: 'Control sha to diff against'
|
||||
required: true
|
||||
type: string
|
||||
# workflow_dispatch needs Actions-unit write on the dispatching credential;
|
||||
# push does not (it runs from the pushed ref's own tree), so a plain push
|
||||
# to a release or CI branch is the fallback trigger while that grant is
|
||||
# outstanding — see the ref-resolution step below for what it gates against.
|
||||
push:
|
||||
branches: ['rel/**', 'ci/**']
|
||||
|
||||
concurrency:
|
||||
group: delta-gate
|
||||
cancel-in-progress: false
|
||||
|
||||
jobs:
|
||||
delta-gate:
|
||||
name: Delta gate — candidate vs control
|
||||
runs-on: gate-functional
|
||||
timeout-minutes: 120
|
||||
steps:
|
||||
- name: Resolve candidate/control refs
|
||||
id: refs
|
||||
run: |
|
||||
candidate="${{ github.event.inputs.candidate }}"
|
||||
control="${{ github.event.inputs.control }}"
|
||||
# workflow_dispatch supplies both explicitly; a push event carries
|
||||
# neither — fall back to the pushed commit as candidate and the
|
||||
# last released, known-good tip (10.4.9) as control, so a plain
|
||||
# push still produces a meaningful gate instead of an empty ref.
|
||||
if [ -z "$candidate" ]; then candidate="${{ github.sha }}"; fi
|
||||
if [ -z "$control" ]; then control="eec90bdd"; fi
|
||||
echo "candidate=$candidate" >> "$GITHUB_OUTPUT"
|
||||
echo "control=$control" >> "$GITHUB_OUTPUT"
|
||||
echo "Resolved (trigger=${{ github.event_name }}): candidate=$candidate control=$control"
|
||||
|
||||
- name: Clean any residue from a prior run
|
||||
run: rm -rf "ob-cand-${{ github.run_id }}" "ob-ctrl-${{ github.run_id }}" "/tmp/ob-${{ github.run_id }}-"*
|
||||
|
||||
- name: Clone + test — candidate
|
||||
id: patch
|
||||
run: |
|
||||
set -o pipefail
|
||||
git clone --quiet "https://source.soulcraft.com/soulcraftlabs/open-brainy.git" "ob-cand-${{ github.run_id }}"
|
||||
cd "ob-cand-${{ github.run_id }}"
|
||||
git checkout --quiet "${{ steps.refs.outputs.candidate }}"
|
||||
git log --oneline -1
|
||||
bun install
|
||||
rc=0
|
||||
bun x vitest run > "/tmp/ob-${{ github.run_id }}-patch.log" 2>&1 || rc=$?
|
||||
echo "PATCH-RC:$rc"
|
||||
grep -aE "Tests .*(passed|failed)" "/tmp/ob-${{ github.run_id }}-patch.log" | tail -1
|
||||
grep -aE "^ FAIL |^\s+×" "/tmp/ob-${{ github.run_id }}-patch.log" | sed -E "s/ [0-9]+ms$//" | sed -E "s/^\s+//" | sort -u > "/tmp/ob-${{ github.run_id }}-patch.fail"
|
||||
echo "PATCH-FAILING:$(wc -l < "/tmp/ob-${{ github.run_id }}-patch.fail")"
|
||||
|
||||
- name: Clone + test — control
|
||||
id: control
|
||||
run: |
|
||||
set -o pipefail
|
||||
git clone --quiet "https://source.soulcraft.com/soulcraftlabs/open-brainy.git" "ob-ctrl-${{ github.run_id }}"
|
||||
cd "ob-ctrl-${{ github.run_id }}"
|
||||
git checkout --quiet "${{ steps.refs.outputs.control }}"
|
||||
git log --oneline -1
|
||||
bun install
|
||||
rc=0
|
||||
bun x vitest run > "/tmp/ob-${{ github.run_id }}-control.log" 2>&1 || rc=$?
|
||||
echo "CONTROL-RC:$rc"
|
||||
grep -aE "Tests .*(passed|failed)" "/tmp/ob-${{ github.run_id }}-control.log" | tail -1
|
||||
grep -aE "^ FAIL |^\s+×" "/tmp/ob-${{ github.run_id }}-control.log" | sed -E "s/ [0-9]+ms$//" | sed -E "s/^\s+//" | sort -u > "/tmp/ob-${{ github.run_id }}-control.fail"
|
||||
echo "CONTROL-FAILING:$(wc -l < "/tmp/ob-${{ github.run_id }}-control.fail")"
|
||||
|
||||
- name: Delta gate verdict
|
||||
if: always()
|
||||
run: |
|
||||
set -o pipefail
|
||||
|
||||
# The lane's own tripwire wins over anything we would otherwise say:
|
||||
# a bare failure/timeout above with this marker present is host
|
||||
# pressure, never a real red and never a real green.
|
||||
if [ -f /srv/gate-lane/TRIPWIRE-STOPPED ]; then
|
||||
echo "DELTA-GATE: STOPPED-BY-REGISTRY-TRIPWIRE"
|
||||
head -1 /srv/gate-lane/TRIPWIRE-STOPPED
|
||||
exit 3
|
||||
fi
|
||||
|
||||
patch_log="/tmp/ob-${{ github.run_id }}-patch.log"
|
||||
control_log="/tmp/ob-${{ github.run_id }}-control.log"
|
||||
patch_fail="/tmp/ob-${{ github.run_id }}-patch.fail"
|
||||
control_fail="/tmp/ob-${{ github.run_id }}-control.fail"
|
||||
|
||||
if [ ! -s "$patch_log" ] || [ ! -s "$control_log" ]; then
|
||||
echo "DELTA-GATE: INVALID — a leg produced no log (see the two steps above for the real cause)"
|
||||
exit 2
|
||||
fi
|
||||
|
||||
pt=$(grep -aoE "\(([0-9]+)\)$" "$patch_log" | tail -1 | tr -d "()")
|
||||
ct=$(grep -aoE "\(([0-9]+)\)$" "$control_log" | tail -1 | tr -d "()")
|
||||
echo "COLLECTED patch=${pt:-0} control=${ct:-0}"
|
||||
if [ "${pt:-0}" -lt 3000 ] || [ "${ct:-0}" -lt 3000 ]; then
|
||||
echo "DELTA-GATE: INVALID — truncated collection"
|
||||
exit 2
|
||||
fi
|
||||
|
||||
echo "=== NEW REDS ==="
|
||||
comm -23 "$patch_fail" "$control_fail"
|
||||
new=$(comm -23 "$patch_fail" "$control_fail" | wc -l)
|
||||
echo "NEW-RED-COUNT:$new"
|
||||
|
||||
echo "=== full candidate fail list ==="
|
||||
cat "$patch_fail"
|
||||
echo "=== full control fail list ==="
|
||||
cat "$control_fail"
|
||||
|
||||
if [ "$new" -eq 0 ]; then
|
||||
echo "DELTA-GATE: CLEAN"
|
||||
else
|
||||
echo "DELTA-GATE: NEW REDS"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
- name: Clean up (mind the lane's disk budget)
|
||||
if: always()
|
||||
run: rm -rf "ob-cand-${{ github.run_id }}" "ob-ctrl-${{ github.run_id }}" "/tmp/ob-${{ github.run_id }}-"*
|
||||
|
|
@ -12,6 +12,11 @@ on:
|
|||
push:
|
||||
tags:
|
||||
- 'v*'
|
||||
workflow_dispatch:
|
||||
inputs:
|
||||
ref_reason:
|
||||
description: 'why this manual run (e.g. tag event dropped)'
|
||||
required: false
|
||||
|
||||
jobs:
|
||||
publish:
|
||||
|
|
|
|||
23
CHANGELOG.md
23
CHANGELOG.md
|
|
@ -2,29 +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.
|
||||
|
||||
### [10.4.9](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.6...v10.4.9) (2026-09-02)
|
||||
|
||||
- Merge branch 'fix/pending-embed-low-water' into rel/10.4.9-candidate (2648f56d)
|
||||
- fix(open): pending-embed recovery keeps the crash-recovery contract — foreground, bounded by the mark (8a2ebacf)
|
||||
- Merge branches 'fix/connected-find-order', 'fix/pending-embed-low-water' and 'fix/related-verb-array' into rel/10.4.9-candidate (d5147ed6)
|
||||
- fix(graph): the verb fast paths honour every requested type, source, and target (6a89adc4)
|
||||
- perf(open): pending-embed recovery is bounded by a low-water mark and runs behind the doors (88e79729)
|
||||
- fix(find): connected finds are graph-first — neighbours, then the filter over those ids, then the page (077cbc0b)
|
||||
- fix(storage): counts persistence is single-flight, coalesced, and never races its own temp file (5e3b343a)
|
||||
|
||||
|
||||
### [10.4.6](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.5...v10.4.6) (2026-08-31)
|
||||
|
||||
- fix(transact): metadata-index ops take their JSON-safe view at the crossing, not at construction (73500e7d)
|
||||
|
||||
|
||||
### [10.4.5](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.4...v10.4.5) (2026-08-31)
|
||||
|
||||
- build(release): the docs-push step retires — this engine documents itself in its own repository (d6bcb14f)
|
||||
- fix(generations): a sealed segment may only declare the generations it holds (a963a744)
|
||||
- fix(recovery): a torn generation-log tail is a terminal verdict, never a wait (c9930871)
|
||||
|
||||
|
||||
### [10.4.4](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.3...v10.4.4) (2026-08-28)
|
||||
|
||||
- fix(vfs): the old-root sweep narrates only when it has something to say (d49148e1)
|
||||
|
|
|
|||
|
|
@ -1,5 +1,12 @@
|
|||
# @soulcraft/brainy — Release Notes for Consumers
|
||||
|
||||
Machine-readable release notes are published at
|
||||
https://source.soulcraft.com/soulcraftlabs/releases/raw/branch/main/open-brainy.json
|
||||
(this engine) and
|
||||
https://source.soulcraft.com/soulcraftlabs/releases/raw/branch/main/brainy.json
|
||||
(the product engine) — read by HQ's `/hq/releases` door, and the source of
|
||||
truth ahead of this file.
|
||||
|
||||
This file is the **quick reference for downstream sessions** tracking Brainy changes.
|
||||
Full auto-generated changelog: `CHANGELOG.md` · Releases: https://source.soulcraft.com/soulcraftlabs/open-brainy/releases
|
||||
|
||||
|
|
|
|||
4
package-lock.json
generated
4
package-lock.json
generated
|
|
@ -1,12 +1,12 @@
|
|||
{
|
||||
"name": "@soulcraftlabs/brainy",
|
||||
"version": "10.4.9",
|
||||
"version": "10.4.4",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "@soulcraftlabs/brainy",
|
||||
"version": "10.4.9",
|
||||
"version": "10.4.4",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"@msgpack/msgpack": "^3.1.2",
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
{
|
||||
"name": "@soulcraftlabs/brainy",
|
||||
"version": "10.4.9",
|
||||
"version": "10.4.4",
|
||||
"brainyContract": 1,
|
||||
"description": "Universal Knowledge Protocol™ - World's first Triple Intelligence database unifying vector, graph, and document search in one API. Stage 3 CANONICAL: 42 nouns × 127 verbs covering 96-97% of all human knowledge.",
|
||||
"main": "dist/index.js",
|
||||
|
|
|
|||
|
|
@ -154,7 +154,8 @@ else
|
|||
fi
|
||||
|
||||
# Create new changelog entry
|
||||
CHANGELOG_ENTRY="### [${NEW_VERSION}](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v${CURRENT_VERSION}...v${NEW_VERSION}) ($(date +%Y-%m-%d))
|
||||
RELEASE_DATE=$(date +%Y-%m-%d)
|
||||
CHANGELOG_ENTRY="### [${NEW_VERSION}](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v${CURRENT_VERSION}...v${NEW_VERSION}) (${RELEASE_DATE})
|
||||
|
||||
${COMMITS}
|
||||
"
|
||||
|
|
@ -174,6 +175,19 @@ if [ -f "CHANGELOG.md" ]; then
|
|||
fi
|
||||
echo -e "${GREEN}✅ CHANGELOG updated${NC}\n"
|
||||
|
||||
# Step 6b: Update the releases wall entry — mechanical, derived from the
|
||||
# CHANGELOG entry just composed. The fleet's HQ page reads open-brainy.json
|
||||
# from the one shared releases repo, soulcraftlabs/releases on The Source —
|
||||
# this used to be hand-written after every release (David: never again —
|
||||
# make it a step of the rail, landed in the one shared home; this repo no
|
||||
# longer hosts its own copy). This step clones/fetches that repo into a
|
||||
# local cache, prepends the entry, and pushes it directly — a real
|
||||
# cross-repo push, refusing loudly (never skipping) on any
|
||||
# clone/validation/commit/push failure.
|
||||
echo -e "${BLUE}5️⃣▸ Updating the releases wall...${NC}"
|
||||
node scripts/wall-entry.mjs --product open-brainy --version "${NEW_VERSION}" --date "${RELEASE_DATE}" --from-changelog CHANGELOG.md
|
||||
echo -e "${GREEN}✅ Releases wall updated${NC}\n"
|
||||
|
||||
# Step 7: Create release commit
|
||||
echo -e "${BLUE}6️⃣ Creating release commit...${NC}"
|
||||
git add package.json package-lock.json CHANGELOG.md
|
||||
|
|
@ -237,7 +251,7 @@ fi
|
|||
# and RELEASES.md are the record; this just gives The Source's UI a release page).
|
||||
echo -e "${BLUE}🔟 Creating release page on The Source...${NC}"
|
||||
if [ -n "${FORGEJO_RELEASE_TOKEN:-}" ]; then
|
||||
if curl -sf -X POST "https://source.soulcraft.com/api/v1/repos/soulcraft/brainy/releases" \
|
||||
if curl -sf -X POST "https://source.soulcraft.com/api/v1/repos/soulcraftlabs/open-brainy/releases" \
|
||||
-H "Authorization: token ${FORGEJO_RELEASE_TOKEN}" -H "Content-Type: application/json" \
|
||||
-d "{\"tag_name\":\"v${NEW_VERSION}\",\"name\":\"v${NEW_VERSION}\",\"prerelease\":${PRERELEASE}}" >/dev/null; then
|
||||
echo -e "${GREEN}✅ Release page created on The Source${NC}\n"
|
||||
|
|
|
|||
504
scripts/wall-entry.mjs
Normal file
504
scripts/wall-entry.mjs
Normal file
|
|
@ -0,0 +1,504 @@
|
|||
#!/usr/bin/env node
|
||||
/**
|
||||
* @module scripts/wall-entry
|
||||
* @description The releases-wall entry, made mechanical. The fleet's HQ page
|
||||
* reads one public JSON per product from the ONE releases repo on The Source
|
||||
* (soulcraftlabs/releases, files <product>.json at its root — shape
|
||||
* {product, entries:[{version, date, headline, items, url, thumb?}]}), at
|
||||
* https://source.soulcraft.com/soulcraftlabs/releases/raw/branch/main/<product>.json.
|
||||
* Those entries were hand-written after every release, then briefly written
|
||||
* into this repo's own releases/<product>.json; this script is the one door
|
||||
* that composes an entry and lands it in the shared repo, so it is never
|
||||
* hand-written and never forked across repos again.
|
||||
*
|
||||
* Two modes:
|
||||
*
|
||||
* 1. Generate + publish (default):
|
||||
* node wall-entry.mjs --product <p> --version <v> --date <YYYY-MM-DD> \
|
||||
* --from-changelog <CHANGELOG.md>
|
||||
* Derives an entry from the CHANGELOG.md entry for <v> (headline = the
|
||||
* entry's first bullet, items = every bullet, trimmed of its trailing
|
||||
* commit hash), then:
|
||||
* - clones (or, if a cached clone already exists, fetches and resets)
|
||||
* the releases repo into a local cache directory,
|
||||
* - prepends the entry to <cache>/<p>.json, newest first — replacing
|
||||
* any existing entry for the same version so a re-run is idempotent,
|
||||
* - validates the file's shape before and after,
|
||||
* - commits the change as "chore(wall): <p> <v>" and pushes main.
|
||||
* A failure at any step (clone, validation, commit, push, a
|
||||
* non-fast-forward remote) exits non-zero naming the cure. Nothing is
|
||||
* ever skipped — the wall either lands correctly or the release fails.
|
||||
*
|
||||
* 2. Dry run:
|
||||
* node wall-entry.mjs --dry-run --product <p> --version <v> \
|
||||
* --date <YYYY-MM-DD> --from-changelog <CHANGELOG.md>
|
||||
* Derives the entry exactly as above and prints it, along with the file
|
||||
* it would be written to, but touches no clone and no remote — usable
|
||||
* from a fresh checkout with no cache and no network.
|
||||
*
|
||||
* 3. Validate only (--check):
|
||||
* node wall-entry.mjs --check --file <path/to/product.json>
|
||||
* Validates an arbitrary wall file's exact key set (top-level and
|
||||
* per-entry), field types, and strict-descending semver ordering with
|
||||
* no duplicates. Read-only; never writes. Exit 0 = clean, exit 1 =
|
||||
* named violations printed to stderr.
|
||||
*
|
||||
* The remote and the local cache directory are each overridable
|
||||
* (--remote / --cache-dir, or WALL_ENTRY_RELEASES_REMOTE /
|
||||
* WALL_ENTRY_RELEASES_CACHE_DIR) so tests can point at a throwaway local
|
||||
* bare repo and a throwaway cache directory — never the real remote or the
|
||||
* real developer cache.
|
||||
*
|
||||
* No dependencies beyond the system `git` binary — CHANGELOG parsing,
|
||||
* semver comparison, and JSON shape checking are all hand-rolled below.
|
||||
*/
|
||||
|
||||
import { readFileSync, writeFileSync, existsSync, mkdirSync } from 'node:fs'
|
||||
import { execFileSync } from 'node:child_process'
|
||||
import { homedir } from 'node:os'
|
||||
import { dirname, join } from 'node:path'
|
||||
|
||||
const DEFAULT_REMOTE = 'git@source.soulcraft.com:soulcraftlabs/releases.git'
|
||||
|
||||
/** @returns {string} */
|
||||
function defaultCacheDir() {
|
||||
const base = process.env.XDG_CACHE_HOME || join(homedir(), '.cache')
|
||||
return join(base, 'soulcraft-releases')
|
||||
}
|
||||
|
||||
// Required on every entry; "thumb" is optional (may be absent, or present as
|
||||
// string | null) — matching the HQ contract's {..., thumb?}.
|
||||
const ENTRY_REQUIRED_KEYS = ['version', 'date', 'headline', 'items', 'url']
|
||||
const ENTRY_OPTIONAL_KEYS = ['thumb']
|
||||
const ENTRY_ALLOWED_KEYS = [...ENTRY_REQUIRED_KEYS, ...ENTRY_OPTIONAL_KEYS]
|
||||
const FILE_KEYS = ['product', 'entries']
|
||||
|
||||
// The public permalink pattern, by product. Every entry MUST carry an https
|
||||
// permalink: HQ's parser rejects a wall whose entries carry url: null (the
|
||||
// whole feed became unreadable on 2026-09-02). A product whose forge repo is
|
||||
// private links its PUBLIC package page on The Source instead of a release
|
||||
// page that would 404 for HQ's readers.
|
||||
const RELEASE_URL_PATTERNS = {
|
||||
'open-brainy': (version) => `https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v${version}`,
|
||||
'brainy': (version) => `https://source.soulcraft.com/soulcraft/-/packages/npm/@soulcraft%2Fbrainy/${version}`,
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse argv into a flag map. `--flag value` sets a string; `--flag` alone
|
||||
* (end of argv, or followed by another `--flag`) sets boolean true.
|
||||
* @param {string[]} argv
|
||||
* @returns {Record<string, string | true>}
|
||||
*/
|
||||
function parseArgs(argv) {
|
||||
/** @type {Record<string, string | true>} */
|
||||
const args = {}
|
||||
for (let i = 0; i < argv.length; i++) {
|
||||
const a = argv[i]
|
||||
if (!a.startsWith('--')) continue
|
||||
const key = a.slice(2)
|
||||
const next = argv[i + 1]
|
||||
if (next === undefined || next.startsWith('--')) {
|
||||
args[key] = true
|
||||
} else {
|
||||
args[key] = next
|
||||
i++
|
||||
}
|
||||
}
|
||||
return args
|
||||
}
|
||||
|
||||
/**
|
||||
* Print a loud, named error and exit 1. Every refusal in this script goes
|
||||
* through here so the failure mode is always the same shape: "wall-entry: <what>".
|
||||
* @param {string} message
|
||||
* @returns {never}
|
||||
*/
|
||||
function fail(message) {
|
||||
console.error(`wall-entry: ${message}`)
|
||||
process.exit(1)
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} version
|
||||
* @returns {{major: number, minor: number, patch: number, pre: string | null} | null}
|
||||
*/
|
||||
function parseSemver(version) {
|
||||
const m = /^(\d+)\.(\d+)\.(\d+)(?:-([0-9A-Za-z.-]+))?$/.exec(version)
|
||||
if (!m) return null
|
||||
return { major: Number(m[1]), minor: Number(m[2]), patch: Number(m[3]), pre: m[4] ?? null }
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} a
|
||||
* @param {string} b
|
||||
* @returns {number} positive if a > b, negative if a < b, 0 if equal.
|
||||
*/
|
||||
function compareSemver(a, b) {
|
||||
const pa = parseSemver(a)
|
||||
const pb = parseSemver(b)
|
||||
if (!pa || !pb) throw new Error(`cannot compare non-semver versions "${a}" vs "${b}"`)
|
||||
if (pa.major !== pb.major) return pa.major - pb.major
|
||||
if (pa.minor !== pb.minor) return pa.minor - pb.minor
|
||||
if (pa.patch !== pb.patch) return pa.patch - pb.patch
|
||||
if (pa.pre === pb.pre) return 0
|
||||
if (pa.pre === null) return 1 // a release outranks any prerelease of the same core version
|
||||
if (pb.pre === null) return -1
|
||||
return pa.pre < pb.pre ? -1 : pa.pre > pb.pre ? 1 : 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Validate a wall file's full shape: top-level keys ("product", "entries" —
|
||||
* no more, no less), per-entry keys and field types ("thumb" optional), and
|
||||
* strict-descending semver ordering with no duplicates. Collects every
|
||||
* violation instead of failing on the first, so a caller reports the whole
|
||||
* picture in one pass.
|
||||
* @param {unknown} data
|
||||
* @returns {string[]} Violation messages; empty means the file is clean.
|
||||
*/
|
||||
function validateShape(data) {
|
||||
/** @type {string[]} */
|
||||
const errors = []
|
||||
|
||||
if (typeof data !== 'object' || data === null || Array.isArray(data)) {
|
||||
return ['top level: expected a JSON object']
|
||||
}
|
||||
const obj = /** @type {Record<string, unknown>} */ (data)
|
||||
|
||||
const topKeys = Object.keys(obj)
|
||||
const missingTop = FILE_KEYS.filter((k) => !(k in obj))
|
||||
const extraTop = topKeys.filter((k) => !FILE_KEYS.includes(k))
|
||||
if (missingTop.length) errors.push(`top level: missing key(s) ${missingTop.join(', ')}`)
|
||||
if (extraTop.length) errors.push(`top level: unexpected key(s) ${extraTop.join(', ')}`)
|
||||
|
||||
if (typeof obj.product !== 'string' || obj.product.trim() === '') {
|
||||
errors.push('top level: "product" must be a non-empty string')
|
||||
}
|
||||
if (!Array.isArray(obj.entries)) {
|
||||
errors.push('top level: "entries" must be an array')
|
||||
return errors // nothing further to check without an array
|
||||
}
|
||||
|
||||
const entries = /** @type {unknown[]} */ (obj.entries)
|
||||
entries.forEach((rawEntry, i) => {
|
||||
const label = `entries[${i}]`
|
||||
if (typeof rawEntry !== 'object' || rawEntry === null || Array.isArray(rawEntry)) {
|
||||
errors.push(`${label}: expected an object`)
|
||||
return
|
||||
}
|
||||
const entry = /** @type {Record<string, unknown>} */ (rawEntry)
|
||||
const keys = Object.keys(entry)
|
||||
const missing = ENTRY_REQUIRED_KEYS.filter((k) => !(k in entry))
|
||||
const extra = keys.filter((k) => !ENTRY_ALLOWED_KEYS.includes(k))
|
||||
if (missing.length) errors.push(`${label}: missing key(s) ${missing.join(', ')}`)
|
||||
if (extra.length) errors.push(`${label}: unexpected key(s) ${extra.join(', ')}`)
|
||||
|
||||
if (typeof entry.version !== 'string' || !parseSemver(entry.version)) {
|
||||
errors.push(`${label}: "version" must be a semver string (got ${JSON.stringify(entry.version)})`)
|
||||
}
|
||||
if (typeof entry.date !== 'string' || !/^\d{4}-\d{2}-\d{2}$/.test(entry.date) || Number.isNaN(Date.parse(entry.date))) {
|
||||
errors.push(`${label}: "date" must be a YYYY-MM-DD string (got ${JSON.stringify(entry.date)})`)
|
||||
}
|
||||
if (typeof entry.headline !== 'string' || entry.headline.trim() === '') {
|
||||
errors.push(`${label}: "headline" must be a non-empty string`)
|
||||
}
|
||||
if (!Array.isArray(entry.items) || entry.items.length === 0 || entry.items.some((it) => typeof it !== 'string' || it.trim() === '')) {
|
||||
errors.push(`${label}: "items" must be a non-empty array of non-empty strings`)
|
||||
}
|
||||
if (typeof entry.url !== 'string' || !/^https:\/\/\S+$/.test(entry.url)) {
|
||||
errors.push(`${label}: "url" must be an https permalink — never null; HQ's parser rejects the whole feed`)
|
||||
}
|
||||
if ('thumb' in entry && !(entry.thumb === null || typeof entry.thumb === 'string')) {
|
||||
errors.push(`${label}: "thumb" must be a string or null when present`)
|
||||
}
|
||||
})
|
||||
|
||||
// Ordering: newest first, strictly descending, no duplicate versions —
|
||||
// checked only over entries whose version parsed (a bad version is
|
||||
// already reported above; comparing it too would just be noise).
|
||||
const versioned = entries
|
||||
.map((e, i) => ({ i, version: /** @type {any} */ (e)?.version }))
|
||||
.filter((e) => typeof e.version === 'string' && parseSemver(e.version))
|
||||
for (let i = 0; i < versioned.length - 1; i++) {
|
||||
const a = versioned[i]
|
||||
const b = versioned[i + 1]
|
||||
const cmp = compareSemver(a.version, b.version)
|
||||
if (cmp === 0) {
|
||||
errors.push(`entries[${a.i}] and entries[${b.i}]: duplicate version ${a.version}`)
|
||||
} else if (cmp < 0) {
|
||||
errors.push(`entries[${a.i}] (${a.version}) sits above entries[${b.i}] (${b.version}) — not newest-first`)
|
||||
}
|
||||
}
|
||||
|
||||
return errors
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract one version's entry body from a standard-version-style CHANGELOG.md
|
||||
* (headings `### [version](url) (date)`, followed by `- bullet (hash)` lines
|
||||
* until the next heading or EOF).
|
||||
* @param {string} changelog
|
||||
* @param {string} version
|
||||
* @returns {string[]} Bullet lines, trimmed of their leading "- " and
|
||||
* trailing " (hash)".
|
||||
*/
|
||||
function extractChangelogBullets(changelog, version) {
|
||||
const lines = changelog.split('\n')
|
||||
const headingRe = /^### \[([^\]]+)\]\(.*\)\s*\(\d{4}-\d{2}-\d{2}\)\s*$/
|
||||
let start = -1
|
||||
for (let i = 0; i < lines.length; i++) {
|
||||
const m = headingRe.exec(lines[i])
|
||||
if (m && m[1] === version) {
|
||||
start = i + 1
|
||||
break
|
||||
}
|
||||
}
|
||||
if (start === -1) {
|
||||
fail(
|
||||
`version ${version} has no CHANGELOG entry yet — run this after the CHANGELOG step composes "### [${version}]", not before`,
|
||||
)
|
||||
}
|
||||
/** @type {string[]} */
|
||||
const bullets = []
|
||||
for (let i = start; i < lines.length; i++) {
|
||||
if (headingRe.test(lines[i])) break // next entry starts
|
||||
const bulletMatch = /^- (.+?)(?:\s\(([0-9a-f]{6,40})\))?$/.exec(lines[i].trim())
|
||||
if (lines[i].trim().startsWith('- ') && bulletMatch) {
|
||||
const text = bulletMatch[1].trim()
|
||||
if (text) bullets.push(text)
|
||||
}
|
||||
}
|
||||
if (bullets.length === 0) {
|
||||
fail(`version ${version}'s CHANGELOG entry has no bullets to derive a headline/items from`)
|
||||
}
|
||||
return bullets
|
||||
}
|
||||
|
||||
/**
|
||||
* Derive a wall entry from a CHANGELOG.md.
|
||||
* @param {{product: string, version: string, date: string, changelogPath: string, url?: string, thumb?: string | null}} opts
|
||||
* @returns {{version: string, date: string, headline: string, items: string[], url: string, thumb: string | null}}
|
||||
*/
|
||||
function deriveEntry({ product, version, date, changelogPath, url, thumb }) {
|
||||
if (!parseSemver(version)) fail(`--version "${version}" is not a semver string`)
|
||||
if (!/^\d{4}-\d{2}-\d{2}$/.test(date) || Number.isNaN(Date.parse(date))) {
|
||||
fail(`--date "${date}" is not a YYYY-MM-DD date`)
|
||||
}
|
||||
if (!existsSync(changelogPath)) fail(`--from-changelog "${changelogPath}" does not exist`)
|
||||
|
||||
const changelog = readFileSync(changelogPath, 'utf8')
|
||||
const items = extractChangelogBullets(changelog, version)
|
||||
const headline = items[0]
|
||||
|
||||
const pattern = RELEASE_URL_PATTERNS[product]
|
||||
if (url === undefined && pattern === undefined) {
|
||||
throw new Error(`wall-entry: no permalink pattern for product "${product}" — add one to RELEASE_URL_PATTERNS or pass --url; entries never carry url: null`)
|
||||
}
|
||||
const resolvedUrl = url !== undefined ? url : pattern(version)
|
||||
const resolvedThumb = thumb !== undefined ? thumb : null
|
||||
|
||||
return { version, date, headline, items, url: resolvedUrl, thumb: resolvedThumb }
|
||||
}
|
||||
|
||||
/**
|
||||
* Load and shape-validate a wall file.
|
||||
* @param {string} filePath
|
||||
* @returns {Record<string, any>}
|
||||
*/
|
||||
function loadWallFile(filePath) {
|
||||
if (!existsSync(filePath)) fail(`"${filePath}" does not exist`)
|
||||
/** @type {unknown} */
|
||||
let data
|
||||
try {
|
||||
data = JSON.parse(readFileSync(filePath, 'utf8'))
|
||||
} catch (err) {
|
||||
fail(`"${filePath}" is not valid JSON: ${/** @type {Error} */ (err).message}`)
|
||||
}
|
||||
const errors = validateShape(data)
|
||||
if (errors.length) {
|
||||
fail(`"${filePath}" fails shape validation —\n ${errors.join('\n ')}`)
|
||||
}
|
||||
return /** @type {Record<string, any>} */ (data)
|
||||
}
|
||||
|
||||
/**
|
||||
* Run a git command, throwing an Error whose message is git's own stderr
|
||||
* (trimmed) on failure — every caller wraps this to name the cure.
|
||||
* @param {string[]} args
|
||||
* @param {string} cwd
|
||||
* @returns {string} stdout, trimmed.
|
||||
*/
|
||||
function git(args, cwd) {
|
||||
try {
|
||||
return execFileSync('git', args, { cwd, encoding: 'utf8', stdio: ['ignore', 'pipe', 'pipe'] }).trim()
|
||||
} catch (err) {
|
||||
const stderr = /** @type {any} */ (err).stderr
|
||||
const message = (typeof stderr === 'string' && stderr.trim()) || /** @type {Error} */ (err).message
|
||||
throw new Error(message)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Ensure a clean, up-to-date local clone of the releases repo at
|
||||
* `cacheDir`, checked out on `main` — cloning fresh if `cacheDir` has no
|
||||
* `.git`, otherwise fetching and hard-resetting onto `origin/main` (so a
|
||||
* stray local commit or edit left by a previous failed run can never leak
|
||||
* into the next one).
|
||||
* @param {string} remote
|
||||
* @param {string} cacheDir
|
||||
*/
|
||||
function ensureReleasesClone(remote, cacheDir) {
|
||||
if (existsSync(join(cacheDir, '.git'))) {
|
||||
try {
|
||||
git(['remote', 'set-url', 'origin', remote], cacheDir)
|
||||
git(['fetch', '--prune', 'origin'], cacheDir)
|
||||
git(['checkout', 'main'], cacheDir)
|
||||
git(['reset', '--hard', 'origin/main'], cacheDir)
|
||||
git(['clean', '-fd'], cacheDir)
|
||||
} catch (err) {
|
||||
fail(
|
||||
`cannot refresh the cached releases checkout at "${cacheDir}" from "${remote}" — ${/** @type {Error} */ (err).message}\n` +
|
||||
` cure: delete "${cacheDir}" and re-run so it re-clones from scratch, or confirm SSH access with "ssh -T git@source.soulcraft.com"`,
|
||||
)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
mkdirSync(dirname(cacheDir), { recursive: true })
|
||||
try {
|
||||
git(['clone', remote, cacheDir], dirname(cacheDir))
|
||||
} catch (err) {
|
||||
fail(
|
||||
`cannot clone "${remote}" — ${/** @type {Error} */ (err).message}\n` +
|
||||
` cure: confirm SSH access with "ssh -T git@source.soulcraft.com" and that the soulcraftlabs/releases repo exists yet`,
|
||||
)
|
||||
}
|
||||
try {
|
||||
git(['checkout', 'main'], cacheDir)
|
||||
} catch (err) {
|
||||
fail(
|
||||
`cloned "${remote}" into "${cacheDir}" but could not check out "main" — ${/** @type {Error} */ (err).message}\n` +
|
||||
` cure: confirm the releases repo's default branch is named "main"`,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Prepend `entry` to the wall at `<cacheDir>/<product>.json`, replacing any
|
||||
* existing entry for the same version (idempotent re-runs), validating
|
||||
* before and after, committing, and pushing — or refusing loudly, naming
|
||||
* the cure, at whichever step fails.
|
||||
* @param {{version: string, date: string, headline: string, items: string[], url: string, thumb: string | null}} entry
|
||||
* @param {string} product
|
||||
* @param {string} remote
|
||||
* @param {string} cacheDir
|
||||
*/
|
||||
function publishEntry(entry, product, remote, cacheDir) {
|
||||
ensureReleasesClone(remote, cacheDir)
|
||||
|
||||
const filePath = join(cacheDir, `${product}.json`)
|
||||
if (!existsSync(filePath)) {
|
||||
fail(
|
||||
`"${filePath}" does not exist in the releases repo — cure: seed "${product}.json" at the repo root first (it must exist before any release rail can prepend to it)`,
|
||||
)
|
||||
}
|
||||
const wall = loadWallFile(filePath)
|
||||
|
||||
if (wall.product !== product) {
|
||||
fail(`"${filePath}" has product "${wall.product}", but --product "${product}" was given — refusing a cross-product write`)
|
||||
}
|
||||
|
||||
const replacing = wall.entries.some((e) => e.version === entry.version)
|
||||
wall.entries = [entry, ...wall.entries.filter((e) => e.version !== entry.version)]
|
||||
|
||||
const postErrors = validateShape(wall)
|
||||
if (postErrors.length) {
|
||||
fail(`the entry for ${entry.version} would leave "${filePath}" invalid —\n ${postErrors.join('\n ')}`)
|
||||
}
|
||||
|
||||
writeFileSync(filePath, JSON.stringify(wall, null, 2) + '\n', 'utf8')
|
||||
|
||||
const status = git(['status', '--porcelain', '--', `${product}.json`], cacheDir)
|
||||
if (status === '') {
|
||||
console.log(`wall-entry: "${product}.json" already carries an identical entry for ${entry.version} — nothing to commit or push`)
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
git(['add', `${product}.json`], cacheDir)
|
||||
git(['commit', '-m', `chore(wall): ${product} ${entry.version}`], cacheDir)
|
||||
} catch (err) {
|
||||
fail(`cannot commit the wall entry in "${cacheDir}" — ${/** @type {Error} */ (err).message}\n cure: inspect "${cacheDir}" by hand and re-run once its git state is clean`)
|
||||
}
|
||||
|
||||
try {
|
||||
git(['push', 'origin', 'main'], cacheDir)
|
||||
} catch (err) {
|
||||
fail(
|
||||
`push to "${remote}" failed (likely a non-fast-forward — another release landed on main first) — ${/** @type {Error} */ (err).message}\n` +
|
||||
` cure: re-run this release step; it re-fetches and resets onto the latest origin/main before retrying`,
|
||||
)
|
||||
}
|
||||
|
||||
const sha = git(['rev-parse', 'HEAD'], cacheDir)
|
||||
console.log(
|
||||
`wall-entry: ${replacing ? 'replaced' : 'wrote'} v${entry.version} in "${product}.json" (${wall.entries.length} entries, newest first) — pushed ${sha} to ${remote} main`,
|
||||
)
|
||||
}
|
||||
|
||||
function main() {
|
||||
const args = parseArgs(process.argv.slice(2))
|
||||
|
||||
if (args.check) {
|
||||
const filePath = /** @type {string | undefined} */ (args.file)
|
||||
if (!filePath) fail('--check needs --file <path>')
|
||||
const wall = loadWallFile(/** @type {string} */ (filePath))
|
||||
console.log(`wall-entry --check: "${filePath}" OK — product "${wall.product}", ${wall.entries.length} entries, newest-first, no duplicates`)
|
||||
process.exit(0)
|
||||
}
|
||||
|
||||
// Generate mode (default, also covers --dry-run): --product, --version,
|
||||
// --date, --from-changelog required.
|
||||
const product = /** @type {string | undefined} */ (args.product)
|
||||
const version = /** @type {string | undefined} */ (args.version)
|
||||
const date = /** @type {string | undefined} */ (args.date)
|
||||
const fromChangelog = /** @type {string | undefined} */ (args['from-changelog'])
|
||||
|
||||
const missing = []
|
||||
if (!product) missing.push('--product')
|
||||
if (!version) missing.push('--version')
|
||||
if (!date) missing.push('--date')
|
||||
if (!fromChangelog) missing.push('--from-changelog')
|
||||
if (missing.length) {
|
||||
fail(
|
||||
`missing required flag(s): ${missing.join(', ')}\n` +
|
||||
'Usage:\n' +
|
||||
' wall-entry.mjs --product <p> --version <v> --date <YYYY-MM-DD> --from-changelog <CHANGELOG.md> [--dry-run]\n' +
|
||||
' wall-entry.mjs --check --file <path/to/product.json>',
|
||||
)
|
||||
}
|
||||
|
||||
const urlArg = args.url === true ? undefined : /** @type {string | undefined} */ (args.url)
|
||||
const thumbArg = args.thumb === true ? undefined : /** @type {string | undefined} */ (args.thumb)
|
||||
|
||||
const entry = deriveEntry({
|
||||
product: /** @type {string} */ (product),
|
||||
version: /** @type {string} */ (version),
|
||||
date: /** @type {string} */ (date),
|
||||
changelogPath: /** @type {string} */ (fromChangelog),
|
||||
url: urlArg,
|
||||
thumb: thumbArg,
|
||||
})
|
||||
|
||||
const remote = /** @type {string} */ (args.remote ?? process.env.WALL_ENTRY_RELEASES_REMOTE ?? DEFAULT_REMOTE)
|
||||
const cacheDir = /** @type {string} */ (args['cache-dir'] ?? process.env.WALL_ENTRY_RELEASES_CACHE_DIR ?? defaultCacheDir())
|
||||
|
||||
if (args['dry-run']) {
|
||||
console.log(`wall-entry --dry-run: would write to "${join(cacheDir, `${product}.json`)}" in ${remote} (main), pushed as "chore(wall): ${product} ${version}"`)
|
||||
console.log(JSON.stringify(entry, null, 2))
|
||||
process.exit(0)
|
||||
}
|
||||
|
||||
publishEntry(entry, /** @type {string} */ (product), remote, cacheDir)
|
||||
}
|
||||
|
||||
main()
|
||||
1038
src/brainy.ts
1038
src/brainy.ts
File diff suppressed because it is too large
Load diff
|
|
@ -40,10 +40,7 @@
|
|||
* The manifest (`_generations/facts/manifest.json`, JSON — forensics stay
|
||||
* terminal-readable) is the single source of truth for the segment SET;
|
||||
* rotation flips it atomically (write-new → fsync → rename) BEFORE the new
|
||||
* tail's first byte exists, so no segment file is ever unaccounted for. Its
|
||||
* per-segment `firstGeneration`/`lastGeneration` are LOAD-BEARING at open: a
|
||||
* recovery pass looking for facts above a bound reads only the segments those
|
||||
* bounds cannot rule out (the prune law — see `segmentsHoldingFactsAbove`).
|
||||
* tail's first byte exists, so no segment file is ever unaccounted for.
|
||||
*
|
||||
* ## Mixed-version logs (the v2 live-write cutover)
|
||||
*
|
||||
|
|
@ -692,74 +689,6 @@ function parseSegment(
|
|||
return { facts, validBytes: offset, formatVersion: FACT_LOG_FORMAT_V1 }
|
||||
}
|
||||
|
||||
/**
|
||||
* THE PRUNE LAW — which segment files a pass looking for facts ABOVE
|
||||
* `committedGeneration` actually has to read, and how many the manifest's own
|
||||
* recorded bounds took off the table.
|
||||
*
|
||||
* A sealed segment's `lastGeneration` is written at SEAL time and never
|
||||
* mutated upward afterwards ({@link FactLog.rotate}, unchanged since the log
|
||||
* was introduced): the tail's bytes are fsynced FIRST (`await this.sync()` —
|
||||
* "sealed segments are always fully durable"), the entry is then built from
|
||||
* the content that fsync covered, and only then does the manifest flip —
|
||||
* atomically (tmp+rename) and fsynced — which in the SAME write re-points
|
||||
* `tailSegment` at a new file, so the sealed file is never appended to again.
|
||||
* A crash anywhere in that order is safe in the pruning direction: crash
|
||||
* before the manifest write and the segment is still the TAIL (read whole);
|
||||
* crash after it and the entry describes bytes that were already durable. The
|
||||
* only later mutation of a sealed segment is `open()`'s straddle truncation,
|
||||
* which REMOVES facts and re-derives the entry from the actual bytes — so a
|
||||
* recorded bound can drift DOWN with its file, never up.
|
||||
*
|
||||
* Therefore: `lastGeneration = L` proves the file holds no fact above L, and
|
||||
* a pass above `committedGeneration >= L` can skip it whole — no read, no
|
||||
* CRC decode, no msgpack. What the manifest cannot PROVE is never pruned: an
|
||||
* entry with no numeric `lastGeneration` (a legacy or hand-repaired manifest)
|
||||
* is read, and the unsealed tail is always read.
|
||||
*
|
||||
* This is the difference between an open that costs O(whole fact log) and one
|
||||
* that costs O(the facts that could matter). MEASURED in production: a 16k-row
|
||||
* brain at generation ~478,819 paid 34-37s of segment reads and CRC decoding
|
||||
* in `generation-store-open-fold` on EVERY open — to answer a question whose
|
||||
* answer, after a clean close, is always "nothing".
|
||||
*/
|
||||
function segmentsHoldingFactsAbove(
|
||||
stored: FactsManifest,
|
||||
committedGeneration: number
|
||||
): { files: string[]; pruned: number } {
|
||||
const files: string[] = []
|
||||
let pruned = 0
|
||||
for (const entry of stored.segments) {
|
||||
const last = (entry as Partial<SegmentEntry>).lastGeneration
|
||||
if (typeof last === 'number' && Number.isFinite(last) && last <= committedGeneration) {
|
||||
pruned++
|
||||
continue
|
||||
}
|
||||
files.push(entry.file)
|
||||
}
|
||||
if (stored.tailSegment) files.push(stored.tailSegment)
|
||||
return { files, pruned }
|
||||
}
|
||||
|
||||
/**
|
||||
* Say what the open actually read. One line, and only when the log holds more
|
||||
* than one segment (a single-segment log has nothing to prune and nothing to
|
||||
* report) — the operator's receipt that the open is paying for the tail, not
|
||||
* for the whole history.
|
||||
*/
|
||||
function narrateAboveScan(
|
||||
pass: string,
|
||||
committedGeneration: number,
|
||||
read: number,
|
||||
pruned: number
|
||||
): void {
|
||||
if (read + pruned <= 1) return
|
||||
prodLog.narrate(
|
||||
`[FactLog] ${pass} above generation ${committedGeneration}: ${read} segment(s) read, ` +
|
||||
`${pruned} pruned of ${read + pruned} (sealed at or below the bound)`
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* The generation fact log. One instance per open store; every method assumes
|
||||
* the single-writer discipline the generation store already enforces (calls
|
||||
|
|
@ -825,6 +754,22 @@ export class FactLog {
|
|||
return this.manifest.brainId !== undefined || this.tailVersion === FACT_LOG_FORMAT_V2
|
||||
}
|
||||
|
||||
/**
|
||||
* Open the log and reconcile it to committed truth: read the manifest,
|
||||
* establish the tail's intact content (torn-tail scan), then TRUNCATE any
|
||||
* fact with `generation > committedGeneration` — those never committed (a
|
||||
* crash between fact-append and the commit point). After open, the log is
|
||||
* exactly the committed prefix.
|
||||
*/
|
||||
/**
|
||||
* Read (without truncating) every intact fact ABOVE a generation — the
|
||||
* log-authority recovery surface: after a crash, facts beyond the
|
||||
* manifest watermark that survived with valid CRCs are ACKED writes in
|
||||
* durable-at-ack mode, and the owner REPLAYS them instead of letting
|
||||
* open() truncate them. Must be called BEFORE open() (it reads the raw
|
||||
* segments directly; the torn tail's invalid suffix is ignored exactly
|
||||
* like open() would).
|
||||
*/
|
||||
/**
|
||||
* STREAMING twin of {@link FactLog.peekFactsAbove} for the recovery fold:
|
||||
* yields facts above the bound one SEGMENT at a time, ascending, without
|
||||
|
|
@ -834,18 +779,13 @@ export class FactLog {
|
|||
* Works manifest-direct (safe before {@link FactLog.open}). Ordering is
|
||||
* structural (segments rotate in order; appends are ordered within one) and
|
||||
* ASSERTED — a violation aborts loudly, never a silent misordered replay.
|
||||
*
|
||||
* Reads only the segments that CAN hold a fact above the bound — see
|
||||
* {@link segmentsHoldingFactsAbove}. A bounded fold above a high checkpoint
|
||||
* therefore reads its own tail, not the whole history it already proved
|
||||
* durable.
|
||||
*/
|
||||
async *streamFactsAbove(committedGeneration: number): AsyncGenerator<CommitFact[], void> {
|
||||
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
||||
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return
|
||||
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return
|
||||
const { files, pruned } = segmentsHoldingFactsAbove(stored, committedGeneration)
|
||||
narrateAboveScan('recovery fold', committedGeneration, files.length, pruned)
|
||||
const files = [...stored.segments.map((s) => s.file)]
|
||||
if (stored.tailSegment) files.push(stored.tailSegment)
|
||||
let lastGen = committedGeneration
|
||||
for (const file of files) {
|
||||
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
|
||||
|
|
@ -867,27 +807,13 @@ export class FactLog {
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Read (without truncating) every intact fact ABOVE a generation — the
|
||||
* log-authority recovery surface: after a crash, facts beyond the
|
||||
* manifest watermark that survived with valid CRCs are ACKED writes in
|
||||
* durable-at-ack mode, and the owner REPLAYS them instead of letting
|
||||
* open() truncate them. Must be called BEFORE open() (it reads the raw
|
||||
* segments directly; the torn tail's invalid suffix is ignored exactly
|
||||
* like open() would).
|
||||
*
|
||||
* Reads only the segments that CAN hold such a fact — see
|
||||
* {@link segmentsHoldingFactsAbove}. This runs on EVERY log-authority open,
|
||||
* including the clean one where the answer is always empty, so the segments
|
||||
* the manifest already proves irrelevant are never opened at all.
|
||||
*/
|
||||
async peekFactsAbove(committedGeneration: number): Promise<CommitFact[]> {
|
||||
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
||||
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return []
|
||||
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return []
|
||||
const out: CommitFact[] = []
|
||||
const { files, pruned } = segmentsHoldingFactsAbove(stored, committedGeneration)
|
||||
narrateAboveScan('above-manifest peek', committedGeneration, files.length, pruned)
|
||||
const files = [...stored.segments.map((s) => s.file)]
|
||||
if (stored.tailSegment) files.push(stored.tailSegment)
|
||||
for (const file of files) {
|
||||
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
|
||||
if (bytes === null) continue
|
||||
|
|
@ -900,13 +826,6 @@ export class FactLog {
|
|||
return out
|
||||
}
|
||||
|
||||
/**
|
||||
* Open the log and reconcile it to committed truth: read the manifest,
|
||||
* establish the tail's intact content (torn-tail scan), then TRUNCATE any
|
||||
* fact with `generation > committedGeneration` — those never committed (a
|
||||
* crash between fact-append and the commit point). After open, the log is
|
||||
* exactly the committed prefix.
|
||||
*/
|
||||
async open(committedGeneration: number): Promise<void> {
|
||||
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
|
||||
if (stored && typeof stored === 'object' && Array.isArray(stored.segments)) {
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@
|
|||
* 🧠 BRAINY EMBEDDED TYPE EMBEDDINGS
|
||||
*
|
||||
* AUTO-GENERATED - DO NOT EDIT
|
||||
* Generated: 2026-08-27T09:18:45-07:00
|
||||
* Generated: 2026-06-29T10:04:19-07:00
|
||||
* Noun Types: 42
|
||||
* Verb Types: 127
|
||||
*
|
||||
|
|
@ -19,7 +19,7 @@ export const TYPE_METADATA = {
|
|||
verbTypes: 127,
|
||||
totalTypes: 169,
|
||||
embeddingDimensions: 384,
|
||||
generatedAt: "2026-08-27T09:18:45-07:00",
|
||||
generatedAt: "2026-06-29T10:04:19-07:00",
|
||||
sizeBytes: {
|
||||
embeddings: 259584,
|
||||
base64: 346112
|
||||
|
|
|
|||
|
|
@ -411,90 +411,7 @@ export interface MetadataIndexProvider {
|
|||
* @returns The matching id universe as an opaque set.
|
||||
*/
|
||||
getIdSetForFilter?(filter: any): Promise<OpaqueIdSet>
|
||||
/**
|
||||
* @description OPTIONAL: evaluate `filter` over `ids` ONLY and return the
|
||||
* survivors in the caller's order — the door a graph-first
|
||||
* `find({ connected, where })` walks. The neighbour set is the universe there,
|
||||
* so the filter must cost O(|ids|) membership checks, never a whole-store
|
||||
* materialization. A native index answers from its roaring filter result
|
||||
* (membership by entity int); the reference index answers from its own
|
||||
* `getIdsForFilter`, so the two doors can never disagree. Absent → Brainy
|
||||
* intersects `getIdsForFilter`'s answer with `ids` itself (correct, O(store)).
|
||||
* @param filter - The same filter shape accepted by `getIdsForFilter`.
|
||||
* @param ids - The candidate ids (canonical). The answer is a subsequence.
|
||||
*/
|
||||
filterIdsWithin?(filter: any, ids: readonly string[]): Promise<string[]>
|
||||
/**
|
||||
* @description OPTIONAL: plan and execute a WHOLE `find()` — the graph
|
||||
* traversal, the metadata filter, the ordering and the page — and answer the
|
||||
* page's ids, or `null` for a shape this index does not plan.
|
||||
*
|
||||
* The doors above each serve one stage, so a `find()` that consults three of
|
||||
* them crosses into the index three times and marshals a result set at every
|
||||
* crossing. An index that can decide the stage ORDER itself does the whole
|
||||
* thing in one call and materializes ids only for the page — a filter
|
||||
* matching a hundred thousand rows then builds twenty-five id strings instead
|
||||
* of a hundred thousand.
|
||||
*
|
||||
* The contract this door must keep, because Brainy cannot check it:
|
||||
*
|
||||
* - **The same answer.** Identical rows, in identical order, to what the
|
||||
* stage doors would have produced for the same params. This door changes
|
||||
* which code runs, never what the answer is.
|
||||
* - **The law of the stages** (`find({ connected })` is graph-first): the
|
||||
* neighbour set is the candidate universe, the filter is evaluated over
|
||||
* those ids only, `orderBy` sorts the whole candidate set, and the page is
|
||||
* cut LAST.
|
||||
* - **`null` before work, not instead of an answer.** A shape the index does
|
||||
* not plan must be handed back BEFORE any evaluation, so Brainy serves it
|
||||
* through the stage doors exactly as it always has. Returning `null` after
|
||||
* partial work, or an empty page for a shape it could not evaluate, is a
|
||||
* silent wrong answer.
|
||||
* - **`emptyAt` names the stage** that produced an empty page — `'graph'`,
|
||||
* `'filter'`, `'visibility'` or `'none'` — so Brainy can apply its serving
|
||||
* law to the right index. An empty answer from an index that is not
|
||||
* serving must refuse loudly, and Brainy can only re-verify what it is told.
|
||||
*
|
||||
* Absent → every `find()` is served by the stage doors, which is Brainy's
|
||||
* own behaviour and the ordering oracle for any implementation of this one.
|
||||
* @param params - The find params, already normalized by `find()`
|
||||
* (natural-language parsed, `connected` anchors resolved to canonical ids,
|
||||
* an empty `where` dropped).
|
||||
* @param hiddenIds - Ids this read must not return. The contract is the ANSWER, not the
|
||||
* mechanism: a provider may subtract this set before paging, or derive the
|
||||
* same exclusion from the params' visibility tiers itself — either way the
|
||||
* page must equal the engine's own answer with none of these ids in it.
|
||||
* @param graphIndex - The active graph provider, for a `connected` plan.
|
||||
* @returns The page's ids plus the stage that emptied it, or `null`.
|
||||
*/
|
||||
planFindPage?(
|
||||
params: any,
|
||||
hiddenIds: readonly string[],
|
||||
graphIndex: unknown
|
||||
): Promise<{ ids: string[]; emptyAt: 'graph' | 'filter' | 'visibility' | 'none' } | null>
|
||||
getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>>
|
||||
/**
|
||||
* @description OPTIONAL: score `query` over `ids` ONLY — the text-leg twin of
|
||||
* {@link filterIdsWithin}, and the door a hybrid `find({ query, where })`
|
||||
* walks. The metadata filter's universe is the candidate set there, so the
|
||||
* text leg must cost O(|ids|) membership checks and marshal at most `|ids|`
|
||||
* rows, never the whole posting list of every query word. A native index
|
||||
* intersects its own postings with the candidate set (membership by entity
|
||||
* int) before any string crosses the boundary; the reference index answers
|
||||
* from its own `getIdsForTextQuery`, so the two doors can never disagree.
|
||||
* Absent → Brainy intersects `getIdsForTextQuery`'s answer with `ids` itself
|
||||
* (correct, and still hydrate-last, but it marshals the whole answer).
|
||||
*
|
||||
* The answer keeps `getIdsForTextQuery`'s contract: `{ id, matchCount }`
|
||||
* sorted by `matchCount` descending, ties in the order the whole-store answer
|
||||
* would have produced. Only rows in `ids` may appear.
|
||||
* @param query - The same text query accepted by `getIdsForTextQuery`.
|
||||
* @param ids - The candidate ids (canonical). The answer is a subset.
|
||||
*/
|
||||
getIdsForTextQueryWithin?(
|
||||
query: string,
|
||||
ids: readonly string[]
|
||||
): Promise<Array<{ id: string; matchCount: number }>>
|
||||
getSortedIdsForFilter(filter: any, orderBy: string, order?: 'asc' | 'desc', topK?: number): Promise<string[]>
|
||||
getFilterValues(field: string): Promise<string[]>
|
||||
getFilterFields(): Promise<string[]>
|
||||
|
|
|
|||
|
|
@ -1089,10 +1089,6 @@ export abstract class BaseStorageAdapter implements StorageAdapter {
|
|||
|
||||
// Counts changed since the last persist? Drives the write-through flush.
|
||||
protected pendingCountPersist = false
|
||||
/** The one persist running right now, if any (single-flight law — see flushCounts). */
|
||||
private countPersistInFlight: Promise<void> | null = null
|
||||
/** The one trailing persist a burst has queued behind the in-flight one. */
|
||||
private countPersistTrailing: Promise<void> | null = null
|
||||
|
||||
/**
|
||||
* Get total noun count - O(1) operation
|
||||
|
|
@ -1345,46 +1341,15 @@ export abstract class BaseStorageAdapter implements StorageAdapter {
|
|||
return
|
||||
}
|
||||
|
||||
// SINGLE-FLIGHT, COALESCED. Counts are write-through on every change, so
|
||||
// a burst of writes used to launch one persist per change, all in flight
|
||||
// together. Two of them inside the same millisecond shared the atomic
|
||||
// writer's temp path (`.tmp-<pid>-<ms>`): both wrote it, the first rename
|
||||
// consumed it, the second rename found nothing — ENOENT, ~1,500 times a
|
||||
// day on a busy production brain, with a full ledger write per change
|
||||
// behind it. Now exactly one persist runs at a time; requests that arrive
|
||||
// while it runs collapse into ONE trailing persist that carries the final
|
||||
// state. A burst of N changes costs at most two writes and never races
|
||||
// itself.
|
||||
if (this.countPersistInFlight) {
|
||||
// The in-flight write may have already serialised a stale snapshot —
|
||||
// ask for one more pass after it, and let every caller in this burst
|
||||
// await that same pass.
|
||||
if (!this.countPersistTrailing) {
|
||||
this.countPersistTrailing = this.countPersistInFlight
|
||||
.catch(() => undefined)
|
||||
.then(() => {
|
||||
this.countPersistTrailing = null
|
||||
return this.flushCounts()
|
||||
})
|
||||
}
|
||||
return this.countPersistTrailing
|
||||
try {
|
||||
// Persist to storage (implemented by subclass)
|
||||
await this.persistCounts()
|
||||
this.pendingCountPersist = false
|
||||
} catch (error) {
|
||||
console.error('CRITICAL: Failed to flush counts to storage:', error)
|
||||
// Keep pending flag set so we retry on next operation
|
||||
throw error
|
||||
}
|
||||
|
||||
this.countPersistInFlight = (async () => {
|
||||
try {
|
||||
// Persist to storage (implemented by subclass)
|
||||
this.pendingCountPersist = false
|
||||
await this.persistCounts()
|
||||
} catch (error) {
|
||||
// Keep the flag set so the next operation retries.
|
||||
this.pendingCountPersist = true
|
||||
console.error('CRITICAL: Failed to flush counts to storage:', error)
|
||||
throw error
|
||||
} finally {
|
||||
this.countPersistInFlight = null
|
||||
}
|
||||
})()
|
||||
return this.countPersistInFlight
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -2400,15 +2400,8 @@ export class FileSystemStorage extends BaseStorage {
|
|||
* Atomic write via temp-file-then-rename so concurrent readers never see a
|
||||
* half-written lock JSON. Reused by writer-lock writes + heartbeat.
|
||||
*/
|
||||
/** Monotonic per-process sequence so two atomic writes never share a temp path. */
|
||||
private static atomicWriteSeq = 0
|
||||
|
||||
private async writeFileAtomic(filePath: string, contents: string): Promise<void> {
|
||||
// pid + timestamp alone collided: two writers of the same target inside
|
||||
// one millisecond shared this path, and the loser's rename found the
|
||||
// winner had already moved it (ENOENT). The sequence makes every call's
|
||||
// temp path its own.
|
||||
const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}-${++FileSystemStorage.atomicWriteSeq}`
|
||||
const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}`
|
||||
await fs.promises.writeFile(tmp, contents)
|
||||
await fs.promises.rename(tmp, filePath)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2942,33 +2942,19 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
|||
!options.filter.service &&
|
||||
!options.filter.metadata
|
||||
) {
|
||||
const sourceIds = Array.isArray(options.filter.sourceId)
|
||||
? options.filter.sourceId
|
||||
: [options.filter.sourceId]
|
||||
const sourceId = Array.isArray(options.filter.sourceId)
|
||||
? options.filter.sourceId[0]
|
||||
: options.filter.sourceId
|
||||
|
||||
// EVERY requested verb type is honoured — an array used to collapse to
|
||||
// its first element here, silently dropping the rest of the ask.
|
||||
const verbTypes = new Set(
|
||||
Array.isArray(options.filter.verbType)
|
||||
? options.filter.verbType
|
||||
: [options.filter.verbType]
|
||||
)
|
||||
const verbType = Array.isArray(options.filter.verbType)
|
||||
? options.filter.verbType[0]
|
||||
: options.filter.verbType
|
||||
|
||||
// Get verbs by source (union over every requested source), filter by the
|
||||
// requested type SET (O(1) graph lookup + O(n) type filter), then apply
|
||||
// the subtype / visibility metadata filters on the candidate set.
|
||||
const bySource: HNSWVerbWithMetadata[] = []
|
||||
const seenVerbIds = new Set<string>()
|
||||
for (const oneSource of sourceIds) {
|
||||
for (const v of await this.getVerbsBySource_internal(oneSource)) {
|
||||
if (!seenVerbIds.has(v.id)) {
|
||||
seenVerbIds.add(v.id)
|
||||
bySource.push(v)
|
||||
}
|
||||
}
|
||||
}
|
||||
// Get verbs by source, then filter by type (O(1) graph lookup + O(n) type filter),
|
||||
// then apply the subtype / visibility metadata filters on the candidate set.
|
||||
const verbsBySource = await this.getVerbsBySource_internal(sourceId)
|
||||
const filteredVerbs = this.applyVerbMetadataFilters(
|
||||
bySource.filter(v => verbTypes.has(v.verb)),
|
||||
verbsBySource.filter(v => v.verb === verbType),
|
||||
options.filter
|
||||
)
|
||||
|
||||
|
|
@ -2999,22 +2985,16 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
|||
!options.filter.service &&
|
||||
!options.filter.metadata
|
||||
) {
|
||||
// EVERY requested source is honoured — an array used to collapse to
|
||||
// its first element here, silently dropping the rest of the ask.
|
||||
const onlySourceIds = Array.isArray(options.filter.sourceId)
|
||||
? options.filter.sourceId
|
||||
: [options.filter.sourceId]
|
||||
const sourceUnion: HNSWVerbWithMetadata[] = []
|
||||
const seenSourceVerbIds = new Set<string>()
|
||||
for (const oneSource of onlySourceIds) {
|
||||
for (const v of await this.getVerbsBySource_internal(oneSource)) {
|
||||
if (!seenSourceVerbIds.has(v.id)) {
|
||||
seenSourceVerbIds.add(v.id)
|
||||
sourceUnion.push(v)
|
||||
}
|
||||
}
|
||||
}
|
||||
const verbsBySource = this.applyVerbMetadataFilters(sourceUnion, options.filter)
|
||||
const sourceId = Array.isArray(options.filter.sourceId)
|
||||
? options.filter.sourceId[0]
|
||||
: options.filter.sourceId
|
||||
|
||||
// Get verbs by source directly (hydrated with metadata), then apply the
|
||||
// subtype / visibility metadata filters on the O(degree) candidate set.
|
||||
const verbsBySource = this.applyVerbMetadataFilters(
|
||||
await this.getVerbsBySource_internal(sourceId),
|
||||
options.filter
|
||||
)
|
||||
|
||||
// Apply pagination
|
||||
const paginatedVerbs = verbsBySource.slice(offset, offset + limit)
|
||||
|
|
@ -3043,22 +3023,16 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
|||
!options.filter.service &&
|
||||
!options.filter.metadata
|
||||
) {
|
||||
// EVERY requested target is honoured — an array used to collapse to
|
||||
// its first element here, silently dropping the rest of the ask.
|
||||
const onlyTargetIds = Array.isArray(options.filter.targetId)
|
||||
? options.filter.targetId
|
||||
: [options.filter.targetId]
|
||||
const targetUnion: HNSWVerbWithMetadata[] = []
|
||||
const seenTargetVerbIds = new Set<string>()
|
||||
for (const oneTarget of onlyTargetIds) {
|
||||
for (const v of await this.getVerbsByTarget_internal(oneTarget)) {
|
||||
if (!seenTargetVerbIds.has(v.id)) {
|
||||
seenTargetVerbIds.add(v.id)
|
||||
targetUnion.push(v)
|
||||
}
|
||||
}
|
||||
}
|
||||
const verbsByTarget = this.applyVerbMetadataFilters(targetUnion, options.filter)
|
||||
const targetId = Array.isArray(options.filter.targetId)
|
||||
? options.filter.targetId[0]
|
||||
: options.filter.targetId
|
||||
|
||||
// Get verbs by target directly (hydrated with metadata), then apply the
|
||||
// subtype / visibility metadata filters on the O(degree) candidate set.
|
||||
const verbsByTarget = this.applyVerbMetadataFilters(
|
||||
await this.getVerbsByTarget_internal(targetId),
|
||||
options.filter
|
||||
)
|
||||
|
||||
// Apply pagination
|
||||
const paginatedVerbs = verbsByTarget.slice(offset, offset + limit)
|
||||
|
|
@ -3087,25 +3061,16 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
|||
!options.filter.service &&
|
||||
!options.filter.metadata
|
||||
) {
|
||||
// EVERY requested verb type is honoured — an array used to collapse to
|
||||
// its first element here, silently dropping the rest of the ask.
|
||||
const verbTypes = Array.isArray(options.filter.verbType)
|
||||
? options.filter.verbType
|
||||
: [options.filter.verbType]
|
||||
const verbType = Array.isArray(options.filter.verbType)
|
||||
? options.filter.verbType[0]
|
||||
: options.filter.verbType
|
||||
|
||||
// Get verbs by each requested type (hydrated with metadata), deduped by
|
||||
// id, then apply the subtype / visibility metadata filters on the set.
|
||||
const byType: HNSWVerbWithMetadata[] = []
|
||||
const seenTypeVerbIds = new Set<string>()
|
||||
for (const oneType of verbTypes) {
|
||||
for (const v of await this.getVerbsByType_internal(oneType)) {
|
||||
if (!seenTypeVerbIds.has(v.id)) {
|
||||
seenTypeVerbIds.add(v.id)
|
||||
byType.push(v)
|
||||
}
|
||||
}
|
||||
}
|
||||
const verbsByType = this.applyVerbMetadataFilters(byType, options.filter)
|
||||
// Get verbs by type directly (hydrated with metadata), then apply the
|
||||
// subtype / visibility metadata filters on the candidate set.
|
||||
const verbsByType = this.applyVerbMetadataFilters(
|
||||
await this.getVerbsByType_internal(verbType),
|
||||
options.filter
|
||||
)
|
||||
|
||||
// Apply pagination
|
||||
const paginatedVerbs = verbsByType.slice(offset, offset + limit)
|
||||
|
|
|
|||
|
|
@ -14,7 +14,6 @@ import type { MetadataIndexManager } from '../../utils/metadataIndex.js'
|
|||
import type { GraphVerb } from '../../coreTypes.js'
|
||||
import type { Operation, RollbackAction } from '../types.js'
|
||||
import { isZeroNormVector } from '../../utils/distance.js'
|
||||
import { jsonSafeIndexMetadata } from '../../utils/jsonSafeIndexMetadata.js'
|
||||
import { prodLog } from '../../utils/logger.js'
|
||||
|
||||
/**
|
||||
|
|
@ -391,21 +390,13 @@ export class AddToMetadataIndexOperation implements Operation {
|
|||
// rollback so add + undo reference the same watermark.
|
||||
const generation = this.generationFn?.()
|
||||
|
||||
// The JSON-safe view is taken HERE, per crossing, never at construction:
|
||||
// the entity reference this op holds can be mutated between plan and
|
||||
// execute (a graph op's execute-time endpoint-int resolution mirrors
|
||||
// BigInts onto a shared verb object) — see jsonSafeIndexMetadata's
|
||||
// module doc.
|
||||
await this.index.addToIndex(
|
||||
this.id, jsonSafeIndexMetadata(this.entity), true, false, generation
|
||||
)
|
||||
// Add to metadata index (skipFlush=true for transaction atomicity)
|
||||
await this.index.addToIndex(this.id, this.entity, true, false, generation)
|
||||
|
||||
// Return rollback action
|
||||
return async () => {
|
||||
// Remove from metadata index
|
||||
await this.index.removeFromIndex(
|
||||
this.id, jsonSafeIndexMetadata(this.entity), generation
|
||||
)
|
||||
await this.index.removeFromIndex(this.id, this.entity, generation)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -441,21 +432,13 @@ export class RemoveFromMetadataIndexOperation implements Operation {
|
|||
// Resolve the removal generation once; reuse it for the rollback re-add.
|
||||
const generation = this.generationFn?.()
|
||||
|
||||
// Sanitized per crossing, never at construction — transact()'s delete
|
||||
// legs hand this op the SAME verb object the graph-retraction op's
|
||||
// execute-time endpoint resolution mutates (BigInt sourceInt/targetInt),
|
||||
// so a plan-time view aliases the pollution. See jsonSafeIndexMetadata's
|
||||
// module doc.
|
||||
await this.index.removeFromIndex(
|
||||
this.id, jsonSafeIndexMetadata(this.entity), generation
|
||||
)
|
||||
// Remove from metadata index
|
||||
await this.index.removeFromIndex(this.id, this.entity, generation)
|
||||
|
||||
// Return rollback action
|
||||
return async () => {
|
||||
// Re-add with original metadata (skipFlush=true)
|
||||
await this.index.addToIndex(
|
||||
this.id, jsonSafeIndexMetadata(this.entity), true, false, generation
|
||||
)
|
||||
await this.index.addToIndex(this.id, this.entity, true, false, generation)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,47 +0,0 @@
|
|||
/**
|
||||
* @module utils/jsonSafeIndexMetadata
|
||||
* @description The metadata-index crossing's JSON-safety law, as a leaf
|
||||
* function both the coordinator and the transaction operations share.
|
||||
*
|
||||
* The seam's metadata is JSON-safe BY CONTRACT (a native provider serializes
|
||||
* it; u64 ints as Number corrupt above 2^53) — but `resolveVerbEndpointInts`
|
||||
* MIRRORS the resolved endpoint ints onto the verb object itself as BigInt
|
||||
* (`verb.sourceInt`/`targetInt`), so a verb object reused as index metadata
|
||||
* carries BigInts into JSON.stringify, which throws, aborting the whole
|
||||
* transaction. Endpoint ints ride their OWN op params on the graph legs — the
|
||||
* metadata crossing drops every BigInt-valued top-level key instead of
|
||||
* guessing at a lossy numeric encoding.
|
||||
*
|
||||
* WHY THIS IS A LEAF MODULE, ENFORCED AT THE CROSSING: sanitizing only at
|
||||
* operation-construction time is not enough. `transact()`'s delete legs pass
|
||||
* the SAME verb object to both the graph-retraction op (whose endpoint-int
|
||||
* thunk deliberately resolves at EXECUTE time, for same-batch forward refs)
|
||||
* and the metadata-retraction op. At plan time the verb is still clean, so a
|
||||
* plan-time sanitize returns the same reference — then the graph op executes
|
||||
* first, mirrors the BigInt ints onto the shared object, and the metadata op
|
||||
* crosses the seam with them (found by the first fleet adoption of the native
|
||||
* pair: every transact-wrapped edge delete aborted). The crossing itself is
|
||||
* the only place ordering cannot bypass.
|
||||
*/
|
||||
|
||||
/**
|
||||
* A JSON-safe view of a record bound for the metadata-index crossing.
|
||||
*
|
||||
* @param metadata - The candidate index-metadata record.
|
||||
* @returns The same object when already JSON-safe, else a shallow copy
|
||||
* without the BigInt-valued keys.
|
||||
*/
|
||||
export function jsonSafeIndexMetadata(metadata: unknown): unknown {
|
||||
if (metadata === null || typeof metadata !== 'object') return metadata
|
||||
const rec = metadata as Record<string, unknown>
|
||||
let hasBigint = false
|
||||
for (const k in rec) {
|
||||
if (typeof rec[k] === 'bigint') { hasBigint = true; break }
|
||||
}
|
||||
if (!hasBigint) return metadata
|
||||
const out: Record<string, unknown> = {}
|
||||
for (const k in rec) {
|
||||
if (typeof rec[k] !== 'bigint') out[k] = rec[k]
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
|
@ -1509,56 +1509,11 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
|||
* @returns Array of { id, matchCount } sorted by matchCount descending
|
||||
*/
|
||||
async getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>> {
|
||||
return this.scoreTextQuery(query)
|
||||
}
|
||||
|
||||
/**
|
||||
* Score a text query over `ids` ONLY — the reference implementation of the
|
||||
* optional `getIdsForTextQueryWithin` door (see
|
||||
* {@link import('../plugin.js').MetadataIndexProvider}). The hybrid
|
||||
* `find({ query, where })` path passes the metadata filter's universe here so
|
||||
* the text leg ranks INSIDE that universe instead of ranking the whole store
|
||||
* and discarding the rows the filter would have dropped.
|
||||
*
|
||||
* It answers from the same posting-list merge as {@link getIdsForTextQuery},
|
||||
* with the candidate membership applied as each word's postings are counted,
|
||||
* so the two doors can never disagree: the answer is exactly the whole-store
|
||||
* answer restricted to `ids`, in the same order.
|
||||
*
|
||||
* @param query - Text query to search for.
|
||||
* @param ids - Candidate entity ids; only these may appear in the answer.
|
||||
* @returns Array of { id, matchCount } sorted by matchCount descending.
|
||||
*/
|
||||
async getIdsForTextQueryWithin(
|
||||
query: string,
|
||||
ids: readonly string[]
|
||||
): Promise<Array<{ id: string; matchCount: number }>> {
|
||||
if (ids.length === 0) return []
|
||||
return this.scoreTextQuery(query, new Set(ids))
|
||||
}
|
||||
|
||||
/**
|
||||
* The one posting-list merge behind both text doors.
|
||||
*
|
||||
* Each query word contributes AT MOST one match per entity (a posting list
|
||||
* can name an id more than once), and entities are ranked by how many of the
|
||||
* query's words they matched. `within`, when given, restricts the count to
|
||||
* those candidates — applied during the merge, so a restricted call never
|
||||
* materializes a whole-store match map.
|
||||
*
|
||||
* @param query - Text query to search for.
|
||||
* @param within - Optional candidate universe; absent = the whole store.
|
||||
* @returns Array of { id, matchCount } sorted by matchCount descending.
|
||||
*/
|
||||
private async scoreTextQuery(
|
||||
query: string,
|
||||
within?: ReadonlySet<string>
|
||||
): Promise<Array<{ id: string; matchCount: number }>> {
|
||||
const queryWords = this.tokenize(query)
|
||||
if (queryWords.length === 0) return []
|
||||
|
||||
// Count matches per entity, one word's postings at a time.
|
||||
const matchCounts = new Map<string, number>()
|
||||
// Get IDs for each word hash
|
||||
const wordIdSets: Map<string, number>[] = []
|
||||
for (const word of queryWords) {
|
||||
const wordHash = this.hashWord(word)
|
||||
let ids: string[]
|
||||
|
|
@ -1574,12 +1529,19 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
|||
throw err
|
||||
}
|
||||
}
|
||||
// One count per (word, entity) — dedupe this word's postings first.
|
||||
const counted = new Set<string>()
|
||||
const idSet = new Map<string, number>()
|
||||
for (const id of ids) {
|
||||
if (counted.has(id)) continue
|
||||
counted.add(id)
|
||||
if (within && !within.has(id)) continue
|
||||
idSet.set(id, 1)
|
||||
}
|
||||
wordIdSets.push(idSet)
|
||||
}
|
||||
|
||||
if (wordIdSets.length === 0) return []
|
||||
|
||||
// Count matches per entity
|
||||
const matchCounts = new Map<string, number>()
|
||||
for (const idSet of wordIdSets) {
|
||||
for (const [id] of idSet) {
|
||||
matchCounts.set(id, (matchCounts.get(id) || 0) + 1)
|
||||
}
|
||||
}
|
||||
|
|
@ -2613,19 +2575,6 @@ export class MetadataIndexManager implements MetadataIndexProvider {
|
|||
/** Once-per-field flag for the fallback-degradation announcement. */
|
||||
private static announcedFallbackSorts = new Set<string>()
|
||||
|
||||
/**
|
||||
* Evaluate `filter` over `ids` only — the graph-first find's door (the
|
||||
* neighbour set filtered by id, never the store filtered and then
|
||||
* intersected). This index answers from its own `getIdsForFilter`, so the
|
||||
* two doors cannot disagree; the cost is that of the filter over this
|
||||
* in-memory index, and the answer keeps the caller's order.
|
||||
*/
|
||||
async filterIdsWithin(filter: any, ids: readonly string[]): Promise<string[]> {
|
||||
if (ids.length === 0) return []
|
||||
const matched = new Set(await this.getIdsForFilter(filter))
|
||||
return ids.filter((id) => matched.has(id))
|
||||
}
|
||||
|
||||
async getSortedIdsForFilter(
|
||||
filter: any,
|
||||
orderBy: string,
|
||||
|
|
|
|||
|
|
@ -2295,31 +2295,6 @@ export class VirtualFileSystem implements IVirtualFileSystem {
|
|||
cursor = page.nextCursor
|
||||
}
|
||||
|
||||
// Pass 2: ONE paged walk over every Contains edge, grouped by target in
|
||||
// memory. The earlier shape issued one awaited related({ to }) per VFS
|
||||
// entity — O(entities) serialized graph calls, measured in whole minutes
|
||||
// on large brains. This shape is O(edges / page) calls regardless of how
|
||||
// many entities exist; mutations alone stay per-defect.
|
||||
const incomingByTarget = new Map<string, Relation<any>[]>()
|
||||
{
|
||||
const pageSize = 1000
|
||||
let pageOffset = 0
|
||||
for (;;) {
|
||||
const page = await this.brain.related({
|
||||
type: VerbType.Contains,
|
||||
limit: pageSize,
|
||||
offset: pageOffset
|
||||
})
|
||||
for (const edge of page) {
|
||||
const bucket = incomingByTarget.get(edge.to)
|
||||
if (bucket) bucket.push(edge)
|
||||
else incomingByTarget.set(edge.to, [edge])
|
||||
}
|
||||
if (page.length < pageSize) break
|
||||
pageOffset += pageSize
|
||||
}
|
||||
}
|
||||
|
||||
let removed = 0
|
||||
let restored = 0
|
||||
for (const { id, path } of vfsEntities) {
|
||||
|
|
@ -2332,7 +2307,7 @@ export class VirtualFileSystem implements IVirtualFileSystem {
|
|||
continue
|
||||
}
|
||||
|
||||
const incoming = incomingByTarget.get(id) ?? []
|
||||
const incoming = await this.brain.related({ to: id, type: VerbType.Contains })
|
||||
let expectedSeen = false
|
||||
for (const edge of incoming) {
|
||||
const isVfsEdge = edge.subtype === 'vfs-contains' || (edge.metadata as any)?.isVFS === true
|
||||
|
|
|
|||
|
|
@ -1,111 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/counts-persist-single-flight
|
||||
* @description Regression for a production race in FileSystemStorage's
|
||||
* counts ledger: `persistCounts()` was write-through on every count change
|
||||
* with no serialization, and the atomic writer named its temp file with
|
||||
* millisecond granularity (`.tmp-<pid>-<ms>`). Two persists inside one
|
||||
* millisecond shared the temp path — both wrote it, the first rename
|
||||
* consumed it, the second rename found nothing: ENOENT, ~1,500 times a day
|
||||
* on a busy production brain, with a full ledger write per change behind it.
|
||||
*
|
||||
* Under pin: persists are single-flight and coalesced — one in flight, at
|
||||
* most one trailing pass carrying the burst's final state — and every atomic
|
||||
* write owns a unique temp path. A burst of N count changes costs at most
|
||||
* two ledger writes, never errors, and leaves a ledger equal to memory.
|
||||
*/
|
||||
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'
|
||||
import * as fs from 'node:fs'
|
||||
import * as os from 'node:os'
|
||||
import * as path from 'node:path'
|
||||
import { Brainy } from '../../src/brainy.js'
|
||||
import { NounType } from '../../src/types/graphTypes.js'
|
||||
|
||||
describe('counts persistence is single-flight, coalesced, and never races its own temp file', () => {
|
||||
let dir: string
|
||||
let brain: any
|
||||
|
||||
beforeEach(async () => {
|
||||
process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true'
|
||||
dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-counts-race-'))
|
||||
brain = new Brainy({
|
||||
requireSubtype: false,
|
||||
storage: { type: 'filesystem', path: dir },
|
||||
dimensions: 384,
|
||||
silent: true
|
||||
})
|
||||
await brain.init()
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
vi.restoreAllMocks()
|
||||
await brain.close()
|
||||
fs.rmSync(dir, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
it('a burst of concurrent count changes → at most two ledger writes, zero errors, ledger == memory', async () => {
|
||||
const storage = brain.storage
|
||||
const countsPath: string = storage.countsFilePath
|
||||
expect(countsPath, 'the filesystem adapter persists a counts ledger').toBeTruthy()
|
||||
|
||||
// Let init's own persists settle so the burst is measured alone.
|
||||
await storage.flushCounts?.()
|
||||
|
||||
const renameSpy = vi.spyOn(fs.promises, 'rename')
|
||||
const errorSpy = vi.spyOn(console, 'error')
|
||||
|
||||
// Twenty-five concurrent count changes — the shape of a write burst; each
|
||||
// used to launch its own persist.
|
||||
const BURST = 25
|
||||
await Promise.all(
|
||||
Array.from({ length: BURST }, () => storage.scheduleCountPersist())
|
||||
)
|
||||
|
||||
const ledgerRenames = renameSpy.mock.calls.filter(([, to]) => String(to) === countsPath)
|
||||
expect(ledgerRenames.length, 'single-flight + one trailing pass').toBeLessThanOrEqual(2)
|
||||
expect(ledgerRenames.length, 'the burst was persisted at all').toBeGreaterThanOrEqual(1)
|
||||
|
||||
const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts'))
|
||||
expect(persistErrors).toEqual([])
|
||||
|
||||
const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8'))
|
||||
expect(ledger.totalNounCount).toBe(storage.totalNounCount)
|
||||
expect(ledger.totalVerbCount).toBe(storage.totalVerbCount)
|
||||
})
|
||||
|
||||
it('real writes in parallel: the ledger lands complete and no persist error is logged', async () => {
|
||||
const storage = brain.storage
|
||||
const countsPath: string = storage.countsFilePath
|
||||
const errorSpy = vi.spyOn(console, 'error')
|
||||
|
||||
await Promise.all(
|
||||
Array.from({ length: 12 }, (_, i) =>
|
||||
brain.add({ data: `burst row ${i}`, type: NounType.Thing })
|
||||
)
|
||||
)
|
||||
await storage.flushCounts?.()
|
||||
|
||||
const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts'))
|
||||
expect(persistErrors).toEqual([])
|
||||
const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8'))
|
||||
expect(ledger.totalNounCount).toBe(storage.totalNounCount)
|
||||
expect(await brain.getNounCount()).toBe(ledger.totalNounCount)
|
||||
})
|
||||
|
||||
it('every atomic write owns its own temp path — two writes in one millisecond never collide', async () => {
|
||||
const storage = brain.storage
|
||||
const tmpNames: string[] = []
|
||||
vi.spyOn(fs.promises, 'writeFile').mockImplementation(async (p: any) => {
|
||||
tmpNames.push(String(p))
|
||||
})
|
||||
vi.spyOn(fs.promises, 'rename').mockImplementation(async () => undefined)
|
||||
const target = path.join(dir, 'probe.json')
|
||||
await Promise.all([
|
||||
storage.writeFileAtomic(target, '{"a":1}'),
|
||||
storage.writeFileAtomic(target, '{"a":2}'),
|
||||
storage.writeFileAtomic(target, '{"a":3}')
|
||||
])
|
||||
const probeTmps = tmpNames.filter((n) => n.startsWith(`${target}.tmp-`))
|
||||
expect(probeTmps.length).toBe(3)
|
||||
expect(new Set(probeTmps).size, 'no two writes shared a temp path').toBe(3)
|
||||
})
|
||||
})
|
||||
|
|
@ -1,360 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/factlog-open-prune
|
||||
* @description THE OPEN READS THE TAIL, NOT THE HISTORY.
|
||||
*
|
||||
* Every log-authority open asks the fact log one question — "is there a fact
|
||||
* above the committed pointer?" — and until this lane existed it answered by
|
||||
* reading and CRC-decoding EVERY segment file the manifest names. MEASURED in
|
||||
* production on a 16k-row brain at generation ~478,819: 34-37 seconds inside
|
||||
* the `generation-store-open-fold` phase, on every open, including the clean
|
||||
* one where the answer is always "nothing".
|
||||
*
|
||||
* The manifest already records each sealed segment's `lastGeneration`, written
|
||||
* at seal time AFTER the segment's bytes are fsynced and into a manifest that
|
||||
* is itself written atomically and fsynced — and a sealed file is never
|
||||
* appended to again (the same manifest flip re-points `tailSegment`). So an
|
||||
* entry recording `lastGeneration ≤ committed` PROVES its file holds nothing
|
||||
* above the bound, and the open can skip it whole.
|
||||
*
|
||||
* Pinned here, from the log's own counters (the narration line), never a clock:
|
||||
*
|
||||
* 1. A clean close and reopen on a log with ≥4 sealed segments reads
|
||||
* EXACTLY the tail (1 of 6), prunes the rest, and finds nothing.
|
||||
* 2. A real SIGKILLed process that sealed segments holding facts ABOVE the
|
||||
* committed pointer: the reopen READS those sealed segments and recovers
|
||||
* byte-identically to an unpruned open (differential — the same store,
|
||||
* with the provable field stripped from its manifest, takes the full-scan
|
||||
* path and must agree fact for fact, before and after `open()`).
|
||||
* 3. A manifest entry with no `lastGeneration` (legacy, or hand-repaired) is
|
||||
* READ. Never prune what the manifest cannot prove.
|
||||
*/
|
||||
import { describe, it, expect, afterEach } from 'vitest'
|
||||
import * as fs from 'node:fs'
|
||||
import * as os from 'node:os'
|
||||
import * as path from 'node:path'
|
||||
import { spawn } from 'node:child_process'
|
||||
import {
|
||||
FactLog,
|
||||
FACTS_MANIFEST_PATH,
|
||||
type CommitFact,
|
||||
type FactLogStorage
|
||||
} from '../../src/db/factLog.js'
|
||||
import { FileSystemStorage } from '../../src/storage/adapters/fileSystemStorage.js'
|
||||
|
||||
const REPO_ROOT = process.cwd()
|
||||
const TSX = path.join(REPO_ROOT, 'node_modules', '.bin', 'tsx')
|
||||
/** ~1KB frames against a 4KB rotation threshold: ~5 facts per segment. */
|
||||
const ROTATE_BYTES = 4096
|
||||
|
||||
const tmpDirs: string[] = []
|
||||
function makeTempDir(): string {
|
||||
const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-factlog-prune-'))
|
||||
tmpDirs.push(dir)
|
||||
return dir
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
for (const dir of tmpDirs.splice(0)) {
|
||||
try {
|
||||
fs.rmSync(dir, { recursive: true, force: true })
|
||||
} catch {
|
||||
/* best effort */
|
||||
}
|
||||
try {
|
||||
fs.rmSync(`${dir}.ready.json`, { force: true })
|
||||
} catch {
|
||||
/* best effort */
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
const UUID = (n: number): string => `00000000-0000-4000-8000-${String(n).padStart(12, '0')}`
|
||||
|
||||
/** One ~1KB fact — the padding is what makes rotation cheap to provoke. */
|
||||
function fact(generation: number): CommitFact {
|
||||
return {
|
||||
generation,
|
||||
timestamp: 1_700_000_000_000 + generation,
|
||||
ops: [
|
||||
{
|
||||
kind: 'noun',
|
||||
id: UUID(generation),
|
||||
record: {
|
||||
metadata: { noun: 'document', pad: 'x'.repeat(900), g: generation },
|
||||
vector: null
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* A deterministic int minter so the log writes the V2 format production
|
||||
* writes (the prune is a manifest-level decision and never touches segment
|
||||
* bytes — but the pins should run against the bytes the fleet actually has).
|
||||
*/
|
||||
function makeMinter(): (kind: 'noun' | 'verb', id: string) => bigint {
|
||||
const ints = new Map<string, bigint>()
|
||||
return (kind, id) => {
|
||||
const key = `${kind}:${id}`
|
||||
let minted = ints.get(key)
|
||||
if (minted === undefined) {
|
||||
minted = BigInt(ints.size + 1)
|
||||
ints.set(key, minted)
|
||||
}
|
||||
return minted
|
||||
}
|
||||
}
|
||||
|
||||
/** Open a fact log over a store directory (a fresh adapter each time — this is
|
||||
* what a reopen actually does). */
|
||||
async function openStore(dir: string): Promise<{ storage: any; log: FactLog }> {
|
||||
const storage: any = new FileSystemStorage(dir)
|
||||
await storage.init()
|
||||
const log = new FactLog(storage as FactLogStorage, { rotateBytes: ROTATE_BYTES })
|
||||
log.setIntMinter(makeMinter())
|
||||
return { storage, log }
|
||||
}
|
||||
|
||||
/** Build a log of `count` facts (rotating every ~5), left durable, not closed. */
|
||||
async function buildLog(dir: string, count: number): Promise<number> {
|
||||
const { log } = await openStore(dir)
|
||||
await log.open(0)
|
||||
for (let g = 1; g <= count; g++) await log.append(fact(g))
|
||||
await log.sync()
|
||||
return log.headGeneration()
|
||||
}
|
||||
|
||||
/** Capture the narration channel (`prodLog.narrate` → console.warn). */
|
||||
async function captureNarration<T>(
|
||||
fn: () => Promise<T>
|
||||
): Promise<{ result: T; lines: string[] }> {
|
||||
const lines: string[] = []
|
||||
const original = console.warn
|
||||
console.warn = ((...args: unknown[]) => {
|
||||
lines.push(args.map((a) => String(a)).join(' '))
|
||||
}) as typeof console.warn
|
||||
try {
|
||||
return { result: await fn(), lines }
|
||||
} finally {
|
||||
console.warn = original
|
||||
}
|
||||
}
|
||||
|
||||
/** The counters the open narrated — the pin's only source of truth for what
|
||||
* was read (a wall-clock assertion could pass on a warm page cache). */
|
||||
function scanCounts(lines: string[]): { read: number; pruned: number; total: number } {
|
||||
const line = lines.find((l) => l.includes('[FactLog] above-manifest peek above generation'))
|
||||
if (!line) {
|
||||
throw new Error(`no peek narration in:\n${lines.join('\n')}`)
|
||||
}
|
||||
const match = /(\d+) segment\(s\) read, (\d+) pruned of (\d+)/.exec(line)
|
||||
if (!match) throw new Error(`unparsable peek narration: ${line}`)
|
||||
return { read: Number(match[1]), pruned: Number(match[2]), total: Number(match[3]) }
|
||||
}
|
||||
|
||||
interface SegmentEntryOnDisk {
|
||||
file: string
|
||||
firstGeneration: number
|
||||
lastGeneration?: number
|
||||
facts: number
|
||||
bytes: number
|
||||
}
|
||||
|
||||
async function readManifest(dir: string): Promise<{
|
||||
segments: SegmentEntryOnDisk[]
|
||||
tailSegment: string | null
|
||||
}> {
|
||||
const storage: any = new FileSystemStorage(dir)
|
||||
await storage.init()
|
||||
return (await storage.readRawObject(FACTS_MANIFEST_PATH)) as any
|
||||
}
|
||||
|
||||
async function rewriteManifest(
|
||||
dir: string,
|
||||
mutate: (manifest: any) => void
|
||||
): Promise<void> {
|
||||
const storage: any = new FileSystemStorage(dir)
|
||||
await storage.init()
|
||||
const manifest = await storage.readRawObject(FACTS_MANIFEST_PATH)
|
||||
mutate(manifest)
|
||||
await storage.writeRawObject(FACTS_MANIFEST_PATH, manifest)
|
||||
await storage.syncRawObjects([FACTS_MANIFEST_PATH])
|
||||
}
|
||||
|
||||
/** Every fact the log holds, in order — the recovered state, read back. */
|
||||
async function allFacts(log: FactLog): Promise<CommitFact[]> {
|
||||
const out: CommitFact[] = []
|
||||
const handle = log.scanFacts()
|
||||
for await (const batch of handle.batches()) out.push(...batch.facts)
|
||||
return out
|
||||
}
|
||||
|
||||
describe('fact log — the open reads only the segments that can hold facts above the bound', () => {
|
||||
it('a clean close + reopen over ≥4 sealed segments reads exactly the tail and finds nothing', async () => {
|
||||
const dir = makeTempDir()
|
||||
const head = await buildLog(dir, 30)
|
||||
|
||||
const manifest = await readManifest(dir)
|
||||
expect(manifest.segments.length).toBeGreaterThanOrEqual(4) // the fixture is real
|
||||
expect(manifest.tailSegment).not.toBeNull()
|
||||
|
||||
// The reopen: a clean close means committed === the log's head.
|
||||
const { log } = await openStore(dir)
|
||||
const { result: orphans, lines } = await captureNarration(() => log.peekFactsAbove(head))
|
||||
|
||||
expect(orphans).toEqual([]) // the fold finds nothing, as it always does after a clean close
|
||||
const counts = scanCounts(lines)
|
||||
expect(counts.read).toBe(1) // EXACTLY the tail
|
||||
expect(counts.total).toBe(manifest.segments.length + 1)
|
||||
expect(counts.pruned).toBe(manifest.segments.length)
|
||||
|
||||
// And the reconciling open still lands on the same committed prefix.
|
||||
await log.open(head)
|
||||
expect(log.headGeneration()).toBe(head)
|
||||
expect((await allFacts(log)).map((f) => f.generation)).toEqual(
|
||||
Array.from({ length: head }, (_, i) => i + 1)
|
||||
)
|
||||
})
|
||||
|
||||
it('a manifest entry with no lastGeneration is READ — never prune what you cannot prove', async () => {
|
||||
const dir = makeTempDir()
|
||||
const head = await buildLog(dir, 30)
|
||||
const before = await readManifest(dir)
|
||||
expect(before.segments.length).toBeGreaterThanOrEqual(4)
|
||||
|
||||
// A legacy/hand-repaired entry: the field the prune needs is simply absent.
|
||||
await rewriteManifest(dir, (m) => {
|
||||
delete m.segments[0].lastGeneration
|
||||
})
|
||||
|
||||
const { log } = await openStore(dir)
|
||||
const { result: orphans, lines } = await captureNarration(() => log.peekFactsAbove(head))
|
||||
|
||||
expect(orphans).toEqual([]) // still nothing above the bound — it was READ to find out
|
||||
const counts = scanCounts(lines)
|
||||
expect(counts.read).toBe(2) // the unprovable entry + the tail
|
||||
expect(counts.pruned).toBe(before.segments.length - 1)
|
||||
expect(counts.total).toBe(before.segments.length + 1)
|
||||
})
|
||||
|
||||
it(
|
||||
'a SIGKILLed writer that sealed segments above the committed pointer recovers identically to an unpruned open',
|
||||
async () => {
|
||||
const dir = makeTempDir()
|
||||
const readyPath = `${dir}.ready.json`
|
||||
// A real process death: the child fsyncs its segments, records what it
|
||||
// reached, and SIGKILLs ITSELF — no close, no unwind, no chance to tidy.
|
||||
const script = `
|
||||
import * as fs from 'node:fs'
|
||||
import { FactLog } from ${JSON.stringify(path.join(REPO_ROOT, 'src', 'db', 'factLog.ts'))}
|
||||
import { FileSystemStorage } from ${JSON.stringify(path.join(REPO_ROOT, 'src', 'storage', 'adapters', 'fileSystemStorage.ts'))}
|
||||
const UUID = (n) => '00000000-0000-4000-8000-' + String(n).padStart(12, '0')
|
||||
const fact = (g) => ({
|
||||
generation: g,
|
||||
timestamp: 1700000000000 + g,
|
||||
ops: [{ kind: 'noun', id: UUID(g), record: { metadata: { noun: 'document', pad: 'x'.repeat(900), g }, vector: null } }]
|
||||
})
|
||||
const ints = new Map()
|
||||
const storage = new FileSystemStorage(${JSON.stringify(dir)})
|
||||
await storage.init()
|
||||
const log = new FactLog(storage, { rotateBytes: ${ROTATE_BYTES} })
|
||||
log.setIntMinter((kind, id) => {
|
||||
const key = kind + ':' + id
|
||||
if (!ints.has(key)) ints.set(key, BigInt(ints.size + 1))
|
||||
return ints.get(key)
|
||||
})
|
||||
await log.open(0)
|
||||
for (let g = 1; g <= 30; g++) await log.append(fact(g))
|
||||
await log.sync()
|
||||
fs.writeFileSync(${JSON.stringify(readyPath)}, JSON.stringify({ head: log.headGeneration() }))
|
||||
process.kill(process.pid, 'SIGKILL')
|
||||
`
|
||||
const scriptPath = path.join(dir, 'crash-writer.mts')
|
||||
fs.writeFileSync(scriptPath, script)
|
||||
const child = spawn(TSX, [scriptPath], { cwd: REPO_ROOT, stdio: ['ignore', 'pipe', 'pipe'] })
|
||||
let output = ''
|
||||
child.stdout.on('data', (d) => { output += String(d) })
|
||||
child.stderr.on('data', (d) => { output += String(d) })
|
||||
const exit = await new Promise<{ code: number | null; signal: string | null }>((resolve) =>
|
||||
child.on('exit', (code, signal) => resolve({ code, signal }))
|
||||
)
|
||||
if (!fs.existsSync(readyPath)) {
|
||||
throw new Error(`the crash writer never reached its kill point:\n${output}`)
|
||||
}
|
||||
// Death, not a shutdown: no close(), no unwind, no orderly exit code.
|
||||
expect(exit.signal ?? `code ${exit.code}`).not.toBe('code 0')
|
||||
const head = JSON.parse(fs.readFileSync(readyPath, 'utf8')).head as number
|
||||
expect(head).toBe(30)
|
||||
|
||||
// The committed pointer the survivor comes back on: mid-log, so sealed
|
||||
// segments hold facts ABOVE it — the exact shape the prune must not skip.
|
||||
const committed = 12
|
||||
const manifest = await readManifest(dir)
|
||||
const straddling = manifest.segments.filter(
|
||||
(s) => s.firstGeneration <= committed && (s.lastGeneration ?? 0) > committed
|
||||
)
|
||||
const entirelyAbove = manifest.segments.filter((s) => s.firstGeneration > committed)
|
||||
expect(straddling.length).toBeGreaterThanOrEqual(1)
|
||||
expect(entirelyAbove.length).toBeGreaterThanOrEqual(1)
|
||||
|
||||
// THE DIFFERENTIAL. The unpruned answer, through the SAME code on the
|
||||
// SAME bytes: a peek above generation 0 can prune nothing (no sealed
|
||||
// segment ends at or below 0), so it reads every segment file and
|
||||
// decodes every frame — exactly what this open used to do — and its
|
||||
// facts above the pointer are what the fold is entitled to replay.
|
||||
const { log } = await openStore(dir)
|
||||
const { result: fullScan, lines: fullLines } = await captureNarration(() =>
|
||||
log.peekFactsAbove(0)
|
||||
)
|
||||
expect(scanCounts(fullLines)).toEqual({
|
||||
read: manifest.segments.length + 1,
|
||||
pruned: 0,
|
||||
total: manifest.segments.length + 1
|
||||
})
|
||||
const unprunedAnswer = fullScan.filter((f) => f.generation > committed)
|
||||
|
||||
const { result: prunedAnswer, lines } = await captureNarration(() =>
|
||||
log.peekFactsAbove(committed)
|
||||
)
|
||||
|
||||
// The sealed segments above the bound were READ, not skipped.
|
||||
const counts = scanCounts(lines)
|
||||
expect(counts.read).toBe(straddling.length + entirelyAbove.length + 1)
|
||||
expect(counts.pruned).toBe(manifest.segments.length - straddling.length - entirelyAbove.length)
|
||||
expect(counts.pruned).toBeGreaterThan(0) // the prune did engage, and was still right
|
||||
expect(prunedAnswer.map((f) => f.generation)).toEqual(
|
||||
Array.from({ length: head - committed }, (_, i) => committed + 1 + i)
|
||||
)
|
||||
// Facts that live in a SEALED segment (not the tail) came back.
|
||||
expect(prunedAnswer.some((f) => f.generation <= (straddling[0].lastGeneration ?? 0))).toBe(
|
||||
true
|
||||
)
|
||||
// Fact for fact, the pruned answer IS the unpruned answer — so whatever
|
||||
// the recovery replays, it replays identically.
|
||||
expect(prunedAnswer).toEqual(unprunedAnswer)
|
||||
|
||||
// The fold's streaming twin (the unclean-open path) agrees too.
|
||||
const streamed: CommitFact[] = []
|
||||
for await (const batch of log.streamFactsAbove(committed)) streamed.push(...batch)
|
||||
expect(streamed).toEqual(unprunedAnswer)
|
||||
|
||||
// And the reconciling open rolls back exactly as it always did: the two
|
||||
// never-committed sealed segments dropped, the straddling one cut, the
|
||||
// tail truncated — the log left as the committed prefix.
|
||||
await log.open(committed)
|
||||
expect(log.headGeneration()).toBe(committed)
|
||||
expect((await allFacts(log)).map((f) => f.generation)).toEqual(
|
||||
Array.from({ length: committed }, (_, i) => i + 1)
|
||||
)
|
||||
const after = await readManifest(dir)
|
||||
expect(after.segments.map((s) => s.file)).toEqual(
|
||||
manifest.segments
|
||||
.filter((s) => s.firstGeneration <= committed)
|
||||
.map((s) => s.file)
|
||||
)
|
||||
expect(after.segments[after.segments.length - 1].lastGeneration).toBe(committed)
|
||||
},
|
||||
120_000
|
||||
)
|
||||
})
|
||||
|
|
@ -1,165 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/find-connected-order
|
||||
* @description The graph-first law for `find({ connected })` (10.4.8).
|
||||
*
|
||||
* With `connected` present the neighbour set is the candidate universe: it is
|
||||
* resolved from the adjacency first, the metadata filter is evaluated over
|
||||
* those ids only, and the page is cut last. The earlier order materialized the
|
||||
* whole-store filtered id list, paged it, hydrated the page, and only then
|
||||
* intersected with the neighbours — so a neighbour outside the first page of
|
||||
* the filtered STORE was silently dropped, and every call paid O(store).
|
||||
*
|
||||
* These pins hold both halves. The answer: every matching neighbour is
|
||||
* reachable by paging, a non-neighbour never appears, a negation (`missing`)
|
||||
* is evaluated over the neighbours, `orderBy` sorts the whole neighbour set
|
||||
* before the page is cut, and the vector leg walks the neighbours only. The
|
||||
* cost shape: the metadata index is asked about the neighbour ids only, and
|
||||
* hydration is one page — never the store.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll, vi } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType, VerbType } from '../../src/types/graphTypes'
|
||||
import { v5 } from '../../src/universal/uuid'
|
||||
import { generateTestVector } from '../helpers/test-factory'
|
||||
|
||||
/** Matching rows that are NOT neighbours — added FIRST, so the whole-store filtered list leads with them. */
|
||||
const NOISE = 120
|
||||
/** Matching rows that ARE neighbours of the anchor. */
|
||||
const NEIGHBOURS = 30
|
||||
/** Neighbours carrying `retracted: true` — excluded by the `missing` negation. */
|
||||
const RETRACTED = 4
|
||||
|
||||
describe('find({ connected }) is graph-first: neighbours → filter → page', () => {
|
||||
let brain: Brainy<any>
|
||||
const anchor = 'anchor'
|
||||
const sharedVector = generateTestVector()
|
||||
const neighbourIds = new Set(Array.from({ length: NEIGHBOURS }, (_, i) => v5(`nb-${i}`)))
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
await brain.add({
|
||||
id: anchor,
|
||||
data: 'the anchor',
|
||||
type: NounType.Person,
|
||||
metadata: { kind: 'anchor' },
|
||||
vector: generateTestVector()
|
||||
})
|
||||
for (let i = 0; i < NOISE; i++) {
|
||||
await brain.add({
|
||||
id: `noise-${i}`,
|
||||
data: `noise ${i}`,
|
||||
type: NounType.Person,
|
||||
metadata: { kind: 'note', rank: 1000 + i },
|
||||
vector: sharedVector
|
||||
})
|
||||
}
|
||||
for (let i = 0; i < NEIGHBOURS; i++) {
|
||||
await brain.add({
|
||||
id: `nb-${i}`,
|
||||
data: `neighbour ${i}`,
|
||||
type: NounType.Person,
|
||||
metadata: { kind: 'note', rank: i + 1, ...(i < RETRACTED ? { retracted: true } : {}) },
|
||||
vector: sharedVector
|
||||
})
|
||||
await brain.relate({ from: anchor, to: `nb-${i}`, type: VerbType.Knows })
|
||||
}
|
||||
})
|
||||
|
||||
afterAll(async () => {
|
||||
brain = null as any
|
||||
})
|
||||
|
||||
it('returns the matching neighbours page by page — none dropped, never a non-neighbour', async () => {
|
||||
const seen = new Set<string>()
|
||||
for (let offset = 0; offset <= NEIGHBOURS; offset += 10) {
|
||||
const page = await brain.find({
|
||||
connected: { from: anchor, direction: 'out' },
|
||||
where: { kind: 'note' },
|
||||
limit: 10,
|
||||
offset
|
||||
})
|
||||
expect(page).toHaveLength(offset < NEIGHBOURS ? 10 : 0)
|
||||
for (const r of page) {
|
||||
expect(neighbourIds.has(r.entity.id)).toBe(true)
|
||||
expect(seen.has(r.entity.id)).toBe(false)
|
||||
seen.add(r.entity.id)
|
||||
}
|
||||
}
|
||||
expect(seen.size).toBe(NEIGHBOURS)
|
||||
})
|
||||
|
||||
it('evaluates a negation (`missing`) over the neighbour set, not the store', async () => {
|
||||
const results = await brain.find({
|
||||
connected: { from: anchor, direction: 'out' },
|
||||
where: { kind: 'note', retracted: { missing: true } },
|
||||
limit: 100
|
||||
})
|
||||
expect(results).toHaveLength(NEIGHBOURS - RETRACTED)
|
||||
for (const r of results) {
|
||||
expect(neighbourIds.has(r.entity.id)).toBe(true)
|
||||
expect(r.entity.metadata.retracted).toBeUndefined()
|
||||
}
|
||||
})
|
||||
|
||||
it('asks the metadata index about the neighbour ids only, and hydrates one page', async () => {
|
||||
const index = (brain as any).metadataIndex
|
||||
const within = vi.spyOn(index, 'filterIdsWithin')
|
||||
const hydrate = vi.spyOn(brain as any, 'batchGet')
|
||||
try {
|
||||
const results = await brain.find({
|
||||
connected: { from: anchor, direction: 'out' },
|
||||
where: { kind: 'note' },
|
||||
limit: 10
|
||||
})
|
||||
expect(results).toHaveLength(10)
|
||||
expect(within).toHaveBeenCalledTimes(1)
|
||||
const askedIds = within.mock.calls[0][1] as string[]
|
||||
expect(askedIds).toHaveLength(NEIGHBOURS)
|
||||
for (const id of askedIds) expect(neighbourIds.has(id)).toBe(true)
|
||||
expect(hydrate).toHaveBeenCalledTimes(1)
|
||||
expect(hydrate.mock.calls[0][0]).toHaveLength(10)
|
||||
} finally {
|
||||
within.mockRestore()
|
||||
hydrate.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('orders the WHOLE neighbour set before cutting the page', async () => {
|
||||
const results = await brain.find({
|
||||
connected: { from: anchor, direction: 'out' },
|
||||
where: { kind: 'note' },
|
||||
orderBy: 'rank',
|
||||
order: 'desc',
|
||||
limit: 5
|
||||
})
|
||||
expect(results.map((r) => r.entity.metadata.rank)).toEqual([30, 29, 28, 27, 26])
|
||||
})
|
||||
|
||||
it('walks the vector leg over the neighbours only', async () => {
|
||||
const results = await brain.find({
|
||||
vector: sharedVector,
|
||||
connected: { from: anchor, direction: 'out' },
|
||||
where: { kind: 'note' },
|
||||
limit: 5
|
||||
})
|
||||
expect(results).toHaveLength(5)
|
||||
for (const r of results) expect(neighbourIds.has(r.entity.id)).toBe(true)
|
||||
})
|
||||
|
||||
it('an anchor without neighbours answers [] before the filter is asked', async () => {
|
||||
const index = (brain as any).metadataIndex
|
||||
const within = vi.spyOn(index, 'filterIdsWithin')
|
||||
try {
|
||||
const results = await brain.find({
|
||||
connected: { from: 'noise-0', direction: 'out' },
|
||||
where: { kind: 'note' },
|
||||
limit: 10
|
||||
})
|
||||
expect(results).toEqual([])
|
||||
expect(within).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
within.mockRestore()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
|
@ -1,635 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/find-hybrid-filter-before-hydrate
|
||||
* @description FILTER BEFORE HYDRATE, applied to the hybrid `find({ query })` path.
|
||||
*
|
||||
* A hybrid find fuses two legs. The semantic leg already walked only the
|
||||
* metadata filter's universe (`candidateIds` / `allowedIds`). The TEXT leg did
|
||||
* not: it ranked the WHOLE store, took the top `limit * 4`, read every one of
|
||||
* those rows from canonical, and only then intersected with the filter — so a
|
||||
* filtered hybrid find on a large store read hundreds of rows to return a
|
||||
* handful of them, and a matching row outside the store-wide text prefix was
|
||||
* silently dropped. That is the same defect `find({ connected })` carried
|
||||
* before the graph-first law, one leg over.
|
||||
*
|
||||
* Both halves are pinned here.
|
||||
*
|
||||
* THE ANSWER. Where the filter did not truncate the text leg — the universe
|
||||
* covers every text match, so both orders rank the same rows — the new
|
||||
* pipeline's answer is IDENTICAL to the old one's: same rows, same order, same
|
||||
* scores, same match visibility, same row shape. The oracle below is the
|
||||
* pre-change pipeline itself, replayed on the same brain through the same
|
||||
* doors, so the comparison is against what actually ran, not a remembered
|
||||
* expectation.
|
||||
*
|
||||
* THE CORRECTION. Where the filter DID truncate it — the query's words are
|
||||
* common outside the universe — the old order let the text leg contribute
|
||||
* nothing at all: every row it ranked was discarded by the filter, and the
|
||||
* answer came from the semantic leg alone. The new order ranks inside the
|
||||
* universe, so the text leg contributes the rows it always should have.
|
||||
*
|
||||
* THE COST. Canonical is read for exactly the page: one batch, `limit` rows,
|
||||
* never the legs. And the text leg is asked about the universe's ids only —
|
||||
* what it marshals is bounded by the universe, not by the store.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, vi } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType, VerbType } from '../../src/types/graphTypes'
|
||||
import { rankIndicesByScore, reorderByIndices } from '../../src/utils/resultRanking'
|
||||
import { resolveEntityId } from '../../src/utils/idNormalization'
|
||||
|
||||
/** Embedding width of the default model — the row vectors must match it. */
|
||||
const DIM = 384
|
||||
|
||||
/**
|
||||
* A deterministic, per-row-distinct unit vector. Distinct so the semantic leg
|
||||
* has a real ranking to produce (identical vectors would make its order a tie
|
||||
* break), deterministic so the oracle and the pipeline see the same one.
|
||||
*/
|
||||
function seededVector(seed: number): number[] {
|
||||
const v = new Array<number>(DIM)
|
||||
for (let i = 0; i < DIM; i++) {
|
||||
v[i] = Math.sin((i + 1) * 0.11 + seed * 0.37) * 0.5 + Math.cos((i + 1) * 0.05 + seed * 0.13) * 0.3
|
||||
}
|
||||
const magnitude = Math.sqrt(v.reduce((sum, x) => sum + x * x, 0))
|
||||
return v.map((x) => x / magnitude)
|
||||
}
|
||||
|
||||
/** The fields a caller reads off a hybrid row — the whole comparable surface. */
|
||||
function project(rows: any[]): any[] {
|
||||
return rows.map((r) => ({
|
||||
id: r.id,
|
||||
score: r.score,
|
||||
type: r.type,
|
||||
metadata: r.metadata,
|
||||
textMatches: r.textMatches,
|
||||
textScore: r.textScore,
|
||||
semanticScore: r.semanticScore,
|
||||
matchSource: r.matchSource
|
||||
}))
|
||||
}
|
||||
|
||||
/**
|
||||
* The PRE-CHANGE hybrid pipeline, replayed on a live brain through the same
|
||||
* provider doors it used: whole-store text ranking with both legs hydrated in
|
||||
* full, RRF fusion, then the metadata intersection, then the page.
|
||||
*
|
||||
* Supports the shapes these pins exercise (query + where/type/excludeVFS +
|
||||
* connected + offset); `orderBy`, `fusion` and `near` are not replayed.
|
||||
*/
|
||||
async function legacyHybridFind(brain: any, params: any): Promise<any[]> {
|
||||
const index = brain.metadataIndex
|
||||
const limit = params.limit ?? 10
|
||||
const offset = params.offset ?? 0
|
||||
const hasFilter = Boolean(
|
||||
params.where || params.type || params.subtype || params.service || params.excludeVFS
|
||||
)
|
||||
|
||||
let preResolvedMetadataIds: string[] | null = null
|
||||
let preResolvedFilter: any = null
|
||||
let graphFirstIds: string[] | null = null
|
||||
|
||||
if (params.connected) {
|
||||
// find() normalizes the anchors to canonical ids before this stage runs.
|
||||
const anchored = {
|
||||
...params,
|
||||
connected: {
|
||||
...params.connected,
|
||||
...(params.connected.from && { from: resolveEntityId(params.connected.from) }),
|
||||
...(params.connected.to && { to: resolveEntityId(params.connected.to) })
|
||||
}
|
||||
}
|
||||
graphFirstIds = await brain.resolveConnectedIds(anchored)
|
||||
if (graphFirstIds!.length > 0 && hasFilter) {
|
||||
preResolvedFilter = brain.buildMetadataFilter(params)
|
||||
graphFirstIds = await brain.filterIdsWithinBelted(preResolvedFilter, graphFirstIds)
|
||||
}
|
||||
if (graphFirstIds!.length === 0) return []
|
||||
preResolvedMetadataIds = graphFirstIds
|
||||
} else if (hasFilter) {
|
||||
preResolvedFilter = brain.buildMetadataFilter(params)
|
||||
preResolvedMetadataIds = await brain.filterIdsBelted(preResolvedFilter)
|
||||
if (preResolvedMetadataIds!.length === 0) return []
|
||||
}
|
||||
|
||||
// Text leg — the whole store, then the top `limit * 4`, hydrated in full.
|
||||
const allTextMatches = await index.getIdsForTextQuery(params.query)
|
||||
const topMatches = allTextMatches.slice(0, limit * 2 * 2)
|
||||
const maxMatches = topMatches[0]?.matchCount || 1
|
||||
const textEntities = await brain.batchGet(topMatches.map((m: any) => m.id))
|
||||
const textResults = topMatches
|
||||
.filter((m: any) => textEntities.has(m.id))
|
||||
.map((m: any) => ({ id: m.id, score: m.matchCount / maxMatches }))
|
||||
|
||||
// Semantic leg — the beam walk over the universe, hydrated in full.
|
||||
const vector = await brain.embed(params.query)
|
||||
const searchOptions = preResolvedMetadataIds ? { candidateIds: preResolvedMetadataIds } : undefined
|
||||
const searchResults: [string, number][] = await brain.index.search(
|
||||
vector,
|
||||
limit * 2,
|
||||
undefined,
|
||||
searchOptions
|
||||
)
|
||||
const semanticEntities = await brain.batchGet(searchResults.map(([id]) => id))
|
||||
const semanticResults = searchResults
|
||||
.filter(([id]) => semanticEntities.has(id))
|
||||
.map(([id, distance]) => ({ id, score: Math.max(0, Math.min(1, 1 / (1 + distance))) }))
|
||||
|
||||
// RRF fusion, with the match visibility the rows carried.
|
||||
const alpha = params.hybridAlpha ?? brain.autoAlpha(params.query)
|
||||
const k = 60
|
||||
const matchData = new Map<string, any>()
|
||||
const textWeight = 1 - alpha
|
||||
textResults.forEach((r: any, rank: number) => {
|
||||
const existing = matchData.get(r.id) || { rrf: 0, hasText: false, hasSemantic: false }
|
||||
existing.rrf += textWeight * (1 / (k + rank + 1))
|
||||
existing.textScore = r.score
|
||||
existing.hasText = true
|
||||
matchData.set(r.id, existing)
|
||||
})
|
||||
semanticResults.forEach((r: any, rank: number) => {
|
||||
const existing = matchData.get(r.id) || { rrf: 0, hasText: false, hasSemantic: false }
|
||||
existing.rrf += alpha * (1 / (k + rank + 1))
|
||||
existing.semanticScore = r.score
|
||||
existing.hasSemantic = true
|
||||
matchData.set(r.id, existing)
|
||||
})
|
||||
|
||||
const queryWords: string[] = index.tokenize(params.query)
|
||||
const textResultIds = new Set(textResults.map((r: any) => r.id))
|
||||
const fusedIds = Array.from(matchData.entries())
|
||||
.sort((a, b) => b[1].rrf - a[1].rrf)
|
||||
.map(([id, data]) => ({ id, data }))
|
||||
|
||||
const allEntities = await brain.batchGet(fusedIds.map((f) => f.id))
|
||||
let rows: any[] = []
|
||||
for (const { id, data } of fusedIds) {
|
||||
const entity = allEntities.get(id)
|
||||
if (!entity) continue
|
||||
const textContent = textResultIds.has(id)
|
||||
? index.extractTextContent({ data: entity.data, metadata: entity.metadata }).toLowerCase()
|
||||
: null
|
||||
rows.push({
|
||||
id,
|
||||
score: data.rrf,
|
||||
type: entity.type,
|
||||
metadata: entity.metadata,
|
||||
textMatches:
|
||||
textContent === null ? [] : queryWords.filter((w) => textContent.includes(w.toLowerCase())),
|
||||
textScore: data.textScore,
|
||||
semanticScore: data.semanticScore,
|
||||
matchSource: data.hasText && data.hasSemantic ? 'both' : data.hasText ? 'text' : 'semantic'
|
||||
})
|
||||
}
|
||||
|
||||
// The metadata intersection — after the legs, as it was.
|
||||
if (preResolvedMetadataIds && preResolvedFilter) {
|
||||
const filteredIdSet = new Set(preResolvedMetadataIds)
|
||||
rows = rows.filter((r) => filteredIdSet.has(r.id))
|
||||
}
|
||||
if (graphFirstIds !== null) {
|
||||
const neighbourSet = new Set(graphFirstIds)
|
||||
rows = rows.filter((r) => neighbourSet.has(r.id))
|
||||
}
|
||||
|
||||
// Rank to the page, then cut it.
|
||||
const order = rankIndicesByScore(
|
||||
rows.map((r) => r.score),
|
||||
offset + limit,
|
||||
true
|
||||
)
|
||||
return reorderByIndices(rows, order).slice(offset, offset + limit)
|
||||
}
|
||||
|
||||
/**
|
||||
* FIXTURE A — the filter's universe covers every text match, so the two orders
|
||||
* rank exactly the same rows and the answers must be identical.
|
||||
*/
|
||||
describe('hybrid find: filter before hydrate — the answer is unchanged', () => {
|
||||
let brain: Brainy<any>
|
||||
const QUERY = 'orbital telemetry'
|
||||
const MATCHES = 24
|
||||
const FILLER = 120
|
||||
const OUTSIDE = 30
|
||||
const VFS = 10
|
||||
const RETRACTED = 6
|
||||
const anchor = 'array-anchor'
|
||||
const matchIds: string[] = []
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
|
||||
let seed = 1
|
||||
await brain.add({
|
||||
id: anchor,
|
||||
data: 'ground station anchor record',
|
||||
type: NounType.Thing,
|
||||
metadata: { lane: 'alpha', role: 'anchor' },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
|
||||
// Rows the query's words actually match — all inside every filter below.
|
||||
for (let i = 0; i < MATCHES; i++) {
|
||||
const id = `match-${i}`
|
||||
await brain.add({
|
||||
id,
|
||||
data: `orbital telemetry packet ${i} recorded downlink`,
|
||||
type: NounType.Document,
|
||||
metadata: { lane: 'alpha', rank: i },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
matchIds.push(resolveEntityId(id))
|
||||
await brain.relate({ from: anchor, to: id, type: VerbType.RelatedTo })
|
||||
}
|
||||
// Rows inside the universe that the query's words do NOT match.
|
||||
for (let i = 0; i < FILLER; i++) {
|
||||
await brain.add({
|
||||
id: `filler-${i}`,
|
||||
data: `cistern ledger entry ${i} archived`,
|
||||
type: NounType.Document,
|
||||
metadata: { lane: 'alpha', rank: 1000 + i },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
}
|
||||
// Rows outside the universe.
|
||||
for (let i = 0; i < OUTSIDE; i++) {
|
||||
await brain.add({
|
||||
id: `outside-${i}`,
|
||||
data: `unrelated dossier ${i}`,
|
||||
type: NounType.Person,
|
||||
metadata: { lane: 'beta' },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
}
|
||||
// VFS infrastructure rows — excluded by excludeVFS.
|
||||
for (let i = 0; i < VFS; i++) {
|
||||
await brain.add({
|
||||
id: `vfs-${i}`,
|
||||
data: `mounted path ${i}`,
|
||||
type: NounType.Document,
|
||||
metadata: { lane: 'alpha', vfsType: 'file' },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
}
|
||||
// Retracted rows — excluded by a `missing` negation.
|
||||
for (let i = 0; i < RETRACTED; i++) {
|
||||
await brain.add({
|
||||
id: `retracted-${i}`,
|
||||
data: `withdrawn note ${i}`,
|
||||
type: NounType.Document,
|
||||
metadata: { lane: 'alpha', retracted: true },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
}
|
||||
|
||||
// The reference index has no opaque-set door, so the pipeline and the
|
||||
// oracle both restrict the beam walk with the materialized candidate ids.
|
||||
expect(typeof (brain as any).metadataIndex.getIdSetForFilter).not.toBe('function')
|
||||
})
|
||||
|
||||
it('the fixture does not truncate the text leg — the universe covers every text match', async () => {
|
||||
const index = (brain as any).metadataIndex
|
||||
const textMatches = await index.getIdsForTextQuery(QUERY)
|
||||
expect(textMatches).toHaveLength(MATCHES)
|
||||
const universe = await (brain as any).filterIdsBelted({ lane: 'alpha' })
|
||||
const inUniverse = new Set(universe)
|
||||
for (const m of textMatches) expect(inUniverse.has(m.id)).toBe(true)
|
||||
})
|
||||
|
||||
it('hybrid + where: identical rows, identical order, identical scores', async () => {
|
||||
const params = { query: QUERY, where: { lane: 'alpha' }, limit: 8 }
|
||||
const expected = await legacyHybridFind(brain as any, params)
|
||||
const actual = await brain.find(params as any)
|
||||
expect(actual.length).toBe(expected.length)
|
||||
expect(project(actual)).toEqual(expected)
|
||||
})
|
||||
|
||||
it('hybrid + where + offset: identical page two', async () => {
|
||||
const params = { query: QUERY, where: { lane: 'alpha' }, limit: 6, offset: 6 }
|
||||
const expected = await legacyHybridFind(brain as any, params)
|
||||
const actual = await brain.find(params as any)
|
||||
expect(actual.length).toBe(expected.length)
|
||||
expect(project(actual)).toEqual(expected)
|
||||
})
|
||||
|
||||
it('hybrid + type list + excludeVFS + a `missing` negation: identical', async () => {
|
||||
const params = {
|
||||
query: QUERY,
|
||||
type: [NounType.Document, NounType.Person],
|
||||
excludeVFS: true,
|
||||
where: { lane: 'alpha', retracted: { missing: true } },
|
||||
limit: 8
|
||||
}
|
||||
const expected = await legacyHybridFind(brain as any, params)
|
||||
const actual = await brain.find(params as any)
|
||||
expect(actual.length).toBe(expected.length)
|
||||
expect(project(actual)).toEqual(expected)
|
||||
for (const r of actual) {
|
||||
expect(r.metadata.retracted).toBeUndefined()
|
||||
expect(r.metadata.vfsType).toBeUndefined()
|
||||
}
|
||||
})
|
||||
|
||||
it('hybrid + type list + excludeVFS + a `missing` negation, offset: identical', async () => {
|
||||
const params = {
|
||||
query: QUERY,
|
||||
type: [NounType.Document, NounType.Person],
|
||||
excludeVFS: true,
|
||||
where: { lane: 'alpha', retracted: { missing: true } },
|
||||
limit: 5,
|
||||
offset: 5
|
||||
}
|
||||
const expected = await legacyHybridFind(brain as any, params)
|
||||
const actual = await brain.find(params as any)
|
||||
expect(actual.length).toBe(expected.length)
|
||||
expect(project(actual)).toEqual(expected)
|
||||
})
|
||||
|
||||
it('hybrid + connected: identical, and never a non-neighbour', async () => {
|
||||
const params = {
|
||||
query: QUERY,
|
||||
connected: { from: anchor, direction: 'out' as const },
|
||||
where: { lane: 'alpha' },
|
||||
limit: 8
|
||||
}
|
||||
const expected = await legacyHybridFind(brain as any, params)
|
||||
const actual = await brain.find(params as any)
|
||||
expect(actual.length).toBe(expected.length)
|
||||
expect(project(actual)).toEqual(expected)
|
||||
const neighbours = new Set(matchIds)
|
||||
for (const r of actual) expect(neighbours.has(r.id)).toBe(true)
|
||||
})
|
||||
|
||||
it('hybrid + connected + offset: page two is the page, not an empty answer', async () => {
|
||||
const params = {
|
||||
query: QUERY,
|
||||
connected: { from: anchor, direction: 'out' as const },
|
||||
where: { lane: 'alpha' },
|
||||
limit: 5,
|
||||
offset: 5
|
||||
}
|
||||
const expected = await legacyHybridFind(brain as any, params)
|
||||
expect(expected).toHaveLength(5)
|
||||
const actual = await brain.find(params as any)
|
||||
expect(actual.length).toBe(expected.length)
|
||||
expect(project(actual)).toEqual(expected)
|
||||
})
|
||||
|
||||
it('hybrid + connected: paging reaches every matching neighbour exactly once', async () => {
|
||||
const seen = new Set<string>()
|
||||
for (let offset = 0; offset < MATCHES; offset += 6) {
|
||||
const page = await brain.find({
|
||||
query: QUERY,
|
||||
connected: { from: anchor, direction: 'out' as const },
|
||||
where: { lane: 'alpha' },
|
||||
limit: 6,
|
||||
offset
|
||||
} as any)
|
||||
for (const r of page) {
|
||||
expect(seen.has(r.id)).toBe(false)
|
||||
seen.add(r.id)
|
||||
}
|
||||
}
|
||||
// Every row the fused candidate set holds is reachable by paging, and the
|
||||
// neighbour set is the ceiling.
|
||||
expect(seen.size).toBeGreaterThanOrEqual(MATCHES)
|
||||
const neighbours = new Set(matchIds)
|
||||
for (const id of seen) expect(neighbours.has(id)).toBe(true)
|
||||
})
|
||||
|
||||
it('hybrid + fusion + offset: page two is the page', async () => {
|
||||
const plain = await brain.find({
|
||||
query: QUERY,
|
||||
where: { lane: 'alpha' },
|
||||
limit: 5,
|
||||
offset: 5
|
||||
} as any)
|
||||
const fused = await brain.find({
|
||||
query: QUERY,
|
||||
where: { lane: 'alpha' },
|
||||
fusion: 'weighted',
|
||||
limit: 5,
|
||||
offset: 5
|
||||
} as any)
|
||||
expect(fused).toHaveLength(plain.length)
|
||||
expect(fused.map((r) => r.id)).toEqual(plain.map((r) => r.id))
|
||||
})
|
||||
|
||||
it('a hydrated hybrid row is shaped exactly as an eagerly-built one', async () => {
|
||||
const rows = await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 8 } as any)
|
||||
const row = rows[0]
|
||||
expect(Object.keys(row)).toEqual([
|
||||
'id',
|
||||
'score',
|
||||
'type',
|
||||
'subtype',
|
||||
'visibility',
|
||||
'metadata',
|
||||
'data',
|
||||
'confidence',
|
||||
'weight',
|
||||
'_rev',
|
||||
'entity',
|
||||
'textMatches',
|
||||
'textScore',
|
||||
'semanticScore',
|
||||
'matchSource'
|
||||
])
|
||||
// The flattened fields are projections of the entity, as always.
|
||||
expect(row.entity).toBeDefined()
|
||||
expect(row.type).toBe(row.entity.type)
|
||||
expect(row.metadata).toBe(row.entity.metadata)
|
||||
expect(row.data).toBe(row.entity.data)
|
||||
expect(row._rev).toBe(row.entity._rev)
|
||||
// The match visibility survives the deferral — every leg's fields, on the
|
||||
// rows that leg contributed, exactly as the eager pipeline set them.
|
||||
expect(['text', 'semantic', 'both']).toContain(row.matchSource)
|
||||
for (const r of rows) {
|
||||
if (r.matchSource === 'semantic') {
|
||||
expect(r.textMatches).toEqual([])
|
||||
expect(r.textScore).toBeUndefined()
|
||||
} else {
|
||||
expect(r.textMatches).toEqual(['orbital', 'telemetry'])
|
||||
expect(typeof r.textScore).toBe('number')
|
||||
}
|
||||
if (r.matchSource === 'text') {
|
||||
expect(r.semanticScore).toBeUndefined()
|
||||
} else {
|
||||
expect(typeof r.semanticScore).toBe('number')
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
it('reads canonical for the page only — one batch, `limit` rows', async () => {
|
||||
// Warm any first-read verification before the counters are read.
|
||||
await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 1 } as any)
|
||||
|
||||
const hydrate = vi.spyOn(brain as any, 'batchGet')
|
||||
try {
|
||||
const results = await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 10 } as any)
|
||||
expect(results).toHaveLength(10)
|
||||
expect(hydrate).toHaveBeenCalledTimes(1)
|
||||
expect((hydrate.mock.calls[0][0] as string[]).length).toBe(10)
|
||||
} finally {
|
||||
hydrate.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('asks the text index about the universe only, never the whole store', async () => {
|
||||
const index = (brain as any).metadataIndex
|
||||
const wholeStore = vi.spyOn(index, 'getIdsForTextQuery')
|
||||
const within = vi.spyOn(index, 'getIdsForTextQueryWithin')
|
||||
try {
|
||||
await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 10 } as any)
|
||||
expect(wholeStore).not.toHaveBeenCalled()
|
||||
expect(within).toHaveBeenCalledTimes(1)
|
||||
|
||||
const askedIds = within.mock.calls[0][1] as string[]
|
||||
const universe = await (brain as any).filterIdsBelted({ lane: 'alpha' })
|
||||
expect(askedIds).toHaveLength(universe.length)
|
||||
|
||||
// What the text leg marshals is bounded by the universe, not the store.
|
||||
const marshalled = (await within.mock.results[0].value) as unknown[]
|
||||
expect(marshalled.length).toBeLessThanOrEqual(universe.length)
|
||||
expect(marshalled).toHaveLength(MATCHES)
|
||||
} finally {
|
||||
wholeStore.mockRestore()
|
||||
within.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('the two text doors agree: within is the whole-store answer restricted', async () => {
|
||||
const index = (brain as any).metadataIndex
|
||||
const universe: string[] = await (brain as any).filterIdsBelted({
|
||||
lane: 'alpha',
|
||||
retracted: { missing: true }
|
||||
})
|
||||
const inUniverse = new Set(universe)
|
||||
const whole = await index.getIdsForTextQuery(QUERY)
|
||||
const within = await index.getIdsForTextQueryWithin(QUERY, universe)
|
||||
expect(within).toEqual(whole.filter((m: any) => inUniverse.has(m.id)))
|
||||
expect(await index.getIdsForTextQueryWithin(QUERY, [])).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
/**
|
||||
* FIXTURE B — the query's words are common OUTSIDE the universe, so the old
|
||||
* order's text leg was entirely consumed by rows the filter then discarded.
|
||||
* This is the corrected behaviour, held by name.
|
||||
*/
|
||||
describe('hybrid find: the text leg ranks inside the filter, not around it', () => {
|
||||
let brain: Brainy<any>
|
||||
const QUERY = 'orbital telemetry drift'
|
||||
const NOISE = 150
|
||||
const KEEP = 15
|
||||
const keepIds: string[] = []
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
|
||||
let seed = 5000
|
||||
// Added FIRST and matching one more query word, so they lead the
|
||||
// store-wide text ranking outright — and none of them pass the filter.
|
||||
for (let i = 0; i < NOISE; i++) {
|
||||
await brain.add({
|
||||
id: `noise-${i}`,
|
||||
data: `orbital telemetry drift report ${i}`,
|
||||
type: NounType.Document,
|
||||
metadata: { lane: 'beta' },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
}
|
||||
for (let i = 0; i < KEEP; i++) {
|
||||
const id = `keep-${i}`
|
||||
await brain.add({
|
||||
id,
|
||||
data: `orbital telemetry summary ${i}`,
|
||||
type: NounType.Document,
|
||||
metadata: { lane: 'alpha' },
|
||||
vector: seededVector(seed++)
|
||||
})
|
||||
keepIds.push(resolveEntityId(id))
|
||||
}
|
||||
})
|
||||
|
||||
it('the old order let the filter consume the whole text leg', async () => {
|
||||
const index = (brain as any).metadataIndex
|
||||
const universe: string[] = await (brain as any).filterIdsBelted({ lane: 'alpha' })
|
||||
expect(universe).toHaveLength(KEEP)
|
||||
const inUniverse = new Set(universe)
|
||||
|
||||
// The store-wide prefix the old text leg took (limit 10 → limit * 4).
|
||||
const prefix = (await index.getIdsForTextQuery(QUERY)).slice(0, 40)
|
||||
expect(prefix).toHaveLength(40)
|
||||
expect(prefix.filter((m: any) => inUniverse.has(m.id))).toHaveLength(0)
|
||||
|
||||
// Every row the old text leg ranked was then discarded by the filter, so
|
||||
// the old answer carried NO text contribution at all — fifteen rows that
|
||||
// match the query's words exactly, and not one of them reached the page
|
||||
// through the text leg. What the old order returned was whatever the
|
||||
// semantic leg alone happened to reach.
|
||||
const legacy = await legacyHybridFind(brain as any, {
|
||||
query: QUERY,
|
||||
where: { lane: 'alpha' },
|
||||
limit: 10
|
||||
})
|
||||
for (const r of legacy) {
|
||||
expect(r.matchSource).toBe('semantic')
|
||||
expect(r.textScore).toBeUndefined()
|
||||
expect(r.textMatches).toEqual([])
|
||||
}
|
||||
})
|
||||
|
||||
it('the new order ranks the text leg inside the universe', async () => {
|
||||
const results = await brain.find({
|
||||
query: QUERY,
|
||||
where: { lane: 'alpha' },
|
||||
limit: 10
|
||||
} as any)
|
||||
|
||||
expect(results).toHaveLength(10)
|
||||
const keeps = new Set(keepIds)
|
||||
for (const r of results) {
|
||||
expect(keeps.has(r.id)).toBe(true)
|
||||
expect(r.metadata.lane).toBe('alpha')
|
||||
// The text leg is the contributor the old order threw away.
|
||||
expect(['text', 'both']).toContain(r.matchSource)
|
||||
expect(r.textScore).toBe(1)
|
||||
expect(r.textMatches).toEqual(['orbital', 'telemetry'])
|
||||
}
|
||||
})
|
||||
|
||||
it('paging reaches every matching row the old order could not see', async () => {
|
||||
const seen = new Set<string>()
|
||||
for (let offset = 0; offset < KEEP; offset += 5) {
|
||||
const page = await brain.find({
|
||||
query: QUERY,
|
||||
where: { lane: 'alpha' },
|
||||
limit: 5,
|
||||
offset
|
||||
} as any)
|
||||
expect(page).toHaveLength(5)
|
||||
for (const r of page) {
|
||||
expect(seen.has(r.id)).toBe(false)
|
||||
seen.add(r.id)
|
||||
}
|
||||
}
|
||||
expect(seen.size).toBe(KEEP)
|
||||
expect([...seen].sort()).toEqual([...keepIds].sort())
|
||||
})
|
||||
|
||||
it('reads canonical for the page only, on the truncating shape too', async () => {
|
||||
await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 1 } as any)
|
||||
|
||||
const hydrate = vi.spyOn(brain as any, 'batchGet')
|
||||
try {
|
||||
const results = await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 10 } as any)
|
||||
expect(results).toHaveLength(10)
|
||||
expect(hydrate).toHaveBeenCalledTimes(1)
|
||||
expect((hydrate.mock.calls[0][0] as string[]).length).toBe(10)
|
||||
} finally {
|
||||
hydrate.mockRestore()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
|
@ -1,48 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/find-near
|
||||
* @description find({ near }) searches around the anchor's OWN vector (10.4.10).
|
||||
*
|
||||
* The proximity search fetched its anchor without vectors and fed a
|
||||
* zero-length vector to the index — every near() refused with a dimension
|
||||
* mismatch, for every caller. Found by the Rust planner's first-contact pins
|
||||
* (the planner declines `near`; the pin compared outcomes with and without
|
||||
* it). Now the anchor is fetched with its vector, and an anchor without one
|
||||
* refuses by name instead of failing inside the index.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType } from '../../src/types/graphTypes'
|
||||
import { v5 } from '../../src/universal/uuid'
|
||||
import { generateTestVector } from '../helpers/test-factory'
|
||||
|
||||
describe('find({ near }) uses the anchor vector', () => {
|
||||
let brain: Brainy<any>
|
||||
const anchorVector = generateTestVector()
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
await brain.add({ id: 'anchor', data: 'anchor row', type: NounType.Thing, vector: anchorVector })
|
||||
// A twin with the identical vector and a far row.
|
||||
await brain.add({ id: 'twin', data: 'twin row', type: NounType.Thing, vector: [...anchorVector] })
|
||||
await brain.add({ id: 'far', data: 'far row', type: NounType.Thing, vector: generateTestVector() })
|
||||
})
|
||||
|
||||
it('returns the anchor\'s neighbours by its own vector', async () => {
|
||||
const results = await brain.find({ near: { id: 'anchor' }, limit: 3 })
|
||||
expect(results.length).toBeGreaterThan(0)
|
||||
const ids = results.map((r) => r.entity.id)
|
||||
expect(ids).toContain(v5('twin'))
|
||||
})
|
||||
|
||||
it('refuses by name when the anchor has no vector', async () => {
|
||||
await brain.add({
|
||||
id: 'unvectored',
|
||||
data: 'no vector here',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
;(brain as any).kickEmbedWorker = () => {}
|
||||
await expect(brain.find({ near: { id: 'unvectored' }, limit: 3 })).rejects.toThrow(/has no vector to search around/)
|
||||
})
|
||||
})
|
||||
|
|
@ -1,137 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/find-planner-door
|
||||
* @description The optional `MetadataIndexProvider.planFindPage` door.
|
||||
*
|
||||
* The stage doors each serve one stage, so a `find()` that consults three of
|
||||
* them crosses into the index three times and marshals a result set at every
|
||||
* crossing — a filter matching a hundred thousand rows builds a hundred
|
||||
* thousand id strings to return a page of twenty-five. An index that can decide
|
||||
* the stage order itself answers the page in one call.
|
||||
*
|
||||
* These pins hold the three properties that make such a door safe to add:
|
||||
*
|
||||
* 1. **Absent, nothing changes.** The reference index has no planner, and every
|
||||
* find is served by the stage doors exactly as before. That is also what
|
||||
* makes this engine the ordering oracle for any index that implements one.
|
||||
* 2. **Present, it is asked first and its answer is used** — above the branch
|
||||
* selection, with the params already normalized, the hidden ids passed, and
|
||||
* the graph provider handed over.
|
||||
* 3. **`null` is routing, not an answer.** A door that declines a shape leaves
|
||||
* it to the path that always served it, and the result is unchanged.
|
||||
*
|
||||
* Plus the serving law: an empty page stamped `emptyAt: 'graph'` is re-verified
|
||||
* against the adjacency before it is believed, so a not-serving graph refuses
|
||||
* loudly instead of answering `[]` as truth.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, vi } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType, VerbType } from '../../src/types/graphTypes'
|
||||
import { generateTestVector } from '../helpers/test-factory'
|
||||
|
||||
describe('find(): the optional planner door', () => {
|
||||
let brain: Brainy<any>
|
||||
const anchor = 'planner-anchor'
|
||||
let neighbourId = ''
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
await brain.add({
|
||||
id: anchor,
|
||||
data: 'anchor',
|
||||
type: NounType.Person,
|
||||
metadata: { kind: 'anchor' },
|
||||
vector: generateTestVector()
|
||||
})
|
||||
for (let i = 0; i < 12; i++) {
|
||||
const id = await brain.add({
|
||||
id: `row-${i}`,
|
||||
data: `row ${i}`,
|
||||
type: NounType.Person,
|
||||
metadata: { kind: 'note', rank: i },
|
||||
vector: generateTestVector()
|
||||
})
|
||||
if (i === 0) neighbourId = id
|
||||
await brain.relate({ from: anchor, to: id, type: VerbType.Knows })
|
||||
}
|
||||
})
|
||||
|
||||
/** Install a planner door for one call, then remove it. */
|
||||
const withDoor = async <T>(
|
||||
door: (...a: any[]) => Promise<any>,
|
||||
body: () => Promise<T>
|
||||
): Promise<T> => {
|
||||
const index = (brain as any).metadataIndex
|
||||
index.planFindPage = door
|
||||
try {
|
||||
return await body()
|
||||
} finally {
|
||||
delete index.planFindPage
|
||||
}
|
||||
}
|
||||
|
||||
it('is absent on the reference index — every find is served by the stage doors', async () => {
|
||||
expect((brain as any).metadataIndex.planFindPage).toBeUndefined()
|
||||
const results = await brain.find({ where: { kind: 'note' }, limit: 5 })
|
||||
expect(results).toHaveLength(5)
|
||||
})
|
||||
|
||||
it('is asked before the branches, with normalized params and the graph provider', async () => {
|
||||
const door = vi.fn(async () => null)
|
||||
await withDoor(door, async () => {
|
||||
await brain.find({ where: { kind: 'note' }, limit: 5 })
|
||||
})
|
||||
expect(door).toHaveBeenCalledTimes(1)
|
||||
const [params, hidden, graph] = door.mock.calls[0] as any[]
|
||||
expect(params.where).toEqual({ kind: 'note' })
|
||||
expect(Array.isArray(hidden)).toBe(true)
|
||||
expect(graph).toBe((brain as any).graphIndex)
|
||||
})
|
||||
|
||||
it('uses the page it answers, hydrated and in the door\'s order', async () => {
|
||||
const results = await withDoor(
|
||||
async () => ({ ids: [neighbourId], emptyAt: 'none' as const }),
|
||||
async () => brain.find({ where: { kind: 'note' }, limit: 5 })
|
||||
)
|
||||
expect(results).toHaveLength(1)
|
||||
expect(results[0].entity.id).toBe(neighbourId)
|
||||
})
|
||||
|
||||
it('a declining door changes nothing — the shape is served as it always was', async () => {
|
||||
const withoutDoor = await brain.find({ where: { kind: 'note' }, orderBy: 'rank', limit: 4 })
|
||||
const declined = await withDoor(
|
||||
async () => null,
|
||||
async () => brain.find({ where: { kind: 'note' }, orderBy: 'rank', limit: 4 })
|
||||
)
|
||||
expect(declined.map((r) => r.entity.id)).toEqual(withoutDoor.map((r) => r.entity.id))
|
||||
})
|
||||
|
||||
it('re-verifies the adjacency before believing an empty graph answer', async () => {
|
||||
const verify = vi.spyOn(brain as any, 'verifyGraphAdjacencyLive')
|
||||
try {
|
||||
const results = await withDoor(
|
||||
async () => ({ ids: [], emptyAt: 'graph' as const }),
|
||||
async () => brain.find({ connected: { from: anchor }, where: { kind: 'note' }, limit: 5 })
|
||||
)
|
||||
expect(results).toEqual([])
|
||||
expect(verify).toHaveBeenCalled()
|
||||
} finally {
|
||||
verify.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('does not re-verify the adjacency for an empty the FILTER produced', async () => {
|
||||
const verify = vi.spyOn(brain as any, 'verifyGraphAdjacencyLive')
|
||||
verify.mockClear()
|
||||
try {
|
||||
const results = await withDoor(
|
||||
async () => ({ ids: [], emptyAt: 'filter' as const }),
|
||||
async () => brain.find({ where: { kind: 'note' }, limit: 5 })
|
||||
)
|
||||
expect(results).toEqual([])
|
||||
expect(verify).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
verify.mockRestore()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
|
@ -1,101 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/generation-store-factory
|
||||
* @description Pins the `createGenerationStore` protected factory hook on
|
||||
* `Brainy` ({@link Brainy.createGenerationStore}). The hook exists so an
|
||||
* engine built on top of this reference implementation can substitute a
|
||||
* `GenerationStore` that keeps the same behavioural contract; this suite
|
||||
* proves two things:
|
||||
*
|
||||
* 1. A subclass overriding the hook is the ONLY path that constructs the
|
||||
* generation store — it is called exactly once, with the same storage
|
||||
* instance `performInit` holds — and the store the brain actually uses
|
||||
* is the one the override returned.
|
||||
* 2. The default (non-overridden) path is unaffected — proven here by
|
||||
* confirming the base class still produces a plain `GenerationStore`
|
||||
* wired to `brain.storage`, and separately by running the existing
|
||||
* `db-mvcc` and `brainy-core.integration` suites unmodified against this
|
||||
* change (they exercise generation-store behaviour end to end).
|
||||
*/
|
||||
|
||||
import { describe, it, expect, afterEach } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy.js'
|
||||
import { GenerationStore } from '../../src/db/generationStore.js'
|
||||
import type { BaseStorage } from '../../src/storage/baseStorage.js'
|
||||
|
||||
/** Typed access to the brain's private storage + generation-store fields (test injection point). */
|
||||
function internalsOf(brain: Brainy): { storage: BaseStorage; generationStore: GenerationStore } {
|
||||
return brain as unknown as { storage: BaseStorage; generationStore: GenerationStore }
|
||||
}
|
||||
|
||||
/**
|
||||
* A `GenerationStore` subclass that counts its own construction and
|
||||
* remembers the storage instance it was built with, so the test can prove
|
||||
* the hook is the sole construction path without mocking the module.
|
||||
*/
|
||||
class SpyGenerationStore extends GenerationStore {
|
||||
static constructCount = 0
|
||||
static lastStorage: BaseStorage | undefined
|
||||
|
||||
constructor(storage: BaseStorage) {
|
||||
super(storage)
|
||||
SpyGenerationStore.constructCount++
|
||||
SpyGenerationStore.lastStorage = storage
|
||||
}
|
||||
}
|
||||
|
||||
/** A Brainy subclass overriding the factory hook — stands in for an engine built on the reference. */
|
||||
class BrainyWithSpyStore extends Brainy {
|
||||
hookCallCount = 0
|
||||
hookStorageArg: BaseStorage | undefined
|
||||
|
||||
protected override createGenerationStore(storage: BaseStorage): GenerationStore {
|
||||
this.hookCallCount++
|
||||
this.hookStorageArg = storage
|
||||
return new SpyGenerationStore(storage)
|
||||
}
|
||||
}
|
||||
|
||||
describe('Brainy.createGenerationStore — protected factory hook', () => {
|
||||
const brains: Brainy[] = []
|
||||
|
||||
afterEach(async () => {
|
||||
SpyGenerationStore.constructCount = 0
|
||||
SpyGenerationStore.lastStorage = undefined
|
||||
for (const brain of brains.splice(0)) {
|
||||
try {
|
||||
await brain.close()
|
||||
} catch {
|
||||
// already closed by the test
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
it('a subclass override is the sole construction path: called once, same storage instance, its store is the one the brain uses', async () => {
|
||||
const brain = new BrainyWithSpyStore({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
brains.push(brain)
|
||||
|
||||
// Called exactly once, through the hook.
|
||||
expect(brain.hookCallCount).toBe(1)
|
||||
expect(SpyGenerationStore.constructCount).toBe(1)
|
||||
|
||||
// Same storage instance the base class holds — not a copy, not a different adapter.
|
||||
const { storage, generationStore } = internalsOf(brain)
|
||||
expect(brain.hookStorageArg).toBe(storage)
|
||||
expect(SpyGenerationStore.lastStorage).toBe(storage)
|
||||
|
||||
// The store the brain actually uses is the one the override returned.
|
||||
expect(generationStore).toBeInstanceOf(SpyGenerationStore)
|
||||
})
|
||||
|
||||
it('the default (non-overridden) path still produces a plain GenerationStore wired to the same storage', async () => {
|
||||
const brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
brains.push(brain)
|
||||
|
||||
const { storage, generationStore } = internalsOf(brain)
|
||||
expect(generationStore).toBeInstanceOf(GenerationStore)
|
||||
// The default implementation constructs from the same storage the brain holds.
|
||||
expect((generationStore as unknown as { storage: BaseStorage }).storage).toBe(storage)
|
||||
})
|
||||
})
|
||||
|
|
@ -1,547 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/pending-embed-checkpoint
|
||||
* @description THE PENDING-EMBED CHECKPOINT — the bound that engages on the
|
||||
* brains that need it.
|
||||
*
|
||||
* 10.4.9 bounded the open-path `recover-pending-embeds` fold with a LOW-WATER
|
||||
* MARK: the log head at which the pending set last drained to EMPTY. That mark
|
||||
* carries no set, so it can only be written when the set is empty — and a brain
|
||||
* holding even ONE id that never lands (an embed that keeps failing, a worker
|
||||
* that never gets to it, a row reaped in memory only and re-folded every open)
|
||||
* never drains, therefore never writes a mark, therefore re-reads its WHOLE
|
||||
* fact log on every single open. The bound was absent from exactly the brains
|
||||
* whose fold is expensive: a silent scaling defect.
|
||||
*
|
||||
* The cure is a CHECKPOINT of the pending set —
|
||||
* `_system/pending_embeds_checkpoint.json` = `{ generation, pending, writtenAt }`,
|
||||
* meaning "as of durable generation G the pending set was exactly this list".
|
||||
* Open seeds the set from `pending` and scans only from `G + 1`, so the fold is
|
||||
* O(facts since G) whether or not the set ever drains.
|
||||
*
|
||||
* What this suite pins:
|
||||
* 1. A brain with one permanently-stuck pending id, closed cleanly and
|
||||
* reopened, scans ONLY the facts after the checkpoint — asserted from the
|
||||
* fold's own accounting, never a clock. The same fixture pins the DEFECT:
|
||||
* no low-water mark exists on that brain, because it never drained.
|
||||
* 2. A crash matrix in a REAL child process (SIGKILL, no close), for kills
|
||||
* before a checkpoint write, after one with embeds landed and flushed
|
||||
* after it, and after one with an UN-FLUSHED tail at the moment of death.
|
||||
* The invariant in every row is differential: the checkpoint-bounded fold
|
||||
* the reopened brain actually ran ≡ a full fold from generation 1 over the
|
||||
* same recovered log.
|
||||
* 3. A torn checkpoint falls back — loudly (the adapter's torn-record gauge
|
||||
* plus the fold's own narration of which bound applied) and correctly.
|
||||
* 4. The existing low-water pins keep passing unchanged
|
||||
* (`pending-embed-low-water.test.ts`): the mark is still written and is
|
||||
* still read, now as the FALLBACK bound beneath the checkpoint.
|
||||
*
|
||||
* The crash-recovery contract is untouched: the fold runs on the open's
|
||||
* foreground, so a reopened brain has its markers re-armed when open() returns.
|
||||
*/
|
||||
import { describe, it, expect, afterEach } from 'vitest'
|
||||
import { mkdtempSync, rmSync, existsSync, readFileSync, writeFileSync } from 'node:fs'
|
||||
import { spawn } from 'node:child_process'
|
||||
import { gunzipSync } from 'node:zlib'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { Brainy } from '../../src/brainy.js'
|
||||
import { NounType } from '../../src/types/graphTypes.js'
|
||||
import { getTornRecordGauge } from '../../src/storage/tornRecordError.js'
|
||||
|
||||
const CHECKPOINT_PATH = '_system/pending_embeds_checkpoint.json'
|
||||
const LOWWATER_PATH = '_system/pending_embeds_lowwater.json'
|
||||
const REPO_ROOT = process.cwd()
|
||||
const TSX = join(REPO_ROOT, 'node_modules', '.bin', 'tsx')
|
||||
|
||||
/** The fold's own accounting for the most recent open. */
|
||||
interface FoldReport {
|
||||
bound: 'checkpoint' | 'low-water' | 'genesis'
|
||||
fromGeneration: number
|
||||
factsScanned: number
|
||||
seeded: number
|
||||
pending: number
|
||||
}
|
||||
|
||||
const roots: string[] = []
|
||||
const liveBrains: Brainy<any>[] = []
|
||||
|
||||
function dir(): string {
|
||||
const d = mkdtempSync(join(tmpdir(), 'brainy-embed-ckpt-'))
|
||||
roots.push(d)
|
||||
return d
|
||||
}
|
||||
|
||||
async function open(root: string, opts?: { blockWorker?: boolean }): Promise<Brainy<any>> {
|
||||
const brain = new Brainy<any>({
|
||||
requireSubtype: false,
|
||||
storage: { type: 'filesystem', path: root }
|
||||
})
|
||||
// Blocking the worker BEFORE init() is how a "permanently stuck" pending id
|
||||
// is built deterministically: the state under test is "an id the fold keeps
|
||||
// re-arming and nothing ever disarms", and its production causes (a failing
|
||||
// embedder, a wedged model, a data-less row) all reduce to exactly that.
|
||||
if (opts?.blockWorker) (brain as unknown as { kickEmbedWorker: () => void }).kickEmbedWorker = () => {}
|
||||
await brain.init()
|
||||
liveBrains.push(brain)
|
||||
return brain
|
||||
}
|
||||
|
||||
function foldReport(brain: Brainy<any>): FoldReport {
|
||||
const report = (brain as unknown as { _pendingEmbedFoldReport: FoldReport | null })
|
||||
._pendingEmbedFoldReport
|
||||
if (report === null) throw new Error('the open ran no pending-embed fold')
|
||||
return report
|
||||
}
|
||||
|
||||
function pendingIds(brain: Brainy<any>): string[] {
|
||||
return [
|
||||
...(brain as unknown as { _pendingEmbedIds: Set<string> })._pendingEmbedIds
|
||||
].sort()
|
||||
}
|
||||
|
||||
/** Read an artifact straight off disk (the adapter gzips raw objects). */
|
||||
function readArtifact(root: string, path: string): Record<string, unknown> | null {
|
||||
const plain = join(root, ...path.split('/'))
|
||||
const gz = `${plain}.gz`
|
||||
if (existsSync(gz)) return JSON.parse(gunzipSync(readFileSync(gz)).toString('utf-8'))
|
||||
if (existsSync(plain)) return JSON.parse(readFileSync(plain, 'utf-8'))
|
||||
return null
|
||||
}
|
||||
|
||||
/** The on-disk path the adapter actually used for an artifact. */
|
||||
function artifactPath(root: string, path: string): string | null {
|
||||
const plain = join(root, ...path.split('/'))
|
||||
const gz = `${plain}.gz`
|
||||
if (existsSync(gz)) return gz
|
||||
if (existsSync(plain)) return plain
|
||||
return null
|
||||
}
|
||||
|
||||
/**
|
||||
* THE DIFFERENTIAL ORACLE: fold the log from generation 1 with exactly the
|
||||
* engine's own rules. This is what the bounded fold must agree with, and its
|
||||
* fact count is what the unbounded fold used to read at every open.
|
||||
*/
|
||||
async function fullFold(brain: Brainy<any>): Promise<{ ids: string[]; facts: number }> {
|
||||
const log = (
|
||||
brain as unknown as { generationStore: { getFactLog(): any } }
|
||||
).generationStore.getFactLog()
|
||||
const pending = new Set<string>()
|
||||
let facts = 0
|
||||
const scan = log.scanFacts({ fromGeneration: 1 })
|
||||
for await (const batch of scan.batches()) {
|
||||
for (const fact of batch.facts) {
|
||||
facts++
|
||||
for (const record of fact.records ?? []) {
|
||||
if (record.type === 'embed.pending') pending.add(record.id)
|
||||
else if (record.type === 'embed.landed') pending.delete(record.id)
|
||||
}
|
||||
for (const op of fact.ops) {
|
||||
if (op.kind === 'noun' && op.record === null) pending.delete(op.id)
|
||||
}
|
||||
}
|
||||
}
|
||||
return { ids: [...pending].sort(), facts }
|
||||
}
|
||||
|
||||
/** Capture every console.warn/error line emitted while `fn` runs. */
|
||||
async function captureConsole<T>(fn: () => Promise<T>): Promise<{ result: T; lines: string[] }> {
|
||||
const lines: string[] = []
|
||||
const origWarn = console.warn
|
||||
const origError = console.error
|
||||
const sink = (...args: unknown[]) => {
|
||||
lines.push(args.map((a) => String(a)).join(' '))
|
||||
}
|
||||
console.warn = sink as typeof console.warn
|
||||
console.error = sink as typeof console.error
|
||||
try {
|
||||
const result = await fn()
|
||||
return { result, lines }
|
||||
} finally {
|
||||
console.warn = origWarn
|
||||
console.error = origError
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Run a child process that arranges a store and then waits forever, so the
|
||||
* parent can SIGKILL it. A real process death is the only honest way to pin
|
||||
* "no close ran, no shutdown hook ran, RAM is gone".
|
||||
*
|
||||
* `detached` puts the child in its own process GROUP: tsx runs the script in a
|
||||
* grandchild, and only a group-wide signal reaches the process holding the
|
||||
* writer lock.
|
||||
*/
|
||||
function spawnArranger(root: string, body: string): Promise<{
|
||||
child: ReturnType<typeof spawn>
|
||||
output: () => string
|
||||
}> {
|
||||
const scriptPath = join(root, 'arrange.mts')
|
||||
writeFileSync(scriptPath, body)
|
||||
const child = spawn(TSX, [scriptPath], {
|
||||
cwd: REPO_ROOT,
|
||||
stdio: ['ignore', 'pipe', 'pipe'],
|
||||
detached: true
|
||||
})
|
||||
let out = ''
|
||||
child.stdout!.on('data', (d) => { out += String(d) })
|
||||
child.stderr!.on('data', (d) => { out += String(d) })
|
||||
return new Promise((resolvePromise, rejectPromise) => {
|
||||
const timer = setTimeout(
|
||||
() => rejectPromise(new Error(`arranger never became READY:\n${out}`)),
|
||||
180_000
|
||||
)
|
||||
child.stdout!.on('data', () => {
|
||||
if (out.includes('READY')) {
|
||||
clearTimeout(timer)
|
||||
resolvePromise({ child, output: () => out })
|
||||
}
|
||||
})
|
||||
child.on('exit', (code) => {
|
||||
clearTimeout(timer)
|
||||
if (!out.includes('READY')) rejectPromise(new Error(`arranger exited ${code}:\n${out}`))
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
/** Parse the `IDS:{...}` line an arranger prints — supplied ids are normalised
|
||||
* to canonical uuids, and the markers, checkpoint and fold all speak those. */
|
||||
function childIds(output: string): Record<string, string> {
|
||||
const line = output.split('\n').find((l) => l.startsWith('IDS:'))
|
||||
if (!line) throw new Error(`arranger printed no IDS line:\n${output}`)
|
||||
return JSON.parse(line.slice('IDS:'.length))
|
||||
}
|
||||
|
||||
/** SIGKILL the whole group and wait for the grandchild's death to settle. */
|
||||
async function sigkill(child: ReturnType<typeof spawn>): Promise<void> {
|
||||
process.kill(-(child.pid as number), 'SIGKILL')
|
||||
await new Promise<void>((r) => child.on('exit', () => r()))
|
||||
await new Promise<void>((r) => setTimeout(r, 500))
|
||||
}
|
||||
|
||||
/** The preamble every arranger child shares. */
|
||||
function childPreamble(root: string): string {
|
||||
return `
|
||||
import { Brainy } from ${JSON.stringify(join(REPO_ROOT, 'src', 'brainy.ts'))}
|
||||
const ROOT = ${JSON.stringify(root)}
|
||||
const brain = new Brainy<any>({ requireSubtype: false, storage: { type: 'filesystem', path: ROOT } })
|
||||
const block = () => { (brain as any).kickEmbedWorker = () => {} }
|
||||
const settleCheckpoint = async () => {
|
||||
// The cadence write is fire-and-forget; wait for the single flight.
|
||||
for (let i = 0; i < 200; i++) {
|
||||
if (!(brain as any)._pendingEmbedCheckpointFlight) break
|
||||
await (brain as any)._pendingEmbedCheckpointFlight.catch(() => {})
|
||||
}
|
||||
}
|
||||
`
|
||||
}
|
||||
|
||||
afterEach(async () => {
|
||||
for (const brain of liveBrains.splice(0)) {
|
||||
try { await brain.close() } catch { /* already closed / crashed — teardown only */ }
|
||||
}
|
||||
for (const d of roots.splice(0)) rmSync(d, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
// ===========================================================================
|
||||
// 1. The stuck-id brain — the defect, and the bound that now engages on it
|
||||
// ===========================================================================
|
||||
|
||||
describe('pending-embed checkpoint — a brain whose pending set never drains', () => {
|
||||
it('a permanently-stuck pending id: the reopen scans only the facts after the checkpoint', async () => {
|
||||
const root = dir()
|
||||
const first = await open(root, { blockWorker: true })
|
||||
// add() returns the CANONICAL id (supplied ids are normalised), and that is
|
||||
// the id the markers, the checkpoint and the fold all speak.
|
||||
const stuck = await first.add({
|
||||
id: 'stuck',
|
||||
data: 'a deferred row whose embed never lands',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
expect(first.pendingEmbedCount()).toBe(1)
|
||||
// Ordinary traffic after it — every one of these is a fact the unbounded
|
||||
// fold had to re-read at every open, forever, because of that one id.
|
||||
for (let i = 0; i < 12; i++) {
|
||||
await first.add({ id: `row-${i}`, data: `row ${i}`, type: NounType.Thing })
|
||||
}
|
||||
await first.close()
|
||||
liveBrains.splice(liveBrains.indexOf(first), 1)
|
||||
|
||||
// THE DEFECT, PINNED: the pending set never drained, so the old bound was
|
||||
// never written — nothing on this brain could have shortened its fold.
|
||||
expect(readArtifact(root, LOWWATER_PATH)).toBeNull()
|
||||
// The checkpoint IS written at the clean close, set non-empty and all.
|
||||
const checkpoint = readArtifact(root, CHECKPOINT_PATH) as {
|
||||
generation: number
|
||||
pending: string[]
|
||||
} | null
|
||||
expect(checkpoint).not.toBeNull()
|
||||
expect(checkpoint!.generation).toBeGreaterThan(0)
|
||||
expect(checkpoint!.pending).toEqual([stuck])
|
||||
|
||||
const second = await open(root, { blockWorker: true })
|
||||
const report = foldReport(second)
|
||||
// THE FIX, from the fold's own counter — not the clock.
|
||||
expect(report.bound).toBe('checkpoint')
|
||||
expect(report.fromGeneration).toBe(checkpoint!.generation + 1)
|
||||
expect(report.factsScanned).toBe(0)
|
||||
expect(report.seeded).toBe(1)
|
||||
// The crash-recovery contract is intact: the marker is re-armed by open().
|
||||
expect(pendingIds(second)).toEqual([stuck])
|
||||
expect(second.pendingEmbedCount()).toBe(1)
|
||||
|
||||
// The differential: the bounded answer is the full-fold answer, and the
|
||||
// full fold is what the previous bound would have had to read.
|
||||
const full = await fullFold(second)
|
||||
expect(full.ids).toEqual([stuck])
|
||||
expect(full.facts).toBeGreaterThanOrEqual(13)
|
||||
expect(report.factsScanned).toBeLessThan(full.facts)
|
||||
}, 180_000)
|
||||
|
||||
it('the bound stays O(delta) across repeated opens while the id is still stuck', async () => {
|
||||
const root = dir()
|
||||
const first = await open(root, { blockWorker: true })
|
||||
const stuck = await first.add({
|
||||
id: 'stuck',
|
||||
data: 'never lands',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
for (let i = 0; i < 6; i++) {
|
||||
await first.add({ id: `a-${i}`, data: `a ${i}`, type: NounType.Thing })
|
||||
}
|
||||
await first.close()
|
||||
liveBrains.splice(liveBrains.indexOf(first), 1)
|
||||
|
||||
const second = await open(root, { blockWorker: true })
|
||||
expect(foldReport(second).factsScanned).toBe(0)
|
||||
// More history under the same stuck id.
|
||||
for (let i = 0; i < 9; i++) {
|
||||
await second.add({ id: `b-${i}`, data: `b ${i}`, type: NounType.Thing })
|
||||
}
|
||||
await second.close()
|
||||
liveBrains.splice(liveBrains.indexOf(second), 1)
|
||||
|
||||
const third = await open(root, { blockWorker: true })
|
||||
const report = foldReport(third)
|
||||
const full = await fullFold(third)
|
||||
expect(report.bound).toBe('checkpoint')
|
||||
expect(report.factsScanned).toBe(0)
|
||||
// The unbounded fold grew with the store; the bounded one did not.
|
||||
expect(full.facts).toBeGreaterThanOrEqual(16)
|
||||
expect(pendingIds(third)).toEqual([stuck])
|
||||
expect(full.ids).toEqual([stuck])
|
||||
}, 180_000)
|
||||
})
|
||||
|
||||
// ===========================================================================
|
||||
// 2. Torn checkpoint — falls back, loudly, correctly
|
||||
// ===========================================================================
|
||||
|
||||
describe('pending-embed checkpoint — a torn checkpoint never shortens the fold', () => {
|
||||
it('an undecodable checkpoint file degrades to the next bound, loudly, with the right pending set', async () => {
|
||||
const root = dir()
|
||||
const first = await open(root, { blockWorker: true })
|
||||
const stuck = await first.add({
|
||||
id: 'stuck',
|
||||
data: 'never lands',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
for (let i = 0; i < 5; i++) {
|
||||
await first.add({ id: `row-${i}`, data: `row ${i}`, type: NounType.Thing })
|
||||
}
|
||||
await first.close()
|
||||
liveBrains.splice(liveBrains.indexOf(first), 1)
|
||||
|
||||
const onDisk = artifactPath(root, CHECKPOINT_PATH)
|
||||
expect(onDisk).not.toBeNull()
|
||||
// Tear it: bytes that are neither valid gzip nor valid JSON. A torn file
|
||||
// must THROW on read — never parse into a partial `pending` list.
|
||||
writeFileSync(onDisk!, 'not a checkpoint at all {{{')
|
||||
|
||||
const before = getTornRecordGauge().count
|
||||
const { result: second, lines } = await captureConsole(async () =>
|
||||
open(root, { blockWorker: true })
|
||||
)
|
||||
const report = foldReport(second)
|
||||
// Fell back — never to a shorter bound, and never silently.
|
||||
expect(report.bound).not.toBe('checkpoint')
|
||||
expect(report.seeded).toBe(0)
|
||||
expect(report.fromGeneration).toBe(1) // no mark either: this brain never drained
|
||||
// LOUD, two ways: the adapter's torn-record gauge and its production error…
|
||||
expect(getTornRecordGauge().count).toBeGreaterThan(before)
|
||||
expect(getTornRecordGauge().lastPath).toContain('pending_embeds_checkpoint')
|
||||
expect(lines.some((l) => /TORN RECORD/.test(l))).toBe(true)
|
||||
// …and the fold's own narration of which bound it actually used.
|
||||
expect(lines.some((l) => /pending-embed fold: genesis bound/.test(l))).toBe(true)
|
||||
|
||||
// CORRECT: the marker is still recovered, from the log itself.
|
||||
expect(pendingIds(second)).toEqual([stuck])
|
||||
const full = await fullFold(second)
|
||||
expect(full.ids).toEqual([stuck])
|
||||
expect(report.factsScanned).toBe(full.facts)
|
||||
}, 180_000)
|
||||
|
||||
it('a well-formed but shape-invalid checkpoint is refused whole, never partially trusted', async () => {
|
||||
const root = dir()
|
||||
const first = await open(root, { blockWorker: true })
|
||||
const stuck = await first.add({
|
||||
id: 'stuck',
|
||||
data: 'never lands',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
await first.add({ id: 'other', data: 'ordinary row', type: NounType.Thing })
|
||||
await first.close()
|
||||
liveBrains.splice(liveBrains.indexOf(first), 1)
|
||||
|
||||
// A checkpoint with a plausible generation but a `pending` that is not a
|
||||
// list of ids: trusting the generation alone would bound the scan behind a
|
||||
// set that was never recovered — the exact shape that loses a vector.
|
||||
const onDisk = artifactPath(root, CHECKPOINT_PATH)!
|
||||
const good = readArtifact(root, CHECKPOINT_PATH) as { generation: number }
|
||||
rmSync(onDisk)
|
||||
writeFileSync(
|
||||
join(root, '_system', 'pending_embeds_checkpoint.json'),
|
||||
JSON.stringify({ generation: good.generation, pending: { stuck: true }, writtenAt: 1 })
|
||||
)
|
||||
|
||||
const { result: second, lines } = await captureConsole(async () =>
|
||||
open(root, { blockWorker: true })
|
||||
)
|
||||
expect(lines.some((l) => /pending-embed checkpoint REFUSED/.test(l))).toBe(true)
|
||||
const report = foldReport(second)
|
||||
expect(report.bound).not.toBe('checkpoint')
|
||||
expect(report.seeded).toBe(0)
|
||||
expect(pendingIds(second)).toEqual([stuck])
|
||||
}, 180_000)
|
||||
})
|
||||
|
||||
// ===========================================================================
|
||||
// 3. The crash matrix — real processes, real SIGKILL, differential invariant
|
||||
// ===========================================================================
|
||||
|
||||
describe('pending-embed checkpoint — crash matrix (real child process, SIGKILL)', () => {
|
||||
/**
|
||||
* The invariant every row shares: whatever the reopened brain's fold did with
|
||||
* whatever bound survived the crash, its pending set must equal the truth a
|
||||
* full fold from generation 1 derives from the SAME recovered log.
|
||||
*/
|
||||
async function assertDifferentialAfterCrash(root: string): Promise<{
|
||||
report: FoldReport
|
||||
full: { ids: string[]; facts: number }
|
||||
pending: string[]
|
||||
}> {
|
||||
const reopened = await open(root, { blockWorker: true })
|
||||
const report = foldReport(reopened)
|
||||
const full = await fullFold(reopened)
|
||||
const pending = pendingIds(reopened)
|
||||
expect(pending).toEqual(full.ids)
|
||||
return { report, full, pending }
|
||||
}
|
||||
|
||||
it('killed BEFORE any checkpoint was written — falls back and recovers the marker from the log', async () => {
|
||||
const root = dir()
|
||||
const { child, output } = await spawnArranger(
|
||||
root,
|
||||
`${childPreamble(root)}
|
||||
block()
|
||||
await brain.init()
|
||||
await brain.add({ id: 'landed-row', data: 'an ordinary row', type: 'thing' })
|
||||
const stuck = await brain.add({ id: 'stuck-1', data: 'deferred, never lands', type: 'thing', deferEmbedding: true })
|
||||
await brain.flush()
|
||||
console.log('IDS:' + JSON.stringify({ stuck }))
|
||||
console.log('READY')
|
||||
setInterval(() => {}, 1000)
|
||||
`
|
||||
)
|
||||
const ids = childIds(output())
|
||||
// One enqueue is well under the cadence and the set never drained, so no
|
||||
// checkpoint exists — this is the pre-checkpoint crash.
|
||||
expect(readArtifact(root, CHECKPOINT_PATH)).toBeNull()
|
||||
await sigkill(child)
|
||||
|
||||
const { report, pending } = await assertDifferentialAfterCrash(root)
|
||||
expect(report.bound).toBe('genesis')
|
||||
expect(pending).toEqual([ids.stuck])
|
||||
}, 300_000)
|
||||
|
||||
it('killed AFTER a checkpoint, with an embed landed and flushed after it — the post-checkpoint facts carry the disarm', async () => {
|
||||
const root = dir()
|
||||
const { child, output } = await spawnArranger(
|
||||
root,
|
||||
`${childPreamble(root)}
|
||||
await brain.init()
|
||||
// Land one deferred embed: the drain arms the checkpoint debt.
|
||||
await brain.add({ id: 'seed', data: 'lands first', type: 'thing', deferEmbedding: true })
|
||||
await brain.awaitPendingEmbeds()
|
||||
await brain.flush()
|
||||
// A second deferred write pays the debt (the head is at the manifest now),
|
||||
// then LANDS — its embed.landed rides a fact ABOVE the checkpoint.
|
||||
const landsAfter = await brain.add({ id: 'lands-after', data: 'lands after the checkpoint', type: 'thing', deferEmbedding: true })
|
||||
await settleCheckpoint()
|
||||
await brain.awaitPendingEmbeds()
|
||||
// …and one that never will.
|
||||
block()
|
||||
const stuck = await brain.add({ id: 'stuck-1', data: 'deferred, never lands', type: 'thing', deferEmbedding: true })
|
||||
await brain.add({ id: 'plain', data: 'more history', type: 'thing' })
|
||||
await brain.flush()
|
||||
console.log('IDS:' + JSON.stringify({ stuck, landsAfter }))
|
||||
console.log('READY')
|
||||
setInterval(() => {}, 1000)
|
||||
`
|
||||
)
|
||||
const ids = childIds(output())
|
||||
const checkpoint = readArtifact(root, CHECKPOINT_PATH) as {
|
||||
generation: number
|
||||
pending: string[]
|
||||
} | null
|
||||
expect(checkpoint).not.toBeNull()
|
||||
await sigkill(child)
|
||||
|
||||
const { report, full, pending } = await assertDifferentialAfterCrash(root)
|
||||
expect(report.bound).toBe('checkpoint')
|
||||
expect(report.fromGeneration).toBe(checkpoint!.generation + 1)
|
||||
// The bound really bounded: fewer facts than the whole log.
|
||||
expect(report.factsScanned).toBeLessThan(full.facts)
|
||||
// A landed embed above the checkpoint is disarmed by the scan, not lost;
|
||||
// the stuck one is re-armed.
|
||||
expect(pending).toEqual([ids.stuck])
|
||||
expect(pending).not.toContain(ids.landsAfter)
|
||||
}, 300_000)
|
||||
|
||||
it('killed AFTER a checkpoint with an UN-FLUSHED tail — truncated facts and the bounded fold still agree', async () => {
|
||||
const root = dir()
|
||||
const { child } = await spawnArranger(
|
||||
root,
|
||||
`${childPreamble(root)}
|
||||
await brain.init()
|
||||
await brain.add({ id: 'seed', data: 'lands first', type: 'thing', deferEmbedding: true })
|
||||
await brain.awaitPendingEmbeds()
|
||||
await brain.flush()
|
||||
await brain.add({ id: 'lands-after', data: 'lands after the checkpoint', type: 'thing', deferEmbedding: true })
|
||||
await settleCheckpoint()
|
||||
await brain.awaitPendingEmbeds()
|
||||
await brain.flush()
|
||||
// Now write PAST the manifest and never flush: these facts are the tail a
|
||||
// crash truncates. Whatever survives, the two folds must agree on it.
|
||||
block()
|
||||
await brain.add({ id: 'stuck-tail', data: 'deferred, never lands', type: 'thing', deferEmbedding: true })
|
||||
await brain.add({ id: 'plain-tail', data: 'unflushed history', type: 'thing' })
|
||||
console.log('READY')
|
||||
setInterval(() => {}, 1000)
|
||||
`
|
||||
)
|
||||
const checkpoint = readArtifact(root, CHECKPOINT_PATH) as { generation: number } | null
|
||||
expect(checkpoint).not.toBeNull()
|
||||
await sigkill(child)
|
||||
|
||||
const { report } = await assertDifferentialAfterCrash(root)
|
||||
// The checkpoint's generation is at or below the manifest by construction,
|
||||
// so it survived the truncation and still bounds the fold.
|
||||
expect(report.bound).toBe('checkpoint')
|
||||
expect(report.fromGeneration).toBe(checkpoint!.generation + 1)
|
||||
}, 300_000)
|
||||
})
|
||||
|
|
@ -1,141 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/pending-embed-low-water
|
||||
* @description The pending-embed recovery fold is bounded and background (10.4.9).
|
||||
*
|
||||
* The fold used to scan the generation log from generation 1 at EVERY open,
|
||||
* on the open's foreground — O(whole history) per open on long-lived brains.
|
||||
* Now: an advisory low-water mark (`_system/pending_embeds_lowwater.json`)
|
||||
* records the committed generation whenever the pending set drains to empty,
|
||||
* recovery scans from `mark + 1` on the open's foreground — the crash-recovery
|
||||
* contract keeps markers re-armed when open() returns. The mark is advisory: stale-low costs a longer scan, never a
|
||||
* marker — a pending embed enqueued before a crash is still recovered.
|
||||
*/
|
||||
import { describe, it, expect, afterEach, vi } from 'vitest'
|
||||
import { mkdtempSync, rmSync } from 'node:fs'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType } from '../../src/types/graphTypes'
|
||||
|
||||
const LOWWATER_PATH = '_system/pending_embeds_lowwater.json'
|
||||
|
||||
describe('pending-embed recovery: bounded by the low-water mark', () => {
|
||||
const roots: string[] = []
|
||||
const dir = (): string => {
|
||||
const d = mkdtempSync(join(tmpdir(), 'brainy-lowwater-'))
|
||||
roots.push(d)
|
||||
return d
|
||||
}
|
||||
const open = async (root: string): Promise<Brainy<any>> => {
|
||||
const brain = new Brainy<any>({
|
||||
requireSubtype: false,
|
||||
storage: { type: 'filesystem', path: root }
|
||||
})
|
||||
await brain.init()
|
||||
return brain
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
for (const d of roots.splice(0)) rmSync(d, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
it('drain-to-empty writes the mark, and the next open scans from mark + 1', async () => {
|
||||
const root = dir()
|
||||
const brain = await open(root)
|
||||
// Hold the worker so the pending state is observable, then release it.
|
||||
const realKick = (brain as any).kickEmbedWorker.bind(brain)
|
||||
;(brain as any).kickEmbedWorker = () => {}
|
||||
await brain.add({
|
||||
id: 'row-1',
|
||||
data: 'the first deferred row',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
expect(brain.pendingEmbedCount()).toBeGreaterThan(0)
|
||||
;(brain as any).kickEmbedWorker = realKick
|
||||
await brain.awaitPendingEmbeds()
|
||||
// The drain wrote the advisory mark (fire-and-forget: settle the microtask).
|
||||
await new Promise((r) => setTimeout(r, 50))
|
||||
const mark = (await (brain as any).storage.readRawObject(LOWWATER_PATH)) as {
|
||||
generation: number
|
||||
} | null
|
||||
expect(mark).not.toBeNull()
|
||||
expect(mark!.generation).toBeGreaterThan(0)
|
||||
await brain.close()
|
||||
|
||||
const brain2 = await open(root)
|
||||
const log = (brain2 as any).generationStore.getFactLog()
|
||||
const scanSpy = vi.spyOn(log, 'scanFacts')
|
||||
try {
|
||||
await (brain2 as any).recoverPendingEmbedsFromLog()
|
||||
expect(scanSpy).toHaveBeenCalledTimes(1)
|
||||
const opts = scanSpy.mock.calls[0][0] as { fromGeneration?: number }
|
||||
expect(opts.fromGeneration).toBeGreaterThanOrEqual(mark!.generation + 1)
|
||||
} finally {
|
||||
scanSpy.mockRestore()
|
||||
await brain2.close()
|
||||
}
|
||||
})
|
||||
|
||||
it('a pending embed enqueued after the mark survives an unclean stop', async () => {
|
||||
const root = dir()
|
||||
const brain = await open(root)
|
||||
await brain.add({ id: 'settled', data: 'lands before the mark', type: NounType.Thing })
|
||||
await brain.awaitPendingEmbeds()
|
||||
await new Promise((r) => setTimeout(r, 50))
|
||||
|
||||
// A deferred write whose embed never lands: block the worker, then drop
|
||||
// the instance without close() — the unclean-stop shape.
|
||||
;(brain as any).kickEmbedWorker = () => {}
|
||||
await brain.add({
|
||||
id: 'orphan',
|
||||
data: 'enqueued then abandoned',
|
||||
type: NounType.Thing,
|
||||
deferEmbedding: true
|
||||
})
|
||||
expect(brain.pendingEmbedCount()).toBeGreaterThan(0)
|
||||
// No close(): simulate the crash by releasing only the writer lock so the
|
||||
// next open can proceed.
|
||||
await (brain as any).storage.releaseWriterLock()
|
||||
|
||||
const brain2 = await open(root)
|
||||
expect(brain2.pendingEmbedCount()).toBeGreaterThan(0)
|
||||
await brain2.awaitPendingEmbeds()
|
||||
expect(brain2.pendingEmbedCount()).toBe(0)
|
||||
await brain2.close()
|
||||
// Reap the crashed instance: its fence is gone, so close() fails loudly —
|
||||
// swallow that here; the point is clearing its watchers and registry entry.
|
||||
await brain.close().catch(() => undefined)
|
||||
})
|
||||
|
||||
it('a reopened brain has its pending set settled when open() returns', async () => {
|
||||
const root = dir()
|
||||
const brain = await open(root)
|
||||
await brain.add({ id: 'a-row', data: 'some data', type: NounType.Thing })
|
||||
await brain.awaitPendingEmbeds()
|
||||
await brain.close()
|
||||
|
||||
const brain2 = await open(root)
|
||||
// The crash-recovery contract: markers are re-armed by open itself —
|
||||
// no latch, no background race. (Here the drain landed, so zero.)
|
||||
expect(brain2.pendingEmbedCount()).toBe(0)
|
||||
await brain2.close()
|
||||
})
|
||||
|
||||
it('a clean close with an empty set writes the mark even if no drain happened', async () => {
|
||||
const root = dir()
|
||||
const brain = await open(root)
|
||||
await brain.add({ id: 'r1', data: 'row one', type: NounType.Thing })
|
||||
await brain.awaitPendingEmbeds()
|
||||
await brain.close()
|
||||
// Read the mark back through the storage door (the adapter owns the
|
||||
// on-disk encoding), on a fresh instance.
|
||||
const brain2 = await open(root)
|
||||
const mark = (await (brain2 as any).storage.readRawObject(LOWWATER_PATH)) as {
|
||||
generation: number
|
||||
} | null
|
||||
expect(mark).not.toBeNull()
|
||||
expect(mark!.generation).toBeGreaterThan(0)
|
||||
await brain2.close()
|
||||
})
|
||||
})
|
||||
|
|
@ -1,89 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/related-verb-array
|
||||
* @description related() honours EVERY verb type in an array (10.4.9).
|
||||
*
|
||||
* The storage fast paths for `sourceId + verbType` and `verbType` collapsed a
|
||||
* verb-type ARRAY to its first element — `related({ from, type: [a, b] })`
|
||||
* silently returned only `a` edges, whichever order the array came in. The
|
||||
* same quiet-loss class as the graph-first paging defect, one seam over.
|
||||
* These pins seed a store where the SECOND requested type's edge must come
|
||||
* back, on every path the collapse lived in.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType, VerbType } from '../../src/types/graphTypes'
|
||||
import { v5 } from '../../src/universal/uuid'
|
||||
|
||||
describe('related() with a verb-type array returns every requested type', () => {
|
||||
let brain: Brainy<any>
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
for (const id of ['a', 'b', 'c', 'd']) {
|
||||
await brain.add({ id, data: `node ${id}`, type: NounType.Person })
|
||||
}
|
||||
await brain.relate({ from: 'a', to: 'b', type: VerbType.Supports })
|
||||
await brain.relate({ from: 'a', to: 'c', type: VerbType.RelatedTo })
|
||||
await brain.relate({ from: 'a', to: 'd', type: VerbType.Knows })
|
||||
await brain.relate({ from: 'b', to: 'c', type: VerbType.RelatedTo })
|
||||
})
|
||||
|
||||
afterAll(async () => {
|
||||
brain = null as any
|
||||
})
|
||||
|
||||
it('from + type array: the second type\'s edge comes back, both orders', async () => {
|
||||
for (const types of [
|
||||
[VerbType.Supports, VerbType.RelatedTo],
|
||||
[VerbType.RelatedTo, VerbType.Supports]
|
||||
]) {
|
||||
const edges = await brain.related({ from: 'a', type: types })
|
||||
const targets = new Set(edges.map((e) => e.to))
|
||||
expect(targets.has(v5('b')), `types [${types}] missing Supports edge`).toBe(true)
|
||||
expect(targets.has(v5('c')), `types [${types}] missing RelatedTo edge`).toBe(true)
|
||||
expect(targets.has(v5('d'))).toBe(false)
|
||||
expect(edges).toHaveLength(2)
|
||||
}
|
||||
})
|
||||
|
||||
it('a single-element array behaves exactly like the scalar', async () => {
|
||||
const scalar = await brain.related({ from: 'a', type: VerbType.Supports })
|
||||
const array = await brain.related({ from: 'a', type: [VerbType.Supports] })
|
||||
expect(array.map((e) => e.id).sort()).toEqual(scalar.map((e) => e.id).sort())
|
||||
expect(array).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('no duplicate edges when types overlap the same edge set', async () => {
|
||||
const edges = await brain.related({
|
||||
from: 'a',
|
||||
type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows]
|
||||
})
|
||||
const ids = edges.map((e) => e.id)
|
||||
expect(new Set(ids).size).toBe(ids.length)
|
||||
expect(edges).toHaveLength(3)
|
||||
})
|
||||
|
||||
it('type-only asks (no anchor) honour the whole array too', async () => {
|
||||
const edges = await brain.related({ type: [VerbType.Supports, VerbType.Knows] })
|
||||
const verbs = new Set(edges.map((e) => e.type))
|
||||
expect(verbs.has(VerbType.Supports)).toBe(true)
|
||||
expect(verbs.has(VerbType.Knows)).toBe(true)
|
||||
expect(edges).toHaveLength(2)
|
||||
})
|
||||
|
||||
it('to + type array: the target side honours every type too', async () => {
|
||||
const edges = await brain.related({ to: 'c', type: [VerbType.RelatedTo, VerbType.Supports] })
|
||||
const froms = new Set(edges.map((e) => e.from))
|
||||
expect(froms.has(v5('a'))).toBe(true)
|
||||
expect(froms.has(v5('b'))).toBe(true)
|
||||
expect(edges).toHaveLength(2)
|
||||
})
|
||||
|
||||
it('pagination stays consistent across the union', async () => {
|
||||
const page1 = await brain.related({ from: 'a', type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows], limit: 2 })
|
||||
const page2 = await brain.related({ from: 'a', type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows], limit: 2, offset: 2 })
|
||||
const all = [...page1, ...page2].map((e) => e.id)
|
||||
expect(new Set(all).size).toBe(3)
|
||||
})
|
||||
})
|
||||
|
|
@ -1,184 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/transact-edge-delete-bigint-aliasing
|
||||
* @description Regression for a fleet-adoption blocker: ANY edge delete
|
||||
* inside `transact()` — a direct unrelate or a noun-remove's cascade —
|
||||
* aborted with the metadata seam's BigInt JSON-guard error on a strict
|
||||
* (native) metadata provider.
|
||||
*
|
||||
* The aliasing chain: `planTxUnrelate`/the remove-cascade pass the SAME verb
|
||||
* object to the graph-retraction op and the metadata-retraction op. The
|
||||
* metadata leg's JSON-safe wrap ran at PLAN time, when the verb was still
|
||||
* clean — so it returned the same reference. At EXECUTE time the graph op
|
||||
* runs first and `resolveVerbEndpointInts` mirrors BigInt
|
||||
* `sourceInt`/`targetInt` onto the shared object (deliberately deferred for
|
||||
* same-batch forward refs — see transact-forward-ref-graph.test.ts); the
|
||||
* metadata op then crossed the seam with the polluted object. Direct
|
||||
* `unrelate()` resolves ints at BUILD time, before its sanitize, which is why
|
||||
* only the transact() shapes ever hit it.
|
||||
*
|
||||
* Fix under pin: the JSON-safe view is taken AT THE CROSSING — inside the
|
||||
* metadata-index operations' execute/rollback — so no plan-vs-execute
|
||||
* ordering can bypass it. The JS baseline index tolerates BigInts (it would
|
||||
* mask the bug), so these pins SPY on the seam and assert what actually
|
||||
* crossed, exactly as a strict native provider would judge it.
|
||||
*/
|
||||
import { describe, it, expect, beforeEach, afterEach } from 'vitest'
|
||||
import * as fs from 'node:fs'
|
||||
import * as os from 'node:os'
|
||||
import * as path from 'node:path'
|
||||
import { Brainy } from '../../src/brainy.js'
|
||||
import { NounType, VerbType } from '../../src/types/graphTypes.js'
|
||||
import {
|
||||
AddToMetadataIndexOperation,
|
||||
RemoveFromMetadataIndexOperation
|
||||
} from '../../src/transaction/operations/index.js'
|
||||
|
||||
let seq = 0
|
||||
const freshId = (): string =>
|
||||
`00000000-0000-4000-8000-${(++seq).toString(16).padStart(12, '0')}`
|
||||
|
||||
/** Top-level BigInt-valued keys of a candidate seam crossing (the guard's law). */
|
||||
const bigintKeys = (metadata: unknown): string[] => {
|
||||
if (metadata === null || typeof metadata !== 'object') return []
|
||||
return Object.entries(metadata as Record<string, unknown>)
|
||||
.filter(([, v]) => typeof v === 'bigint')
|
||||
.map(([k]) => k)
|
||||
}
|
||||
|
||||
describe('transact() edge deletes never carry BigInt across the metadata seam', () => {
|
||||
let dir: string
|
||||
let brain: any
|
||||
let crossings: Array<{ door: string; id: string; keys: string[] }>
|
||||
|
||||
beforeEach(async () => {
|
||||
process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true'
|
||||
dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-tx-bigint-'))
|
||||
brain = new Brainy({
|
||||
requireSubtype: false,
|
||||
storage: { type: 'filesystem', path: dir },
|
||||
dimensions: 384,
|
||||
silent: true
|
||||
})
|
||||
await brain.init()
|
||||
|
||||
// Spy on the seam the way a strict native provider judges it: record the
|
||||
// BigInt-valued top-level keys of every metadata argument that crosses.
|
||||
// The JS baseline index tolerates BigInts, so without this the baseline
|
||||
// run would green a shape the native pair aborts on.
|
||||
crossings = []
|
||||
const index = brain.metadataIndex
|
||||
for (const door of ['addToIndex', 'removeFromIndex'] as const) {
|
||||
const real = index[door].bind(index)
|
||||
index[door] = (id: string, metadata: unknown, ...rest: unknown[]) => {
|
||||
crossings.push({ door, id, keys: bigintKeys(metadata) })
|
||||
return real(id, metadata, ...rest)
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await brain.close()
|
||||
fs.rmSync(dir, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
it('CASE 1 (the fleet repro): relate, then transact([{op: unrelate}])', async () => {
|
||||
const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing })
|
||||
const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing })
|
||||
const verbId = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo })
|
||||
|
||||
crossings.length = 0
|
||||
await brain.transact([{ op: 'unrelate', id: verbId }])
|
||||
|
||||
const polluted = crossings.filter((c) => c.keys.length > 0)
|
||||
expect(polluted).toEqual([])
|
||||
expect(await brain.storage.getVerb(verbId)).toBeFalsy()
|
||||
})
|
||||
|
||||
it('CASE 2 (the cascade shape): transact([{op: remove}]) cascading edge deletes', async () => {
|
||||
const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing })
|
||||
const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing })
|
||||
const c = await brain.add({ id: freshId(), data: 'c', type: NounType.Thing })
|
||||
const ab = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo })
|
||||
const ca = await brain.relate({ from: c, to: a, type: VerbType.RelatedTo })
|
||||
|
||||
crossings.length = 0
|
||||
await brain.transact([{ op: 'remove', id: a }])
|
||||
|
||||
const polluted = crossings.filter((c2) => c2.keys.length > 0)
|
||||
expect(polluted).toEqual([])
|
||||
expect(await brain.get(a)).toBeFalsy()
|
||||
expect(await brain.storage.getVerb(ab)).toBeFalsy()
|
||||
expect(await brain.storage.getVerb(ca)).toBeFalsy()
|
||||
})
|
||||
|
||||
it('CASE 3 (one batch, both legs): adds + relate + unrelate of a pre-existing edge', async () => {
|
||||
const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing })
|
||||
const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing })
|
||||
const old = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo })
|
||||
|
||||
const x = freshId()
|
||||
crossings.length = 0
|
||||
await brain.transact([
|
||||
{ op: 'add', id: x, data: 'x', type: NounType.Thing },
|
||||
{ op: 'relate', from: a, to: x, type: VerbType.RelatedTo },
|
||||
{ op: 'unrelate', id: old }
|
||||
])
|
||||
|
||||
const polluted = crossings.filter((c) => c.keys.length > 0)
|
||||
expect(polluted).toEqual([])
|
||||
expect(await brain.storage.getVerb(old)).toBeFalsy()
|
||||
const edges = await brain.related({ from: a })
|
||||
expect(edges.length).toBe(1)
|
||||
expect(edges[0].id).not.toBe(old)
|
||||
})
|
||||
})
|
||||
|
||||
describe('the metadata-index operations sanitize at the crossing, not at construction', () => {
|
||||
/** A strict seam: refuses BigInts exactly as the native provider does. */
|
||||
const strictIndex = () => {
|
||||
const seen: Array<{ door: string; keys: string[] }> = []
|
||||
const judge = (door: string, metadata: unknown) => {
|
||||
const keys = bigintKeys(metadata)
|
||||
seen.push({ door, keys })
|
||||
if (keys.length > 0) {
|
||||
throw new Error(
|
||||
`${door}: the metadata object violates the provider seam's JSON ` +
|
||||
`contract — BigInt at ${keys.join(', ')}.`
|
||||
)
|
||||
}
|
||||
}
|
||||
return {
|
||||
seen,
|
||||
addToIndex: async (_id: string, metadata: unknown) => judge('addToIndex', metadata),
|
||||
removeFromIndex: async (_id: string, metadata: unknown) => judge('removeFromIndex', metadata)
|
||||
}
|
||||
}
|
||||
|
||||
it('RemoveFromMetadataIndexOperation: entity mutated AFTER construction still crosses clean', async () => {
|
||||
const index = strictIndex()
|
||||
const verb: Record<string, unknown> = { id: 'v1', sourceId: 'a', targetId: 'b' }
|
||||
const op = new RemoveFromMetadataIndexOperation(index as any, 'v1', verb, () => 7n)
|
||||
|
||||
// The graph leg's execute-time endpoint resolution, simulated: the shared
|
||||
// object is polluted between plan and execute.
|
||||
verb.sourceInt = 800_000n
|
||||
verb.targetInt = 800_001n
|
||||
|
||||
const rollback = await op.execute()
|
||||
await rollback()
|
||||
expect(index.seen.map((s) => s.keys)).toEqual([[], []])
|
||||
})
|
||||
|
||||
it('AddToMetadataIndexOperation: same law on the add leg and its rollback', async () => {
|
||||
const index = strictIndex()
|
||||
const verb: Record<string, unknown> = { id: 'v2', sourceId: 'a', targetId: 'b' }
|
||||
const op = new AddToMetadataIndexOperation(index as any, 'v2', verb, () => 7n)
|
||||
|
||||
verb.sourceInt = 800_000n
|
||||
verb.targetInt = 800_001n
|
||||
|
||||
const rollback = await op.execute()
|
||||
await rollback()
|
||||
expect(index.seen.map((s) => s.keys)).toEqual([[], []])
|
||||
})
|
||||
})
|
||||
|
|
@ -1,115 +0,0 @@
|
|||
/**
|
||||
* @module tests/integration/vfs-containment-batched
|
||||
* @description repairContainment costs O(edges/page) graph calls, not O(entities) (10.4.9 train).
|
||||
*
|
||||
* Pass 2 used to issue one awaited `related({ to })` per VFS entity — minutes
|
||||
* of serialized graph calls on large brains. Now one paged walk over every
|
||||
* Contains edge feeds an in-memory group-by-target, and only actual defects
|
||||
* mutate. These pins hold the verdicts (duplicate removed, stale parent
|
||||
* removed, missing edge restored, user knowledge edges untouched) AND the
|
||||
* cost shape (related() call count independent of the entity count).
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll, vi } from 'vitest'
|
||||
import { Brainy } from '../../src/brainy'
|
||||
import { NounType, VerbType } from '../../src/types/graphTypes'
|
||||
|
||||
const FILES = 60
|
||||
|
||||
describe('repairContainment: batched pass 2', () => {
|
||||
let brain: Brainy<any>
|
||||
let result: { removed: number; restored: number }
|
||||
let relatedCalls = 0
|
||||
|
||||
beforeAll(async () => {
|
||||
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
|
||||
await brain.init()
|
||||
const vfs = (brain as any).vfs ?? (brain as any)._vfs
|
||||
expect(vfs).toBeTruthy()
|
||||
await vfs.init()
|
||||
|
||||
// A directory and FILES entries under it, wired as real VFS rows.
|
||||
const mkNode = async (id: string, path: string, vfsType: string): Promise<void> => {
|
||||
await brain.add({
|
||||
id,
|
||||
data: `vfs node ${path}`,
|
||||
type: NounType.File,
|
||||
visibility: 'system',
|
||||
metadata: { vfsType, path }
|
||||
})
|
||||
}
|
||||
await mkNode('dir', '/docs', 'directory')
|
||||
const rootId = vfs.rootEntityId ?? (await vfs.initializeRoot?.())
|
||||
if (rootId) {
|
||||
await brain.relate({
|
||||
from: rootId,
|
||||
to: 'dir',
|
||||
type: VerbType.Contains,
|
||||
subtype: 'vfs-contains',
|
||||
metadata: { isVFS: true }
|
||||
})
|
||||
}
|
||||
for (let i = 0; i < FILES; i++) {
|
||||
await mkNode(`f-${i}`, `/docs/f-${i}.md`, 'file')
|
||||
if (i === 0) continue // f-0: MISSING edge — must be restored
|
||||
await brain.relate({
|
||||
from: 'dir',
|
||||
to: `f-${i}`,
|
||||
type: VerbType.Contains,
|
||||
subtype: 'vfs-contains',
|
||||
metadata: { isVFS: true }
|
||||
})
|
||||
}
|
||||
// NOTE: relate() is idempotent for an identical from/to/type, so a true
|
||||
// duplicate (a concurrent-writer artifact) cannot be seeded through the
|
||||
// public API — the duplicate branch is covered by the tree-correctness
|
||||
// pin below, which proves at most one vfs edge survives per file.
|
||||
// f-2: STALE parent edge (from a sibling file) — must be removed.
|
||||
await brain.relate({
|
||||
from: 'f-3',
|
||||
to: 'f-2',
|
||||
type: VerbType.Contains,
|
||||
subtype: 'vfs-contains',
|
||||
metadata: { isVFS: true }
|
||||
})
|
||||
// A USER knowledge Contains edge (not vfs-flagged) — must be untouched.
|
||||
await brain.relate({ from: 'f-4', to: 'f-5', type: VerbType.Contains })
|
||||
|
||||
const spy = vi.spyOn(brain, 'related')
|
||||
result = await vfs.repairContainment()
|
||||
relatedCalls = spy.mock.calls.length
|
||||
spy.mockRestore()
|
||||
})
|
||||
|
||||
afterAll(async () => {
|
||||
brain = null as any
|
||||
})
|
||||
|
||||
it('restores the missing edge and removes the stale parent — exactly', () => {
|
||||
expect(result.restored).toBe(1) // f-0's missing edge
|
||||
expect(result.removed).toBe(1) // f-2's stale parent (f-3 → f-2)
|
||||
})
|
||||
|
||||
it('the repaired tree is correct: every file has exactly one vfs edge from its dir', async () => {
|
||||
for (let i = 0; i < 6; i++) {
|
||||
const incoming = await brain.related({ to: `f-${i}`, type: VerbType.Contains })
|
||||
const vfsEdges = incoming.filter(
|
||||
(e) => e.subtype === 'vfs-contains' || (e.metadata as any)?.isVFS === true
|
||||
)
|
||||
expect(vfsEdges, `f-${i}`).toHaveLength(1)
|
||||
}
|
||||
})
|
||||
|
||||
it('never touches user knowledge edges', async () => {
|
||||
const incoming = await brain.related({ to: 'f-5', type: VerbType.Contains })
|
||||
const user = incoming.filter(
|
||||
(e) => e.subtype !== 'vfs-contains' && (e.metadata as any)?.isVFS !== true
|
||||
)
|
||||
expect(user).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('cost shape: related() calls do not scale with the entity count', () => {
|
||||
// One paged type-only walk (~E/1000 pages) — with 60+ entities the old
|
||||
// shape issued 60+ calls; the new one a handful. Bound generously.
|
||||
expect(relatedCalls).toBeLessThanOrEqual(5)
|
||||
})
|
||||
})
|
||||
395
tests/unit/release/wall-entry.test.ts
Normal file
395
tests/unit/release/wall-entry.test.ts
Normal file
|
|
@ -0,0 +1,395 @@
|
|||
/**
|
||||
* scripts/wall-entry.mjs — the mechanical releases-wall entry.
|
||||
*
|
||||
* The script's only real interface is its CLI (it has no importable
|
||||
* exports by design — one door, no parallel API to drift from it), so
|
||||
* these tests spawn it exactly as scripts/release.sh does: as a child
|
||||
* process, against a fixture CHANGELOG and a throwaway local bare repo
|
||||
* standing in for git@source.soulcraft.com:soulcraftlabs/releases.git
|
||||
* (--remote) plus a throwaway cache directory (--cache-dir) standing in
|
||||
* for ~/.cache/soulcraft-releases — never the real remote, never the
|
||||
* real developer cache.
|
||||
*/
|
||||
import { describe, it, expect, beforeEach, afterEach } from 'vitest'
|
||||
import { execFileSync } from 'node:child_process'
|
||||
import { mkdtempSync, rmSync, writeFileSync, readFileSync, chmodSync } from 'node:fs'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
|
||||
const SCRIPT = join(process.cwd(), 'scripts/wall-entry.mjs')
|
||||
|
||||
/** Run the script and capture the outcome without throwing on a non-zero exit. */
|
||||
function run(args: string[], cwd: string): { status: number; stdout: string; stderr: string } {
|
||||
try {
|
||||
const stdout = execFileSync('node', [SCRIPT, ...args], { cwd, encoding: 'utf8' })
|
||||
return { status: 0, stdout, stderr: '' }
|
||||
} catch (err: any) {
|
||||
return { status: err.status ?? 1, stdout: err.stdout ?? '', stderr: err.stderr ?? '' }
|
||||
}
|
||||
}
|
||||
|
||||
function git(args: string[], cwd: string): string {
|
||||
return execFileSync('git', ['-C', cwd, ...args], { encoding: 'utf8' }).trim()
|
||||
}
|
||||
|
||||
const CHANGELOG_HEADER = '# Changelog\n\nAll notable changes, in this fixture.\n'
|
||||
|
||||
/** Build a CHANGELOG.md with one entry per [version, bullets[]] pair, newest first. */
|
||||
function buildChangelog(entries: Array<{ version: string; date: string; bullets: string[] }>): string {
|
||||
const body = entries
|
||||
.map(
|
||||
(e) =>
|
||||
`### [${e.version}](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/vX...v${e.version}) (${e.date})\n\n` +
|
||||
e.bullets.map((b) => `- ${b} (abc1234)`).join('\n') +
|
||||
'\n',
|
||||
)
|
||||
.join('\n')
|
||||
return CHANGELOG_HEADER + '\n' + body
|
||||
}
|
||||
|
||||
function wallFile(product: string, entries: unknown[]): string {
|
||||
return JSON.stringify({ product, entries }, null, 2) + '\n'
|
||||
}
|
||||
|
||||
const BASE_ENTRY = {
|
||||
version: '10.4.11',
|
||||
date: '2026-09-02',
|
||||
headline: 'A faster open',
|
||||
items: ['A faster open.'],
|
||||
url: 'https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.11',
|
||||
thumb: null,
|
||||
}
|
||||
|
||||
/** A throwaway bare repo standing in for the real soulcraftlabs/releases remote. */
|
||||
function initBareRemote(): string {
|
||||
const remoteDir = mkdtempSync(join(tmpdir(), 'wall-remote-'))
|
||||
execFileSync('git', ['init', '--bare', '-b', 'main', remoteDir])
|
||||
return remoteDir
|
||||
}
|
||||
|
||||
/** Seed the bare remote with an initial <product>.json, via a throwaway clone. */
|
||||
function seedRemote(remoteDir: string, product: string, entries: unknown[]): void {
|
||||
const seedDir = mkdtempSync(join(tmpdir(), 'wall-seed-'))
|
||||
execFileSync('git', ['clone', remoteDir, seedDir], { stdio: 'ignore' })
|
||||
git(['config', 'user.email', 'seed@example.com'], seedDir)
|
||||
git(['config', 'user.name', 'Seed'], seedDir)
|
||||
writeFileSync(join(seedDir, `${product}.json`), wallFile(product, entries))
|
||||
git(['add', `${product}.json`], seedDir)
|
||||
git(['commit', '-m', 'seed'], seedDir)
|
||||
git(['push', 'origin', 'main'], seedDir)
|
||||
rmSync(seedDir, { recursive: true, force: true })
|
||||
}
|
||||
|
||||
/** Read <product>.json back out of the bare remote's main tip, via a throwaway clone. */
|
||||
function readRemote(remoteDir: string, product: string): any {
|
||||
const readDir = mkdtempSync(join(tmpdir(), 'wall-read-'))
|
||||
execFileSync('git', ['clone', remoteDir, readDir], { stdio: 'ignore' })
|
||||
const data = JSON.parse(readFileSync(join(readDir, `${product}.json`), 'utf8'))
|
||||
rmSync(readDir, { recursive: true, force: true })
|
||||
return data
|
||||
}
|
||||
|
||||
/** Reject every push — stands in for any push failure (including a genuine
|
||||
* non-fast-forward raced by a concurrent release rail), which this script
|
||||
* treats identically: refuse loudly, name the cure, touch nothing further. */
|
||||
function makeRemoteRejectPushes(remoteDir: string): void {
|
||||
const hookPath = join(remoteDir, 'hooks', 'pre-receive')
|
||||
writeFileSync(hookPath, '#!/bin/sh\necho "remote: simulated push rejection" >&2\nexit 1\n')
|
||||
chmodSync(hookPath, 0o755)
|
||||
}
|
||||
|
||||
let dir: string
|
||||
let remoteDir: string
|
||||
let cacheDir: string
|
||||
|
||||
beforeEach(() => {
|
||||
dir = mkdtempSync(join(tmpdir(), 'wall-entry-test-'))
|
||||
remoteDir = initBareRemote()
|
||||
cacheDir = join(mkdtempSync(join(tmpdir(), 'wall-cache-')), 'soulcraft-releases')
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
rmSync(dir, { recursive: true, force: true })
|
||||
rmSync(remoteDir, { recursive: true, force: true })
|
||||
rmSync(cacheDir, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('wall-entry.mjs — generate + publish', () => {
|
||||
it('derives headline from the first bullet and items from every bullet, hashes stripped, and pushes it to the remote', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [BASE_ENTRY])
|
||||
writeFileSync(
|
||||
join(dir, 'CHANGELOG.md'),
|
||||
buildChangelog([{ version: '10.4.12', date: '2026-09-03', bullets: ['fix(wall): mechanize the entry', 'test(wall): pin the shape'] }]),
|
||||
)
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '10.4.12', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
expect(result.status).toBe(0)
|
||||
expect(result.stdout).toMatch(/wrote v10\.4\.12.*pushed/i)
|
||||
|
||||
const wall = readRemote(remoteDir, 'open-brainy')
|
||||
expect(wall.entries).toHaveLength(2)
|
||||
expect(wall.entries[0]).toEqual({
|
||||
version: '10.4.12',
|
||||
date: '2026-09-03',
|
||||
headline: 'fix(wall): mechanize the entry',
|
||||
items: ['fix(wall): mechanize the entry', 'test(wall): pin the shape'],
|
||||
url: 'https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.12',
|
||||
thumb: null,
|
||||
})
|
||||
// the older entry stays put, still second
|
||||
expect(wall.entries[1].version).toBe('10.4.11')
|
||||
})
|
||||
|
||||
it('prepends newest-first — the new entry lands at index 0 ahead of every existing one', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [BASE_ENTRY, { ...BASE_ENTRY, version: '10.4.10' }])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.5.0', date: '2026-09-03', bullets: ['feat: ten five'] }]))
|
||||
|
||||
run(['--product', 'open-brainy', '--version', '10.5.0', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir], dir)
|
||||
|
||||
const wall = readRemote(remoteDir, 'open-brainy')
|
||||
expect(wall.entries.map((e: any) => e.version)).toEqual(['10.5.0', '10.4.11', '10.4.10'])
|
||||
})
|
||||
|
||||
it('replaces an entry with the same version instead of duplicating it — idempotent re-runs', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [
|
||||
{ ...BASE_ENTRY, headline: 'stale headline, pre-fix' },
|
||||
{ ...BASE_ENTRY, version: '10.4.10' },
|
||||
])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.11', date: '2026-09-02', bullets: ['fix: the corrected headline'] }]))
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '10.4.11', '--date', '2026-09-02', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
expect(result.status).toBe(0)
|
||||
expect(result.stdout).toMatch(/replaced v10\.4\.11/i)
|
||||
|
||||
const wall = readRemote(remoteDir, 'open-brainy')
|
||||
expect(wall.entries).toHaveLength(2) // not 3 — replaced, not duplicated
|
||||
expect(wall.entries[0].version).toBe('10.4.11')
|
||||
expect(wall.entries[0].headline).toBe('fix: the corrected headline')
|
||||
expect(wall.entries[1].version).toBe('10.4.10')
|
||||
})
|
||||
|
||||
it('a re-run with byte-identical content commits nothing and still succeeds', () => {
|
||||
// headline always equals items[0] for a derived entry, so this fixture
|
||||
// (unlike BASE_ENTRY, whose headline/items intentionally diverge for the
|
||||
// shape-only tests below) has to keep the two in lockstep to ever roundtrip.
|
||||
const stableEntry = { ...BASE_ENTRY, headline: 'A faster open.', items: ['A faster open.'] }
|
||||
seedRemote(remoteDir, 'open-brainy', [stableEntry])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.11', date: '2026-09-02', bullets: ['A faster open.'] }]))
|
||||
const before = readRemote(remoteDir, 'open-brainy')
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '10.4.11', '--date', '2026-09-02', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
expect(result.status).toBe(0)
|
||||
expect(result.stdout).toMatch(/nothing to commit/i)
|
||||
expect(readRemote(remoteDir, 'open-brainy')).toEqual(before)
|
||||
})
|
||||
|
||||
it('derives the public package-page permalink for the product engine (private repo, never null)', () => {
|
||||
seedRemote(remoteDir, 'brainy', [{ ...BASE_ENTRY, version: '11.0.5', url: 'https://source.soulcraft.com/soulcraft/-/packages/npm/@soulcraft%2Fbrainy/11.0.5' }])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '11.0.6', date: '2026-09-03', bullets: ['fix: a native-only fix'] }]))
|
||||
|
||||
const result = run(
|
||||
['--product', 'brainy', '--version', '11.0.6', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
expect(result.status).toBe(0)
|
||||
|
||||
const wall = readRemote(remoteDir, 'brainy')
|
||||
expect(wall.entries[0].url).toBe('https://source.soulcraft.com/soulcraft/-/packages/npm/@soulcraft%2Fbrainy/11.0.6')
|
||||
expect(wall.entries[0].thumb).toBeNull()
|
||||
})
|
||||
|
||||
it('refuses a product with no permalink pattern, naming the cure', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [BASE_ENTRY])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '1.0.0', date: '2026-09-03', bullets: ['feat: first'] }]))
|
||||
|
||||
const result = run(['--product', 'mystery', '--version', '1.0.0', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir], dir)
|
||||
expect(result.status).not.toBe(0)
|
||||
expect(result.stderr).toMatch(/no permalink pattern for product "mystery"/)
|
||||
expect(result.stderr).toMatch(/never carry url: null/)
|
||||
})
|
||||
|
||||
it('refuses when the CHANGELOG has no entry yet for the target version, and touches no remote', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.11', date: '2026-09-02', bullets: ['fix: whatever'] }]))
|
||||
const beforeSha = git(['rev-parse', 'main'], remoteDir)
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '99.0.0', '--date', '2026-09-02', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/no CHANGELOG entry yet/i)
|
||||
expect(git(['rev-parse', 'main'], remoteDir)).toBe(beforeSha)
|
||||
})
|
||||
|
||||
it('refuses by naming the cure when the remote cannot be cloned', () => {
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.12', date: '2026-09-03', bullets: ['fix: whatever'] }]))
|
||||
const noSuchRemote = join(tmpdir(), 'wall-remote-does-not-exist-' + Date.now())
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '10.4.12', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', noSuchRemote, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/cannot clone/i)
|
||||
expect(result.stderr).toMatch(/cure:/i)
|
||||
})
|
||||
|
||||
it('refuses by naming the cure, and touches no remote, when the fetched wall fails shape validation', () => {
|
||||
const seedDir = mkdtempSync(join(tmpdir(), 'wall-seed-broken-'))
|
||||
execFileSync('git', ['clone', remoteDir, seedDir], { stdio: 'ignore' })
|
||||
git(['config', 'user.email', 'seed@example.com'], seedDir)
|
||||
git(['config', 'user.name', 'Seed'], seedDir)
|
||||
writeFileSync(
|
||||
join(seedDir, 'open-brainy.json'),
|
||||
JSON.stringify({ product: 'open-brainy', entries: [{ version: '10.4.11', date: '2026-09-02', items: ['x'], url: null }] }, null, 2),
|
||||
)
|
||||
git(['add', 'open-brainy.json'], seedDir)
|
||||
git(['commit', '-m', 'seed broken'], seedDir)
|
||||
git(['push', 'origin', 'main'], seedDir)
|
||||
rmSync(seedDir, { recursive: true, force: true })
|
||||
const beforeSha = git(['rev-parse', 'main'], remoteDir)
|
||||
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.12', date: '2026-09-03', bullets: ['fix: whatever'] }]))
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '10.4.12', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/fails shape validation/i)
|
||||
expect(result.stderr).toMatch(/missing key\(s\) headline/i)
|
||||
expect(git(['rev-parse', 'main'], remoteDir)).toBe(beforeSha)
|
||||
})
|
||||
|
||||
it('refuses by naming the cure when the remote rejects the push (stands in for a raced non-fast-forward)', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [BASE_ENTRY])
|
||||
makeRemoteRejectPushes(remoteDir)
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.12', date: '2026-09-03', bullets: ['fix: whatever'] }]))
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '10.4.12', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/push to .* failed/i)
|
||||
expect(result.stderr).toMatch(/cure:/i)
|
||||
})
|
||||
|
||||
it('refuses a cross-product write when the file\'s "product" field does not match --product', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [BASE_ENTRY])
|
||||
const seedDir = mkdtempSync(join(tmpdir(), 'wall-seed-mismatch-'))
|
||||
execFileSync('git', ['clone', remoteDir, seedDir], { stdio: 'ignore' })
|
||||
git(['config', 'user.email', 'seed@example.com'], seedDir)
|
||||
git(['config', 'user.name', 'Seed'], seedDir)
|
||||
const corrupted = JSON.parse(readFileSync(join(seedDir, 'open-brainy.json'), 'utf8'))
|
||||
corrupted.product = 'brainy'
|
||||
writeFileSync(join(seedDir, 'open-brainy.json'), JSON.stringify(corrupted, null, 2) + '\n')
|
||||
git(['add', 'open-brainy.json'], seedDir)
|
||||
git(['commit', '-m', 'corrupt product field'], seedDir)
|
||||
git(['push', 'origin', 'main'], seedDir)
|
||||
rmSync(seedDir, { recursive: true, force: true })
|
||||
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '1.0.0', date: '2026-09-03', bullets: ['fix: wrong repo'] }]))
|
||||
|
||||
const result = run(
|
||||
['--product', 'open-brainy', '--version', '1.0.0', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/product "brainy".*--product "open-brainy"/i)
|
||||
})
|
||||
})
|
||||
|
||||
describe('wall-entry.mjs — --dry-run', () => {
|
||||
it('prints the entry and the target path, and touches neither the cache dir nor the remote', () => {
|
||||
seedRemote(remoteDir, 'open-brainy', [BASE_ENTRY])
|
||||
writeFileSync(join(dir, 'CHANGELOG.md'), buildChangelog([{ version: '10.4.12', date: '2026-09-03', bullets: ['fix: a dry run'] }]))
|
||||
const beforeSha = git(['rev-parse', 'main'], remoteDir)
|
||||
|
||||
const result = run(
|
||||
['--dry-run', '--product', 'open-brainy', '--version', '10.4.12', '--date', '2026-09-03', '--from-changelog', 'CHANGELOG.md', '--remote', remoteDir, '--cache-dir', cacheDir],
|
||||
dir,
|
||||
)
|
||||
|
||||
expect(result.status).toBe(0)
|
||||
expect(result.stdout).toMatch(/would write to/i)
|
||||
expect(result.stdout).toMatch(/"version": "10\.4\.12"/)
|
||||
expect(git(['rev-parse', 'main'], remoteDir)).toBe(beforeSha)
|
||||
})
|
||||
})
|
||||
|
||||
describe('wall-entry.mjs — --check', () => {
|
||||
it('passes a well-formed, newest-first file with no duplicates', () => {
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [BASE_ENTRY, { ...BASE_ENTRY, version: '10.4.10' }]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(0)
|
||||
expect(result.stdout).toMatch(/OK/)
|
||||
})
|
||||
|
||||
it('passes a file where "thumb" is entirely absent (optional per the HQ contract)', () => {
|
||||
const { thumb, ...noThumb } = BASE_ENTRY as any
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [noThumb]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(0)
|
||||
})
|
||||
|
||||
it('catches a missing entry key', () => {
|
||||
const broken = { version: '1.0.0', date: '2026-09-03', headline: 'h', items: ['i'] } // no "url"
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [broken]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/missing key\(s\) url/)
|
||||
})
|
||||
|
||||
it('catches an unexpected top-level key (e.g. the retired "history" field)', () => {
|
||||
const raw = JSON.parse(wallFile('open-brainy', [BASE_ENTRY]))
|
||||
raw.history = 'retired field'
|
||||
writeFileSync(join(dir, 'wall.json'), JSON.stringify(raw))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/unexpected key\(s\) history/)
|
||||
})
|
||||
|
||||
it('catches entries that are not newest-first', () => {
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [{ ...BASE_ENTRY, version: '10.4.10' }, BASE_ENTRY]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/not newest-first/)
|
||||
})
|
||||
|
||||
it('catches a duplicate version even with identical entries', () => {
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [BASE_ENTRY, { ...BASE_ENTRY }]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/duplicate version 10\.4\.11/)
|
||||
})
|
||||
|
||||
it('catches an empty items array', () => {
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [{ ...BASE_ENTRY, items: [] }]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/"items" must be a non-empty array/)
|
||||
})
|
||||
|
||||
it('catches a malformed date', () => {
|
||||
writeFileSync(join(dir, 'wall.json'), wallFile('open-brainy', [{ ...BASE_ENTRY, date: '09/03/2026' }]))
|
||||
const result = run(['--check', '--file', 'wall.json'], dir)
|
||||
expect(result.status).toBe(1)
|
||||
expect(result.stderr).toMatch(/"date" must be a YYYY-MM-DD string/)
|
||||
})
|
||||
})
|
||||
Loading…
Add table
Add a link
Reference in a new issue