brainy/src/storage/baseStorage.ts
David Snelling a3467e1f9b feat: temporal VFS — file content joins the Model-B immutability model
The temporal model had a hole exactly where files were concerned: every
entity write is an immutable generation with before-images, but VFS content
BYTES lived under an eager refCount GC left over from the pre-8.0 design —
unlink could physically destroy bytes that in-window history still
referenced, and overwrite never released the old hash at all (an unbounded
silent leak whose accidental byproduct was the only thing "preserving"
history). Reading the past could therefore return a stale field, a dangling
hash, or nothing, depending on luck.

Fix: blob reclamation becomes a HISTORY decision instead of a LIVENESS
decision. Each blob's metadata now carries historyRefCount alongside the
live refCount:

- The commit seam counts one history reference per persisted before-image
  record carrying a content hash (commitTransaction staging and the
  group-commit flush), recorded BEFORE the record-set persists and carried
  in the generation delta (blobHashes — always present on new deltas, so
  compaction only falls back to reading records for pre-contract
  generations). An aborted transaction compensates best-effort.
- unlink/rmdir/overwrite drop ONLY the live reference (BlobStorage.delete →
  release; overwrite finally releases the superseded hash — cancelling the
  dedup increment on same-content rewrites and closing the leak), and only
  AFTER the canonical mutation commits, so a failed delete can never leave a
  live file whose bytes compaction might reclaim.
- History compaction is the ONE reclamation point: after deleting a
  generation's record-set it releases that set's references and physically
  reclaims any hash at zero live AND zero history references. Pins are
  exempt automatically. Crash ordering is over-count-only in every path
  (record before persist, release after delete), so a crash can leak until
  the scrub recounts but can never reclaim bytes a retained generation
  needs. scrubBlobHistoryRefCounts() restores exactness; existing stores get
  a one-time marker-gated backfill on open, failing into leak-safe mode
  (reclamation disabled) rather than guessing.

On top of the protected history, the temporal API the generational model
always implied:

- vfs.readFile(path, { asOf }) — the exact bytes as of a generation or Date,
  materialized from the history (pinned view released so compaction is
  never blocked by a read).
- vfs.history(path) — FileVersion[] ascending ({ generation, timestamp,
  hash, size, mimeType? }), the newest entry being the live state.
- Overwrites now refresh the file entity's data/embedding text — semantic
  search and the data field previously served the FIRST version's text
  forever (the stale-field defect a consumer's incident recovery depended
  on by luck).

Integration suite (temporal-vfs.test.ts): per-version exact reads +
history listing, leak-fix + history protection on overwrite, rm keeps bytes
readable, compaction reclaims past-window bytes and preserves in-window
(including the cross-file dedup case where an old file's history and a
newer file's removal share one hash), data freshness, and scrub exactness.
2026-07-10 16:43:48 -07:00

4569 lines
175 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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