/** * Metadata Index System * Maintains inverted indexes for fast metadata filtering * Automatically updates indexes when data changes */ import { StorageAdapter, resolveEntityField, NounMetadata, VerbMetadata } from '../coreTypes.js' import { SYSTEM_ENTITY_SCALARS, parseFieldAddress, UnresolvableFieldError, type FieldAddress } from '../db/fieldAddressing.js' import { splitNounMetadataRecord } from '../types/reservedFields.js' import { ColumnStore } from '../indexes/columnStore/ColumnStore.js' import type { MetadataIndexProvider } from '../plugin.js' import { MetadataIndexCache, MetadataIndexCacheConfig } from './metadataIndexCache.js' import { compareCodePoints } from './collation.js' import { prodLog } from './logger.js' import { getGlobalCache, UnifiedCache } from './unifiedCache.js' import { computeWatermarkVerdict, makeProjectionStamp, readStampedWatermark, type WatermarkVerdict, type WatermarkVerdictResult } from './projectionWatermark.js' import type { FactScanHandle } from '../db/factLog.js' import { NounType, VerbType, TypeUtils, NOUN_TYPE_COUNT, VERB_TYPE_COUNT } from '../types/graphTypes.js' import { SparseIndex, ChunkManager, AdaptiveChunkingStrategy, ChunkData, ChunkDescriptor, ZoneMap, compareNormalizedValues } from './metadataIndexChunking.js' import { EntityIdMapper } from './entityIdMapper.js' import { RoaringBitmap32, roaringLibraryInitialize } from './roaring/index.js' import { FieldTypeInference, FieldType } from './fieldTypeInference.js' import { BrainyError } from '../errors/brainyError.js' /** * Fields whose values are stored in the sparse index as BUCKETED values * (rounded to a coarser granularity to keep the index compact). Sorting * and any precision-sensitive comparison on these fields must bypass the * index and read the actual value directly from entity storage. * * Currently only timestamps are bucketed — they round to 1-minute windows * via `Math.floor(ts / 60000) * 60000` in the chunking layer. If any new * bucketed field is added (e.g. a compressed float), add it here too. */ const BUCKETED_INDEX_FIELDS: ReadonlySet = new Set([ 'system.createdAt', 'system.updatedAt' ]) export interface MetadataIndexEntry { field: string value: string | number | boolean ids: Set lastUpdated: number } export interface FieldIndexData { // Maps value -> count for quick filter discovery values: Record lastUpdated: number } export interface MetadataIndexStats { totalEntries: number totalIds: number fieldsIndexed: string[] lastRebuild: number indexSize: number // in bytes } /** * @description What {@link MetadataIndexManager.applyWatermarkCatchup} did, * for the caller's narration. * - `'noop'` — the verdict was `null`/`'adopt'`: the artifact already * reflects committed truth. Zero index writes. * - `'rescan'` — the verdict was `'rescan'`, OR a `'catchup'` verdict was * demoted (no window, or no fact log to scan) — either way a full * {@link MetadataIndexManager.rebuild} already ran; `reason` names why. * - `'caught-up'` — the `(from, to]` window folded successfully; the * artifact is stamped and flushed at `to`. */ export interface CatchupApplyResult { action: 'noop' | 'rescan' | 'caught-up' /** Present on `'rescan'` — why the fold could not proceed as a catchup. */ reason?: string /** Present on `'caught-up'` — the fact-log window that was folded. */ window?: { from: number; to: number } /** Present on `'caught-up'` — noun ops applied (add/update/delete). */ nounsApplied?: number /** Present on `'caught-up'` — verb ops applied (add/update/delete). */ verbsApplied?: number /** Present on `'caught-up'` — distinct committed generations folded. */ factsApplied?: number } export interface MetadataIndexConfig { maxIndexSize?: number // Max number of entries per field value (default: 10000) rebuildThreshold?: number // Rebuild if index is this % stale (default: 0.1) autoOptimize?: boolean // Auto-cleanup unused entries (default: true) // NOTE: the name-based indexedFields/excludeFields knobs died with the // field-addressing law ("no special names"): EVERY user field indexes, // whatever its name. Bulk-payload protection is value-SHAPE based and // uniform across all names (large arrays never become posting scalars; // long values index hashed) — shape is not a name carve-out. } export interface MetadataIndexOptions { entityIdMapper?: EntityIdMapper // Optional pre-configured EntityIdMapper (e.g., native from cor) } /** * Manages metadata indexes for fast filtering * Maintains inverted indexes: field+value -> list of IDs */ // Cardinality tracking for optimization decisions interface CardinalityInfo { uniqueValues: number totalValues: number distribution: 'uniform' | 'skewed' | 'sparse' updateFrequency: number lastAnalyzed: number } // Field statistics for smart optimization interface FieldStats { cardinality: CardinalityInfo queryCount: number rangeQueryCount: number exactQueryCount: number avgQueryTime: number indexType: 'hash' // Only 'hash' since all fields use chunked sparse indices with zone maps normalizationStrategy?: 'none' | 'precision' | 'bucket' } /** * Storage key for the metadata projection's watermark stamp — a sidecar * record beside the artifact (field registry + field indexes + chunked * sparse indexes + column-store segments + id-mapper records). Written LAST * in {@link MetadataIndexManager.flush} so stamp-after-data ordering holds * for every byte the stamp certifies. */ export const METADATA_INDEX_STAMP_KEY = '__index_metadata_watermark__' /** * Implements {@link MetadataIndexProvider}: the metadata-index surface Brainy * calls on whatever the `'metadataIndex'` provider resolves to (its own * manager, or Cor's native Rust engine). */ export class MetadataIndexManager implements MetadataIndexProvider { private storage: StorageAdapter private config: Required private isRebuilding = false private metadataCache: MetadataIndexCache private fieldIndexes = new Map() private dirtyFields = new Set() private lastFlushTime = Date.now() private autoFlushThreshold = 10 // Start with 10 for more frequent non-blocking flushes // --- Watermark stamp state (see utils/projectionWatermark for the law) --- /** Generation handed in via {@link stampWatermark}, awaiting the next flush. */ private pendingWatermark: number | null = null /** Last watermark durably stamped by this instance or loaded at init. */ private stampedWatermark: number | null = null /** The three-way verdict computed at init; null until init runs. */ private loadVerdict: WatermarkVerdictResult | null = null /** * Set only when {@link loadVerdict}.verdict is `'rescan'`: whether a * persisted artifact existed at load (even an unstamped/unverifiable * one) — distinguishes genuine first boot (nothing here yet, routine) * from an artifact whose watermark is unverifiable (the loud case). The * verdict value alone doesn't carry this distinction; see {@link * watermarkArtifactPresent}. */ private rescanArtifactPresent = false /** * @description THE BUILD-BESIDE SEAM (B3 Deliverable 3): when set (via * {@link beginShadow}), every live `addToIndex`/`removeFromIndex` call on * THIS instance also applies to the shadow instance — so a caller building * a fresh replacement manager beside this one (walking canonical into it) * never misses a write that lands during the build. This is the ONE seam * that makes build-beside possible without touching every call site: every * existing `AddToMetadataIndexOperation`/`RemoveFromMetadataIndexOperation` * (and the JS manager's own `rebuild()`/catchup fold) keep calling the SAME * serving instance exactly as before; only THIS instance knows it is also * mirroring to a shadow. Null = no build in flight (the overwhelmingly * common case; the check costs one property read per write). */ private shadow: MetadataIndexManager | null = null /** * @description Start mirroring every `addToIndex`/`removeFromIndex` call on * this instance to `shadow` too — see {@link shadow}'s JSDoc. The caller * owns sequencing: writes mirrored WHILE a canonical walk is populating * `shadow` may be clobbered by the walk's own (possibly stale) reads for * the same id; the caller closes that window with a bounded fact-log fold * AFTER the walk (the same mechanism {@link applyWatermarkCatchup} uses) * before treating `shadow` as authoritative. * @param shadow - The manager to mirror writes to. */ beginShadow(shadow: MetadataIndexManager): void { this.shadow = shadow } /** * @description Stop mirroring writes to a shadow (see {@link beginShadow}). * Idempotent; a no-op when no shadow is attached. */ endShadow(): void { this.shadow = null } // Cardinality and field statistics tracking private fieldStats = new Map() private cardinalityUpdateInterval = 100 // Update cardinality every N operations private operationCount = 0 // Smart normalization thresholds private readonly HIGH_CARDINALITY_THRESHOLD = 1000 private readonly TIMESTAMP_PRECISION_MS = 60000 // 1 minute buckets private readonly FLOAT_PRECISION = 2 // decimal places // Type-Field Affinity Tracking for intelligent NLP private typeFieldAffinity = new Map>() // nounType -> field -> count private totalEntitiesByType = new Map() // nounType -> total count // Phase 1b: Fixed-size type tracking (Stage 3 CANONICAL: 99.2% memory reduction vs Maps) // Uint32Array provides O(1) access via type enum index // 42 noun types × 4 bytes = 168 bytes (vs ~20KB with Map overhead) // 127 verb types × 4 bytes = 508 bytes (vs ~62KB with Map overhead) // Total: 676 bytes (vs ~85KB) = 99.2% memory reduction private entityCountsByTypeFixed = new Uint32Array(NOUN_TYPE_COUNT) // 168 bytes (Stage 3 CANONICAL: 42 types) private verbCountsByTypeFixed = new Uint32Array(VERB_TYPE_COUNT) // 508 bytes (Stage 3 CANONICAL: 127 types) // Unified cache for coordinated memory management private unifiedCache: UnifiedCache // File locking for concurrent write protection (prevents race conditions) private activeLocks = new Map() private lockPromises = new Map>() private lockTimers = new Map() // Track timers for cleanup // Adaptive Chunked Sparse Indexing // Reduces file count from 560k → 89 files (630x reduction) // ALL fields now use chunking - no more flat files // Removed sparseIndices Map - now lazy-loaded via UnifiedCache only // PROJECTED: Reduces metadata memory from 35GB → 5GB @ 1B scale (86% reduction from chunking strategy, not yet benchmarked) private chunkManager: ChunkManager private chunkingStrategy: AdaptiveChunkingStrategy // (Removed in 7.22.0) `dirtyChunks` and `dirtySparseIndices` Maps — // never populated since the sparse-index write path was deleted in 7.20.0 // (commit 11be039). The associated `flushDirtyMetadata()` no-op was also // removed. Column store is the single source of truth for indexed writes. // Roaring Bitmap Support // EntityIdMapper for UUID ↔ integer conversion private idMapper: EntityIdMapper // Field Type Inference (Production-ready value-based type detection) // Replaces unreliable pattern matching with DuckDB-inspired value analysis private fieldTypeInference: FieldTypeInference /** * Unified Column Store — replaces sparse index internals for filtering + sorting. * Created in the constructor (no storage needed for writes), storage discovery * happens in init(). Public so brainy.ts can call sortTopK directly for * unfiltered sort. */ public columnStore: ColumnStore constructor(storage: StorageAdapter, config: MetadataIndexConfig = {}, options: MetadataIndexOptions = {}) { this.storage = storage this.config = { maxIndexSize: config.maxIndexSize ?? 10000, rebuildThreshold: config.rebuildThreshold ?? 0.1, autoOptimize: config.autoOptimize ?? true // No name-based exclude/allow lists — the field-addressing law: every // user field indexes, whatever its name ('content', 'data', 'id', // 'vector', … included). Bulk payloads are kept out by uniform value- // SHAPE rules in extractIndexableFields (arrays >10 never become // posting scalars; >100-char values index hashed), never by name. } // Initialize metadata cache with similar config to search cache this.metadataCache = new MetadataIndexCache({ maxAge: 5 * 60 * 1000, // 5 minutes maxSize: 500, // 500 entries (field indexes + value chunks) enabled: true }) // Get global unified cache for coordinated memory management this.unifiedCache = getGlobalCache() // Use injected EntityIdMapper (e.g., native from cor) or create JS fallback this.idMapper = options.entityIdMapper ?? new EntityIdMapper({ storage, storageKey: 'brainy:entityIdMapper' }) // Initialize chunking system with roaring bitmap support this.chunkManager = new ChunkManager(storage, this.idMapper) this.chunkingStrategy = new AdaptiveChunkingStrategy() // Initialize Field Type Inference this.fieldTypeInference = new FieldTypeInference(storage) // Create column store — works immediately for writes (in-memory tail buffers). // Storage discovery (loading existing segments) happens in init(). this.columnStore = new ColumnStore() // Removed lazyLoadCounts() call from constructor // It was a race condition (not awaited) and read from wrong source. // Now properly called in init() after warmCache() loads the sparse index. } /** * Get the shared EntityIdMapper instance. Used by the ColumnStore and * other subsystems that need UUID ↔ u32 mapping without creating a * second mapper that could diverge. */ getIdMapper(): EntityIdMapper { return this.idMapper } /** * Initialize the metadata index manager * This must be called after construction and before any queries */ async init(): Promise { // Initialize roaring-wasm library (browser bundle requires async init) await roaringLibraryInitialize() // Load field registry to discover persisted indices // Must run first to populate fieldIndexes directory before warming cache await this.loadFieldRegistry() // Compute the watermark verdict for the persisted artifact BEFORE any // early return below — the verdict is recorded for every open, whether // the workspace is empty, rebuilding, or warm. Computed and exposed // only: today's rebuild triggers are unchanged (acting on 'catchup' — // the incremental fold — lands with the coordinator's wiring). await this.loadWatermarkVerdict() // Initialize EntityIdMapper (loads UUID ↔ integer mappings from storage) await this.idMapper.init() // Initialize column store storage discovery (load existing segment manifests). // The column store was created in the constructor for immediate writes; // this step loads persisted segments so queries can find existing data. try { await this.columnStore.init(this.storage, this.idMapper) } catch (err) { prodLog.warn('[MetadataIndex] Column store storage discovery failed:', err) } // Check if field registry was loaded successfully const hasFields = this.fieldIndexes.size > 0 if (!hasFields) { // Don't trust "empty" — field registry may be missing due to interrupted flush. // Probe storage for actual entities before concluding the workspace is empty. try { const probe = await this.storage.getNouns({ pagination: { limit: 1, offset: 0 } }) const hasEntities = (probe.totalCount ?? 0) > 0 || probe.items.length > 0 if (hasEntities) { console.warn( `[MetadataIndex] Field registry missing but ${probe.totalCount ?? 'unknown'} entities exist on disk — rebuilding index` ) await this.rebuild() return // rebuild handles warmCache + lazyLoadCounts internally } } catch { // Storage probe failed — genuinely empty or storage not ready } return // Truly empty workspace — nothing to warm } // Warm the cache with common fields (lazy loading optimization) // This loads the type column ('system.type') needed for type counts await this.warmCache() // Load type counts AFTER warmCache (sparse index is now cached) await this.lazyLoadCounts() // Phase 1b: Sync loaded counts to fixed-size arrays this.syncTypeCountsToFixed() } /** * Detect index corruption and automatically repair via rebuild * This catches the update() field asymmetry bug that causes 7 fields to accumulate per update * Corruption threshold: 100 avg metadata entries/entity, excluding __words__ (expected ~30) * * Removed from init() hot path for performance. Call explicitly via: * - brain.checkHealth() — returns health status * - brain.repairIndex() — runs detection + auto-repair */ async detectAndRepairCorruption(): Promise { const validation = await this.validateConsistency() if (!validation.healthy) { prodLog.warn(`⚠️ Index corruption detected (${validation.avgEntriesPerEntity.toFixed(1)} avg entries/entity)`) prodLog.warn('🔄 Auto-rebuilding index to repair...') // Clear and rebuild await this.clearAllIndexData() await this.rebuild() // Re-validate after rebuild const postRebuild = await this.validateConsistency() if (postRebuild.healthy) { prodLog.info(`✅ Index rebuilt successfully (${postRebuild.avgEntriesPerEntity.toFixed(1)} avg entries/entity)`) } else { prodLog.error( `❌ Index still appears corrupted after rebuild (${postRebuild.avgEntriesPerEntity.toFixed(1)} avg entries/entity). ` + `This may indicate a different issue.` ) } } } /** * Warm the cache by preloading common field sparse indices * This improves cache hit rates by loading frequently-accessed fields at startup * Target: >80% cache hit rate for typical workloads */ async warmCache(): Promise { // Common columns used in most queries — the frozen system keys, plus // legacy spellings for a pre-epoch-3 brain read before its rebuild runs. const commonFields = ['system.type', 'system.service', 'system.createdAt', 'noun'] prodLog.debug(`🔥 Warming metadata cache with common fields: ${commonFields.join(', ')}`) // Preload in parallel for speed await Promise.all( commonFields.map(async field => { try { await this.loadSparseIndex(field) } catch (error) { // Silently ignore if field doesn't exist yet // This maintains zero-configuration principle prodLog.debug(`Cache warming: field '${field}' not yet indexed`) } }) ) prodLog.debug('✅ Metadata cache warmed successfully') // Phase 1b: Also warm cache for top types (type-aware optimization) await this.warmCacheForTopTypes(3) } /** * Phase 1b: Warm cache for top types (type-aware optimization) * Preloads metadata indices for the most common entity types and their top fields * This significantly improves query performance for the most frequently accessed data * * @param topN Number of top types to warm (default: 3) */ async warmCacheForTopTypes(topN: number = 3): Promise { // Get top noun types by entity count const topTypes = this.getTopNounTypes(topN) if (topTypes.length === 0) { prodLog.debug('⏭️ Skipping type-aware cache warming: no types found yet') return } prodLog.debug(`🔥 Warming cache for top ${topTypes.length} types: ${topTypes.join(', ')}`) // For each top type, warm cache for its top fields for (const type of topTypes) { // Get fields with high affinity to this type const typeFields = this.typeFieldAffinity.get(type) if (!typeFields) continue // Sort fields by count (most common first) const topFields = Array.from(typeFields.entries()) .sort((a, b) => b[1] - a[1]) .slice(0, 5) // Top 5 fields per type .map(([field]) => field) if (topFields.length === 0) continue prodLog.debug(` 📊 Type '${type}' - warming fields: ${topFields.join(', ')}`) // Preload sparse indices for these fields in parallel await Promise.all( topFields.map(async field => { try { await this.loadSparseIndex(field) } catch (error) { // Silently ignore if field doesn't exist yet prodLog.debug(` ⏭️ Field '${field}' not yet indexed for type '${type}'`) } }) ) } prodLog.debug('✅ Type-aware cache warming completed') } /** * Full hydration — the {@link Brainy.warm} readiness seam for the metadata * index. Unlike {@link warmCache} / {@link warmCacheForTopTypes} (which * warm only a heuristic subset: common fields plus the top-N types' top * fields), this loads EVERY field's sparse index the field registry knows * about — a real read through {@link loadSparseIndex} into the unified * cache for each field, not a stat/existence check. Idempotent: an * already-cached field's `loadSparseIndex` call is a cheap cache hit. * * Re-reads the field registry first when `fieldIndexes` is empty (a warm() * call issued before `init()` populated it would otherwise hydrate * nothing), then loads every discovered field in parallel. */ async hydrateAll(): Promise { if (this.fieldIndexes.size === 0) { await this.loadFieldRegistry() } const fields = Array.from(this.fieldIndexes.keys()) if (fields.length === 0) { prodLog.debug('[MetadataIndex] hydrateAll: no persisted fields to hydrate') return } prodLog.debug(`[MetadataIndex] hydrateAll: loading ${fields.length} field(s) — ${fields.join(', ')}`) await Promise.all( fields.map(async field => { try { await this.loadSparseIndex(field) } catch (error) { // A single field's load failure doesn't abort the rest of the // hydration — warm() is a best-effort readiness step, never a // correctness gate (queries still demand-load on miss). prodLog.debug(`[MetadataIndex] hydrateAll: field '${field}' failed to load:`, error) } }) ) } /** * Acquire an in-memory lock for coordinating concurrent metadata index writes * Uses in-memory locks since MetadataIndexManager doesn't have direct file system access * @param lockKey The key to lock on (e.g., 'field_noun', 'sorted_timestamp') * @param ttl Time to live for the lock in milliseconds (default: 10 seconds) * @returns Promise that resolves to true if lock was acquired, false otherwise */ private async acquireLock( lockKey: string, ttl: number = 10000 ): Promise { const lockValue = `${Date.now()}_${Math.random()}` const expiresAt = Date.now() + ttl // Check if lock already exists and is still valid const existingLock = this.activeLocks.get(lockKey) if (existingLock && existingLock.expiresAt > Date.now()) { // Lock exists and is still valid - wait briefly and retry once await new Promise(resolve => setTimeout(resolve, 50)) // Check again after wait const recheckLock = this.activeLocks.get(lockKey) if (recheckLock && recheckLock.expiresAt > Date.now()) { return false // Lock still held } } // Acquire the lock this.activeLocks.set(lockKey, { expiresAt, lockValue }) // Schedule automatic cleanup when lock expires const timer = setTimeout(() => { this.releaseLock(lockKey, lockValue).catch((error) => { prodLog.debug(`Failed to auto-release expired lock ${lockKey}:`, error) }) }, ttl) this.lockTimers.set(lockKey, timer) return true } /** * Release an in-memory lock * @param lockKey The key to unlock * @param lockValue The value used when acquiring the lock (for verification) * @returns Promise that resolves when lock is released */ private async releaseLock( lockKey: string, lockValue?: string ): Promise { // If lockValue is provided, verify it matches before releasing if (lockValue) { const existingLock = this.activeLocks.get(lockKey) if (existingLock && existingLock.lockValue !== lockValue) { // Lock was acquired by someone else, don't release it return } } // Clear the timeout timer if it exists const timer = this.lockTimers.get(lockKey) if (timer) { clearTimeout(timer) this.lockTimers.delete(lockKey) } // Remove the lock this.activeLocks.delete(lockKey) } /** * Lazy load entity counts from the type column (O(n) where n = number of * types). The frozen key is 'system.type' (epoch 3); the legacy 'noun' * column is read as a fallback for a pre-epoch-3 brain observed before its * rebuild has run (e.g. a reader-mode open against an old writer). * FIX: Previously read from stats.nounCount which was SERVICE-keyed, not TYPE-keyed */ private async lazyLoadCounts(): Promise { try { // CRITICAL FIX - Clear counts before loading to prevent accumulation // Previously, counts accumulated across restarts causing 100x inflation this.totalEntitiesByType.clear() this.entityCountsByTypeFixed.fill(0) this.verbCountsByTypeFixed.fill(0) // PRIMARY (8.0+): rehydrate per-type counts from the column store's // type column — the authoritative on-disk source after a cold reopen. // Frozen key first ('system.type', epoch 3), legacy 'noun' as the // pre-rebuild fallback. // // The chunked sparse-index WRITE path was removed in 7.20.0 (commit // 11be039): new workspaces persist the type column ONLY to the column // store, never to a sparse-index blob. So the legacy sparse path below // finds nothing and leaves every count at 0 — which is exactly why // counts.byType/byTypeEnum/topTypes/allNounTypeCounts all read empty // after close()+reopen while find()/getNounCount() (different sources) // stay correct. The column store's per-value cardinality matches the warm // `updateTypeFieldAffinity` counts EXACTLY because both are driven from the // same `addToIndex` field set, in lockstep, with no visibility gate on // either — so this rehydration reproduces the warm values precisely. const indexedCols = this.columnStore ? this.columnStore.getIndexedFields() : [] const typeCol = indexedCols.includes('system.type') ? 'system.type' : indexedCols.includes('noun') ? 'noun' : null if (this.columnStore && typeCol) { const nounValues = await this.columnStore.getFilterValues(typeCol) for (const value of nounValues) { const bitmap = await this.columnStore.filter(typeCol, value) if (bitmap.size > 0) { // Use the stored value directly as the key (the legacy sparse path // did the same): it is already the normalized type string that // getNounFromIndex/getEntityCountByType expect, so syncTypeCountsToFixed // — called immediately after lazyLoadCounts in init() — copies it into // entityCountsByTypeFixed without re-normalization drift. this.totalEntitiesByType.set(value, bitmap.size) } } prodLog.debug(`✅ Rehydrated type counts from column store: ${this.totalEntitiesByType.size} types`) return } // LEGACY FALLBACK (pre-7.20.0 workspaces still on the chunked sparse index). const sparseCol = (await this.loadSparseIndex('system.type')) ? 'system.type' : 'noun' const nounSparseIndex = await this.loadSparseIndex(sparseCol) if (!nounSparseIndex) { // No column-store type column and no sparse index yet — counts will be // populated as entities are added. return } // Iterate through all chunks and sum up bitmap sizes by type for (const chunkId of nounSparseIndex.getAllChunkIds()) { const chunk = await this.chunkManager.loadChunk(sparseCol, chunkId) if (chunk) { for (const [type, bitmap] of chunk.entries) { const currentCount = this.totalEntitiesByType.get(type) || 0 this.totalEntitiesByType.set(type, currentCount + bitmap.size) } } } prodLog.debug(`✅ Loaded type counts from sparse index: ${this.totalEntitiesByType.size} types`) } catch (error) { // Silently fail - counts will be populated as entities are added // This maintains zero-configuration principle prodLog.debug('Could not load type counts:', error) } } /** * Phase 1b: Sync Map-based counts to fixed-size Uint32Arrays * This enables gradual migration from Maps to arrays while maintaining backward compatibility * Called periodically and on demand to keep both representations in sync */ private syncTypeCountsToFixed(): void { // Sync noun counts from totalEntitiesByType Map to entityCountsByTypeFixed array for (let i = 0; i < NOUN_TYPE_COUNT; i++) { const type = TypeUtils.getNounFromIndex(i) const count = this.totalEntitiesByType.get(type) || 0 this.entityCountsByTypeFixed[i] = count } // Sync verb counts from totalEntitiesByType Map to verbCountsByTypeFixed array // Note: Verb counts are currently tracked alongside noun counts in totalEntitiesByType // In the future, we may want a separate Map for verb counts for (let i = 0; i < VERB_TYPE_COUNT; i++) { const type = TypeUtils.getVerbFromIndex(i) const count = this.totalEntitiesByType.get(type) || 0 this.verbCountsByTypeFixed[i] = count } } /** * Phase 1b: Sync from fixed-size arrays back to Maps (reverse direction) * Used when Uint32Arrays are the source of truth and need to update Maps */ private syncTypeCountsFromFixed(): void { // Sync noun counts from array to Map for (let i = 0; i < NOUN_TYPE_COUNT; i++) { const count = this.entityCountsByTypeFixed[i] if (count > 0) { const type = TypeUtils.getNounFromIndex(i) this.totalEntitiesByType.set(type, count) } } // Sync verb counts from array to Map for (let i = 0; i < VERB_TYPE_COUNT; i++) { const count = this.verbCountsByTypeFixed[i] if (count > 0) { const type = TypeUtils.getVerbFromIndex(i) this.totalEntitiesByType.set(type, count) } } } /** * Update cardinality statistics for a field */ private updateCardinalityStats(field: string, value: any, operation: 'add' | 'remove'): void { // Initialize field stats if needed if (!this.fieldStats.has(field)) { this.fieldStats.set(field, { cardinality: { uniqueValues: 0, totalValues: 0, distribution: 'uniform', updateFrequency: 0, lastAnalyzed: Date.now() }, queryCount: 0, rangeQueryCount: 0, exactQueryCount: 0, avgQueryTime: 0, indexType: 'hash' }) } const stats = this.fieldStats.get(field)! const cardinality = stats.cardinality // Track unique values by checking fieldIndex counts const fieldIndex = this.fieldIndexes.get(field) const normalizedValue = this.normalizeValue(value, field) const currentCount = fieldIndex?.values[normalizedValue] || 0 if (operation === 'add') { // If this is a new value (count is 0), increment unique values if (currentCount === 0) { cardinality.uniqueValues++ } cardinality.totalValues++ } else if (operation === 'remove') { // If count will become 0, decrement unique values if (currentCount === 1) { cardinality.uniqueValues = Math.max(0, cardinality.uniqueValues - 1) } cardinality.totalValues = Math.max(0, cardinality.totalValues - 1) } // Update frequency tracking cardinality.updateFrequency++ // Periodically analyze distribution if (++this.operationCount % this.cardinalityUpdateInterval === 0) { this.analyzeFieldDistribution(field) } // Determine optimal index type based on cardinality this.updateIndexStrategy(field, stats) } /** * Analyze field distribution for optimization */ private analyzeFieldDistribution(field: string): void { const stats = this.fieldStats.get(field) if (!stats) return const cardinality = stats.cardinality const ratio = cardinality.uniqueValues / Math.max(1, cardinality.totalValues) // Determine distribution type if (ratio > 0.9) { cardinality.distribution = 'sparse' // High uniqueness (like IDs, timestamps) } else if (ratio < 0.1) { cardinality.distribution = 'skewed' // Low uniqueness (like status, type) } else { cardinality.distribution = 'uniform' // Balanced distribution } cardinality.lastAnalyzed = Date.now() } /** * Update index strategy based on field statistics */ private updateIndexStrategy(field: string, stats: FieldStats): void { const hasHighCardinality = stats.cardinality.uniqueValues > this.HIGH_CARDINALITY_THRESHOLD // All fields use chunked sparse indexing with zone maps stats.indexType = 'hash' // Determine normalization strategy for high cardinality NON-temporal fields // (Temporal fields are already bucketed in normalizeValue from the start!) if (hasHighCardinality) { // Check if field looks numeric (for float precision reduction) const fieldLower = field.toLowerCase() const looksNumeric = fieldLower.includes('count') || fieldLower.includes('score') || fieldLower.includes('value') || fieldLower.includes('amount') if (looksNumeric) { stats.normalizationStrategy = 'precision' // Reduce float precision } else { stats.normalizationStrategy = 'none' // Keep as-is for strings } } else { stats.normalizationStrategy = 'none' } } // ============================================================================ // Adaptive Chunked Sparse Indexing // All fields use chunking - simplified implementation // ============================================================================ /** * Load sparse index from storage */ private async loadSparseIndex(field: string): Promise { const indexPath = `__sparse_index__${field}` const unifiedKey = `metadata:sparse:${field}` return await this.unifiedCache.get(unifiedKey, async () => { try { const data = await this.storage.getMetadata(indexPath) if (data) { const sparseIndex = SparseIndex.fromJSON(data) // CRITICAL: Initialize chunk ID counter from existing chunks to prevent ID conflicts this.chunkManager.initializeNextChunkId(field, sparseIndex) // Add to unified cache (sparse indices are expensive to rebuild) const size = JSON.stringify(data).length this.unifiedCache.set(unifiedKey, sparseIndex, 'metadata', size, 200) return sparseIndex } } catch (error) { prodLog.debug(`Failed to load sparse index for field '${field}':`, error) } return undefined }) } /** * Save sparse index to storage */ private async saveSparseIndex(field: string, sparseIndex: SparseIndex): Promise { const indexPath = `__sparse_index__${field}` const unifiedKey = `metadata:sparse:${field}` const data = sparseIndex.toJSON() await this.storage.saveMetadata(indexPath, data) // Update unified cache const size = JSON.stringify(data).length this.unifiedCache.set(unifiedKey, sparseIndex, 'metadata', size, 200) } // flushDirtyMetadata — DELETED in 7.22.0. The dirtyChunks / dirtySparseIndices // accumulators it drained were never populated after the 7.20.0 column-store // refactor (commit 11be039). Column store flush happens in flush() directly. /** * Get IDs for a value using the legacy chunked sparse index. * * **This path is only for pre-7.20.0 workspaces** still being migrated to * the column store. The write path for sparse indices was removed in * commit `11be039` — new workspaces never get them. * * If neither the column store nor a sparse index covers the field, the * function throws `BrainyError(FIELD_NOT_INDEXED)`. Returning `[]` for a * genuinely unindexed field was a long-standing silent-empty bug class — * an empty result indistinguishable from "the data really isn't there." */ private async getIdsFromChunks(field: string, value: any): Promise { // Load sparse index via UnifiedCache (lazy loading) const sparseIndex = await this.loadSparseIndex(field) if (!sparseIndex) { // No column store match (we'd have returned in getIds()) AND no legacy // sparse index for this field — the field is genuinely not indexed. // Throw so find()-evaluation can log and translate to []. throw BrainyError.fieldNotIndexed(field) } // Find candidate chunks using zone maps and bloom filters const normalizedValue = this.normalizeValue(value, field) const candidateChunkIds = sparseIndex.findChunksForValue(normalizedValue) if (candidateChunkIds.length === 0) { return [] // No chunks contain this value } // Load chunks and collect integer IDs from roaring bitmaps const allIntIds = new Set() for (const chunkId of candidateChunkIds) { const chunk = await this.chunkManager.loadChunk(field, chunkId) if (chunk) { const bitmap = chunk.entries.get(normalizedValue) if (bitmap) { // Iterate through roaring bitmap integers for (const intId of bitmap) { allIntIds.add(intId) } } } } // Convert integer IDs back to UUIDs return this.idMapper.intsIterableToUuids(allIntIds) } /** * Get IDs for a range using chunked sparse index with zone maps and roaring bitmaps * Now fully lazy-loaded via UnifiedCache (no local sparseIndices Map) * Normalize min/max for timestamp bucketing before comparison */ private async getIdsFromChunksForRange( field: string, min?: any, max?: any, includeMin: boolean = true, includeMax: boolean = true ): Promise { // Load sparse index via UnifiedCache (lazy loading) const sparseIndex = await this.loadSparseIndex(field) if (!sparseIndex) { return [] // No chunked index exists yet } // Normalize min/max for consistent comparison with indexed values // (indexed values are bucketed for timestamps, so we must bucket the query bounds too) const normalizedMin = min !== undefined ? this.normalizeValue(min, field) : undefined const normalizedMax = max !== undefined ? this.normalizeValue(max, field) : undefined // Find candidate chunks using zone maps const candidateChunkIds = sparseIndex.findChunksForRange(normalizedMin, normalizedMax) if (candidateChunkIds.length === 0) { return [] } // Load chunks and filter by range, collecting integer IDs from roaring bitmaps const allIntIds = new Set() for (const chunkId of candidateChunkIds) { const chunk = await this.chunkManager.loadChunk(field, chunkId) if (chunk) { for (const [value, bitmap] of chunk.entries) { // Check if value is in range using numeric-aware comparison // (normalizeValue converts numbers to strings, so we must compare numerically) let inRange = true if (normalizedMin !== undefined) { const cmp = compareNormalizedValues(value, normalizedMin) inRange = inRange && (includeMin ? cmp >= 0 : cmp > 0) } if (normalizedMax !== undefined) { const cmp = compareNormalizedValues(value, normalizedMax) inRange = inRange && (includeMax ? cmp <= 0 : cmp < 0) } if (inRange) { // Iterate through roaring bitmap integers for (const intId of bitmap) { allIntIds.add(intId) } } } } } // Convert integer IDs back to UUIDs return this.idMapper.intsIterableToUuids(allIntIds) } /** * Get roaring bitmap for a field-value pair without converting to UUIDs * This is used for fast multi-field intersection queries using hardware-accelerated bitmap AND * Now fully lazy-loaded via UnifiedCache (no local sparseIndices Map) * @returns RoaringBitmap32 containing integer IDs, or null if no matches */ private async getBitmapFromChunks(field: string, value: any): Promise { // Load sparse index via UnifiedCache (lazy loading) const sparseIndex = await this.loadSparseIndex(field) if (!sparseIndex) { return null // No chunked index exists yet } // Find candidate chunks using zone maps and bloom filters const normalizedValue = this.normalizeValue(value, field) const candidateChunkIds = sparseIndex.findChunksForValue(normalizedValue) if (candidateChunkIds.length === 0) { return null // No chunks contain this value } // If only one chunk, return its bitmap directly if (candidateChunkIds.length === 1) { const chunk = await this.chunkManager.loadChunk(field, candidateChunkIds[0]) if (chunk) { const bitmap = chunk.entries.get(normalizedValue) return bitmap || null } return null } // Multiple chunks: collect all bitmaps and combine with OR const bitmaps: RoaringBitmap32[] = [] for (const chunkId of candidateChunkIds) { const chunk = await this.chunkManager.loadChunk(field, chunkId) if (chunk) { const bitmap = chunk.entries.get(normalizedValue) if (bitmap && bitmap.size > 0) { bitmaps.push(bitmap) } } } if (bitmaps.length === 0) { return null } if (bitmaps.length === 1) { return bitmaps[0] } // Combine multiple bitmaps with OR operation return RoaringBitmap32.orMany(bitmaps) } /** * Get IDs for multiple field-value pairs using fast roaring bitmap intersection * * This method provides 500-900x faster multi-field queries by: * - Using hardware-accelerated bitmap AND operations (SIMD: AVX2/SSE4.2) * - Avoiding intermediate UUID array allocations * - Converting integers to UUIDs only once at the end * * Example: { status: 'active', role: 'admin', verified: true } * Instead of: fetch 3 UUID arrays → convert to Sets → filter intersection * We do: fetch 3 bitmaps → hardware AND → convert final bitmap to UUIDs * * @param fieldValuePairs Array of field-value pairs to intersect * @returns Array of UUID strings matching ALL criteria */ /** * Multi-field intersection query: find entities matching ALL field-value pairs. * * Collects roaring bitmaps for each pair via the column store, then * intersects them using hardware-accelerated AND operations. * * @param fieldValuePairs - Array of { field, value } to intersect * @returns Array of entity UUID strings matching ALL pairs */ async getIdsForMultipleFields(fieldValuePairs: Array<{ field: string; value: any }>): Promise { if (fieldValuePairs.length === 0) return [] if (fieldValuePairs.length === 1) { return await this.getIds(fieldValuePairs[0].field, fieldValuePairs[0].value) } // Collect roaring bitmaps for each field-value pair via column store const bitmaps: RoaringBitmap32[] = [] for (const { field, value } of fieldValuePairs) { const bitmap = this.columnStore.hasField(field) ? await this.columnStore.filter(field, value) : await this.getBitmapFromChunks(field, value) ?? new RoaringBitmap32() if (bitmap.size === 0) return [] // Short circuit: empty intersection bitmaps.push(bitmap) } // Intersect all bitmaps (SIMD-accelerated roaring AND) let result = bitmaps[0] for (let i = 1; i < bitmaps.length; i++) { result = RoaringBitmap32.and(result, bitmaps[i]) } return result.size > 0 ? this.idMapper.intsIterableToUuids(result) : [] } // addToChunkedIndex — DELETED. Column store handles all writes. // removeFromChunkedIndex — DELETED. Column store handles removes via global deleted bitmap. /** * Get IDs matching a range query using zone maps */ /** * Range query: find all entity IDs where a field's value falls within a range. * * Routes through the column store for O(log n) binary search when available, * falls back to sparse index zone-map scan for fields not yet in the column store. * * @param field - Field name to query * @param min - Lower bound (undefined = no lower bound) * @param max - Upper bound (undefined = no upper bound) * @param includeMin - Whether to include the lower bound (default: true) * @param includeMax - Whether to include the upper bound (default: true) * @returns Array of matching entity UUID strings */ private async getIdsForRange( field: string, min?: any, max?: any, includeMin: boolean = true, includeMax: boolean = true ): Promise { // Track range query for field statistics if (this.fieldStats.has(field)) { const stats = this.fieldStats.get(field)! stats.rangeQueryCount++ } // Column store path: O(log n) binary search on sorted column. // Use raw values (no normalization) — column store stores exact values. // Thread includeMin/includeMax so strict lessThan/greaterThan stay strict // (the column store is no longer inclusive-only). if (this.columnStore && this.columnStore.hasField(field)) { const bitmap = await this.columnStore.rangeQuery(field, min, max, includeMin, includeMax) return this.idMapper.intsIterableToUuids(bitmap) } // Fallback: sparse index zone-map scan (legacy path) return await this.getIdsFromChunksForRange(field, min, max, includeMin, includeMax) } /** * Generate field index filename for filter discovery */ private getFieldIndexFilename(field: string): string { return `field_${field}` } // getValueChunkFilename, makeSafeFilename — DELETED. Sparse index file naming no longer needed. /** * Normalize value for consistent indexing with VALUE-BASED temporal detection * * Replaced unreliable field name pattern matching with production-ready * value-based detection (DuckDB-inspired). Analyzes actual data values, not names. * * NO FALLBACKS - Pure value-based detection only. */ private normalizeValue(value: any, field?: string): string { if (value === null || value === undefined) return '__NULL__' if (typeof value === 'boolean') return value ? '__TRUE__' : '__FALSE__' // VALUE-BASED temporal detection (no pattern matching!) // Analyze the VALUE itself to determine if it's a timestamp if (typeof value === 'number') { // Check if value looks like a Unix timestamp (2000-01-01 to 2100-01-01) const MIN_TIMESTAMP_S = 946684800 // 2000-01-01 in seconds const MAX_TIMESTAMP_S = 4102444800 // 2100-01-01 in seconds const MIN_TIMESTAMP_MS = MIN_TIMESTAMP_S * 1000 const MAX_TIMESTAMP_MS = MAX_TIMESTAMP_S * 1000 const isTimestampSeconds = value >= MIN_TIMESTAMP_S && value <= MAX_TIMESTAMP_S const isTimestampMilliseconds = value >= MIN_TIMESTAMP_MS && value <= MAX_TIMESTAMP_MS if (isTimestampSeconds || isTimestampMilliseconds) { // VALUE is a timestamp! Apply 1-minute bucketing const bucketSize = this.TIMESTAMP_PRECISION_MS // 60000ms = 1 minute const bucketed = Math.floor(value / bucketSize) * bucketSize return bucketed.toString() } } // Check if string value is ISO 8601 datetime if (typeof value === 'string') { // ISO 8601 pattern: YYYY-MM-DDTHH:MM:SS... const iso8601Pattern = /^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}/ if (iso8601Pattern.test(value)) { // VALUE is an ISO 8601 datetime! Convert to timestamp and bucket try { const timestamp = new Date(value).getTime() if (!isNaN(timestamp)) { const bucketSize = this.TIMESTAMP_PRECISION_MS const bucketed = Math.floor(timestamp / bucketSize) * bucketSize return bucketed.toString() } } catch { // Not a valid date, treat as string } } } // Apply smart normalization based on field statistics (for non-temporal fields) if (field && this.fieldStats.has(field)) { const stats = this.fieldStats.get(field)! const strategy = stats.normalizationStrategy if (strategy === 'precision' && typeof value === 'number') { // Reduce float precision for high cardinality numeric fields const rounded = Math.round(value * Math.pow(10, this.FLOAT_PRECISION)) / Math.pow(10, this.FLOAT_PRECISION) return rounded.toString() } } // Default normalization if (typeof value === 'number') return value.toString() if (Array.isArray(value)) { const joined = value.map(v => this.normalizeValue(v, field)).join(',') // Hash very long array values to avoid filesystem limits if (joined.length > 100) { return this.hashValue(joined) } return joined } const stringValue = String(value).toLowerCase().trim() // Hash very long string values to avoid filesystem limits if (stringValue.length > 100) { return this.hashValue(stringValue) } return stringValue } /** * Create a short hash for long values to avoid filesystem filename limits */ private hashValue(value: string): string { // Simple hash function to create shorter keys let hash = 0 for (let i = 0; i < value.length; i++) { const char = value.charCodeAt(i) hash = ((hash << 5) - hash) + char hash = hash & hash // Convert to 32-bit integer } return `__HASH_${Math.abs(hash).toString(36)}` } /** * Extract indexable field-value pairs from entity or metadata * * Handles BOTH entity structure (with top-level fields) AND record shapes * - Record-frame system scalars index under literal 'system.' keys * - The user's metadata bag indexes under bare keys — EVERY name (the * field-addressing law: no special names; 'level', 'data', 'id', * 'content', 'vector' in a bag are ordinary user fields) * - Record-frame plumbing (vector, connections, level, data, _rev, id) * never indexes — that is namespace routing, not a name carve-out * - Value-SHAPE rules apply uniformly to all names: arrays >10 never * become posting scalars; purely numeric key names (array indices) * skip; >100-char values index hashed (normalizeValue) */ private extractIndexableFields(data: any): Array<{ field: string, value: any }> { const fields: Array<{ field: string, value: any }> = [] // RECORD-FRAME-ONLY plumbing guard: on an entity/stored-record frame // these keys are the engine's structural payloads (the 384-dim vector, // embeddings, the adjacency list, the identity field) and never index. // This set is NEVER applied inside the user's metadata bag — under the // field-addressing law every user name indexes; a real vector-sized // value in a bag is kept out by the uniform array-size shape guard, not // by its name. const RECORD_PLUMBING = new Set(['vector', 'embedding', 'embeddings', 'connections', 'id']) // THE FROZEN INDEX KEY FORMAT (cross-engine, sealed 2026-08-03; the native // accelerator keys identically — epoch 3 rebuilds every brain onto it): // user fields index under BARE keys exactly as the caller wrote them; // the ten system scalars index under literal 'system.' keys — the // key IS the query address, so the two namespaces can never collide // inside the index again. // Frame kinds: 'entity-record' = entityForIndexing shape / v2 nested-bag // stored record (user fields nested under `metadata`; stray top-level // keys are DROPPED, not guessed); 'flat-record' = the LEGACY stored // metadata-record shape (user fields flat beside the engine's — sound to // split by name because the pre-law write door refused user metadata // carrying engine names, so a flat key matching a system name IS the // system value); 'user' = inside the metadata bag, where EVERY key is // the user's and indexes bare — collider names included. type Frame = 'entity-record' | 'flat-record' | 'user' const extract = (obj: any, prefix = '', frame: Frame = 'entity-record'): void => { for (const [key, value] of Object.entries(obj)) { let fullKey = prefix ? `${prefix}.${key}` : key if (!prefix && frame !== 'user') { if (key === 'metadata' && typeof value === 'object' && value !== null && !Array.isArray(value)) { extract(value, '', 'user') // the user's namespace: bare keys continue } if (key === 'type' || key === 'noun') { fullKey = 'system.type' // legacy 'noun' spelling folds into the frozen key } else if (SYSTEM_ENTITY_SCALARS.has(key) && key !== 'id') { fullKey = `system.${key}` } else if ( key === 'data' || key === '_rev' || key === 'level' || key === '_fmt' || RECORD_PLUMBING.has(key) ) { continue // plumbing / identity / format stamp — never indexed from a record frame } else if (frame === 'entity-record') { continue // stray entity-frame key: dropped, not guessed } // flat-record fallthrough: a non-system, non-plumbing key IS a user // field (flat beside the engine's, legacy shape) — indexes bare. } // User frame: NO name-based skips — every user field indexes, whatever // its name (the field-addressing law). Only the uniform value-shape // guards below apply. // Skip purely numeric field names (array indices converted to object keys) // Legitimate field names should never be purely numeric // This catches vectors stored as objects: {0: 0.1, 1: 0.2, ...} if (/^\d+$/.test(key)) continue // Skip large arrays (> 10 elements) - likely vectors or bulk data if (Array.isArray(value) && value.length > 10) continue if (value && typeof value === 'object' && !Array.isArray(value)) { // Recurse into nested objects (but not arrays), keeping the frame extract(value, fullKey, frame) } else if (Array.isArray(value) && value.length <= 10) { // Small arrays: index as multi-value field (all with same field name) // Example: tags: ["javascript", "node"] → field="tags", value="javascript" + field="tags", value="node" for (const item of value) { // Only index primitive values (not nested objects/arrays) if (item !== null && typeof item !== 'object') { fields.push({ field: fullKey, value: item }) } } } else { // Primitive value: index it under the frozen key computed above. // (The legacy 'type'→'noun' remap is gone — 'noun' columns die at // the epoch-3 rebuild; system.type is the one spelling.) fields.push({ field: fullKey, value }) } } } if (data && typeof data === 'object') { // Shape detection for the top frame: an object carrying a nested // `metadata` bag is the entityForIndexing shape; anything else is the // flat stored-record shape (user fields flat beside reserved ones). const entityShaped = 'metadata' in data && typeof data.metadata === 'object' && data.metadata !== null extract(data, '', entityShaped ? 'entity-record' : 'flat-record') } // Extract words for hybrid text search // Production-scale word limit (5000 words) // - Handles articles, chapters, and large documents // - Roaring Bitmaps + Chunked Sparse Index + LRU caching // - Int32 hashes store words as 4-byte values, not strings // // Memory managed by existing optimizations: // - Roaring Bitmaps: 90%+ compression for sparse data // - Chunked Sparse Index: ~50 values per chunk, lazy-loaded // - UnifiedCache LRU: Only hot chunks in memory // // A Bloom-filter hybrid could lift the per-entity word cap entirely if // full-document indexing at billion-entity scale ever becomes a need. const textContent = this.extractTextContent(data) if (textContent) { const MAX_WORDS_PER_ENTITY = 5000 // Handles articles/chapters, memory-safe at scale const allWords = this.tokenize(textContent) const words = allWords.slice(0, MAX_WORDS_PER_ENTITY) if (allWords.length > MAX_WORDS_PER_ENTITY) { // Log once per entity, not per word - avoids log spam prodLog.debug( `Entity text has ${allWords.length} words, indexing first ${MAX_WORDS_PER_ENTITY} for hybrid search` ) } for (const word of words) { // Hash word to int32 for memory efficiency (saves ~10GB at 1B scale) const wordHash = this.hashWord(word) fields.push({ field: '__words__', value: wordHash }) } } return fields } /** * Extract text content from entity data for word indexing * * Recursively extracts string values from data, excluding: * - vector, embedding, connections, level, id (internal fields) * - Arrays with more than 10 elements (likely vectors/bulk data) * - Numeric-only keys (array indices) * * @param data - Entity data or metadata * @returns Concatenated text content */ extractTextContent(data: any): string { if (data === null || data === undefined) return '' if (typeof data === 'string') return data if (typeof data === 'number' || typeof data === 'boolean') return String(data) if (Array.isArray(data)) { // Skip numeric arrays (vectors/embeddings), allow object/string arrays if (data.length > 0 && typeof data[0] === 'number') return '' return data.map(d => this.extractTextContent(d)).filter(Boolean).join(' ') } if (typeof data === 'object') { // Mirror of NEVER_INDEX for the text-extraction path: bulk structural // payloads only. `level` removed for the same reason (it silently dropped // a real user field from hybrid text search too). const skipKeys = new Set(['vector', 'embedding', 'embeddings', 'connections', 'id']) const texts: string[] = [] for (const [key, value] of Object.entries(data)) { // Skip internal fields and numeric keys (array indices) if (skipKeys.has(key) || /^\d+$/.test(key)) continue const text = this.extractTextContent(value) if (text) texts.push(text) } return texts.join(' ') } return '' } /** * Tokenize text into words for indexing * * - Converts to lowercase * - Removes punctuation * - Splits on whitespace * - Filters by length (2-50 chars) * - Deduplicates per entity * * @param text - Text content to tokenize * @returns Array of unique words */ tokenize(text: string): string[] { if (!text) return [] return text .toLowerCase() .replace(/[^\w\s]/g, ' ') // Remove punctuation .split(/\s+/) // Split on whitespace .filter(w => w.length >= 2 && w.length <= 50) // Length filter .filter((w, i, arr) => arr.indexOf(w) === i) // Dedupe per entity } /** * Hash word to int32 using FNV-1a * * FNV-1a is fast with low collision rate, suitable for word hashing. * Saves ~10GB at billion scale by avoiding string storage. * * @param word - Word to hash * @returns Int32 hash value */ hashWord(word: string): number { let hash = 2166136261 // FNV offset basis for (let i = 0; i < word.length; i++) { hash ^= word.charCodeAt(i) hash = Math.imul(hash, 16777619) // FNV prime } return hash | 0 // Convert to signed int32 } /** * Get entity IDs matching a text query * * Performs word-based text search using the __words__ index. * Returns IDs ranked by match count (entities with more matching words first). * * @param query - Text query to search for * @returns Array of { id, matchCount } sorted by matchCount descending */ async getIdsForTextQuery(query: string): Promise> { const queryWords = this.tokenize(query) if (queryWords.length === 0) return [] // Get IDs for each word hash const wordIdSets: Map[] = [] for (const word of queryWords) { const wordHash = this.hashWord(word) let ids: string[] try { ids = await this.getIds('__words__', wordHash) } catch (err) { // `__words__` is not yet indexed (e.g. no text content has been // added). Treat as no matches and continue — text search against // an empty workspace should return [], not throw. if (err instanceof BrainyError && err.type === 'FIELD_NOT_INDEXED') { ids = [] } else { throw err } } const idSet = new Map() for (const id of ids) { idSet.set(id, 1) } wordIdSets.push(idSet) } if (wordIdSets.length === 0) return [] // Count matches per entity const matchCounts = new Map() for (const idSet of wordIdSets) { for (const [id] of idSet) { matchCounts.set(id, (matchCounts.get(id) || 0) + 1) } } // Sort by match count descending return Array.from(matchCounts.entries()) .map(([id, matchCount]) => ({ id, matchCount })) .sort((a, b) => b.matchCount - a.matchCount) } /** * Add item to metadata indexes * * Now accepts either entity structure or plain metadata * - Entity structure: { id, type, confidence, weight, createdAt, metadata: {...} } * - Plain metadata: { noun, confidence, weight, createdAt, ... } * * @param id - Entity ID * @param entityOrMetadata - Either full entity structure or plain metadata (backward compat) * @param skipFlush - Skip automatic flush (used during batch operations) * @param deferWrites - Batch mode: buffer postings for a later flush * @param generation - Brainy's commit generation for this write (see the * {@link import('../plugin.js').MetadataIndexProvider} contract). This JS * manager keeps a single live view with no per-record delta log, so it * has no slot to store it — the value is accepted for contract parity * and forwarded to the shared id mapper (an injected native mapper * stamps its assignment records with it; the JS mapper ignores it). * The JS twin adopts full per-write stamping with the watermark train. */ async addToIndex(id: string, entityOrMetadata: any, skipFlush: boolean = false, deferWrites: boolean = false, generation?: bigint): Promise { const fields = this.extractIndexableFields(entityOrMetadata) // Sanity check for excessive indexed fields (indicates possible data issue) // Separate threshold for metadata fields vs word fields // - Metadata fields: warn if > 100 (indicates deeply nested metadata) // - Word fields: expected to be many for large documents, warn only for extreme cases const metadataFields = fields.filter(f => f.field !== '__words__') const wordFields = fields.filter(f => f.field === '__words__') if (metadataFields.length > 100) { prodLog.warn( `Entity ${id} has ${metadataFields.length} metadata fields (expected ~30). ` + `Possible deeply nested metadata. First 10 fields: ${metadataFields.slice(0, 10).map(f => f.field).join(', ')}` ) } // Words are expected to be many for large documents - only log for extreme cases if (wordFields.length > 5000) { prodLog.debug(`Entity ${id} has ${wordFields.length} indexed words (large document)`) } // Sort fields to process the type column first for type-field affinity // tracking ('system.type' is the frozen key; 'noun' died at epoch 3). fields.sort((a, b) => { if (a.field === 'system.type') return -1 if (b.field === 'system.type') return 1 return 0 }) // Update statistics and tracking for each field for (let i = 0; i < fields.length; i++) { const { field, value } = fields[i] this.updateCardinalityStats(field, value, 'add') this.updateTypeFieldAffinity(id, field, value, 'add', entityOrMetadata) await this.updateFieldIndex(field, value, 1) } // Write to column store — the single write path for all indexed data. // Converts extracted fields into a map and feeds the column store. // // extractIndexableFields emits ONE {field,value} entry per array element for // a multi-valued field (e.g. tags:['a','b','c'] → three 'tags' entries). We // must accumulate every repeated field into an array, not just __words__: // columnStore.addEntity expands an array value to one indexed entry per // element, so a scalar overwrite (last-value-wins) would index only the final // element and `contains` would miss the rest. if (this.columnStore) { // Thread the commit generation into the mint: an injected native mapper // stamps the assignment record's delta log with the real watermark // instead of a literal 0 (the JS mapper accepts and ignores it). const entityIntId = this.idMapper.getOrAssign(id, generation) const fieldsMap: Record = {} for (const { field, value } of fields) { if (field === '__words__') { // Always an array (keeps the field's multiValue manifest flag set even // for a single-word document). if (!fieldsMap.__words__) fieldsMap.__words__ = [] ;(fieldsMap.__words__ as unknown[]).push(value) } else if (field in fieldsMap) { // Repeated field → multi-valued. Promote the scalar to an array on the // second occurrence, then accumulate. const existing = fieldsMap[field] if (Array.isArray(existing)) { ;(existing as unknown[]).push(value) } else { fieldsMap[field] = [existing, value] } } else { fieldsMap[field] = value } } this.columnStore.addEntity(BigInt(entityIntId), fieldsMap) } // Adaptive auto-flush based on usage patterns if (!skipFlush) { const timeSinceLastFlush = Date.now() - this.lastFlushTime const shouldAutoFlush = this.dirtyFields.size >= this.autoFlushThreshold || // Size threshold (this.dirtyFields.size > 10 && timeSinceLastFlush > 5000) // Time threshold (5 seconds) if (shouldAutoFlush) { const startTime = Date.now() await this.flush() const flushTime = Date.now() - startTime // Adapt threshold based on flush performance if (flushTime < 50) { // Fast flush, can handle more entries this.autoFlushThreshold = Math.min(200, this.autoFlushThreshold * 1.2) } else if (flushTime > 200) { // Slow flush, reduce batch size this.autoFlushThreshold = Math.max(20, this.autoFlushThreshold * 0.8) } // Yield to event loop after flush to prevent blocking await this.yieldToEventLoop() } } // Invalidate cache for these fields for (const { field } of fields) { this.metadataCache.invalidatePattern(`field_values_${field}`) } // THE BUILD-BESIDE SEAM — see `shadow`'s JSDoc. Mirrors this write to a // shadow manager under construction, if one is attached. `skipFlush: // true` always: the shadow's own persistence is the build orchestrator's // job (it flushes once, after the swap — never mid-build, to avoid // colliding with this instance's own persisted keys). if (this.shadow) { await this.shadow.addToIndex(id, entityOrMetadata, true, false, generation) } } /** * Update field index with value count */ private async updateFieldIndex(field: string, value: any, delta: number): Promise { let fieldIndex = this.fieldIndexes.get(field) if (!fieldIndex) { // Load from storage if not in memory fieldIndex = await this.loadFieldIndex(field) ?? { values: {}, lastUpdated: Date.now() } this.fieldIndexes.set(field, fieldIndex) } const normalizedValue = this.normalizeValue(value, field) // Pass field for bucketing! fieldIndex.values[normalizedValue] = (fieldIndex.values[normalizedValue] || 0) + delta // Remove if count drops to 0 if (fieldIndex.values[normalizedValue] <= 0) { delete fieldIndex.values[normalizedValue] } fieldIndex.lastUpdated = Date.now() this.dirtyFields.add(field) } /** * Remove item from metadata indexes * * Now accepts either entity structure or plain metadata (same as addToIndex) * - Entity structure: { id, type, confidence, weight, createdAt, metadata: {...} } * - Plain metadata: { noun, confidence, weight, createdAt, ... } * * @param id - Entity ID to remove * @param metadata - Optional entity or metadata structure (if not provided, requires scanning all fields - slow!) * @param generation - Brainy's commit generation for this removal (see the * {@link import('../plugin.js').MetadataIndexProvider} contract). Accepted * for contract parity — this JS manager removes immediately (no tombstone * chain) and forwards it to the shared id mapper's `remove`, where an * injected native mapper tombstones the mapping at this generation. */ async removeFromIndex(id: string, metadata?: any, generation?: bigint): Promise { if (metadata) { const fields = this.extractIndexableFields(metadata) // Update statistics and tracking for (const { field, value } of fields) { this.updateCardinalityStats(field, value, 'remove') this.updateTypeFieldAffinity(id, field, value, 'remove', metadata) await this.updateFieldIndex(field, value, -1) this.metadataCache.invalidatePattern(`field_values_${field}`) } } // Remove from column store (global deleted bitmap) if (this.columnStore) { const intId = this.idMapper.getInt(id) if (intId !== undefined) { this.columnStore.removeEntity(BigInt(intId)) } } // Clean up ID mapper — must happen AFTER column store removal since it uses // idMapper.getInt(id). Prevents deleted IDs from persisting in the mapper // universe, which would cause ne/exists:false queries to return deleted entities. // The generation rides along so a native mapper tombstones the mapping at // the real commit watermark (the JS mapper ignores it). this.idMapper.remove(id, generation) await this.idMapper.flush() // THE BUILD-BESIDE SEAM — see `shadow`'s JSDoc. if (this.shadow) { await this.shadow.removeFromIndex(id, metadata, generation) } } /** * Get all IDs in the index */ async getAllIds(): Promise { // Use storage as the source of truth const allIds = new Set() // Storage.getNouns() is the definitive source of all entity IDs if (this.storage && typeof this.storage.getNouns === 'function') { try { const result = await this.storage.getNouns({ pagination: { limit: 100000 } }) if (result && result.items) { result.items.forEach((item) => { if (item.id) allIds.add(item.id) }) } } catch (e) { // If storage method fails, return empty array prodLog.warn('Failed to get all IDs from storage:', e) return [] } } return Array.from(allIds) } /** * Get IDs for a specific field-value combination using chunked sparse index */ /** * Point query: find all entity IDs where a field has a specific value. * * Routes through the column store for O(log n) binary search when available, * falls back to sparse index scan for fields not yet in the column store. * * @param field - Field name to query * @param value - Exact value to match * @returns Array of matching entity UUID strings */ /** * Report which index path a `where` clause on `field` will hit. Used by * `brain.explain()` so an operator can see *before* running a query whether * the field has any index entries at all. A `find({ where: { someField: ... } })` * against a field with no index entries returns `[]` silently — `explainField` * surfaces that as `path: 'none'` so the empty result has an explanation. */ async explainField(field: string): Promise<{ path: 'column-store' | 'sparse-chunked' | 'none' notes?: string }> { if (this.columnStore && this.columnStore.hasField(field)) { return { path: 'column-store', notes: 'O(log n) binary search + roaring bitmap. Best path.' } } const sparse = await this.loadSparseIndex(field) if (sparse) { return { path: 'sparse-chunked', notes: 'Chunked sparse index with zone maps and bloom filters.' } } return { path: 'none', notes: `No index entries for field "${field}". A find({ where: { ${field}: ... } }) ` + `will return an empty result regardless of whether matching entities exist on disk. ` + `Likely causes: (1) the writer registered the field in memory but has not flushed; ` + `(2) the field name does not match what was written (typo or casing); ` + `(3) the field is genuinely absent from all entities. Call requestFlush() on the ` + `writer or call brain.flush() before relying on the result.` } } /** * Resolve a `where: { field: value }` clause to entity UUIDs. * * Lookup order: * 1. **Column store** — the post-7.20.0 single source of truth. Fast. * 2. **Legacy sparse index** — only consulted for pre-7.20.0 workspaces * that haven't been migrated. Returns `[]` if no sparse data either. * * Throws `BrainyError(FIELD_NOT_INDEXED)` if the field has no entries in * either store. Callers in find()-evaluation catch this and translate to * an empty result with a logged warning. The throw aligns the production * `find()` path with the `brain.explain()` diagnostic, so a silently * empty result for an unindexed field is no longer possible. */ async getIds(field: string, value: any): Promise { // Track exact query for field statistics if (this.fieldStats.has(field)) { const stats = this.fieldStats.get(field)! stats.exactQueryCount++ } // Column store path: O(log n) binary search + roaring bitmap. // Use raw value (no normalization) — the column store stores exact values, // not the bucketed/stringified format the sparse index uses. if (this.columnStore && this.columnStore.hasField(field)) { const bitmap = await this.columnStore.filter(field, value) return this.idMapper.intsIterableToUuids(bitmap) } // Fallback: sparse index scan (legacy path during migration). If neither // store has the field, throw FIELD_NOT_INDEXED so the caller knows it's a // genuine "no such index" rather than "the value isn't there". return await this.getIdsFromChunks(field, value) } /** * Get all available values for a field (for filter discovery) */ async getFilterValues(field: string): Promise { // Check cache first const cacheKey = `field_values_${field}` const cachedValues = this.metadataCache.get(cacheKey) if (cachedValues) { return cachedValues } // Check in-memory field indexes first let fieldIndex = this.fieldIndexes.get(field) // If not in memory, load from storage if (!fieldIndex) { const loaded = await this.loadFieldIndex(field) if (loaded) { fieldIndex = loaded this.fieldIndexes.set(field, loaded) } } if (!fieldIndex) { return [] } const values = Object.keys(fieldIndex.values) // Cache the result this.metadataCache.set(cacheKey, values) return values } /** * Get all indexed fields (for filter discovery) */ async getFilterFields(): Promise { // Check cache first const cacheKey = 'all_filter_fields' const cachedFields = this.metadataCache.get(cacheKey) if (cachedFields) { return cachedFields } // Get fields from in-memory indexes and storage const fields = new Set(this.fieldIndexes.keys()) // Also scan storage for persisted field indexes (in case not loaded) // This would require a new storage method to list field indexes // For now, just use in-memory fields const fieldsArray = Array.from(fields) // Cache the result this.metadataCache.set(cacheKey, fieldsArray) return fieldsArray } /** * Convert Brainy Field Operator filter to simple field-value criteria for indexing */ private convertFilterToCriteria(filter: any): Array<{ field: string, values: any[] }> { const criteria: Array<{ field: string, values: any[] }> = [] if (!filter || typeof filter !== 'object') { return criteria } for (const [key, value] of Object.entries(filter)) { // Skip logical operators for now - handle them separately if (key === 'allOf' || key === 'anyOf' || key === 'not') continue if (value && typeof value === 'object' && !Array.isArray(value)) { // Handle Brainy Field Operators for (const [op, operand] of Object.entries(value)) { switch (op) { case 'oneOf': if (Array.isArray(operand)) { criteria.push({ field: key, values: operand }) } break case 'equals': case 'eq': criteria.push({ field: key, values: [operand] }) break case 'contains': // For contains, the operand is the value we're looking for in an array field criteria.push({ field: key, values: [operand] }) break case 'greaterThan': case 'lessThan': case 'between': // Range queries will be handled separately // Sorted index will be created/loaded when needed in getIdsForRange break default: break } } } else { // Direct value or array const values = Array.isArray(value) ? value : [value] criteria.push({ field: key, values }) } } return criteria } /** * Get IDs matching a Brainy Field Operator metadata filter using indexes where possible. * The optional `_opts` page bound is part of the provider contract for the native * index (early-stop at `offset+limit`); the JS index returns ALL matches and lets * the caller window them, so `_opts` is intentionally ignored here. */ /** Once-per-field throttle for the sparse-store did-you-mean WARN. */ private readonly warnedNeverCarried = new Set() /** * THE SPARSE-STORE CUT (ruled 2026-08-12): a WHERE filter naming a field * no row carries is SERVED OPERATOR-TRUTHFULLY (eq/range/contains → []; * ne/exists:false → all rows; exists:true → []) — the JS evaluator below * already computes exactly these truths via complements — with the * did-you-mean demoted to this throttled WARN. A fresh store's first * filtered read is a correct empty answer, never a refusal. orderBy and * ambiguous addresses KEEP their hard refusals (no truthful order * exists; ambiguity is a contract error — absence is data). */ /** Is this field known to the index at all (any row ever carried it)? */ private fieldRegistryHas(field: string): boolean { return this.fieldStats.has(field) } private warnNeverCarriedOnce(field: string): void { if (this.warnedNeverCarried.has(field)) return this.warnedNeverCarried.add(field) prodLog.warn( `[MetadataIndex] filter names field '${field}' which no row carries — ` + `serving the operator-truthful answer (empty for positive matches; ` + `the complement for ne/exists:false). If this is a typo, check the ` + `field name; refusals remain on orderBy.` ) } async getIdsForFilter(filter: any, _opts?: { limit?: number; offset?: number }): Promise { if (!filter || Object.keys(filter).length === 0) { return [] } // Handle logical operators if (filter.allOf && Array.isArray(filter.allOf)) { // For allOf, we need intersection of all sub-filters const allIds: string[][] = [] for (const subFilter of filter.allOf) { const subIds = await this.getIdsForFilter(subFilter) allIds.push(subIds) } if (allIds.length === 0) return [] if (allIds.length === 1) return allIds[0] // Set-based intersection O(n) — start with smallest set for optimal perf const sorted = allIds.sort((a, b) => a.length - b.length) let result = new Set(sorted[0]) for (let i = 1; i < sorted.length; i++) { const current = new Set(sorted[i]) result = new Set([...result].filter(id => current.has(id))) } return Array.from(result) } if (filter.anyOf && Array.isArray(filter.anyOf)) { // For anyOf, we need union of all sub-filters const unionIds = new Set() for (const subFilter of filter.anyOf) { const subIds = await this.getIdsForFilter(subFilter) subIds.forEach(id => unionIds.add(id)) } // Fix - Check for outer-level field conditions that need AND application // This handles cases like { anyOf: [...], vfsType: { exists: false } } // where the anyOf results must be intersected with other field conditions const outerFields = Object.keys(filter).filter( (k) => k !== 'anyOf' && k !== 'allOf' && k !== 'not' ) if (outerFields.length > 0) { // Build filter with just outer fields and get matching IDs const outerFilter: any = {} for (const field of outerFields) { outerFilter[field] = filter[field] } const outerIds = await this.getIdsForFilter(outerFilter) const outerIdSet = new Set(outerIds) // Intersect: anyOf union AND outer field conditions return Array.from(unionIds).filter((id) => outerIdSet.has(id)) } return Array.from(unionIds) } // Process field filters with range support const idSets: string[][] = [] // Capture field-not-indexed warnings so we log once per find() call, // not once per AND-clause inside it. const unindexedFields: string[] = [] for (const [rawField, condition] of Object.entries(filter)) { // Skip logical operators if (rawField === 'allOf' || rawField === 'anyOf' || rawField === 'not') continue // THE ONE ADDRESSING LAW (sealed 2026-08-03): every filter key routes // through parseFieldAddress — bare and 'metadata.'-prefixed spellings // address the user's fields (indexed under BARE keys), 'system.' // addresses the ten engine scalars (indexed under their literal // 'system.' keys). A malformed address (system., // plumbing in the system spelling) throws typed BEFORE any index read — // an accepted name either works or refuses. const address = parseFieldAddress(rawField, 'entity') const field = address.scope === 'system' ? `system.${address.field}` : address.field // Sparse-store cut: a user field no row carries serves operator-truth // below (the evaluators' complements are already correct) — announce // it once so a typo is findable without breaking a fresh store. if ( address.scope !== 'system' && !(this.columnStore && this.columnStore.hasField(field)) && !this.fieldRegistryHas(field) ) { this.warnNeverCarriedOnce(field) } let fieldResults: string[] = [] try { // The block below evaluates one field clause. If `getIds()` throws // FIELD_NOT_INDEXED (no column-store and no legacy sparse index for // this field), we treat the clause as matching zero entities. This // makes the production `find()` path consistent with the // `brain.explain()` diagnostic: an unindexed field returns no // results AND logs a warning, instead of silently returning []. if (condition && typeof condition === 'object' && !Array.isArray(condition)) { // Handle Brainy Field Operators (canonical operators defined) // See docs/api/README.md for complete operator reference // // Multiple operators on ONE field are AND-combined (intersected): e.g. // { greaterThan: 2009, lessThan: 2020 } requires BOTH bounds to hold. Each // operator computes its own match set, then intersects with the running set. let opIndex = 0 for (const [op, operand] of Object.entries(condition)) { const prevOpResults = fieldResults fieldResults = [] switch (op) { // ===== EQUALITY OPERATORS ===== // Canonical: 'eq' | Alias: 'equals' case 'equals': // Alias for 'eq' case 'eq': fieldResults = await this.getIds(field, operand) break // ===== NEGATION OPERATORS ===== // Canonical: 'ne' | Alias: 'notEquals' case 'notEquals': // Alias for 'ne' case 'ne': { // All ids EXCEPT those matching the value. Important for soft delete: // `deleted !== true` must include items WITHOUT a deleted field. The // excluded set is typically small (the matching value); compute the // complement as a bitmap difference over the int-id universe rather // than materializing the whole corpus as UUID strings to filter it. const excludeInts: number[] = [] // Sparse-store truth: a never-carried field has NOTHING to // exclude — the complement of nothing is EVERYTHING. getIds // throws FIELD_NOT_INDEXED there; the clause-level catch // would wrongly zero this NEGATIVE operator, so absorb it // here as the empty exclude set (the ruled operator-truth). let neMatches: string[] = [] try { neMatches = await this.getIds(field, operand) } catch { neMatches = [] } for (const uuid of neMatches) { const intId = this.idMapper.getInt(uuid) if (intId !== undefined) excludeInts.push(intId) } fieldResults = this.complementIds(excludeInts) break } // ===== MULTI-VALUE OPERATORS ===== // Canonical: 'in' | Alias: 'oneOf' case 'oneOf': // Alias for 'in' case 'in': if (Array.isArray(operand)) { const unionIds = new Set() for (const value of operand) { const ids = await this.getIds(field, value) ids.forEach(id => unionIds.add(id)) } fieldResults = Array.from(unionIds) } break // ===== GREATER THAN OPERATORS ===== // Canonical: 'gt' | Alias: 'greaterThan' case 'greaterThan': // Alias for 'gt' case 'gt': fieldResults = await this.getIdsForRange(field, operand, undefined, false, true) break // ===== GREATER THAN OR EQUAL OPERATORS ===== // Canonical: 'gte' | Alias: 'greaterThanOrEqual' case 'greaterThanOrEqual': // Alias for 'gte' case 'gte': fieldResults = await this.getIdsForRange(field, operand, undefined, true, true) break // ===== LESS THAN OPERATORS ===== // Canonical: 'lt' | Alias: 'lessThan' case 'lessThan': // Alias for 'lt' case 'lt': fieldResults = await this.getIdsForRange(field, undefined, operand, true, false) break // ===== LESS THAN OR EQUAL OPERATORS ===== // Canonical: 'lte' | Alias: 'lessThanOrEqual' case 'lessThanOrEqual': // Alias for 'lte' case 'lte': fieldResults = await this.getIdsForRange(field, undefined, operand, true, true) break // ===== RANGE OPERATOR ===== // between: [min, max] - inclusive range query case 'between': if (Array.isArray(operand) && operand.length === 2) { fieldResults = await this.getIdsForRange(field, operand[0], operand[1], true, true) } break // ===== ARRAY CONTAINS OPERATOR ===== // contains: value - check if array field contains value case 'contains': fieldResults = await this.getIds(field, operand) break // ===== EXISTENCE OPERATOR ===== // exists: boolean - check if field exists (any value) case 'exists': { // Column store path: rangeQuery with no bounds returns all IDs for the field const existsBitmap = (this.columnStore && this.columnStore.hasField(field)) ? await this.columnStore.rangeQuery(field) : await this.getExistsBitmapLegacy(field) if (operand) { // exists: true — entities that HAVE this field fieldResults = this.idMapper.intsIterableToUuids(existsBitmap) } else { // exists: false — entities that DON'T have this field (universe \ has-field) fieldResults = this.complementIds(existsBitmap) } break } // ===== MISSING OPERATOR ===== // missing: boolean - equivalent to exists: !boolean case 'missing': { const missingBitmap = (this.columnStore && this.columnStore.hasField(field)) ? await this.columnStore.rangeQuery(field) : await this.getExistsBitmapLegacy(field) if (operand) { // missing: true — entities that DON'T have this field (universe \ has-field) fieldResults = this.complementIds(missingBitmap) } else { // missing: false — entities that HAVE this field (same as exists: true) fieldResults = this.idMapper.intsIterableToUuids(missingBitmap) } break } } // Intersect this operator's matches with the running set (AND semantics // for multiple operators on the same field). if (opIndex > 0) { const prevSet = new Set(prevOpResults) fieldResults = fieldResults.filter((id) => prevSet.has(id)) } opIndex++ } } else { // Direct value match (shorthand for 'eq' operator) fieldResults = await this.getIds(field, condition) } } catch (err) { if (err instanceof BrainyError && err.type === 'FIELD_NOT_INDEXED') { unindexedFields.push(field) fieldResults = [] } else { throw err } } if (fieldResults.length > 0) { idSets.push(fieldResults) } else { // If any field has no matches, intersection will be empty if (unindexedFields.length > 0) { prodLog.warn( `[brainy] find() where-clause referenced unindexed field(s): ` + `${unindexedFields.join(', ')}. Returning []. Use ` + `brain.explain({ where: {...} }) for diagnostics.` ) } return [] } } if (idSets.length === 0) return [] if (unindexedFields.length > 0) { prodLog.warn( `[brainy] find() where-clause referenced unindexed field(s) ` + `${unindexedFields.join(', ')}; their clauses contributed no rows. ` + `Use brain.explain({ where: {...} }) for diagnostics.` ) } if (idSets.length === 1) return idSets[0] // Set-based intersection O(n) — start with smallest set for optimal perf const sortedSets = idSets.sort((a, b) => a.length - b.length) let resultSet = new Set(sortedSets[0]) for (let i = 1; i < sortedSets.length; i++) { const current = new Set(sortedSets[i]) resultSet = new Set([...resultSet].filter(id => current.has(id))) } return Array.from(resultSet) } /** * Legacy helper: get all entity int IDs that have any value for a field, * using the sparse index. Used by exists/missing operators when the * column store doesn't have data for the field. * * @param field - Field name * @returns Roaring bitmap of entity int IDs (or iterable for compatibility) * @private */ /** * All ids EXCEPT the excluded int-id set, computed as a roaring-bitmap difference * over the int-id universe. Used by the negation/absence operators (`ne`, * `exists:false`, `missing:true`) so they don't first materialize the ENTIRE * corpus as an array of UUID strings (plus a Set, plus an O(N) filter pass) just * to remove a small subset — only the final result is converted back to UUIDs. */ private complementIds(excludeInts: Iterable): string[] { const universe = new RoaringBitmap32() for (const intId of this.idMapper.getAllIntIds()) universe.add(intId) const exclude = new RoaringBitmap32() for (const intId of excludeInts) exclude.add(intId) return this.idMapper.intsIterableToUuids(RoaringBitmap32.andNot(universe, exclude)) } private async getExistsBitmapLegacy(field: string): Promise> { const allIntIds = new Set() const sparseIndex = await this.loadSparseIndex(field) if (sparseIndex) { for (const chunkId of sparseIndex.getAllChunkIds()) { const chunk = await this.chunkManager.loadChunk(field, chunkId) if (chunk) { for (const bitmap of chunk.entries.values()) { for (const intId of bitmap) { allIntIds.add(intId) } } } } } return allIntIds } /** * Get filtered IDs sorted by a field (production-scale sorting) * * **Performance Characteristics** (designed for billions of entities): * - **Filtering**: O(log n) using roaring bitmaps with SIMD acceleration * - **Field Loading**: O(k) where k = filtered result count (NOT O(n)) * - **Sorting**: O(k log k) in-memory (IDs + sort values only, NOT full entities) * - **Memory**: O(k) for k filtered results, independent of total entity count * * **Scalability**: * - Total entities: Billions (memory usage unaffected) * - Filtered set: Up to 10M (reasonable for in-memory sort of ID+value pairs) * - Pagination: Happens AFTER sorting, so only page entities are loaded * * **Example**: * ```typescript * // Production-scale: 1B entities, 100K match filter, sort by createdAt * const sortedIds = await metadataIndex.getSortedIdsForFilter( * { status: 'published', category: 'AI' }, * 'createdAt', * 'desc' * ) * // Returns: 100K sorted IDs * // Memory: ~5MB (100K IDs + 100K timestamps) * // Then caller paginates: sortedIds.slice(0, 20) and loads only 20 entities * ``` * * @param filter - Metadata filter criteria (uses roaring bitmaps) * @param orderBy - Field name to sort by (e.g., 'createdAt', 'title') * @param order - Sort direction: 'asc' (default) or 'desc' * @param topK - Optional page bound: produce only the top `K` sorted ids * (`offset + limit`) instead of the full sorted match set. A broad filter + * orderBy that returns one page no longer materializes hundreds of millions of * sorted ids at billion scale — the column store's top-K heap produces only the * page. Omit for the full sorted set. * @returns Promise - Entity IDs sorted by specified field * */ /** * Resolve the orderBy value for MANY entities in BATCHED metadata-record * reads — the sort path's one sanctioned value source (BRAINY-PROD-LATENCY-TRIAD). * * THE ASYMPTOTIC LAW THIS ENFORCES: an ordered read never does per-row * storage round-trips. The previous shape — `await getFieldValueForEntity` * per id, each opening the VECTOR record serially — cost 62–98ms × N on a * production filesystem brain: 3,224 rows took 199–317 SECONDS, silently. * The metadata RECORD (smaller, cached, batch-readable) carries everything * a sort can address: the ten system scalars top-level — EXACT values, no * bucketing loss — and the user's bag (v2 nested or legacy flat, resolved * through the shape-aware split). One batched read pass serves any N. * * The call-shape is pinned by tests (zero per-row reads, batch calls only) * so the serial loop cannot quietly return. * * @param ids - Entity ids to resolve (any size; reads are chunk-batched). * @param orderAddress - The parsed orderBy address (system or metadata scope). * @returns id → value map; ids whose record is missing map to `undefined` * (they sort LAST per the ordering contract — never dropped). */ private async resolveOrderValuesBatch( ids: string[], orderAddress: FieldAddress ): Promise> { const values = new Map() if (ids.length === 0) return values // Batch door, best first: BaseStorage's getNounMetadataBatch (native // batch or parallel reads inside), then the adapter-optional // getMetadataBatch, then chunked-parallel single reads — NEVER serial. const storage = this.storage as StorageAdapter & { getNounMetadataBatch?(ids: string[]): Promise> } const CHUNK = 500 const records = new Map() for (let i = 0; i < ids.length; i += CHUNK) { const chunk = ids.slice(i, i + CHUNK) if (typeof storage.getNounMetadataBatch === 'function') { const batch = await storage.getNounMetadataBatch(chunk) for (const [id, rec] of batch) records.set(id, rec) } else if (typeof storage.getMetadataBatch === 'function') { const batch = await storage.getMetadataBatch(chunk) for (const [id, rec] of batch) records.set(id, rec) } else { const loaded = await Promise.all( chunk.map(async (id) => [id, await storage.getNounMetadata(id)] as const) ) for (const [id, rec] of loaded) if (rec) records.set(id, rec) } } for (const id of ids) { const record = records.get(id) if (!record) { values.set(id, undefined) continue } // Shape-aware split serves both record eras: engine scalars from the // reserved half (EXACT timestamps — the bucketed index is never // consulted here), user fields from the bag. const { reserved, custom } = splitNounMetadataRecord( record as Record ) if (orderAddress.scope === 'system') { values.set( id, orderAddress.field === 'type' ? reserved.noun : (reserved as Record)[orderAddress.field] ) } else { let value: unknown = custom[orderAddress.field] if (value === undefined && orderAddress.field.includes('.')) { // Dotted user path: traverse INSIDE the bag. value = orderAddress.field .split('.') .reduce( (o, seg) => o && typeof o === 'object' ? (o as Record)[seg] : undefined, custom ) } values.set(id, value) } } return values } /** Once-per-field flag for the fallback-degradation announcement. */ private static announcedFallbackSorts = new Set() async getSortedIdsForFilter( filter: any, orderBy: string, order: 'asc' | 'desc' = 'asc', topK?: number ): Promise { // THE ONE ADDRESSING LAW — the orderBy address routes through the same // parse the filter path uses (the historical asymmetry where the filter // path understood 'metadata.' but the sorted path never did is dead). // Bare / 'metadata.' → the user's bare index key; 'system.' → the // literal frozen key; malformed addresses throw typed before any read. const orderAddress = parseFieldAddress(orderBy, 'entity') const orderKey = orderAddress.scope === 'system' ? `system.${orderAddress.field}` : orderAddress.field // DATA-AWARE REFUSAL (the did-you-mean): a bare address no user field // carries cannot mean anything as a sort key — and when the name collides // with a system scalar the caller almost certainly meant system.. // Refusing loudly with both candidates beats silently sorting nothing. if ( orderAddress.scope === 'metadata' && !(this.columnStore && this.columnStore.hasField(orderKey)) && !(await this.loadSparseIndex(orderKey)) ) { throw new UnresolvableFieldError(orderAddress.raw, 'entity') } // Column store path: O(K log S) sort via k-way merge across segments. // No per-entity storage reads, no precision loss from bucketing. if (this.columnStore && this.columnStore.hasField(orderKey)) { // Get filtered IDs from existing roaring bitmap path const hasFilter = filter && Object.keys(filter).length > 0 const filteredIds = hasFilter ? await this.getIdsForFilter(filter) : [] if (hasFilter && filteredIds.length === 0) return [] let sortedIntIds: bigint[] if (hasFilter) { // Build filter bitmap for the column store const filterBitmap = new RoaringBitmap32() for (const id of filteredIds) { const intId = this.idMapper.getInt(id) if (intId !== undefined) filterBitmap.add(intId) } // Page-bounded: produce only the top `topK` (offset+limit), not every // match, so a broad filter + orderBy returning one page stays O(matches // log K) heap, not a full sort materialization. const k = topK !== undefined ? Math.min(topK, filteredIds.length) : filteredIds.length sortedIntIds = await this.columnStore.filteredSortTopK( filterBitmap, orderKey, order, k ) } else { // Unfiltered sort — column store handles the full entity set efficiently sortedIntIds = await this.columnStore.sortTopK( orderKey, order, topK !== undefined ? Math.min(topK, this.idMapper.size) : this.idMapper.size ) } // Convert int IDs back to UUIDs. Number() narrowing is lossless — the // shipped EntityIdSpaceExceeded guard caps the JS mapper at u32. const sortedUuids = sortedIntIds .map(intId => this.idMapper.getUuid(Number(intId))) .filter((uuid): uuid is string => uuid !== undefined) // ORDERING CONTRACT (cross-engine, sealed): rows missing the field are // NEVER dropped — they sort LAST in both directions — and ties break by // id ascending. The column only contains rows that HAVE the field, so // (1) re-sort the page deterministically (value, then id) via ONE // batched value resolution — never per-row reads — and (2) append the // filtered rows the column omitted, id-ascending, filling any // remaining page budget. const pageValues = await this.resolveOrderValuesBatch(sortedUuids, orderAddress) const page = sortedUuids.map(id => ({ id, value: pageValues.get(id) })) page.sort((a, b) => this.compareAddressedValues(a.value, b.value, a.id, b.id, order)) let result = page.map(p => p.id) if (hasFilter) { const present = new Set(sortedUuids) if (topK === undefined || result.length < topK) { const missing = filteredIds.filter(id => !present.has(id)).sort() result = result.concat(missing) } } return topK !== undefined ? result.slice(0, topK) : result } // Fallback: no column serves this field. BOUNDED + ANNOUNCED, never // silent (the B2 no-silent-degradation law, BRAINY-PROD-LATENCY-TRIAD): // O(N) in row count but served by BATCHED metadata-record reads — the // serial per-row getNoun loop that turned 3,224 rows into a 199–317s // scan is dead, and the call-shape pin keeps it dead. const filteredIds = await this.getIdsForFilter(filter) if (filteredIds.length === 0) { return [] } if ( filteredIds.length > 500 && !MetadataIndexManager.announcedFallbackSorts.has(orderKey) ) { MetadataIndexManager.announcedFallbackSorts.add(orderKey) prodLog.warn( `[brainy] ordered read on '${orderKey}' has no column index — served by the ` + `batched fallback over ${filteredIds.length} rows (bounded, one batch pass; ` + `announced once per field). A native column for this field makes it O(K).` ) } const fallbackValues = await this.resolveOrderValuesBatch(filteredIds, orderAddress) const idValuePairs = filteredIds.map(id => ({ id, value: fallbackValues.get(id) })) idValuePairs.sort((a, b) => this.compareAddressedValues(a.value, b.value, a.id, b.id, order)) const sorted = idValuePairs.map(p => p.id) return topK !== undefined ? sorted.slice(0, topK) : sorted } /** * Get field value for a specific entity (helper for sorted queries) * * Three-path lookup: * * 1. **Bucketed fields** (timestamps) — the sparse index stores values * rounded to 1-minute buckets to keep the index compact for range * queries. That bucketing loses precision, so sorting must read the * actual value directly from entity storage. * * 2. **Custom fields with no sparse index** — VFS fields like `modified` * and `accessed`, plus any user custom field whose sparse index was * never built. Resolved from entity storage via `resolveEntityField`, * which knows the top-level-vs-metadata shape contract. * * 3. **Indexed fields** — strings, enums, and low-cardinality ints live * in the sparse roaring index. O(chunks) lookup, typically 1-10 chunks. * * **Performance**: * - Paths 1 & 2: O(1) entity load from storage (cached) * - Path 3: O(chunks) roaring bitmap lookup * * @param entityId - Entity UUID to get field value for * @param field - Field name to retrieve (e.g., 'createdAt', 'title') * @returns Promise - Field value or undefined if not found * * @public (called from brainy.ts for sorted queries) */ /** * The cross-engine ordering contract in one comparator (sealed 2026-08-03): * missing/null values sort LAST in BOTH directions — the direction flip * never moves them to the front — and ties break by id ascending, so an * ordered read is deterministic and identical on both engines. Numbers * compare numerically; everything else by code-point (UTF-8 byte) order, * matching the native column store exactly. */ private compareAddressedValues( aVal: any, bVal: any, aId: string, bId: string, order: 'asc' | 'desc' ): number { const aNull = aVal == null const bNull = bVal == null if (aNull || bNull) { if (aNull && bNull) return aId < bId ? -1 : aId > bId ? 1 : 0 return aNull ? 1 : -1 } let comparison = 0 if (aVal !== bVal) { if (typeof aVal === 'number' && typeof bVal === 'number') { comparison = aVal < bVal ? -1 : 1 } else { comparison = compareCodePoints(String(aVal), String(bVal)) } } if (comparison === 0) return aId < bId ? -1 : aId > bId ? 1 : 0 return order === 'asc' ? comparison : -comparison } async getFieldValueForEntity(entityId: string, field: string): Promise { // `field` arrives as a FROZEN INDEX KEY (bare = user metadata; // 'system.' = engine scalar). Storage fallbacks read the matching // side of the record — a system key reads the record scalar, a bare key // reads the user's metadata bag; the two can never shadow each other. const systemInner = field.startsWith('system.') ? field.slice('system.'.length) : null // Path 1: Bucketed fields need the actual (un-bucketed) value from storage. if (BUCKETED_INDEX_FIELDS.has(field)) { const noun = await this.storage.getNoun(entityId) if (!noun) return undefined return (noun as unknown as Record)[systemInner as string] } // Path 3 precondition: entity must be in the id mapper for bitmap lookup. const intId = this.idMapper.getInt(entityId) if (intId === undefined) { return undefined } // Load sparse index for this field (cached via UnifiedCache). const sparseIndex = await this.loadSparseIndex(field) // Path 2: No sparse index exists — fall back to entity storage. // Covers VFS custom fields (modified, accessed) and user fields not // yet indexed. resolveEntityField handles the shape contract. if (!sparseIndex) { const noun = await this.storage.getNoun(entityId) if (!noun) return undefined if (systemInner !== null) { return (noun as unknown as Record)[systemInner] } return (noun as { metadata?: Record }).metadata?.[field] } // Path 3: Search sparse index chunks for this entity's value. // Typically 1-10 chunks per field, so this is fast. for (const chunkId of sparseIndex.getAllChunkIds()) { const chunk = await this.chunkManager.loadChunk(field, chunkId) if (!chunk) continue // Check each value's roaring bitmap for our entity ID. // Roaring bitmap .has() is O(1) with SIMD optimization. for (const [value, bitmap] of chunk.entries) { if (bitmap.has(intId)) { return this.denormalizeValue(value, field) } } } return undefined } /** * Denormalize a value (reverse of normalizeValue) * * Converts normalized/stringified values back to their original type. * For most fields, this just parses numbers or returns strings as-is. * * **NOTE**: This is NOT used for timestamp sorting! Timestamp fields * (createdAt, updatedAt) are loaded directly from entity metadata by * getFieldValueForEntity() to avoid precision loss from bucketing. * * **Timestamp Bucketing (for range queries only)**: * - Indexed as: Math.floor(timestamp / 60000) * 60000 * - Used for: Range queries (gte, lte) where 1-minute precision is acceptable * - NOT used for: Sorting (requires exact millisecond precision) * * @param normalized - Normalized value string from index * @param field - Field name (used for type inference) * @returns Denormalized value in original type * * @private */ private denormalizeValue(normalized: string, field: string): any { // Try parsing as number (timestamps, integers, floats) const asNumber = Number(normalized) if (!isNaN(asNumber)) { return asNumber } // For strings, return as-is (already denormalized) return normalized } /** * Flush dirty entries to storage (non-blocking version) * NOTE: Sparse indices are flushed immediately in add/remove operations */ async flush(): Promise { // Always save field registry — even with no dirty fields. This tiny file // (list of field names) is the critical link that init() needs to discover // persisted indices. Without it, the index appears empty after restart. if (this.fieldIndexes.size > 0) { await this.saveFieldRegistry() } // Also always flush the EntityIdMapper — prevents ID collisions on restart await this.idMapper.flush() // Check if we have anything else to flush if (this.dirtyFields.size === 0) { // Nothing dirty — but a pending watermark still stamps (the registry // + id-mapper writes above are the only bytes this pass touched, and // they are durable at this point). Stamp-after-data holds. await this.writePendingStamp() return // No dirty field indexes to flush } // Process in smaller batches to avoid blocking const BATCH_SIZE = 20 const allPromises: Promise[] = [] // Flush field indexes in batches const dirtyFieldsArray = Array.from(this.dirtyFields) for (let i = 0; i < dirtyFieldsArray.length; i += BATCH_SIZE) { const batch = dirtyFieldsArray.slice(i, i + BATCH_SIZE) const batchPromises = batch.map(field => { const fieldIndex = this.fieldIndexes.get(field) return fieldIndex ? this.saveFieldIndex(field, fieldIndex) : Promise.resolve() }) allPromises.push(...batchPromises) // Yield to event loop between batches if (i + BATCH_SIZE < dirtyFieldsArray.length) { await this.yieldToEventLoop() } } // Wait for all operations to complete await Promise.all(allPromises) // Flush EntityIdMapper (UUID ↔ integer mappings) await this.idMapper.flush() // Save field registry for fast cold-start discovery await this.saveFieldRegistry() this.dirtyFields.clear() this.lastFlushTime = Date.now() // Flush column store tail buffers to L0 segments if (this.columnStore) { await this.columnStore.flush() } // STAMP-AFTER-DATA: the watermark stamp is the LAST write of the flush — // every byte it certifies (field indexes, registry, id-mapper records, // column-store segments) is durable before the stamp lands. A crash // anywhere above leaves the artifact behind-stamped or unstamped, which // verdicts as catchup/rescan on the next open — never a wrong adopt. await this.writePendingStamp() } /** * @description Record the committed generation this projection reflects. * The stamp is NOT written here — it is written as the final storage write * of the next {@link flush} (stamp-after-data ordering is a module * guarantee, not a caller obligation). The coordinator calls this with the * store's committed generation right before flushing. * @param generation - The committed generation every flushed byte reflects. */ stampWatermark(generation: number): void { this.pendingWatermark = generation } /** * @description The projection's current watermark: the stamp loaded at * init (or the last stamp durably written by this instance). Null = * unstamped (legacy artifact, first boot, or stamping never wired). */ watermark(): number | null { return this.stampedWatermark } /** * @description The three-way adoption verdict computed at init — * `'adopt'` (stamped == committed, zero work), `'catchup'` (stamped < * committed; the gap from {@link watermarkGap} awaits an incremental * fold), `'rescan'` (unstamped or stamped above committed — never * trusted). Null until init() has run. The coordinator (`Brainy.open()`) * consumes this via {@link applyWatermarkCatchup} right after init. */ watermarkVerdict(): WatermarkVerdict | null { return this.loadVerdict?.verdict ?? null } /** * @description The catch-up window `(from, to]` when the init verdict was * `'catchup'`; null otherwise. */ watermarkGap(): { from: number; to: number } | null { return this.loadVerdict?.gap ?? null } /** * @description Meaningful only when {@link watermarkVerdict} is * `'rescan'`: `true` when a persisted artifact existed at load (even an * unstamped/unverifiable one — real prior state, worth narrating loudly); * `false` for a genuine first boot (nothing persisted yet — a caller * should narrate this at a routine log level, not as an alarm, even * though the verdict value is the same `'rescan'` either way). */ watermarkArtifactPresent(): boolean { return this.rescanArtifactPresent } /** * @description Consume the three-way watermark verdict {@link * watermarkVerdict} computed at init — the cure for a crash-recovered * store whose canonical reads/counts recover every acked write but whose * metadata projection (flushed only periodically, not per-commit) keeps * serving the pre-crash state. Call once, right after `init()`, before * anything reads from this projection. * * - `null`/`'adopt'` → the artifact already reflects the store's * committed generation. Zero index writes. * - `'catchup'` → the caller-supplied `scan` (expected already opened * over `(watermarkGap().from, watermarkGap().to]`) is folded in, ONE * op at a time, through the SAME two legs {@link rebuild} uses (ADR-007 * A4 — one mechanism, never a second hand-rolled add/update shape): a * tombstone (`op.record === null`) retracts id-keyed (this projection * keeps no per-record delta log, so the pre-crash metadata for that id * — if any — is what a value-precise removal would need, and it isn't * available; the same tradeoff `remove()`'s null-metadata closure * already accepts elsewhere); an after-image retracts-then-reposts, so * an update never leaves stale postings under the old field values. A * fact outside the window is skipped defensively (belt: the scan is * already opened to the window; suspenders: this loop never trusts an * over-run). On success the artifact is stamped at `to` and flushed — * the same STAMP-AFTER-DATA door {@link flush} always writes through. * - `'rescan'` (or a `'catchup'` verdict with no window, or no `scan` to * fold — the store hosts no fact log) → the persisted artifact is * unverifiable; this method runs the existing {@link rebuild} itself * rather than leave the caller to notice and trigger it separately. * * @param scan - An open fact scan covering the catchup window (see * {@link Brainy.scanFacts}), or `null` when none is available/needed. * Ignored when the verdict is not `'catchup'`. * @returns What happened — see {@link CatchupApplyResult}. */ async applyWatermarkCatchup(scan: FactScanHandle | null): Promise { const verdict = this.watermarkVerdict() if (verdict === null || verdict === 'adopt') return { action: 'noop' } if (verdict === 'rescan') { await this.rebuild() return { action: 'rescan', reason: 'persisted artifact is unverifiable (unstamped, or stamped ABOVE the ' + "store's committed generation) — never adopting unverifiable state" } } // verdict === 'catchup' const window = this.watermarkGap() if (window === null) { await this.rebuild() return { action: 'rescan', reason: "'catchup' verdict exposed no window — cannot bound a fold" } } if (scan === null) { await this.rebuild() return { action: 'rescan', reason: `no fact log available to fold the (${window.from}, ${window.to}] catchup window` } } const { nounsApplied, verbsApplied, factsApplied } = await this.foldFactWindow(scan, window.from, window.to) this.stampWatermark(window.to) await this.flush() return { action: 'caught-up', window, nounsApplied, verbsApplied, factsApplied } } /** * @description Fold an open fact scan's `(fromGeneration, toGeneration]` * window into this projection, ONE op at a time, through the SAME two legs * {@link rebuild} uses (ADR-007 A4 — one mechanism, never a second * hand-rolled add/update shape): a tombstone retracts id-keyed; an * after-image retracts-then-reposts. THE CORE LOOP shared by {@link * applyWatermarkCatchup} (which stamps + flushes after) and {@link * buildBeside} (which does neither — persistence is the caller's job, * exactly once, after a swap). Never stamps, never flushes, never touches * storage beyond what `addToIndex`/`removeFromIndex` do internally * (skipFlush is always forced true). * @param scan - An open fact scan. * @param fromGeneration - Window lower bound (exclusive). * @param toGeneration - Window upper bound (inclusive). * @returns Counts for the caller's narration. */ private async foldFactWindow( scan: FactScanHandle, fromGeneration: number, toGeneration: number ): Promise<{ nounsApplied: number; verbsApplied: number; factsApplied: number }> { let nounsApplied = 0 let verbsApplied = 0 let factsApplied = 0 for await (const batch of scan.batches()) { for (const fact of batch.facts) { // Defensive containment: the scan is already opened to the window, // but a fact outside it is never applied regardless. if (fact.generation <= fromGeneration || fact.generation > toGeneration) continue const generation = BigInt(fact.generation) for (const op of fact.ops) { if (op.record === null) { // TOMBSTONE — the id-keyed removal path (no per-record delta // log to recover the old field values from). await this.removeFromIndex(op.id, undefined, generation) } else { // AFTER-IMAGE — retract any stale posting for this id, then // repost the new shape. Covers both a fresh add (nothing to // retract; a no-op-ish remove) and an update, through the same // two calls. await this.removeFromIndex(op.id, undefined, generation) await this.indexStoredRecord(op.id, op.record.metadata, { skipFlush: true, deferWrites: false, generation }) } if (op.kind === 'noun') nounsApplied++ else verbsApplied++ } factsApplied++ } } return { nounsApplied, verbsApplied, factsApplied } } /** * @description B3 Deliverable 3 — the shadow-build lifecycle's init: the * MINIMUM setup {@link buildBeside} needs, deliberately NOT the general * {@link init} sequence. Two reasons general `init()` is unsafe for a * build-beside shadow: * 1. `init()` unconditionally re-initializes the id mapper from storage * (`idMapper.init()`) — safe for a FRESH mapper, but this instance is * constructed with the CURRENTLY-SERVING manager's SHARED, already-live * mapper (identity is shared, never a second mapper — this train's own * law). Re-running its init() would DISCARD every not-yet-flushed * UUID↔int assignment sitting in memory, breaking the live manager's * own serving mid-build. * 2. `init()` loads the field registry and, on a registry that's * missing/empty while canonical has entities (exactly the shape a * rebuild is often invoked to FIX), triggers `rebuild()` itself — * WITHOUT `inMemoryOnly`, which would touch the shared storage keys * the live manager depends on. * What this DOES run: the WASM roaring-bitmap library init (idempotent; * needed before any column-store write) and the column store's OWN * segment-manifest discovery (read-only against shared storage; needed so * THIS instance's eventual post-swap flush continues segment numbering * correctly instead of colliding with the retiring manager's segments). */ private async initForShadowBuild(): Promise { await roaringLibraryInitialize() try { await this.columnStore.init(this.storage, this.idMapper) } catch (err) { prodLog.warn('[MetadataIndex] shadow build: column store storage discovery failed:', err) } } /** * @description B3 Deliverable 3 — THE ONLINE REBUILD's manager-side half: * populate THIS instance (expected fresh/empty, constructed with the SAME * storage + idMapper as the manager it will replace — see {@link * initForShadowBuild}) from canonical storage without ever touching the * shared storage keys the currently-serving manager depends on — no chunk * deletion, no flush, anywhere in this call. The caller (the brain's * rebuild-beside orchestrator) is responsible for: * 1. Attaching this instance as a {@link beginShadow} target on the OLD * manager BEFORE calling this, so live writes during the walk mirror * here too (best-effort — the walk below may still clobber a mirrored * write with a stale read for the same id; the fold after the walk is * what makes the final state authoritative, not the mirror). * 2. Swapping its own reference to this instance once this resolves. * 3. Calling {@link stampWatermark} + {@link flush} EXACTLY ONCE, after * the swap — this instance never persists itself. * @param committedGenerationAtStart - The store's committed generation * captured by the caller BEFORE this call — the fold's lower bound. * @returns The generation this instance's canonical data reflects once the * walk + fold settle — the fold's upper bound (writes committed after * this point but before the swap only reach this instance via the live * {@link beginShadow} mirror, so the caller re-reads the store's * committed generation right before stamping, rather than trusting this * return value as final). * @throws If canonical advanced during the walk but no fact log is * available to fold the gap — never a silently incomplete shadow. */ async buildBeside(committedGenerationAtStart: number): Promise { await this.initForShadowBuild() await this.rebuild({ inMemoryOnly: true }) const committedAfterWalk = this.storage.committedGeneration?.() ?? committedGenerationAtStart if (committedAfterWalk > committedGenerationAtStart) { const scan = this.storage.scanFacts?.({ fromGeneration: committedGenerationAtStart + 1, toGeneration: committedAfterWalk }) ?? null if (scan === null) { throw new Error( `MetadataIndexManager.buildBeside: canonical advanced from generation ` + `${committedGenerationAtStart} to ${committedAfterWalk} during the walk, but this ` + `store hosts no fact log to fold the gap — refusing a silently incomplete shadow` ) } await this.foldFactWindow(scan, committedGenerationAtStart, committedAfterWalk) } return committedAfterWalk } /** * @description Write the pending watermark stamp as a sidecar record — * always called AFTER the data it certifies is durable. A stamp-write * failure is fail-safe (the artifact stays unstamped/behind → rescan or * catchup on next open, never a wrong adopt) but is said out loud and the * pending stamp is retained for the next flush. */ private async writePendingStamp(): Promise { if (this.pendingWatermark === null) return const watermark = this.pendingWatermark try { await this.storage.saveMetadata(METADATA_INDEX_STAMP_KEY, { noun: 'IndexWatermark', ...makeProjectionStamp(watermark) }) this.stampedWatermark = watermark this.pendingWatermark = null } catch (error) { prodLog.error( `[MetadataIndex] failed to write watermark stamp (generation ${watermark}) — ` + `artifact stays behind-stamped (safe: verdicts catchup/rescan, never wrong-adopt); ` + `retrying on next flush:`, error ) } } /** * @description Read the artifact's stamp and compute the three-way verdict * against the store's committed generation. Unstamped state on a stamped * store verdicts `'rescan'` LOUDLY — never a silent adopt. * * MIGRATION COST: existing pre-stamp brains verdict `'rescan'` exactly * once (this open re-derives from source as it already does today); the * next flush stamps them, and every later open adopts. */ private async loadWatermarkVerdict(): Promise { const committed = this.storage.committedGeneration?.() ?? null let stamped: number | null = null try { const record = await this.storage.getMetadata(METADATA_INDEX_STAMP_KEY) stamped = readStampedWatermark(record) } catch { // An unreadable stamp is unstamped — the fail-safe direction. stamped = null } const result = computeWatermarkVerdict(stamped, committed) this.loadVerdict = result this.stampedWatermark = stamped if (result.verdict === 'rescan') { const artifactPresent = this.fieldIndexes.size > 0 || stamped !== null this.rescanArtifactPresent = artifactPresent if (artifactPresent) { prodLog.warn( `[MetadataIndex] watermark verdict: RESCAN — persisted index is ` + (stamped === null ? 'unstamped (legacy pre-stamp artifact, or a crash between data and stamp)' : `stamped at generation ${stamped}, ABOVE the store's committed generation ${committed}`) + ` — never adopting unverifiable state` ) } else { prodLog.debug( '[MetadataIndex] watermark verdict: rescan (no persisted artifact — first boot)' ) } } else if (result.verdict === 'catchup') { prodLog.info( `[MetadataIndex] watermark verdict: catchup — index stamped at generation ` + `${stamped}, store committed at ${committed}; the (${stamped}, ${committed}] ` + `window awaits an incremental fold (verdict exposed; the fold lands with the ` + `coordinator's wiring)` ) } } /** * Yield control back to the Node.js event loop * Prevents blocking during long-running operations */ private async yieldToEventLoop(): Promise { return new Promise(resolve => setImmediate(resolve)) } /** * Load field index from storage */ private async loadFieldIndex(field: string): Promise { const filename = this.getFieldIndexFilename(field) const unifiedKey = `metadata:field:${filename}` // Check unified cache first with loader function return await this.unifiedCache.get(unifiedKey, async () => { try { const cacheKey = `field_index_${filename}` // Check old cache for migration const cached = this.metadataCache.get(cacheKey) if (cached) { // Add to unified cache const size = JSON.stringify(cached).length this.unifiedCache.set(unifiedKey, cached, 'metadata', size, 1) // Low rebuild cost return cached } // Load from storage const indexId = `__metadata_field_index__${filename}` const data = await this.storage.getMetadata(indexId) if (data) { const fieldIndex = { values: data.values || {}, lastUpdated: data.lastUpdated || Date.now() } // Add to unified cache const size = JSON.stringify(fieldIndex).length this.unifiedCache.set(unifiedKey, fieldIndex, 'metadata', size, 1) // Also keep in old cache for now (transition period) this.metadataCache.set(cacheKey, fieldIndex) return fieldIndex } } catch (error) { // Field index doesn't exist yet } return null }) } /** * Save field index to storage with file locking */ private async saveFieldIndex(field: string, fieldIndex: FieldIndexData): Promise { const filename = this.getFieldIndexFilename(field) const lockKey = `field_index_${field}` const lockAcquired = await this.acquireLock(lockKey, 5000) // 5 second timeout if (!lockAcquired) { prodLog.warn( `Failed to acquire lock for field index '${field}', proceeding without lock` ) } try { const indexId = `__metadata_field_index__${filename}` const unifiedKey = `metadata:field:${filename}` // Add required 'noun' property for NounMetadata await this.storage.saveMetadata(indexId, { noun: 'MetadataFieldIndex', values: fieldIndex.values, lastUpdated: fieldIndex.lastUpdated }) // Update unified cache const size = JSON.stringify(fieldIndex).length this.unifiedCache.set(unifiedKey, fieldIndex, 'metadata', size, 1) // Invalidate old cache this.metadataCache.invalidatePattern(`field_index_${filename}`) } finally { if (lockAcquired) { await this.releaseLock(lockKey) } } } /** * Save field registry to storage for fast cold-start discovery * Solves 100x performance regression by persisting field directory * * This enables instant cold starts by discovering which fields have persisted indices * without needing to rebuild from scratch. Similar to how HNSW persists system metadata. * * Registry size: ~4-8KB for typical deployments (50-200 fields) * Scales: O(log N) - field count grows logarithmically with entity count */ private async saveFieldRegistry(): Promise { // Nothing to save if no fields indexed yet if (this.fieldIndexes.size === 0) { return } try { const registry = { noun: 'FieldRegistry', fields: Array.from(this.fieldIndexes.keys()), version: 1, lastUpdated: Date.now(), totalFields: this.fieldIndexes.size } await this.storage.saveMetadata('__metadata_field_registry__', registry) prodLog.debug(`📝 Saved field registry: ${registry.totalFields} fields`) } catch (error) { // Non-critical: Log warning but don't throw // System will rebuild registry on next cold start if needed prodLog.warn('Failed to save field registry:', error) } } /** * Load field registry from storage to populate fieldIndexes directory * Enables O(1) discovery of persisted sparse indices * * Called during init() to discover which fields have persisted indices. * Populates fieldIndexes Map with skeleton entries - actual sparse indices * are lazy-loaded via UnifiedCache when first accessed. * * Gracefully handles missing registry (first run or corrupted data). */ private async loadFieldRegistry(): Promise { try { const registry = await this.storage.getMetadata('__metadata_field_registry__') if (!registry?.fields || !Array.isArray(registry.fields)) { // Registry doesn't exist or is invalid - not an error, just first run prodLog.debug('📂 No field registry found - will build on first flush') return } // Populate fieldIndexes Map from discovered fields // Skeleton entries with empty values - sparse indices loaded lazily const lastUpdated = typeof registry.lastUpdated === 'number' ? registry.lastUpdated : Date.now() for (const field of registry.fields) { if (typeof field === 'string' && field.length > 0) { this.fieldIndexes.set(field, { values: {}, lastUpdated }) } } prodLog.info( `✅ Loaded field registry: ${registry.fields.length} persisted fields discovered\n` + ` Fields: ${registry.fields.slice(0, 5).join(', ')}${registry.fields.length > 5 ? '...' : ''}` ) } catch (error) { // Silent failure - registry not critical, will rebuild if needed prodLog.debug('Could not load field registry:', error) } } /** * Get list of persisted fields from storage (not in-memory) * Used during rebuild to discover which chunk files need deletion * * @returns Array of field names that have persisted sparse indices */ private async getPersistedFieldList(): Promise { try { const registry = await this.storage.getMetadata('__metadata_field_registry__') if (!registry?.fields || !Array.isArray(registry.fields)) { return [] } return registry.fields.filter((f: unknown) => typeof f === 'string' && f.length > 0) } catch (error) { prodLog.debug('Could not load persisted field list:', error) return [] } } /** * Delete all chunk files for a specific field * Used during rebuild to ensure clean slate * * @param field Field name whose chunks should be deleted */ private async deleteFieldChunks(field: string): Promise { try { // Load sparse index to get chunk IDs const indexPath = `__sparse_index__${field}` const sparseData = await this.storage.getMetadata(indexPath) if (sparseData) { const sparseIndex = SparseIndex.fromJSON(sparseData) // Delete all chunk files for this field for (const chunkId of sparseIndex.getAllChunkIds()) { await this.chunkManager.deleteChunk(field, chunkId) } // Delete the sparse index file itself. // Typed boundary: the storage metadata channel doubles as the delete // path — writing a JSON `null` tombstone clears the entry, but the // adapter signature only models real payloads. await this.storage.saveMetadata(indexPath, null as unknown as NounMetadata) } } catch (error) { // Silent failure - if we can't delete old chunks, rebuild will still work // (new chunks will be created, old ones become orphaned) prodLog.debug(`Could not clear chunks for field '${field}':`, error) } } /** * Clear ALL metadata index data from storage (for recovery) * Nuclear option for recovering from corrupted index state * * WARNING: This deletes all indexed data - requires full rebuild after! * Use when index is corrupted beyond normal rebuild repair. */ public async clearAllIndexData(): Promise { prodLog.warn('🗑️ Clearing ALL metadata index data from storage...') // Get all persisted fields const fields = await this.getPersistedFieldList() // Delete chunks and sparse indices for each field let deletedCount = 0 for (const field of fields) { await this.deleteFieldChunks(field) deletedCount++ } // Delete field registry. // Typed boundary: writing a JSON `null` tombstone clears the entry, but // the adapter signature only models real payloads. try { await this.storage.saveMetadata('__metadata_field_registry__', null as unknown as NounMetadata) } catch (error) { prodLog.debug('Could not delete field registry:', error) } // Clear in-memory state this.fieldIndexes.clear() this.dirtyFields.clear() this.unifiedCache.clear('metadata') this.totalEntitiesByType.clear() this.entityCountsByTypeFixed.fill(0) this.verbCountsByTypeFixed.fill(0) this.typeFieldAffinity.clear() // Clear EntityIdMapper. This is the explicit destructive path: the caller // asked for nuclear recovery of a corrupted index, so renumbering UUIDs is // intentional. Persisted int-keyed data (vector-mmap slots, graph // link-compression encodings) is invalidated by this op — the warning // below makes that explicit. Rebuild on its own does NOT clear the mapper. await this.idMapper.clear() // Clear chunk manager cache this.chunkManager.clearCache() prodLog.info(`✅ Cleared ${deletedCount} field indexes and all in-memory state`) prodLog.warn('⚠️ EntityIdMapper was cleared — any persisted int-keyed data ' + '(vector mmap slots, graph link-compression encodings, etc.) is now stale ' + 'and must be rebuilt from canonical sources.') prodLog.info('⚠️ Run brain.index.rebuild() to recreate the index from entity data') } /** * Get count of entities by type - O(1) operation using existing tracking * This exposes the production-ready counting that's already maintained */ getEntityCountByType(type: string): number { return this.totalEntitiesByType.get(type) || 0 } /** * Get total count of all entities - O(1) operation */ getTotalEntityCount(): number { let total = 0 for (const count of this.totalEntitiesByType.values()) { total += count } return total } /** * Get all entity types and their counts - O(1) operation. * `totalEntitiesByType` is populated by `updateTypeFieldAffinity` during add * operations (warm path) and rehydrated from the column store's 'noun' field * by `lazyLoadCounts` on init (cold reopen), so this is accurate both within a * session and after close()+reopen. */ getAllEntityCounts(): Map { return new Map(this.totalEntitiesByType) } // ============================================================================ // VFS Statistics Methods (uses existing Roaring bitmap infrastructure) // ============================================================================ /** * Read the type column's bitmap for one type value — frozen key first * ('system.type', epoch 3), legacy 'noun' as the pre-rebuild fallback. */ private async getTypeBitmap(type: string): Promise { return ( (await this.getBitmapFromChunks('system.type', type)) ?? (await this.getBitmapFromChunks('noun', type)) ) } /** * Get VFS entity count for a specific type using Roaring bitmap intersection * Uses hardware-accelerated SIMD operations (AVX2/SSE4.2) * @param type The noun type to query * @returns Count of VFS entities of this type */ async getVFSEntityCountByType(type: string): Promise { const vfsBitmap = await this.getBitmapFromChunks('isVFSEntity', true) const typeBitmap = await this.getTypeBitmap(type) if (!vfsBitmap || !typeBitmap) return 0 // Hardware-accelerated intersection + O(1) cardinality const intersection = RoaringBitmap32.and(vfsBitmap, typeBitmap) return intersection.size } /** * Get all VFS entity counts by type using Roaring bitmap operations * @returns Map of type -> VFS entity count */ async getAllVFSEntityCounts(): Promise> { const vfsBitmap = await this.getBitmapFromChunks('isVFSEntity', true) if (!vfsBitmap || vfsBitmap.size === 0) { return new Map() } const result = new Map() // Iterate through all known types and compute VFS count via intersection for (const type of this.totalEntitiesByType.keys()) { const typeBitmap = await this.getTypeBitmap(type) if (typeBitmap) { const intersection = RoaringBitmap32.and(vfsBitmap, typeBitmap) if (intersection.size > 0) { result.set(type, intersection.size) } } } return result } /** * Get total count of VFS entities - O(1) using Roaring bitmap cardinality * @returns Total VFS entity count */ async getTotalVFSEntityCount(): Promise { const vfsBitmap = await this.getBitmapFromChunks('isVFSEntity', true) return vfsBitmap?.size ?? 0 } // ============================================================================ // Phase 1b: Type Enum Methods (O(1) access via Uint32Arrays) // ============================================================================ /** * Get entity count for a noun type using type enum (O(1) array access) * More efficient than Map-based getEntityCountByType * @param type Noun type from NounTypeEnum * @returns Count of entities of this type */ getEntityCountByTypeEnum(type: NounType): number { const index = TypeUtils.getNounIndex(type) return this.entityCountsByTypeFixed[index] } /** * Get verb count for a verb type using type enum (O(1) array access) * @param type Verb type from VerbTypeEnum * @returns Count of verbs of this type */ getVerbCountByTypeEnum(type: VerbType): number { const index = TypeUtils.getVerbIndex(type) return this.verbCountsByTypeFixed[index] } /** * Get top N noun types by entity count (using fixed-size arrays) * Useful for type-aware cache warming and query optimization * @param n Number of top types to return * @returns Array of noun types sorted by count (highest first) */ getTopNounTypes(n: number): NounType[] { const types: Array<{ type: NounType; count: number }> = [] // Iterate through all noun types for (let i = 0; i < NOUN_TYPE_COUNT; i++) { const count = this.entityCountsByTypeFixed[i] if (count > 0) { const type = TypeUtils.getNounFromIndex(i) types.push({ type, count }) } } // Sort by count (descending) and return top N return types .sort((a, b) => b.count - a.count) .slice(0, n) .map(t => t.type) } /** * Get top N verb types by count (using fixed-size arrays) * @param n Number of top types to return * @returns Array of verb types sorted by count (highest first) */ getTopVerbTypes(n: number): VerbType[] { const types: Array<{ type: VerbType; count: number }> = [] // Iterate through all verb types for (let i = 0; i < VERB_TYPE_COUNT; i++) { const count = this.verbCountsByTypeFixed[i] if (count > 0) { const type = TypeUtils.getVerbFromIndex(i) types.push({ type, count }) } } // Sort by count (descending) and return top N return types .sort((a, b) => b.count - a.count) .slice(0, n) .map(t => t.type) } /** * Get all noun type counts as a Map (using fixed-size arrays) * More efficient than getAllEntityCounts for type-aware queries * @returns Map of noun type to count */ getAllNounTypeCounts(): Map { const counts = new Map() for (let i = 0; i < NOUN_TYPE_COUNT; i++) { const count = this.entityCountsByTypeFixed[i] if (count > 0) { const type = TypeUtils.getNounFromIndex(i) counts.set(type, count) } } return counts } /** * Get all verb type counts as a Map (using fixed-size arrays) * @returns Map of verb type to count */ getAllVerbTypeCounts(): Map { const counts = new Map() for (let i = 0; i < VERB_TYPE_COUNT; i++) { const count = this.verbCountsByTypeFixed[i] if (count > 0) { const type = TypeUtils.getVerbFromIndex(i) counts.set(type, count) } } return counts } /** * Get count of entities matching field-value criteria - queries chunked sparse index */ async getCountForCriteria(field: string, value: any): Promise { // Use chunked sparse indexing const ids = await this.getIds(field, value) return ids.length } /** * Get index statistics. * * Source-of-truth precedence (post-7.20.0 column-store-first architecture): * 1. **EntityIdMapper** — `idMapper.size` is the canonical entity count. * Every indexed entity gets a UUID→int mapping; nothing else is * consistent across instances. * 2. **ColumnStore** — `getIndexedFields()` is the canonical list of * indexed fields. `getFieldSizeSummary()` provides segment / tail * bookkeeping per field. * 3. **Legacy sparse-index registry** — only for pre-7.20.0 workspaces * whose data hasn't been migrated. `getPersistedFieldList()` may know * fields the column store doesn't yet, so we union them in. * * Prior implementation read from `this.fieldIndexes` + lazy-loaded sparse * indices, which silently returned `0` entries for any workspace written * after sparse-index writes were deleted in commit `11be039`. That * silent-zero defect is why this reads the column store first. */ async getStats(): Promise { const entityCount = this.idMapper.size // Field set: union of column-store fields and any legacy sparse-index // fields registered on disk. Exclude the `__words__` text index by // convention (it's not a metadata field in the public sense). const fields = new Set() if (this.columnStore) { for (const f of this.columnStore.getIndexedFields()) { if (f !== '__words__') fields.add(f) } } // Legacy fallback: pre-7.20.0 workspaces may have sparse-index registry // entries the column store doesn't know about yet. Surfacing them in the // field list lets the rest of the system migrate them on read. try { const legacyFields = await this.getPersistedFieldList() for (const f of legacyFields) { if (f !== '__words__') fields.add(f) } } catch { // Registry missing — nothing to add. } // `totalEntries` semantically means "distinct entities tracked by this // index". That's `idMapper.size`. `totalIds` is the sum of all // (field, value) → entityId postings — proxied by the segment/tail size // summary so we don't have to scan every bitmap. let totalIds = 0 if (this.columnStore) { for (const summary of this.columnStore.getFieldSizeSummary()) { if (summary.field === '__words__') continue totalIds += summary.tailSize // Segment count is a proxy; for a coarser-grained number we'd open // each segment cursor. Avoided here because stats() is on the hot // path for `brain.stats()` / health checks. totalIds += summary.segmentCount * 1 // segments contribute at least 1 posting } } return { totalEntries: entityCount, totalIds, fieldsIndexed: Array.from(fields).sort(), lastRebuild: Date.now(), indexSize: entityCount * 100 // rough estimate } } /** * Validate index consistency and detect corruption * Returns health status and recommendations for repair * * Counts metadata field entries only (excludes __words__ keyword index). * Corruption typically manifests as high avg entries/entity (expected ~30, corrupted can be 100+) * caused by the update() field asymmetry bug */ async validateConsistency(): Promise<{ healthy: boolean avgEntriesPerEntity: number entityCount: number indexEntryCount: number recommendation: string | null }> { const entityCount = this.idMapper.size // If no entities, index is trivially healthy if (entityCount === 0) { return { healthy: true, avgEntriesPerEntity: 0, entityCount: 0, indexEntryCount: 0, recommendation: null } } // Count total index entries across all fields (excluding keyword index) let indexEntryCount = 0 for (const field of this.fieldIndexes.keys()) { if (field === '__words__') continue // Keyword entries are expected to be high-volume const sparseIndex = await this.loadSparseIndex(field) if (sparseIndex) { for (const chunkId of sparseIndex.getAllChunkIds()) { const chunk = await this.chunkManager.loadChunk(field, chunkId) if (chunk) { for (const ids of chunk.entries.values()) { indexEntryCount += ids.size } } } } } const avgEntriesPerEntity = indexEntryCount / entityCount // Threshold: 100 metadata entries/entity is clearly corrupted (expected ~30) // __words__ keyword entries are excluded from this count since they can be 50-5000 per entity // This catches the update() asymmetry bug which causes 7 fields to accumulate per update const CORRUPTION_THRESHOLD = 100 const healthy = avgEntriesPerEntity <= CORRUPTION_THRESHOLD let recommendation: string | null = null if (!healthy) { recommendation = `Index corruption detected (${avgEntriesPerEntity.toFixed(1)} avg entries/entity, expected ~30). ` + `Run brain.index.clearAllIndexData() followed by brain.index.rebuild() to repair.` } return { healthy, avgEntriesPerEntity, entityCount, indexEntryCount, recommendation } } /** * @description Index one raw stored noun/verb record — THE ONE add leg * shared by {@link rebuild}'s canonical walk and {@link * applyWatermarkCatchup}'s fact-log fold (ADR-007 A4: one mechanism, * never a second hand-rolled shape). No conversion step is needed here: * a raw stored record (`storage.getNounMetadata`/`getVerbMetadata`, or a * fact's after-image `record.metadata`) is byte-identical — both read the * exact same canonical path — and already the v2 nested-bag * ("entity-record") shape {@link extractIndexableFields} expects. * @param id - Entity/relationship id. * @param storedMetadata - The raw stored metadata record. * @param opts.skipFlush - Forwarded to {@link addToIndex}. * @param opts.deferWrites - Forwarded to {@link addToIndex}. * @param opts.generation - Forwarded to {@link addToIndex}. */ private async indexStoredRecord( id: string, storedMetadata: unknown, opts: { skipFlush: boolean; deferWrites: boolean; generation?: bigint } ): Promise { await this.addToIndex(id, storedMetadata, opts.skipFlush, opts.deferWrites, opts.generation) } /** * Rebuild entire index from scratch using pagination * Non-blocking version that yields control back to event loop * Sparse indices now lazy-loaded via UnifiedCache (no need to clear Map) * * @param options.inMemoryOnly - B3 Deliverable 3 (build-beside): when * `true`, this call never touches the shared storage keys another, * currently-serving `MetadataIndexManager` over the SAME storage may * depend on — it skips deleting persisted legacy chunk files AND skips * the final `flush()` (which would otherwise write field indexes AND * flush the column store's tail buffers to shared segment keys, * colliding with a live manager's own writes). The caller ({@link * buildBeside}) owns persistence entirely — exactly once, after this * instance becomes the sole owner via an atomic swap. Default `false` * (every other caller keeps today's clear-then-persist behavior). */ async rebuild(options?: { inMemoryOnly?: boolean }): Promise { if (this.isRebuilding) return const inMemoryOnly = options?.inMemoryOnly ?? false this.isRebuilding = true try { prodLog.info('🔄 Starting non-blocking metadata index rebuild with batch processing...') prodLog.info(`📊 Storage adapter: ${this.storage.constructor.name}`) prodLog.info(`🔧 Batch processing available: ${!!this.storage.getMetadataBatch}`) // Clear existing indexes // No sparseIndices Map to clear - UnifiedCache handles eviction this.fieldIndexes.clear() this.dirtyFields.clear() // CRITICAL FIX - Clear type counts to prevent accumulation // Previously, counts accumulated across rebuilds causing incorrect values this.totalEntitiesByType.clear() this.entityCountsByTypeFixed.fill(0) this.verbCountsByTypeFixed.fill(0) this.typeFieldAffinity.clear() // Clear all cached sparse indices in UnifiedCache // This ensures rebuild starts fresh this.unifiedCache.clear('metadata') // Clear existing chunk files from storage to prevent overcounting. // Chunks are deleted first, then rebuilt. The field registry is NOT deleted // here — it's always saved at the end of rebuild via flush(). This ensures // that if rebuild fails partway, the next init() can still discover fields // and trigger another rebuild attempt. // // SKIPPED for inMemoryOnly: these are the SHARED storage keys a live // manager over the same storage may still be reading (see this // method's JSDoc) — deleting them before the swap is a live-read // hazard, not a cleanup. if (!inMemoryOnly) { prodLog.info('Clearing existing metadata index chunks from storage...') const existingFields = await this.getPersistedFieldList() if (existingFields.length > 0) { for (const field of existingFields) { await this.deleteFieldChunks(field) } prodLog.info(`Cleared ${existingFields.length} field indexes from storage`) } } // EntityIdMapper is intentionally NOT cleared here. Rebuild re-iterates // every entity in storage and calls idMapper.getOrAssign(uuid), which // returns the existing int for known UUIDs (no renumbering). This is the // foundational stability guarantee — vector-mmap slot indices, graph // link-compression encodings, and any other persisted int-keyed data // remain valid across a rebuild. Previously this line reset nextId to 1 // and renumbered every UUID by re-insertion order, silently breaking // any consumer that had persisted int-keyed data against the old map. // Stale entries for UUIDs no longer in storage persist (harmless memory // overhead); a dedicated prune step can be added if it ever matters. // The destructive wipe is still available via clearAllIndexData() → // idMapper.clear(), which is the explicit "recovery" path with the // appropriate warning about invalidating persisted int-keyed data. // Clear chunk manager cache this.chunkManager.clearCache() // Brainy 8.0 ships filesystem + memory storage only. Load all nouns // at once — the cloud-storage paginated branch was deleted alongside // the cloud adapters in step 7. let totalNounsProcessed = 0 { prodLog.info(`⚡ Loading all nouns at once (local storage)`) const result = await this.storage.getNouns({ pagination: { offset: 0, limit: 1000000 } // Effectively unlimited }) prodLog.info(`📦 Loading ${result.items.length} nouns with metadata...`) // Get all metadata in one batch if available const nounIds = result.items.map(noun => noun.id) let metadataBatch: Map if (this.storage.getMetadataBatch) { metadataBatch = await this.storage.getMetadataBatch(nounIds) prodLog.info(`✅ Loaded ${metadataBatch.size}/${nounIds.length} metadata objects`) } else { metadataBatch = new Map() for (const id of nounIds) { try { const metadata = await this.storage.getNounMetadata(id) if (metadata) metadataBatch.set(id, metadata) } catch (error) { prodLog.debug(`Failed to read metadata for ${id}:`, error) } } } for (const noun of result.items) { const metadata = metadataBatch.get(noun.id) if (metadata) { await this.indexStoredRecord(noun.id, metadata, { skipFlush: true, deferWrites: true }) } } totalNounsProcessed = result.items.length prodLog.info(`✅ Indexed ${totalNounsProcessed} nouns`) } // Rebuild verb metadata indexes — same single-pass local strategy. let totalVerbsProcessed = 0 { prodLog.info(`⚡ Loading all verbs at once (local storage)`) const result = await this.storage.getVerbs({ pagination: { offset: 0, limit: 1000000 } // Effectively unlimited }) prodLog.info(`📦 Loading ${result.items.length} verbs with metadata...`) const verbIds = result.items.map(verb => verb.id) let verbMetadataBatch: Map // Optional adapter capability: batched verb-metadata reads. Not part of // the StorageAdapter contract, so it is probed structurally. const batchCapableStorage = this.storage as StorageAdapter & { getVerbMetadataBatch?: (ids: string[]) => Promise> } if (batchCapableStorage.getVerbMetadataBatch) { verbMetadataBatch = await batchCapableStorage.getVerbMetadataBatch(verbIds) prodLog.info(`✅ Loaded ${verbMetadataBatch.size}/${verbIds.length} verb metadata objects`) } else { verbMetadataBatch = new Map() for (const id of verbIds) { try { const metadata = await this.storage.getVerbMetadata(id) if (metadata) verbMetadataBatch.set(id, metadata) } catch (error) { prodLog.debug(`Failed to read verb metadata for ${id}:`, error) } } } for (const verb of result.items) { const metadata = verbMetadataBatch.get(verb.id) if (metadata) { await this.indexStoredRecord(verb.id, metadata, { skipFlush: true, deferWrites: true }) } } totalVerbsProcessed = result.items.length prodLog.info(`✅ Indexed ${totalVerbsProcessed} verbs`) } // Flush to storage. The column store's flush() handles tail-buffer-to- // segment promotion + manifest persistence. // // SKIPPED for inMemoryOnly — see this method's JSDoc: flush() writes // the shared field-index keys AND flushes the column store's tail // buffers to shared segment keys, which would race a live manager's // own flushes over the SAME storage. The caller flushes exactly once, // after the swap. if (!inMemoryOnly) { prodLog.debug('💾 Flushing metadata index to storage...') await this.flush() } prodLog.info(`✅ Metadata index rebuild completed! Processed ${totalNounsProcessed} nouns and ${totalVerbsProcessed} verbs`) } finally { this.isRebuilding = false } } /** * Get field statistics for optimization and discovery */ async getFieldStatistics(): Promise> { // Initialize stats for fields we haven't seen yet for (const field of this.fieldIndexes.keys()) { if (!this.fieldStats.has(field)) { this.fieldStats.set(field, { cardinality: { uniqueValues: 0, totalValues: 0, distribution: 'uniform', updateFrequency: 0, lastAnalyzed: Date.now() }, queryCount: 0, rangeQueryCount: 0, exactQueryCount: 0, avgQueryTime: 0, indexType: 'hash' }) } } return new Map(this.fieldStats) } /** * Get field cardinality information */ async getFieldCardinality(field: string): Promise { const stats = this.fieldStats.get(field) return stats ? stats.cardinality : null } /** * Get all field names with their cardinality (for query optimization) */ async getFieldsWithCardinality(): Promise> { const fields: Array<{ field: string; cardinality: number; distribution: string }> = [] for (const [field, stats] of this.fieldStats) { fields.push({ field, cardinality: stats.cardinality.uniqueValues, distribution: stats.cardinality.distribution }) } // Sort by cardinality (low cardinality fields are better for filtering) fields.sort((a, b) => a.cardinality - b.cardinality) return fields } /** * Get optimal query plan based on field statistics */ async getOptimalQueryPlan(filters: Record): Promise<{ strategy: 'exact' | 'range' | 'hybrid' fieldOrder: string[] estimatedCost: number }> { const fieldOrder: string[] = [] let hasRangeQueries = false let totalEstimatedCost = 0 // Analyze each filter for (const [field, value] of Object.entries(filters)) { const stats = this.fieldStats.get(field) if (!stats) continue // Check if this is a range query if (typeof value === 'object' && value !== null && !Array.isArray(value)) { hasRangeQueries = true } // Estimate cost based on cardinality const cardinality = stats.cardinality.uniqueValues const estimatedCost = Math.log2(Math.max(1, cardinality)) totalEstimatedCost += estimatedCost fieldOrder.push(field) } // Sort fields by cardinality (process low cardinality first) fieldOrder.sort((a, b) => { const statsA = this.fieldStats.get(a) const statsB = this.fieldStats.get(b) if (!statsA || !statsB) return 0 return statsA.cardinality.uniqueValues - statsB.cardinality.uniqueValues }) return { strategy: hasRangeQueries ? 'hybrid' : 'exact', fieldOrder, estimatedCost: totalEstimatedCost } } /** * Export field statistics for analysis */ async exportFieldStats(): Promise { const stats: any = { fields: {}, summary: { totalFields: this.fieldStats.size, highCardinalityFields: 0, sparseFields: 0, skewedFields: 0, uniformFields: 0 } } for (const [field, fieldStats] of this.fieldStats) { stats.fields[field] = { cardinality: fieldStats.cardinality, queryStats: { total: fieldStats.queryCount, exact: fieldStats.exactQueryCount, range: fieldStats.rangeQueryCount, avgTime: fieldStats.avgQueryTime }, indexType: fieldStats.indexType, normalization: fieldStats.normalizationStrategy } // Update summary if (fieldStats.cardinality.uniqueValues > this.HIGH_CARDINALITY_THRESHOLD) { stats.summary.highCardinalityFields++ } switch (fieldStats.cardinality.distribution) { case 'sparse': stats.summary.sparseFields++ break case 'skewed': stats.summary.skewedFields++ break case 'uniform': stats.summary.uniformFields++ break } } return stats } /** * Update type-field affinity tracking for intelligent NLP * Tracks which fields commonly appear with which entity types */ private updateTypeFieldAffinity(entityId: string, field: string, value: any, operation: 'add' | 'remove', metadata?: any): void { // Only track affinity for user fields (plus the type column itself, // which drives detection). Engine columns carry the literal 'system.' // prefix under the frozen key format. if (field.startsWith('system.') && field !== 'system.type') return // For the type column ('system.type'), the value IS the entity type let entityType: string | null = null if (field === 'system.type') { // This is the type definition itself entityType = this.normalizeValue(value, field) // Pass field for bucketing! } else if (metadata && (metadata.noun ?? metadata.type)) { // Extract entity type from the source shape: stored records carry it // under 'noun', entity-for-indexing views under 'type'. entityType = this.normalizeValue(metadata.noun ?? metadata.type, 'system.type') } else { // No type information available, skip affinity tracking return } if (!entityType) return // No type found, skip affinity tracking // Initialize affinity tracking for this type if (!this.typeFieldAffinity.has(entityType)) { this.typeFieldAffinity.set(entityType, new Map()) } if (!this.totalEntitiesByType.has(entityType)) { this.totalEntitiesByType.set(entityType, 0) } const typeFields = this.typeFieldAffinity.get(entityType)! if (operation === 'add') { // Increment field count for this type const currentCount = typeFields.get(field) || 0 typeFields.set(field, currentCount + 1) // Update total entities of this type (only count once per entity — // the type column appears exactly once per entity) if (field === 'system.type') { const newCount = this.totalEntitiesByType.get(entityType)! + 1 this.totalEntitiesByType.set(entityType, newCount) // Phase 1b: Also update fixed-size array // Try to parse as noun type - if it matches a known type, update the array try { const nounTypeIndex = TypeUtils.getNounIndex(entityType as NounType) this.entityCountsByTypeFixed[nounTypeIndex] = newCount } catch { // Not a recognized noun type, skip fixed-size array update } } } else if (operation === 'remove') { // Decrement field count for this type const currentCount = typeFields.get(field) || 0 if (currentCount > 1) { typeFields.set(field, currentCount - 1) } else { typeFields.delete(field) } // Update total entities of this type if (field === 'system.type') { const total = this.totalEntitiesByType.get(entityType)! if (total > 1) { const newCount = total - 1 this.totalEntitiesByType.set(entityType, newCount) // Phase 1b: Also update fixed-size array try { const nounTypeIndex = TypeUtils.getNounIndex(entityType as NounType) this.entityCountsByTypeFixed[nounTypeIndex] = newCount } catch { // Not a recognized noun type, skip fixed-size array update } } else { this.totalEntitiesByType.delete(entityType) this.typeFieldAffinity.delete(entityType) // Phase 1b: Also zero out fixed-size array try { const nounTypeIndex = TypeUtils.getNounIndex(entityType as NounType) this.entityCountsByTypeFixed[nounTypeIndex] = 0 } catch { // Not a recognized noun type, skip fixed-size array update } } } } } /** * Get fields that commonly appear with a specific entity type * Returns fields with their affinity scores (0-1) */ async getFieldsForType(nounType: NounType): Promise> { const typeFields = this.typeFieldAffinity.get(nounType) const totalEntities = this.totalEntitiesByType.get(nounType) if (!typeFields || !totalEntities) { return [] } const fieldsWithAffinity: Array<{ field: string affinity: number occurrences: number totalEntities: number }> = [] for (const [field, count] of typeFields.entries()) { const affinity = count / totalEntities // 0-1 score fieldsWithAffinity.push({ field, affinity, occurrences: count, totalEntities }) } // Sort by affinity (most common fields first) fieldsWithAffinity.sort((a, b) => b.affinity - a.affinity) return fieldsWithAffinity } /** * Get type-field affinity statistics for analysis */ async getTypeFieldAffinityStats(): Promise<{ totalTypes: number averageFieldsPerType: number typeBreakdown: Record }> }> { const typeBreakdown: Record = {} let totalFields = 0 for (const [nounType, fieldsMap] of this.typeFieldAffinity.entries()) { const totalEntities = this.totalEntitiesByType.get(nounType) || 0 const fields = Array.from(fieldsMap.entries()) // Get top 5 fields for this type const topFields = fields .map(([field, count]) => ({ field, affinity: count / totalEntities })) .sort((a, b) => b.affinity - a.affinity) .slice(0, 5) typeBreakdown[nounType] = { totalEntities, uniqueFields: fieldsMap.size, topFields } totalFields += fieldsMap.size } return { totalTypes: this.typeFieldAffinity.size, averageFieldsPerType: totalFields / Math.max(1, this.typeFieldAffinity.size), typeBreakdown } } }