diff --git a/CHANGELOG.md b/CHANGELOG.md index 7154d5a2..a54d609e 100644 --- a/CHANGELOG.md +++ b/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. -### [10.4.6](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.5...v10.4.6) (2026-08-31) - -- fix(transact): metadata-index ops take their JSON-safe view at the crossing, not at construction (73500e7d) - - -### [10.4.5](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.4...v10.4.5) (2026-08-31) - -- build(release): the docs-push step retires — this engine documents itself in its own repository (d6bcb14f) -- fix(generations): a sealed segment may only declare the generations it holds (a963a744) -- fix(recovery): a torn generation-log tail is a terminal verdict, never a wait (c9930871) - - ### [10.4.4](https://source.soulcraft.com/soulcraftlabs/open-brainy/compare/v10.4.3...v10.4.4) (2026-08-28) - fix(vfs): the old-root sweep narrates only when it has something to say (d49148e1) diff --git a/package-lock.json b/package-lock.json index 9e573da3..c4f66561 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@soulcraftlabs/brainy", - "version": "10.4.6", + "version": "10.4.4", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@soulcraftlabs/brainy", - "version": "10.4.6", + "version": "10.4.4", "license": "MIT", "dependencies": { "@msgpack/msgpack": "^3.1.2", diff --git a/package.json b/package.json index 51322998..06ce0253 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@soulcraftlabs/brainy", - "version": "10.4.6", + "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", diff --git a/releases/brainy.json b/releases/brainy.json new file mode 100644 index 00000000..8f61c7f2 --- /dev/null +++ b/releases/brainy.json @@ -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." +} diff --git a/releases/open-brainy.json b/releases/open-brainy.json new file mode 100644 index 00000000..582f4847 --- /dev/null +++ b/releases/open-brainy.json @@ -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." +} diff --git a/scripts/release.sh b/scripts/release.sh index 1a4fe575..5d434320 100755 --- a/scripts/release.sh +++ b/scripts/release.sh @@ -237,7 +237,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" diff --git a/src/brainy.ts b/src/brainy.ts index 464f689b..da04577e 100644 --- a/src/brainy.ts +++ b/src/brainy.ts @@ -15,7 +15,6 @@ import { JsHnswVectorIndex } from './hnsw/hnswIndex.js' import { createStorage, resolveFilesystemRoot } from './storage/storageFactory.js' import type { StorageOptions } from './storage/storageFactory.js' import { rebuildCounts } from './utils/rebuildCounts.js' -import { jsonSafeIndexMetadata } from './utils/jsonSafeIndexMetadata.js' import type { MetadataWriteBuffer } from './utils/metadataWriteBuffer.js' import { BaseStorage } from './storage/baseStorage.js' import { @@ -1820,31 +1819,31 @@ export class Brainy implements BrainyInterface { // a deferred write's ack and its background embed DELAYED a vector; // this is where it lands. if (!this.isReadOnly) { - // BEHIND THE DOORS (the open pays nothing here): the bridge + the - // recovery fold run as one latched background task; the embed worker - // starts when it settles. A pending embed's outcome was always - // eventual — moving its recovery off the open's foreground changes - // when the worker starts, never whether a marker is honored. - // awaitPendingEmbeds() and close() wait on the latch first. - this._pendingEmbedRecovery = (async () => { - try { - await this.bridgeLegacyPendingEmbedSidecars() - await this.recoverPendingEmbedsFromLog() - if (this._pendingEmbedIds.size > 0) { - prodLog.info( - `[Brainy] ${this._pendingEmbedIds.size} deferred embed(s) pending from a previous ` + - `session — resuming in the background` - ) - const t = setTimeout(() => this.kickEmbedWorker(), 0) - ;(t as { unref?: () => void }).unref?.() - } - } catch (err) { - prodLog.warn( - `[Brainy] pending-embed recovery failed: ${(err as Error).message} — ` + - `the log's markers remain durable; recovery retries next open` + try { + await step( + 'bridge-pending-embed-sidecars', + 'migrating any pre-log deferred-embed marker files into the generation log', + () => this.bridgeLegacyPendingEmbedSidecars() + ) + await step( + 'recover-pending-embeds', + 'folding the generation log\'s deferred-embed markers back into the pending set', + () => this.recoverPendingEmbedsFromLog() + ) + if (this._pendingEmbedIds.size > 0) { + prodLog.info( + `[Brainy] ${this._pendingEmbedIds.size} deferred embed(s) pending from a previous ` + + `session — resuming in the background` ) + const t = setTimeout(() => this.kickEmbedWorker(), 0) + ;(t as { unref?: () => void }).unref?.() } - })() + } catch (err) { + prodLog.warn( + `[Brainy] pending-embed recovery failed: ${(err as Error).message} — ` + + `the log's markers remain durable; recovery retries next open` + ) + } } // PHASE 4 of 5 — "VFS bootstrap": shutdown-hook registration, blob @@ -2408,19 +2407,6 @@ export class Brainy implements BrainyInterface { */ private static readonly PENDING_EMBED_PREFIX = '_system/pending_embeds/' - /** - * Storage-root-relative path of the ADVISORY pending-embed low-water mark: - * `{ generation, writtenAt }`, written whenever the pending set drains to - * empty (and at clean close when empty). Every marker in facts at or below - * `generation` is consumed, so recovery scans from `generation + 1`. The - * mark is advisory and monotone-safe: stale-low costs a longer scan, never - * a lost marker; it is never required for correctness. - */ - private static readonly PENDING_EMBED_LOWWATER_PATH = '_system/pending_embeds_lowwater.json' - - /** Resolves when the background pending-embed recovery fold has settled (open arms it). */ - private _pendingEmbedRecovery: Promise | null = null - /** * @description Mark a deferred embed pending (MT5): the id joins the * in-memory fast-path set and the returned `embed.pending` record is @@ -2448,40 +2434,6 @@ export class Brainy implements BrainyInterface { */ private clearPendingEmbed(id: string): void { this._pendingEmbedIds.delete(id) - if (this._pendingEmbedIds.size === 0) this.maybeWriteEmbedLowWater() - } - - /** - * @description Advance the advisory low-water mark: called at drain-to-empty - * (and at clean close when empty), it records the fact log's CURRENT head — - * with the set empty, every marker at or below the head has been consumed, - * so the next open's recovery fold scans only what comes after. Fire-and- - * forget at the drain (close() awaits the core); loud on failure: a missed - * write costs the next open a longer scan, never a marker. No-op without a - * fact log (no durable markers exist there) and on read-only opens. - */ - private maybeWriteEmbedLowWater(): void { - void this.writeEmbedLowWater() - } - - /** The awaitable core of {@link maybeWriteEmbedLowWater} — close() awaits it. */ - private async writeEmbedLowWater(): Promise { - if (this.isReadOnly) return - const log = this.generationStore ? this.generationStore.getFactLog() : null - if (!log) return - const generation = log.headGeneration() - if (!(generation > 0)) return - try { - await this.storage.writeRawObject(Brainy.PENDING_EMBED_LOWWATER_PATH, { - generation, - writtenAt: Date.now() - }) - } catch (err) { - prodLog.warn( - `[Brainy] pending-embed low-water write failed at generation ${generation}: ` + - `${(err as Error).message} — the next open scans from the previous mark` - ) - } } /** @@ -2492,14 +2444,9 @@ export class Brainy implements BrainyInterface { * survives the fold is exactly the set of acknowledged deferred writes * whose vectors have not landed. * - * BOUND: the scan starts at the advisory low-water mark - * ({@link Brainy.PENDING_EMBED_LOWWATER_PATH}) — the log head at which the - * pending set last drained to empty — so a settled brain reads only the - * facts since then, not its whole history. Without a mark (first open - * after upgrade) it scans from generation 1, once; a stale-low mark costs - * a longer scan, never a marker. The fold runs BEHIND the doors (open - * arms it as a background task and the embed worker starts when it - * settles); {@link awaitPendingEmbeds} and close() wait for it first. + * BOUND (honest): no durable low-water mark exists for the earliest + * unconsumed pending, so the fold scans the log's committed facts from + * generation 1 — a sequential read of the log at open, O(log bytes). * It is SKIPPED WHOLESALE when the log has never had a v2 tail * ({@link FactLog.hasV2History} — v1 facts cannot carry marker records), * so pre-cutover brains pay nothing; on a mixed log the scan still reads @@ -2512,18 +2459,7 @@ export class Brainy implements BrainyInterface { private async recoverPendingEmbedsFromLog(): Promise { const log = this.generationStore.getFactLog() if (!log || !log.hasV2History()) return - let fromGeneration = 1 - try { - const mark = (await this.storage.readRawObject(Brainy.PENDING_EMBED_LOWWATER_PATH)) as { - generation?: number - } | null - if (mark && typeof mark.generation === 'number' && mark.generation > 0) { - fromGeneration = mark.generation + 1 - } - } catch { - // No mark (or unreadable): scan from 1 — correctness over cost. - } - const scan = log.scanFacts({ fromGeneration }) + const scan = log.scanFacts({ fromGeneration: 1 }) for await (const batch of scan.batches()) { for (const fact of batch.facts) { for (const record of fact.records ?? []) { @@ -2710,7 +2646,6 @@ export class Brainy implements BrainyInterface { * before I proceed" callers use this; nothing else ever needs to wait. */ public async awaitPendingEmbeds(): Promise { - if (this._pendingEmbedRecovery) await this._pendingEmbedRecovery while (this._pendingEmbedIds.size > 0 || this._embedWorkerFlight) { this.kickEmbedWorker() await (this._embedWorkerFlight ?? Promise.resolve()) @@ -4268,19 +4203,32 @@ export class Brainy implements BrainyInterface { */ /** * @description A JSON-safe view of a record bound for the metadata-index - * crossing — delegates to the shared {@link jsonSafeIndexMetadata} leaf, - * which the metadata-index transaction operations ALSO apply at execute - * and rollback time. This plan-time wrap alone proved insufficient: it - * returns the same reference when the record is clean, and `transact()`'s - * delete legs share that reference with a graph-retraction op whose - * execute-time endpoint resolution mirrors BigInt ints onto it (the full - * aliasing story lives on the leaf module's doc). + * crossing. The seam's metadata is JSON-safe BY CONTRACT (a native provider + * serializes it; u64 ints as Number corrupt above 2^53) — but + * {@link resolveVerbEndpointInts} MIRRORS the resolved endpoint ints onto + * the verb object itself as BigInt (`verb.sourceInt`/`targetInt`), so a + * verb object reused as index metadata carried BigInts into + * JSON.stringify, which throws, aborting the whole transaction (found by + * the first joint pair gate). Endpoint ints ride their OWN op params on the + * graph legs — the metadata crossing drops every BigInt-valued top-level + * key instead of guessing at a lossy numeric encoding. * @param metadata - The candidate index-metadata record. * @returns The same object when already JSON-safe, else a shallow copy * without the BigInt-valued keys. */ private static jsonSafeIndexMetadata(metadata: unknown): unknown { - return jsonSafeIndexMetadata(metadata) + if (metadata === null || typeof metadata !== 'object') return metadata + const rec = metadata as Record + let hasBigint = false + for (const k in rec) { + if (typeof rec[k] === 'bigint') { hasBigint = true; break } + } + if (!hasBigint) return metadata + const out: Record = {} + for (const k in rec) { + if (typeof rec[k] !== 'bigint') out[k] = rec[k] + } + return out } private metadataIndexRetractionOp( @@ -7344,47 +7292,6 @@ export class Brainy implements BrainyInterface { await this.verifyMetadataLive() } - // PLANNED FIND (optional provider door, `MetadataIndexProvider.planFindPage`). - // - // The stage doors below 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 and materializes ids only for the page. - // - // The hook sits ABOVE the branch selection because the branches are what - // decide stage order per call site; an index that plans has to be asked - // before that choice is made, not inside one of its arms. - // - // Optional and additive: a provider without the door, and any shape the - // door hands back, take exactly the path they always took. `null` is a - // routing decision the door must make BEFORE doing any work — never a - // partial answer. Every guard above still ran (readiness, the migration - // gate, the where-clause validation, the metadata cold-read guard), and - // the serving law is applied here on the way out: an empty answer is - // re-verified against the index that produced it before it is believed. - const planningIndex = this.metadataIndex as unknown as MetadataIndexProvider - if (typeof planningIndex.planFindPage === 'function') { - const planned = await planningIndex.planFindPage(params, [...hiddenIds], this.graphIndex) - if (planned !== null && planned !== undefined) { - if (planned.ids.length === 0) { - // A cold adjacency can report a size yet hold no edges, so an empty - // graph answer is not truth until the adjacency verifies live. A - // genuinely edgeless anchor verifies and the empty result stands. - if (planned.emptyAt === 'graph') await this.verifyGraphAdjacencyLive() - return [] - } - const plannedEntities = await this.batchGet(planned.ids) - const plannedResults: Result[] = [] - for (const id of planned.ids) { - const entity = plannedEntities.get(id) - if (entity) plannedResults.push(this.createResult(id, 1.0, entity)) - } - return plannedResults - } - } - // Handle metadata-only queries (no vector search needed) if (!hasVectorSearchCriteria && !hasGraphCriteria && hasFilterCriteria) { // Build filter for metadata index @@ -7575,37 +7482,7 @@ export class Brainy implements BrainyInterface { // JS path — there the materialized `candidateIds` restricts the walk instead. let preResolvedAllowedIds: OpaqueIdSet | undefined - // Graph-first law (10.4.8, BRAINY-PROD-LATENCY-TRIAD rounds 44/45): with - // `connected` present the NEIGHBOUR SET is the candidate universe. It is - // resolved first from the adjacency (O(neighbours)), the metadata filter - // is evaluated over those ids only, and paging happens LAST. The earlier - // order materialized the whole-store filtered id list, paged it, hydrated - // the page, and only then intersected with the neighbours — O(store) per - // call, and a neighbour outside the first page was silently dropped. - let graphFirstIds: string[] | null = null - if (hasGraphCriteria) { - graphFirstIds = await this.resolveConnectedIds(params) - if (hiddenIds.size > 0) { - graphFirstIds = graphFirstIds.filter((id) => !hiddenIds.has(id)) - } - if ( - graphFirstIds.length > 0 && - (params.where || params.type || params.subtype || params.service || params.excludeVFS) - ) { - preResolvedFilter = this.buildMetadataFilter(params) - graphFirstIds = await this.filterIdsWithinBelted(preResolvedFilter, graphFirstIds) - } - if (graphFirstIds.length === 0) { - return [] - } - if (!hasVectorSearchCriteria) { - return await this.pageConnectedIds(params, graphFirstIds) - } - // The vector leg walks ONLY the neighbours (its candidate walk). The - // filter is already applied above, so no opaque universe is produced — - // it would describe the whole store, not the neighbour set. - preResolvedMetadataIds = graphFirstIds - } else if (params.where || params.type || params.subtype || params.service || params.excludeVFS) { + if (params.where || params.type || params.subtype || params.service || params.excludeVFS) { preResolvedFilter = this.buildMetadataFilter(params) preResolvedMetadataIds = await this.filterIdsBelted(preResolvedFilter) @@ -7794,11 +7671,9 @@ export class Brainy implements BrainyInterface { } } - // The text leg of a hybrid find has no candidate door, so its hits are - // held to the neighbour set here; the vector leg walked only the neighbours. - if (graphFirstIds !== null && results.length > 0) { - const neighbourSet = new Set(graphFirstIds) - results = results.filter((r) => neighbourSet.has(r.id)) + // Graph search component with O(1) traversal + if (params.connected) { + results = await this.executeGraphSearch(params, results) } // Apply fusion scoring if requested @@ -12913,29 +12788,6 @@ export class Brainy implements BrainyInterface { } } - /** - * The id-scoped twin of {@link filterIdsBelted}: evaluate `filter` over `ids` - * only, through the provider's own evaluation so the answer can never drift - * from `getIdsForFilter`'s. A provider without the door is served by its - * whole-store answer intersected here (the reference index implements the - * door itself). Same belt: field refusals cross as `BrainyFieldRefusal`. - */ - private async filterIdsWithinBelted(filter: unknown, ids: readonly string[]): Promise { - this.ensureIndexesLoaded(['metadata']) - const mip = this.metadataIndex as unknown as MetadataIndexProvider - try { - if (typeof mip.filterIdsWithin === 'function') { - return await mip.filterIdsWithin(filter, ids) - } - const matched = new Set(await this.metadataIndex.getIdsForFilter(filter)) - return ids.filter((id) => matched.has(id)) - } catch (err) { - const normalized = asBrainyFieldRefusal(err) - if (normalized) throw normalized - throw err - } - } - async getIndexStatus(): Promise<{ initialized: boolean /** `true` once open()'s index-build-if-needed step has run. Named for API @@ -15919,16 +15771,16 @@ export class Brainy implements BrainyInterface { } /** - * Resolve `params.connected` to the neighbour id set — the graph-first - * find's candidate universe (deterministic traversal order, anchors excluded). + * Execute graph search component. * * Honors the full `GraphConstraints` contract: multi-hop `depth` (breadth-first via - * `neighbors()`), `via`/`type` verb-type filtering, and `direction`. An empty set - * is re-verified against the adjacency before it is believed — a not-serving - * adjacency throws rather than answering `[]` as truth. + * `neighbors()`), `via`/`type` verb-type filtering, and `direction`. Previously this read + * only `from`/`to`/`direction` and did a single 1-hop `getNeighbors()`, so `depth` and `via` + * were silently ignored — `find({ connected: { from, depth: 3 } })` returned only the + * immediate neighbour at every depth. */ - private async resolveConnectedIds(params: FindParams): Promise { - if (!params.connected) return [] + private async executeGraphSearch(params: FindParams, existingResults: Result[]): Promise[]> { + if (!params.connected) return existingResults const { from, to, depth, direction = 'both' } = params.connected const via = params.connected.via ?? params.connected.type @@ -15982,8 +15834,8 @@ export class Brainy implements BrainyInterface { if (anchorInt === undefined) return new Set() // unmapped → no relations const verbTypeIndex = TypeUtils.getVerbIndex(via as VerbType) - // No limit: match the JS BFS exactly — the page is cut downstream, - // after the metadata filter, by pageConnectedIds / the candidate walk. + // No limit: match the JS BFS exactly — overall result limiting happens + // downstream against existingResults. const reachedInts = await provider.findConnectedSubtype( anchorInt, verbTypeIndex, subtypeArr[0], effectiveDepth, null ) @@ -16068,44 +15920,22 @@ export class Brainy implements BrainyInterface { await this.verifyGraphAdjacencyLive() } - return [...connectedIds] - } - - /** - * Page and hydrate an already-filtered neighbour set — the pure graph (and - * graph + metadata) find's tail. `orderBy` sorts the WHOLE set by field value - * before the page is cut (never the page after), null values last on `asc` - * and first on `desc`; without `orderBy` the traversal order stands. - */ - private async pageConnectedIds(params: FindParams, ids: string[]): Promise[]> { - const limit = params.limit || 10 - const offset = params.offset || 0 - let ordered = ids - if (params.orderBy) { - const field = params.orderBy - const asc = (params.order || 'asc') === 'asc' - const valued = await Promise.all( - ids.map(async (id) => ({ id, value: await this.metadataIndex.getFieldValueForEntity(id, field) })) - ) - valued.sort((a, b) => { - if (a.value == null && b.value == null) return 0 - if (a.value == null) return asc ? 1 : -1 - if (b.value == null) return asc ? -1 : 1 - if (a.value === b.value) return 0 - const comparison = a.value < b.value ? -1 : 1 - return asc ? comparison : -comparison - }) - ordered = valued.map((v) => v.id) + // Filter existing results to only connected entities + if (existingResults.length > 0) { + return existingResults.filter(r => connectedIds.has(r.id)) } - const pageIds = ordered.slice(offset, offset + limit) - const entitiesMap = await this.batchGet(pageIds) + + // Batch-load connected entities for fast cloud-storage performance const results: Result[] = [] - for (const id of pageIds) { + const ids = [...connectedIds] + const entitiesMap = await this.batchGet(ids) + for (const id of ids) { const entity = entitiesMap.get(id) if (entity) { results.push(this.createResult(id, 1.0, entity)) } } + return results } @@ -19625,19 +19455,6 @@ export class Brainy implements BrainyInterface { * terminal releases have run. */ async close(): Promise { - if (this._pendingEmbedRecovery) { - // Settle the background marker fold before the durable steps — its scan - // is bounded by the low-water mark (a full scan happens at most once, - // on the first open after upgrade). - const settleStart = Date.now() - await this._pendingEmbedRecovery - const settleMs = Date.now() - settleStart - if (settleMs >= 1000) { - prodLog.info(`[Brainy] close: pending-embed recovery settled in ${settleMs}ms`) - } - this._pendingEmbedRecovery = null - } - if (this._pendingEmbedIds.size === 0) await this.writeEmbedLowWater() let closeFailure: unknown = null try { await this.closeDurableSteps() diff --git a/src/neural/embeddedTypeEmbeddings.ts b/src/neural/embeddedTypeEmbeddings.ts index f4cdd632..5b10116c 100644 --- a/src/neural/embeddedTypeEmbeddings.ts +++ b/src/neural/embeddedTypeEmbeddings.ts @@ -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 diff --git a/src/plugin.ts b/src/plugin.ts index fdad42f0..b1aef8e0 100644 --- a/src/plugin.ts +++ b/src/plugin.ts @@ -411,65 +411,6 @@ export interface MetadataIndexProvider { * @returns The matching id universe as an opaque set. */ getIdSetForFilter?(filter: any): Promise - /** - * @description OPTIONAL: evaluate `filter` over `ids` ONLY and return the - * survivors in the caller's order — the door a graph-first - * `find({ connected, where })` walks. The neighbour set is the universe there, - * so the filter must cost O(|ids|) membership checks, never a whole-store - * materialization. A native index answers from its roaring filter result - * (membership by entity int); the reference index answers from its own - * `getIdsForFilter`, so the two doors can never disagree. Absent → Brainy - * intersects `getIdsForFilter`'s answer with `ids` itself (correct, O(store)). - * @param filter - The same filter shape accepted by `getIdsForFilter`. - * @param ids - The candidate ids (canonical). The answer is a subsequence. - */ - filterIdsWithin?(filter: any, ids: readonly string[]): Promise - /** - * @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; apply BEFORE paging so - * `limit` stays exact. - * @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> getSortedIdsForFilter(filter: any, orderBy: string, order?: 'asc' | 'desc', topK?: number): Promise getFilterValues(field: string): Promise diff --git a/src/storage/adapters/baseStorageAdapter.ts b/src/storage/adapters/baseStorageAdapter.ts index a90adb93..cabe2e30 100644 --- a/src/storage/adapters/baseStorageAdapter.ts +++ b/src/storage/adapters/baseStorageAdapter.ts @@ -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 | null = null - /** The one trailing persist a burst has queued behind the in-flight one. */ - private countPersistTrailing: Promise | null = null /** * Get total noun count - O(1) operation @@ -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--`): 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 } /** diff --git a/src/storage/adapters/fileSystemStorage.ts b/src/storage/adapters/fileSystemStorage.ts index 87b6406f..5ec1d88e 100644 --- a/src/storage/adapters/fileSystemStorage.ts +++ b/src/storage/adapters/fileSystemStorage.ts @@ -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 { - // 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) } diff --git a/src/storage/baseStorage.ts b/src/storage/baseStorage.ts index a1cc2e35..d8bcb780 100644 --- a/src/storage/baseStorage.ts +++ b/src/storage/baseStorage.ts @@ -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() - 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() - 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() - 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() - 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) diff --git a/src/transaction/operations/IndexOperations.ts b/src/transaction/operations/IndexOperations.ts index 0142dc54..1bbbca88 100644 --- a/src/transaction/operations/IndexOperations.ts +++ b/src/transaction/operations/IndexOperations.ts @@ -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) } } } diff --git a/src/utils/jsonSafeIndexMetadata.ts b/src/utils/jsonSafeIndexMetadata.ts deleted file mode 100644 index d3b1be5f..00000000 --- a/src/utils/jsonSafeIndexMetadata.ts +++ /dev/null @@ -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 - let hasBigint = false - for (const k in rec) { - if (typeof rec[k] === 'bigint') { hasBigint = true; break } - } - if (!hasBigint) return metadata - const out: Record = {} - for (const k in rec) { - if (typeof rec[k] !== 'bigint') out[k] = rec[k] - } - return out -} diff --git a/src/utils/metadataIndex.ts b/src/utils/metadataIndex.ts index 0fd312e2..3e0e3d17 100644 --- a/src/utils/metadataIndex.ts +++ b/src/utils/metadataIndex.ts @@ -2575,19 +2575,6 @@ export class MetadataIndexManager implements MetadataIndexProvider { /** Once-per-field flag for the fallback-degradation announcement. */ private static announcedFallbackSorts = new Set() - /** - * Evaluate `filter` over `ids` only — the graph-first find's door (the - * neighbour set filtered by id, never the store filtered and then - * intersected). This index answers from its own `getIdsForFilter`, so the - * two doors cannot disagree; the cost is that of the filter over this - * in-memory index, and the answer keeps the caller's order. - */ - async filterIdsWithin(filter: any, ids: readonly string[]): Promise { - if (ids.length === 0) return [] - const matched = new Set(await this.getIdsForFilter(filter)) - return ids.filter((id) => matched.has(id)) - } - async getSortedIdsForFilter( filter: any, orderBy: string, diff --git a/tests/integration/counts-persist-single-flight.test.ts b/tests/integration/counts-persist-single-flight.test.ts deleted file mode 100644 index 5acbdcc3..00000000 --- a/tests/integration/counts-persist-single-flight.test.ts +++ /dev/null @@ -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--`). Two persists inside one - * millisecond shared the temp path — both wrote it, the first rename - * consumed it, the second rename found nothing: ENOENT, ~1,500 times a day - * on a busy production brain, with a full ledger write per change behind it. - * - * Under pin: persists are single-flight and coalesced — one in flight, at - * most one trailing pass carrying the burst's final state — and every atomic - * write owns a unique temp path. A burst of N count changes costs at most - * two ledger writes, never errors, and leaves a ledger equal to memory. - */ -import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest' -import * as fs from 'node:fs' -import * as os from 'node:os' -import * as path from 'node:path' -import { Brainy } from '../../src/brainy.js' -import { NounType } from '../../src/types/graphTypes.js' - -describe('counts persistence is single-flight, coalesced, and never races its own temp file', () => { - let dir: string - let brain: any - - beforeEach(async () => { - process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true' - dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-counts-race-')) - brain = new Brainy({ - requireSubtype: false, - storage: { type: 'filesystem', path: dir }, - dimensions: 384, - silent: true - }) - await brain.init() - }) - - afterEach(async () => { - vi.restoreAllMocks() - await brain.close() - fs.rmSync(dir, { recursive: true, force: true }) - }) - - it('a burst of concurrent count changes → at most two ledger writes, zero errors, ledger == memory', async () => { - const storage = brain.storage - const countsPath: string = storage.countsFilePath - expect(countsPath, 'the filesystem adapter persists a counts ledger').toBeTruthy() - - // Let init's own persists settle so the burst is measured alone. - await storage.flushCounts?.() - - const renameSpy = vi.spyOn(fs.promises, 'rename') - const errorSpy = vi.spyOn(console, 'error') - - // Twenty-five concurrent count changes — the shape of a write burst; each - // used to launch its own persist. - const BURST = 25 - await Promise.all( - Array.from({ length: BURST }, () => storage.scheduleCountPersist()) - ) - - const ledgerRenames = renameSpy.mock.calls.filter(([, to]) => String(to) === countsPath) - expect(ledgerRenames.length, 'single-flight + one trailing pass').toBeLessThanOrEqual(2) - expect(ledgerRenames.length, 'the burst was persisted at all').toBeGreaterThanOrEqual(1) - - const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts')) - expect(persistErrors).toEqual([]) - - const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8')) - expect(ledger.totalNounCount).toBe(storage.totalNounCount) - expect(ledger.totalVerbCount).toBe(storage.totalVerbCount) - }) - - it('real writes in parallel: the ledger lands complete and no persist error is logged', async () => { - const storage = brain.storage - const countsPath: string = storage.countsFilePath - const errorSpy = vi.spyOn(console, 'error') - - await Promise.all( - Array.from({ length: 12 }, (_, i) => - brain.add({ data: `burst row ${i}`, type: NounType.Thing }) - ) - ) - await storage.flushCounts?.() - - const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts')) - expect(persistErrors).toEqual([]) - const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8')) - expect(ledger.totalNounCount).toBe(storage.totalNounCount) - expect(await brain.getNounCount()).toBe(ledger.totalNounCount) - }) - - it('every atomic write owns its own temp path — two writes in one millisecond never collide', async () => { - const storage = brain.storage - const tmpNames: string[] = [] - vi.spyOn(fs.promises, 'writeFile').mockImplementation(async (p: any) => { - tmpNames.push(String(p)) - }) - vi.spyOn(fs.promises, 'rename').mockImplementation(async () => undefined) - const target = path.join(dir, 'probe.json') - await Promise.all([ - storage.writeFileAtomic(target, '{"a":1}'), - storage.writeFileAtomic(target, '{"a":2}'), - storage.writeFileAtomic(target, '{"a":3}') - ]) - const probeTmps = tmpNames.filter((n) => n.startsWith(`${target}.tmp-`)) - expect(probeTmps.length).toBe(3) - expect(new Set(probeTmps).size, 'no two writes shared a temp path').toBe(3) - }) -}) diff --git a/tests/integration/find-connected-order.test.ts b/tests/integration/find-connected-order.test.ts deleted file mode 100644 index b04e7f99..00000000 --- a/tests/integration/find-connected-order.test.ts +++ /dev/null @@ -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 - const anchor = 'anchor' - const sharedVector = generateTestVector() - const neighbourIds = new Set(Array.from({ length: NEIGHBOURS }, (_, i) => v5(`nb-${i}`))) - - beforeAll(async () => { - brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } }) - await brain.init() - await brain.add({ - id: anchor, - data: 'the anchor', - type: NounType.Person, - metadata: { kind: 'anchor' }, - vector: generateTestVector() - }) - for (let i = 0; i < NOISE; i++) { - await brain.add({ - id: `noise-${i}`, - data: `noise ${i}`, - type: NounType.Person, - metadata: { kind: 'note', rank: 1000 + i }, - vector: sharedVector - }) - } - for (let i = 0; i < NEIGHBOURS; i++) { - await brain.add({ - id: `nb-${i}`, - data: `neighbour ${i}`, - type: NounType.Person, - metadata: { kind: 'note', rank: i + 1, ...(i < RETRACTED ? { retracted: true } : {}) }, - vector: sharedVector - }) - await brain.relate({ from: anchor, to: `nb-${i}`, type: VerbType.Knows }) - } - }) - - afterAll(async () => { - brain = null as any - }) - - it('returns the matching neighbours page by page — none dropped, never a non-neighbour', async () => { - const seen = new Set() - for (let offset = 0; offset <= NEIGHBOURS; offset += 10) { - const page = await brain.find({ - connected: { from: anchor, direction: 'out' }, - where: { kind: 'note' }, - limit: 10, - offset - }) - expect(page).toHaveLength(offset < NEIGHBOURS ? 10 : 0) - for (const r of page) { - expect(neighbourIds.has(r.entity.id)).toBe(true) - expect(seen.has(r.entity.id)).toBe(false) - seen.add(r.entity.id) - } - } - expect(seen.size).toBe(NEIGHBOURS) - }) - - it('evaluates a negation (`missing`) over the neighbour set, not the store', async () => { - const results = await brain.find({ - connected: { from: anchor, direction: 'out' }, - where: { kind: 'note', retracted: { missing: true } }, - limit: 100 - }) - expect(results).toHaveLength(NEIGHBOURS - RETRACTED) - for (const r of results) { - expect(neighbourIds.has(r.entity.id)).toBe(true) - expect(r.entity.metadata.retracted).toBeUndefined() - } - }) - - it('asks the metadata index about the neighbour ids only, and hydrates one page', async () => { - const index = (brain as any).metadataIndex - const within = vi.spyOn(index, 'filterIdsWithin') - const hydrate = vi.spyOn(brain as any, 'batchGet') - try { - const results = await brain.find({ - connected: { from: anchor, direction: 'out' }, - where: { kind: 'note' }, - limit: 10 - }) - expect(results).toHaveLength(10) - expect(within).toHaveBeenCalledTimes(1) - const askedIds = within.mock.calls[0][1] as string[] - expect(askedIds).toHaveLength(NEIGHBOURS) - for (const id of askedIds) expect(neighbourIds.has(id)).toBe(true) - expect(hydrate).toHaveBeenCalledTimes(1) - expect(hydrate.mock.calls[0][0]).toHaveLength(10) - } finally { - within.mockRestore() - hydrate.mockRestore() - } - }) - - it('orders the WHOLE neighbour set before cutting the page', async () => { - const results = await brain.find({ - connected: { from: anchor, direction: 'out' }, - where: { kind: 'note' }, - orderBy: 'rank', - order: 'desc', - limit: 5 - }) - expect(results.map((r) => r.entity.metadata.rank)).toEqual([30, 29, 28, 27, 26]) - }) - - it('walks the vector leg over the neighbours only', async () => { - const results = await brain.find({ - vector: sharedVector, - connected: { from: anchor, direction: 'out' }, - where: { kind: 'note' }, - limit: 5 - }) - expect(results).toHaveLength(5) - for (const r of results) expect(neighbourIds.has(r.entity.id)).toBe(true) - }) - - it('an anchor without neighbours answers [] before the filter is asked', async () => { - const index = (brain as any).metadataIndex - const within = vi.spyOn(index, 'filterIdsWithin') - try { - const results = await brain.find({ - connected: { from: 'noise-0', direction: 'out' }, - where: { kind: 'note' }, - limit: 10 - }) - expect(results).toEqual([]) - expect(within).not.toHaveBeenCalled() - } finally { - within.mockRestore() - } - }) -}) diff --git a/tests/integration/find-planner-door.test.ts b/tests/integration/find-planner-door.test.ts deleted file mode 100644 index 964b13f9..00000000 --- a/tests/integration/find-planner-door.test.ts +++ /dev/null @@ -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 - 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 ( - door: (...a: any[]) => Promise, - body: () => Promise - ): Promise => { - 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() - } - }) -}) diff --git a/tests/integration/pending-embed-low-water.test.ts b/tests/integration/pending-embed-low-water.test.ts deleted file mode 100644 index ff01b349..00000000 --- a/tests/integration/pending-embed-low-water.test.ts +++ /dev/null @@ -1,145 +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`, and the fold runs behind the doors as a - * latched background task the worker, `awaitPendingEmbeds()` and `close()` - * wait on. 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, behind the doors', () => { - const roots: string[] = [] - const dir = (): string => { - const d = mkdtempSync(join(tmpdir(), 'brainy-lowwater-')) - roots.push(d) - return d - } - const open = async (root: string): Promise> => { - const brain = new Brainy({ - requireSubtype: false, - storage: { type: 'filesystem', path: root } - }) - await brain.init() - return brain - } - - afterEach(() => { - for (const d of roots.splice(0)) rmSync(d, { recursive: true, force: true }) - }) - - it('drain-to-empty writes the mark, and the next open scans from mark + 1', async () => { - const root = dir() - const brain = await open(root) - // Hold the worker so the pending state is observable, then release it. - const realKick = (brain as any).kickEmbedWorker.bind(brain) - ;(brain as any).kickEmbedWorker = () => {} - await brain.add({ - id: 'row-1', - data: 'the first deferred row', - type: NounType.Thing, - deferEmbedding: true - }) - expect(brain.pendingEmbedCount()).toBeGreaterThan(0) - ;(brain as any).kickEmbedWorker = realKick - await brain.awaitPendingEmbeds() - // The drain wrote the advisory mark (fire-and-forget: settle the microtask). - await new Promise((r) => setTimeout(r, 50)) - const mark = (await (brain as any).storage.readRawObject(LOWWATER_PATH)) as { - generation: number - } | null - expect(mark).not.toBeNull() - expect(mark!.generation).toBeGreaterThan(0) - await brain.close() - - const brain2 = await open(root) - const log = (brain2 as any).generationStore.getFactLog() - const scanSpy = vi.spyOn(log, 'scanFacts') - try { - await (brain2 as any).recoverPendingEmbedsFromLog() - expect(scanSpy).toHaveBeenCalledTimes(1) - const opts = scanSpy.mock.calls[0][0] as { fromGeneration?: number } - expect(opts.fromGeneration).toBeGreaterThanOrEqual(mark!.generation + 1) - } finally { - scanSpy.mockRestore() - await brain2.close() - } - }) - - it('a pending embed enqueued after the mark survives an unclean stop', async () => { - const root = dir() - const brain = await open(root) - await brain.add({ id: 'settled', data: 'lands before the mark', type: NounType.Thing }) - await brain.awaitPendingEmbeds() - await new Promise((r) => setTimeout(r, 50)) - - // A deferred write whose embed never lands: block the worker, then drop - // the instance without close() — the unclean-stop shape. - ;(brain as any).kickEmbedWorker = () => {} - await brain.add({ - id: 'orphan', - data: 'enqueued then abandoned', - type: NounType.Thing, - deferEmbedding: true - }) - expect(brain.pendingEmbedCount()).toBeGreaterThan(0) - // No close(): simulate the crash by releasing only the writer lock so the - // next open can proceed. - await (brain as any).storage.releaseWriterLock() - - const brain2 = await open(root) - await (brain2 as any)._pendingEmbedRecovery - 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('open arms the fold as a background latch; awaitPendingEmbeds waits on it', 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 latch exists the moment init() returns (writable filesystem brain)… - expect((brain2 as any)._pendingEmbedRecovery).not.toBeNull() - // …and the barrier settles it before answering. - await brain2.awaitPendingEmbeds() - expect(brain2.pendingEmbedCount()).toBe(0) - await brain2.close() - }) - - it('a clean close with an empty set writes the mark even if no drain happened', async () => { - const root = dir() - const brain = await open(root) - await brain.add({ id: 'r1', data: 'row one', type: NounType.Thing }) - await brain.awaitPendingEmbeds() - await brain.close() - // Read the mark back through the storage door (the adapter owns the - // on-disk encoding), on a fresh instance. - const brain2 = await open(root) - const mark = (await (brain2 as any).storage.readRawObject(LOWWATER_PATH)) as { - generation: number - } | null - expect(mark).not.toBeNull() - expect(mark!.generation).toBeGreaterThan(0) - await brain2.close() - }) -}) diff --git a/tests/integration/related-verb-array.test.ts b/tests/integration/related-verb-array.test.ts deleted file mode 100644 index 36a49850..00000000 --- a/tests/integration/related-verb-array.test.ts +++ /dev/null @@ -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 - - beforeAll(async () => { - brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } }) - await brain.init() - for (const id of ['a', 'b', 'c', 'd']) { - await brain.add({ id, data: `node ${id}`, type: NounType.Person }) - } - await brain.relate({ from: 'a', to: 'b', type: VerbType.Supports }) - await brain.relate({ from: 'a', to: 'c', type: VerbType.RelatedTo }) - await brain.relate({ from: 'a', to: 'd', type: VerbType.Knows }) - await brain.relate({ from: 'b', to: 'c', type: VerbType.RelatedTo }) - }) - - afterAll(async () => { - brain = null as any - }) - - it('from + type array: the second type\'s edge comes back, both orders', async () => { - for (const types of [ - [VerbType.Supports, VerbType.RelatedTo], - [VerbType.RelatedTo, VerbType.Supports] - ]) { - const edges = await brain.related({ from: 'a', type: types }) - const targets = new Set(edges.map((e) => e.to)) - expect(targets.has(v5('b')), `types [${types}] missing Supports edge`).toBe(true) - expect(targets.has(v5('c')), `types [${types}] missing RelatedTo edge`).toBe(true) - expect(targets.has(v5('d'))).toBe(false) - expect(edges).toHaveLength(2) - } - }) - - it('a single-element array behaves exactly like the scalar', async () => { - const scalar = await brain.related({ from: 'a', type: VerbType.Supports }) - const array = await brain.related({ from: 'a', type: [VerbType.Supports] }) - expect(array.map((e) => e.id).sort()).toEqual(scalar.map((e) => e.id).sort()) - expect(array).toHaveLength(1) - }) - - it('no duplicate edges when types overlap the same edge set', async () => { - const edges = await brain.related({ - from: 'a', - type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows] - }) - const ids = edges.map((e) => e.id) - expect(new Set(ids).size).toBe(ids.length) - expect(edges).toHaveLength(3) - }) - - it('type-only asks (no anchor) honour the whole array too', async () => { - const edges = await brain.related({ type: [VerbType.Supports, VerbType.Knows] }) - const verbs = new Set(edges.map((e) => e.type)) - expect(verbs.has(VerbType.Supports)).toBe(true) - expect(verbs.has(VerbType.Knows)).toBe(true) - expect(edges).toHaveLength(2) - }) - - it('to + type array: the target side honours every type too', async () => { - const edges = await brain.related({ to: 'c', type: [VerbType.RelatedTo, VerbType.Supports] }) - const froms = new Set(edges.map((e) => e.from)) - expect(froms.has(v5('a'))).toBe(true) - expect(froms.has(v5('b'))).toBe(true) - expect(edges).toHaveLength(2) - }) - - it('pagination stays consistent across the union', async () => { - const page1 = await brain.related({ from: 'a', type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows], limit: 2 }) - const page2 = await brain.related({ from: 'a', type: [VerbType.Supports, VerbType.RelatedTo, VerbType.Knows], limit: 2, offset: 2 }) - const all = [...page1, ...page2].map((e) => e.id) - expect(new Set(all).size).toBe(3) - }) -}) diff --git a/tests/integration/transact-edge-delete-bigint-aliasing.test.ts b/tests/integration/transact-edge-delete-bigint-aliasing.test.ts deleted file mode 100644 index b3902571..00000000 --- a/tests/integration/transact-edge-delete-bigint-aliasing.test.ts +++ /dev/null @@ -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) - .filter(([, v]) => typeof v === 'bigint') - .map(([k]) => k) -} - -describe('transact() edge deletes never carry BigInt across the metadata seam', () => { - let dir: string - let brain: any - let crossings: Array<{ door: string; id: string; keys: string[] }> - - beforeEach(async () => { - process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true' - dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-tx-bigint-')) - brain = new Brainy({ - requireSubtype: false, - storage: { type: 'filesystem', path: dir }, - dimensions: 384, - silent: true - }) - await brain.init() - - // Spy on the seam the way a strict native provider judges it: record the - // BigInt-valued top-level keys of every metadata argument that crosses. - // The JS baseline index tolerates BigInts, so without this the baseline - // run would green a shape the native pair aborts on. - crossings = [] - const index = brain.metadataIndex - for (const door of ['addToIndex', 'removeFromIndex'] as const) { - const real = index[door].bind(index) - index[door] = (id: string, metadata: unknown, ...rest: unknown[]) => { - crossings.push({ door, id, keys: bigintKeys(metadata) }) - return real(id, metadata, ...rest) - } - } - }) - - afterEach(async () => { - await brain.close() - fs.rmSync(dir, { recursive: true, force: true }) - }) - - it('CASE 1 (the fleet repro): relate, then transact([{op: unrelate}])', async () => { - const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing }) - const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing }) - const verbId = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo }) - - crossings.length = 0 - await brain.transact([{ op: 'unrelate', id: verbId }]) - - const polluted = crossings.filter((c) => c.keys.length > 0) - expect(polluted).toEqual([]) - expect(await brain.storage.getVerb(verbId)).toBeFalsy() - }) - - it('CASE 2 (the cascade shape): transact([{op: remove}]) cascading edge deletes', async () => { - const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing }) - const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing }) - const c = await brain.add({ id: freshId(), data: 'c', type: NounType.Thing }) - const ab = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo }) - const ca = await brain.relate({ from: c, to: a, type: VerbType.RelatedTo }) - - crossings.length = 0 - await brain.transact([{ op: 'remove', id: a }]) - - const polluted = crossings.filter((c2) => c2.keys.length > 0) - expect(polluted).toEqual([]) - expect(await brain.get(a)).toBeFalsy() - expect(await brain.storage.getVerb(ab)).toBeFalsy() - expect(await brain.storage.getVerb(ca)).toBeFalsy() - }) - - it('CASE 3 (one batch, both legs): adds + relate + unrelate of a pre-existing edge', async () => { - const a = await brain.add({ id: freshId(), data: 'a', type: NounType.Thing }) - const b = await brain.add({ id: freshId(), data: 'b', type: NounType.Thing }) - const old = await brain.relate({ from: a, to: b, type: VerbType.RelatedTo }) - - const x = freshId() - crossings.length = 0 - await brain.transact([ - { op: 'add', id: x, data: 'x', type: NounType.Thing }, - { op: 'relate', from: a, to: x, type: VerbType.RelatedTo }, - { op: 'unrelate', id: old } - ]) - - const polluted = crossings.filter((c) => c.keys.length > 0) - expect(polluted).toEqual([]) - expect(await brain.storage.getVerb(old)).toBeFalsy() - const edges = await brain.related({ from: a }) - expect(edges.length).toBe(1) - expect(edges[0].id).not.toBe(old) - }) -}) - -describe('the metadata-index operations sanitize at the crossing, not at construction', () => { - /** A strict seam: refuses BigInts exactly as the native provider does. */ - const strictIndex = () => { - const seen: Array<{ door: string; keys: string[] }> = [] - const judge = (door: string, metadata: unknown) => { - const keys = bigintKeys(metadata) - seen.push({ door, keys }) - if (keys.length > 0) { - throw new Error( - `${door}: the metadata object violates the provider seam's JSON ` + - `contract — BigInt at ${keys.join(', ')}.` - ) - } - } - return { - seen, - addToIndex: async (_id: string, metadata: unknown) => judge('addToIndex', metadata), - removeFromIndex: async (_id: string, metadata: unknown) => judge('removeFromIndex', metadata) - } - } - - it('RemoveFromMetadataIndexOperation: entity mutated AFTER construction still crosses clean', async () => { - const index = strictIndex() - const verb: Record = { id: 'v1', sourceId: 'a', targetId: 'b' } - const op = new RemoveFromMetadataIndexOperation(index as any, 'v1', verb, () => 7n) - - // The graph leg's execute-time endpoint resolution, simulated: the shared - // object is polluted between plan and execute. - verb.sourceInt = 800_000n - verb.targetInt = 800_001n - - const rollback = await op.execute() - await rollback() - expect(index.seen.map((s) => s.keys)).toEqual([[], []]) - }) - - it('AddToMetadataIndexOperation: same law on the add leg and its rollback', async () => { - const index = strictIndex() - const verb: Record = { id: 'v2', sourceId: 'a', targetId: 'b' } - const op = new AddToMetadataIndexOperation(index as any, 'v2', verb, () => 7n) - - verb.sourceInt = 800_000n - verb.targetInt = 800_001n - - const rollback = await op.execute() - await rollback() - expect(index.seen.map((s) => s.keys)).toEqual([[], []]) - }) -})