/** * Base Storage Adapter * Provides common functionality for all storage adapters */ import { GraphAdjacencyIndex } from '../graph/graphAdjacencyIndex.js' import type { GraphEntityIdResolver } from '../graph/graphAdjacencyIndex.js' import { GraphVerb, HNSWNoun, HNSWVerb, NounMetadata, VerbMetadata, HNSWNounWithMetadata, HNSWVerbWithMetadata, StatisticsData } from '../coreTypes.js' import { BaseStorageAdapter } from './adapters/baseStorageAdapter.js' import { validateNounType, validateVerbType } from '../utils/typeValidation.js' import { NounType, VerbType, TypeUtils, NOUN_TYPE_COUNT, VERB_TYPE_COUNT } from '../types/graphTypes.js' import { getShardId } from './sharding.js' import { BlobStorage, type BlobStoreAdapter } from './blobStorage.js' import { unwrapBinaryData } from './binaryDataCodec.js' import { prodLog } from '../utils/logger.js' import { BrainyError } from '../errors/brainyError.js' import { MetadataWriteBuffer } from '../utils/metadataWriteBuffer.js' import { splitNounMetadataRecord, splitVerbMetadataRecord } from '../types/reservedFields.js' /** * Normalize a stored timestamp value to epoch milliseconds. Brainy 8.0 writes * plain numbers; records written by pre-8.0 cloud adapters may carry the * `{ seconds, nanoseconds }` object form. Anything else falls back to now — * matching the long-standing `|| Date.now()` combine behavior. * @param value - The raw `createdAt`/`updatedAt` value from a stored metadata record. * @returns Epoch milliseconds. */ function normalizeStoredTimestamp(value: unknown): number { if (typeof value === 'number' && value > 0) { return value } if ( value !== null && typeof value === 'object' && typeof (value as { seconds?: unknown }).seconds === 'number' ) { return (value as { seconds: number }).seconds * 1000 } return Date.now() } /** * Storage key analysis result * Used to determine whether a key is a system key or entity key, and its storage path */ interface StorageKeyInfo { original: string isEntity: boolean shardId: string | null directory: string fullPath: string } /** * Storage adapter batch configuration profile * Each storage adapter declares its optimal batch behavior for rate limiting * and performance optimization * */ export interface StorageBatchConfig { /** Maximum items per batch */ maxBatchSize: number /** Delay between batches in milliseconds (for rate limiting) */ batchDelayMs: number /** Maximum concurrent operations this storage can handle */ maxConcurrent: number /** Whether storage can handle parallel writes efficiently */ supportsParallelWrites: boolean /** Rate limit characteristics of this storage adapter */ rateLimit: { /** Approximate operations per second this storage can handle */ operationsPerSecond: number /** Maximum burst capacity before throttling occurs */ burstCapacity: number } } // Clean directory structure // All storage adapters use this consistent structure export const NOUNS_METADATA_DIR = 'entities/nouns/metadata' export const VERBS_METADATA_DIR = 'entities/verbs/metadata' export const SYSTEM_DIR = '_system' export const STATISTICS_KEY = 'statistics' /** * Metadata persisted in the writer lock file. Used by stale-lock detection * (PID liveness + hostname + heartbeat freshness) and by `brain.stats()` for * operator-facing diagnostics. */ export interface WriterLockInfo { pid: number hostname: string startedAt: string // ISO timestamp when the lock was first acquired lastHeartbeat: string // ISO timestamp of the most recent heartbeat update version: string // Brainy version that wrote the lock rootDir?: string // Convenience for log lines / error messages } /** * FNV-1a hash returning a 2-char hex bucket (00-ff). * Distributes system keys across 256 sub-prefixes to avoid * cloud storage per-prefix rate limits. */ function systemKeyBucket(key: string): string { let hash = 2166136261 for (let i = 0; i < key.length; i++) { hash ^= key.charCodeAt(i) hash += (hash << 1) + (hash << 4) + (hash << 7) + (hash << 8) + (hash << 24) } return ((hash >>> 0) & 0xff).toString(16).padStart(2, '0') } const SINGLETON_SYSTEM_KEYS = new Set([ '__metadata_field_registry__', 'brainy:entityIdMapper', 'statistics', 'counts', 'hnsw-system', 'type-statistics', ]) const SINGLETON_SYSTEM_PREFIXES = [ 'statistics_', ] function isSingletonSystemKey(key: string): boolean { if (SINGLETON_SYSTEM_KEYS.has(key)) return true return SINGLETON_SYSTEM_PREFIXES.some(p => key.startsWith(p)) } /** * Type-first path generators * Built-in type-aware organization for all storage adapters */ /** * Get ID-first path for noun vectors * No type parameter needed - direct O(1) lookup by ID */ function getNounVectorPath(id: string): string { const shard = getShardId(id) return `entities/nouns/${shard}/${id}/vectors.json` } /** * Get ID-first path for noun metadata * No type parameter needed - direct O(1) lookup by ID */ function getNounMetadataPath(id: string): string { const shard = getShardId(id) return `entities/nouns/${shard}/${id}/metadata.json` } /** * Get ID-first path for verb vectors * No type parameter needed - direct O(1) lookup by ID */ function getVerbVectorPath(id: string): string { const shard = getShardId(id) return `entities/verbs/${shard}/${id}/vectors.json` } /** * @description Extract the entity id embedded in a vector path * (`entities/{nouns|verbs}/{shard}/{id}/vectors.json`). Used by the cursored * noun/verb walks to order and skip candidates by id WITHOUT reading each file, * which is what keeps a full cursored pagination O(N) instead of O(N²). * @param path - A vector path (full or prefix-relative; must end with `/vectors.json`). * @returns The entity id (the path segment immediately before `/vectors.json`). */ function idFromVectorPath(path: string): string { const withoutSuffix = path.replace(/\/vectors\.json$/, '') const lastSlash = withoutSuffix.lastIndexOf('/') return lastSlash >= 0 ? withoutSuffix.slice(lastSlash + 1) : withoutSuffix } /** * Get ID-first path for verb metadata * No type parameter needed - direct O(1) lookup by ID */ function getVerbMetadataPath(id: string): string { const shard = getShardId(id) return `entities/verbs/${shard}/${id}/metadata.json` } /** * Optional count capabilities probed via duck typing by getNouns()/getVerbs(). * Adapters with a native O(1) count API may implement these; they are not part * of the BaseStorageAdapter contract, so BaseStorage feature-detects them at * runtime before falling back to scan-based counting. */ interface OptionalCountCapabilities { countNouns?: (filter?: { nounType?: string | string[] service?: string | string[] metadata?: Record }) => Promise countVerbs?: (filter?: { verbType?: string | string[] sourceId?: string | string[] targetId?: string | string[] service?: string | string[] metadata?: Record }) => Promise } /** * Whether an entity/relationship with the given visibility tier counts toward * the user-facing counts (`counts.json` totals, `nounCountsByType`, `stats()`). * Public entities (absent tier === `'public'`) are counted; `'internal'` and * `'system'` are excluded. Single source of truth for the count-exclusion rule. * * @param visibility - The stored visibility value (may be `undefined`/`unknown`). * @returns `true` when the entity should be counted, `false` for internal/system. */ function isCountedVisibility(visibility: unknown): boolean { return visibility !== 'internal' && visibility !== 'system' } /** * Base storage adapter that implements common functionality * This is an abstract class that should be extended by specific storage adapters */ export abstract class BaseStorage extends BaseStorageAdapter { protected isInitialized = false protected graphIndex?: GraphAdjacencyIndex protected graphIndexPromise?: Promise /** * Shared UUID ↔ int resolver for the graph index's BigInt boundary. * Wired by Brainy via {@link setGraphEntityIdResolver}; until then the verb * read paths fall back to shard iteration. */ protected graphEntityIdResolver?: GraphEntityIdResolver protected readOnly = false // Write-through cache for read-after-write consistency // Extended lifetime - persists until explicit flush() call // Guarantees that immediately after writeCanonicalObject(), readCanonicalObject() returns the data // Cache key: storage-root-relative object path // Cache lifetime: write start → flush() call (provides safety net for batch operations) // Memory footprint: Bounded by batch size (typically <1000 items during imports) private writeCache = new Map() /** * Clear the write-through cache * MUST be called by all storage adapter clear() implementations to ensure * read-after-write consistency cache doesn't return stale data after clear. * @protected - Available to subclasses for clear() implementation */ protected clearWriteCache(): void { this.writeCache.clear() } /** * Content-addressed blob store backing VFS file content (deduplicated, * zstd-compressed, SHA-256 addressed). Lazily created by * {@link initializeBlobStorage}; lives under the `_cas/` storage area. */ public blobStorage?: BlobStorage // Type-first indexing support // Built into all storage adapters for billion-scale efficiency protected nounCountsByType = new Uint32Array(NOUN_TYPE_COUNT) // 168 bytes (Stage 3: 42 types) protected verbCountsByType = new Uint32Array(VERB_TYPE_COUNT) // 508 bytes (Stage 3: 127 types) /** * Per-NounType-per-subtype counts. Outer key is the NounType index (matching * `nounCountsByType` indexing); inner key is the subtype string. Populated * incrementally as entities are saved and decremented on delete; persisted * to `_system/subtype-statistics.json` alongside type-statistics. * * Memory: one Map entry per (type, subtype) pair actually used — typically * tens of inner entries per NounType in production. Sparse by design — types * with no subtype-bearing entities have no outer entry. */ protected subtypeCountsByType = new Map>() /** * Per-VerbType-per-subtype counts. Verb-side mirror of `subtypeCountsByType`. * Outer key is the VerbType index (matching `verbCountsByType` indexing); * inner key is the subtype string. Populated incrementally as relationships * are saved and decremented on delete; persisted to * `_system/verb-subtype-statistics.json` (same shape as the noun-side rollup * — `{ counts: { [verbTypeIdx]: { [subtype]: count } }, updatedAt }`). */ protected verbSubtypeCountsByType = new Map>() // Count attribution (type / subtype / visibility) is sourced directly from the // canonical metadata RECORD, never from id-keyed in-memory caches. The metadata // save/delete paths already read the prior record (`existingMetadata` on write, // read-before-delete on remove), so the entity's `noun`/`verb` type, `subtype`, // and `visibility` are in hand exactly where a count must change — there is no // need for a parallel O(N) `id → type/subtype/visibility` map resident on the // writer. The five such caches that used to live here were removed in the 8.0 // billion-scale RAM pass; `nounCountsByType` / `verbCountsByType` (fixed // type-indexed Uint32Arrays) and `subtypeCountsByType` / `verbSubtypeCountsByType` // (bounded by distinct subtype labels) remain because they are NOT id-keyed. // Type caches REMOVED - ID-first paths eliminate need for type lookups! // With ID-first architecture, we construct paths directly from IDs: {SHARD}/{ID}/metadata.json // Type is just a field in the metadata, indexed by MetadataIndexManager for queries // Track if type counts have been rebuilt (prevent repeated rebuilds) private typeCountsRebuilt = false // Write buffer for cloud storage adapters — deduplicates rapid writes to the same path // FileSystem adapter does NOT use this (local writes are already fast) // Initialized by cloud adapters in their init() method protected metadataWriteBuffer: MetadataWriteBuffer | null = null /** * Analyze a storage key to determine its routing and path * @param id - The key to analyze (UUID or system key) * @param context - The context for the key (noun-metadata, verb-metadata, or system) * @returns Storage key information including path and shard ID * @private */ private analyzeKey(id: string, context: 'noun-metadata' | 'verb-metadata' | 'system'): StorageKeyInfo { // Guard against undefined/null IDs if (!id || typeof id !== 'string') { throw new Error(`Invalid storage key: ${id} (must be a non-empty string)`) } // System resource detection const isSystemKey = id.startsWith('__metadata_') || id.startsWith('__index_') || id.startsWith('__system_') || id.startsWith('statistics_') || id === 'statistics' || id.startsWith('__chunk__') || // Metadata index chunks (roaring bitmap data) id.startsWith('__sparse_index__') // Metadata sparse indices (zone maps + bloom filters) if (isSystemKey) { if (isSingletonSystemKey(id)) { return { original: id, isEntity: false, shardId: null, directory: SYSTEM_DIR, fullPath: `${SYSTEM_DIR}/${id}.json` } } const bucket = systemKeyBucket(id) return { original: id, isEntity: false, shardId: bucket, directory: `${SYSTEM_DIR}/idx/${bucket}`, fullPath: `${SYSTEM_DIR}/idx/${bucket}/${id}.json` } } // UUID validation for entity keys const uuidRegex = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i if (!uuidRegex.test(id)) { prodLog.warn(`[Storage] Unknown key format: ${id} - treating as system resource`) if (isSingletonSystemKey(id)) { return { original: id, isEntity: false, shardId: null, directory: SYSTEM_DIR, fullPath: `${SYSTEM_DIR}/${id}.json` } } const bucket = systemKeyBucket(id) return { original: id, isEntity: false, shardId: bucket, directory: `${SYSTEM_DIR}/idx/${bucket}`, fullPath: `${SYSTEM_DIR}/idx/${bucket}/${id}.json` } } // Valid entity UUID - apply sharding const shardId = getShardId(id) if (context === 'noun-metadata') { return { original: id, isEntity: true, shardId, directory: `${NOUNS_METADATA_DIR}/${shardId}`, fullPath: `${NOUNS_METADATA_DIR}/${shardId}/${id}.json` } } else if (context === 'verb-metadata') { return { original: id, isEntity: true, shardId, directory: `${VERBS_METADATA_DIR}/${shardId}`, fullPath: `${VERBS_METADATA_DIR}/${shardId}/${id}.json` } } else { // system context - but UUID format return { original: id, isEntity: false, shardId: null, directory: SYSTEM_DIR, fullPath: `${SYSTEM_DIR}/${id}.json` } } } /** * Initialize the storage adapter * Loads type statistics for built-in type-aware indexing * * IMPORTANT: If your adapter overrides init(), call await super.init() first! */ public async init(): Promise { // CRITICAL FIX - Set flag FIRST to prevent infinite recursion // If any code path during initialization calls ensureInitialized(), it would // trigger init() again. Setting the flag immediately breaks the recursion cycle. this.isInitialized = true try { // Load type statistics from storage (if they exist) await this.loadTypeStatistics() await this.loadSubtypeStatistics() await this.loadVerbSubtypeStatistics() // GraphAdjacencyIndex is now SINGLETON via getGraphIndex() // - Removed direct creation here to fix dual-ownership bug // - GraphAdjacencyIndex will be created lazily on first getGraphIndex() call // - This ensures there's only ONE instance per storage adapter // - See: https://github.com/soulcraftlabs/brainy/issues/vfs-corruption prodLog.debug('[BaseStorage] init() complete - GraphAdjacencyIndex will be created via getGraphIndex()') } catch (error) { // Reset flag on failure to allow retry this.isInitialized = false throw error } } /** * Rebuild GraphAdjacencyIndex from existing verbs * Call this manually if you have existing verb data that needs to be indexed * @public */ public async rebuildGraphIndex(): Promise { const index = await this.getGraphIndex() prodLog.info('[BaseStorage] Rebuilding graph index from existing data...') await index.rebuild() prodLog.info('[BaseStorage] Graph index rebuild complete') } /** * Invalidate GraphAdjacencyIndex * Call this when clearing data to force re-creation * The next getGraphIndex() call will create a fresh instance and rebuild * @public */ public invalidateGraphIndex(): void { if (this.graphIndex) { prodLog.info('[BaseStorage] Invalidating GraphAdjacencyIndex for clear()') // Stop any pending operations. stopAutoFlush is an optional capability // of plugin-provided graph indexes (duck-typed; the built-in // GraphAdjacencyIndex does not implement it). const flushable = this.graphIndex as GraphAdjacencyIndex & { stopAutoFlush?: () => void } if (typeof flushable.stopAutoFlush === 'function') { flushable.stopAutoFlush() } this.graphIndex = undefined this.graphIndexPromise = undefined } } /** * Set the graph index instance (used by Brainy to wire plugin-provided graph indexes). * This ensures getVerbsBySource() uses the fast GraphAdjacencyIndex path * instead of falling back to O(n) shard iteration. */ public setGraphIndex(index: GraphAdjacencyIndex): void { this.graphIndex = index this.graphIndexPromise = Promise.resolve(index) } /** * @description Wire the shared UUID ↔ int resolver used at the graph index's * BigInt boundary (8.0 u64 contract). Brainy calls this with * `metadataIndex.getIdMapper()` after the graph index is resolved (init, * fork, and checkout paths). The storage layer needs it to convert UUIDs to * entity ints before `getVerbIdsBySource`/`getVerbIdsByTarget` calls and to * resolve returned verb ints back to verb-id strings. Until it's wired, the * verb read paths fall back to shard iteration (correct, just slower). * @param resolver - The shared entity-id resolver. * @returns Nothing. */ public setGraphEntityIdResolver(resolver: GraphEntityIdResolver): void { this.graphEntityIdResolver = resolver // Thread the resolver into the JS graph index too, when present — native // providers carry their own mapper and don't expose this setter. if (this.graphIndex && typeof (this.graphIndex as GraphAdjacencyIndex).setEntityIdMapper === 'function') { (this.graphIndex as GraphAdjacencyIndex).setEntityIdMapper(resolver) } } /** * Ensure the storage adapter is initialized */ protected async ensureInitialized(): Promise { if (!this.isInitialized) { await this.init() } } /** * Whether this storage adapter enforces multi-process writer exclusion. * Filesystem storage returns true; cloud and memory adapters return false. * Brainy.init() checks this to decide whether to log a multi-process warning * in writer mode. */ public supportsMultiProcessLocking(): boolean { return false } /** * Attempt to acquire the process-level writer lock at init time. Default * implementation is a no-op (multi-process safety not enforced). Filesystem * storage overrides this with real locking semantics. * * @param options.force - If true, overwrite any existing writer lock. Use * only when stale detection cannot prove the existing lock is dead. * @returns Metadata about the acquired lock (or `null` if no lock was needed). * @throws If another live writer holds the lock and `force` is not set. */ public async acquireWriterLock(options?: { force?: boolean }): Promise { return null } /** * Release the writer lock acquired by `acquireWriterLock()`. No-op if no lock * was held. Filesystem storage overrides this to delete the lock file and * stop the heartbeat timer. */ public async releaseWriterLock(): Promise { // No-op by default } /** * Read the current writer-lock metadata if one is held by any process. Used * by `brain.stats()` for diagnostics. Default returns null (no lock model). */ public async readWriterLock(): Promise { return null } /** * Start watching for cross-process flush requests. The writer Brainy * instance calls this so that out-of-process inspectors can ask for a * synchronous flush before they open the store read-only. Default is a * no-op (non-filesystem backends have no shared filesystem to poll). * * @param onRequest - Callback invoked when a request file appears. Should * call `brain.flush()` and resolve when persistence is complete. */ public startFlushRequestWatcher(onRequest: () => Promise): void { // No-op by default } /** * Stop the flush-request watcher started by `startFlushRequestWatcher`. */ public stopFlushRequestWatcher(): void { // No-op by default } /** * Write a flush-request file and wait for the writer to acknowledge by * writing the corresponding response file. Returns true if a response was * received before the timeout, false if it timed out. * * Cross-platform RPC over the shared filesystem — no signals, so this works * on Windows, Linux, macOS, and inside containers without IPC config. */ public async requestFlushOverFilesystem(_timeoutMs: number): Promise { return false } /** * @description Initialize the content-addressed blob store (idempotent). * Creates a {@link BlobStorage} over this adapter's object primitives, * rooted at the `_cas/` storage area. The blob store backs VFS file * content: deduplicated, SHA-256 addressed, zstd-compressed where it pays. * * Called automatically during `brain.init()` (before the VFS is built) and * again after `clear()` re-creates the storage area. * * @returns Promise that resolves when the blob store is ready. */ public async initializeBlobStorage(): Promise { if (this.blobStorage) { return } // Key-value bridge: adapts this adapter's object primitives to the // BlobStoreAdapter interface. Key naming is an explicit type contract: // - 'blob-meta:' → JSON (BlobStorage metadata) // - 'blob:' → binary (possibly compressed blob bytes) const casAdapter: BlobStoreAdapter = { get: async (key: string): Promise => { try { const data = await this.readObjectFromPath(`_cas/${key}`) if (data === null) { return undefined } // Unwraps binary data stored as {_binary: true, data: "base64..."} // — hash verification must run on the original content bytes. return unwrapBinaryData(data) } catch (error) { return undefined } }, put: async (key: string, data: Buffer): Promise => { // Metadata keys are always JSON; blob keys are always binary. The key // format decides — no content sniffing (JSON.parse guessing corrupted // compressed blobs that happened to parse as JSON). const obj = key.includes('-meta:') ? JSON.parse(data.toString()) : { _binary: true, data: data.toString('base64') } await this.writeObjectToPath(`_cas/${key}`, obj) }, delete: async (key: string): Promise => { try { await this.deleteObjectFromPath(`_cas/${key}`) } catch (error) { // Ignore if doesn't exist } }, list: async (prefix: string): Promise => { try { // Keys are stored as files like `_cas/blob:`, so listing is // prefix filtering over the `_cas/` area with the area prefix // stripped from the returned keys. const allPaths = await this.listObjectsUnderPath('_cas/') return allPaths .map(p => p.replace(/^_cas\//, '')) .filter(key => key.startsWith(prefix)) } catch (error: any) { // `_cas/` doesn't exist yet — empty store. return [] } } } this.blobStorage = new BlobStorage(casAdapter) } /** * @description Write a canonical object (write-cache coherent). The object * lands in the write-through cache before the asynchronous write starts, so * {@link readCanonicalObject} returns it immediately — read-after-write * consistency within the process. The cache persists until `flush()`. * * @param path - Storage-root-relative object path. * @param data - JSON-serializable object to write. * @protected */ protected async writeCanonicalObject(path: string, data: any): Promise { this.writeCache.set(path, data) // Use write buffer if available, otherwise write directly (filesystem) if (this.metadataWriteBuffer) { await this.metadataWriteBuffer.write(path, data) } else { await this.writeObjectToPath(path, data) } // Cache is NOT cleared here - persists until flush() // This provides a safety net for immediate queries after batch writes } /** * @description Read a canonical object: write cache first (synchronous, * guarantees read-after-write consistency), then the adapter. * * @param path - Storage-root-relative object path. * @returns The object, or `null` when absent. * @protected */ protected async readCanonicalObject(path: string): Promise { const cachedData = this.writeCache.get(path) if (cachedData !== undefined) { return cachedData } return this.readObjectFromPath(path) } /** * @description Delete a canonical object, evicting the write-cache entry * first so subsequent reads never return stale cached data. * * @param path - Storage-root-relative object path. * @protected */ protected async deleteCanonicalObject(path: string): Promise { this.writeCache.delete(path) return this.deleteObjectFromPath(path) } /** * @description List canonical objects under a storage-root-relative prefix. * * @param prefix - Storage-root-relative directory prefix. * @returns Storage-root-relative paths of the objects found. * @protected */ protected async listCanonicalObjects(prefix: string): Promise { return this.listObjectsUnderPath(prefix) } // ============================================================================ // GENERATIONAL RECORD LAYER PRIMITIVES (8.0 MVCC) // // The narrow surface `GenerationStore` (src/db/generationStore.ts) needs from // an adapter — see the `GenerationStorage` contract in src/db/types.ts. // // Raw-object methods operate on storage-root-relative paths and deliberately // bypass the write cache: the record layer owns the `_system/` + // `_generations/` areas outright. Entity-raw methods, by contrast, go // through the write-cache-coherent canonical helpers so before-images // capture exactly the bytes the live read paths see. // ============================================================================ /** * Hook invoked after every entity-visible single-operation write (noun/verb * metadata save or delete). Registered by the generation store so * `brain.generation()` advances on writes performed outside `transact()`. * See {@link BaseStorage.setGenerationBumpHook}. */ protected generationBumpHook?: () => void /** * Register (or detach, with `undefined`) the generation-bump hook. * * The hook fires once per entity-visible metadata mutation — noun/verb * metadata saves and deletes, the one storage write every logical Brainy * mutation (`add`, `update`, `remove`, `relate`, `updateRelation`, * `unrelate`) performs exactly once per entity it touches. It does NOT fire * for derived-index writes (HNSW node data, metadata-index chunks, LSM * segments, statistics), so the generation counter tracks *data* mutations, * not index maintenance. The counter is a monotonic watermark, not an * operation count: a cascade delete bumps once per removed record. * * @param hook - Callback invoked synchronously after each qualifying write, * or `undefined` to detach. */ public setGenerationBumpHook(hook: (() => void) | undefined): void { this.generationBumpHook = hook } /** * Read a raw object at a storage-root-relative path. Bypasses the write * cache (record-layer files are written through * {@link BaseStorage.writeRawObject} only). * * @param path - Storage-root-relative object path (e.g. `_system/manifest.json`). * @returns The parsed object, or `null` if absent. */ public async readRawObject(path: string): Promise { await this.ensureInitialized() return this.readObjectFromPath(path) } /** * Write a raw object at a storage-root-relative path. On disk this is an * atomic tmp+rename (the filesystem adapter's primitive), which is what * makes the manifest rename a valid commit point. * * @param path - Storage-root-relative object path. * @param data - JSON-serializable object to persist. */ public async writeRawObject(path: string, data: any): Promise { await this.ensureInitialized() await this.writeObjectToPath(path, data) } /** * Delete a raw object at a storage-root-relative path (no-op if absent). * * @param path - Storage-root-relative object path. */ public async deleteRawObject(path: string): Promise { await this.ensureInitialized() await this.deleteObjectFromPath(path) } /** * List raw object paths under a storage-root-relative prefix (normalized, * `.gz`-stripped — the adapter primitives already normalize). * * @param prefix - Storage-root-relative directory prefix. * @returns Normalized object paths under the prefix (empty when none). */ public async listRawObjects(prefix: string): Promise { await this.ensureInitialized() return this.listObjectsUnderPath(prefix) } /** * Remove every object under a storage-root-relative prefix. The filesystem * adapter overrides this with a recursive directory removal; this default * lists and deletes individually (exactly what the in-memory adapter needs). * * @param prefix - Storage-root-relative directory prefix to remove. */ public async removeRawPrefix(prefix: string): Promise { await this.ensureInitialized() const paths = await this.listObjectsUnderPath(prefix) for (const p of paths) { await this.deleteObjectFromPath(p) } } /** * Durability barrier for the commit protocol: ensure the listed raw-object * paths are durable before the caller proceeds. The base implementation is * a no-op (in-memory writes are durable-by-definition within the process); * the filesystem adapter overrides it with real `fsync` of the files and * their parent directories. * * @param paths - Storage-root-relative object paths previously written via * {@link BaseStorage.writeRawObject}. */ public async syncRawObjects(paths: string[]): Promise { void paths } /** * Read an entity's raw stored objects — the exact bytes at its canonical * metadata + vector paths (write-cache coherent). Used by the generation * store to capture before-images. * * @param id - The entity id. * @returns The raw stored metadata and vector objects (`null` per part when * the corresponding file is absent). */ public async readNounRaw(id: string): Promise<{ metadata: any | null; vector: any | null }> { await this.ensureInitialized() const [metadata, vector] = await Promise.all([ this.readCanonicalObject(getNounMetadataPath(id)), this.readCanonicalObject(getNounVectorPath(id)) ]) return { metadata: metadata ?? null, vector: vector ?? null } } /** * Restore an entity's raw stored objects byte-for-byte (a `null` part * deletes that file). Used by crash recovery and transaction aborts to * restore before-images. * * Bypasses the statistics/count bookkeeping of the normal save paths on * purpose: restores must reproduce the exact prior bytes, and the count * rollups are derived state with their own rebuild paths * (`rebuildTypeCounts()` / `rebuildSubtypeCounts()`). * * @param id - The entity id. * @param record - Raw stored objects as returned by {@link BaseStorage.readNounRaw}. */ public async writeNounRaw(id: string, record: { metadata: any | null; vector: any | null }): Promise { await this.ensureInitialized() if (record.metadata === null) { await this.deleteCanonicalObject(getNounMetadataPath(id)) } else { await this.writeCanonicalObject(getNounMetadataPath(id), record.metadata) } if (record.vector === null) { await this.deleteCanonicalObject(getNounVectorPath(id)) } else { await this.writeCanonicalObject(getNounVectorPath(id), record.vector) } } /** * Read a relationship's raw stored objects (verb-side mirror of * {@link BaseStorage.readNounRaw}). * * @param id - The relationship id. * @returns The raw stored metadata and vector objects (`null` per part when * the corresponding file is absent). */ public async readVerbRaw(id: string): Promise<{ metadata: any | null; vector: any | null }> { await this.ensureInitialized() const [metadata, vector] = await Promise.all([ this.readCanonicalObject(getVerbMetadataPath(id)), this.readCanonicalObject(getVerbVectorPath(id)) ]) return { metadata: metadata ?? null, vector: vector ?? null } } /** * Restore a relationship's raw stored objects byte-for-byte (verb-side * mirror of {@link BaseStorage.writeNounRaw}; same bookkeeping caveats). * * @param id - The relationship id. * @param record - Raw stored objects as returned by {@link BaseStorage.readVerbRaw}. */ public async writeVerbRaw(id: string, record: { metadata: any | null; vector: any | null }): Promise { await this.ensureInitialized() if (record.metadata === null) { await this.deleteCanonicalObject(getVerbMetadataPath(id)) } else { await this.writeCanonicalObject(getVerbMetadataPath(id), record.metadata) } if (record.vector === null) { await this.deleteCanonicalObject(getVerbVectorPath(id)) } else { await this.writeCanonicalObject(getVerbVectorPath(id), record.vector) } } /** * Append one line to the transaction log (`_system/tx-log.jsonl`). The * filesystem adapter appends to a real JSONL file; the in-memory adapter * keeps a line array (serialized by `snapshotToDirectory`). * * @param line - One complete JSON document, without trailing newline. */ public abstract appendTxLogLine(line: string): Promise /** * Read all transaction-log lines, oldest first (empty array when no log * exists). A torn trailing line from a crashed append is returned as-is — * callers tolerate unparseable lines. */ public abstract readTxLogLines(): Promise /** * Snapshot the entire store into `targetPath`. The filesystem adapter * builds a hard-link farm (instant, space-shared — safe because data files * are immutable-by-rename); the in-memory adapter serializes its object * store to a filesystem-storage-compatible directory. The result is a * self-contained store openable via `Brainy.load(path)`. * * @param targetPath - Absolute directory path for the snapshot (created if * missing; must be empty or absent). */ public abstract snapshotToDirectory(targetPath: string): Promise /** * Replace the entire store's contents from a snapshot directory previously * produced by {@link BaseStorage.snapshotToDirectory}. Implementations * clear current contents (preserving live lock files), copy the snapshot * in (byte copy — never hard links, so the snapshot stays independent), * and then call {@link BaseStorage.reloadDerivedState}. * * @param sourcePath - Absolute path of the snapshot directory. */ public abstract restoreFromDirectory(sourcePath: string): Promise /** * Reset and reload every piece of adapter-internal derived state after the * underlying objects changed wholesale (restore-from-snapshot): the * write-through cache, type/subtype statistics, total counts, and the * graph-index singleton (invalidated so the next accessor rebuilds from the * restored verbs). Count attribution reads the canonical metadata record, so * there are no id-keyed caches to clear here. */ protected async reloadDerivedState(): Promise { this.clearWriteCache() this.nounCountsByType.fill(0) this.verbCountsByType.fill(0) this.subtypeCountsByType.clear() this.verbSubtypeCountsByType.clear() this.statisticsCache = null this.statisticsModified = false this.invalidateGraphIndex() // Re-create the blob store: a restore replaced the `_cas/` area wholesale, // and the old instance's LRU cache could serve blobs the restored store no // longer contains. if (this.blobStorage) { this.blobStorage = undefined await this.initializeBlobStorage() } await this.loadTypeStatistics() await this.loadSubtypeStatistics() await this.loadVerbSubtypeStatistics() await this.initializeCounts() } /** * Save a noun to storage (vector only, metadata saved separately) * @param noun Pure HNSW vector data (no metadata) */ public async saveNoun(noun: HNSWNoun): Promise { await this.ensureInitialized() // Save the HNSWNoun vector data only // Metadata must be saved separately via saveNounMetadata() await this.saveNoun_internal(noun) } /** * Hydrate a deserialized noun (pure HNSW vector data) with its stored flat * metadata record — THE canonical noun combine for every storage read path. * The record is split through `splitNounMetadataRecord` (single source of * truth: src/types/reservedFields.ts): reserved fields surface ONLY at * top level and `metadata` carries ONLY the consumer's custom fields. * Adding a combine site that bypasses this helper reintroduces the * reserved-field echo bug — don't. * * @param noun - The deserialized HNSW noun (id/vector/connections/level). * @param metadata - The stored flat metadata record (reserved + custom keys). * @returns The combined noun with reserved fields top-level, custom fields in `metadata`. */ protected hydrateNounWithMetadata( noun: HNSWNoun, metadata: Record | null | undefined ): HNSWNounWithMetadata { const { reserved, custom } = splitNounMetadataRecord(metadata) return { ...noun, // Standard fields at top-level type: (reserved.noun as NounType) || NounType.Thing, subtype: reserved.subtype as string | undefined, // visibility is a reserved top-level field (absent === 'public'). Surfacing it // here lets visibility-aware reads (export node streaming, count/find candidate // filters) see a noun's tier from getNouns without re-reading metadata — the // noun mirror of the verb-hydration fix. visibility: reserved.visibility as HNSWNounWithMetadata['visibility'], createdAt: normalizeStoredTimestamp(reserved.createdAt), updatedAt: normalizeStoredTimestamp(reserved.updatedAt), confidence: reserved.confidence as number | undefined, weight: reserved.weight as number | undefined, service: reserved.service as string | undefined, data: reserved.data as Record | undefined, createdBy: reserved.createdBy as HNSWNounWithMetadata['createdBy'], _rev: typeof reserved._rev === 'number' ? reserved._rev : 1, // Only custom user fields remain in metadata metadata: custom } } /** * Hydrate a deserialized verb (structural core) with its stored flat * metadata record — THE canonical verb combine, the relationship mirror of * {@link hydrateNounWithMetadata}. Splitting through * `splitVerbMetadataRecord` extracts `verb` too, so the type key never * echoes inside `metadata`. * * @param verb - The deserialized HNSW verb (id/vector/connections/verb/sourceId/targetId). * @param metadata - The stored flat metadata record (reserved + custom keys). * @returns The combined verb with reserved fields top-level, custom fields in `metadata`. */ protected hydrateVerbWithMetadata( verb: HNSWVerb, metadata: Record | null | undefined ): HNSWVerbWithMetadata { const { reserved, custom } = splitVerbMetadataRecord(metadata) return { ...verb, // Standard fields at top-level subtype: reserved.subtype as string | undefined, // visibility is a reserved top-level field (the verb mirror of Entity.visibility). // Surfacing it here lets the graph-index fast paths apply visibility filtering on // their already-hydrated results (so default related() stays O(degree) instead of // falling through to a full scan), and lets related({ includeInternal }) results // actually report which edges are internal. visibility: reserved.visibility as HNSWVerbWithMetadata['visibility'], createdAt: normalizeStoredTimestamp(reserved.createdAt), updatedAt: normalizeStoredTimestamp(reserved.updatedAt), confidence: reserved.confidence as number | undefined, weight: reserved.weight as number | undefined, service: reserved.service as string | undefined, data: reserved.data as Record | undefined, createdBy: reserved.createdBy as HNSWVerbWithMetadata['createdBy'], // Only custom user fields remain in metadata metadata: custom } } /** * @description Apply the metadata-derived verb filters (`subtype`, * `excludeVisibility`) that the graph-index fast paths can satisfy on their * already-hydrated O(degree) / O(type) candidate set — matching the semantics * of the full-scan fallback in {@link getVerbsWithPagination}. This is what * keeps default `related({ from/to })` (which always sets the visibility * exclusion) on the fast adjacency path instead of forcing a full O(E) scan. * * - `subtype`: keeps only verbs carrying a matching subtype; a verb with no * subtype is excluded when a subtype filter is set (same as the fallback). * - `excludeVisibility`: drops verbs whose stored visibility tier is in the * excluded set; an absent tier is `'public'` and is always kept. * * @param verbs - Hydrated candidate verbs from a fast-path lookup. * @param filter - The `getVerbs` filter (only `subtype` / `excludeVisibility` * are read here; the structural match was already done by the caller). * @returns The candidates with the metadata filters applied (input order preserved). */ protected applyVerbMetadataFilters( verbs: HNSWVerbWithMetadata[], filter?: { subtype?: string | string[] excludeVisibility?: string[] } ): HNSWVerbWithMetadata[] { let out = verbs const subtypeFilter = filter?.subtype if (subtypeFilter) { const allowed = new Set(Array.isArray(subtypeFilter) ? subtypeFilter : [subtypeFilter]) out = out.filter((v) => v.subtype !== undefined && allowed.has(v.subtype)) } const excludeVisibility = filter?.excludeVisibility if (excludeVisibility && excludeVisibility.length > 0) { const excluded = new Set(excludeVisibility) out = out.filter((v) => !(v.visibility && excluded.has(v.visibility))) } return out } /** * Get a noun from storage (returns combined HNSWNounWithMetadata) * @param id Entity ID * @returns Combined vector + metadata or null */ public async getNoun(id: string): Promise { await this.ensureInitialized() // Load vector and metadata separately const vector = await this.getNoun_internal(id) if (!vector) { return null } // Load metadata const metadata = await this.getNounMetadata(id) if (!metadata) { prodLog.warn(`[Storage] Noun ${id} has vector but no metadata - this should not happen`) return null } return this.hydrateNounWithMetadata(vector, metadata) } /** * Get nouns by noun type * @param nounType The noun type to filter by * @returns Promise that resolves to an array of nouns of the specified noun type */ public async getNounsByNounType(nounType: string): Promise { await this.ensureInitialized() // Internal method returns HNSWNoun[], need to combine with metadata const nouns = await this.getNounsByNounType_internal(nounType) // Combine each noun with its metadata via the canonical hydration helper const nounsWithMetadata: HNSWNounWithMetadata[] = [] for (const noun of nouns) { const metadata = await this.getNounMetadata(noun.id) if (metadata) { nounsWithMetadata.push(this.hydrateNounWithMetadata(noun, metadata)) } } return nounsWithMetadata } /** * Delete a noun from storage */ public async deleteNoun(id: string): Promise { await this.ensureInitialized() // Delete both the vector file and metadata file (2-file system) await this.deleteNoun_internal(id) // Delete metadata file (if it exists) try { await this.deleteNounMetadata(id) } catch (error) { // Ignore if metadata file doesn't exist prodLog.debug(`No metadata file to delete for noun ${id}`) } } /** * Save a verb to storage (verb only, metadata saved separately) * * @param verb Pure HNSW verb with core relational fields (verb, sourceId, targetId) */ public async saveVerb(verb: HNSWVerb): Promise { await this.ensureInitialized() // Validate verb type before saving - storage boundary protection validateVerbType(verb.verb) // Save the HNSWVerb vector and core fields only // Metadata must be saved separately via saveVerbMetadata() await this.saveVerb_internal(verb) } /** * Get a verb from storage (returns combined HNSWVerbWithMetadata) * @param id Entity ID * @returns Combined verb + metadata or null */ public async getVerb(id: string): Promise { await this.ensureInitialized() // Load verb vector and core fields const verb = await this.getVerb_internal(id) if (!verb) { return null } // Load metadata const metadata = await this.getVerbMetadata(id) if (!metadata) { prodLog.warn(`[Storage] Verb ${id} has vector but no metadata - this should not happen`) return null } return this.hydrateVerbWithMetadata(verb, metadata) } /** * Batch get multiple verbs * * **Performance**: Eliminates N+1 pattern for verb loading * - Current: N × getVerb() = N × 50ms on GCS = 250ms for 5 verbs * - Batched: 1 × getVerbsBatch() = 1 × 50ms on GCS = 50ms (**5x faster**) * * **Use cases:** * - graphIndex.getVerbsBatchCached() for relate() duplicate checking * - Loading relationships in batch operations * - Pre-loading verbs for graph traversal * * @param ids Array of verb IDs to fetch * @returns Map of id → HNSWVerbWithMetadata (only successful reads included) * */ public async getVerbsBatch(ids: string[]): Promise> { await this.ensureInitialized() const results = new Map() if (ids.length === 0) return results // Batch-fetch vectors and metadata in parallel // Build paths for vectors const vectorPaths: Array<{ path: string; id: string }> = ids.map(id => ({ path: getVerbVectorPath(id), id })) // Build paths for metadata const metadataPaths: Array<{ path: string; id: string }> = ids.map(id => ({ path: getVerbMetadataPath(id), id })) // Batch read vectors and metadata in parallel const [vectorResults, metadataResults] = await Promise.all([ this.readCanonicalObjectBatch(vectorPaths.map(p => p.path)), this.readCanonicalObjectBatch(metadataPaths.map(p => p.path)) ]) // Combine vectors + metadata into HNSWVerbWithMetadata for (const { path: vectorPath, id } of vectorPaths) { const vectorData = vectorResults.get(vectorPath) const metadataPath = getVerbMetadataPath(id) const metadataData = metadataResults.get(metadataPath) if (vectorData && metadataData) { // Deserialize, then combine via the canonical hydration helper const verb = this.deserializeVerb(vectorData) results.set(id, this.hydrateVerbWithMetadata(verb, metadataData)) } } return results } /** * Internal method for loading all verbs - used by performance optimizations * @internal - Do not use directly, use getVerbs() with pagination instead */ protected async _loadAllVerbsForOptimization(): Promise { await this.ensureInitialized() // Only use this for internal optimizations when safe const result = await this.getVerbs({ pagination: { limit: Number.MAX_SAFE_INTEGER } }) // Convert HNSWVerbWithMetadata to HNSWVerb (strip metadata) const hnswVerbs: HNSWVerb[] = result.items.map(verbWithMetadata => ({ id: verbWithMetadata.id, vector: verbWithMetadata.vector, connections: verbWithMetadata.connections, verb: verbWithMetadata.verb, sourceId: verbWithMetadata.sourceId, targetId: verbWithMetadata.targetId })) return hnswVerbs } /** * Get verbs by source */ public async getVerbsBySource(sourceId: string): Promise { await this.ensureInitialized() // CRITICAL: Fetch ALL verbs for this source, not just first page // This is needed for delete operations to clean up all relationships const result = await this.getVerbs({ filter: { sourceId }, pagination: { limit: Number.MAX_SAFE_INTEGER } }) return result.items } /** * Get verbs by target */ public async getVerbsByTarget(targetId: string): Promise { await this.ensureInitialized() // CRITICAL: Fetch ALL verbs for this target, not just first page // This is needed for delete operations to clean up all relationships const result = await this.getVerbs({ filter: { targetId }, pagination: { limit: Number.MAX_SAFE_INTEGER } }) return result.items } /** * Get verbs by type */ public async getVerbsByType(type: string): Promise { await this.ensureInitialized() // Fetch ALL verbs of this type (no pagination limit) const result = await this.getVerbs({ filter: { verbType: type }, pagination: { limit: Number.MAX_SAFE_INTEGER } }) return result.items } /** * Internal method for loading all nouns - used by performance optimizations * @internal - Do not use directly, use getNouns() with pagination instead */ protected async _loadAllNounsForOptimization(): Promise { await this.ensureInitialized() // Only use this for internal optimizations when safe const result = await this.getNouns({ pagination: { limit: Number.MAX_SAFE_INTEGER } }) return result.items } /** * Get nouns with pagination and filtering * @param options Pagination and filtering options * @returns Promise that resolves to a paginated result of nouns */ public async getNouns(options?: { pagination?: { offset?: number limit?: number cursor?: string } filter?: { nounType?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWNounWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string }> { await this.ensureInitialized() // Set default pagination values const pagination = options?.pagination || {} const limit = pagination.limit || 100 const offset = pagination.offset || 0 const cursor = pagination.cursor // Optimize for common filter cases to avoid loading all nouns if (options?.filter) { // If filtering by nounType only, use the optimized method if ( options.filter.nounType && !options.filter.service && !options.filter.metadata ) { const nounType = Array.isArray(options.filter.nounType) ? options.filter.nounType[0] : options.filter.nounType // Get nouns by type directly (already combines with metadata) const nounsByType = await this.getNounsByNounType(nounType) // Apply pagination const paginatedNouns = nounsByType.slice(offset, offset + limit) const hasMore = offset + limit < nounsByType.length // Set next cursor if there are more items let nextCursor: string | undefined = undefined if (hasMore && paginatedNouns.length > 0) { const lastItem = paginatedNouns[paginatedNouns.length - 1] nextCursor = lastItem.id } return { items: paginatedNouns, totalCount: nounsByType.length, hasMore, nextCursor } } } // For more complex filtering or no filtering, use a paginated approach // that avoids loading all nouns into memory at once try { // First, try to get a count of total nouns (if the adapter supports it) let totalCount: number | undefined = undefined try { // This is an optional method that adapters may implement (duck-typed — // see OptionalCountCapabilities) const adapter = this as BaseStorage & OptionalCountCapabilities if (typeof adapter.countNouns === 'function') { totalCount = await adapter.countNouns(options?.filter) } } catch (countError) { // Ignore errors from count method, it's optional prodLog.warn('Error getting noun count:', countError) } // Check if the adapter has a paginated method for getting nouns if (typeof this.getNounsWithPagination === 'function') { // Use the adapter's paginated method - pass offset directly to adapter. // The annotation widens totalCount to optional: adapter overrides // follow BaseStorageAdapter's contract, where totalCount may be absent. const result: { items: HNSWNounWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string } = await this.getNounsWithPagination({ limit, offset, // Let the adapter handle offset for O(1) operation cursor, filter: options?.filter }) // Don't slice here - the adapter should handle offset efficiently const items = result.items // CRITICAL SAFETY CHECK: Prevent infinite loops // If we have no items but hasMore is true, force hasMore to false // This prevents pagination bugs from causing infinite loops const safeHasMore = items.length > 0 ? result.hasMore : false // VALIDATION: Ensure adapter returns totalCount (prevents restart bugs) // If adapter forgets to return totalCount, log warning and use pre-calculated count let finalTotalCount = result.totalCount || totalCount if (result.totalCount === undefined && this.totalNounCount > 0) { prodLog.warn( `⚠️ Storage adapter missing totalCount in getNounsWithPagination result! ` + `Using pre-calculated count (${this.totalNounCount}) as fallback. ` + `Please ensure your storage adapter returns totalCount: this.totalNounCount` ) finalTotalCount = this.totalNounCount } return { items, totalCount: finalTotalCount, hasMore: safeHasMore, nextCursor: result.nextCursor } } // Storage adapter does not support pagination. This is a hard // misconfiguration — find()/rebuild/aggregation all read through this path, // so returning an empty page would silently present a misconfigured adapter // as an empty database. Fail loud instead. throw BrainyError.storage( 'Storage adapter does not implement getNounsWithPagination(). The deprecated getAllNouns_internal() fallback has been removed; implement pagination in your storage adapter.' ) } catch (error) { // Never convert a genuine read failure into a success-shaped empty page: // this is the highest-fan-in read path (find() fallback, cold-start rebuild // gating, aggregation backfill, getNounsByNounType), so a swallowed error // would propagate as "zero rows" everywhere — indistinguishable from a truly // empty store, and the exact silent-failure class the 8.0 contract forbids. if (error instanceof BrainyError) throw error throw BrainyError.storage( 'getNouns pagination read failed', error instanceof Error ? error : undefined ) } } /** * Get nouns with pagination (Type-first implementation) * * CRITICAL: This method is required for brain.find() to work! * Iterates through noun types with billion-scale optimizations. * * ARCHITECTURE: Reads storage directly (not indexes) to avoid circular dependencies. * Storage → Indexes (one direction only). GraphAdjacencyIndex built FROM storage. * * OPTIMIZATIONS: * - Skip empty types using nounCountsByType[] tracking (O(1) check) * - Early termination when offset + limit entities collected * - Memory efficient: Never loads full dataset */ public override async getNounsWithPagination(options: { limit: number offset: number cursor?: string // Opaque resume token from a prior page's nextCursor; when set it supersedes offset (O(N) walk, no re-scan) filter?: { nounType?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWNounWithMetadata[] totalCount: number hasMore: boolean nextCursor?: string }> { await this.ensureInitialized() const { limit, offset = 0, filter } = options // Cursor (8.0): resume token carrying the (shard, nounId) of the last returned // noun — the noun mirror of getVerbsWithPagination. When present it supersedes // `offset` and resumes the shard walk immediately AFTER that position, so a full // walk is O(N) instead of the O(N²) of offset paging. Malformed/foreign tokens // decode to null → offset fallback. (Previously the cursor was ignored, which // was latent — the only multi-page consumer used a single big page — until small // chunk sizes needed page 2 and an offset-0-on-every-call walk never terminated.) const cursor = this.decodeNounWalkCursor(options.cursor) const collected: Array<{ noun: HNSWNounWithMetadata; shard: number }> = [] // Peek one past the window so `hasMore` is decidable. Cursor mode collects one // page (+1); offset mode keeps the full [0, offset+limit] window (+1). const peekCount = cursor ? limit + 1 : offset + limit + 1 const startShard = cursor ? cursor.shard : 0 // Iterate by shards (0x00-0xFF), early-terminating at peekCount. for (let shard = startShard; shard < 256 && collected.length < peekCount; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/nouns/${shardHex}` try { const nounFiles = await this.listCanonicalObjects(shardDir) // Stable within-shard order (by noun id) so offset windows and cursor resume // are deterministic; ids come from the path so skipped nouns are never read. const entries = nounFiles .filter((p) => p.includes('/vectors.json')) .map((p) => ({ path: p, id: idFromVectorPath(p) })) .sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) for (const { path: nounPath, id: nounId } of entries) { if (collected.length >= peekCount) break // Resume: in the cursor's own shard, skip up to AND INCLUDING the cursor id. if (cursor && shard === cursor.shard && nounId <= cursor.id) continue try { const noun = await this.readCanonicalObject(nounPath) if (!noun) continue const deserialized = this.deserializeNoun(noun) const metadata = await this.getNounMetadata(deserialized.id) if (!metadata) continue // Apply type filter if (filter?.nounType && metadata.noun) { const types = Array.isArray(filter.nounType) ? filter.nounType : [filter.nounType] if (!types.includes(metadata.noun)) { continue } } // Apply service filter if (filter?.service) { const services = Array.isArray(filter.service) ? filter.service : [filter.service] if (metadata.service && !services.includes(metadata.service)) { continue } } // Combine noun + metadata via the canonical hydration helper — // reserved fields top-level, ONLY custom fields in `metadata`. collected.push({ noun: this.hydrateNounWithMetadata(deserialized, metadata), shard }) } catch (error) { // Skip nouns that fail to load } } } catch (error) { // Skip shards that have no data } } // Window selection. Cursor mode already starts at the resume point (window // [0, limit)); offset mode slices [offset, offset+limit). The peeked extra // entry (if any) is dropped — its existence is exactly what makes hasMore true. const windowStart = cursor ? 0 : offset const pagePairs = collected.slice(windowStart, windowStart + limit) const paginatedNouns = pagePairs.map((p) => p.noun) const hasMore = collected.length > windowStart + limit // totalCount must be the TRUE dataset total, not this peeked page. For the // unfiltered case the authoritative total is the O(1) counter maintained on // every add/delete (rehydrated on init); `Math.max` guards a stale counter. A // filtered scan has no cheap exact total, so it keeps the collected length. const totalCount = filter ? collected.length : Math.max(this.totalNounCount, collected.length) // nextCursor = the (shard, id) of the last RETURNED noun, so the next call // resumes immediately after it (works for both cursor and offset callers). let nextCursor: string | undefined = undefined if (hasMore && pagePairs.length > 0) { const lastPair = pagePairs[pagePairs.length - 1] nextCursor = this.encodeNounWalkCursor(lastPair.shard, lastPair.noun.id) } return { items: paginatedNouns, totalCount, hasMore, nextCursor } } /** * @description Encode a noun-walk resume cursor — the `(shard, nounId)` of the * last returned noun — as an opaque, version-tagged token (`cn1:` prefix lets * {@link decodeNounWalkCursor} reject foreign tokens, e.g. a bare-id cursor). * `nounId` is placed last and the decoder re-joins on `:` so any id format survives. * @param shard - The shard (0–255) the noun lives in. * @param id - The noun id. * @returns The opaque cursor token. */ private encodeNounWalkCursor(shard: number, id: string): string { return `cn1:${shard}:${id}` } /** * @description Decode a noun-walk cursor from {@link encodeNounWalkCursor}; * returns `null` for an absent / malformed / foreign token (caller falls back * to offset paging rather than mis-resuming). * @param cursor - The opaque cursor token, or undefined. * @returns `{ shard, id }` resume position, or `null`. */ private decodeNounWalkCursor(cursor?: string): { shard: number; id: string } | null { if (!cursor) return null const parts = cursor.split(':') if (parts.length < 3 || parts[0] !== 'cn1') return null const shard = Number(parts[1]) if (!Number.isInteger(shard) || shard < 0 || shard > 255) return null const id = parts.slice(2).join(':') if (id.length === 0) return null return { shard, id } } /** * Get verbs with pagination (Type-first implementation with billion-scale optimizations) * * CRITICAL: This method is required for brain.related() to work! * Iterates through verb types with the same optimizations as nouns. * * ARCHITECTURE: Reads storage directly (not indexes) to avoid circular dependencies. * Storage → Indexes (one direction only). GraphAdjacencyIndex built FROM storage. * * OPTIMIZATIONS: * - Skip empty types using verbCountsByType[] tracking (O(1) check) * - Early termination when offset + limit verbs collected * - Memory efficient: Never loads full dataset * - Inline filtering for sourceId, targetId, verbType */ public override async getVerbsWithPagination(options: { limit: number offset: number cursor?: string // Opaque resume token from a prior page's nextCursor; when set it supersedes offset (O(N) walk, no re-scan) filter?: { verbType?: string | string[] sourceId?: string | string[] targetId?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWVerbWithMetadata[] totalCount: number hasMore: boolean nextCursor?: string }> { await this.ensureInitialized() const { limit, offset = 0, filter } = options // Cursor (8.0): an opaque resume token (see encodeVerbWalkCursor) carrying the // (shard, verbId) of the last returned verb. When present it SUPERSEDES `offset` // and resumes the shard walk immediately AFTER that position, so a full walk is // O(N) total instead of the O(N²) of offset paging (which re-scans from shard 0 // every page). Malformed / foreign tokens decode to null → offset fallback. const cursor = this.decodeVerbWalkCursor(options.cursor) // Each collected entry remembers its shard so nextCursor can point at the exact // (shard, id) resume position. const collected: Array<{ verb: HNSWVerbWithMetadata; shard: number }> = [] // Prepare filter sets for efficient lookup const filterVerbTypes = filter?.verbType ? new Set(Array.isArray(filter.verbType) ? filter.verbType : [filter.verbType]) : null const filterSourceIds = filter?.sourceId ? new Set(Array.isArray(filter.sourceId) ? filter.sourceId : [filter.sourceId]) : null const filterTargetIds = filter?.targetId ? new Set(Array.isArray(filter.targetId) ? filter.targetId : [filter.targetId]) : null // `subtype` rides alongside the declared filter fields (Brainy's // related() path passes it through); it's applied after metadata // loads below, since subtype lives in verb metadata, not on the raw verb. const subtypeFilterValue = ( filter as { verbType?: string | string[]; subtype?: string | string[] } | undefined )?.subtype const filterSubtypes = subtypeFilterValue ? new Set( Array.isArray(subtypeFilterValue) ? subtypeFilterValue : [subtypeFilterValue] ) : null // 8.0 visibility exclusion — applied after metadata load (verb visibility lives in // metadata), so hidden edges are skipped BEFORE the pagination window fills. const excludeVisibility = (filter as { excludeVisibility?: string[] } | undefined)?.excludeVisibility const filterExcludeVisibility = excludeVisibility && excludeVisibility.length > 0 ? new Set(excludeVisibility) : null // Peek one past the window so hasMore is decidable. Cursor mode collects exactly // one page (+1); offset mode keeps the full [0, offset+limit] window (+1). const peekCount = cursor ? limit + 1 : offset + limit + 1 // Cursor resume skips every shard BEFORE the cursor's shard outright (the core of // the O(N) win); offset mode always starts at shard 0. const startShard = cursor ? cursor.shard : 0 // Iterate by shards (0x00-0xFF) — single pass, early-terminating at peekCount. for (let shard = startShard; shard < 256 && collected.length < peekCount; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { const verbFiles = await this.listCanonicalObjects(shardDir) // Stable within-shard order (by verb id) so offset windows and cursor resume // are deterministic and consistent across calls. Ids come from the path, so // verbs skipped by the cursor are never read. const entries = verbFiles .filter((p) => p.includes('/vectors.json')) .map((p) => ({ path: p, id: idFromVectorPath(p) })) .sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) for (const { path: verbPath, id: verbId } of entries) { if (collected.length >= peekCount) break // Resume: in the cursor's own shard, skip up to AND INCLUDING the cursor id // (later shards are processed in full). No read for skipped verbs. if (cursor && shard === cursor.shard && verbId <= cursor.id) continue try { const rawVerb = await this.readCanonicalObject(verbPath) if (!rawVerb) continue // Deserialize connections Map from JSON storage format const verb = this.deserializeVerb(rawVerb) // Apply type filter if (filterVerbTypes && !filterVerbTypes.has(verb.verb)) { continue } // Apply sourceId filter if (filterSourceIds && !filterSourceIds.has(verb.sourceId)) { continue } // Apply targetId filter if (filterTargetIds && !filterTargetIds.has(verb.targetId)) { continue } // Load metadata const metadata = await this.getVerbMetadata(verb.id) // Apply subtype filter (requires metadata — checked AFTER load) if (filterSubtypes) { const subtype = metadata?.subtype as string | undefined if (!subtype || !filterSubtypes.has(subtype)) { continue } } // Apply visibility exclusion (8.0). Absent === 'public' (kept); a stored // 'internal'/'system' value is dropped when its tier is excluded. if (filterExcludeVisibility) { const visibility = metadata?.visibility as string | undefined if (visibility && filterExcludeVisibility.has(visibility)) { continue } } // Combine verb + metadata via the canonical hydration helper — // reserved fields top-level, ONLY custom fields in `metadata`. collected.push({ verb: this.hydrateVerbWithMetadata(verb, metadata), shard }) } catch (error) { // Skip verbs that fail to load } } } catch (error) { // Skip shards that have no data } } // Window selection. Cursor mode already starts at the resume point, so its window // is [0, limit); offset mode slices [offset, offset+limit). The peeked extra entry // (if any) is dropped here — its existence is exactly what makes hasMore true. const windowStart = cursor ? 0 : offset const pagePairs = collected.slice(windowStart, windowStart + limit) const paginatedVerbs = pagePairs.map((p) => p.verb) const hasMore = collected.length > windowStart + limit // totalCount must be the TRUE dataset total, not this peeked page. For the // unfiltered scan the authoritative total is the O(1) `totalVerbCount` counter // (isNew-gated, visibility-filtered, rehydrated on init); `Math.max` guards a // stale counter from under-reporting. A filtered scan has no cheap exact total, // so it keeps the collected length (a lower bound). const totalCount = filter ? collected.length : Math.max(this.totalVerbCount, collected.length) // nextCursor encodes the (shard, id) of the LAST RETURNED verb so the next call // resumes immediately after it — for both cursor and offset callers (an offset // caller can switch to cursor paging to escape the O(N²)). let nextCursor: string | undefined = undefined if (hasMore && pagePairs.length > 0) { const lastPair = pagePairs[pagePairs.length - 1] nextCursor = this.encodeVerbWalkCursor(lastPair.shard, lastPair.verb.id) } return { items: paginatedVerbs, totalCount, hasMore, nextCursor } } /** * @description Encode a verb-walk resume cursor — the `(shard, verbId)` of the * last returned verb — as an opaque, version-tagged token. The `cv1:` prefix * lets {@link decodeVerbWalkCursor} reject foreign tokens (e.g. the bare-id * cursors the graph-index fast paths emit). `verbId` is placed last and the * decoder re-joins on `:` so any id format survives the round-trip. * @param shard - The shard (0–255) the verb lives in. * @param id - The verb id. * @returns The opaque cursor token. */ private encodeVerbWalkCursor(shard: number, id: string): string { return `cv1:${shard}:${id}` } /** * @description Decode a verb-walk cursor produced by {@link encodeVerbWalkCursor}. * Returns `null` for an absent, malformed, or foreign token so the caller falls * back to offset-based paging rather than mis-resuming. * @param cursor - The opaque cursor token, or undefined. * @returns `{ shard, id }` resume position, or `null`. */ private decodeVerbWalkCursor(cursor?: string): { shard: number; id: string } | null { if (!cursor) return null const parts = cursor.split(':') if (parts.length < 3 || parts[0] !== 'cv1') return null const shard = Number(parts[1]) if (!Number.isInteger(shard) || shard < 0 || shard > 255) return null const id = parts.slice(2).join(':') if (id.length === 0) return null return { shard, id } } /** * Get verbs with pagination and filtering * @param options Pagination and filtering options * @returns Promise that resolves to a paginated result of verbs */ public async getVerbs(options?: { pagination?: { offset?: number limit?: number cursor?: string } filter?: { verbType?: string | string[] sourceId?: string | string[] targetId?: string | string[] service?: string | string[] metadata?: Record /** * 8.0 visibility: tiers to exclude from results (e.g. `['internal','system']`). * Applied after metadata load in the full scan, so it disqualifies the * metadata-less graph-index fast paths (same as `subtype`). */ excludeVisibility?: string[] } }): Promise<{ items: HNSWVerbWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string }> { await this.ensureInitialized() // Set default pagination values const pagination = options?.pagination || {} const limit = pagination.limit || 100 const offset = pagination.offset || 0 const cursor = pagination.cursor // Optimize for common filter cases to avoid loading all verbs. // The graph-index fast paths (getVerbsBySource/Target/Type_internal) hydrate // each candidate's metadata, so `subtype` and the 8.0 `excludeVisibility` // filters — both metadata fields — are applied on the small O(degree) / // O(type) candidate set via applyVerbMetadataFilters() BEFORE pagination. // (They used to disqualify the fast paths and force a full O(E) shard scan, // which made default related({ from/to }) — which always sets the visibility // exclusion — scan the entire graph per node.) An arbitrary `metadata` filter // still falls through, since each block guards `!options.filter.metadata`. if (options?.filter) { // CRITICAL VFS FIX: If filtering by sourceId + verbType (most common VFS pattern!) // This is the query PathResolver.getChildren() uses: related({ from: dirId, type: VerbType.Contains }) if ( options.filter.sourceId && options.filter.verbType && !options.filter.targetId && !options.filter.service && !options.filter.metadata ) { const sourceId = Array.isArray(options.filter.sourceId) ? options.filter.sourceId[0] : options.filter.sourceId const verbType = Array.isArray(options.filter.verbType) ? options.filter.verbType[0] : options.filter.verbType // 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( verbsBySource.filter(v => v.verb === verbType), options.filter ) // Apply pagination const paginatedVerbs = filteredVerbs.slice(offset, offset + limit) const hasMore = offset + limit < filteredVerbs.length // Set next cursor if there are more items let nextCursor: string | undefined = undefined if (hasMore && paginatedVerbs.length > 0) { const lastItem = paginatedVerbs[paginatedVerbs.length - 1] nextCursor = lastItem.id } return { items: paginatedVerbs, totalCount: filteredVerbs.length, hasMore, nextCursor } } // If filtering by sourceId only, use the optimized method if ( options.filter.sourceId && !options.filter.verbType && !options.filter.targetId && !options.filter.service && !options.filter.metadata ) { 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) const hasMore = offset + limit < verbsBySource.length // Set next cursor if there are more items let nextCursor: string | undefined = undefined if (hasMore && paginatedVerbs.length > 0) { const lastItem = paginatedVerbs[paginatedVerbs.length - 1] nextCursor = lastItem.id } return { items: paginatedVerbs, totalCount: verbsBySource.length, hasMore, nextCursor } } // If filtering by targetId only, use the optimized method if ( options.filter.targetId && !options.filter.verbType && !options.filter.sourceId && !options.filter.service && !options.filter.metadata ) { 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) const hasMore = offset + limit < verbsByTarget.length // Set next cursor if there are more items let nextCursor: string | undefined = undefined if (hasMore && paginatedVerbs.length > 0) { const lastItem = paginatedVerbs[paginatedVerbs.length - 1] nextCursor = lastItem.id } return { items: paginatedVerbs, totalCount: verbsByTarget.length, hasMore, nextCursor } } // If filtering by verbType only, use the optimized method if ( options.filter.verbType && !options.filter.sourceId && !options.filter.targetId && !options.filter.service && !options.filter.metadata ) { const verbType = Array.isArray(options.filter.verbType) ? options.filter.verbType[0] : options.filter.verbType // 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) const hasMore = offset + limit < verbsByType.length // Set next cursor if there are more items let nextCursor: string | undefined = undefined if (hasMore && paginatedVerbs.length > 0) { const lastItem = paginatedVerbs[paginatedVerbs.length - 1] nextCursor = lastItem.id } return { items: paginatedVerbs, totalCount: verbsByType.length, hasMore, nextCursor } } // Fast path for SINGLE sourceId + verbType combo (common VFS pattern) // This avoids the slow type-iteration fallback for VFS operations // NOTE: Only use fast path for single sourceId to avoid incomplete results const isSingleSourceId = options.filter.sourceId && !Array.isArray(options.filter.sourceId) if ( isSingleSourceId && options.filter.verbType && !options.filter.targetId && !options.filter.service && !options.filter.metadata ) { const sourceId = options.filter.sourceId as string const verbTypes = Array.isArray(options.filter.verbType) ? options.filter.verbType : [options.filter.verbType] prodLog.debug(`[BaseStorage] getVerbs: Using fast path for sourceId=${sourceId}, verbTypes=${verbTypes.join(',')}`) // Get verbs by source (uses GraphAdjacencyIndex if available) const verbsBySource = await this.getVerbsBySource_internal(sourceId) // Filter by verbType in memory (fast - usually small number of verbs per source), // then apply the subtype / visibility metadata filters on the candidate set. const filtered = this.applyVerbMetadataFilters( verbsBySource.filter(v => verbTypes.includes(v.verb)), options.filter ) // Apply pagination const paginatedVerbs = filtered.slice(offset, offset + limit) const hasMore = offset + limit < filtered.length // Set next cursor if there are more items let nextCursor: string | undefined = undefined if (hasMore && paginatedVerbs.length > 0) { const lastItem = paginatedVerbs[paginatedVerbs.length - 1] nextCursor = lastItem.id } prodLog.debug(`[BaseStorage] getVerbs: Fast path returned ${filtered.length} verbs (${paginatedVerbs.length} after pagination)`) return { items: paginatedVerbs, totalCount: filtered.length, hasMore, nextCursor } } } // For more complex filtering or no filtering, use a paginated approach // that avoids loading all verbs into memory at once try { // First, try to get a count of total verbs (if the adapter supports it) let totalCount: number | undefined = undefined try { // This is an optional method that adapters may implement (duck-typed — // see OptionalCountCapabilities) const adapter = this as BaseStorage & OptionalCountCapabilities if (typeof adapter.countVerbs === 'function') { totalCount = await adapter.countVerbs(options?.filter) } } catch (countError) { // Ignore errors from count method, it's optional prodLog.warn('Error getting verb count:', countError) } // Check if the adapter has a paginated method for getting verbs if (typeof this.getVerbsWithPagination === 'function') { // getVerbsWithPagination honors `offset` directly (it slices the // [offset, offset+limit) window; `cursor` is not yet implemented there). // Pass the real offset through. Previously offset was zeroed and smuggled // via an ignored `cursor`, so every page returned items [0, limit) and // offset was silently dropped (related({ offset }) paginated incorrectly). const result: { items: HNSWVerbWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string } = await this.getVerbsWithPagination({ limit, offset, cursor, filter: options?.filter }) const items = result.items // CRITICAL SAFETY CHECK: Prevent infinite loops // If we have no items but hasMore is true, force hasMore to false // This prevents pagination bugs from causing infinite loops const safeHasMore = items.length > 0 ? result.hasMore : false // VALIDATION: Ensure adapter returns totalCount (prevents restart bugs) // If adapter forgets to return totalCount, log warning and use pre-calculated count let finalTotalCount = result.totalCount || totalCount if (result.totalCount === undefined && this.totalVerbCount > 0) { prodLog.warn( `⚠️ Storage adapter missing totalCount in getVerbsWithPagination result! ` + `Using pre-calculated count (${this.totalVerbCount}) as fallback. ` + `Please ensure your storage adapter returns totalCount: this.totalVerbCount` ) finalTotalCount = this.totalVerbCount } return { items, totalCount: finalTotalCount, hasMore: safeHasMore, nextCursor: result.nextCursor } } // UNIVERSAL FALLBACK: Iterate through verb types with early termination (billion-scale safe) // This approach works for ALL storage adapters without requiring adapter-specific pagination prodLog.warn( 'Using universal type-iteration strategy for getVerbs(). ' + 'This works for all adapters but may be slower than native pagination. ' + 'For optimal performance at scale, storage adapters can implement getVerbsWithPagination().' ) const collectedVerbs: HNSWVerbWithMetadata[] = [] let totalScanned = 0 const targetCount = offset + limit // We need this many verbs total (including offset) // BUG FIX: Check if optimization should be used // Only use type-skipping optimization if counts are non-zero (reliable) const totalVerbCountFromArray = this.verbCountsByType.reduce((sum, c) => sum + c, 0) const useOptimization = totalVerbCountFromArray > 0 // BUG FIX: Pre-compute requested verb types to avoid skipping them // When a specific verbType filter is provided, we MUST check that type // even if verbCountsByType shows 0 (counts can be stale after restart) const requestedVerbTypes = options?.filter?.verbType const requestedVerbTypesSet = requestedVerbTypes ? new Set(Array.isArray(requestedVerbTypes) ? requestedVerbTypes : [requestedVerbTypes]) : null // Iterate through all 127 verb types (Stage 3 CANONICAL) with early termination // OPTIMIZATION: Skip types with zero count (only if counts are reliable) for (let i = 0; i < VERB_TYPE_COUNT && collectedVerbs.length < targetCount; i++) { const type = TypeUtils.getVerbFromIndex(i) // FIX: Never skip a type that's explicitly requested in the filter // This fixes VFS bug where Contains relationships were skipped after restart // when verbCountsByType[Contains] was 0 due to stale statistics const isRequestedType = requestedVerbTypesSet?.has(type) ?? false const countIsZero = this.verbCountsByType[i] === 0 // Skip empty types for performance (but only if optimization is enabled AND not requested) if (useOptimization && countIsZero && !isRequestedType) { continue } // Log when we DON'T skip a requested type that would have been skipped // This helps diagnose stale statistics issues in production if (useOptimization && countIsZero && isRequestedType) { prodLog.debug( `[BaseStorage] getVerbs: NOT skipping type=${type} despite count=0 (type was explicitly requested). ` + `Statistics may be stale - consider running rebuildTypeCounts().` ) } try { const verbsOfType = await this.getVerbsByType_internal(type) // Apply filtering inline (memory efficient) for (const verb of verbsOfType) { // Apply filters if specified if (options?.filter) { // Filter by sourceId if (options.filter.sourceId) { const sourceIds = Array.isArray(options.filter.sourceId) ? options.filter.sourceId : [options.filter.sourceId] if (!sourceIds.includes(verb.sourceId)) { continue } } // Filter by targetId if (options.filter.targetId) { const targetIds = Array.isArray(options.filter.targetId) ? options.filter.targetId : [options.filter.targetId] if (!targetIds.includes(verb.targetId)) { continue } } // Filter by verbType if (options.filter.verbType) { const verbTypes = Array.isArray(options.filter.verbType) ? options.filter.verbType : [options.filter.verbType] if (!verbTypes.includes(verb.verb)) { continue } } } // Verb passed filters - add to collection collectedVerbs.push(verb) // Early termination: stop when we have enough for offset + limit if (collectedVerbs.length >= targetCount) { break } } totalScanned += verbsOfType.length } catch (error) { // Ignore errors for types with no verbs (directory may not exist) // This is expected for types that haven't been used yet } } // Apply pagination (slice for offset) const paginatedVerbs = collectedVerbs.slice(offset, offset + limit) const hasMore = collectedVerbs.length > targetCount // Fixed >= to > (was causing infinite loop) return { items: paginatedVerbs, totalCount: collectedVerbs.length, // Accurate count of filtered results hasMore, nextCursor: hasMore && paginatedVerbs.length > 0 ? paginatedVerbs[paginatedVerbs.length - 1].id : undefined } } catch (error) { // Same no-silent-failure contract as getNouns: a genuine read failure must // surface as a named, catchable error, not a success-shaped empty page that // callers read as "this graph has no verbs." if (error instanceof BrainyError) throw error throw BrainyError.storage( 'getVerbs pagination read failed', error instanceof Error ? error : undefined ) } } /** * Delete a verb from storage */ public async deleteVerb(id: string): Promise { await this.ensureInitialized() // Delete both the vector file and metadata file (2-file system) await this.deleteVerb_internal(id) // Delete metadata file (if it exists) try { await this.deleteVerbMetadata(id) } catch (error) { // Ignore if metadata file doesn't exist prodLog.debug(`No metadata file to delete for verb ${id}`) } } /** * Get graph index (lazy initialization with concurrent access protection) * Fixed race condition where concurrent calls could trigger multiple rebuilds */ async getGraphIndex(): Promise { // If already initialized, return immediately if (this.graphIndex) { return this.graphIndex } // If initialization in progress, wait for it if (this.graphIndexPromise) { return this.graphIndexPromise } // Start initialization (only first caller reaches here) this.graphIndexPromise = this._initializeGraphIndex() try { const index = await this.graphIndexPromise return index } finally { // Clear promise after completion (success or failure) this.graphIndexPromise = undefined } } /** * Internal method to initialize graph index (called once by getGraphIndex) * @private */ private async _initializeGraphIndex(): Promise { prodLog.info('Initializing GraphAdjacencyIndex...') // Thread the shared entity-id resolver when already wired (re-init after // invalidateGraphIndex); on first init Brainy wires it right after. this.graphIndex = new GraphAdjacencyIndex(this, {}, this.graphEntityIdResolver) // Check if we need to rebuild from existing data const sampleVerbs = await this.getVerbs({ pagination: { limit: 1 } }) if (sampleVerbs.items.length > 0) { prodLog.info('Found existing verbs, rebuilding graph index...') await this.graphIndex.rebuild() } return this.graphIndex } /** * Clear all data from storage * This method should be implemented by each specific adapter */ public abstract override clear(): Promise /** * Get information about storage usage and capacity * This method should be implemented by each specific adapter */ public abstract override getStorageStatus(): Promise<{ type: string used: number quota: number | null details?: Record }> /** * Write a JSON object to a specific path in storage * This is a primitive operation that all adapters must implement * @param path - Full path including filename (e.g., "_system/statistics.json" or "entities/nouns/metadata/3f/3fa85f64-....json") * @param data - Data to write (will be JSON.stringify'd) * @protected */ protected abstract writeObjectToPath(path: string, data: any): Promise /** * Read a JSON object from a specific path in storage * This is a primitive operation that all adapters must implement * @param path - Full path including filename * @returns The parsed JSON object, or null if not found * @protected */ protected abstract readObjectFromPath(path: string): Promise /** * Delete an object from a specific path in storage * This is a primitive operation that all adapters must implement * @param path - Full path including filename * @protected */ protected abstract deleteObjectFromPath(path: string): Promise /** * List all object paths under a given prefix * This is a primitive operation that all adapters must implement * @param prefix - Directory prefix to list (e.g., "entities/nouns/metadata/3f/") * @returns Array of full paths * @protected */ protected abstract listObjectsUnderPath(prefix: string): Promise // Raw binary-blob primitive (saveBinaryBlob/loadBinaryBlob/deleteBinaryBlob/ // getBinaryBlobPath) is declared abstract on BaseStorageAdapter — the class // that `implements StorageAdapter` — alongside the other public storage // methods. Each concrete adapter implements it. See BaseStorageAdapter for the // contract and the shared `_blobs/.bin` key→location convention. /** * Save metadata to storage (now typed) * Routes to correct location (system or entity) based on key format */ public async saveMetadata(id: string, metadata: NounMetadata): Promise { await this.ensureInitialized() const keyInfo = this.analyzeKey(id, 'system') return this.writeCanonicalObject(keyInfo.fullPath, metadata) } /** * Get metadata from storage (now typed) * Routes to correct location (system or entity) based on key format */ public async getMetadata(id: string): Promise { await this.ensureInitialized() const keyInfo = this.analyzeKey(id, 'system') // Try the new shard-prefixed path first const data = await this.readCanonicalObject(keyInfo.fullPath) if (data !== null) return data // Backward compat: if the key was sharded, fall back to the legacy flat path if (keyInfo.shardId !== null) { const legacyPath = `${SYSTEM_DIR}/${id}.json` return this.readCanonicalObject(legacyPath) } return null } /** * Delete a system/metadata object previously written with {@link saveMetadata}. * The exact inverse of save/get: routes through `analyzeKey(id, 'system')` and * removes the canonical object (plus the legacy flat path that `getMetadata` * also reads, so a sharded key leaves nothing behind). Lets keyed-payload * owners — e.g. the LSM graph store reclaiming its compacted-away SSTables — * free storage instead of orphaning it. Idempotent: deleting a missing path is * a no-op. */ public async deleteMetadata(id: string): Promise { await this.ensureInitialized() const keyInfo = this.analyzeKey(id, 'system') await this.deleteCanonicalObject(keyInfo.fullPath) // Mirror getMetadata's legacy fallback so an older flat-path payload for a // sharded key is also reclaimed. if (keyInfo.shardId !== null) { await this.deleteCanonicalObject(`${SYSTEM_DIR}/${id}.json`) } } /** * Save noun metadata to storage (now typed) * Routes to correct sharded location based on UUID */ public async saveNounMetadata(id: string, metadata: NounMetadata): Promise { // Validate noun type in metadata - storage boundary protection validateNounType(metadata.noun) return this.saveNounMetadata_internal(id, metadata) } /** * Internal method for saving noun metadata (now typed) * Uses routing logic to handle both UUIDs (sharded) and system keys (unsharded) * * CRITICAL: Count synchronization happens here * This ensures counts are updated AFTER metadata exists, fixing the race condition * where storage adapters tried to read metadata before it was saved. * * @protected */ protected async saveNounMetadata_internal(id: string, metadata: NounMetadata): Promise { await this.ensureInitialized() // ID-first path - no type needed! const path = getNounMetadataPath(id) // Determine if this is a new entity by checking if metadata already exists const existingMetadata = await this.readCanonicalObject(path) const isNew = !existingMetadata // Save the metadata (write-cache coherent canonical write) await this.writeCanonicalObject(path, metadata) // Track subtype changes: on type or subtype change via update(), decrement // the prior bucket before incrementing the new one. The prior (type, subtype) // comes straight from the canonical record (`existingMetadata`, already loaded // above) — there is no id-keyed subtype cache. Symmetric with the delete-path // decrement in `deleteNounMetadata()`. const priorSubtype = isNew ? undefined : (typeof existingMetadata?.subtype === 'string' && (existingMetadata.subtype as string).length > 0 ? (existingMetadata.subtype as string) : undefined) const priorTypeForSubtype = isNew ? undefined : (existingMetadata?.noun as NounType | undefined) const newSubtype = typeof metadata.subtype === 'string' && metadata.subtype.length > 0 ? metadata.subtype as string : undefined const newType = metadata.noun as NounType | undefined if (priorSubtype && priorTypeForSubtype && (priorSubtype !== newSubtype || priorTypeForSubtype !== newType)) { this.decrementSubtypeCount(priorTypeForSubtype, priorSubtype) } if (newSubtype && newType && (isNew || priorSubtype !== newSubtype || priorTypeForSubtype !== newType)) { this.incrementSubtypeCount(newType, newSubtype) } // Visibility (8.0): only public entities count toward the user-facing totals. // The gate reads `metadata.visibility` (new) and `existingMetadata?.visibility` // (prior) directly off the record — no id-keyed visibility cache. const newVisibility = metadata.visibility const wasCounted = isNew ? false : isCountedVisibility(existingMetadata?.visibility) const isCounted = isCountedVisibility(newVisibility) // CRITICAL FIX: Increment count for new entities // This runs AFTER metadata is saved, guaranteeing type information is available // Uses synchronous increment since storage operations are already serialized // Fixes Bug #1: Count synchronization failure during add() and import() // 8.0: skip the user-facing total for internal/system entities (counts.json + getNounCount()). if (isNew && metadata.noun && isCounted) { this.incrementEntityCount(metadata.noun) // Per-type counter (stats().entitiesByType / counts.byTypeEnum) is maintained // HERE — gated on isNew + visibility, exactly parallel to the total above. It // used to be bumped unconditionally in saveNoun_internal(), but the HNSW index // re-saves a node on every neighbor-link change, so that inflated the per-type // counts with graph connectivity (e.g. 8 documents could read as 44). const typeIdx = TypeUtils.getNounIndex(metadata.noun as NounType) this.nounCountsByType[typeIdx]++ // Persist counts asynchronously (fire and forget) this.scheduleCountPersist().catch(() => { // Ignore persist errors - will retry on next operation }) // Persist type-statistics on the first entity of a type and every 100th // thereafter. This trigger used to live in saveNoun_internal(), which had to // call getNounType() purely to recover the type index; sourcing the type from // the metadata record here keeps the hot vector-save path free of any type // lookup. The "only when counted" half of the heuristic holds by construction // inside this branch. if (this.nounCountsByType[typeIdx] === 1 || this.nounCountsByType[typeIdx] % 100 === 0) { await this.saveTypeStatistics() } } else if (!isNew && metadata.noun && wasCounted !== isCounted) { // Visibility flipped on update(): move the entity in/out of the user-facing // total (counts.json / getNounCount()) AND the per-type counter together, so // stats().entitiesByType stays consistent with getNounCount(). const typeIdx = TypeUtils.getNounIndex(metadata.noun as NounType) if (isCounted) { this.incrementEntityCount(metadata.noun) this.nounCountsByType[typeIdx]++ // Same cadence-gated type-statistics persist as the fresh-add branch — only // fires when the entity is now counted (public), matching the original // `counted && (count === 1 || count % 100 === 0)` heuristic. if (this.nounCountsByType[typeIdx] === 1 || this.nounCountsByType[typeIdx] % 100 === 0) { await this.saveTypeStatistics() } } else { this.decrementEntityCount(metadata.noun) if (this.nounCountsByType[typeIdx] > 0) this.nounCountsByType[typeIdx]-- } this.scheduleCountPersist().catch(() => {}) } // 8.0 MVCC: entity-visible write — advance the generation watermark // (suppressed inside transact batches by the generation store). this.generationBumpHook?.() } /** * Get noun metadata from storage (METADATA-ONLY, NO VECTORS) * * **Performance**: Direct O(1) ID-first lookup - NO type search needed! * - **All lookups**: 1 read, ~500ms on cloud (consistent performance) * - **No cache needed**: Type is in the metadata, not the path * - **No type search**: ID-first paths eliminate 42-type search entirely * * **Clean architecture**: * - Path: `entities/nouns/{SHARD}/{ID}/metadata.json` * - Type is just a field in metadata (`noun: "document"`) * - MetadataIndex handles type queries (no path scanning needed) * - Scales to billions without any overhead * * **Performance**: Fast path for metadata-only reads * - **Speed**: 10ms vs 43ms (76-81% faster than getNoun) * - **Bandwidth**: 300 bytes vs 6KB (95% less) * - **Memory**: 300 bytes vs 6KB (87% less) * * **What's included**: * - All entity metadata (data, type, timestamps, confidence, weight) * - Custom user fields * - VFS metadata (_vfs.path, _vfs.size, etc.) * * **What's excluded**: * - 384-dimensional vector embeddings * - HNSW graph connections * * **Usage**: * - VFS operations (readFile, stat, readdir) - 100% of cases * - Existence checks: `if (await storage.getNounMetadata(id))` * - Metadata inspection: `metadata.data`, `metadata.noun` (type) * - Relationship traversal: Just need IDs, not vectors * * **When to use getNoun() instead**: * - Computing similarity on this specific entity * - Manual vector operations * - HNSW graph traversal * * @param id - Entity ID to retrieve metadata for * @returns Metadata or null if not found * * @performance * - O(1) direct ID lookup - always 1 read (~10ms on local disk) * - No caching complexity * - No type search fallbacks * * Type-first paths (removed) * Promoted to fast path for brain.get() optimization * CLEAN FIX: ID-first paths eliminate all type-search complexity */ public async getNounMetadata(id: string): Promise { await this.ensureInitialized() // Clean, simple, O(1) lookup - no type needed! const path = getNounMetadataPath(id) return this.readCanonicalObject(path) } /** * Batch fetch noun metadata from storage * * **Performance**: Reduces N sequential calls → 1-2 batch calls * - Local storage: N × 10ms → 1 × 10ms parallel (N× faster) * - Cloud storage: N × 300ms → 1 × 300ms batch (N× faster) * * **Use cases:** * - VFS tree traversal (fetch all children at once) * - brain.find() result hydration (batch load entities) * - brain.related() target entities (eliminate N+1) * - Import operations (batch existence checks) * * @param ids Array of entity IDs to fetch * @returns Map of id → metadata (only successful fetches included) * * @example * ```typescript * // Before (N+1 pattern) * for (const id of ids) { * const metadata = await storage.getNounMetadata(id) // N calls * } * * // After (batched) * const metadataMap = await storage.getNounMetadataBatch(ids) // 1 call * for (const id of ids) { * const metadata = metadataMap.get(id) * } * ``` * */ public async getNounMetadataBatch(ids: string[]): Promise> { await this.ensureInitialized() const results = new Map() if (ids.length === 0) return results // ID-first paths - no type grouping or search needed! // Build direct paths for all IDs const pathsToFetch: Array<{ path: string; id: string }> = ids.map(id => ({ path: getNounMetadataPath(id), id })) // Batch read all paths (uses adapter's native batch API or parallel fallback) const batchResults = await this.readCanonicalObjectBatch(pathsToFetch.map(p => p.path)) // Map results back to IDs for (const { path, id } of pathsToFetch) { const metadata = batchResults.get(path) if (metadata) { results.set(id, metadata) } } return results } /** * Batch get multiple nouns with vectors * * **Performance**: Eliminates N+1 pattern for vector loading * - Current: N × getNoun() = N × 50ms on GCS = 500ms for 10 entities * - Batched: 1 × getNounBatch() = 1 × 50ms on GCS = 50ms (**10x faster**) * * **Use cases:** * - batchGet() with includeVectors: true * - Loading entities for similarity computation * - Pre-loading vectors for batch processing * * @param ids Array of entity IDs to fetch (with vectors) * @returns Map of id → HNSWNounWithMetadata (only successful reads included) * */ public async getNounBatch(ids: string[]): Promise> { await this.ensureInitialized() const results = new Map() if (ids.length === 0) return results // Batch-fetch vectors and metadata in parallel // Build paths for vectors const vectorPaths: Array<{ path: string; id: string }> = ids.map(id => ({ path: getNounVectorPath(id), id })) // Build paths for metadata const metadataPaths: Array<{ path: string; id: string }> = ids.map(id => ({ path: getNounMetadataPath(id), id })) // Batch read vectors and metadata in parallel const [vectorResults, metadataResults] = await Promise.all([ this.readCanonicalObjectBatch(vectorPaths.map(p => p.path)), this.readCanonicalObjectBatch(metadataPaths.map(p => p.path)) ]) // Combine vectors + metadata into HNSWNounWithMetadata for (const { path: vectorPath, id } of vectorPaths) { const vectorData = vectorResults.get(vectorPath) const metadataPath = getNounMetadataPath(id) const metadataData = metadataResults.get(metadataPath) if (vectorData && metadataData) { // Deserialize, then combine via the canonical hydration helper const noun = this.deserializeNoun(vectorData) results.set(id, this.hydrateNounWithMetadata(noun, metadataData)) } } return results } /** * Batch read multiple canonical storage paths * * Core batching primitive that all batch operations build upon. * Handles the write cache and adapter-specific batching. * * **Performance**: * - Uses adapter's native batch API when available * - Falls back to parallel reads for non-batch adapters * - Respects rate limits via StorageBatchConfig * * @param paths Array of storage-root-relative paths to read * @returns Map of path → data (only successful reads included) * * @protected - Available to subclasses and batch operations */ protected async readCanonicalObjectBatch(paths: string[]): Promise> { if (paths.length === 0) return new Map() const results = new Map() // Step 1: Check write cache first (synchronous, instant) const pathsToFetch: string[] = [] for (const path of paths) { const cachedData = this.writeCache.get(path) if (cachedData !== undefined) { results.set(path, cachedData) } else { pathsToFetch.push(path) } } if (pathsToFetch.length === 0) { return results // All in write cache } // Step 2: Batch read from adapter const batchData = await this.readBatchFromAdapter(pathsToFetch) for (const [path, data] of batchData.entries()) { if (data !== null) { results.set(path, data) } } return results } /** * Adapter-level batch read with automatic batching strategy * * Uses adapter's native batch API when available: * - GCS: batch API (100 ops) * - S3/R2: batch operations (1000 ops) * - Azure: batch API (100 ops) * - Others: parallel reads via Promise.all() * * Automatically chunks large batches based on adapter's maxBatchSize. * * @param paths Array of resolved storage paths * @returns Map of path → data * * @private */ private async readBatchFromAdapter(paths: string[]): Promise> { if (paths.length === 0) return new Map() // Check if this class implements batch operations (will be added to cloud // adapters). Duck-typed optional capability — readBatch is not part of the // BaseStorageAdapter contract. const selfWithBatch = this as BaseStorage & { readBatch?: (paths: string[]) => Promise> } if (typeof selfWithBatch.readBatch === 'function') { // Adapter has native batch support - use it try { return await selfWithBatch.readBatch(paths) } catch (error) { // Fall back to parallel reads on batch failure prodLog.warn(`Batch read failed, falling back to parallel: ${error}`) } } // Fallback: Parallel individual reads // Respect adapter's maxConcurrent limit const batchConfig = this.getBatchConfig() const chunkSize = batchConfig.maxConcurrent || 50 const results = new Map() for (let i = 0; i < paths.length; i += chunkSize) { const chunk = paths.slice(i, i + chunkSize) const chunkResults = await Promise.allSettled( chunk.map(async path => ({ path, data: await this.readObjectFromPath(path) })) ) for (const result of chunkResults) { if (result.status === 'fulfilled' && result.value.data !== null) { results.set(result.value.path, result.value.data) } } } return results } /** * Get batch configuration for this storage adapter * * Override in subclasses to provide adapter-specific batch limits. * Defaults to conservative limits for safety. * * @public - Inherited from BaseStorageAdapter */ public override getBatchConfig(): StorageBatchConfig { // Conservative defaults - adapters should override with their actual limits return { maxBatchSize: 100, batchDelayMs: 0, maxConcurrent: 50, supportsParallelWrites: true, rateLimit: { operationsPerSecond: 1000, burstCapacity: 5000 } } } /** * Delete noun metadata from storage (ID-first, O(1) delete) */ public async deleteNounMetadata(id: string): Promise { await this.ensureInitialized() // Direct O(1) delete with ID-first path. Read the canonical record BEFORE // removing it: the per-type and subtype decrements are sourced from the // entity's own metadata (`noun` type, `subtype`, `visibility`) rather than an // id-keyed cache, keeping type-statistics honest across deletes — symmetric // with the increments in `saveNounMetadata_internal()`. const path = getNounMetadataPath(id) const record = await this.readCanonicalObject(path) await this.deleteCanonicalObject(path) const priorType = record?.noun as NounType | undefined // 8.0 visibility: an internal/system entity was never added to `nounCountsByType` // (gated in `saveNounMetadata_internal()`), so it must not be decremented here either. const priorCounted = isCountedVisibility(record?.visibility) if (priorType) { if (priorCounted) { const idx = TypeUtils.getNounIndex(priorType) if (this.nounCountsByType[idx] > 0) { this.nounCountsByType[idx]-- } } // Symmetric subtype decrement — same non-empty-string guard as the write path. const priorSubtype = typeof record?.subtype === 'string' && (record.subtype as string).length > 0 ? (record.subtype as string) : undefined if (priorSubtype) { this.decrementSubtypeCount(priorType, priorSubtype) } } // 8.0 MVCC: entity-visible write — advance the generation watermark // (suppressed inside transact batches by the generation store). this.generationBumpHook?.() } /** * Save verb metadata to storage (now typed) * Routes to correct sharded location based on UUID */ public async saveVerbMetadata(id: string, metadata: VerbMetadata): Promise { // Note: verb type is in HNSWVerb, not metadata return this.saveVerbMetadata_internal(id, metadata) } /** * Internal method for saving verb metadata (now typed) * Uses ID-first paths (must match getVerbMetadata) * * CRITICAL: Count synchronization happens here * This ensures verb counts are updated AFTER metadata exists, fixing the race condition * where storage adapters tried to read metadata before it was saved. * * Note: Verb type is now stored in both HNSWVerb (vector file) and VerbMetadata for count tracking * * @protected */ protected async saveVerbMetadata_internal(id: string, metadata: VerbMetadata): Promise { await this.ensureInitialized() // Extract verb type from metadata for ID-first path const verbType = metadata.verb as VerbType | undefined if (!verbType) { // Backward compatibility: fallback to old path if no verb type const keyInfo = this.analyzeKey(id, 'verb-metadata') await this.writeCanonicalObject(keyInfo.fullPath, metadata) // 8.0 MVCC: still an entity-visible write — advance the watermark. this.generationBumpHook?.() return } // Use ID-first path const path = getVerbMetadataPath(id) // Determine if this is a new verb by checking if metadata already exists const existingMetadata = await this.readCanonicalObject(path) const isNew = !existingMetadata // Save the metadata (write-cache coherent canonical write) await this.writeCanonicalObject(path, metadata) // Track verb subtype changes: on type or subtype change via updateRelation(), // decrement the prior bucket before incrementing the new one. The prior // (verb, subtype) is read straight from the canonical record (`existingMetadata`, // loaded above) — there is no id-keyed verb-subtype cache. Symmetric with the // delete-path decrement in `deleteVerbMetadata()`. const priorVerbForSubtype = isNew ? undefined : (existingMetadata?.verb as VerbType | undefined) const priorSubtype = isNew ? undefined : (typeof existingMetadata?.subtype === 'string' && (existingMetadata.subtype as string).length > 0 ? (existingMetadata.subtype as string) : undefined) const newSubtype = typeof metadata.subtype === 'string' && metadata.subtype.length > 0 ? metadata.subtype as string : undefined if (priorSubtype && priorVerbForSubtype && (priorSubtype !== newSubtype || priorVerbForSubtype !== verbType)) { this.decrementVerbSubtypeCount(priorVerbForSubtype, priorSubtype) } if (newSubtype && (isNew || priorSubtype !== newSubtype || priorVerbForSubtype !== verbType)) { this.incrementVerbSubtypeCount(verbType, newSubtype) } // Visibility (8.0): verb mirror of the noun count gating. The gate reads // `metadata.visibility` (new) and `existingMetadata?.visibility` (prior) // directly off the record — no id-keyed visibility cache. // // NOTE on `verbCountsByType`: unlike the noun path, `saveVerb_internal()` runs // BEFORE this method (relate() saves the verb vector first) and has already done // an UNCONDITIONAL `verbCountsByType[idx]++`. We therefore COMPENSATE here: for a // new hidden edge, undo that bump. `updateRelation()` does not re-run // `saveVerb_internal()`, so on a visibility flip we adjust the bucket directly. const newVisibility = metadata.visibility const wasCounted = isNew ? false : isCountedVisibility(existingMetadata?.visibility) const isCounted = isCountedVisibility(newVisibility) const verbTypeIdx = TypeUtils.getVerbIndex(verbType) // CRITICAL FIX: Increment verb count for new relationships // This runs AFTER metadata is saved // Uses synchronous increment since storage operations are already serialized // Fixes Bug #2: Count synchronization failure during relate() and import() // 8.0: skip the user-facing total for internal/system edges (counts.json + getVerbCount()). if (isNew) { if (isCounted) { this.incrementVerbCount(verbType) } else { // Hidden edge: undo the unconditional bump from saveVerb_internal(). if (this.verbCountsByType[verbTypeIdx] > 0) this.verbCountsByType[verbTypeIdx]-- } // Persist counts asynchronously (fire and forget) this.scheduleCountPersist().catch(() => { // Ignore persist errors - will retry on next operation }) } else if (wasCounted !== isCounted) { // Visibility flipped on updateRelation() (saveVerb_internal did not run): move the // edge in/out of both the user-facing total and the per-type bucket. if (isCounted) { this.incrementVerbCount(verbType) this.verbCountsByType[verbTypeIdx]++ } else { this.decrementVerbCount(verbType) if (this.verbCountsByType[verbTypeIdx] > 0) this.verbCountsByType[verbTypeIdx]-- } this.scheduleCountPersist().catch(() => {}) } // 8.0 MVCC: entity-visible write — advance the generation watermark // (suppressed inside transact batches by the generation store). this.generationBumpHook?.() } /** * Get verb metadata from storage (now typed) * Uses ID-first paths (must match saveVerbMetadata_internal) */ public async getVerbMetadata(id: string): Promise { await this.ensureInitialized() // Direct O(1) lookup with ID-first paths - no type search needed! // Symmetric with getNounMetadata: readCanonicalObject already returns null for // a genuine not-found, so a real storage fault (permission/corruption/IO) must // propagate rather than be masked as "this verb has no metadata". const path = getVerbMetadataPath(id) return this.readCanonicalObject(path) } /** * Delete verb metadata from storage (ID-first, O(1) delete) */ public async deleteVerbMetadata(id: string): Promise { await this.ensureInitialized() // Direct O(1) delete with ID-first path. Read the canonical record BEFORE // removing it so the verb-subtype decrement is sourced from the edge's own // metadata (`verb` type + `subtype`) rather than an id-keyed cache — symmetric // with the increment in `saveVerbMetadata_internal()`. Verb deletes do not // touch `verbCountsByType` in this path (matching prior behavior), so no // visibility read is needed here. const path = getVerbMetadataPath(id) const record = await this.readCanonicalObject(path) await this.deleteCanonicalObject(path) const priorVerb = record?.verb as VerbType | undefined const priorSubtype = typeof record?.subtype === 'string' && (record.subtype as string).length > 0 ? (record.subtype as string) : undefined if (priorVerb && priorSubtype) { this.decrementVerbSubtypeCount(priorVerb, priorSubtype) } // 8.0 MVCC: entity-visible write — advance the generation watermark // (suppressed inside transact batches by the generation store). this.generationBumpHook?.() } // ============================================================================ // ID-FIRST HELPER METHODS // Direct O(1) ID lookups - no type needed! // Clean, simple architecture for billion-scale performance // ============================================================================ /** * Load type statistics from storage * Rebuilds type counts if needed (called during init) * * Auto-detects the 7.20.0–7.21.0 poisoned-statistics signature: every * noun attributed to `'thing'` because the old `getNounType()` was hardcoded. * If detected, runs `rebuildTypeCounts()` once to rewrite the file with * correct per-type counts derived from on-disk metadata. Logged loudly. */ protected async loadTypeStatistics(): Promise { try { const stats = await this.readObjectFromPath(`${SYSTEM_DIR}/type-statistics.json`) if (stats) { // Restore counts from saved statistics if (stats.nounCounts && stats.nounCounts.length === NOUN_TYPE_COUNT) { this.nounCountsByType = new Uint32Array(stats.nounCounts) } if (stats.verbCounts && stats.verbCounts.length === VERB_TYPE_COUNT) { this.verbCountsByType = new Uint32Array(stats.verbCounts) } if (await this.detectPoisonedTypeStatistics()) { prodLog.warn( '[BaseStorage] Detected poisoned type-statistics.json signature ' + '(all nouns attributed to \'thing\' — symptom of pre-7.22 getNounType ' + 'hardcode). Rebuilding counts from on-disk metadata.' ) await this.rebuildTypeCounts() } } } catch (error) { // No existing type statistics, starting fresh } } /** * Detect the 7.20.0–7.21.0 poisoned-statistics signature. * * The defect: the old `getNounType()` returned hardcoded `'thing'`, so * `_system/type-statistics.json` ended up with `nounCounts[thing] === N` * and every other index `=== 0`, regardless of the actual mix on disk. * * Heuristic: nounCounts has at least 2 non-thing entities on disk (so the * file *should* show multiple non-zero buckets) but only the `thing` * bucket is non-zero. We bound the on-disk check with `limit: 3` so this * never costs more than a couple of cheap reads. * * Returns false (no self-heal needed) for genuinely thing-only stores or * empty stores. */ private async detectPoisonedTypeStatistics(): Promise { const thingIdx = TypeUtils.getNounIndex('thing' as NounType) if (thingIdx < 0) return false let nonZeroBuckets = 0 let thingCount = 0 for (let i = 0; i < this.nounCountsByType.length; i++) { if (this.nounCountsByType[i] > 0) { nonZeroBuckets++ if (i === thingIdx) thingCount = this.nounCountsByType[i] } } // Only one bucket populated, and it's 'thing' with ≥2 entities — suspicious. if (nonZeroBuckets !== 1 || thingCount < 2) return false // Cross-check: sample a few metadata files. If any non-'thing' types // show up, we're confirmed poisoned. If only 'thing' types appear in the // sample, this is a genuine thing-only store and we leave the file alone. try { const sample = await this.getNouns({ pagination: { offset: 0, limit: 3 } }) for (const noun of sample.items) { const type = await this.getNounTypeFromStorageAsync(noun.id) if (type && type !== ('thing' as NounType)) { return true } } } catch { // If we can't read the metadata, don't trigger a rebuild — leave state alone. } return false } /** * Save type statistics to storage * Periodically called when counts are updated */ protected async saveTypeStatistics(): Promise { const stats = { nounCounts: Array.from(this.nounCountsByType), verbCounts: Array.from(this.verbCountsByType), updatedAt: Date.now() } await this.writeObjectToPath(`${SYSTEM_DIR}/type-statistics.json`, stats) } /** * Increment the (type, subtype) count, creating the inner map on first use. * Indexed by NounType index so it lines up with `nounCountsByType` and can * be reduced into per-type totals without re-keying. */ protected incrementSubtypeCount(type: NounType, subtype: string): void { const typeIdx = TypeUtils.getNounIndex(type) if (typeIdx < 0) return let inner = this.subtypeCountsByType.get(typeIdx) if (!inner) { inner = new Map() this.subtypeCountsByType.set(typeIdx, inner) } inner.set(subtype, (inner.get(subtype) || 0) + 1) } /** * Decrement the (type, subtype) count. Deletes the inner key when it reaches * 0 and the outer entry when its inner map is empty, so the persisted shape * stays compact across heavy churn. */ protected decrementSubtypeCount(type: NounType, subtype: string): void { const typeIdx = TypeUtils.getNounIndex(type) if (typeIdx < 0) return const inner = this.subtypeCountsByType.get(typeIdx) if (!inner) return const next = (inner.get(subtype) || 0) - 1 if (next <= 0) { inner.delete(subtype) if (inner.size === 0) this.subtypeCountsByType.delete(typeIdx) } else { inner.set(subtype, next) } } /** * Load `_system/subtype-statistics.json` into `subtypeCountsByType`. * Persisted shape: `{ counts: { [typeIdx]: { [subtype]: count } }, updatedAt }`. * Missing file or parse error → start from empty (matches loadTypeStatistics). */ protected async loadSubtypeStatistics(): Promise { try { const stats = await this.readObjectFromPath(`${SYSTEM_DIR}/subtype-statistics.json`) if (stats && stats.counts && typeof stats.counts === 'object') { this.subtypeCountsByType.clear() for (const [typeKey, subtypeMap] of Object.entries(stats.counts as Record>)) { const typeIdx = Number(typeKey) if (!Number.isInteger(typeIdx) || typeIdx < 0 || typeIdx >= NOUN_TYPE_COUNT) continue const inner = new Map() for (const [subtype, count] of Object.entries(subtypeMap)) { if (typeof count === 'number' && count > 0) inner.set(subtype, count) } if (inner.size > 0) this.subtypeCountsByType.set(typeIdx, inner) } } } catch { // No existing subtype statistics, starting fresh. } } /** * Save subtype statistics to storage. Mirrors the type-statistics persistence * cadence (called from `flushCounts()` and the periodic save in * `saveNoun_internal`). */ protected async saveSubtypeStatistics(): Promise { const counts: Record> = {} for (const [typeIdx, inner] of this.subtypeCountsByType.entries()) { const innerObj: Record = {} for (const [subtype, count] of inner.entries()) innerObj[subtype] = count counts[String(typeIdx)] = innerObj } await this.writeObjectToPath(`${SYSTEM_DIR}/subtype-statistics.json`, { counts, updatedAt: Date.now() }) } /** * Rebuild subtype counts from on-disk metadata. Companion to `rebuildTypeCounts()` * — used for poison recovery and explicit repair via `brainy inspect repair`. * O(N) over all nouns. */ public async rebuildSubtypeCounts(): Promise { prodLog.info('[BaseStorage] Rebuilding subtype counts from storage...') this.subtypeCountsByType.clear() for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/nouns/${shardHex}` try { const paths = await this.listCanonicalObjects(shardDir) for (const path of paths) { if (!path.includes('/metadata.json')) continue try { const metadata = await this.readCanonicalObject(path) if (metadata && metadata.noun && typeof metadata.subtype === 'string' && metadata.subtype.length > 0) { this.incrementSubtypeCount(metadata.noun as NounType, metadata.subtype) } } catch { /* skip unreadable entities */ } } } catch { /* skip missing shards */ } } await this.saveSubtypeStatistics() const totals = Array.from(this.subtypeCountsByType.values()) .reduce((sum, inner) => sum + Array.from(inner.values()).reduce((s, n) => s + n, 0), 0) prodLog.info(`[BaseStorage] Rebuilt subtype counts: ${totals} entities across ${this.subtypeCountsByType.size} NounTypes`) } /** * Increment the (verb, subtype) count, creating the inner map on first use. * Verb-side mirror of `incrementSubtypeCount`. Indexed by VerbType index so it * lines up with `verbCountsByType` and can be reduced into per-verb totals * without re-keying. */ protected incrementVerbSubtypeCount(verb: VerbType, subtype: string): void { const verbIdx = TypeUtils.getVerbIndex(verb) if (verbIdx < 0) return let inner = this.verbSubtypeCountsByType.get(verbIdx) if (!inner) { inner = new Map() this.verbSubtypeCountsByType.set(verbIdx, inner) } inner.set(subtype, (inner.get(subtype) || 0) + 1) } /** * Decrement the (verb, subtype) count. Mirror of `decrementSubtypeCount`. * Deletes the inner key when it reaches 0 and the outer entry when its inner * map is empty, so the persisted shape stays compact across heavy churn. */ protected decrementVerbSubtypeCount(verb: VerbType, subtype: string): void { const verbIdx = TypeUtils.getVerbIndex(verb) if (verbIdx < 0) return const inner = this.verbSubtypeCountsByType.get(verbIdx) if (!inner) return const next = (inner.get(subtype) || 0) - 1 if (next <= 0) { inner.delete(subtype) if (inner.size === 0) this.verbSubtypeCountsByType.delete(verbIdx) } else { inner.set(subtype, next) } } /** * Load `_system/verb-subtype-statistics.json` into `verbSubtypeCountsByType`. * Persisted shape mirrors the noun-side rollup: * `{ counts: { [verbIdx]: { [subtype]: count } }, updatedAt }`. Missing file * or parse error → start from empty. */ protected async loadVerbSubtypeStatistics(): Promise { try { const stats = await this.readObjectFromPath(`${SYSTEM_DIR}/verb-subtype-statistics.json`) if (stats && stats.counts && typeof stats.counts === 'object') { this.verbSubtypeCountsByType.clear() for (const [verbKey, subtypeMap] of Object.entries(stats.counts as Record>)) { const verbIdx = Number(verbKey) if (!Number.isInteger(verbIdx) || verbIdx < 0 || verbIdx >= VERB_TYPE_COUNT) continue const inner = new Map() for (const [subtype, count] of Object.entries(subtypeMap)) { if (typeof count === 'number' && count > 0) inner.set(subtype, count) } if (inner.size > 0) this.verbSubtypeCountsByType.set(verbIdx, inner) } } } catch { // No existing verb subtype statistics, starting fresh. } } /** * Save verb subtype statistics to storage. Mirrors `saveSubtypeStatistics`. * Same persistence cadence (flushCounts + periodic saves alongside other * type-statistics). */ protected async saveVerbSubtypeStatistics(): Promise { const counts: Record> = {} for (const [verbIdx, inner] of this.verbSubtypeCountsByType.entries()) { const innerObj: Record = {} for (const [subtype, count] of inner.entries()) innerObj[subtype] = count counts[String(verbIdx)] = innerObj } await this.writeObjectToPath(`${SYSTEM_DIR}/verb-subtype-statistics.json`, { counts, updatedAt: Date.now() }) } /** * Rebuild verb subtype counts from on-disk metadata. Companion to * `rebuildSubtypeCounts()`. O(N) over all verbs. Used for poison recovery * and explicit repair via `brainy inspect repair`. */ public async rebuildVerbSubtypeCounts(): Promise { prodLog.info('[BaseStorage] Rebuilding verb subtype counts from storage...') this.verbSubtypeCountsByType.clear() for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { const paths = await this.listCanonicalObjects(shardDir) for (const path of paths) { if (!path.includes('/metadata.json')) continue try { const metadata = await this.readCanonicalObject(path) if (metadata && metadata.verb && typeof metadata.subtype === 'string' && metadata.subtype.length > 0) { const verb = metadata.verb as VerbType const subtype = metadata.subtype as string this.incrementVerbSubtypeCount(verb, subtype) } } catch { /* skip unreadable verbs */ } } } catch { /* skip missing shards */ } } await this.saveVerbSubtypeStatistics() const totals = Array.from(this.verbSubtypeCountsByType.values()) .reduce((sum, inner) => sum + Array.from(inner.values()).reduce((s, n) => s + n, 0), 0) prodLog.info(`[BaseStorage] Rebuilt verb subtype counts: ${totals} relationships across ${this.verbSubtypeCountsByType.size} VerbTypes`) } /** * Persist both counter systems atomically when an explicit flush is * requested. `super.flushCounts()` writes `entityCounts` (the Map in * `BaseStorageAdapter`); we additionally write `nounCountsByType` / * `verbCountsByType` so a reader opening the same directory sees the * exact same counts the writer holds in memory. * * Before this override the Uint32Array counters were only persisted * on a heuristic schedule inside `saveNoun_internal` (first-of-type or * every-100th), which left readers seeing stale counts after a clean * writer flush — the same silent-stale failure mode that once produced * zero counts from `brain.stats()`. */ public override async flushCounts(): Promise { await super.flushCounts() await this.saveTypeStatistics() await this.saveSubtypeStatistics() await this.saveVerbSubtypeStatistics() } /** * Get noun counts by type (O(1) access to type statistics) * Exposed for MetadataIndexManager to use as single source of truth * @returns Uint32Array indexed by NounType enum value (42 types) */ public getNounCountsByType(): Uint32Array { return this.nounCountsByType } /** * Get subtype counts by NounType (O(1) access to subtype statistics). * Returned map is the live in-memory view — callers must treat it as * read-only. Outer key: NounType index. Inner map: subtype → count. * * @returns Map keyed by NounType index → Map of subtype → count */ public getSubtypeCountsByType(): Map> { return this.subtypeCountsByType } /** * Get verb subtype counts (O(1) access). Verb-side mirror of * `getSubtypeCountsByType`. Outer key: VerbType index. Inner map: subtype → count. * Returned map is the live in-memory view — callers must treat it as read-only. */ public getVerbSubtypeCountsByType(): Map> { return this.verbSubtypeCountsByType } /** * Get verb counts by type (O(1) access to type statistics) * Exposed for MetadataIndexManager to use as single source of truth * @returns Uint32Array indexed by VerbType enum value (127 types) */ public getVerbCountsByType(): Uint32Array { return this.verbCountsByType } /** * Rebuild type counts from actual storage. Scans every shard, reads each * entity's metadata, and reconstructs `nounCountsByType` / `verbCountsByType` * from ground truth. Persists the corrected counts to `type-statistics.json`. * * Public so it can be triggered by `brainy inspect repair` and by the * `brain.health()` reconciliation path. Called automatically by * `loadTypeStatistics()` when the persisted state matches the poisoned * 7.20.0–7.21.0 signature (all entities attributed to `'thing'`). * * For very large stores this is O(N) where N is total entities; intended * for diagnostic / repair use, not on every init. */ public async rebuildTypeCounts(): Promise { prodLog.info('[BaseStorage] Rebuilding type counts from storage...') // Rebuild by scanning shards (0x00-0xFF) and reading metadata this.nounCountsByType = new Uint32Array(NOUN_TYPE_COUNT) this.verbCountsByType = new Uint32Array(VERB_TYPE_COUNT) // Scan noun shards for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/nouns/${shardHex}` try { const paths = await this.listCanonicalObjects(shardDir) for (const path of paths) { if (!path.includes('/metadata.json')) continue try { const metadata = await this.readCanonicalObject(path) if (metadata && metadata.noun) { // 8.0 visibility: rebuild only public entities into the user-facing per-type stat. if (isCountedVisibility(metadata.visibility)) { const typeIndex = TypeUtils.getNounIndex(metadata.noun) if (typeIndex >= 0 && typeIndex < NOUN_TYPE_COUNT) { this.nounCountsByType[typeIndex]++ } } } } catch (error) { // Skip entities that fail to load } } } catch (error) { // Skip shards that don't exist } } // Scan verb shards for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { const paths = await this.listCanonicalObjects(shardDir) for (const path of paths) { if (!path.includes('/metadata.json')) continue try { const metadata = await this.readCanonicalObject(path) if (metadata && metadata.verb) { // 8.0 visibility: rebuild only public edges into the user-facing per-type stat. if (isCountedVisibility(metadata.visibility)) { const typeIndex = TypeUtils.getVerbIndex(metadata.verb) if (typeIndex >= 0 && typeIndex < VERB_TYPE_COUNT) { this.verbCountsByType[typeIndex]++ } } } } catch (error) { // Skip entities that fail to load } } } catch (error) { // Skip shards that don't exist } } // Save rebuilt counts to storage await this.saveTypeStatistics() const totalVerbs = this.verbCountsByType.reduce((sum, count) => sum + count, 0) const totalNouns = this.nounCountsByType.reduce((sum, count) => sum + count, 0) prodLog.info(`[BaseStorage] Rebuilt counts: ${totalNouns} nouns, ${totalVerbs} verbs`) } /** * Resolve a noun's `NounType` straight from its canonical metadata record on * disk. Used by the poisoned-statistics detector (`detectPoisonedTypeStatistics()`) * to confirm whether on-disk types disagree with a `'thing'`-only persisted * rollup before triggering a full `rebuildTypeCounts()`. Returns `null` when * metadata genuinely doesn't exist or the read fails; the caller decides whether * to skip the entity. There is no in-memory id→type cache — the record is the * single source of truth. * * @param id - The noun id whose type to resolve. * @returns The stored `NounType`, or `null` if absent/unreadable. */ protected async getNounTypeFromStorageAsync(id: string): Promise { try { const metadataPath = getNounMetadataPath(id) const metadata = await this.readCanonicalObject(metadataPath) if (metadata && (metadata as NounMetadata).noun) { return (metadata as NounMetadata).noun as NounType } } catch { // Storage error — treat as unknown. } return null } /** * Get verb type from verb object * Verb type is a required field in HNSWVerb */ protected getVerbType(verb: HNSWVerb | GraphVerb): VerbType { // verb is a required field in HNSWVerb if ('verb' in verb && verb.verb) { return verb.verb as VerbType } // Fallback for GraphVerb (type alias) if ('type' in verb && verb.type) { return verb.type as VerbType } // This should never happen with current data prodLog.warn(`[BaseStorage] Verb missing type field for ${verb.id}, defaulting to 'relatedTo'`) return 'relatedTo' } // ============================================================================ // DESERIALIZATION HELPERS // Centralized Map/Set reconstruction from JSON storage format // ============================================================================ /** * Deserialize HNSW connections from JSON storage format * * Converts plain object { "0": ["id1"], "1": ["id2"] } * into Map> * * Central helper to fix serialization bug across all code paths * Root cause: JSON.stringify(Map) = {} (empty object), must reconstruct on read */ protected deserializeConnections(connections: any): Map> { const result = new Map>() if (!connections || typeof connections !== 'object') { return result } // Already a Map (in-memory, not from JSON) if (connections instanceof Map) { return connections } // Deserialize from plain object for (const [levelStr, ids] of Object.entries(connections)) { if (Array.isArray(ids)) { result.set(parseInt(levelStr, 10), new Set(ids)) } else if (ids && typeof ids === 'object') { // Handle Set-like or array-like objects result.set(parseInt(levelStr, 10), new Set(Object.values(ids))) } } return result } /** * Deserialize HNSWNoun from JSON storage format * * Ensures connections are properly reconstructed from Map → object → Map * Fixes: "TypeError: noun.connections.entries is not a function" */ protected deserializeNoun(data: any): HNSWNoun { return { ...data, connections: this.deserializeConnections(data.connections) } } /** * Deserialize HNSWVerb from JSON storage format * * Ensures connections are properly reconstructed from Map → object → Map * Fixes same serialization bug for verbs */ protected deserializeVerb(data: any): HNSWVerb { return { ...data, connections: this.deserializeConnections(data.connections) } } // ============================================================================ // ABSTRACT METHOD IMPLEMENTATIONS // Converted from abstract to concrete - all adapters now have built-in type-aware // ============================================================================ /** * Save a noun to storage (ID-first path) */ protected async saveNoun_internal(noun: HNSWNoun): Promise { const path = getNounVectorPath(noun.id) // Hot path: write the vector record only. Per-type counters // (`nounCountsByType`, read by stats().entitiesByType) AND the periodic // type-statistics persist trigger are both maintained in // saveNounMetadata_internal(), which holds the canonical metadata record — // and therefore the NounType + visibility — at the exact point a count // changes. saveNoun_internal() also re-runs on every HNSW neighbor-link // re-save, so it deliberately performs NO type lookup and NO count work here. await this.writeCanonicalObject(path, noun) } /** * Get a noun from storage (ID-first path) */ protected async getNoun_internal(id: string): Promise { // Direct O(1) lookup with ID-first paths - no type search needed! const path = getNounVectorPath(id) try { // Write-cache coherent canonical read const noun = await this.readCanonicalObject(path) if (noun) { // Deserialize connections Map from JSON storage format return this.deserializeNoun(noun) } } catch (error) { // Entity not found return null } return null } /** * Get nouns by noun type (Shard-based iteration!) */ protected async getNounsByNounType_internal( nounType: string ): Promise { // Iterate by shards (0x00-0xFF) instead of types // Type is stored in metadata.noun field, we filter as we load const nouns: HNSWNoun[] = [] for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/nouns/${shardHex}` try { const nounFiles = await this.listCanonicalObjects(shardDir) for (const nounPath of nounFiles) { if (!nounPath.includes('/vectors.json')) continue try { const noun = await this.readCanonicalObject(nounPath) if (noun) { const deserialized = this.deserializeNoun(noun) // Check type from metadata const metadata = await this.getNounMetadata(deserialized.id) if (metadata && metadata.noun === nounType) { nouns.push(deserialized) } } } catch (error) { // Skip nouns that fail to load } } } catch (error) { // Skip shards that have no data } } return nouns } /** * Delete a noun from storage (ID-first, O(1) delete) */ protected async deleteNoun_internal(id: string): Promise { // Direct O(1) delete with ID-first path const path = getNounVectorPath(id) await this.deleteCanonicalObject(path) // Note: Type-specific counts will be decremented via metadata tracking // The real type is in metadata, accessible if needed via getNounMetadata(id) } /** * Save a verb to storage (ID-first path) */ protected async saveVerb_internal(verb: HNSWVerb): Promise { // Type is now a first-class field in HNSWVerb - no caching needed! const type = verb.verb as VerbType const path = getVerbVectorPath(verb.id) prodLog.debug(`[BaseStorage] saveVerb_internal: id=${verb.id}, sourceId=${verb.sourceId}, targetId=${verb.targetId}, type=${type}`) // Update type tracking const typeIndex = TypeUtils.getVerbIndex(type) this.verbCountsByType[typeIndex]++ // Write-cache coherent canonical write await this.writeCanonicalObject(path, verb) // GraphAdjacencyIndex updates are now handled EXCLUSIVELY by Brainy.relate() // via AddToGraphIndexOperation in the transaction system. This provides: // 1. Singleton pattern - only one graphIndex instance exists (via getGraphIndex()) // 2. Transaction rollback - if relate() fails, index update is rolled back // 3. No double-counting - prevents duplicate addVerb() calls // REMOVED: Direct graphIndex.addVerb() call that caused dual-ownership bugs // Periodically save statistics // Also save on first verb of each type to ensure low-count types are tracked // This prevents stale statistics after restart for types with < 100 verbs (common for VFS) const shouldSave = this.verbCountsByType[typeIndex] === 1 || // First verb of type this.verbCountsByType[typeIndex] % 100 === 0 // Every 100th if (shouldSave) { await this.saveTypeStatistics() } } /** * Get a verb from storage (ID-first path) */ protected async getVerb_internal(id: string): Promise { // Direct O(1) lookup with ID-first paths - no type search needed! const path = getVerbVectorPath(id) try { // Write-cache coherent canonical read const verb = await this.readCanonicalObject(path) if (verb) { // Deserialize connections Map from JSON storage format return this.deserializeVerb(verb) } } catch (error) { // Entity not found return null } return null } /** * Get verbs by source (Uses GraphAdjacencyIndex when available) * Falls back to shard iteration during initialization to avoid circular dependency */ protected async getVerbsBySource_internal( sourceId: string ): Promise { await this.ensureInitialized() prodLog.debug(`[BaseStorage] getVerbsBySource_internal: sourceId=${sourceId}, graphIndex=${!!this.graphIndex}, isInitialized=${this.graphIndex?.isInitialized}`) // Fast path - use GraphAdjacencyIndex if available (lazy-loaded). // 8.0 BigInt boundary: convert the UUID to an entity int up front and // resolve returned verb ints back to verb-id strings. if (this.graphIndex && this.graphIndex.isInitialized && this.graphEntityIdResolver) { try { const sourceInt = this.graphEntityIdResolver.getInt(sourceId) if (sourceInt === undefined) { // Never-mapped UUID — the entity has no relations by definition. return [] } const verbInts = await this.graphIndex.getVerbIdsBySource(BigInt(sourceInt)) const verbIds = (await this.graphIndex.verbIntsToIds(verbInts)) .filter((id): id is string => id !== null) prodLog.debug(`[BaseStorage] GraphAdjacencyIndex found ${verbIds.length} verb IDs for sourceId=${sourceId}`) // PERFORMANCE FIX - Batch fetch verbs + metadata (eliminates N+1 pattern) // Before: N sequential calls (10 children = 20 × 300ms = 6000ms on GCS) // After: 2 parallel batch calls (10 children = 2 × 300ms = 600ms on GCS) // 10x improvement for cloud storage (GCS, S3, Azure) const verbPaths = verbIds.map(id => getVerbVectorPath(id)) const metadataPaths = verbIds.map(id => getVerbMetadataPath(id)) const [verbsMap, metadataMap] = await Promise.all([ this.readCanonicalObjectBatch(verbPaths), this.readCanonicalObjectBatch(metadataPaths) ]) const results: HNSWVerbWithMetadata[] = [] for (const verbId of verbIds) { const verbPath = getVerbVectorPath(verbId) const metadataPath = getVerbMetadataPath(verbId) const rawVerb = verbsMap.get(verbPath) const metadata = metadataMap.get(metadataPath) if (rawVerb && metadata) { // CRITICAL - Deserialize connections Map from JSON storage format, // then combine via the canonical hydration helper (reserved fields // top-level, ONLY custom fields in `metadata`). const verb = this.deserializeVerb(rawVerb) results.push(this.hydrateVerbWithMetadata(verb, metadata)) } } prodLog.debug(`[BaseStorage] GraphAdjacencyIndex + batch fetch returned ${results.length} verbs`) return results } catch (error) { prodLog.warn('[BaseStorage] GraphAdjacencyIndex lookup failed, falling back to shard iteration:', error) } } // Fallback - iterate by shards (WITH deserialization fix!) prodLog.debug(`[BaseStorage] Using shard iteration fallback for sourceId=${sourceId}`) const results: HNSWVerbWithMetadata[] = [] let shardsScanned = 0 let verbsFound = 0 for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { const verbFiles = await this.listCanonicalObjects(shardDir) shardsScanned++ for (const verbPath of verbFiles) { if (!verbPath.includes('/vectors.json')) continue try { const rawVerb = await this.readCanonicalObject(verbPath) if (!rawVerb) continue verbsFound++ // CRITICAL - Deserialize connections Map from JSON storage format const verb = this.deserializeVerb(rawVerb) if (verb.sourceId === sourceId) { const metadataPath = getVerbMetadataPath(verb.id) const metadata = await this.readCanonicalObject(metadataPath) // Canonical hydration — reserved fields top-level, ONLY custom // fields in `metadata`. results.push(this.hydrateVerbWithMetadata(verb, metadata)) } } catch (error) { // Skip verbs that fail to load prodLog.debug(`[BaseStorage] Failed to load verb from ${verbPath}:`, error) } } } catch (error) { // Skip shards that have no data } } prodLog.debug(`[BaseStorage] Shard iteration: scanned ${shardsScanned} shards, found ${verbsFound} total verbs, matched ${results.length} for sourceId=${sourceId}`) return results } /** * Batch get verbs by source IDs * * **Performance**: Eliminates N+1 query pattern for relationship lookups * - Current: N × getVerbsBySource() = N × (list all verbs + filter) * - Batched: 1 × list all verbs + filter by N sourceIds * * **Use cases:** * - VFS tree traversal (get Contains edges for multiple directories) * - brain.related() for multiple entities * - Graph traversal (fetch neighbors of multiple nodes) * * @param sourceIds Array of source entity IDs * @param verbType Optional verb type filter (e.g., VerbType.Contains for VFS) * @returns Map of sourceId → verbs[] * * @example * ```typescript * // Before (N+1 pattern) * for (const dirId of dirIds) { * const children = await storage.getVerbsBySource(dirId) // N calls * } * * // After (batched) * const childrenByDir = await storage.getVerbsBySourceBatch(dirIds, VerbType.Contains) // 1 scan * for (const dirId of dirIds) { * const children = childrenByDir.get(dirId) || [] * } * ``` * */ public async getVerbsBySourceBatch( sourceIds: string[], verbType?: VerbType ): Promise> { await this.ensureInitialized() const results = new Map() if (sourceIds.length === 0) return results // Initialize empty arrays for all requested sourceIds for (const sourceId of sourceIds) { results.set(sourceId, []) } // Convert sourceIds to Set for O(1) lookup const sourceIdSet = new Set(sourceIds) // Iterate by shards (0x00-0xFF) instead of types for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { // List all verb files in this shard const verbFiles = await this.listCanonicalObjects(shardDir) // Build paths for batch read const verbPaths: string[] = [] const metadataPaths: string[] = [] const pathToId = new Map() for (const verbPath of verbFiles) { if (!verbPath.includes('/vectors.json')) continue verbPaths.push(verbPath) // Extract ID from path: "entities/verbs/{shard}/{id}/vector.json" const parts = verbPath.split('/') const verbId = parts[parts.length - 2] // ID is second-to-last segment pathToId.set(verbPath, verbId) // Prepare metadata path metadataPaths.push(getVerbMetadataPath(verbId)) } // Batch read all verb files for this shard const verbDataMap = await this.readCanonicalObjectBatch(verbPaths) const metadataMap = await this.readCanonicalObjectBatch(metadataPaths) // Process results for (const [verbPath, rawVerbData] of verbDataMap.entries()) { if (!rawVerbData || !rawVerbData.sourceId) continue // Deserialize connections Map from JSON storage format const verbData = this.deserializeVerb(rawVerbData) // Check if this verb's source is in our requested set if (!sourceIdSet.has(verbData.sourceId)) continue // If verbType specified, filter by type if (verbType && verbData.verb !== verbType) continue // Found matching verb - hydrate with metadata const verbId = pathToId.get(verbPath)! const metadataPath = getVerbMetadataPath(verbId) const metadata = metadataMap.get(metadataPath) || {} // Canonical hydration — reserved fields top-level (including // subtype/data, which this site previously dropped), ONLY custom // fields in `metadata`. const hydratedVerb = this.hydrateVerbWithMetadata(verbData, metadata) // Add to results for this sourceId const sourceVerbs = results.get(verbData.sourceId)! sourceVerbs.push(hydratedVerb) } } catch (error) { // Skip shards that have no data } } return results } /** * Get verbs by target * Reverted to fix circular dependency deadlock * Fixed to directly list verb files instead of directories */ protected async getVerbsByTarget_internal( targetId: string ): Promise { await this.ensureInitialized() // Fast path - use GraphAdjacencyIndex if available (lazy-loaded). // 8.0 BigInt boundary: convert the UUID to an entity int up front and // resolve returned verb ints back to verb-id strings. if (this.graphIndex && this.graphIndex.isInitialized && this.graphEntityIdResolver) { try { const targetInt = this.graphEntityIdResolver.getInt(targetId) if (targetInt === undefined) { // Never-mapped UUID — the entity has no relations by definition. return [] } const verbInts = await this.graphIndex.getVerbIdsByTarget(BigInt(targetInt)) const verbIds = (await this.graphIndex.verbIntsToIds(verbInts)) .filter((id): id is string => id !== null) const results: HNSWVerbWithMetadata[] = [] for (const verbId of verbIds) { const verb = await this.getVerb_internal(verbId) const metadata = await this.getVerbMetadata(verbId) if (verb && metadata) { // Canonical hydration — reserved fields top-level, ONLY custom // fields in `metadata`. results.push(this.hydrateVerbWithMetadata(verb, metadata)) } } return results } catch (error) { prodLog.warn('[BaseStorage] GraphAdjacencyIndex lookup failed, falling back to shard iteration:', error) } } // Fallback - iterate by shards (WITH deserialization fix!) const results: HNSWVerbWithMetadata[] = [] for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { const verbFiles = await this.listCanonicalObjects(shardDir) for (const verbPath of verbFiles) { if (!verbPath.includes('/vectors.json')) continue try { const rawVerb = await this.readCanonicalObject(verbPath) if (!rawVerb) continue // CRITICAL - Deserialize connections Map from JSON storage format const verb = this.deserializeVerb(rawVerb) if (verb.targetId === targetId) { const metadataPath = getVerbMetadataPath(verb.id) const metadata = await this.readCanonicalObject(metadataPath) // Canonical hydration — reserved fields top-level, ONLY custom // fields in `metadata`. results.push(this.hydrateVerbWithMetadata(verb, metadata)) } } catch (error) { // Skip verbs that fail to load } } } catch (error) { // Skip shards that have no data } } return results } /** * Get verbs by type (Shard iteration with type filtering) */ protected async getVerbsByType_internal(verbType: string): Promise { // Iterate by shards (0x00-0xFF) instead of type-first paths const verbs: HNSWVerbWithMetadata[] = [] for (let shard = 0; shard < 256; shard++) { const shardHex = shard.toString(16).padStart(2, '0') const shardDir = `entities/verbs/${shardHex}` try { const verbFiles = await this.listCanonicalObjects(shardDir) for (const verbPath of verbFiles) { if (!verbPath.includes('/vectors.json')) continue try { const rawVerb = await this.readCanonicalObject(verbPath) if (!rawVerb) continue // Deserialize connections Map from JSON storage format const hnswVerb = this.deserializeVerb(rawVerb) // Filter by verb type if (hnswVerb.verb !== verbType) continue // Load metadata separately (optional), then combine via the // canonical hydration helper (defensive vector copy preserved) const metadata = await this.getVerbMetadata(hnswVerb.id) verbs.push( this.hydrateVerbWithMetadata( { ...hnswVerb, vector: [...hnswVerb.vector] }, metadata ) ) } catch (error) { // Skip verbs that fail to load } } } catch (error) { // Skip shards that have no data } } return verbs } /** * Delete a verb from storage (ID-first, O(1) delete) */ protected async deleteVerb_internal(id: string): Promise { // Direct O(1) delete with ID-first path const path = getVerbVectorPath(id) await this.deleteCanonicalObject(path) // Note: Type-specific counts will be decremented via metadata tracking // The real type is in metadata, accessible if needed via getVerbMetadata(id) } /** * Helper method to convert a Map to a plain object for serialization */ protected mapToObject( map: Map, valueTransformer: (value: V) => any = (v) => v ): Record { const obj: Record = {} for (const [key, value] of map.entries()) { obj[key.toString()] = valueTransformer(value) } return obj } /** * Save statistics data to storage (public interface) * @param statistics The statistics data to save */ public override async saveStatistics(statistics: StatisticsData): Promise { return this.saveStatisticsData(statistics) } /** * Get statistics data from storage (public interface) * @returns Promise that resolves to the statistics data or null if not found */ public override async getStatistics(): Promise { return this.getStatisticsData() } /** * Save statistics data to storage * This method should be implemented by each specific adapter * @param statistics The statistics data to save */ protected abstract override saveStatisticsData( statistics: StatisticsData ): Promise /** * Get statistics data from storage * This method should be implemented by each specific adapter * @returns Promise that resolves to the statistics data or null if not found */ protected abstract override getStatisticsData(): Promise }