/** * Base Storage Adapter * Provides common functionality for all storage adapters, including statistics tracking */ import { StatisticsData, StorageAdapter, HNSWNoun, HNSWVerb, GraphVerb, HNSWNounWithMetadata, HNSWVerbWithMetadata, NounMetadata, VerbMetadata } from '../../coreTypes.js' import { StorageBatchConfig } from '../baseStorage.js' import { extractFieldNamesFromJson, mapToStandardField } from '../../utils/fieldNameTracking.js' import { getGlobalMutex, cleanupMutexes } from '../../utils/mutex.js' // ============================================================================= // Progressive Initialization Types // ============================================================================= /** * Initialization mode for cloud storage adapters. * * Controls how storage adapters initialize during startup, enabling * fast cold starts in serverless environments while maintaining strict * validation for local development. * * * | Mode | Description | Use Case | * |------|-------------|----------| * | `'auto'` | Auto-detect environment (progressive in cloud, strict locally) | Default, zero-config | * | `'progressive'` | Fast init (<200ms), validate lazily on first write | Serverless, Cloud Run | * | `'strict'` | Blocking validation during init (traditional behavior) | Local dev, testing | * * @example * ```typescript * // Zero-config: auto-detects Cloud Run, Lambda, etc. * const storage = new FileSystemStorage({ bucket: 'my-bucket' }) * * // Force progressive mode for all environments * const storage = new FileSystemStorage({ * bucket: 'my-bucket', * initMode: 'progressive' * }) * * // Force strict validation (useful for testing) * const storage = new FileSystemStorage({ * bucket: 'my-bucket', * initMode: 'strict' * }) * ``` */ export type InitMode = 'progressive' | 'strict' | 'auto' /** * Base class for storage adapters that implements statistics tracking */ export abstract class BaseStorageAdapter implements StorageAdapter { // Abstract methods that must be implemented by subclasses abstract init(): Promise abstract saveNoun(noun: HNSWNoun): Promise abstract saveNounMetadata(id: string, metadata: NounMetadata): Promise abstract deleteNounMetadata(id: string): Promise abstract getNoun(id: string): Promise abstract getNounsByNounType(nounType: string): Promise abstract deleteNoun(id: string): Promise abstract saveVerb(verb: HNSWVerb): Promise abstract saveVerbMetadata(id: string, metadata: VerbMetadata): Promise abstract deleteVerbMetadata(id: string): Promise abstract getVerb(id: string): Promise abstract getVerbsBySource(sourceId: string): Promise abstract getVerbsByTarget(targetId: string): Promise abstract getVerbsByType(type: string): Promise abstract deleteVerb(id: string): Promise abstract saveMetadata(id: string, metadata: any): Promise abstract getMetadata(id: string): Promise abstract getNounMetadata(id: string): Promise abstract getVerbMetadata(id: string): Promise // Vector Index Persistence (Brainy 8.0 — algorithm-neutral surface). // Concrete adapters persist the per-entity HNSW graph node (level + connections) // under these methods. Cortex's DiskANN-style native vector index doesn't use // them (it persists its own single mmap'd `.dkann` file under // `_system/vector-index/.dkann` per BRAINY-8.0-RENAME-COORDINATION § B.2). abstract getNounVector(id: string): Promise abstract saveVectorIndexData(nounId: string, data: { level: number connections: Record }): Promise abstract getVectorIndexData(nounId: string): Promise<{ level: number connections: Record } | null> abstract saveHNSWSystem(systemData: { entryPointId: string | null maxLevel: number }): Promise abstract getHNSWSystem(): Promise<{ entryPointId: string | null maxLevel: number } | null> abstract clear(): Promise abstract getStorageStatus(): Promise<{ type: string used: number quota: number | null details?: Record }> // =========================================================================== // Raw binary-blob primitive // =========================================================================== // // A first-class storage primitive for opaque byte payloads that must NOT be // wrapped in a JSON envelope. The JSON object path base64-encodes binary data, // inflating it ~33% and forcing full in-memory materialization on every read. // Blobs side-step that: bytes are written and read verbatim. // // This unblocks zero-copy, mmap-able column-store segments and batch vector // I/O at billion scale. On filesystem-backed adapters `getBinaryBlobPath` // returns a real local path so native code (Rust) can mmap the file directly. // // Blob key → location convention (shared by every adapter so cross-language // and cross-adapter consumers agree on exactly where a blob lives): // // key location // "graph-lsm/source/sstable-123" "/_blobs/graph-lsm/source/sstable-123.bin" // // i.e. the key's "/"-separated segments are nested under a `_blobs/` prefix and // suffixed with `.bin`. Blobs are deliberately NOT branch-scoped (COW): they // are immutable, content-addressed segments managed by their producer, // mirroring how COW metadata under `_cow/` is also kept global. /** * Persist a raw binary blob under `key`, writing the bytes verbatim (no JSON * envelope, no base64). Overwrites any existing blob at the same key. Where the * backend is a real filesystem the write is atomic (temp file + rename) so a * concurrent reader never observes a torn write. * * @param key - Logical blob key. "/"-separated segments map to nested * directories under the adapter's `_blobs/` prefix (e.g. * `"graph-lsm/source/sstable-123"`). * @param data - The exact bytes to store. * @returns Resolves once the blob is durably written. * @example * await storage.saveBinaryBlob('graph-lsm/source/sstable-7', segmentBytes) */ abstract saveBinaryBlob(key: string, data: Buffer): Promise /** * Load the raw bytes previously stored under `key`, or `null` if no blob * exists at that key. The returned buffer is byte-identical to what was passed * to {@link saveBinaryBlob}. * * @param key - The blob key used when saving. * @returns The blob bytes, or `null` if absent. * @example * const bytes = await storage.loadBinaryBlob('graph-lsm/source/sstable-7') * if (bytes) decodeSegment(bytes) */ abstract loadBinaryBlob(key: string): Promise /** * Delete the blob stored under `key`. Missing blobs are ignored (no error) so * delete is idempotent. * * @param key - The blob key to delete. * @returns Resolves once the blob is gone (or was already absent). */ abstract deleteBinaryBlob(key: string): Promise /** * Resolve `key` to a real local filesystem path that native code can `mmap` * directly, or `null` when this backend has no local file for the blob. * * Filesystem-backed adapters return the on-disk path (whether or not the file * currently exists — the caller is expected to write before mapping). Remote * object stores (S3, R2, GCS, Azure), in-memory storage, browser storage * (OPFS), and the read-only historical adapter genuinely have no local path * and therefore return `null`. A `null` here is correct behavior, not a * fallback: callers must use {@link loadBinaryBlob} when no path is available. * * @param key - The blob key. * @returns An absolute local filesystem path, or `null` if none exists. */ abstract getBinaryBlobPath(key: string): string | null /** * Get optimal batch configuration for this storage adapter * Override in subclasses to provide storage-specific optimization * * This method allows each storage adapter to declare its optimal batch behavior * for rate limiting and performance. The configuration is used by addMany(), * relateMany(), and import operations to automatically adapt to storage capabilities. * * @returns Batch configuration optimized for this storage type */ public getBatchConfig(): StorageBatchConfig { // Conservative defaults that work safely across all storage types // Cloud storage adapters should override with higher throughput values // Local storage adapters should override with no delays return { maxBatchSize: 50, batchDelayMs: 100, maxConcurrent: 50, supportsParallelWrites: false, rateLimit: { operationsPerSecond: 100, burstCapacity: 200 } } } // NOTE: getAllNouns and getAllVerbs have been removed to prevent expensive full scans. // Use getNouns() and getVerbs() with pagination instead. /** * Get nouns with pagination and filtering * @param options Pagination and filtering options * @returns Promise that resolves to a paginated result of nouns */ abstract getNouns(options?: { pagination?: { offset?: number limit?: number cursor?: string } filter?: { nounType?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWNounWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string }> /** * Get verbs with pagination and filtering * @param options Pagination and filtering options * @returns Promise that resolves to a paginated result of verbs */ abstract getVerbs(options?: { pagination?: { offset?: number limit?: number cursor?: string } filter?: { verbType?: string | string[] sourceId?: string | string[] targetId?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWVerbWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string }> /** * Get nouns with pagination (internal implementation) * This method should be implemented by storage adapters to support efficient pagination * @param options Pagination options * @returns Promise that resolves to a paginated result of nouns */ getNounsWithPagination?(options: { limit?: number cursor?: string filter?: { nounType?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWNounWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string }> /** * Get verbs with pagination (internal implementation) * This method should be implemented by storage adapters to support efficient pagination * @param options Pagination options * @returns Promise that resolves to a paginated result of verbs */ getVerbsWithPagination?(options: { limit?: number cursor?: string filter?: { verbType?: string | string[] sourceId?: string | string[] targetId?: string | string[] service?: string | string[] metadata?: Record } }): Promise<{ items: HNSWVerbWithMetadata[] totalCount?: number hasMore: boolean nextCursor?: string }> // Statistics cache protected statisticsCache: StatisticsData | null = null // Batch update timer ID protected statisticsBatchUpdateTimerId: NodeJS.Timeout | null = null // Flag to indicate if statistics have been modified since last save protected statisticsModified = false // Time of last statistics flush to storage protected lastStatisticsFlushTime = 0 // Minimum time between statistics flushes (5 seconds) protected readonly MIN_FLUSH_INTERVAL_MS = 5000 // Maximum time to wait before flushing statistics (30 seconds) protected readonly MAX_FLUSH_DELAY_MS = 30000 // Throttling tracking properties protected throttlingDetected = false protected throttlingBackoffMs = 1000 // Start with 1 second protected maxBackoffMs = 30000 // Max 30 seconds protected consecutiveThrottleEvents = 0 protected lastThrottleTime = 0 protected totalThrottleEvents = 0 protected throttleEventsByHour: number[] = new Array(24).fill(0) protected throttleReasons: Record = {} protected lastThrottleHourIndex = -1 // Operation impact tracking protected delayedOperations = 0 protected retriedOperations = 0 protected failedDueToThrottling = 0 protected totalDelayMs = 0 // Service-level throttling protected serviceThrottling: Map = new Map() // ============================================= // Progressive Initialization State // ============================================= // These properties enable fast cold starts in cloud environments // by deferring validation and count loading to background tasks. /** * Initialization mode for this storage adapter. * * - `'auto'` (default): Progressive in cloud environments, strict locally * - `'progressive'`: Always use fast init with lazy validation * - `'strict'`: Always validate during init (traditional behavior) * * @protected */ protected initMode: InitMode = 'auto' /** * Whether the bucket/container has been validated to exist. * * In progressive mode, this starts as `false` and is set to `true` * after the first successful write operation or background validation. * * @protected */ protected bucketValidated = false /** * Error from bucket validation, if any. * * Stored here so subsequent operations can fail fast with the same error. * * @protected */ protected bucketValidationError: Error | null = null /** * Whether counts have been loaded from storage. * * In progressive mode, counts start at 0 and are loaded in the background. * Operations work immediately; counts become accurate after background load. * * @protected */ protected countsLoaded = false /** * Whether all background initialization tasks have completed. * * Useful for tests and diagnostics to ensure full initialization. * * @protected */ protected backgroundTasksComplete = false /** * Promise that resolves when background initialization completes. * * Can be awaited by callers who need to ensure full initialization * before proceeding (e.g., in tests). * * @protected */ protected backgroundInitPromise: Promise | null = null // Statistics-specific methods that must be implemented by subclasses protected abstract saveStatisticsData( statistics: StatisticsData ): Promise protected abstract getStatisticsData(): Promise /** * Save statistics data * @param statistics The statistics data to save */ async saveStatistics(statistics: StatisticsData): Promise { // Update the cache with a deep copy to avoid reference issues this.statisticsCache = { nounCount: { ...statistics.nounCount }, verbCount: { ...statistics.verbCount }, metadataCount: { ...statistics.metadataCount }, hnswIndexSize: statistics.hnswIndexSize, lastUpdated: statistics.lastUpdated, // Include serviceActivity if present ...(statistics.serviceActivity && { serviceActivity: Object.fromEntries( Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}]) ) }), // Include services if present ...(statistics.services && { services: statistics.services.map(s => ({...s})) }) } // Schedule a batch update instead of saving immediately this.scheduleBatchUpdate() } /** * Get statistics data * @returns Promise that resolves to the statistics data */ async getStatistics(): Promise { // If we have cached statistics, return a deep copy if (this.statisticsCache) { return { nounCount: { ...this.statisticsCache.nounCount }, verbCount: { ...this.statisticsCache.verbCount }, metadataCount: { ...this.statisticsCache.metadataCount }, hnswIndexSize: this.statisticsCache.hnswIndexSize, lastUpdated: this.statisticsCache.lastUpdated } } // Otherwise, get from storage const statistics = await this.getStatisticsData() // If we found statistics, update the cache if (statistics) { // Update the cache with a deep copy this.statisticsCache = { nounCount: { ...statistics.nounCount }, verbCount: { ...statistics.verbCount }, metadataCount: { ...statistics.metadataCount }, hnswIndexSize: statistics.hnswIndexSize, lastUpdated: statistics.lastUpdated } } return statistics } /** * Schedule a batch update of statistics */ protected scheduleBatchUpdate(): void { // Mark statistics as modified this.statisticsModified = true // If a timer is already set, don't set another one if (this.statisticsBatchUpdateTimerId !== null) { return } // Calculate time since last flush const now = Date.now() const timeSinceLastFlush = now - this.lastStatisticsFlushTime // If we've recently flushed, wait longer before the next flush const delayMs = timeSinceLastFlush < this.MIN_FLUSH_INTERVAL_MS ? this.MAX_FLUSH_DELAY_MS : this.MIN_FLUSH_INTERVAL_MS // Schedule the batch update this.statisticsBatchUpdateTimerId = setTimeout(() => { this.flushStatistics() }, delayMs) } /** * Flush statistics to storage */ protected async flushStatistics(): Promise { // Clear the timer if (this.statisticsBatchUpdateTimerId !== null) { clearTimeout(this.statisticsBatchUpdateTimerId) this.statisticsBatchUpdateTimerId = null } // If statistics haven't been modified, no need to flush if (!this.statisticsModified || !this.statisticsCache) { return } try { // Save the statistics to storage await this.saveStatisticsData(this.statisticsCache) // Update the last flush time this.lastStatisticsFlushTime = Date.now() // Reset the modified flag this.statisticsModified = false } catch (error) { console.error('Failed to flush statistics data:', error) // Mark as still modified so we'll try again later this.statisticsModified = true // Don't throw the error to avoid disrupting the application } } /** * Increment a statistic counter * @param type The type of statistic to increment ('noun', 'verb', 'metadata') * @param service The service that inserted the data * @param amount The amount to increment by (default: 1) */ async incrementStatistic( type: 'noun' | 'verb' | 'metadata', service: string, amount: number = 1 ): Promise { // Get current statistics from cache or storage let statistics = this.statisticsCache if (!statistics) { statistics = await this.getStatisticsData() if (!statistics) { statistics = this.createDefaultStatistics() } // Update the cache this.statisticsCache = { nounCount: { ...statistics.nounCount }, verbCount: { ...statistics.verbCount }, metadataCount: { ...statistics.metadataCount }, hnswIndexSize: statistics.hnswIndexSize, lastUpdated: statistics.lastUpdated, // Include serviceActivity if present ...(statistics.serviceActivity && { serviceActivity: Object.fromEntries( Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}]) ) }), // Include services if present ...(statistics.services && { services: statistics.services.map(s => ({...s})) }) } } // Increment the appropriate counter const counterMap = { noun: this.statisticsCache!.nounCount, verb: this.statisticsCache!.verbCount, metadata: this.statisticsCache!.metadataCount } const counter = counterMap[type] counter[service] = (counter[service] || 0) + amount // Track service activity this.trackServiceActivity(service, 'add') // Update timestamp this.statisticsCache!.lastUpdated = new Date().toISOString() // Schedule a batch update instead of saving immediately this.scheduleBatchUpdate() } /** * Track service activity (first/last activity, operation counts) * @param service The service name * @param operation The operation type */ protected trackServiceActivity( service: string, operation: 'add' | 'update' | 'delete' ): void { if (!this.statisticsCache) { return } // Initialize serviceActivity if it doesn't exist if (!this.statisticsCache.serviceActivity) { this.statisticsCache.serviceActivity = {} } const now = new Date().toISOString() const activity = this.statisticsCache.serviceActivity[service] if (!activity) { // First activity for this service this.statisticsCache.serviceActivity[service] = { firstActivity: now, lastActivity: now, totalOperations: 1 } } else { // Update existing activity activity.lastActivity = now activity.totalOperations++ } } /** * Decrement a statistic counter * @param type The type of statistic to decrement ('noun', 'verb', 'metadata') * @param service The service that inserted the data * @param amount The amount to decrement by (default: 1) */ async decrementStatistic( type: 'noun' | 'verb' | 'metadata', service: string, amount: number = 1 ): Promise { // Get current statistics from cache or storage let statistics = this.statisticsCache if (!statistics) { statistics = await this.getStatisticsData() if (!statistics) { statistics = this.createDefaultStatistics() } // Update the cache this.statisticsCache = { nounCount: { ...statistics.nounCount }, verbCount: { ...statistics.verbCount }, metadataCount: { ...statistics.metadataCount }, hnswIndexSize: statistics.hnswIndexSize, lastUpdated: statistics.lastUpdated, // Include serviceActivity if present ...(statistics.serviceActivity && { serviceActivity: Object.fromEntries( Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}]) ) }), // Include services if present ...(statistics.services && { services: statistics.services.map(s => ({...s})) }) } } // Decrement the appropriate counter const counterMap = { noun: this.statisticsCache!.nounCount, verb: this.statisticsCache!.verbCount, metadata: this.statisticsCache!.metadataCount } const counter = counterMap[type] counter[service] = Math.max(0, (counter[service] || 0) - amount) // Track service activity this.trackServiceActivity(service, 'delete') // Update timestamp this.statisticsCache!.lastUpdated = new Date().toISOString() // Schedule a batch update instead of saving immediately this.scheduleBatchUpdate() } /** * Update the HNSW index size statistic * @param size The new size of the HNSW index */ async updateHnswIndexSize(size: number): Promise { // Get current statistics from cache or storage let statistics = this.statisticsCache if (!statistics) { statistics = await this.getStatisticsData() if (!statistics) { statistics = this.createDefaultStatistics() } // Update the cache this.statisticsCache = { nounCount: { ...statistics.nounCount }, verbCount: { ...statistics.verbCount }, metadataCount: { ...statistics.metadataCount }, hnswIndexSize: statistics.hnswIndexSize, lastUpdated: statistics.lastUpdated, // Include serviceActivity if present ...(statistics.serviceActivity && { serviceActivity: Object.fromEntries( Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}]) ) }), // Include services if present ...(statistics.services && { services: statistics.services.map(s => ({...s})) }) } } // Update HNSW index size this.statisticsCache!.hnswIndexSize = size // Update timestamp this.statisticsCache!.lastUpdated = new Date().toISOString() // Schedule a batch update instead of saving immediately this.scheduleBatchUpdate() } /** * Force an immediate flush of statistics to storage * This ensures that any pending statistics updates are written to persistent storage */ async flushStatisticsToStorage(): Promise { // If there are no statistics in cache or they haven't been modified, nothing to flush if (!this.statisticsCache || !this.statisticsModified) { return } // Call the protected flushStatistics method to immediately write to storage await this.flushStatistics() } /** * Track field names from a JSON document * @param jsonDocument The JSON document to extract field names from * @param service The service that inserted the data */ async trackFieldNames(jsonDocument: any, service: string): Promise { // Skip if not a JSON object if (typeof jsonDocument !== 'object' || jsonDocument === null || Array.isArray(jsonDocument)) { return } // Get current statistics from cache or storage let statistics = this.statisticsCache if (!statistics) { statistics = await this.getStatisticsData() if (!statistics) { statistics = this.createDefaultStatistics() } // Update the cache this.statisticsCache = { ...statistics, nounCount: { ...statistics.nounCount }, verbCount: { ...statistics.verbCount }, metadataCount: { ...statistics.metadataCount }, fieldNames: { ...statistics.fieldNames }, standardFieldMappings: { ...statistics.standardFieldMappings } } } // Ensure fieldNames exists if (!this.statisticsCache!.fieldNames) { this.statisticsCache!.fieldNames = {} } // Ensure standardFieldMappings exists if (!this.statisticsCache!.standardFieldMappings) { this.statisticsCache!.standardFieldMappings = {} } // Extract field names from the JSON document const fieldNames = extractFieldNamesFromJson(jsonDocument) // Initialize service entry if it doesn't exist if (!this.statisticsCache!.fieldNames[service]) { this.statisticsCache!.fieldNames[service] = [] } // Add new field names to the service's list for (const fieldName of fieldNames) { if (!this.statisticsCache!.fieldNames[service].includes(fieldName)) { this.statisticsCache!.fieldNames[service].push(fieldName) } // Map to standard field if possible const standardField = mapToStandardField(fieldName) if (standardField) { // Initialize standard field entry if it doesn't exist if (!this.statisticsCache!.standardFieldMappings[standardField]) { this.statisticsCache!.standardFieldMappings[standardField] = {} } // Initialize service entry if it doesn't exist if (!this.statisticsCache!.standardFieldMappings[standardField][service]) { this.statisticsCache!.standardFieldMappings[standardField][service] = [] } // Add field name to standard field mapping if not already there if (!this.statisticsCache!.standardFieldMappings[standardField][service].includes(fieldName)) { this.statisticsCache!.standardFieldMappings[standardField][service].push(fieldName) } } } // Update timestamp this.statisticsCache!.lastUpdated = new Date().toISOString() // Schedule a batch update this.statisticsModified = true this.scheduleBatchUpdate() } /** * Get available field names by service * @returns Record of field names by service */ async getAvailableFieldNames(): Promise> { // Get current statistics from cache or storage let statistics = this.statisticsCache if (!statistics) { statistics = await this.getStatisticsData() if (!statistics) { return {} } } // Return field names by service return statistics.fieldNames || {} } /** * Get standard field mappings * @returns Record of standard field mappings */ async getStandardFieldMappings(): Promise>> { // Get current statistics from cache or storage let statistics = this.statisticsCache if (!statistics) { statistics = await this.getStatisticsData() if (!statistics) { return {} } } // Return standard field mappings return statistics.standardFieldMappings || {} } /** * Create default statistics data * @returns Default statistics data */ protected createDefaultStatistics(): StatisticsData { return { nounCount: {}, verbCount: {}, metadataCount: {}, hnswIndexSize: 0, fieldNames: {}, standardFieldMappings: {}, lastUpdated: new Date().toISOString() } } /** * Detect if an error is a throttling error * Override this method in specific adapters for custom detection */ protected isThrottlingError(error: any): boolean { const statusCode = error.$metadata?.httpStatusCode || error.statusCode || error.code const message = error.message?.toLowerCase() || '' return ( statusCode === 429 || // Too Many Requests statusCode === 503 || // Service Unavailable / Slow Down statusCode === 'ECONNRESET' || // Connection reset statusCode === 'ETIMEDOUT' || // Timeout message.includes('throttl') || message.includes('slow down') || message.includes('rate limit') || message.includes('too many requests') || message.includes('quota exceeded') ) } /** * Track a throttling event * @param error The error that caused throttling * @param service Optional service that was throttled */ protected trackThrottlingEvent(error: any, service?: string): void { this.throttlingDetected = true this.consecutiveThrottleEvents++ this.lastThrottleTime = Date.now() this.totalThrottleEvents++ // Track by hour const hourIndex = new Date().getHours() if (hourIndex !== this.lastThrottleHourIndex) { // Reset hour tracking if we've moved to a new hour this.throttleEventsByHour = new Array(24).fill(0) this.lastThrottleHourIndex = hourIndex } this.throttleEventsByHour[hourIndex]++ // Track throttle reason const reason = this.getThrottleReason(error) this.throttleReasons[reason] = (this.throttleReasons[reason] || 0) + 1 // Track service-level throttling if (service) { const serviceInfo = this.serviceThrottling.get(service) || { throttleCount: 0, lastThrottle: 0, status: 'normal' as const } serviceInfo.throttleCount++ serviceInfo.lastThrottle = Date.now() serviceInfo.status = 'throttled' this.serviceThrottling.set(service, serviceInfo) } // Exponential backoff this.throttlingBackoffMs = Math.min( this.throttlingBackoffMs * 2, this.maxBackoffMs ) } /** * Get the reason for throttling from an error */ protected getThrottleReason(error: any): string { const statusCode = error.$metadata?.httpStatusCode || error.statusCode || error.code if (statusCode === 429) return '429_TooManyRequests' if (statusCode === 503) return '503_ServiceUnavailable' if (statusCode === 'ECONNRESET') return 'ConnectionReset' if (statusCode === 'ETIMEDOUT') return 'Timeout' const message = error.message?.toLowerCase() || '' if (message.includes('throttl')) return 'Throttled' if (message.includes('slow down')) return 'SlowDown' if (message.includes('rate limit')) return 'RateLimit' if (message.includes('quota exceeded')) return 'QuotaExceeded' return 'Unknown' } /** * Clear throttling state after successful operations */ protected clearThrottlingState(): void { if (this.consecutiveThrottleEvents > 0) { this.consecutiveThrottleEvents = 0 this.throttlingBackoffMs = 1000 // Reset to initial backoff if (this.throttlingDetected) { this.throttlingDetected = false // Update service statuses for (const [service, info] of this.serviceThrottling) { if (info.status === 'throttled') { info.status = 'recovering' } else if (info.status === 'recovering') { const timeSinceThrottle = Date.now() - info.lastThrottle if (timeSinceThrottle > 60000) { // 1 minute recovery period info.status = 'normal' } } } } } } /** * Handle throttling by implementing exponential backoff * @param error The error that triggered throttling * @param service Optional service that was throttled */ async handleThrottling(error: any, service?: string): Promise { if (this.isThrottlingError(error)) { this.trackThrottlingEvent(error, service) // Add delay for retry const delayMs = this.throttlingBackoffMs this.totalDelayMs += delayMs this.delayedOperations++ await new Promise(resolve => setTimeout(resolve, delayMs)) } else { // Clear throttling state on non-throttling errors this.clearThrottlingState() } } /** * Track a retried operation */ protected trackRetriedOperation(): void { this.retriedOperations++ } /** * Track an operation that failed due to throttling */ protected trackFailedDueToThrottling(): void { this.failedDueToThrottling++ } /** * Get current throttling metrics */ protected getThrottlingMetrics(): StatisticsData['throttlingMetrics'] { const averageDelayMs = this.delayedOperations > 0 ? this.totalDelayMs / this.delayedOperations : 0 // Convert service throttling map to record const serviceThrottlingRecord: Record = {} for (const [service, info] of this.serviceThrottling) { serviceThrottlingRecord[service] = { throttleCount: info.throttleCount, lastThrottle: new Date(info.lastThrottle).toISOString(), status: info.status } } return { storage: { currentlyThrottled: this.throttlingDetected, lastThrottleTime: this.lastThrottleTime > 0 ? new Date(this.lastThrottleTime).toISOString() : undefined, consecutiveThrottleEvents: this.consecutiveThrottleEvents, currentBackoffMs: this.throttlingBackoffMs, totalThrottleEvents: this.totalThrottleEvents, throttleEventsByHour: [...this.throttleEventsByHour], throttleReasons: { ...this.throttleReasons } }, operationImpact: { delayedOperations: this.delayedOperations, retriedOperations: this.retriedOperations, failedDueToThrottling: this.failedDueToThrottling, averageDelayMs, totalDelayMs: this.totalDelayMs }, serviceThrottling: Object.keys(serviceThrottlingRecord).length > 0 ? serviceThrottlingRecord : undefined } } /** * Include throttling metrics in statistics */ async getStatisticsWithThrottling(): Promise { const stats = await this.getStatistics() if (stats) { stats.throttlingMetrics = this.getThrottlingMetrics() } return stats } // ============================================= // Universal O(1) Count Management // ============================================= // Universal count tracking - O(1) operations protected totalNounCount = 0 protected totalVerbCount = 0 protected entityCounts: Map = new Map() // type -> count protected verbCounts: Map = new Map() // verb type -> count protected countCache: Map = new Map() protected readonly COUNT_CACHE_TTL = 60000 // 1 minute cache TTL // ============================================= // Smart Count Batching // ============================================= // Count batching state - mirrors statistics batching pattern protected pendingCountPersist = false // Counts changed since last persist? protected lastCountPersistTime = 0 // Timestamp of last persist protected scheduledCountPersistTimeout: NodeJS.Timeout | null = null // Scheduled persist timer protected pendingCountOperations = 0 // Operations since last persist // Batching configuration (overridable by subclasses for custom strategies) protected countPersistBatchSize = 10 // Operations before forcing persist (cloud storage) protected countPersistInterval = 5000 // Milliseconds before forcing persist (cloud storage) /** * Get total noun count - O(1) operation * @returns Promise that resolves to the total number of nouns */ async getNounCount(): Promise { return this.totalNounCount } /** * Get total verb count - O(1) operation * @returns Promise that resolves to the total number of verbs */ async getVerbCount(): Promise { return this.totalVerbCount } /** * Increment count for entity type - O(1) operation * Protected by storage-specific mechanisms (mutex, distributed consensus, etc.) * @param type The entity type */ protected incrementEntityCount(type: string): void { // For distributed scenarios, this is aggregated across shards // For single-node, this is protected by storage-specific locking this.entityCounts.set(type, (this.entityCounts.get(type) || 0) + 1) this.totalNounCount++ // Update cache this.countCache.set('nouns_count', { count: this.totalNounCount, timestamp: Date.now() }) } /** * Thread-safe increment for concurrent scenarios * Uses mutex for single-node, distributed consensus for multi-node */ protected async incrementEntityCountSafe(type: string): Promise { // Single-node mutex protection (distributed mode handled by coordinator) const mutex = getGlobalMutex() await mutex.runExclusive(`count-entity-${type}`, async () => { this.incrementEntityCount(type) // Smart batching: Adapts to storage type // - Cloud storage (GCS, S3): Batches 10 ops OR 5 seconds // - Local storage (File, Memory): Persists immediately await this.scheduleCountPersist() }) } /** * Decrement count for entity type - O(1) operation * @param type The entity type */ protected decrementEntityCount(type: string): void { const current = this.entityCounts.get(type) || 0 if (current > 1) { this.entityCounts.set(type, current - 1) } else { this.entityCounts.delete(type) } if (this.totalNounCount > 0) { this.totalNounCount-- } // Update cache this.countCache.set('nouns_count', { count: this.totalNounCount, timestamp: Date.now() }) } /** * Thread-safe decrement for concurrent scenarios */ protected async decrementEntityCountSafe(type: string): Promise { const mutex = getGlobalMutex() await mutex.runExclusive(`count-entity-${type}`, async () => { this.decrementEntityCount(type) // Smart batching: Adapts to storage type await this.scheduleCountPersist() }) } /** * Increment verb count - O(1) operation (now synchronous) * Protected by storage-specific mechanisms (mutex, distributed consensus, etc.) * @param type The verb type */ protected incrementVerbCount(type: string): void { this.verbCounts.set(type, (this.verbCounts.get(type) || 0) + 1) this.totalVerbCount++ // Update cache this.countCache.set('verbs_count', { count: this.totalVerbCount, timestamp: Date.now() }) } /** * Thread-safe increment for verb counts * Uses mutex for single-node, distributed consensus for multi-node * @param type The verb type */ protected async incrementVerbCountSafe(type: string): Promise { const mutex = getGlobalMutex() await mutex.runExclusive(`count-verb-${type}`, async () => { this.incrementVerbCount(type) // Smart batching: Adapts to storage type await this.scheduleCountPersist() }) } /** * Decrement verb count - O(1) operation (now synchronous) * @param type The verb type */ protected decrementVerbCount(type: string): void { const current = this.verbCounts.get(type) || 0 if (current > 1) { this.verbCounts.set(type, current - 1) } else { this.verbCounts.delete(type) } if (this.totalVerbCount > 0) { this.totalVerbCount-- } // Update cache this.countCache.set('verbs_count', { count: this.totalVerbCount, timestamp: Date.now() }) } /** * Thread-safe decrement for verb counts * @param type The verb type */ protected async decrementVerbCountSafe(type: string): Promise { const mutex = getGlobalMutex() await mutex.runExclusive(`count-verb-${type}`, async () => { this.decrementVerbCount(type) // Smart batching: Adapts to storage type await this.scheduleCountPersist() }) } // ============================================= // Smart Batching Methods // ============================================= /** * Detect if this storage adapter uses cloud storage (network I/O) * Cloud storage benefits from batching; local storage does not. * * Override this method in subclasses for accurate detection. * Default implementation checks storage type from getStorageStatus(). * * @returns true if cloud storage (GCS, S3, R2), false if local (File, Memory) */ protected isCloudStorage(): boolean { // Default: assume local storage (conservative, prefers reliability over performance) // Subclasses should override this for accurate detection return false } // ============================================= // Progressive Initialization Methods // ============================================= /** * Detect if running in a cloud/serverless environment. * * Cloud environments benefit from progressive initialization to minimize * cold start times. This method checks for common environment variables * set by cloud providers. * * | Provider | Environment Variable | * |----------|---------------------| * | Cloud Run | `K_SERVICE`, `K_REVISION` | * | Lambda | `AWS_LAMBDA_FUNCTION_NAME` | * | Cloud Functions | `FUNCTIONS_TARGET` | * | Azure Functions | `AZURE_FUNCTIONS_ENVIRONMENT` | * * @returns `true` if running in a detected cloud environment * @protected */ protected detectCloudEnvironment(): boolean { return !!( process.env.K_SERVICE || // Google Cloud Run process.env.K_REVISION || // Google Cloud Run (revision) process.env.AWS_LAMBDA_FUNCTION_NAME || // AWS Lambda process.env.FUNCTIONS_TARGET || // Google Cloud Functions process.env.AZURE_FUNCTIONS_ENVIRONMENT // Azure Functions ) } /** * Determine the effective initialization mode. * * Resolves `'auto'` to either `'progressive'` or `'strict'` based on * the detected environment. * * @returns The resolved init mode (`'progressive'` or `'strict'`) * @protected */ protected resolveInitMode(): 'progressive' | 'strict' { if (this.initMode === 'auto') { return this.detectCloudEnvironment() ? 'progressive' : 'strict' } return this.initMode } /** * Schedule background initialization tasks. * * This method is called in progressive mode after the adapter is marked * as initialized. It schedules non-blocking tasks to: * 1. Validate bucket/container existence * 2. Load counts from storage * * Override in subclasses to add storage-specific background tasks. * Always call `super.scheduleBackgroundInit()` first. * * @protected */ protected scheduleBackgroundInit(): void { // Create a promise that tracks all background work this.backgroundInitPromise = this.runBackgroundInit() // Don't await - let it run in background this.backgroundInitPromise .then(() => { this.backgroundTasksComplete = true }) .catch((error) => { // Log but don't throw - progressive mode is best-effort console.error('[Progressive Init] Background initialization failed:', error) this.backgroundTasksComplete = true // Mark complete even on error }) } /** * Run background initialization tasks. * * Override in subclasses to add storage-specific tasks. * The default implementation does nothing (for local storage adapters). * * @protected */ protected async runBackgroundInit(): Promise { // Default implementation: nothing to do for local storage // Cloud storage adapters override this to: // 1. Validate bucket/container // 2. Load counts from storage } /** * Wait for background initialization to complete. * * Useful in tests or when you need to ensure all background tasks * have finished before proceeding. * * @example * ```typescript * const storage = new FileSystemStorage({ bucket: 'my-bucket' }) * await storage.init() * * // Wait for background validation and count loading * await storage.awaitBackgroundInit() * * // Now counts are accurate and bucket is validated * const count = await storage.getNounCount() * ``` * * @public */ public async awaitBackgroundInit(): Promise { if (this.backgroundInitPromise) { await this.backgroundInitPromise } } /** * Check if background initialization has completed. * * @returns `true` if all background tasks are complete * @public */ public isBackgroundInitComplete(): boolean { return this.backgroundTasksComplete } /** * Ensure bucket/container is validated before write operations. * * In progressive mode, bucket validation is deferred until the first * write operation. This method performs lazy validation and caches * the result (or error) for subsequent calls. * * Override in subclasses to implement storage-specific validation. * The base implementation does nothing (for local storage adapters). * * @throws Error if bucket validation fails * @protected */ protected async ensureValidatedForWrite(): Promise { // If already validated, nothing to do if (this.bucketValidated) { return } // If we have a cached validation error, throw it if (this.bucketValidationError) { throw this.bucketValidationError } // Default implementation: no validation needed (local storage) // Cloud storage adapters override this to validate bucket/container this.bucketValidated = true } /** * Schedule a smart batched persist operation. * * Strategy: * - Local Storage: Persist immediately (fast, no network latency) * - Cloud Storage: Batch persist (10 ops OR 5 seconds, whichever first) * * This mirrors the statistics batching pattern for consistency. */ protected async scheduleCountPersist(): Promise { // Mark counts as pending persist this.pendingCountPersist = true this.pendingCountOperations++ // Local storage: persist immediately (fast enough, no benefit from batching) if (!this.isCloudStorage()) { await this.flushCounts() return } // Cloud storage: use smart batching // Persist if we've hit the batch size threshold if (this.pendingCountOperations >= this.countPersistBatchSize) { await this.flushCounts() return } // Otherwise, schedule a time-based persist if not already scheduled if (!this.scheduledCountPersistTimeout) { this.scheduledCountPersistTimeout = setTimeout(() => { this.flushCounts().catch(error => { console.error('Failed to flush counts on timer:', error) }) }, this.countPersistInterval) } } /** * Flush counts immediately to storage. * * Used for: * - Graceful shutdown (SIGTERM handler) * - Forced persist (batch threshold reached) * - Local storage immediate persist * * This is the public API that shutdown hooks can call. */ async flushCounts(): Promise { // Clear any scheduled persist if (this.scheduledCountPersistTimeout) { clearTimeout(this.scheduledCountPersistTimeout) this.scheduledCountPersistTimeout = null } // Nothing to flush? if (!this.pendingCountPersist) { return } try { // Persist to storage (implemented by subclass) await this.persistCounts() // Update state this.lastCountPersistTime = Date.now() this.pendingCountPersist = false this.pendingCountOperations = 0 } catch (error) { console.error('❌ CRITICAL: Failed to flush counts to storage:', error) // Keep pending flag set so we retry on next operation throw error } } /** * Initialize counts from storage - must be implemented by each adapter * @protected */ protected abstract initializeCounts(): Promise /** * Persist counts to storage - must be implemented by each adapter * @protected */ protected abstract persistCounts(): Promise }