diff --git a/.forgejo/workflows/ci.yml b/.forgejo/workflows/ci.yml index da5887f6..5e93cd96 100644 --- a/.forgejo/workflows/ci.yml +++ b/.forgejo/workflows/ci.yml @@ -5,10 +5,6 @@ name: CI # sequential, so tag-triggered matrix jobs (~22 min) would queue AHEAD of the # tag's publish-source run and starve every release (observed on 8.10.3 and # 9.0.0: the publish sat behind the tag's own redundant CI). -concurrency: - group: ci-${{ github.ref }} - cancel-in-progress: true - on: push: branches: ['**'] diff --git a/CHANGELOG.md b/CHANGELOG.md index a54d609e..7154d5a2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,18 @@ 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 c4f66561..9e573da3 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@soulcraftlabs/brainy", - "version": "10.4.4", + "version": "10.4.6", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@soulcraftlabs/brainy", - "version": "10.4.4", + "version": "10.4.6", "license": "MIT", "dependencies": { "@msgpack/msgpack": "^3.1.2", diff --git a/package.json b/package.json index 06ce0253..51322998 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@soulcraftlabs/brainy", - "version": "10.4.4", + "version": "10.4.6", "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 deleted file mode 100644 index 8f61c7f2..00000000 --- a/releases/brainy.json +++ /dev/null @@ -1,76 +0,0 @@ -{ - "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 deleted file mode 100644 index dab25971..00000000 --- a/releases/open-brainy.json +++ /dev/null @@ -1,122 +0,0 @@ -{ - "product": "open-brainy", - "entries": [ - { - "version": "10.4.10", - "date": "2026-09-02", - "headline": "A planner door for indexes, batched containment repair, and a fixed near()", - "items": [ - "An optional planFindPage door lets an index plan a find() and answer it in one call, instead of the engine assembling the plan itself.", - "repairContainment's reconcile pass now walks paged edges once instead of issuing one graph call per file.", - "find({ near }) now searches around the anchor's own vector and refuses by name when none is available, instead of silently querying with no vector at all." - ], - "url": "https://source.soulcraft.com/soulcraftlabs/open-brainy/releases/tag/v10.4.10", - "thumb": null - }, - { - "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 5d434320..1a4fe575 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/soulcraftlabs/open-brainy/releases" \ + if curl -sf -X POST "https://source.soulcraft.com/api/v1/repos/soulcraft/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 da04577e..17fa4ad9 100644 --- a/src/brainy.ts +++ b/src/brainy.ts @@ -15,6 +15,7 @@ import { JsHnswVectorIndex } from './hnsw/hnswIndex.js' import { createStorage, resolveFilesystemRoot } from './storage/storageFactory.js' import type { StorageOptions } from './storage/storageFactory.js' import { rebuildCounts } from './utils/rebuildCounts.js' +import { jsonSafeIndexMetadata } from './utils/jsonSafeIndexMetadata.js' import type { MetadataWriteBuffer } from './utils/metadataWriteBuffer.js' import { BaseStorage } from './storage/baseStorage.js' import { @@ -4203,32 +4204,19 @@ export class Brainy implements BrainyInterface { */ /** * @description A JSON-safe view of a record bound for the metadata-index - * crossing. The seam's metadata is JSON-safe BY CONTRACT (a native provider - * serializes it; u64 ints as Number corrupt above 2^53) — but - * {@link resolveVerbEndpointInts} MIRRORS the resolved endpoint ints onto - * the verb object itself as BigInt (`verb.sourceInt`/`targetInt`), so a - * verb object reused as index metadata carried BigInts into - * JSON.stringify, which throws, aborting the whole transaction (found by - * the first joint pair gate). Endpoint ints ride their OWN op params on the - * graph legs — the metadata crossing drops every BigInt-valued top-level - * key instead of guessing at a lossy numeric encoding. + * crossing — delegates to the shared {@link jsonSafeIndexMetadata} leaf, + * which the metadata-index transaction operations ALSO apply at execute + * and rollback time. This plan-time wrap alone proved insufficient: it + * returns the same reference when the record is clean, and `transact()`'s + * delete legs share that reference with a graph-retraction op whose + * execute-time endpoint resolution mirrors BigInt ints onto it (the full + * aliasing story lives on the leaf module's doc). * @param metadata - The candidate index-metadata record. * @returns The same object when already JSON-safe, else a shallow copy * without the BigInt-valued keys. */ private static jsonSafeIndexMetadata(metadata: unknown): unknown { - if (metadata === null || typeof metadata !== 'object') return metadata - const rec = metadata as Record - 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 + return jsonSafeIndexMetadata(metadata) } private metadataIndexRetractionOp( @@ -7482,7 +7470,37 @@ export class Brainy implements BrainyInterface { // JS path — there the materialized `candidateIds` restricts the walk instead. let preResolvedAllowedIds: OpaqueIdSet | undefined - if (params.where || params.type || params.subtype || params.service || params.excludeVFS) { + // Graph-first law (10.4.8, BRAINY-PROD-LATENCY-TRIAD rounds 44/45): with + // `connected` present the NEIGHBOUR SET is the candidate universe. It is + // resolved first from the adjacency (O(neighbours)), the metadata filter + // is evaluated over those ids only, and paging happens LAST. The earlier + // order materialized the whole-store filtered id list, paged it, hydrated + // the page, and only then intersected with the neighbours — O(store) per + // call, and a neighbour outside the first page was silently dropped. + let graphFirstIds: string[] | null = null + if (hasGraphCriteria) { + graphFirstIds = await this.resolveConnectedIds(params) + if (hiddenIds.size > 0) { + graphFirstIds = graphFirstIds.filter((id) => !hiddenIds.has(id)) + } + if ( + graphFirstIds.length > 0 && + (params.where || params.type || params.subtype || params.service || params.excludeVFS) + ) { + preResolvedFilter = this.buildMetadataFilter(params) + graphFirstIds = await this.filterIdsWithinBelted(preResolvedFilter, graphFirstIds) + } + if (graphFirstIds.length === 0) { + return [] + } + if (!hasVectorSearchCriteria) { + return await this.pageConnectedIds(params, graphFirstIds) + } + // The vector leg walks ONLY the neighbours (its candidate walk). The + // filter is already applied above, so no opaque universe is produced — + // it would describe the whole store, not the neighbour set. + preResolvedMetadataIds = graphFirstIds + } else if (params.where || params.type || params.subtype || params.service || params.excludeVFS) { preResolvedFilter = this.buildMetadataFilter(params) preResolvedMetadataIds = await this.filterIdsBelted(preResolvedFilter) @@ -7671,9 +7689,11 @@ export class Brainy implements BrainyInterface { } } - // Graph search component with O(1) traversal - if (params.connected) { - results = await this.executeGraphSearch(params, results) + // The text leg of a hybrid find has no candidate door, so its hits are + // held to the neighbour set here; the vector leg walked only the neighbours. + if (graphFirstIds !== null && results.length > 0) { + const neighbourSet = new Set(graphFirstIds) + results = results.filter((r) => neighbourSet.has(r.id)) } // Apply fusion scoring if requested @@ -12788,6 +12808,29 @@ 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 @@ -15771,16 +15814,16 @@ export class Brainy implements BrainyInterface { } /** - * Execute graph search component. + * Resolve `params.connected` to the neighbour id set — the graph-first + * find's candidate universe (deterministic traversal order, anchors excluded). * * Honors the full `GraphConstraints` contract: multi-hop `depth` (breadth-first via - * `neighbors()`), `via`/`type` verb-type filtering, and `direction`. Previously this read - * only `from`/`to`/`direction` and did a single 1-hop `getNeighbors()`, so `depth` and `via` - * were silently ignored — `find({ connected: { from, depth: 3 } })` returned only the - * immediate neighbour at every depth. + * `neighbors()`), `via`/`type` verb-type filtering, and `direction`. An empty set + * is re-verified against the adjacency before it is believed — a not-serving + * adjacency throws rather than answering `[]` as truth. */ - private async executeGraphSearch(params: FindParams, existingResults: Result[]): Promise[]> { - if (!params.connected) return existingResults + private async resolveConnectedIds(params: FindParams): Promise { + if (!params.connected) return [] const { from, to, depth, direction = 'both' } = params.connected const via = params.connected.via ?? params.connected.type @@ -15834,8 +15877,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 — overall result limiting happens - // downstream against existingResults. + // No limit: match the JS BFS exactly — the page is cut downstream, + // after the metadata filter, by pageConnectedIds / the candidate walk. const reachedInts = await provider.findConnectedSubtype( anchorInt, verbTypeIndex, subtypeArr[0], effectiveDepth, null ) @@ -15920,22 +15963,44 @@ export class Brainy implements BrainyInterface { await this.verifyGraphAdjacencyLive() } - // Filter existing results to only connected entities - if (existingResults.length > 0) { - return existingResults.filter(r => connectedIds.has(r.id)) - } + return [...connectedIds] + } - // Batch-load connected entities for fast cloud-storage performance + /** + * Page and hydrate an already-filtered neighbour set — the pure graph (and + * graph + metadata) find's tail. `orderBy` sorts the WHOLE set by field value + * before the page is cut (never the page after), null values last on `asc` + * and first on `desc`; without `orderBy` the traversal order stands. + */ + private async pageConnectedIds(params: FindParams, 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) + } + const pageIds = ordered.slice(offset, offset + limit) + const entitiesMap = await this.batchGet(pageIds) const results: Result[] = [] - const ids = [...connectedIds] - const entitiesMap = await this.batchGet(ids) - for (const id of ids) { + for (const id of pageIds) { const entity = entitiesMap.get(id) if (entity) { results.push(this.createResult(id, 1.0, entity)) } } - return results } diff --git a/src/plugin.ts b/src/plugin.ts index b1aef8e0..15b14b4e 100644 --- a/src/plugin.ts +++ b/src/plugin.ts @@ -411,6 +411,19 @@ 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 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 cabe2e30..a90adb93 100644 --- a/src/storage/adapters/baseStorageAdapter.ts +++ b/src/storage/adapters/baseStorageAdapter.ts @@ -1089,6 +1089,10 @@ export abstract class BaseStorageAdapter implements StorageAdapter { // Counts changed since the last persist? Drives the write-through flush. protected pendingCountPersist = false + /** The one persist running right now, if any (single-flight law — see flushCounts). */ + private countPersistInFlight: Promise | 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 @@ -1341,15 +1345,46 @@ export abstract class BaseStorageAdapter implements StorageAdapter { return } - try { - // Persist to storage (implemented by subclass) - await this.persistCounts() - this.pendingCountPersist = false - } catch (error) { - console.error('CRITICAL: Failed to flush counts to storage:', error) - // Keep pending flag set so we retry on next operation - throw error + // SINGLE-FLIGHT, COALESCED. Counts are write-through on every change, so + // a burst of writes used to launch one persist per change, all in flight + // together. Two of them inside the same millisecond shared the atomic + // writer's temp path (`.tmp--`): both wrote it, the first rename + // consumed it, the second rename found nothing — ENOENT, ~1,500 times a + // day on a busy production brain, with a full ledger write per change + // behind it. Now exactly one persist runs at a time; requests that arrive + // while it runs collapse into ONE trailing persist that carries the final + // state. A burst of N changes costs at most two writes and never races + // itself. + if (this.countPersistInFlight) { + // The in-flight write may have already serialised a stale snapshot — + // ask for one more pass after it, and let every caller in this burst + // await that same pass. + if (!this.countPersistTrailing) { + this.countPersistTrailing = this.countPersistInFlight + .catch(() => undefined) + .then(() => { + this.countPersistTrailing = null + return this.flushCounts() + }) + } + return this.countPersistTrailing } + + this.countPersistInFlight = (async () => { + try { + // Persist to storage (implemented by subclass) + this.pendingCountPersist = false + await this.persistCounts() + } catch (error) { + // Keep the flag set so the next operation retries. + this.pendingCountPersist = true + console.error('CRITICAL: Failed to flush counts to storage:', error) + throw error + } finally { + this.countPersistInFlight = null + } + })() + return this.countPersistInFlight } /** diff --git a/src/storage/adapters/fileSystemStorage.ts b/src/storage/adapters/fileSystemStorage.ts index 5ec1d88e..87b6406f 100644 --- a/src/storage/adapters/fileSystemStorage.ts +++ b/src/storage/adapters/fileSystemStorage.ts @@ -2400,8 +2400,15 @@ export class FileSystemStorage extends BaseStorage { * Atomic write via temp-file-then-rename so concurrent readers never see a * half-written lock JSON. Reused by writer-lock writes + heartbeat. */ + /** Monotonic per-process sequence so two atomic writes never share a temp path. */ + private static atomicWriteSeq = 0 + private async writeFileAtomic(filePath: string, contents: string): Promise { - const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}` + // pid + timestamp alone collided: two writers of the same target inside + // one millisecond shared this path, and the loser's rename found the + // winner had already moved it (ENOENT). The sequence makes every call's + // temp path its own. + const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}-${++FileSystemStorage.atomicWriteSeq}` await fs.promises.writeFile(tmp, contents) await fs.promises.rename(tmp, filePath) } diff --git a/src/transaction/operations/IndexOperations.ts b/src/transaction/operations/IndexOperations.ts index 1bbbca88..0142dc54 100644 --- a/src/transaction/operations/IndexOperations.ts +++ b/src/transaction/operations/IndexOperations.ts @@ -14,6 +14,7 @@ import type { MetadataIndexManager } from '../../utils/metadataIndex.js' import type { GraphVerb } from '../../coreTypes.js' import type { Operation, RollbackAction } from '../types.js' import { isZeroNormVector } from '../../utils/distance.js' +import { jsonSafeIndexMetadata } from '../../utils/jsonSafeIndexMetadata.js' import { prodLog } from '../../utils/logger.js' /** @@ -390,13 +391,21 @@ export class AddToMetadataIndexOperation implements Operation { // rollback so add + undo reference the same watermark. const generation = this.generationFn?.() - // Add to metadata index (skipFlush=true for transaction atomicity) - await this.index.addToIndex(this.id, this.entity, true, false, generation) + // The JSON-safe view is taken HERE, per crossing, never at construction: + // the entity reference this op holds can be mutated between plan and + // execute (a graph op's execute-time endpoint-int resolution mirrors + // BigInts onto a shared verb object) — see jsonSafeIndexMetadata's + // module doc. + await this.index.addToIndex( + this.id, jsonSafeIndexMetadata(this.entity), true, false, generation + ) // Return rollback action return async () => { // Remove from metadata index - await this.index.removeFromIndex(this.id, this.entity, generation) + await this.index.removeFromIndex( + this.id, jsonSafeIndexMetadata(this.entity), generation + ) } } } @@ -432,13 +441,21 @@ export class RemoveFromMetadataIndexOperation implements Operation { // Resolve the removal generation once; reuse it for the rollback re-add. const generation = this.generationFn?.() - // Remove from metadata index - await this.index.removeFromIndex(this.id, this.entity, generation) + // Sanitized per crossing, never at construction — transact()'s delete + // legs hand this op the SAME verb object the graph-retraction op's + // execute-time endpoint resolution mutates (BigInt sourceInt/targetInt), + // so a plan-time view aliases the pollution. See jsonSafeIndexMetadata's + // module doc. + await this.index.removeFromIndex( + this.id, jsonSafeIndexMetadata(this.entity), generation + ) // Return rollback action return async () => { // Re-add with original metadata (skipFlush=true) - await this.index.addToIndex(this.id, this.entity, true, false, generation) + await this.index.addToIndex( + this.id, jsonSafeIndexMetadata(this.entity), true, false, generation + ) } } } diff --git a/src/utils/jsonSafeIndexMetadata.ts b/src/utils/jsonSafeIndexMetadata.ts new file mode 100644 index 00000000..d3b1be5f --- /dev/null +++ b/src/utils/jsonSafeIndexMetadata.ts @@ -0,0 +1,47 @@ +/** + * @module utils/jsonSafeIndexMetadata + * @description The metadata-index crossing's JSON-safety law, as a leaf + * function both the coordinator and the transaction operations share. + * + * The seam's metadata is JSON-safe BY CONTRACT (a native provider serializes + * it; u64 ints as Number corrupt above 2^53) — but `resolveVerbEndpointInts` + * MIRRORS the resolved endpoint ints onto the verb object itself as BigInt + * (`verb.sourceInt`/`targetInt`), so a verb object reused as index metadata + * carries BigInts into JSON.stringify, which throws, aborting the whole + * transaction. Endpoint ints ride their OWN op params on the graph legs — the + * metadata crossing drops every BigInt-valued top-level key instead of + * guessing at a lossy numeric encoding. + * + * WHY THIS IS A LEAF MODULE, ENFORCED AT THE CROSSING: sanitizing only at + * operation-construction time is not enough. `transact()`'s delete legs pass + * the SAME verb object to both the graph-retraction op (whose endpoint-int + * thunk deliberately resolves at EXECUTE time, for same-batch forward refs) + * and the metadata-retraction op. At plan time the verb is still clean, so a + * plan-time sanitize returns the same reference — then the graph op executes + * first, mirrors the BigInt ints onto the shared object, and the metadata op + * crosses the seam with them (found by the first fleet adoption of the native + * pair: every transact-wrapped edge delete aborted). The crossing itself is + * the only place ordering cannot bypass. + */ + +/** + * A JSON-safe view of a record bound for the metadata-index crossing. + * + * @param metadata - The candidate index-metadata record. + * @returns The same object when already JSON-safe, else a shallow copy + * without the BigInt-valued keys. + */ +export function jsonSafeIndexMetadata(metadata: unknown): unknown { + if (metadata === null || typeof metadata !== 'object') return metadata + const rec = metadata as Record + 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 3e0e3d17..0fd312e2 100644 --- a/src/utils/metadataIndex.ts +++ b/src/utils/metadataIndex.ts @@ -2575,6 +2575,19 @@ 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 new file mode 100644 index 00000000..5acbdcc3 --- /dev/null +++ b/tests/integration/counts-persist-single-flight.test.ts @@ -0,0 +1,111 @@ +/** + * @module tests/integration/counts-persist-single-flight + * @description Regression for a production race in FileSystemStorage's + * counts ledger: `persistCounts()` was write-through on every count change + * with no serialization, and the atomic writer named its temp file with + * millisecond granularity (`.tmp--`). 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 new file mode 100644 index 00000000..b04e7f99 --- /dev/null +++ b/tests/integration/find-connected-order.test.ts @@ -0,0 +1,165 @@ +/** + * @module tests/integration/find-connected-order + * @description The graph-first law for `find({ connected })` (10.4.8). + * + * With `connected` present the neighbour set is the candidate universe: it is + * resolved from the adjacency first, the metadata filter is evaluated over + * those ids only, and the page is cut last. The earlier order materialized the + * whole-store filtered id list, paged it, hydrated the page, and only then + * intersected with the neighbours — so a neighbour outside the first page of + * the filtered STORE was silently dropped, and every call paid O(store). + * + * These pins hold both halves. The answer: every matching neighbour is + * reachable by paging, a non-neighbour never appears, a negation (`missing`) + * is evaluated over the neighbours, `orderBy` sorts the whole neighbour set + * before the page is cut, and the vector leg walks the neighbours only. The + * cost shape: the metadata index is asked about the neighbour ids only, and + * hydration is one page — never the store. + */ +import { describe, it, expect, beforeAll, afterAll, vi } from 'vitest' +import { Brainy } from '../../src/brainy' +import { NounType, VerbType } from '../../src/types/graphTypes' +import { v5 } from '../../src/universal/uuid' +import { generateTestVector } from '../helpers/test-factory' + +/** Matching rows that are NOT neighbours — added FIRST, so the whole-store filtered list leads with them. */ +const NOISE = 120 +/** Matching rows that ARE neighbours of the anchor. */ +const NEIGHBOURS = 30 +/** Neighbours carrying `retracted: true` — excluded by the `missing` negation. */ +const RETRACTED = 4 + +describe('find({ connected }) is graph-first: neighbours → filter → page', () => { + let brain: Brainy + 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/transact-edge-delete-bigint-aliasing.test.ts b/tests/integration/transact-edge-delete-bigint-aliasing.test.ts new file mode 100644 index 00000000..b3902571 --- /dev/null +++ b/tests/integration/transact-edge-delete-bigint-aliasing.test.ts @@ -0,0 +1,184 @@ +/** + * @module tests/integration/transact-edge-delete-bigint-aliasing + * @description Regression for a fleet-adoption blocker: ANY edge delete + * inside `transact()` — a direct unrelate or a noun-remove's cascade — + * aborted with the metadata seam's BigInt JSON-guard error on a strict + * (native) metadata provider. + * + * The aliasing chain: `planTxUnrelate`/the remove-cascade pass the SAME verb + * object to the graph-retraction op and the metadata-retraction op. The + * metadata leg's JSON-safe wrap ran at PLAN time, when the verb was still + * clean — so it returned the same reference. At EXECUTE time the graph op + * runs first and `resolveVerbEndpointInts` mirrors BigInt + * `sourceInt`/`targetInt` onto the shared object (deliberately deferred for + * same-batch forward refs — see transact-forward-ref-graph.test.ts); the + * metadata op then crossed the seam with the polluted object. Direct + * `unrelate()` resolves ints at BUILD time, before its sanitize, which is why + * only the transact() shapes ever hit it. + * + * Fix under pin: the JSON-safe view is taken AT THE CROSSING — inside the + * metadata-index operations' execute/rollback — so no plan-vs-execute + * ordering can bypass it. The JS baseline index tolerates BigInts (it would + * mask the bug), so these pins SPY on the seam and assert what actually + * crossed, exactly as a strict native provider would judge it. + */ +import { describe, it, expect, beforeEach, afterEach } from 'vitest' +import * as fs from 'node:fs' +import * as os from 'node:os' +import * as path from 'node:path' +import { Brainy } from '../../src/brainy.js' +import { NounType, VerbType } from '../../src/types/graphTypes.js' +import { + AddToMetadataIndexOperation, + RemoveFromMetadataIndexOperation +} from '../../src/transaction/operations/index.js' + +let seq = 0 +const freshId = (): string => + `00000000-0000-4000-8000-${(++seq).toString(16).padStart(12, '0')}` + +/** Top-level BigInt-valued keys of a candidate seam crossing (the guard's law). */ +const bigintKeys = (metadata: unknown): string[] => { + if (metadata === null || typeof metadata !== 'object') return [] + return Object.entries(metadata as Record) + .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([[], []]) + }) +})