Compare commits
11 commits
fix/relate
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| 4f1e27c9a0 | |||
| 297a3d7657 | |||
| 7ab670b525 | |||
| f097cbf6f2 | |||
| e64e2bc175 | |||
| 655aa13ea7 | |||
| 39c71ecdac | |||
| 0759c03a82 | |||
| 9a888c37e9 | |||
|
|
298cb6daca | ||
| b8475cc86a |
15 changed files with 265 additions and 604 deletions
12
CHANGELOG.md
12
CHANGELOG.md
|
|
@ -2,18 +2,6 @@
|
||||||
|
|
||||||
All notable changes to this project will be documented in this file. See [standard-version](https://github.com/conventional-changelog/standard-version) for commit guidelines.
|
All notable changes to this project will be documented in this file. See [standard-version](https://github.com/conventional-changelog/standard-version) for commit guidelines.
|
||||||
|
|
||||||
### [10.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)
|
### [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)
|
- fix(vfs): the old-root sweep narrates only when it has something to say (d49148e1)
|
||||||
|
|
|
||||||
4
package-lock.json
generated
4
package-lock.json
generated
|
|
@ -1,12 +1,12 @@
|
||||||
{
|
{
|
||||||
"name": "@soulcraftlabs/brainy",
|
"name": "@soulcraftlabs/brainy",
|
||||||
"version": "10.4.6",
|
"version": "10.4.4",
|
||||||
"lockfileVersion": 3,
|
"lockfileVersion": 3,
|
||||||
"requires": true,
|
"requires": true,
|
||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"name": "@soulcraftlabs/brainy",
|
"name": "@soulcraftlabs/brainy",
|
||||||
"version": "10.4.6",
|
"version": "10.4.4",
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@msgpack/msgpack": "^3.1.2",
|
"@msgpack/msgpack": "^3.1.2",
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
{
|
{
|
||||||
"name": "@soulcraftlabs/brainy",
|
"name": "@soulcraftlabs/brainy",
|
||||||
"version": "10.4.6",
|
"version": "10.4.4",
|
||||||
"brainyContract": 1,
|
"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.",
|
"description": "Universal Knowledge Protocol™ - World's first Triple Intelligence database unifying vector, graph, and document search in one API. Stage 3 CANONICAL: 42 nouns × 127 verbs covering 96-97% of all human knowledge.",
|
||||||
"main": "dist/index.js",
|
"main": "dist/index.js",
|
||||||
|
|
|
||||||
76
releases/brainy.json
Normal file
76
releases/brainy.json
Normal file
|
|
@ -0,0 +1,76 @@
|
||||||
|
{
|
||||||
|
"product": "brainy",
|
||||||
|
"entries": [
|
||||||
|
{
|
||||||
|
"version": "11.0.5",
|
||||||
|
"date": "2026-09-02",
|
||||||
|
"headline": "Graph-first finds in production, and opens that stop rescanning history",
|
||||||
|
"items": [
|
||||||
|
"find({ connected, where }) now walks the neighbours first and filters only those rows through a native door — correct at every page and O(neighbours), never the whole store.",
|
||||||
|
"related() with a list of verb types returns every requested kind (a fast path had silently kept only the first).",
|
||||||
|
"Deferred-embedding recovery resumes from a low-water mark instead of rescanning the whole generation log at every open — measured at two minutes on a large brain, now milliseconds."
|
||||||
|
],
|
||||||
|
"url": null,
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "11.0.4",
|
||||||
|
"date": "2026-09-01",
|
||||||
|
"headline": "Closes in milliseconds, index rebuilds without the disk-sync storm",
|
||||||
|
"items": [
|
||||||
|
"close() no longer pays deferred compaction or waits out an in-flight rebuild — measured 8 ms against the 4-minute closes it replaces; deferred work resumes at the next open, in the background.",
|
||||||
|
"The metadata index's rebuild syncs to disk per shard instead of per row, and the durability point moved to the publish step — the same guarantee, a fraction of the disk traffic.",
|
||||||
|
"A new native filter door evaluates queries over exactly the candidate rows a graph walk found, never the whole store."
|
||||||
|
],
|
||||||
|
"url": null,
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "11.0.3",
|
||||||
|
"date": "2026-09-01",
|
||||||
|
"headline": "The embedding upgrade ceremony runs on every brain",
|
||||||
|
"items": [
|
||||||
|
"A brain opened through the standard plugin now carries its embedding-model identity, so the full-precision upgrade ceremony can run on it.",
|
||||||
|
"A one-fix release; nothing else changed."
|
||||||
|
],
|
||||||
|
"url": null,
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "11.0.2",
|
||||||
|
"date": "2026-08-31",
|
||||||
|
"headline": "One embedding quality everywhere, 3–4× faster imports",
|
||||||
|
"items": [
|
||||||
|
"Every runtime embeds with the same full-precision model — search quality no longer depends on where you run.",
|
||||||
|
"Bulk embedding measured 3.1–4.2× faster, and an online re-embed ceremony upgrades existing stores without downtime.",
|
||||||
|
"The engine's change feed is documented, with the SSE/WebSocket fan-out pattern for realtime surfaces."
|
||||||
|
],
|
||||||
|
"url": null,
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "11.0.1",
|
||||||
|
"date": "2026-08-31",
|
||||||
|
"headline": "Deletes inside transactions are safe",
|
||||||
|
"items": [
|
||||||
|
"Deleting relations inside a transact() no longer corrupts index bookkeeping.",
|
||||||
|
"A store that deletes its last relation keeps serving instead of refusing."
|
||||||
|
],
|
||||||
|
"url": null,
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "11.0.0",
|
||||||
|
"date": "2026-08-28",
|
||||||
|
"headline": "One install, one engine — Brainy",
|
||||||
|
"items": [
|
||||||
|
"The former two-package pair is one package: the native engine under the familiar API. One import is the whole install.",
|
||||||
|
"A missing native build refuses loudly with its cures named; nothing falls back silently.",
|
||||||
|
"Stores open in place — no migration."
|
||||||
|
],
|
||||||
|
"url": null,
|
||||||
|
"thumb": null
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"history": "The version line continues from the 4.3.x native-engine releases; their record lives in the product repository's CHANGELOG.md."
|
||||||
|
}
|
||||||
110
releases/open-brainy.json
Normal file
110
releases/open-brainy.json
Normal file
|
|
@ -0,0 +1,110 @@
|
||||||
|
{
|
||||||
|
"product": "open-brainy",
|
||||||
|
"entries": [
|
||||||
|
{
|
||||||
|
"version": "10.4.9",
|
||||||
|
"date": "2026-09-02",
|
||||||
|
"headline": "Graph-first finds, honest verb arrays, and opens that stop rescanning history",
|
||||||
|
"items": [
|
||||||
|
"find({ connected, where }) now walks the neighbours first and filters only those rows — correct at every page, and O(neighbours) instead of O(store).",
|
||||||
|
"related() with a list of verb types (or sources, or targets) returns every requested kind — four fast paths silently kept only the first.",
|
||||||
|
"Deferred-embedding recovery resumes from a low-water mark instead of rescanning the whole generation log at every open — measured at two minutes on a large brain, now milliseconds."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.9",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.7",
|
||||||
|
"date": "2026-09-01",
|
||||||
|
"headline": "Count ledgers can no longer race themselves",
|
||||||
|
"items": [
|
||||||
|
"Concurrent count flushes coalesce into one writer with a trailing pass — parallel flushes can no longer corrupt a store's count ledger.",
|
||||||
|
"Atomic writes carry a per-process sequence, so two processes' temp files can never collide."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.7",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.6",
|
||||||
|
"date": "2026-08-31",
|
||||||
|
"headline": "Transactions cross the index seam safely",
|
||||||
|
"items": [
|
||||||
|
"Deleting relations inside a transact() no longer fails against the metadata index — operations take a JSON-safe view at the moment they execute.",
|
||||||
|
"Fixes a class of transaction failures on stores with integer-mapped relation endpoints."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.6",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.5",
|
||||||
|
"date": "2026-08-31",
|
||||||
|
"headline": "Recovery tells the truth, docs live at home",
|
||||||
|
"items": [
|
||||||
|
"A torn generation-log tail is a terminal verdict with a named cure — never an endless wait at open.",
|
||||||
|
"A sealed segment declares only the generations it actually holds.",
|
||||||
|
"The engine's documentation now publishes from its own repository."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.5",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.4",
|
||||||
|
"date": "2026-08-28",
|
||||||
|
"headline": "Faster opens, quieter idle",
|
||||||
|
"items": [
|
||||||
|
"Opening a store discovers generations from directory names instead of walking the log, and answers \"any entities?\" with one directory read.",
|
||||||
|
"The flush-request watch is event-driven; idle stores stop paying a polling heartbeat.",
|
||||||
|
"A slow open now names the exact step it is in, so operators see what is being paid and why."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.4",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.3",
|
||||||
|
"date": "2026-08-27",
|
||||||
|
"headline": "Open Brainy, under its own name",
|
||||||
|
"items": [
|
||||||
|
"The same engine as 10.4.2, now published as @soulcraftlabs/brainy — the MIT reference engine, on The Source.",
|
||||||
|
"No code changes; your imports change once and everything else stays put."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.3",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.2",
|
||||||
|
"date": "2026-08-27",
|
||||||
|
"headline": "Vectors that lie are refused, counts that drift are caught",
|
||||||
|
"items": [
|
||||||
|
"A zero-norm vector is not a vector: the index refuses them, rebuilds skip them, and a sanctioned unvector door removes them cleanly.",
|
||||||
|
"The canonical count ledger derives from identity records and marks legacy-derived ledgers suspect at load.",
|
||||||
|
"Plugin activation failures keep their original error as cause, so the real frame reaches your logs."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.2",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.1",
|
||||||
|
"date": "2026-08-26",
|
||||||
|
"headline": "Writes that change nothing cost nothing",
|
||||||
|
"items": [
|
||||||
|
"The read gate is per index family, and a write carrying unchanged data never re-embeds.",
|
||||||
|
"The vectored-row count joins the ledger, so vector coverage is a number you can read, not a guess."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.1",
|
||||||
|
"thumb": null
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"version": "10.4.0",
|
||||||
|
"date": "2026-08-26",
|
||||||
|
"headline": "Repair routing, the vector ledger, and honest empties",
|
||||||
|
"items": [
|
||||||
|
"Repairs route to the index that owns the damage, and the open gate closes the vector leg until coverage is proven.",
|
||||||
|
"An empty string is real data, not a missing field.",
|
||||||
|
"The metadata crossing never carries raw integer relation endpoints — a whole class of serialization faults closed."
|
||||||
|
],
|
||||||
|
"url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.0",
|
||||||
|
"thumb": null
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"history": "Earlier releases are recorded in CHANGELOG.md in this repository."
|
||||||
|
}
|
||||||
|
|
@ -237,7 +237,7 @@ fi
|
||||||
# and RELEASES.md are the record; this just gives The Source's UI a release page).
|
# 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}"
|
echo -e "${BLUE}🔟 Creating release page on The Source...${NC}"
|
||||||
if [ -n "${FORGEJO_RELEASE_TOKEN:-}" ]; then
|
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" \
|
-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
|
-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"
|
echo -e "${GREEN}✅ Release page created on The Source${NC}\n"
|
||||||
|
|
|
||||||
|
|
@ -15,7 +15,6 @@ import { JsHnswVectorIndex } from './hnsw/hnswIndex.js'
|
||||||
import { createStorage, resolveFilesystemRoot } from './storage/storageFactory.js'
|
import { createStorage, resolveFilesystemRoot } from './storage/storageFactory.js'
|
||||||
import type { StorageOptions } from './storage/storageFactory.js'
|
import type { StorageOptions } from './storage/storageFactory.js'
|
||||||
import { rebuildCounts } from './utils/rebuildCounts.js'
|
import { rebuildCounts } from './utils/rebuildCounts.js'
|
||||||
import { jsonSafeIndexMetadata } from './utils/jsonSafeIndexMetadata.js'
|
|
||||||
import type { MetadataWriteBuffer } from './utils/metadataWriteBuffer.js'
|
import type { MetadataWriteBuffer } from './utils/metadataWriteBuffer.js'
|
||||||
import { BaseStorage } from './storage/baseStorage.js'
|
import { BaseStorage } from './storage/baseStorage.js'
|
||||||
import {
|
import {
|
||||||
|
|
@ -4204,19 +4203,32 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
*/
|
*/
|
||||||
/**
|
/**
|
||||||
* @description A JSON-safe view of a record bound for the metadata-index
|
* @description A JSON-safe view of a record bound for the metadata-index
|
||||||
* crossing — delegates to the shared {@link jsonSafeIndexMetadata} leaf,
|
* crossing. The seam's metadata is JSON-safe BY CONTRACT (a native provider
|
||||||
* which the metadata-index transaction operations ALSO apply at execute
|
* serializes it; u64 ints as Number corrupt above 2^53) — but
|
||||||
* and rollback time. This plan-time wrap alone proved insufficient: it
|
* {@link resolveVerbEndpointInts} MIRRORS the resolved endpoint ints onto
|
||||||
* returns the same reference when the record is clean, and `transact()`'s
|
* the verb object itself as BigInt (`verb.sourceInt`/`targetInt`), so a
|
||||||
* delete legs share that reference with a graph-retraction op whose
|
* verb object reused as index metadata carried BigInts into
|
||||||
* execute-time endpoint resolution mirrors BigInt ints onto it (the full
|
* JSON.stringify, which throws, aborting the whole transaction (found by
|
||||||
* aliasing story lives on the leaf module's doc).
|
* the first joint pair gate). Endpoint ints ride their OWN op params on the
|
||||||
|
* graph legs — the metadata crossing drops every BigInt-valued top-level
|
||||||
|
* key instead of guessing at a lossy numeric encoding.
|
||||||
* @param metadata - The candidate index-metadata record.
|
* @param metadata - The candidate index-metadata record.
|
||||||
* @returns The same object when already JSON-safe, else a shallow copy
|
* @returns The same object when already JSON-safe, else a shallow copy
|
||||||
* without the BigInt-valued keys.
|
* without the BigInt-valued keys.
|
||||||
*/
|
*/
|
||||||
private static jsonSafeIndexMetadata(metadata: unknown): unknown {
|
private static jsonSafeIndexMetadata(metadata: unknown): unknown {
|
||||||
return jsonSafeIndexMetadata(metadata)
|
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
|
||||||
}
|
}
|
||||||
|
|
||||||
private metadataIndexRetractionOp(
|
private metadataIndexRetractionOp(
|
||||||
|
|
|
||||||
|
|
@ -1089,10 +1089,6 @@ export abstract class BaseStorageAdapter implements StorageAdapter {
|
||||||
|
|
||||||
// Counts changed since the last persist? Drives the write-through flush.
|
// Counts changed since the last persist? Drives the write-through flush.
|
||||||
protected pendingCountPersist = false
|
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
|
* Get total noun count - O(1) operation
|
||||||
|
|
@ -1345,46 +1341,15 @@ export abstract class BaseStorageAdapter implements StorageAdapter {
|
||||||
return
|
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
|
|
||||||
}
|
|
||||||
|
|
||||||
this.countPersistInFlight = (async () => {
|
|
||||||
try {
|
try {
|
||||||
// Persist to storage (implemented by subclass)
|
// Persist to storage (implemented by subclass)
|
||||||
this.pendingCountPersist = false
|
|
||||||
await this.persistCounts()
|
await this.persistCounts()
|
||||||
|
this.pendingCountPersist = false
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
// Keep the flag set so the next operation retries.
|
|
||||||
this.pendingCountPersist = true
|
|
||||||
console.error('CRITICAL: Failed to flush counts to storage:', error)
|
console.error('CRITICAL: Failed to flush counts to storage:', error)
|
||||||
|
// Keep pending flag set so we retry on next operation
|
||||||
throw 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
|
* Atomic write via temp-file-then-rename so concurrent readers never see a
|
||||||
* half-written lock JSON. Reused by writer-lock writes + heartbeat.
|
* 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> {
|
private async writeFileAtomic(filePath: string, contents: string): Promise<void> {
|
||||||
// pid + timestamp alone collided: two writers of the same target inside
|
const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}`
|
||||||
// one millisecond shared this path, and the loser's rename found the
|
|
||||||
// winner had already moved it (ENOENT). The sequence makes every call's
|
|
||||||
// temp path its own.
|
|
||||||
const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}-${++FileSystemStorage.atomicWriteSeq}`
|
|
||||||
await fs.promises.writeFile(tmp, contents)
|
await fs.promises.writeFile(tmp, contents)
|
||||||
await fs.promises.rename(tmp, filePath)
|
await fs.promises.rename(tmp, filePath)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2942,33 +2942,19 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
||||||
!options.filter.service &&
|
!options.filter.service &&
|
||||||
!options.filter.metadata
|
!options.filter.metadata
|
||||||
) {
|
) {
|
||||||
const sourceIds = Array.isArray(options.filter.sourceId)
|
const sourceId = Array.isArray(options.filter.sourceId)
|
||||||
? options.filter.sourceId
|
? options.filter.sourceId[0]
|
||||||
: [options.filter.sourceId]
|
: options.filter.sourceId
|
||||||
|
|
||||||
// EVERY requested verb type is honoured — an array used to collapse to
|
const verbType = Array.isArray(options.filter.verbType)
|
||||||
// its first element here, silently dropping the rest of the ask.
|
? options.filter.verbType[0]
|
||||||
const verbTypes = new Set(
|
: options.filter.verbType
|
||||||
Array.isArray(options.filter.verbType)
|
|
||||||
? options.filter.verbType
|
|
||||||
: [options.filter.verbType]
|
|
||||||
)
|
|
||||||
|
|
||||||
// Get verbs by source (union over every requested source), filter by the
|
// Get verbs by source, then filter by type (O(1) graph lookup + O(n) type filter),
|
||||||
// requested type SET (O(1) graph lookup + O(n) type filter), then apply
|
// then apply the subtype / visibility metadata filters on the candidate set.
|
||||||
// the subtype / visibility metadata filters on the candidate set.
|
const verbsBySource = await this.getVerbsBySource_internal(sourceId)
|
||||||
const bySource: HNSWVerbWithMetadata[] = []
|
|
||||||
const seenVerbIds = new Set<string>()
|
|
||||||
for (const oneSource of sourceIds) {
|
|
||||||
for (const v of await this.getVerbsBySource_internal(oneSource)) {
|
|
||||||
if (!seenVerbIds.has(v.id)) {
|
|
||||||
seenVerbIds.add(v.id)
|
|
||||||
bySource.push(v)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
const filteredVerbs = this.applyVerbMetadataFilters(
|
const filteredVerbs = this.applyVerbMetadataFilters(
|
||||||
bySource.filter(v => verbTypes.has(v.verb)),
|
verbsBySource.filter(v => v.verb === verbType),
|
||||||
options.filter
|
options.filter
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -2999,22 +2985,16 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
||||||
!options.filter.service &&
|
!options.filter.service &&
|
||||||
!options.filter.metadata
|
!options.filter.metadata
|
||||||
) {
|
) {
|
||||||
// EVERY requested source is honoured — an array used to collapse to
|
const sourceId = Array.isArray(options.filter.sourceId)
|
||||||
// its first element here, silently dropping the rest of the ask.
|
? options.filter.sourceId[0]
|
||||||
const onlySourceIds = Array.isArray(options.filter.sourceId)
|
: options.filter.sourceId
|
||||||
? options.filter.sourceId
|
|
||||||
: [options.filter.sourceId]
|
// Get verbs by source directly (hydrated with metadata), then apply the
|
||||||
const sourceUnion: HNSWVerbWithMetadata[] = []
|
// subtype / visibility metadata filters on the O(degree) candidate set.
|
||||||
const seenSourceVerbIds = new Set<string>()
|
const verbsBySource = this.applyVerbMetadataFilters(
|
||||||
for (const oneSource of onlySourceIds) {
|
await this.getVerbsBySource_internal(sourceId),
|
||||||
for (const v of await this.getVerbsBySource_internal(oneSource)) {
|
options.filter
|
||||||
if (!seenSourceVerbIds.has(v.id)) {
|
)
|
||||||
seenSourceVerbIds.add(v.id)
|
|
||||||
sourceUnion.push(v)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
const verbsBySource = this.applyVerbMetadataFilters(sourceUnion, options.filter)
|
|
||||||
|
|
||||||
// Apply pagination
|
// Apply pagination
|
||||||
const paginatedVerbs = verbsBySource.slice(offset, offset + limit)
|
const paginatedVerbs = verbsBySource.slice(offset, offset + limit)
|
||||||
|
|
@ -3043,22 +3023,16 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
||||||
!options.filter.service &&
|
!options.filter.service &&
|
||||||
!options.filter.metadata
|
!options.filter.metadata
|
||||||
) {
|
) {
|
||||||
// EVERY requested target is honoured — an array used to collapse to
|
const targetId = Array.isArray(options.filter.targetId)
|
||||||
// its first element here, silently dropping the rest of the ask.
|
? options.filter.targetId[0]
|
||||||
const onlyTargetIds = Array.isArray(options.filter.targetId)
|
: options.filter.targetId
|
||||||
? options.filter.targetId
|
|
||||||
: [options.filter.targetId]
|
// Get verbs by target directly (hydrated with metadata), then apply the
|
||||||
const targetUnion: HNSWVerbWithMetadata[] = []
|
// subtype / visibility metadata filters on the O(degree) candidate set.
|
||||||
const seenTargetVerbIds = new Set<string>()
|
const verbsByTarget = this.applyVerbMetadataFilters(
|
||||||
for (const oneTarget of onlyTargetIds) {
|
await this.getVerbsByTarget_internal(targetId),
|
||||||
for (const v of await this.getVerbsByTarget_internal(oneTarget)) {
|
options.filter
|
||||||
if (!seenTargetVerbIds.has(v.id)) {
|
)
|
||||||
seenTargetVerbIds.add(v.id)
|
|
||||||
targetUnion.push(v)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
const verbsByTarget = this.applyVerbMetadataFilters(targetUnion, options.filter)
|
|
||||||
|
|
||||||
// Apply pagination
|
// Apply pagination
|
||||||
const paginatedVerbs = verbsByTarget.slice(offset, offset + limit)
|
const paginatedVerbs = verbsByTarget.slice(offset, offset + limit)
|
||||||
|
|
@ -3087,25 +3061,16 @@ export abstract class BaseStorage extends BaseStorageAdapter {
|
||||||
!options.filter.service &&
|
!options.filter.service &&
|
||||||
!options.filter.metadata
|
!options.filter.metadata
|
||||||
) {
|
) {
|
||||||
// EVERY requested verb type is honoured — an array used to collapse to
|
const verbType = Array.isArray(options.filter.verbType)
|
||||||
// its first element here, silently dropping the rest of the ask.
|
? options.filter.verbType[0]
|
||||||
const verbTypes = Array.isArray(options.filter.verbType)
|
: options.filter.verbType
|
||||||
? options.filter.verbType
|
|
||||||
: [options.filter.verbType]
|
|
||||||
|
|
||||||
// Get verbs by each requested type (hydrated with metadata), deduped by
|
// Get verbs by type directly (hydrated with metadata), then apply the
|
||||||
// id, then apply the subtype / visibility metadata filters on the set.
|
// subtype / visibility metadata filters on the candidate set.
|
||||||
const byType: HNSWVerbWithMetadata[] = []
|
const verbsByType = this.applyVerbMetadataFilters(
|
||||||
const seenTypeVerbIds = new Set<string>()
|
await this.getVerbsByType_internal(verbType),
|
||||||
for (const oneType of verbTypes) {
|
options.filter
|
||||||
for (const v of await this.getVerbsByType_internal(oneType)) {
|
)
|
||||||
if (!seenTypeVerbIds.has(v.id)) {
|
|
||||||
seenTypeVerbIds.add(v.id)
|
|
||||||
byType.push(v)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
const verbsByType = this.applyVerbMetadataFilters(byType, options.filter)
|
|
||||||
|
|
||||||
// Apply pagination
|
// Apply pagination
|
||||||
const paginatedVerbs = verbsByType.slice(offset, offset + limit)
|
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 { GraphVerb } from '../../coreTypes.js'
|
||||||
import type { Operation, RollbackAction } from '../types.js'
|
import type { Operation, RollbackAction } from '../types.js'
|
||||||
import { isZeroNormVector } from '../../utils/distance.js'
|
import { isZeroNormVector } from '../../utils/distance.js'
|
||||||
import { jsonSafeIndexMetadata } from '../../utils/jsonSafeIndexMetadata.js'
|
|
||||||
import { prodLog } from '../../utils/logger.js'
|
import { prodLog } from '../../utils/logger.js'
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
@ -391,21 +390,13 @@ export class AddToMetadataIndexOperation implements Operation {
|
||||||
// rollback so add + undo reference the same watermark.
|
// rollback so add + undo reference the same watermark.
|
||||||
const generation = this.generationFn?.()
|
const generation = this.generationFn?.()
|
||||||
|
|
||||||
// The JSON-safe view is taken HERE, per crossing, never at construction:
|
// Add to metadata index (skipFlush=true for transaction atomicity)
|
||||||
// the entity reference this op holds can be mutated between plan and
|
await this.index.addToIndex(this.id, this.entity, true, false, generation)
|
||||||
// execute (a graph op's execute-time endpoint-int resolution mirrors
|
|
||||||
// BigInts onto a shared verb object) — see jsonSafeIndexMetadata's
|
|
||||||
// module doc.
|
|
||||||
await this.index.addToIndex(
|
|
||||||
this.id, jsonSafeIndexMetadata(this.entity), true, false, generation
|
|
||||||
)
|
|
||||||
|
|
||||||
// Return rollback action
|
// Return rollback action
|
||||||
return async () => {
|
return async () => {
|
||||||
// Remove from metadata index
|
// Remove from metadata index
|
||||||
await this.index.removeFromIndex(
|
await this.index.removeFromIndex(this.id, this.entity, generation)
|
||||||
this.id, jsonSafeIndexMetadata(this.entity), generation
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -441,21 +432,13 @@ export class RemoveFromMetadataIndexOperation implements Operation {
|
||||||
// Resolve the removal generation once; reuse it for the rollback re-add.
|
// Resolve the removal generation once; reuse it for the rollback re-add.
|
||||||
const generation = this.generationFn?.()
|
const generation = this.generationFn?.()
|
||||||
|
|
||||||
// Sanitized per crossing, never at construction — transact()'s delete
|
// Remove from metadata index
|
||||||
// legs hand this op the SAME verb object the graph-retraction op's
|
await this.index.removeFromIndex(this.id, this.entity, generation)
|
||||||
// execute-time endpoint resolution mutates (BigInt sourceInt/targetInt),
|
|
||||||
// so a plan-time view aliases the pollution. See jsonSafeIndexMetadata's
|
|
||||||
// module doc.
|
|
||||||
await this.index.removeFromIndex(
|
|
||||||
this.id, jsonSafeIndexMetadata(this.entity), generation
|
|
||||||
)
|
|
||||||
|
|
||||||
// Return rollback action
|
// Return rollback action
|
||||||
return async () => {
|
return async () => {
|
||||||
// Re-add with original metadata (skipFlush=true)
|
// Re-add with original metadata (skipFlush=true)
|
||||||
await this.index.addToIndex(
|
await this.index.addToIndex(this.id, this.entity, true, false, generation)
|
||||||
this.id, jsonSafeIndexMetadata(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
|
|
||||||
}
|
|
||||||
|
|
@ -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,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([[], []])
|
|
||||||
})
|
|
||||||
})
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue