/** * File System Storage Adapter * File system storage adapter for Node.js environments */ import { HNSWNoun, HNSWNounWithMetadata, StatisticsData, NounType } from '../../coreTypes.js' import { BaseStorage, StorageBatchConfig, SYSTEM_DIR, STATISTICS_KEY, WriterLockInfo, WriterCloseRecord } from '../baseStorage.js' import { getBrainyVersion } from '../../utils/index.js' import { isAbsentError } from '../../utils/errorClassification.js' import { prodLog } from '../../utils/logger.js' import { isZeroNormVector } from '../../utils/distance.js' import { TornRecordError, isUnparseablePayloadError, registerTornRecordEncounter } from '../tornRecordError.js' // Node.js modules - dynamically imported to avoid issues in browser environments let fs: any let path: any let zlib: any let moduleLoadingPromise: Promise | null = null // Try to load Node.js modules try { // Using dynamic imports to avoid issues in browser environments const fsPromise = import('node:fs') const pathPromise = import('node:path') const zlibPromise = import('node:zlib') moduleLoadingPromise = Promise.all([fsPromise, pathPromise, zlibPromise]) .then(([fsModule, pathModule, zlibModule]) => { fs = fsModule path = pathModule.default zlib = zlibModule }) .catch((error) => { console.error('Failed to load Node.js modules:', error) throw error }) } catch (error) { console.error( 'FileSystemStorage: Failed to load Node.js modules. This adapter is not supported in this environment.', error ) } /** * File system storage adapter for Node.js environments * Uses the file system to store data in the specified directory structure * * Type-aware storage now built into BaseStorage * - Removed 10 *_internal method overrides (now inherit from BaseStorage's type-first implementation) * - Removed 2 pagination method overrides (getNounsWithPagination, getVerbsWithPagination) * - Updated HNSW methods to use BaseStorage's getNoun/saveNoun (type-first paths) * - All operations now use type-first paths: entities/nouns/{type}/vectors/{shard}/{id}.json */ export class FileSystemStorage extends BaseStorage { // FileSystem-specific count persistence private countsFilePath?: string // Will be set after init // Fixed sharding configuration for optimal balance of simplicity and performance // Single-level sharding (depth=1) provides excellent performance for 1-2.5M entities // Structure: nouns/ab/uuid.json where 'ab' = first 2 hex chars of UUID // - 256 shard directories (00-ff) // - Handles 2.5M+ entities with < 10K files per shard // - Eliminates dynamic depth changes that cause path mismatch bugs private readonly SHARDING_DEPTH = 1 as const protected rootDir: string private nounsDir!: string private verbsDir!: string private metadataDir!: string private nounMetadataDir!: string private verbMetadataDir!: string private indexDir!: string // Legacy - for backward compatibility private systemDir!: string private lockDir!: string // Root for raw binary blobs (`/_blobs`). Blobs are stored verbatim // (no JSON envelope, no compression) so native code can mmap them directly via // getBinaryBlobPath(). Set in init() once the path module is loaded. private blobsDir!: string private activeLocks: Set = new Set() private lockTimers: Map = new Map() // Track timers for cleanup private allTimers: Set = new Set() // Track all timers for cleanup // Writer-lock state. The writer lock at `locks/_writer.lock` is acquired // at Brainy.init() in writer mode and released at close(). A heartbeat // timer rewrites the lock every 10s so stale-lock detection can tell a dead // writer from a slow one. The constant name matches the file path used. private static readonly WRITER_LOCK_FILE = '_writer.lock' /** * The clean-close record at `locks/_writer.close` (see * {@link WriterCloseRecord}). Written when the lock is released, consumed by * the next claim, so an open can distinguish "the previous writer left" from * "the previous writer died" without inferring either from a pid. */ private static readonly WRITER_CLOSE_FILE = '_writer.close' /** * How often the lock file's `lastHeartbeat` is rewritten. * * THIS IS OBSERVABILITY ONLY, and the cadence follows from that. Staleness * is decided by PID LIVENESS alone (see isWriterLockStale) and the fence * compares pid + hostname — no decision anywhere reads this timestamp. It * exists so an operator inspecting a lock file, or reading the * BRAINY_WRITER_LOCKED error, can judge liveness themselves. * * At 10s it was a lock-file WRITE every ten seconds per brain, forever: 2.1 * writes/s across a production process holding 21 idle brains, for a * human-readable timestamp nothing computes with. At 60s an operator still * sees a heartbeat inside the minute, at a sixth of the cost. With the * clean-close record now recording orderly releases explicitly, the * heartbeat carries even less weight than it did. */ private static readonly WRITER_HEARTBEAT_MS = 60_000 private static readonly WRITER_STALE_THRESHOLD_MS = 60_000 private writerLockHeartbeat?: NodeJS.Timeout private writerLockInfo?: WriterLockInfo /** * The currently-executing heartbeat refresh, if any. `releaseWriterLock()` * awaits it before unlinking: clearInterval() stops FUTURE ticks but not a * tick already in flight, and a straggler landing after the unlink would * RE-CREATE the lock file — a phantom lock blocking the next writer until * the stale TTL expires (the pool-eviction reopen case). */ private writerHeartbeatInFlight?: Promise /** * The in-flight background count-ledger derivation, if one was needed at * open. See {@link scheduleCountLedgerDerivation} — awaited only by * {@link whenCountLedgerSettled}, never by a read. */ private countLedgerDerivation?: Promise // Flush-request RPC state. The writer polls `locks/_flush_requests/` for // new `.req` files and emits `.ack` files in `locks/_flush_responses/` after // flushing. Inspectors call `requestFlushOverFilesystem` to drop a request // and wait for the ack. Polling interval is short enough to feel synchronous // for operator workflows but doesn't pressure the FS. private static readonly FLUSH_REQUEST_DIR = '_flush_requests' private static readonly FLUSH_RESPONSE_DIR = '_flush_responses' private static readonly FLUSH_WATCH_INTERVAL_MS = 500 /** * The safety sweep behind the fs.watch: catches events an exotic filesystem * dropped, and runs the stale-request GC. See startFlushRequestWatcher. */ private static readonly FLUSH_SAFETY_SWEEP_MS = 30_000 private static readonly FLUSH_POLL_INTERVAL_MS = 100 private static readonly FLUSH_REQUEST_TTL_MS = 60_000 private flushWatcherInterval?: NodeJS.Timeout /** The inotify-backed watch on the request directory, when the FS supports one. */ private flushWatcher?: import('node:fs').FSWatcher private flushWatcherInFlight = false private flushWatcherOnRequest?: () => Promise // CRITICAL FIX: Mutex locks for HNSW concurrency control // Prevents read-modify-write races during concurrent neighbor updates at scale (1000+ ops) // Matches MemoryStorage and OPFSStorage behavior (tested in production) private hnswLocks = new Map>() // Compression configuration private compressionEnabled: boolean = true // Enable gzip compression by default for 60-80% disk savings private compressionLevel: number = 6 // zlib compression level (1-9, default: 6 = balanced) // Transaction durability barrier (see GenerationStorage.beginWriteBarrier). // Non-null ONLY between beginWriteBarrier() and flushWriteBarrier(). Because // canonical writes are tmp+rename (durable only in the page cache until an // fsync), the generation store fsyncs everything a transaction wrote before // it advances the generation counter. `writeBarrierPaths` collects the // root-relative object paths written; `writeBarrierDeleteDirs` the parent // dirs of deleted objects (an unlink is durable only once its directory is // fsync'd). Both are null outside a transaction, so the tracking `add`s below // are free on the single-op and non-transactional write paths. private writeBarrierPaths: Set | null = null private writeBarrierDeleteDirs: Set | null = null /** * Initialize the storage adapter * @param rootDirectory The root directory for storage * @param options Optional configuration */ constructor( rootDirectory: string, options?: { compression?: boolean // Enable gzip compression (default: true) compressionLevel?: number // Compression level 1-9 (default: 6) } ) { super() this.rootDir = rootDirectory // Configure compression if (options?.compression !== undefined) { this.compressionEnabled = options.compression } if (options?.compressionLevel !== undefined) { this.compressionLevel = Math.min(9, Math.max(1, options.compressionLevel)) } // Defer path operations until init() when path module is guaranteed to be loaded } /** * Get FileSystem-optimized batch configuration * * File system storage is I/O bound but not rate limited: * - Large batch sizes (500 items) * - No delays needed (0ms) * - Moderate concurrency (100 operations) - limited by I/O threads * - Parallel processing supported * * @returns FileSystem-optimized batch configuration */ public override getBatchConfig(): StorageBatchConfig { return { maxBatchSize: 500, batchDelayMs: 0, maxConcurrent: 100, supportsParallelWrites: true, // Filesystem handles parallel I/O rateLimit: { operationsPerSecond: 5000, // Depends on disk speed burstCapacity: 2000 } } } /** * Initialize the storage adapter */ public override async init(): Promise { if (this.isInitialized) { return } // Wait for module loading to complete if (moduleLoadingPromise) { try { await moduleLoadingPromise } catch (error) { throw new Error( 'FileSystemStorage requires a Node.js environment, but `fs` and `path` modules could not be loaded.' ) } } // Check if Node.js modules are available if (!fs || !path) { throw new Error( 'FileSystemStorage requires a Node.js environment, but `fs` and `path` modules could not be loaded.' ) } try { // Initialize directory paths now that path module is loaded // Clean directory structure this.nounsDir = path.join(this.rootDir, 'entities/nouns/hnsw') this.verbsDir = path.join(this.rootDir, 'entities/verbs/hnsw') this.metadataDir = path.join(this.rootDir, 'entities/nouns/metadata') // Legacy reference this.nounMetadataDir = path.join(this.rootDir, 'entities/nouns/metadata') this.verbMetadataDir = path.join(this.rootDir, 'entities/verbs/metadata') this.indexDir = path.join(this.rootDir, 'indexes') this.systemDir = path.join(this.rootDir, SYSTEM_DIR) this.lockDir = path.join(this.rootDir, 'locks') this.blobsDir = path.join(this.rootDir, '_blobs') // Create the root directory if it doesn't exist await this.ensureDirectoryExists(this.rootDir) // Finish any restore interrupted by a crash (resume the staged swap, or // discard an uncommitted staging area) BEFORE counts/derived state load, // so the rest of startup sees the completed store. ORDER-DEPENDENT: // `swapStagedRestoreIn()` reads `fs.readdir(rootDir)` and then // removes/renames rootDir's own TOP-LEVEL entries to place the staged // copy — racing that against the directory-creation batch below (which // also touches rootDir's children) could see a half-created directory // mid-swap or a mkdir racing a concurrent rm/rename on the same path. // Stays strictly sequential, never folded into the OPEN-PATH batch. await this.completeInterruptedRestore() // OPEN-PATH FIX: the remaining bootstrap directories are mutually // independent — each is its own subtree under rootDir, and // `fs.mkdir(dir, { recursive: true })` creates every intermediate // segment of ITS OWN path in one call, so it never depends on any // sibling here existing first. Nothing between here and // `initializeCounts()` reads any of them, so batching collapses what // was up to 8 sequential mkdir round-trips (each a real syscall+await) // into one wave — this is what serialized an N-writer restart storm on // filesystem I/O it never structurally needed. `initializeCounts()` // right after DOES depend on `systemDir` (which the batch creates), so // it stays outside, awaited only once every directory has landed. await Promise.all([ // Create the nouns directory if it doesn't exist this.ensureDirectoryExists(this.nounsDir), // Create the verbs directory if it doesn't exist this.ensureDirectoryExists(this.verbsDir), // Create the metadata directory if it doesn't exist this.ensureDirectoryExists(this.metadataDir), // Create the noun metadata directory if it doesn't exist this.ensureDirectoryExists(this.nounMetadataDir), // Create the verb metadata directory if it doesn't exist this.ensureDirectoryExists(this.verbMetadataDir), // Create both directories for backward compatibility this.ensureDirectoryExists(this.systemDir), // Only create legacy directory if it exists (don't create new legacy // dirs) — a read-then-maybe-write, but on its own subtree, so it's // still independent of every other entry in this batch. (async () => { if (await this.directoryExists(this.indexDir)) { await this.ensureDirectoryExists(this.indexDir) } })(), // Create the locks directory if it doesn't exist this.ensureDirectoryExists(this.lockDir), // Create the binary blobs directory if it doesn't exist this.ensureDirectoryExists(this.blobsDir) ]) // Initialize count management — depends on systemDir, created above. this.countsFilePath = path.join(this.systemDir, 'counts.json') await this.initializeCounts() // Boot log: new-vs-established, decided from the canonical layout the // database actually reads and writes (`entities/nouns///`) // plus the known noun count. The legacy hnsw sharding-depth probe and // its depth-migration machinery are gone: the 8.0 write path never // populated the directory they inspected, so the probe concluded "new // installation" for every store on every boot and the migration branch // was unreachable. const established = this.totalNounCount > 0 || (await this.hasCanonicalEntities()) console.log( established ? `📁 Using depth ${this.SHARDING_DEPTH} sharding (${this.totalNounCount} entities)` : `📁 New installation: using depth ${this.SHARDING_DEPTH} sharding (optimal for 1-2.5M entities)` ) // Initialize GraphAdjacencyIndex and type statistics await super.init() } catch (error) { console.error('Error initializing FileSystemStorage:', error) throw error } } /** * Check if a directory exists */ private async directoryExists(dirPath: string): Promise { try { const stats = await fs.promises.stat(dirPath) return stats.isDirectory() } catch (error) { return false } } /** * Ensure a directory exists, creating it if necessary */ private async ensureDirectoryExists(dirPath: string): Promise { try { await fs.promises.mkdir(dirPath, { recursive: true }) } catch (error: any) { // Ignore EEXIST error, which means the directory already exists if (error.code !== 'EEXIST') { throw error } } } /** * Primitive operation: Write object to path * All metadata operations use this internally via base class routing * Supports gzip compression for 60-80% disk savings * CRITICAL FIX: Added atomic write pattern to prevent file corruption during concurrent imports */ protected async writeObjectToPath(pathStr: string, data: any): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, pathStr) await this.ensureDirectoryExists(path.dirname(fullPath)) if (this.compressionEnabled) { // Write compressed data with .gz extension using atomic pattern const compressedPath = `${fullPath}.gz` const tempPath = `${compressedPath}.tmp.${Date.now()}.${Math.random().toString(36).substring(2)}` try { // ATOMIC WRITE SEQUENCE: // 1. Compress and write to temp file const jsonString = JSON.stringify(data, null, 2) const compressed = await new Promise((resolve, reject) => { zlib.gzip(Buffer.from(jsonString, 'utf-8'), { level: this.compressionLevel }, (err: any, result: Buffer) => { if (err) reject(err) else resolve(result) }) }) await fs.promises.writeFile(tempPath, compressed) // 2. Atomic rename temp → final (crash-safe, prevents truncation during concurrent writes) await fs.promises.rename(tempPath, compressedPath) } catch (error: any) { // Clean up temp file on any error try { await fs.promises.unlink(tempPath) } catch (cleanupError) { // Ignore cleanup errors } throw error } // Clean up uncompressed file if it exists (migration from uncompressed) try { await fs.promises.unlink(fullPath) } catch (error: any) { // Ignore if file doesn't exist if (error.code !== 'ENOENT') { console.warn(`Failed to remove uncompressed file ${fullPath}:`, error) } } } else { // Write uncompressed data using atomic pattern const tempPath = `${fullPath}.tmp.${Date.now()}.${Math.random().toString(36).substring(2)}` try { // ATOMIC WRITE SEQUENCE: // 1. Write to temp file await fs.promises.writeFile(tempPath, JSON.stringify(data, null, 2)) // 2. Atomic rename temp → final (crash-safe, prevents truncation during concurrent writes) await fs.promises.rename(tempPath, fullPath) } catch (error: any) { // Clean up temp file on any error try { await fs.promises.unlink(tempPath) } catch (cleanupError) { // Ignore cleanup errors } throw error } } // Transaction durability barrier: record the write so the generation store // can fsync it before advancing the counter. Reached only on a successful // rename; the logical (non-`.gz`) path is stored — flushWriteBarrier's // syncRawObjects resolves the compressed variant. this.writeBarrierPaths?.add(pathStr) } /** * Primitive operation: Read object from path * All metadata operations use this internally via base class routing * Supports reading both compressed (.gz) and uncompressed files for backward compatibility * * Read contract (loud errors, never quiet losses): * - Genuine absence (ENOENT on every variant) → `null`. Only a missing file * is "not found". * - TORN record (a file EXISTS but its bytes cannot be decoded — invalid * JSON, truncated/garbled gzip) → the encounter is registered (production * ERROR log + per-process gauge) and a typed {@link TornRecordError} is * thrown. Corruption must NEVER read as absence: callers that can degrade * (manifest recovery, rebuildable statistics) catch the typed error at * their sites; entity reads surface it. * Legacy dual-format exception: when the `.gz` variant is torn but the * uncompressed fallback decodes, the recovered object is returned — AFTER * the torn `.gz` was logged and counted (loud recovery, not a silent skip). * - Real storage fault (EIO/EACCES/EMFILE/…) → propagates as itself; a * fault is neither absence nor corruption and must not be reshaped. */ protected async readObjectFromPath(pathStr: string): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, pathStr) const compressedPath = `${fullPath}.gz` // Try reading compressed file first (if compression is enabled or file exists). // A torn .gz is remembered so the uncompressed fallback can either recover // (legacy dual-format installs) or surface the corruption typed. let tornCompressed: TornRecordError | null = null try { const compressedData = await fs.promises.readFile(compressedPath) const decompressed = await new Promise((resolve, reject) => { zlib.gunzip(compressedData, (err: any, result: Buffer) => { if (err) reject(err) else resolve(result) }) }) return JSON.parse(decompressed.toString('utf-8')) } catch (error: any) { if (error.code === 'ENOENT') { // No compressed variant — fall through to the uncompressed path. } else if (isUnparseablePayloadError(error)) { // The .gz EXISTS but cannot be decoded (zlib Z_* error or JSON // SyntaxError after gunzip): torn record. Register NOW (log + gauge), // then attempt the uncompressed fallback as a recovery read. tornCompressed = registerTornRecordEncounter(`${pathStr}.gz`, error) } else { // Real storage fault on an existing .gz (EIO/EACCES/…): propagate. throw error } } // Fall back to reading uncompressed file (for backward compatibility) try { const data = await fs.promises.readFile(fullPath, 'utf-8') return JSON.parse(data) } catch (error: any) { if (error.code === 'ENOENT') { // No uncompressed file. If the .gz variant existed but was torn, the // object EXISTS and is unreadable — that must surface typed, never as // "absent". Otherwise this is genuine absence. if (tornCompressed !== null) { throw tornCompressed } return null } // The file EXISTS but its content cannot be parsed: torn record. // Register (production ERROR + gauge) and throw typed — a corrupt row // must be distinguishable from a missing row, or nothing ever heals it. if (isUnparseablePayloadError(error)) { throw registerTornRecordEncounter(pathStr, error) } // A real storage fault (EIO/EACCES/EMFILE/…) is NOT "object absent". The // ENOENT branch (above) already returns null; a genuine fault reaching // here must propagate loudly rather than masquerade as a missing object // — which would corrupt reads and drive needless rebuilds. throw error } } /** * Primitive operation: Delete object from path * All metadata operations use this internally via base class routing * Deletes both compressed and uncompressed versions (for cleanup) */ protected async deleteObjectFromPath(pathStr: string): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, pathStr) const compressedPath = `${fullPath}.gz` // Try deleting both compressed and uncompressed files (for cleanup during migration) let deletedCount = 0 // Delete compressed file try { await fs.promises.unlink(compressedPath) deletedCount++ } catch (error: any) { if (error.code !== 'ENOENT') { console.warn(`Error deleting compressed file ${compressedPath}:`, error) } } // Delete uncompressed file try { await fs.promises.unlink(fullPath) deletedCount++ } catch (error: any) { if (error.code !== 'ENOENT') { console.error(`Error deleting uncompressed file ${pathStr}:`, error) throw error } } // If neither file existed, it's not an error (already deleted) if (deletedCount === 0) { // File doesn't exist - this is fine } // Transaction durability barrier: an unlink is durable only once its parent // directory is fsync'd. Record the dir (root-relative; '.' for a top-level // object) so flushWriteBarrier can sync it before the counter advances. this.writeBarrierDeleteDirs?.add(path.dirname(pathStr)) } /** * @description Remove the entity container directory after both canonical legs * were deleted (full-removal delete). `objectLegPath` is one leg's * storage-root-relative path (e.g. `entities/nouns///vectors.json`); * its parent is the `/` entity directory. `rmdir` removes it ONLY if empty * — both legs are gone by this point, so it should be. A non-empty dir means an * unexpected leftover (a leg whose delete faulted — already surfaced upstream — * or foreign data): we do NOT recursively nuke it (loud errors, never quiet * losses); we log and leave it for the orphan repair sweep (`repairIndex`). */ protected override async removeCanonicalContainer(objectLegPath: string): Promise { const relDir = path.dirname(objectLegPath) const absDir = path.join(this.rootDir, relDir) try { await fs.promises.rmdir(absDir) // The dir removal is durable once its PARENT (the shard dir) is fsync'd. this.writeBarrierDeleteDirs?.add(path.dirname(relDir)) } catch (error: any) { if (error?.code === 'ENOENT') return // container already gone — fine if (error?.code === 'ENOTEMPTY' || error?.code === 'EEXIST') { console.warn( `[FileSystemStorage] entity container ${relDir} not empty after delete — ` + `leaving it for the orphan repair sweep (brain.repairIndex()).` ) return } throw error } } /** * @description Prune orphaned entity directories left by the pre-8.3.1 * partial-delete defect. A delete used to remove the metadata (content) leg * but leave the `vectors.json` leg + the `/` directory (a "ghost"), or — * on the direct storage path — remove both legs but leave an empty directory * (a "scar"). Neither is a live entity (`getNoun` needs the metadata content * leg) yet each lingers on disk, inflating enumerated counts and confusing * locator resolution. * * CONSERVATIVE by design: removes ONLY dirs with NO metadata content leg — a * directory that still holds its content is left untouched. LOUD: logs every * removed orphan. Operator-invoked via {@link Brainy.repairIndex}; never runs * automatically. Returns the pruned container ids so the caller can recompute * counts. */ /** * @description Whether an id directory's file legs include the metadata * CONTENT leg (`metadata.json` or its `.json.gz` variant) — the single * test that decides whether an `entities////` container is * a live entity or a ghost/scar orphan left by the pre-8.3.1 partial-delete * defect (see {@link pruneOrphanedEntities}). Shared by the orphan prune * and {@link scanCanonicalEntities} so the two agree by construction — one * counted entity per identity record, never per bare container. * @param legs - File names in one `entities////` directory. */ private hasMetadataContentLeg(legs: string[]): boolean { return legs.some((f) => f.startsWith('metadata.json')) } public async pruneOrphanedEntities(): Promise<{ nouns: string[]; verbs: string[] }> { await this.ensureInitialized() const pruned: { nouns: string[]; verbs: string[] } = { nouns: [], verbs: [] } for (const [kind, root] of [ ['nouns', 'entities/nouns'], ['verbs', 'entities/verbs'] ] as const) { const rootAbs = path.join(this.rootDir, root) let shards: string[] try { shards = await fs.promises.readdir(rootAbs) } catch (error: any) { if (error?.code === 'ENOENT') continue // no entities of this kind yet throw error } for (const shard of shards) { const shardAbs = path.join(rootAbs, shard) let entries: import('fs').Dirent[] try { entries = await fs.promises.readdir(shardAbs, { withFileTypes: true }) } catch (error: any) { if (error?.code === 'ENOENT') continue throw error } for (const entry of entries) { if (!entry.isDirectory()) continue const idAbs = path.join(shardAbs, entry.name) let legs: string[] try { legs = await fs.promises.readdir(idAbs) } catch (error: any) { if (error?.code === 'ENOENT') continue throw error } // A live entity has its metadata content leg. No content leg → a // vector-only ghost or an empty scar → prune the whole container. if (this.hasMetadataContentLeg(legs)) continue await fs.promises.rm(idAbs, { recursive: true, force: true }) pruned[kind].push(entry.name) console.warn( `[FileSystemStorage] pruned orphaned ${kind === 'nouns' ? 'noun' : 'verb'} ` + `container ${root}/${shard}/${entry.name} (no metadata content leg)` ) } } } return pruned } /** * @description The IMMEDIATE child directory names under a prefix — ONE * `readdir`, no recursion, no file paths. See the seam's JSDoc * (`src/db/types.ts`) for what this replaced: discovering the generations on * disk walked the entire generation log on every open, reading out every * file in every generation, to learn the set of integers the top-level * directory names already spell. * @param prefix - Storage-root-relative directory prefix. * @returns The child directory names (not paths); empty when the prefix does * not exist. */ public override async listRawPrefixes(prefix: string): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, prefix) try { const entries = await fs.promises.readdir(fullPath, { withFileTypes: true }) return entries.filter((e: { isDirectory: () => boolean }) => e.isDirectory()) .map((e: { name: string }) => e.name) } catch (error: any) { if (error?.code === 'ENOENT') return [] throw error } } /** * Primitive operation: List objects under path prefix * All metadata operations use this internally via base class routing * Handles both .json and .json.gz files, normalizes paths */ protected async listObjectsUnderPath(prefix: string): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, prefix) const paths: string[] = [] const seen = new Set() // Track files to avoid duplicates (both .json and .json.gz) try { const entries = await fs.promises.readdir(fullPath, { withFileTypes: true }) for (const entry of entries) { if (entry.isFile()) { // Handle multiple compression formats for broad compatibility // - .json.gz: Standard entity/metadata files (JSON compressed) // - .gz: Raw compressed payloads (e.g. blob-store binary values) // - .json: Uncompressed JSON files if (entry.name.endsWith('.json.gz')) { // Strip .gz extension and add the .json path const normalizedName = entry.name.slice(0, -3) // Remove .gz const normalizedPath = path.join(prefix, normalizedName) if (!seen.has(normalizedPath)) { paths.push(normalizedPath) seen.add(normalizedPath) } } else if (entry.name.endsWith('.gz')) { // Raw payloads stored as .gz (not .json.gz) // Strip .gz extension and return path const normalizedName = entry.name.slice(0, -3) // Remove .gz const normalizedPath = path.join(prefix, normalizedName) if (!seen.has(normalizedPath)) { paths.push(normalizedPath) seen.add(normalizedPath) } } else if (entry.name.endsWith('.json')) { const filePath = path.join(prefix, entry.name) if (!seen.has(filePath)) { paths.push(filePath) seen.add(filePath) } } } else if (entry.isDirectory()) { const subpath = path.join(prefix, entry.name) const subdirPaths = await this.listObjectsUnderPath(subpath) paths.push(...subdirPaths) } } return paths.sort() } catch (error: any) { if (error.code === 'ENOENT') { return [] } throw error } } // =========================================================================== // Generational record layer (8.0 MVCC) — durability + snapshot primitives // =========================================================================== /** * Storage-root-relative paths that are mutated **in place** (appended to) * rather than replaced via atomic tmp+rename. `snapshotToDirectory()` must * byte-copy these instead of hard-linking them: a hard link shares the * inode, so a post-snapshot append to the live file would silently mutate * the snapshot. Every other persisted file in this adapter is written via * tmp+rename (objects, blobs, counts, locks are excluded entirely), which * makes hard links safe — a rewrite swaps in a new inode and the snapshot * keeps the old one. */ private static readonly SNAPSHOT_BYTE_COPY_PATHS = new Set([ `${SYSTEM_DIR}/tx-log.jsonl` ]) /** * Top-level directories whose EVERY file is mutated in place (not tmp+rename) * and must therefore be byte-copied into a snapshot, not hard-linked — the * directory analogue of {@link SNAPSHOT_BYTE_COPY_PATHS}. * * `_id_mapper` holds the native provider's shared mmap `BinaryIdMapper`, which * `flush()` msyncs and a migration rebuild can truncate+re-inject IN PLACE * (confirmed by the native side). A hard link shares the inode, so an in-place * msync/truncate on the live file would reach through into the pre-upgrade * backup — byte-copy keeps the snapshot a faithful, independent copy. It is * bounded (the id map), not the large index files, so the copy cost is small. */ private static readonly SNAPSHOT_BYTE_COPY_DIRS = new Set(['_id_mapper']) /** * Nested path PREFIXES whose files are byte-copied into snapshots, not * hard-linked — for append-in-place files below the top level. The * generation fact log's tail segment is appended in place between rotations; * a hard-linked tail would let post-snapshot appends reach through into the * snapshot. (Sealed segments are immutable and would be link-safe, but the * prefix rule keeps the discipline simple; segments are bounded by the * rotation threshold, so the copy cost is small.) */ private static readonly SNAPSHOT_BYTE_COPY_PREFIXES: string[] = ['_generations/facts/'] /** * Top-level directories excluded from snapshots: process-local lock state * (writer lock, flush-request RPC files) must never travel with the data, and * the restore staging area ({@link RESTORE_STAGING_DIR}) is transient scratch * that must never be captured or restored. */ private static readonly SNAPSHOT_EXCLUDED_TOP_DIRS = new Set(['locks', '_restore_staging']) /** * Transient top-level directory holding a restore-in-progress: the snapshot is * fully copied here (sparse-aware) BEFORE any live data is touched, then an * atomic per-entry swap moves it into place. Its presence + the completion * marker let {@link completeInterruptedRestore} resume a crashed restore. */ private static readonly RESTORE_STAGING_DIR = '_restore_staging' /** * Written+fsync'd inside the staging dir ONLY after the whole snapshot has * copied successfully. Its presence authorizes the swap (and its resume): a * staging dir WITHOUT this marker is an interrupted copy — discardable debris, * live data still authoritative. */ private static readonly RESTORE_MARKER = '.restore-manifest.json' /** Chunk size for sparse-aware copying (holes are preserved at this grain). */ private static readonly SPARSE_CHUNK_BYTES = 4 * 1024 * 1024 /** A zero buffer the size of one sparse chunk, for all-zero (hole) detection. */ private static readonly SPARSE_ZERO_CHUNK = Buffer.alloc(4 * 1024 * 1024) /** * Remove every object under a storage-root-relative prefix — one recursive * directory removal instead of the base class's list+delete loop. * * @param prefix - Storage-root-relative directory prefix to remove. */ public override async removeRawPrefix(prefix: string): Promise { await this.ensureInitialized() // Registered-blob contract: a prefix-nuke must not take out a protected // family member. Throws if the prefix intersects one. await this.assertPrefixNotProtected(prefix) await fs.promises.rm(path.join(this.rootDir, prefix), { recursive: true, force: true }) } /** * Durability barrier: `fsync` each listed object file (resolving the * compressed `.gz` variant when present) and then the set of parent * directories, so both the file contents and the rename directory entries * are durable before the commit protocol proceeds. Paths whose file no * longer exists are skipped (a later step may have replaced them). * * @param paths - Storage-root-relative object paths previously written. */ public override async syncRawObjects(paths: string[]): Promise { await this.ensureInitialized() const parentDirs = new Set() for (const objectPath of paths) { const fullPath = path.join(this.rootDir, objectPath) let synced = false for (const candidate of [`${fullPath}.gz`, fullPath]) { let handle: any try { handle = await fs.promises.open(candidate, 'r') } catch (error: any) { if (error.code === 'ENOENT') continue throw error } try { await handle.sync() } finally { await handle.close() } parentDirs.add(path.dirname(fullPath)) synced = true break } // An absent path is a state too: fsync the parent directory so a // completed unlink is durable (a delete must survive power loss as // surely as a write — otherwise a bounded log fold could let a // tombstoned record resurrect from a lost directory update). if (!synced) parentDirs.add(path.dirname(fullPath)) } for (const dir of parentDirs) { let handle: any try { handle = await fs.promises.open(dir, 'r') } catch { continue // directory vanished or platform disallows opening dirs } try { await handle.sync() } catch { // Some platforms (and some filesystems) reject directory fsync — // file-level fsync above already covers the data itself. } finally { await handle.close() } } } /** * Begin a transaction durability barrier: start recording every canonical * object write and delete so {@link flushWriteBarrier} can fsync them before * the generation counter advances. Resets unconditionally, discarding any * tracking left by a transaction that aborted without flushing. * * @see GenerationStorage.beginWriteBarrier */ public beginWriteBarrier(): void { this.writeBarrierPaths = new Set() this.writeBarrierDeleteDirs = new Set() } /** * Flush the transaction durability barrier: fsync every canonical write since * {@link beginWriteBarrier} (file contents AND the rename directory entries, * via {@link syncRawObjects}), then fsync the parent directory of every * canonical delete so the unlinks are durable too. Clears the tracking. After * this resolves, the transaction's entire canonical footprint is on disk, so * the generation counter/manifest can be advanced without risking a * counter-ahead-of-state torn store on a hard kill. * * @see GenerationStorage.flushWriteBarrier */ public async flushWriteBarrier(): Promise { const paths = this.writeBarrierPaths const deleteDirs = this.writeBarrierDeleteDirs this.writeBarrierPaths = null this.writeBarrierDeleteDirs = null if (paths && paths.size > 0) { // syncRawObjects fsyncs each file and its parent directory. await this.syncRawObjects([...paths]) } if (deleteDirs && deleteDirs.size > 0) { for (const relDir of deleteDirs) { const dirPath = path.join(this.rootDir, relDir) let handle: any try { handle = await fs.promises.open(dirPath, 'r') } catch { continue // directory vanished or platform disallows opening dirs } try { await handle.sync() } catch { // Some platforms/filesystems reject directory fsync — best effort. } finally { await handle.close() } } } } /** * Append one line to `_system/tx-log.jsonl`. Plain `appendFile` — the * tx-log is the one append-in-place file in the store (and is byte-copied, * never hard-linked, by {@link FileSystemStorage.snapshotToDirectory}). * * @param line - One complete JSON document, without trailing newline. */ public async appendTxLogLine(line: string): Promise { await this.ensureInitialized() const logPath = path.join(this.systemDir, 'tx-log.jsonl') await fs.promises.appendFile(logPath, `${line}\n`, 'utf-8') } /** * Read all tx-log lines, oldest first (empty array when no log exists). * Torn trailing lines from a crashed append are returned as-is — callers * tolerate unparseable lines. */ public async readTxLogLines(): Promise { await this.ensureInitialized() const logPath = path.join(this.systemDir, 'tx-log.jsonl') try { const content: string = await fs.promises.readFile(logPath, 'utf-8') return content.split('\n').filter((l: string) => l.length > 0) } catch (error: any) { if (error.code === 'ENOENT') return [] throw error } } // ========================================================================== // Binary raw-byte primitives — the substrate for append-only log-structured // files (the generation fact log's CRC-framed segments). Paths are used // VERBATIM (no .gz/.bin suffixing). Append durability rides syncRawObjects // at the commit barrier, like every other staged write. // ========================================================================== /** * Append bytes to a raw binary file, creating it (and parent directories) * when absent. NOT fsync'd here — the caller batches durability via * `syncRawObjects` at its commit barrier. */ public async appendRawBytes(rawPath: string, bytes: Uint8Array): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, rawPath) await fs.promises.mkdir(path.dirname(fullPath), { recursive: true }) await fs.promises.appendFile(fullPath, bytes) } /** * Read a raw binary file whole. Absent → `null`; a real IO fault throws — * a present-but-unreadable log segment must never read as "no facts". */ public async readRawBytes(rawPath: string): Promise { await this.ensureInitialized() try { const buf: Buffer = await fs.promises.readFile(path.join(this.rootDir, rawPath)) return new Uint8Array(buf.buffer, buf.byteOffset, buf.byteLength) } catch (error: any) { if (isAbsentError(error)) return null throw error } } /** * Replace a raw binary file atomically: write-new → fsync → rename. The * reconcile primitive (e.g. truncating a fact-log tail back to committed * truth after a crash) — a crash mid-replace leaves either the old file or * the new one, never a mix. */ public async writeRawBytes(rawPath: string, bytes: Uint8Array): Promise { await this.ensureInitialized() const fullPath = path.join(this.rootDir, rawPath) await fs.promises.mkdir(path.dirname(fullPath), { recursive: true }) const tmpPath = `${fullPath}.tmp.${Date.now()}.${Math.random().toString(36).slice(2)}` const handle = await fs.promises.open(tmpPath, 'w') try { await handle.writeFile(bytes) await handle.sync() } finally { await handle.close() } await fs.promises.rename(tmpPath, fullPath) } /** * Byte size of a raw binary file, or `null` when absent. */ public async rawByteSize(rawPath: string): Promise { await this.ensureInitialized() try { const stat = await fs.promises.stat(path.join(this.rootDir, rawPath)) return stat.size } catch (error: any) { if (isAbsentError(error)) return null throw error } } /** * Snapshot the entire store into `targetPath` as a hard-link farm * (Cassandra-style: instant, space-shared). Safe because every data file * is immutable-by-rename — rewrites swap in new inodes, leaving the * snapshot's links pointing at the old bytes. The two exceptions are * handled explicitly: append-in-place files * ({@link FileSystemStorage.SNAPSHOT_BYTE_COPY_PATHS}) are byte-copied, and * process-local lock state * ({@link FileSystemStorage.SNAPSHOT_EXCLUDED_TOP_DIRS}) is excluded. * Cross-device targets (where `link(2)` fails with `EXDEV`) fall back to * byte copies per file. * * @param targetPath - Absolute directory for the snapshot. Created if * missing; must be empty or absent (refuses to overwrite). */ public async snapshotToDirectory(targetPath: string): Promise { await this.ensureInitialized() try { const existing = await fs.promises.readdir(targetPath) if (existing.length > 0) { throw new Error( `snapshotToDirectory: target ${targetPath} already exists and is not empty. ` + `Choose a fresh directory per snapshot.` ) } } catch (error: any) { if (error.code !== 'ENOENT') throw error } await fs.promises.mkdir(targetPath, { recursive: true }) const files: string[] = [] await this.collectSnapshotFiles(this.rootDir, '', files) for (const relPath of files) { const sourceFile = path.join(this.rootDir, relPath) const targetFile = path.join(targetPath, relPath) await fs.promises.mkdir(path.dirname(targetFile), { recursive: true }) // Byte-copy list: compare against the normalized (extension-preserving) // relative path with separators unified, so `_system/tx-log.jsonl` // matches on every platform. const normalized = relPath.split(path.sep).join('/') // Byte-copy (never hard-link) files that are mutated in place: the exact // append-in-place paths, and every file under a mmap-mutated directory // (e.g. the native `_id_mapper/*`). A shared inode would let a live // msync/truncate reach through into the snapshot. if ( FileSystemStorage.SNAPSHOT_BYTE_COPY_PATHS.has(normalized) || FileSystemStorage.SNAPSHOT_BYTE_COPY_DIRS.has(normalized.split('/')[0]) || FileSystemStorage.SNAPSHOT_BYTE_COPY_PREFIXES.some((p) => normalized.startsWith(p)) ) { await fs.promises.copyFile(sourceFile, targetFile) continue } try { await fs.promises.link(sourceFile, targetFile) } catch (error: any) { // EXDEV: cross-device; EPERM/ENOTSUP: filesystem forbids links. // ENOENT: the live file was atomically replaced mid-walk — retry as // a copy of whatever is current (single-writer discipline means this // only happens for derived files being flushed concurrently). if (['EXDEV', 'EPERM', 'ENOTSUP', 'ENOENT'].includes(error.code)) { try { await fs.promises.copyFile(sourceFile, targetFile) } catch (copyError: any) { if (copyError.code !== 'ENOENT') throw copyError } } else { throw error } } } } /** * @description Pre-upgrade backup: a hard-link snapshot of the whole store into * a SIBLING directory (`.migration-backup`, outside `rootDir` so the * snapshot never recurses into itself), taken before a 7.x → 8.0 upgrade * rebuilds the derived indexes. Zero-copy (shared inodes; the store is * immutable-by-rename) and instant even at scale. Returns the backup path, or * `null` when the store holds no canonical nouns (nothing to protect). If a * backup already exists from a prior FAILED upgrade, it is REUSED as-is (that * is the true pre-upgrade state; a retry must not overwrite it). */ public async createMigrationBackup(): Promise { await this.ensureInitialized() // Nothing to protect if the store has no canonical nouns (e.g. a brand-new // brain whose marker is simply absent — `epochStale` is true but there is no // 7.x data to migrate). const probe = await this.getNouns({ pagination: { limit: 1 } }) if ((probe.totalCount || 0) === 0 && probe.items.length === 0) { return null } const backupPath = `${this.rootDir}.migration-backup` // Reuse an existing backup (a prior failed upgrade's pre-state) rather than // overwrite it — snapshotToDirectory also refuses a non-empty target. try { const existing = await fs.promises.readdir(backupPath) if (existing.length > 0) return backupPath } catch (error: any) { if (error.code !== 'ENOENT') throw error } await this.snapshotToDirectory(backupPath) return backupPath } /** * @description Remove a {@link createMigrationBackup} snapshot. Best-effort: a * missing path is a no-op, and any error is swallowed (a leftover backup dir is * harmless — the operator can delete it). Removing the hard-links never touches * the live store's bytes (shared inodes; only the extra links go away). */ public async removeMigrationBackup(location: string): Promise { try { await fs.promises.rm(location, { recursive: true, force: true }) } catch { // best-effort — leaving the backup behind is safe } } /** * Recursively collect snapshot-eligible files under `dirAbs`, excluding the * lock directory and in-flight `*.tmp.*` write files. */ private async collectSnapshotFiles(dirAbs: string, relPrefix: string, out: string[]): Promise { let entries: any[] try { entries = await fs.promises.readdir(dirAbs, { withFileTypes: true }) } catch (error: any) { if (error.code === 'ENOENT') return throw error } for (const entry of entries) { const rel = relPrefix ? path.join(relPrefix, entry.name) : entry.name if (entry.isDirectory()) { if (relPrefix === '' && FileSystemStorage.SNAPSHOT_EXCLUDED_TOP_DIRS.has(entry.name)) { continue } await this.collectSnapshotFiles(path.join(dirAbs, entry.name), rel, out) } else if (entry.isFile()) { if (entry.name.includes('.tmp.')) continue // in-flight atomic write out.push(rel) } } } /** * Replace the store's contents from a snapshot directory: every current * top-level entry except `locks/` (the live writer lock must survive) is * removed, the snapshot is byte-copied in (`fs.cp` — never hard-linked, so * the snapshot stays independent of the restored store), and all * adapter-internal derived state is reloaded. * * @param sourcePath - Absolute path of a directory produced by * {@link FileSystemStorage.snapshotToDirectory}. * * Non-destructive: the snapshot is copied into a staging area (sparse-aware, * so a store of mostly-hole mmap blobs cannot balloon and ENOSPC) BEFORE any * live data is touched. Only once the full copy has succeeded and a completion * marker is fsync'd does an atomic per-entry swap move it into place. A copy * failure (including ENOSPC) leaves the live store exactly as it was; a crash * mid-swap is resumed forward on the next {@link init} by * {@link completeInterruptedRestore}. The previous implementation removed the * live store first and then `fs.cp`'d — a copy failure destroyed the brain it * was meant to recover. */ public async restoreFromDirectory(sourcePath: string): Promise { await this.ensureInitialized() const sourceStat = await fs.promises.stat(sourcePath).catch(() => null) if (!sourceStat || !sourceStat.isDirectory()) { throw new Error(`restoreFromDirectory: ${sourcePath} is not a directory`) } const staging = path.join(this.rootDir, FileSystemStorage.RESTORE_STAGING_DIR) // Discard any staging left by a prior aborted restore, then stage fresh. await fs.promises.rm(staging, { recursive: true, force: true }) await fs.promises.mkdir(staging, { recursive: true }) // Copy the snapshot into staging (sparse-aware). ANY failure here — most // importantly ENOSPC — leaves the live store untouched: we remove only the // half-written staging area and re-throw. const staged: string[] = [] try { const sourceEntries = await fs.promises.readdir(sourcePath) for (const entry of sourceEntries) { if (FileSystemStorage.SNAPSHOT_EXCLUDED_TOP_DIRS.has(entry)) continue await this.copyTreeSparse( path.join(sourcePath, entry), path.join(staging, entry) ) staged.push(entry) } // Commit point: fsync a marker naming the fully-staged entries. Only after // this does the swap (and its resume) become authorized. const markerPath = path.join(staging, FileSystemStorage.RESTORE_MARKER) await fs.promises.writeFile(markerPath, JSON.stringify({ entries: staged })) const mfh = await fs.promises.open(markerPath, 'r') try { await mfh.sync() } finally { await mfh.close() } const dfh = await fs.promises.open(staging, 'r').catch(() => null) if (dfh) { try { await dfh.sync() } catch { // platform may reject directory fsync — best effort } finally { await dfh.close() } } } catch (error: any) { await fs.promises.rm(staging, { recursive: true, force: true }).catch(() => {}) throw new Error( `restoreFromDirectory: staging copy failed, live store left untouched: ${error.message}` ) } // Atomic per-entry swap (metadata-only renames within rootDir — cannot ENOSPC). await this.swapStagedRestoreIn() await this.reloadDerivedState() } /** * Move a fully-staged, marker-committed restore into place, then clear the * staging area. Idempotent and resumable: driven by the marker's entry list * and by which staged entries remain, so a crash at any point is completed by * simply calling it again (from {@link completeInterruptedRestore} on the next * open). Every step is a same-filesystem rename or a remove — no operation can * fail for disk space, so once the marker exists the store is guaranteed to * reach the restored state. */ private async swapStagedRestoreIn(): Promise { const staging = path.join(this.rootDir, FileSystemStorage.RESTORE_STAGING_DIR) const markerPath = path.join(staging, FileSystemStorage.RESTORE_MARKER) let staged: string[] try { const marker = JSON.parse(await fs.promises.readFile(markerPath, 'utf-8')) staged = Array.isArray(marker?.entries) ? marker.entries : [] } catch { return // no committed marker — nothing to swap } const stagedSet = new Set(staged) // 1. Remove stale live entries the snapshot does not contain (excluding // process-local dirs and the staging area itself). An already-placed // staged entry is in stagedSet, so it is kept. for (const entry of await fs.promises.readdir(this.rootDir)) { if (FileSystemStorage.SNAPSHOT_EXCLUDED_TOP_DIRS.has(entry)) continue if (entry === FileSystemStorage.RESTORE_STAGING_DIR) continue if (stagedSet.has(entry)) continue await fs.promises.rm(path.join(this.rootDir, entry), { recursive: true, force: true }) } // 2. Place each staged entry (idempotent: a prior attempt that already moved // it leaves staging/entry absent, so we skip). `rm` the old first — a // rename onto an existing non-empty directory is not allowed; the staged // copy is the durable source until it is placed, so a crash between the // rm and the rename is recovered forward on the next call. for (const entry of staged) { const from = path.join(staging, entry) const to = path.join(this.rootDir, entry) const exists = await fs.promises.lstat(from).then(() => true, () => false) if (!exists) continue await fs.promises.rm(to, { recursive: true, force: true }) await fs.promises.rename(from, to) } // 3. Clear the staging area (marker last-standing entry). await fs.promises.rm(staging, { recursive: true, force: true }) } /** * On open, finish any restore interrupted by a crash. A staging dir WITH the * completion marker means the copy had succeeded — resume the swap forward * (loudly). A staging dir WITHOUT the marker is an interrupted copy — pure * debris; the live store is authoritative, so discard it. Called from * {@link init} before counts/derived state load, so recovery is invisible to * the rest of startup. Returns `true` if a swap was resumed. */ private async completeInterruptedRestore(): Promise { const staging = path.join(this.rootDir, FileSystemStorage.RESTORE_STAGING_DIR) const stagingStat = await fs.promises.stat(staging).catch(() => null) if (!stagingStat || !stagingStat.isDirectory()) return false const markerPath = path.join(staging, FileSystemStorage.RESTORE_MARKER) const hasMarker = await fs.promises .stat(markerPath) .then(() => true, () => false) if (!hasMarker) { // Interrupted before the copy committed — live data untouched, discard. await fs.promises.rm(staging, { recursive: true, force: true }).catch(() => {}) return false } console.log('♻️ Resuming an interrupted restore (completing the staged swap)') await this.swapStagedRestoreIn() return true } /** * Recursively copy `src` to `dest`, preserving holes (sparse regions). Files * are copied chunk-by-chunk skipping all-zero chunks, so a store of * mostly-hole mmap blobs restores at its true allocated size instead of * materializing every hole (the failure that made `fs.cp` ENOSPC a restore * that would otherwise fit). Directories recurse; symlinks are recreated. */ private async copyTreeSparse(src: string, dest: string): Promise { const stat = await fs.promises.lstat(src) if (stat.isDirectory()) { await fs.promises.mkdir(dest, { recursive: true }) for (const child of await fs.promises.readdir(src)) { await this.copyTreeSparse(path.join(src, child), path.join(dest, child)) } } else if (stat.isSymbolicLink()) { await fs.promises.symlink(await fs.promises.readlink(src), dest) } else if (stat.isFile()) { await this.copyFileSparse(src, dest, stat.size, stat.mode) } // Other node types (sockets, devices) do not occur in a brain store. } /** Sparse-aware single-file copy — see {@link copyTreeSparse}. */ private async copyFileSparse( src: string, dest: string, size: number, mode: number ): Promise { const CHUNK = FileSystemStorage.SPARSE_CHUNK_BYTES const srcFh = await fs.promises.open(src, 'r') try { const destFh = await fs.promises.open(dest, 'w', mode) try { // Pre-size the destination so unwritten regions are holes. await destFh.truncate(size) const buf = Buffer.allocUnsafe(CHUNK) let pos = 0 while (pos < size) { const { bytesRead } = await srcFh.read(buf, 0, CHUNK, pos) if (bytesRead === 0) break const chunk = buf.subarray(0, bytesRead) // Skip all-zero chunks: leaving them unwritten preserves the hole. const isHole = chunk.equals( FileSystemStorage.SPARSE_ZERO_CHUNK.subarray(0, bytesRead) ) if (!isHole) { await destFh.write(chunk, 0, bytesRead, pos) } pos += bytesRead } } finally { await destFh.close() } } finally { await srcFh.close() } } // =========================================================================== // Raw binary-blob primitive (mmap-friendly) // =========================================================================== /** * Resolve a blob key to its on-disk path under `/_blobs`. * * The key's "/"-separated segments become nested directories and the file is * suffixed with `.bin`, e.g. `"graph-lsm/source/sstable-123"` → * `/_blobs/graph-lsm/source/sstable-123.bin`. This convention is * shared with cor's `MmapFileSystemStorage` so native code and the TS layer * agree on exactly where each blob lives. * * @param key - The blob key. * @returns The absolute on-disk path for the blob. * @private */ private blobPath(key: string): string { // `path` and `blobsDir` are populated in init(). If a caller resolves a blob // path before init() (e.g. native code probing for a mmap target), fall back // to a POSIX-style join off rootDir so this synchronous method never throws. if (path && this.blobsDir) { return path.join(this.blobsDir, ...key.split('/')) + '.bin' } return `${this.rootDir}/_blobs/${key}.bin` } /** * Persist a raw binary blob verbatim under `key`, using an atomic * temp-file + rename so concurrent readers never observe a torn write. * Parent directories are created on demand. * * @param key - The blob key (see {@link getBinaryBlobPath} for the convention). * @param data - The exact bytes to store. */ public async saveBinaryBlob(key: string, data: Buffer): Promise { await this.ensureInitialized() const filePath = this.blobPath(key) await this.ensureDirectoryExists(path.dirname(filePath)) // Atomic write via a UNIQUE per-writer temp suffix — matches the pattern // used at every other atomic-write site in this file (lines 336, 551, 744, // 781, 1529, 2908). The unique suffix (pid+time+random) means NO other // writer ever touches this temp, so two concurrent saveBinaryBlob() calls // for the same key can never collide on the temp path (the shared-`.tmp` // race that once threw ENOENT — column-store compaction vs an explicit // flush() on `_column_index//DELETED.bin` — is gone). // // Crucially, BECAUSE the temp is unique, a rename ENOENT can no longer mean // "a concurrent idempotent writer already renamed it": nobody else has this // temp. It means OUR just-written temp vanished before the rename, so the // bytes did NOT land — returning success would acknowledge a write that // stored nothing (the native provider mmaps these blobs; a phantom-acked // blob is exactly the silent-loss class). The only way that happens with a // unique temp is an external sweeper / crash-cleanup removing it mid-write, // so retry ONCE with a fresh temp; if it vanishes again, FAIL LOUD. const writeOnce = async (): Promise<'ok' | 'temp-vanished'> => { const tmpPath = `${filePath}.tmp.${process.pid}.${Date.now()}.${Math.random().toString(36).slice(2)}` await fs.promises.writeFile(tmpPath, data) try { await fs.promises.rename(tmpPath, filePath) return 'ok' } catch (err) { // Clean up our own temp (best-effort) so no failure path orphans it. await fs.promises.unlink(tmpPath).catch(() => {}) if ((err as NodeJS.ErrnoException).code === 'ENOENT') return 'temp-vanished' throw err } } if ((await writeOnce()) === 'temp-vanished' && (await writeOnce()) === 'temp-vanished') { throw new Error( `saveBinaryBlob('${key}'): the temp file was removed before rename on two ` + `successive attempts — the blob did NOT persist. Some external process is ` + `deleting files under ${this.blobsDir} mid-write. Failing loud rather than ` + `acknowledging a durable write that stored nothing.` ) } } /** * Load the raw bytes stored under `key`, or `null` if the blob does not exist. * * @param key - The blob key. * @returns The blob bytes, or `null` if absent. */ public async loadBinaryBlob(key: string): Promise { await this.ensureInitialized() try { return await fs.promises.readFile(this.blobPath(key)) } catch (err) { // Absent blob → null (the documented contract). A real fault // (EIO/EACCES/EMFILE/…) must NOT be masked as "absent": doing so makes a // present-but-unreadable blob look missing and drives a needless rebuild // or an empty read (the native provider consumes this). Mandate: loud // errors, never quiet losses. Fault-propagation restored lockstep with // cortex 3.0.13, whose two column-store call sites now handle the throw // (they mark the field unavailable + throw a named error) instead of // relying on null-on-error. if (isAbsentError(err)) return null throw err } } /** * Delete the blob stored under `key`. Missing blobs are ignored. * * @param key - The blob key. */ public async deleteBinaryBlob(key: string): Promise { await this.ensureInitialized() // Registered-blob contract: refuse to delete a declared // derived-index family member — an in-process GC/sweeper cannot remove a // load-bearing index file. Throws ProtectedArtifactError; no-op when no // families are registered. await this.assertBlobKeyDeletable(key) try { await fs.promises.unlink(this.blobPath(key)) } catch { /* ignore missing files */ } } /** * Return the real on-disk path for `key` so native code can mmap the file * directly. The path is returned whether or not the file currently exists — * callers are expected to write before mapping. * * @param key - The blob key. * @returns The absolute on-disk path (never `null` for filesystem storage). */ public getBinaryBlobPath(key: string): string | null { return this.blobPath(key) } /** * Get multiple metadata objects in batches (CRITICAL: Prevents socket exhaustion) * FileSystem implementation uses controlled concurrency to prevent too many file reads */ public async getMetadataBatch(ids: string[]): Promise> { await this.ensureInitialized() const results = new Map() const batchSize = 10 // Process 10 files at a time // Process in batches to avoid overwhelming the filesystem for (let i = 0; i < ids.length; i += batchSize) { const batch = ids.slice(i, i + batchSize) const batchPromises = batch.map(async (id) => { try { // CRITICAL: Use getNounMetadata() instead of deprecated getMetadata() // This ensures we fetch from the correct noun metadata store (2-file system) const metadata = await this.getNounMetadata(id) return { id, metadata } } catch (error) { console.debug(`Failed to read metadata for ${id}:`, error) return { id, metadata: null } } }) const batchResults = await Promise.all(batchPromises) for (const { id, metadata } of batchResults) { if (metadata !== null) { results.set(id, metadata) } } // Small yield between batches await new Promise(resolve => setImmediate(resolve)) } return results } /** * Get nouns with pagination support * @param options Pagination options */ // Removed getNounsWithPagination override - now uses BaseStorage's type-first implementation /** * Clear all data from storage */ public async clear(): Promise { await this.ensureInitialized() // Check if fs module is available if (!fs || !fs.promises) { console.warn('FileSystemStorage.clear: fs module not available, skipping clear operation') return } // Helper function to remove all files in a directory const removeDirectoryContents = async (dirPath: string): Promise => { try { const files = await fs.promises.readdir(dirPath) for (const file of files) { const filePath = path.join(dirPath, file) const stats = await fs.promises.stat(filePath) if (stats.isDirectory()) { await removeDirectoryContents(filePath) await fs.promises.rmdir(filePath) } else { await fs.promises.unlink(filePath) } } } catch (error: any) { if (error.code !== 'ENOENT') { console.error(`Error removing directory contents ${dirPath}:`, error) throw error } } } // Clear the canonical entity/verb data area const entitiesDir = path.join(this.rootDir, 'entities') if (await this.directoryExists(entitiesDir)) { await removeDirectoryContents(entitiesDir) } // Remove all files in both system directories await removeDirectoryContents(this.systemDir) if (await this.directoryExists(this.indexDir)) { await removeDirectoryContents(this.indexDir) } // Remove the content-addressed blob store (VFS file content) const casDir = path.join(this.rootDir, '_cas') if (await this.directoryExists(casDir)) { // Delete the entire _cas/ directory (not just contents) await fs.promises.rm(casDir, { recursive: true, force: true }) } // Remove the raw-blob + native-shared + column-index footprint. These // top-level trees were NOT wiped before, so a cleared brain re-read stale // native blobs (HNSW/LSM segments, the native dkann index), a stale native // id-mapper, and — worst — orphaned column manifests. `_column_index` holds // the column-store MANIFEST.json files while their segment bytes live under // `_blobs/_column_index/...`; removing one without the other would strand a // manifest listing segments that no longer exist, which the load path now // (correctly) refuses with ColumnSegmentLoadError. They must fall together, // as a set, exactly as `_cas` does — the complete derived footprint, not a // subset. (`locks/` is deliberately left: it is live coordination state, // not data.) for (const nativeDir of ['_blobs', '_id_mapper', '_column_index']) { const dir = path.join(this.rootDir, nativeDir) if (await this.directoryExists(dir)) { await fs.promises.rm(dir, { recursive: true, force: true }) } } // Reset ALL shared derived state via the same path restore uses: the // write-through cache, id→type/subtype caches, per-type/subtype count // rollups, statistics cache, graph-index singleton, and the BlobStorage // instance (its LRU cache holds deleted blobs), then recount from the // now-empty data area. Without the rollup reset, stats() would keep // reporting per-type counts for deleted entities. await this.reloadDerivedState() } /** * Enhanced clear operation with safety mechanisms and performance optimizations * Provides progress tracking, backup options, and instance name confirmation */ public async clearEnhanced(options: import('../enhancedClearOperations.js').ClearOptions = {}): Promise { await this.ensureInitialized() // Check if fs module is available if (!fs || !fs.promises) { throw new Error('FileSystemStorage.clearEnhanced: fs module not available') } const { EnhancedFileSystemClear } = await import('../enhancedClearOperations.js') const enhancedClear = new EnhancedFileSystemClear(this.rootDir, fs, path) const result = await enhancedClear.clear(options) if (result.success) { // Clear the statistics cache this.statisticsCache = null this.statisticsModified = false } return result } /** * Get information about storage usage and capacity */ public async getStorageStatus(): Promise<{ type: string used: number quota: number | null details?: Record }> { await this.ensureInitialized() // Check if fs module is available if (!fs || !fs.promises) { console.warn('FileSystemStorage.getStorageStatus: fs module not available, returning default values') return { type: 'filesystem', used: 0, quota: null, details: { nounsCount: 0, verbsCount: 0, metadataCount: 0, directorySizes: { nouns: 0, verbs: 0, metadata: 0, index: 0 } } } } try { // Calculate the total size of all files in the storage directories let totalSize = 0 // Helper function to calculate directory size const calculateSize = async (dirPath: string): Promise => { let size = 0 try { const files = await fs.promises.readdir(dirPath) for (const file of files) { const filePath = path.join(dirPath, file) const stats = await fs.promises.stat(filePath) if (stats.isDirectory()) { size += await calculateSize(filePath) } else { size += stats.size } } } catch (error: any) { if (error.code !== 'ENOENT') { console.error( `Error calculating size for directory ${dirPath}:`, error ) } } return size } // Calculate size for each directory const nounsDirSize = await calculateSize(this.nounsDir) const verbsDirSize = await calculateSize(this.verbsDir) const metadataDirSize = await calculateSize(this.metadataDir) const indexDirSize = await calculateSize(this.indexDir) totalSize = nounsDirSize + verbsDirSize + metadataDirSize + indexDirSize // CRITICAL FIX: Use persisted counts instead of directory reads // This is O(1) instead of O(n), and handles sharded structure correctly const nounsCount = this.totalNounCount const verbsCount = this.totalVerbCount // Count metadata files (these are NOT sharded) const metadataCount = ( await fs.promises.readdir(this.metadataDir) ).filter((file: string) => file.endsWith('.json')).length // Use persisted entity counts by type (O(1) instead of scanning all files) const nounTypeCounts: Record = Object.fromEntries(this.entityCounts) // Skip the expensive metadata file scan since we have counts const metadataFiles: string[] = [] // Empty array to skip the loop below for (const file of metadataFiles) { if (file.endsWith('.json')) { try { const filePath = path.join(this.metadataDir, file) const data = await fs.promises.readFile(filePath, 'utf-8') const metadata = JSON.parse(data) if (metadata.noun) { nounTypeCounts[metadata.noun] = (nounTypeCounts[metadata.noun] || 0) + 1 } } catch (error) { console.error(`Error reading metadata file ${file}:`, error) } } } return { type: 'filesystem', used: totalSize, quota: null, // File system doesn't provide quota information details: { rootDirectory: this.rootDir, nounsCount, verbsCount, metadataCount, nounsDirSize, verbsDirSize, metadataDirSize, indexDirSize, nounTypes: nounTypeCounts } } } catch (error) { console.error('Failed to get storage status:', error) return { type: 'filesystem', used: 0, quota: null, details: { error: String(error) } } } } // Removed 10 *_internal method overrides - now inherit from BaseStorage's type-first implementation // Removed 2 pagination methods (getNounsWithPagination, getVerbsWithPagination) - use BaseStorage's implementation public override supportsMultiProcessLocking(): boolean { return true } /** * Acquire the process-level writer lock for this data directory. Called by * `Brainy.init()` in writer mode. Refuses to open if another live writer * holds the lock — the worst possible default is to "succeed" with stale * indexes and silently lie about query results. * * Lock file: `/locks/_writer.lock`. * * Stale-lock detection: a lock is considered stale when ALL of: * 1. Same hostname as the current process (cross-host PID checks are unsafe). * 2. The recorded PID is no longer alive (`process.kill(pid, 0)` → ESRCH), OR * `lastHeartbeat` is older than `WRITER_STALE_THRESHOLD_MS` (60s). * Stale locks are overwritten with a warning. `force: true` overrides even live locks. */ public override async acquireWriterLock(options?: { force?: boolean }): Promise { await this.ensureInitialized() await this.ensureDirectoryExists(this.lockDir) const lockFile = path.join(this.lockDir, FileSystemStorage.WRITER_LOCK_FILE) const os = await import('node:os') const hostname = os.hostname() const myPid = typeof process !== 'undefined' && process.pid ? process.pid : 0 // Bounded acquire loop. The CLAIM itself is an atomic create-exclusive // write (O_EXCL) — two processes racing an ABSENT lock can never both // succeed, which closes the read-then-write window where the loser used // to keep running unlocked, silently. An EEXIST loser loops, re-reads, // and handles whatever it finds honestly (fresh foreign lock → loud // throw; stale/forced → verified takeover). const MAX_ATTEMPTS = 3 for (let attempt = 0; attempt < MAX_ATTEMPTS; attempt++) { const now = new Date().toISOString() const existing = await this.readWriterLock() // TORN-LOCK RECOVERY: power loss can legally leave the lock file // present but EMPTY/unparseable (the claim's non-atomic write died // mid-flight). readWriterLock() reports it as null — but the O_EXCL // claim below would EEXIST forever, a PERMANENT lockout no staleness // check can clear (staleness needs a parsed PID). A torn lock is // stale BY DEFINITION: no live holder has one (a holder either // completed its write or is dead). Unlink loudly and re-loop; a // racer that rewrites a VALID lock first simply wins the next read. if (existing === null) { try { await fs.promises.access(lockFile) console.warn( `[brainy] Writer lock at ${lockFile} exists but is unreadable/unparseable ` + `(torn write from a previous power loss) — treating as stale and removing.` ) try { await fs.promises.unlink(lockFile) } catch (unlinkErr: any) { if (unlinkErr.code !== 'ENOENT') throw unlinkErr } } catch (accessErr: any) { if (accessErr.code !== 'ENOENT') throw accessErr // Absent: the normal fresh-claim path below. } } // THE CLEAN-CLOSE RECORD IS READ BEFORE ANY VERDICT (see // WriterCloseRecord). A lock file whose release was RECORDED is // bookkeeping left by an orderly shutdown, not evidence of anything — // and that is true whether the previous holder was another process or // an earlier instance in THIS one. A production restart reported // "Re-acquiring writer lock ... this is a bug" immediately after a clean // close, sending an operator hunting for a leak that did not exist. const closeRecord = existing ? await this.readWriterCloseRecord() : null const releasedCleanly = existing !== null && closeRecord !== null && this.closeRecordVouchesFor(closeRecord, existing) if (existing) { // Same-process re-open: a second Brainy instance in this Node process // (e.g. test "simulate server restart" patterns, or a consumer that // explicitly re-instantiates without closing first). This isn't the // dangerous cross-process case the lock exists to prevent — the two // instances share a memory space and can't silently diverge from each // other beyond what their callers already see. Warn and take over — // unless the record proves the previous instance already let go, in // which case there is nothing to warn about. if (existing.pid === myPid && existing.hostname === hostname && !options?.force) { if (releasedCleanly) { console.warn( `[brainy] Clearing the leftover writer lock for ${this.rootDir} — an earlier ` + `instance in this process (PID ${existing.pid}) RELEASED it cleanly at ` + `${closeRecord!.closedAt} but could not remove the file. Nothing to recover.` ) } else { console.warn( `[brainy] Re-acquiring writer lock for ${this.rootDir} held by the same process (PID ${existing.pid}). ` + `If you intended to keep the previous Brainy instance alive, this is a bug — close it first.` ) } const info: WriterLockInfo = { pid: myPid, hostname, startedAt: now, lastHeartbeat: now, version: getBrainyVersion(), rootDir: this.rootDir } await this.writeFileAtomic(lockFile, JSON.stringify(info, null, 2)) await this.clearWriterCloseRecord() this.installWriterLock(info) return info } // A cleanly-released lock is stale by RECORD, not by inference. Only // when no record vouches for this lock do we fall back to pid // liveness, and then we say THAT honestly too: an unrecorded lock // means the writer did not complete its close, so the store was not // closed cleanly and this open pays recovery. const stale = releasedCleanly || (!options?.force && (await this.isWriterLockStale(existing))) if (!options?.force && !stale) { // Consumer-facing error contract: callers detect this case via // err.code and read the holder's details from err.lockInfo. throw this.writerLockedError(existing) } console.warn( options?.force ? `[brainy] Force-overwriting writer lock for ${this.rootDir} ` + `(was held by PID ${existing.pid} on ${existing.hostname}).` : releasedCleanly ? `[brainy] Clearing the leftover writer lock for ${this.rootDir} — ` + `PID ${existing.pid} on ${existing.hostname} RELEASED it cleanly at ` + `${closeRecord!.closedAt} but could not remove the file. ` + `Nothing to recover.` : `[brainy] Overwriting stale writer lock for ${this.rootDir} ` + `(PID ${existing.pid} on ${existing.hostname} is gone and left NO ` + `clean-close record — that writer did not finish closing, so this ` + `store was not closed cleanly; open will run crash recovery and ` + `report its wall).` ) // Takeover: verify the file still holds the lock we judged (a live // successor may have claimed meanwhile), then remove it and fall // through to the atomic claim below. A racing claimer who beats us to // the create simply wins — our next loop iteration reads their fresh // lock and throws honestly. (Advisory file locking has no // compare-and-delete; staleness requiring a 60s-old heartbeat keeps // the residual verify-to-unlink window practically unreachable.) const recheck = await this.readWriterLock() if ( recheck && (recheck.pid !== existing.pid || recheck.startedAt !== existing.startedAt || recheck.lastHeartbeat !== existing.lastHeartbeat) ) { continue // the lock changed hands while we deliberated — re-evaluate } try { await fs.promises.unlink(lockFile) } catch (err: any) { if (err.code !== 'ENOENT') throw err } } const info: WriterLockInfo = { pid: myPid, hostname, startedAt: existing && options?.force ? existing.startedAt : now, lastHeartbeat: now, version: getBrainyVersion(), rootDir: this.rootDir } // The atomic claim: write the FULL contents to a temp file, then // hard-link it into place — link(2) fails EEXIST if the target exists, // and the lock file appears with its complete JSON in one atomic step. // (The previous claim was writeFile with O_EXCL, whose open→write→close // is NOT atomic: a concurrent opener could read the file in its empty // window, judge it torn, unlink a LIVE claim, and take the lock — two // live writers. The link claim leaves no empty window to misread.) const claimTmp = `${lockFile}.claim-${myPid}-${Date.now()}` try { await fs.promises.writeFile(claimTmp, JSON.stringify(info, null, 2)) await fs.promises.link(claimTmp, lockFile) } catch (err: any) { if (err.code === 'EEXIST') { continue // someone else claimed between our read and create — re-evaluate } throw err } finally { await fs.promises.unlink(claimTmp).catch(() => {}) } // CONSUME the previous writer's clean-close record. It described the // lock generation that just ended; leaving it in place would let it // vouch for OUR lock if this process later dies without closing — // turning a real crash into a "closed cleanly" verdict. One unlink. await this.clearWriterCloseRecord() this.installWriterLock(info) return info } // Attempts exhausted: something is claiming this directory faster than we // can evaluate it. Read whoever holds it now and fail loudly with their // details rather than degrading into a lockless open. const holder = await this.readWriterLock() if (holder) throw this.writerLockedError(holder) throw new Error( `Failed to acquire the writer lock for ${this.rootDir} after ${MAX_ATTEMPTS} attempts — ` + `the lock file is being contended. Retry, or inspect ${lockFile}.` ) } /** Record lock ownership + start the unref'd heartbeat. */ private installWriterLock(info: WriterLockInfo): void { this.writerLockInfo = info // Heartbeat — rewrite lastHeartbeat every WRITER_HEARTBEAT_MS so other // processes can tell a live writer from one that crashed without releasing. this.writerLockHeartbeat = setInterval(() => { const tick = this.refreshWriterLockHeartbeat().catch((err) => { // ENOENT = the lock (or its directory) vanished mid-refresh — the // store was released or removed under us; the next acquire recreates // it. Benign by construction; anything else stays loud. if ((err as NodeJS.ErrnoException)?.code !== 'ENOENT') { console.warn('[brainy] Failed to refresh writer lock heartbeat:', err) } }) this.writerHeartbeatInFlight = tick.finally(() => { if (this.writerHeartbeatInFlight === tick) { this.writerHeartbeatInFlight = undefined } }) }, FileSystemStorage.WRITER_HEARTBEAT_MS) if (typeof this.writerLockHeartbeat.unref === 'function') { // Don't keep the event loop alive just for the heartbeat. this.writerLockHeartbeat.unref() } } /** * THE FENCE: verify this instance still owns the writer lock before a * commit barrier proceeds. An evicted writer (an operator's * `{ force: true }` takeover, or an operator deleting the lock file) must * fail LOUDLY on its next flush instead of writing on unaware — the * unfenced evicted writer was half of a production split-brain (each * writer flushing its own internally-consistent id-mapper snapshot, * alternating the store between two truths). One small file read per * flush window, never per record. No-op when this instance holds no * writer lock (read-only opens, in-memory stores). * * @throws `BRAINY_WRITER_FENCED` when the lock is gone or held by another. */ public override async assertWriterFenceHeld(): Promise { if (!this.writerLockInfo) return const current = await this.readWriterLock() // Ownership is PER-PROCESS: pid + hostname, deliberately NOT startedAt. // The documented same-process re-open path ("warn and take over" — two // instances in one Node process, the server-restart test pattern) // rewrites the lock with a fresh startedAt; fencing the first instance // on that mismatch latched its background flushes dead while its own // process held the lock (caught by the plant's integration lane, twice). // startedAt adds nothing against pid recycling either: a recycled pid's // victim is a DEAD process — it runs no fence checks. if ( current && current.pid === this.writerLockInfo.pid && current.hostname === this.writerLockInfo.hostname ) { return } const err = new Error( `Writer fence lost for ${this.rootDir}: this process (PID ${this.writerLockInfo.pid}) ` + `no longer holds the writer lock — ` + (current ? `it is now held by PID ${current.pid} on ${current.hostname} (since ${current.startedAt}).` : `the lock file is gone (released or removed by an operator).`) + `\nThis instance refuses to commit further writes: a fenced-out writer continuing to ` + `flush is how split-brain stores are made. Close this instance; if the takeover was a ` + `mistake, close the successor and re-open.` ) as Error & { code: string } err.code = 'BRAINY_WRITER_FENCED' throw err } /** The consumer-facing BRAINY_WRITER_LOCKED error, holder details attached. */ private writerLockedError(existing: WriterLockInfo): Error { const err = new Error( `Another writer holds this Brainy directory.\n` + ` PID: ${existing.pid} on host ${existing.hostname}\n` + ` Started: ${existing.startedAt}\n` + ` Heartbeat: ${existing.lastHeartbeat}\n` + ` Version: ${existing.version}\n` + ` Directory: ${this.rootDir}\n\n` + `For diagnostic queries against this live store, use:\n` + ` const reader = await Brainy.openReadOnly({ storage: { type: 'filesystem', path: '${this.rootDir}' } })\n\n` + `If you have verified the existing lock is stale (e.g. a crashed writer on a different host that PID liveness cannot reach), pass { force: true }.` ) as Error & { code: string; lockInfo: WriterLockInfo } err.code = 'BRAINY_WRITER_LOCKED' err.lockInfo = existing return err } public override async releaseWriterLock(): Promise { if (this.writerLockHeartbeat) { clearInterval(this.writerLockHeartbeat) this.writerLockHeartbeat = undefined } // Drain an in-flight heartbeat tick BEFORE unlinking: clearInterval stops // future ticks only, and a straggler write landing after the unlink would // re-create the lock as a phantom (blocking the next writer until the // stale TTL). After the drain, any refresh is either fully landed (we // unlink its output below) or not started (it sees writerLockInfo // undefined and returns). if (this.writerHeartbeatInFlight) { await this.writerHeartbeatInFlight } if (!this.writerLockInfo) { return } const lockFile = path.join(this.lockDir, FileSystemStorage.WRITER_LOCK_FILE) const released = this.writerLockInfo try { // Only delete if we still own it — avoid clobbering a successor that // claimed the lock via force-override. const current = await this.readWriterLock() const ours = current === null || (current.pid === released.pid && current.hostname === released.hostname) if (current && ours) { await fs.promises.unlink(lockFile) } // THE CLEAN-CLOSE RECORD (see WriterCloseRecord). Written whenever this // instance gives up a lock nobody else has taken — the unlink above // having succeeded OR the file already being gone. The next open reads // it instead of guessing from pid liveness: a recorded release is an // orderly shutdown, an absent record is a writer that never finished // closing. Not written when a successor holds the lock: our release is // then a no-op and a record would slander their live lock. if (ours) { await this.writeWriterCloseRecord(released) } } catch (err: any) { if (err.code !== 'ENOENT') { console.warn('[brainy] Failed to release writer lock file:', err) } } finally { this.writerLockInfo = undefined } } /** * @description Read the clean-close record at `locks/_writer.close`, or * `null` when it is absent or unparseable. A torn record is treated as * absent — the conservative direction, since an unreadable record can * vouch for nothing. * @returns The record, or null. */ public async readWriterCloseRecord(): Promise { await this.ensureInitialized() const recordFile = path.join(this.lockDir, FileSystemStorage.WRITER_CLOSE_FILE) try { const raw = await fs.promises.readFile(recordFile, 'utf-8') const parsed = JSON.parse(raw) as WriterCloseRecord if ( typeof parsed?.pid !== 'number' || typeof parsed?.hostname !== 'string' || typeof parsed?.startedAt !== 'string' || typeof parsed?.closedAt !== 'string' ) { return null } return parsed } catch (err: any) { if (err.code === 'ENOENT') return null return null } } /** * @description Whether a clean-close record describes the very lock * generation `lock` represents. The match is pid + hostname + `startedAt`: * `startedAt` is the lock generation's identity, so a record can never * vouch for a LATER lock taken by the same pid on the same host (the * same-process re-open path mints a fresh `startedAt`). * @param record - The clean-close record read from disk. * @param lock - The lock file's contents. */ private closeRecordVouchesFor(record: WriterCloseRecord, lock: WriterLockInfo): boolean { return ( record.pid === lock.pid && record.hostname === lock.hostname && record.startedAt === lock.startedAt ) } /** * @description Write the clean-close record for a lock this instance just * released. Atomic (temp + rename) so a concurrent opener never reads half * a record. A failure here costs the next open nothing but the honest * fallback (pid liveness), so it warns rather than failing the close. * @param released - The lock info this instance held. */ private async writeWriterCloseRecord(released: WriterLockInfo): Promise { const record: WriterCloseRecord = { pid: released.pid, hostname: released.hostname, startedAt: released.startedAt, closedAt: new Date().toISOString(), version: released.version } const recordFile = path.join(this.lockDir, FileSystemStorage.WRITER_CLOSE_FILE) try { await this.writeFileAtomic(recordFile, JSON.stringify(record, null, 2)) } catch (err) { // ENOENT = the lock directory is gone, i.e. the whole store was removed // under us. There is no next open to inform. if ((err as NodeJS.ErrnoException)?.code === 'ENOENT') return console.warn( `[brainy] Failed to write the writer clean-close record for ${this.rootDir} — ` + `the next open will fall back to pid liveness and may report this orderly ` + `shutdown as a crash:`, err ) } } /** * @description Remove the clean-close record. Called by every successful * lock claim so a record never outlives the lock generation it describes. */ private async clearWriterCloseRecord(): Promise { const recordFile = path.join(this.lockDir, FileSystemStorage.WRITER_CLOSE_FILE) try { await fs.promises.unlink(recordFile) } catch (err: any) { if (err.code !== 'ENOENT') { console.warn('[brainy] Failed to clear the writer clean-close record:', err) } } } public override async readWriterLock(): Promise { await this.ensureInitialized() const lockFile = path.join(this.lockDir, FileSystemStorage.WRITER_LOCK_FILE) try { const raw = await fs.promises.readFile(lockFile, 'utf-8') return JSON.parse(raw) as WriterLockInfo } catch (err: any) { if (err.code === 'ENOENT') return null console.warn('[brainy] Failed to read writer lock file:', err) return null } } /** * Atomically refresh `lastHeartbeat` on the writer lock we own. Skips if the * lock file has been deleted out from under us (e.g. operator removed it). */ private async refreshWriterLockHeartbeat(): Promise { if (!this.writerLockInfo) return const current = await this.readWriterLock() if (!current) return // Defensive: don't overwrite if a successor claimed the lock. if (current.pid !== this.writerLockInfo.pid || current.hostname !== this.writerLockInfo.hostname) { return } const updated: WriterLockInfo = { ...this.writerLockInfo, lastHeartbeat: new Date().toISOString() } const lockFile = path.join(this.lockDir, FileSystemStorage.WRITER_LOCK_FILE) await this.writeFileAtomic(lockFile, JSON.stringify(updated, null, 2)) this.writerLockInfo = updated } /** * Determine whether an existing writer lock is stale (safe to overwrite). * Same hostname and DEAD PID → stale. That is the whole rule: a LIVE * process is never auto-evicted, however old its heartbeat — a >60s * event-loop stall (debugger pause, GC, heavy sync work) is a slow writer, * not a dead one, and heartbeat-age eviction of live writers was the * dominant mechanism behind a production split-brain (two live unaware * writers alternating a store's id-mapper between two truths). A holder * that LOOKS alive but is truly wedged is the operator's call via * `{ force: true }` — and the fence check on every flush * ({@link assertWriterFenceHeld}) guarantees a forced-out holder fails * loudly instead of writing on. Different hostname → cannot prove * anything, treat as live. The heartbeat remains for OBSERVABILITY (the * lock error names it so an operator can judge staleness themselves). */ private async isWriterLockStale(lock: WriterLockInfo): Promise { const os = await import('node:os') if (lock.hostname !== os.hostname()) { return false } return !this.isPidAlive(lock.pid) } /** * `process.kill(pid, 0)` sends signal 0 — no signal is actually delivered; * it just checks whether the kernel still considers `pid` reachable from * this process. ESRCH = no such process. EPERM = exists but we can't signal * it, which still proves it's alive. */ private isPidAlive(pid: number): boolean { if (!pid || pid <= 0) return false try { process.kill(pid, 0) return true } catch (err: any) { if (err.code === 'EPERM') return true return false } } /** * Atomic write via temp-file-then-rename so concurrent readers never see a * half-written lock JSON. Reused by writer-lock writes + heartbeat. */ private async writeFileAtomic(filePath: string, contents: string): Promise { const tmp = `${filePath}.tmp-${process.pid}-${Date.now()}` await fs.promises.writeFile(tmp, contents) await fs.promises.rename(tmp, filePath) } /** * Start watching for cross-process flush requests. Called by Brainy.init() * in writer mode. Each new `.req` file in `locks/_flush_requests/` triggers * the supplied callback (`brain.flush()`), after which an `.ack` is written * to `locks/_flush_responses/` with the same request ID. Stale `.req` files * (>FLUSH_REQUEST_TTL_MS) are garbage-collected on each sweep. * * THE WATCH IS EVENT-DRIVEN, NOT A POLL. It used to `readdir` the request * directory every 500 ms, per brain, for the entire life of every writer — * armed on every non-reader brain whether or not any inspector process * existed. MEASURED on a production process holding 21 brains: 42 directory * reads per second on a completely idle service, plus a stale-request GC * pass on every one of them. The engine does no periodic work without a * cause, and a request that has not been made is not a cause. * * `fs.watch` (inotify on Linux) delivers the arrival itself, so a request is * seen SOONER than the old poll saw it. Two honest concessions ride with it: * - a slow SAFETY SWEEP (FLUSH_SAFETY_SWEEP_MS) still runs, because * `fs.watch` can miss events on network and fuse filesystems and because * the stale-request GC needs some tick of its own. At 30s that is 0.7 * reads/s across 21 brains where the poll cost 42. * - a filesystem that cannot watch at all falls back to the ORIGINAL * 500 ms poll, narrated once, because correctness outranks idle cost: * an inspector whose request is never seen waits forever. */ public override startFlushRequestWatcher(onRequest: () => Promise): void { // Already watching — or already ARMING. The arm is asynchronous (the // request directory is created before it can be watched), so neither the // watcher nor the interval exists yet during that window; the callback is // the flag that covers it. Without this a second call in the window would // leave two watchers and two sweeps running for the life of the store. if (this.flushWatcherInterval || this.flushWatcher || this.flushWatcherOnRequest) return this.flushWatcherOnRequest = onRequest const reqDir = path.join(this.lockDir, FileSystemStorage.FLUSH_REQUEST_DIR) const ackDir = path.join(this.lockDir, FileSystemStorage.FLUSH_RESPONSE_DIR) const sweep = (): void => { if (this.flushWatcherInFlight) return // skip overlapping sweep this.flushWatcherInFlight = true this.processFlushRequests(reqDir, ackDir).finally(() => { this.flushWatcherInFlight = false }) } // Ensure both dirs exist up front so the first .req drop doesn't race with // mkdir — and so there is a directory to watch. void this.ensureDirectoryExists(reqDir) .then(() => this.ensureDirectoryExists(ackDir)) .then(() => { if (this.flushWatcherOnRequest !== onRequest) return // stopped meanwhile try { const watcher = fs.watch(reqDir, () => sweep()) this.flushWatcher = watcher watcher.on('error', (err: Error) => { // A watch that dies mid-life must not leave the door deaf. console.warn( `[brainy] Flush-request watch failed (${err.message}) — falling back to polling.` ) this.flushWatcher?.close() this.flushWatcher = undefined // The SAFETY sweep must go first. It is already armed at 30s, and // startFlushRequestPolling() declines to arm over an existing // interval — so leaving it would quietly leave this store answering // flush requests on a 30s cadence instead of the 500ms one the door // promises. A degrade nobody asked for is still a degrade. if (this.flushWatcherInterval) { clearInterval(this.flushWatcherInterval) this.flushWatcherInterval = undefined } this.startFlushRequestPolling(sweep) }) if (typeof watcher.unref === 'function') watcher.unref() // The safety sweep: missed events on exotic filesystems, and the // stale-request GC. this.flushWatcherInterval = setInterval(sweep, FileSystemStorage.FLUSH_SAFETY_SWEEP_MS) if (typeof this.flushWatcherInterval.unref === 'function') { this.flushWatcherInterval.unref() } // One sweep now: a request may have been dropped before the watch armed. sweep() } catch (err) { console.warn( `[brainy] Flush-request directory cannot be watched on this filesystem ` + `(${(err as Error).message}) — polling every ` + `${FileSystemStorage.FLUSH_WATCH_INTERVAL_MS}ms instead.` ) this.startFlushRequestPolling(sweep) } }) .catch(() => { // The request directory could not be created; nothing to watch. A // cross-process flush request cannot be made either, so there is // nothing to miss. }) } /** The original 500 ms poll — the fallback when a directory cannot be watched. */ private startFlushRequestPolling(sweep: () => void): void { if (this.flushWatcherInterval) return this.flushWatcherInterval = setInterval(sweep, FileSystemStorage.FLUSH_WATCH_INTERVAL_MS) if (typeof this.flushWatcherInterval.unref === 'function') { this.flushWatcherInterval.unref() } } public override stopFlushRequestWatcher(): void { if (this.flushWatcher) { this.flushWatcher.close() this.flushWatcher = undefined } if (this.flushWatcherInterval) { clearInterval(this.flushWatcherInterval) this.flushWatcherInterval = undefined } this.flushWatcherOnRequest = undefined } /** * Process any pending `.req` files: invoke the flush callback once, then * write an `.ack` for each request. Multiple requests that arrived in the * same tick share a single flush — they all see the same ack timestamp. */ private async processFlushRequests(reqDir: string, ackDir: string): Promise { let entries: string[] try { entries = await fs.promises.readdir(reqDir) } catch (err: any) { if (err.code !== 'ENOENT') { console.warn('[brainy] Flush watcher readdir failed:', err) } return } const reqs = entries.filter((e: string) => e.endsWith('.req')) if (reqs.length === 0) return // One flush per tick — N pending requests share it. const callback = this.flushWatcherOnRequest if (!callback) return let flushError: Error | null = null try { await callback() } catch (err: any) { flushError = err instanceof Error ? err : new Error(String(err)) console.warn('[brainy] Flush callback threw inside flush-request watcher:', err) } const ackTimestamp = new Date().toISOString() for (const filename of reqs) { const requestId = filename.replace(/\.req$/, '') const ackPath = path.join(ackDir, `${requestId}.ack`) const ackBody = JSON.stringify({ requestId, completedAt: ackTimestamp, ok: !flushError, error: flushError ? flushError.message : undefined }) try { await this.writeFileAtomic(ackPath, ackBody) await fs.promises.unlink(path.join(reqDir, filename)) } catch (err: any) { if (err.code !== 'ENOENT') { console.warn('[brainy] Failed to write flush ack:', err) } } } // Garbage-collect stale .req files in case a watcher missed them (e.g. // the writer was restarted between request and processing). const now = Date.now() for (const filename of entries) { if (!filename.endsWith('.req')) continue const fp = path.join(reqDir, filename) try { const stat = await fs.promises.stat(fp) if (now - stat.mtimeMs > FileSystemStorage.FLUSH_REQUEST_TTL_MS) { await fs.promises.unlink(fp).catch(() => {}) } } catch { // ignore — file already removed } } } /** * Inspector side: drop a `.req` file and poll for the corresponding `.ack`. * Returns true if the writer acknowledged in time, false on timeout. */ public override async requestFlushOverFilesystem(timeoutMs: number): Promise { await this.ensureInitialized() const reqDir = path.join(this.lockDir, FileSystemStorage.FLUSH_REQUEST_DIR) const ackDir = path.join(this.lockDir, FileSystemStorage.FLUSH_RESPONSE_DIR) await this.ensureDirectoryExists(reqDir) await this.ensureDirectoryExists(ackDir) const requestId = `${Date.now()}-${process.pid}-${Math.random().toString(36).slice(2, 10)}` const reqPath = path.join(reqDir, `${requestId}.req`) const ackPath = path.join(ackDir, `${requestId}.ack`) const body = JSON.stringify({ requestId, requestedAt: new Date().toISOString(), requesterPid: process.pid }) await this.writeFileAtomic(reqPath, body) const deadline = Date.now() + timeoutMs while (Date.now() < deadline) { try { await fs.promises.access(ackPath, fs.constants.F_OK) await fs.promises.unlink(ackPath).catch(() => {}) return true } catch { // not yet } await new Promise((r) => setTimeout(r, FileSystemStorage.FLUSH_POLL_INTERVAL_MS)) } // Timeout — best effort to clean up our request so the writer doesn't act // on stale work later. Watcher will GC anyway after FLUSH_REQUEST_TTL_MS. await fs.promises.unlink(reqPath).catch(() => {}) return false } /** * Acquire a file-based lock for coordinating operations across multiple processes * @param lockKey The key to lock on * @param ttl Time to live for the lock in milliseconds (default: 30 seconds) * @returns Promise that resolves to true if lock was acquired, false otherwise */ private async acquireLock( lockKey: string, ttl: number = 30000 ): Promise { await this.ensureInitialized() // Ensure lock directory exists await this.ensureDirectoryExists(this.lockDir) const lockFile = path.join(this.lockDir, `${lockKey}.lock`) const lockValue = `${Date.now()}_${Math.random()}_${process.pid || 'unknown'}` const expiresAt = Date.now() + ttl try { // Check if lock file already exists and is still valid try { const lockData = await fs.promises.readFile(lockFile, 'utf-8') const lockInfo = JSON.parse(lockData) if (lockInfo.expiresAt > Date.now()) { // Lock exists and is still valid return false } } catch (error: any) { // If file doesn't exist or can't be read, we can proceed to create the lock if (error.code !== 'ENOENT') { console.warn(`Error reading lock file ${lockFile}:`, error) } } // Try to create the lock file const lockInfo = { lockValue, expiresAt, pid: process.pid || 'unknown', timestamp: Date.now() } await fs.promises.writeFile(lockFile, JSON.stringify(lockInfo, null, 2)) // Add to active locks for cleanup this.activeLocks.add(lockKey) // Schedule automatic cleanup when lock expires setTimeout(() => { this.releaseLock(lockKey, lockValue).catch((error) => { console.warn(`Failed to auto-release expired lock ${lockKey}:`, error) }) }, ttl) return true } catch (error) { console.warn(`Failed to acquire lock ${lockKey}:`, error) return false } } /** * Release a file-based lock * @param lockKey The key to unlock * @param lockValue The value used when acquiring the lock (for verification) * @returns Promise that resolves when lock is released */ private async releaseLock( lockKey: string, lockValue?: string ): Promise { await this.ensureInitialized() const lockFile = path.join(this.lockDir, `${lockKey}.lock`) try { // If lockValue is provided, verify it matches before releasing if (lockValue) { try { const lockData = await fs.promises.readFile(lockFile, 'utf-8') const lockInfo = JSON.parse(lockData) if (lockInfo.lockValue !== lockValue) { // Lock was acquired by someone else, don't release it return } } catch (error: any) { // If lock file doesn't exist, that's fine if (error.code === 'ENOENT') { return } throw error } } // Delete the lock file await fs.promises.unlink(lockFile) // Remove from active locks this.activeLocks.delete(lockKey) } catch (error: any) { if (error.code !== 'ENOENT') { console.warn(`Failed to release lock ${lockKey}:`, error) } } } /** * Clean up expired lock files */ private async cleanupExpiredLocks(): Promise { await this.ensureInitialized() try { const lockFiles = await fs.promises.readdir(this.lockDir) const now = Date.now() for (const lockFile of lockFiles) { if (!lockFile.endsWith('.lock')) continue const lockPath = path.join(this.lockDir, lockFile) try { const lockData = await fs.promises.readFile(lockPath, 'utf-8') const lockInfo = JSON.parse(lockData) if (lockInfo.expiresAt <= now) { await fs.promises.unlink(lockPath) const lockKey = lockFile.replace('.lock', '') this.activeLocks.delete(lockKey) } } catch (error) { // If we can't read or parse the lock file, remove it try { await fs.promises.unlink(lockPath) } catch (unlinkError) { console.warn( `Failed to cleanup invalid lock file ${lockPath}:`, unlinkError ) } } } } catch (error) { console.warn('Failed to cleanup expired locks:', error) } } /** * Save statistics data to storage with file-based locking */ protected async saveStatisticsData( statistics: StatisticsData ): Promise { const lockKey = 'statistics' const lockAcquired = await this.acquireLock(lockKey, 10000) // 10 second timeout if (!lockAcquired) { console.warn( 'Failed to acquire lock for statistics update, proceeding without lock' ) } try { // Get existing statistics to merge with new data const existingStats = await this.getStatisticsData() if (existingStats) { // Merge statistics data const mergedStats: StatisticsData = { totalNodes: Math.max( statistics.totalNodes || 0, existingStats.totalNodes || 0 ), totalEdges: Math.max( statistics.totalEdges || 0, existingStats.totalEdges || 0 ), totalMetadata: Math.max( statistics.totalMetadata || 0, existingStats.totalMetadata || 0 ), // Preserve any additional fields from existing stats ...existingStats, // Override with new values where provided ...statistics, // Always update lastUpdated to current time lastUpdated: new Date().toISOString() } await this.saveStatisticsWithBackwardCompat(mergedStats) } else { // No existing statistics, save new ones const newStats: StatisticsData = { ...statistics, lastUpdated: new Date().toISOString() } await this.saveStatisticsWithBackwardCompat(newStats) } } finally { if (lockAcquired) { await this.releaseLock(lockKey) } } } /** * Get statistics data from storage */ protected async getStatisticsData(): Promise { try { const statsPath = path.join(this.systemDir, `${STATISTICS_KEY}.json`) const data = await fs.promises.readFile(statsPath, 'utf-8') return JSON.parse(data) } catch (error: any) { if (error.code !== 'ENOENT') { console.error('Error reading statistics:', error) } return null } } /** * Save statistics to storage */ private async saveStatisticsWithBackwardCompat(statistics: StatisticsData): Promise { const statsPath = path.join(this.systemDir, `${STATISTICS_KEY}.json`) await this.ensureDirectoryExists(this.systemDir) await fs.promises.writeFile(statsPath, JSON.stringify(statistics, null, 2)) } // ============================================= // Count Management for O(1) Scalability // ============================================= /** * Initialize counts from filesystem storage */ protected async initializeCounts(): Promise { if (!this.countsFilePath) return try { if (await this.fileExists(this.countsFilePath)) { const data = await fs.promises.readFile(this.countsFilePath, 'utf-8') const counts = JSON.parse(data) // Restore entity counts this.entityCounts = new Map(Object.entries(counts.entityCounts || {})) this.verbCounts = new Map(Object.entries(counts.verbCounts || {})) this.totalNounCount = counts.totalNounCount || 0 this.totalVerbCount = counts.totalVerbCount || 0 // The ALL-visibility scalars (ledger denominators). A counts.json // written before they existed carries neither key: derive both ONCE // from the canonical id tree (an id-directory listing — O(ids), no // record reads), persist, and never scan again. Absent keys are a // legacy file, not a zero — a zero here would make every provider's // coverage ledger read "over-posted" on a populated store. let needsPersist = false if ( typeof counts.totalNounCountAll === 'number' && typeof counts.totalVerbCountAll === 'number' ) { this.totalNounCountAll = counts.totalNounCountAll this.totalVerbCountAll = counts.totalVerbCountAll if (counts.allCountsDerivedBy === 'identity-record') { // Derived (or recounted) under the honest rule — one counted // entity per metadata content leg. Trust the persisted suspect // flag as-is; an unprovable delete since may still have set it. this.allCountsDerivedBy = 'identity-record' this.allCountsSuspect = counts.allCountsSuspect === true } else { // The ALL scalars exist but predate the identity-record stamp — // they were derived under the legacy rule that counted one // entity per id DIRECTORY, so orphaned ghost/scar containers (a // pre-8.3.1 partial-delete defect — see pruneOrphanedEntities()) // were counted as entities too. O(1) field read, NEVER a walk // here: force suspect and name it loudly. A sanctioned recount // (repairIndex) restores exact denominators and clears this. this.allCountsDerivedBy = undefined this.allCountsSuspect = true needsPersist = true prodLog.narrate( '[FileSystemStorage] canonical count ledger was derived under the legacy ' + 'container rule — it counts one entity per id DIRECTORY, so every ghost/scar ' + 'container inflates it. Marked suspect, and an honest recount is scheduled to ' + 'run in the background after this open; until it lands, do not subtract ' + 'against these ALL scalars.' ) // A suspect ledger used to stay wrong for the life of the store, // waiting for an operator to run repairIndex. A downstream index // heal took its "remaining" figure from these inflated // denominators and reported work that did not exist. The ledger // now HEALS ITSELF — in the background, because a denominator is // a derived scalar and no read is ever served from it. this.scheduleCountLedgerDerivation('legacy container-rule ledger') } } else { // No ALL scalars at all. There is nothing to serve in the meantime — // a zero would read as an empty store — so the scalars stay unknown // and SUSPECT until the background derivation lands. The open does // not wait for it: an id-tree walk is O(ids) and this file has been // the whole reason a 24k-id store opened in silence. this.allCountsSuspect = true this.scheduleCountLedgerDerivation('counts.json predates the ALL-visibility ledger') } // The vectored-noun scalar (shipped after the ALL scalars above — a // counts.json can carry `totalNounCountAll`/`totalVerbCountAll` but // still predate THIS key). Unlike the ALL scalars, presence cannot be // decided from the id-directory listing alone: a deferred-embed // noun's `vectors.json` EXISTS with an empty `vector: []` until its // embed lands, so this derivation reads every noun's `vectors.json` // ONCE (O(nouns) reads, not O(ids) listing) — honest, one-time cost. if (typeof counts.totalVectoredNounCount === 'number') { this.totalVectoredNounCount = counts.totalVectoredNounCount } else { // O(nouns) CONTENT reads — the most expensive derivation of the // three, and the one most likely to have been the silent minutes at // the front of a large store's open. Background, suspect until it // lands, same as the ALL scalars. this.allCountsSuspect = true this.scheduleCountLedgerDerivation('counts.json predates the vectored-noun ledger') } if (needsPersist) { await this.persistCounts() } // Also populate the cache for backward compatibility this.countCache.set('nouns_count', { count: this.totalNounCount, timestamp: Date.now() }) this.countCache.set('verbs_count', { count: this.totalVerbCount, timestamp: Date.now() }) } else { // If no counts file exists, do one initial count await this.initializeCountsFromDisk() } } catch (error) { console.warn('Could not load persisted counts, will initialize from disk:', error) await this.initializeCountsFromDisk() } } /** * Initialize counts by scanning disk (only done once) */ private async initializeCountsFromDisk(): Promise { const startedAt = Date.now() // THIS ONE CANNOT LEAVE THE FOREGROUND, and the reason is worth stating: // it derives `totalNounCount` / `totalVerbCount`, the scalars // `getNounCount()` and `getVerbCount()` RETURN. Backgrounding it would // make a populated store answer "0 entities" until the walk landed — a // wrong answer, not a slow one, and the serving law grades a failure by // whether an answer could be wrong. The ALL-visibility denominators, which // no read is served from, DO run in the background (see // scheduleCountLedgerDerivation). What this walk owes the operator instead // is narration: it announces itself, and reports its wall. prodLog.narrate( `[FileSystemStorage] no usable counts.json — deriving the entity counters from ` + `the canonical id tree now. This is O(ids) listings plus one vectors.json read ` + `per noun, and it BLOCKS the open because getNounCount()/getVerbCount() are ` + `served from it. It runs once; the result is persisted.` ) try { // Count the CANONICAL 8.0 layout (`entities////…`) — // the tree saveNoun/getNouns actually read and write. The previous scan // counted the vestigial `entities/*/hnsw` directories, which the 8.0 // write path never populates, so a store recovering from a lost or // corrupted counts.json re-initialized every counter to ZERO on real // data (wrong stats/boot logs and a mis-sized rebuild-strategy // decision at open). const nouns = await this.scanCanonicalEntities('nouns') this.totalNounCount = nouns.count const verbs = await this.scanCanonicalEntities('verbs') this.totalVerbCount = verbs.count // The id-tree scan counts every tier — it IS the ALL-visibility ledger. this.totalNounCountAll = nouns.count this.totalVerbCountAll = verbs.count this.allCountsSuspect = false this.allCountsDerivedBy = 'identity-record' // Vectored-noun scalar: presence needs each noun's vectors.json CONTENT // (a deferred-embed noun's file exists but holds an empty vector until // its embed lands), so this is a full O(nouns) content scan — see // scanVectoredNounCount()'s JSDoc for the cost note. Paid once, here, // alongside the rest of this from-disk recovery. this.totalVectoredNounCount = await this.scanVectoredNounCount() // Sample some entities for the type distribution (don't read all). // Read the metadata files DIRECTLY with fs — this runs inside init(), // and every guarded accessor (getNounMetadata → ensureInitialized) // re-enters init() from here, deadlocking the open. (Latent in the old // code too: its scan of the empty hnsw dirs just never sampled.) for (const entityDir of nouns.sampleDirs) { const metadata = await this.readEntityMetadataRaw(entityDir) if (metadata) { const type = metadata.noun || 'default' this.entityCounts.set(type, (this.entityCounts.get(type) || 0) + 1) } } // Extrapolate counts if we sampled const sampleSize = nouns.sampleDirs.length if (sampleSize < this.totalNounCount && sampleSize > 0) { const multiplier = this.totalNounCount / sampleSize for (const [type, count] of this.entityCounts.entries()) { this.entityCounts.set(type, Math.round(count * multiplier)) } } await this.persistCounts() prodLog.narrate( `[FileSystemStorage] counter derivation from the canonical id tree finished in ` + `${Date.now() - startedAt}ms: ${this.totalNounCount} nouns, ${this.totalVerbCount} verbs, ` + `${this.totalVectoredNounCount} vectored nouns — persisted, stamped identity-record.` ) } catch (error) { console.error('Error initializing counts from disk:', error) } } /** * Walk the canonical `entities//<2-hex-shard>//` tree, counting * one entity per id directory that holds the metadata CONTENT leg * (`metadata.json` or its `.json.gz` variant — see * {@link hasMetadataContentLeg}). A bare container — a ghost (a stale * `vectors.json` left with no metadata leg) or a scar (an empty directory), * both artifacts of the pre-8.3.1 partial-delete defect — counts ZERO: the * identity record IS the population (ADR-008 G1), never the directory. * This is the ONE-TIME legacy derivation walk (see callers); a prior * version of this scan counted every id directory regardless of content, * over-counting any store carrying orphaned containers — see * `allCountsDerivedBy` for how a counts.json derived under that old rule is * marked suspect on load. Returns up to 100 sampled *counted* entity * directories (absolute paths) — nouns feed the type-distribution estimate * above. An absent tree (fresh store) counts zero. */ /** * @description Derive the ALL-visibility count ledger honestly — one entity * per IDENTITY RECORD, never per id directory — IN THE BACKGROUND, once, * and persist the result stamped `identity-record`. * * Why background: these scalars are DENOMINATORS. No read is served from * them, so deriving them cannot be allowed to hold an open hostage — a * store with 24,898 ids spent minutes of a production restart inside walks * exactly like these, in silence, before serving anything. Why at all: a * ledger derived under the old container rule stayed wrong for the life of * the store, and a downstream index heal subtracted against it and reported * remaining work that did not exist (measured on a real store: 14,231 * derived against 14,056 identity records — precisely the store's 25 noun * scar directories; verbs 72,729 against 72,679, its 50 verb scars). * * Idempotent: a second call while one is in flight joins the first. * @param reason - What made the ledger untrustworthy, quoted in narration. * @returns Nothing; observe completion with {@link whenCountLedgerSettled}. */ private scheduleCountLedgerDerivation(reason: string): void { if (this.countLedgerDerivation) return this.countLedgerDerivation = (async () => { const startedAt = Date.now() prodLog.narrate( `[FileSystemStorage] count-ledger derivation started in the background ` + `(${reason}) — counting identity records, not id directories; the open does ` + `not wait for it and no read is served from these scalars.` ) try { const beforeNouns = this.totalNounCountAll const beforeVerbs = this.totalVerbCountAll const beforeVectored = this.totalVectoredNounCount // A walk that RACED A WRITE cannot prove its number: a row that landed // mid-walk may or may not have been in the shard the walk had already // passed. Rather than persist a figure that might be off by one and // stamp it "exact", the walk is repeated once on a quiet store, and if // the store is never quiet the ledger stays SUSPECT and says so. One // retry, never a spin. let attempt = 0 let derived: { nouns: number; verbs: number; vectored: number } | null = null while (attempt < 2 && derived === null) { attempt++ const activityBefore = this.ledgerActivityStamp() const nouns = await this.scanCanonicalEntities('nouns') const verbs = await this.scanCanonicalEntities('verbs') const vectored = await this.scanVectoredNounCount() if (this.ledgerActivityStamp() === activityBefore) { derived = { nouns: nouns.count, verbs: verbs.count, vectored } } } if (derived === null) { this.allCountsSuspect = true prodLog.narrate( `[FileSystemStorage] count-ledger derivation could not finish on a quiet store ` + `after ${attempt} attempts (${Date.now() - startedAt}ms) — writes landed during ` + `every walk. The ALL-visibility scalars stay SUSPECT and must not be subtracted ` + `against; brain.repairIndex() derives them under a recount barrier.` ) return } this.totalNounCountAll = derived.nouns this.totalVerbCountAll = derived.verbs this.totalVectoredNounCount = derived.vectored this.allCountsDerivedBy = 'identity-record' this.allCountsSuspect = false await this.persistCounts() prodLog.narrate( `[FileSystemStorage] count-ledger derivation finished in ${Date.now() - startedAt}ms: ` + `${derived.nouns} nouns / ${derived.verbs} verbs / ${derived.vectored} vectored nouns` + (beforeNouns !== derived.nouns || beforeVerbs !== derived.verbs || beforeVectored !== derived.vectored ? ` (corrected from ${beforeNouns} / ${beforeVerbs} / ${beforeVectored} — the ` + `difference is ghost and scar containers the old rule counted as entities)` : ' (unchanged)') + ` — persisted, stamped identity-record, no longer suspect.` ) } catch (error) { // The ledger stays suspect and the next open retries. Loud: a // denominator nobody can derive is a fact an operator must have. this.allCountsSuspect = true prodLog.error( `[FileSystemStorage] count-ledger derivation FAILED after ` + `${Date.now() - startedAt}ms — the ALL-visibility scalars remain SUSPECT ` + `and must not be subtracted against; the next open retries:`, error ) } })() } /** * @description A cheap witness that the ledger changed while a walk was * running. Every landed write moves one of these live counters, so an * unchanged stamp across a walk means no write landed during it. * @returns A value that differs whenever the live ALL scalars have moved. */ private ledgerActivityStamp(): string { return `${this.totalNounCountAll}:${this.totalVerbCountAll}:${this.totalVectoredNounCount}` } /** * @description Resolve once any background count-ledger derivation has * settled (succeeded or failed). Resolves immediately when none was needed. * Exists so tests and operators can observe the ledger's honest value rather * than race it; nothing in the read path waits on this. * @returns A promise that settles with the derivation. */ public async whenCountLedgerSettled(): Promise { await this.countLedgerDerivation } private async scanCanonicalEntities( kind: 'nouns' | 'verbs' ): Promise<{ count: number; sampleDirs: string[] }> { const base = path.join(this.rootDir, 'entities', kind) const SAMPLE_MAX = 100 let count = 0 const sampleDirs: string[] = [] try { const shards = await fs.promises.readdir(base, { withFileTypes: true }) for (const shard of shards) { if (!shard.isDirectory() || !/^[0-9a-f]{2}$/i.test(shard.name)) continue const shardPath = path.join(base, shard.name) const ids = await fs.promises.readdir(shardPath, { withFileTypes: true }) for (const entry of ids) { if (!entry.isDirectory()) continue const idAbs = path.join(shardPath, entry.name) let legs: string[] try { legs = await fs.promises.readdir(idAbs) } catch (error: any) { if (error?.code === 'ENOENT') continue throw error } // No metadata content leg → a ghost or scar container → not an // entity. Same test pruneOrphanedEntities() uses, so the two agree // by construction. if (!this.hasMetadataContentLeg(legs)) continue count++ if (sampleDirs.length < SAMPLE_MAX) { sampleDirs.push(idAbs) } } } } catch (error: any) { if (error?.code !== 'ENOENT') throw error } return { count, sampleDirs } } /** * Read one canonical entity's `metadata.json` (or `.json.gz`) directly with * fs — NO guarded accessors. Used only by the init-time count recovery, * where `getNounMetadata`'s `ensureInitialized()` would re-enter `init()`. * @param entityDir - Absolute `entities///` directory. * @returns The parsed metadata, or null when absent/unreadable. */ private async readEntityMetadataRaw(entityDir: string): Promise { const base = path.join(entityDir, 'metadata.json') try { return JSON.parse(await fs.promises.readFile(base, 'utf-8')) } catch { // fall through to the compressed variant } try { const gz = await fs.promises.readFile(`${base}.gz`) return JSON.parse(zlib.gunzipSync(gz).toString('utf-8')) } catch { return null } } /** * Read one canonical noun's `vectors.json` (or `.json.gz`) directly with fs * — the vector-side mirror of {@link readEntityMetadataRaw}, same * reentrancy reason (bypasses `getNoun()`'s `ensureInitialized()`). * @param entityDir - Absolute `entities/nouns//` directory. * @returns The parsed vector record, or null when absent/unreadable. */ private async readEntityVectorRaw(entityDir: string): Promise { const base = path.join(entityDir, 'vectors.json') try { return JSON.parse(await fs.promises.readFile(base, 'utf-8')) } catch { // fall through to the compressed variant } try { const gz = await fs.promises.readFile(`${base}.gz`) return JSON.parse(zlib.gunzipSync(gz).toString('utf-8')) } catch { return null } } /** * Count canonical nouns holding a REAL (non-empty, non-zero-norm) vector — * the vectored-noun ledger scalar. UNLIKE {@link scanCanonicalEntities}, * presence cannot be decided from the id-directory listing alone: a * deferred-embed noun's `vectors.json` EXISTS (written at `add()` time * with `vector: []`) until its embed LANDS, so this walk reads every * noun's `vectors.json` CONTENT — O(nouns) reads, not O(ids) listing. * ZERO-NORM LAW: a real all-zero vector is not a vector — it never counts * here either (Brainy's write paths normalize an explicit zero-norm * vector to `[]` at write time, but a store created before that fix may * still carry legacy all-zero rows on disk; this derivation must agree * with the live ledger's definition of "vectored" regardless of when the * row was written). Used ONLY for a one-time legacy-counts.json derivation * or a lost/corrupted counts.json recovery; the result is persisted so * this scan never repeats. */ private async scanVectoredNounCount(): Promise { const base = path.join(this.rootDir, 'entities', 'nouns') let vectored = 0 try { const shards = await fs.promises.readdir(base, { withFileTypes: true }) for (const shard of shards) { if (!shard.isDirectory() || !/^[0-9a-f]{2}$/i.test(shard.name)) continue const shardPath = path.join(base, shard.name) const ids = await fs.promises.readdir(shardPath, { withFileTypes: true }) for (const entry of ids) { if (!entry.isDirectory()) continue const record = await this.readEntityVectorRaw(path.join(shardPath, entry.name)) if ( record && Array.isArray(record.vector) && record.vector.length > 0 && !isZeroNormVector(record.vector) ) { vectored++ } } } } catch (error: any) { if (error?.code !== 'ENOENT') throw error } return vectored } /** * Persist counts to filesystem storage */ protected async persistCounts(): Promise { if (!this.countsFilePath) return try { const counts = { entityCounts: Object.fromEntries(this.entityCounts), verbCounts: Object.fromEntries(this.verbCounts), totalNounCount: this.totalNounCount, totalVerbCount: this.totalVerbCount, // ALL-visibility ledger scalars (+ the suspect flag) — absent in files // written before the ledger existed; initializeCounts() derives them once. totalNounCountAll: this.totalNounCountAll, totalVerbCountAll: this.totalVerbCountAll, // Vectored-noun ledger scalar — absent in files written before it // existed; initializeCounts() derives it once (a content scan, see // scanVectoredNounCount()'s JSDoc). totalVectoredNounCount: this.totalVectoredNounCount, allCountsSuspect: this.allCountsSuspect, // Derivation-rule stamp for the ALL scalars above — 'identity-record' // when they were counted one-per-metadata-content-leg (the honest // rule); omitted (JSON.stringify drops `undefined`) when the current // in-memory scalars came from a legacy container-rule counts.json // that hasn't been through a sanctioned recount yet, so a future load // keeps naming them suspect rather than trusting an unproven value. allCountsDerivedBy: this.allCountsDerivedBy, lastUpdated: new Date().toISOString() } // ATOMIC (temp + rename), never a plain writeFile. A direct write // truncates the file first, so every persist opened a window — measured // at roughly 750ms after a flush or close on a real store — in which a // concurrent reader saw counts.json EMPTY. An empty file is unparseable, // and an unparseable ledger sends the next open down the full-rescan // path: the cheapest file in the store was costing the most expensive // recovery. The rename is atomic, so a reader sees the old ledger or the // new one, never neither. await this.writeFileAtomic(this.countsFilePath, JSON.stringify(counts, null, 2)) } catch (error) { console.error('Error persisting counts:', error) } } // ============================================= // Intelligent Directory Sharding // ============================================= /** * Migrate a single file atomically */ private async migrateFile( fileInfo: { oldPath: string; id: string; type: 'noun' | 'verb' }, fromDepth: number, toDepth: number ): Promise { const baseDir = fileInfo.type === 'noun' ? this.nounsDir : this.verbsDir // Calculate old path (already known) const oldPath = fileInfo.oldPath // Calculate new path using target depth const shard = fileInfo.id.substring(0, 2).toLowerCase() const newPath = path.join(baseDir, shard, `${fileInfo.id}.json`) // Check if file already exists at new location if (await this.fileExists(newPath)) { // File already migrated or duplicate - skip return } // Atomic rename/move await fs.promises.rename(oldPath, newPath) } /** * Whether this store already holds canonical 8.0 entities. * * 8.0 writes nouns to `entities/nouns///vectors.json` (see * `getNounVectorPath`). This checks that canonical shard tree — the one the * DB actually reads and writes (the same `entities/nouns/` shards * `getNounsWithPagination` walks) — so the new-vs-established boot log is * truthful. (The 7.x hnsw sharding probe this replaced inspected a directory * the 8.0 write path never populated, so it mislabeled every established * store "New installation" on every boot; that probe and its depth-migration * machinery are removed.) * * @returns true if at least one 2-hex shard directory (00–ff) exists under * `entities/nouns/`, i.e. the store has previously persisted entities. */ private async hasCanonicalEntities(): Promise { const canonicalNounsDir = path.join(this.rootDir, 'entities', 'nouns') try { const entries = await fs.promises.readdir(canonicalNounsDir, { withFileTypes: true }) // A populated store has ≥1 two-hex shard dir (00–ff). The vestigial // `hnsw` subdir is 4 chars and correctly excluded by the hex test. return entries.some( (e: any) => e.isDirectory() && /^[0-9a-f]{2}$/i.test(e.name) ) } catch { // Directory absent (fresh store) or unreadable → no persisted entities. return false } } /** * Check if a file exists (handles both sharded and non-sharded) */ private async fileExists(filePath: string): Promise { try { await fs.promises.access(filePath, fs.constants.F_OK) return true } catch { return false } } // ============================================= // HNSW Index Persistence // ============================================= /** * Get vector for a noun * Uses BaseStorage's getNoun (type-first paths) */ public async getNounVector(id: string): Promise { const noun = await this.getNoun(id) return noun ? noun.vector : null } /** * Save HNSW graph data for a noun * * Uses BaseStorage's getNoun/saveNoun (type-first paths) * CRITICAL: Preserves mutex locking to prevent read-modify-write races */ public async saveVectorIndexData(nounId: string, hnswData: { level: number connections: Record }): Promise { const lockKey = `hnsw/${nounId}` // CRITICAL FIX: Mutex lock to prevent read-modify-write races // Problem: Without mutex, concurrent operations can: // 1. Thread A reads noun (connections: [1,2,3]) // 2. Thread B reads noun (connections: [1,2,3]) // 3. Thread A adds connection 4, writes [1,2,3,4] // 4. Thread B adds connection 5, writes [1,2,3,5] ← Connection 4 LOST! // Solution: Mutex serializes operations per entity (like Memory/OPFS adapters) // Production scale: Prevents corruption at 1000+ concurrent operations // Wait for any pending operations on this entity while (this.hnswLocks.has(lockKey)) { await this.hnswLocks.get(lockKey) } // Acquire lock let releaseLock!: () => void const lockPromise = new Promise(resolve => { releaseLock = resolve }) this.hnswLocks.set(lockKey, lockPromise) try { // Use BaseStorage's getNoun (type-first paths) // Read existing noun data (if exists) const existingNoun = await this.getNoun(nounId) if (!existingNoun) { // Noun doesn't exist - cannot update HNSW data for non-existent noun throw new Error(`Cannot save HNSW data: noun ${nounId} not found`) } // Convert connections from Record to Map format for storage const connectionsMap = new Map>() for (const [level, nodeIds] of Object.entries(hnswData.connections)) { connectionsMap.set(Number(level), new Set(nodeIds)) } // Preserve id and vector, update only HNSW graph metadata const updatedNoun: HNSWNoun = { ...existingNoun, level: hnswData.level, connections: connectionsMap } // Use BaseStorage's saveNoun (type-first paths, atomic write) await this.saveNoun(updatedNoun) } finally { // Release lock (ALWAYS runs, even if error thrown) this.hnswLocks.delete(lockKey) releaseLock() } } /** * Get HNSW graph data for a noun * Uses BaseStorage's getNoun (type-first paths) */ public async getVectorIndexData(nounId: string): Promise<{ level: number connections: Record } | null> { const noun = await this.getNoun(nounId) if (!noun) { return null } // Convert connections from Map to Record format const connectionsRecord: Record = {} if (noun.connections) { for (const [level, nodeIds] of noun.connections.entries()) { connectionsRecord[String(level)] = Array.from(nodeIds) } } return { level: noun.level || 0, connections: connectionsRecord } } /** * Save HNSW system data (entry point, max level) * * CRITICAL FIX: Mutex lock + atomic write to prevent race conditions */ public async saveHNSWSystem(systemData: { entryPointId: string | null maxLevel: number }): Promise { await this.ensureInitialized() const lockKey = 'hnsw/system' // CRITICAL FIX: Mutex lock to serialize system updates // System data (entry point, max level) updated frequently during HNSW construction // Without mutex, concurrent updates can lose data (same as entity-level problem) // Wait for any pending system updates while (this.hnswLocks.has(lockKey)) { await this.hnswLocks.get(lockKey) } // Acquire lock let releaseLock!: () => void const lockPromise = new Promise(resolve => { releaseLock = resolve }) this.hnswLocks.set(lockKey, lockPromise) try { const filePath = path.join(this.systemDir, 'hnsw-system.json') const tempPath = `${filePath}.tmp.${Date.now()}.${Math.random().toString(36).substring(2)}` try { // Write to temp file await this.ensureDirectoryExists(path.dirname(tempPath)) await fs.promises.writeFile(tempPath, JSON.stringify(systemData, null, 2)) // Atomic rename temp → final (POSIX atomicity guarantee) await fs.promises.rename(tempPath, filePath) } catch (error: any) { // Clean up temp file on any error try { await fs.promises.unlink(tempPath) } catch (cleanupError) { // Ignore cleanup errors } throw error } } finally { // Release lock this.hnswLocks.delete(lockKey) releaseLock() } } /** * Get HNSW system data */ public async getHNSWSystem(): Promise<{ entryPointId: string | null maxLevel: number } | null> { await this.ensureInitialized() const filePath = path.join(this.systemDir, 'hnsw-system.json') try { const data = await fs.promises.readFile(filePath, 'utf-8') return JSON.parse(data) } catch (error: any) { if (error.code !== 'ENOENT') { console.error('Error reading HNSW system data:', error) } return null } } }