/** * Cloudflare R2 Storage Adapter (Dedicated) * Optimized specifically for Cloudflare R2 with all latest features * * R2-Specific Optimizations: * - Zero egress fees (aggressive caching) * - Cloudflare global network (edge-aware routing) * - Workers integration (optional edge compute) * - High-volume mode for bulk operations * - Smart batching and backpressure * * Based on latest GCS and S3 implementations with R2-specific enhancements */ import { GraphVerb, HNSWNoun, HNSWVerb, NounMetadata, VerbMetadata, HNSWNounWithMetadata, HNSWVerbWithMetadata, StatisticsData, NounType } from '../../coreTypes.js' import { BaseStorage, StorageBatchConfig, NOUNS_DIR, VERBS_DIR, METADATA_DIR, INDEX_DIR, SYSTEM_DIR, STATISTICS_KEY, getDirectoryPath } from '../baseStorage.js' import { BrainyError } from '../../errors/brainyError.js' import { CacheManager } from '../cacheManager.js' import { createModuleLogger, prodLog } from '../../utils/logger.js' import { getGlobalSocketManager } from '../../utils/adaptiveSocketManager.js' import { getGlobalBackpressure } from '../../utils/adaptiveBackpressure.js' import { getWriteBuffer, WriteBuffer } from '../../utils/writeBuffer.js' import { getCoalescer, RequestCoalescer } from '../../utils/requestCoalescer.js' import { getShardIdFromUuid, getAllShardIds, getShardIdByIndex, TOTAL_SHARDS } from '../sharding.js' // Type aliases for better readability type HNSWNode = HNSWNoun type Edge = HNSWVerb // S3 client types - R2 uses S3-compatible API type S3Client = any type S3Command = any // R2 API limits (same as S3) const MAX_R2_PAGE_SIZE = 1000 /** * Dedicated Cloudflare R2 storage adapter * Optimized for R2's unique characteristics and global edge network * * v5.4.0: Type-aware storage now built into BaseStorage * - Removed 10 *_internal method overrides (now inherit from BaseStorage's type-first implementation) * - Removed getNounsWithPagination override * - 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 * * Configuration: * ```typescript * const r2Storage = new R2Storage({ * bucketName: 'my-brainy-data', * accountId: 'YOUR_CLOUDFLARE_ACCOUNT_ID', * accessKeyId: 'YOUR_R2_ACCESS_KEY_ID', * secretAccessKey: 'YOUR_R2_SECRET_ACCESS_KEY' * }) * ``` */ export class R2Storage extends BaseStorage { private s3Client: S3Client | null = null private bucketName: string private accountId: string private accessKeyId: string private secretAccessKey: string // R2-specific endpoint (auto-constructed from account ID) private endpoint: string // Prefixes for different types of data private nounPrefix: string private verbPrefix: string private metadataPrefix: string // Noun metadata private verbMetadataPrefix: string // Verb metadata private systemPrefix: string // System data // Statistics caching for better performance protected statisticsCache: StatisticsData | null = null // Backpressure and performance management private pendingOperations: number = 0 private maxConcurrentOperations: number = 150 // R2 handles more concurrent ops private baseBatchSize: number = 15 // Larger batches for R2 private currentBatchSize: number = 15 private lastMemoryCheck: number = 0 private memoryCheckInterval: number = 5000 // Adaptive backpressure for automatic flow control private backpressure = getGlobalBackpressure() // Write buffers for bulk operations private nounWriteBuffer: WriteBuffer | null = null private verbWriteBuffer: WriteBuffer | null = null // Request coalescer for deduplication private requestCoalescer: RequestCoalescer | null = null // High-volume mode detection (R2-specific thresholds) private highVolumeMode = false private lastVolumeCheck = 0 private volumeCheckInterval = 800 // Check more frequently on R2 private forceHighVolumeMode = false // Multi-level cache manager for efficient data access private nounCacheManager: CacheManager private verbCacheManager: CacheManager // Module logger private logger = createModuleLogger('R2Storage') // v5.4.0: HNSW mutex locks to prevent read-modify-write races private hnswLocks = new Map>() /** * Initialize the R2 storage adapter * @param options Configuration options for Cloudflare R2 */ constructor(options: { bucketName: string accountId: string accessKeyId: string secretAccessKey: string // Optional configuration cacheConfig?: { hotCacheMaxSize?: number hotCacheEvictionThreshold?: number warmCacheTTL?: number } readOnly?: boolean }) { super() this.bucketName = options.bucketName this.accountId = options.accountId this.accessKeyId = options.accessKeyId this.secretAccessKey = options.secretAccessKey this.readOnly = options.readOnly || false // R2-specific endpoint format this.endpoint = `https://${this.accountId}.r2.cloudflarestorage.com` // Set up prefixes for different types of data using entity-based structure this.nounPrefix = `${getDirectoryPath('noun', 'vector')}/` this.verbPrefix = `${getDirectoryPath('verb', 'vector')}/` this.metadataPrefix = `${getDirectoryPath('noun', 'metadata')}/` this.verbMetadataPrefix = `${getDirectoryPath('verb', 'metadata')}/` this.systemPrefix = `${SYSTEM_DIR}/` // Initialize cache managers with R2-optimized settings this.nounCacheManager = new CacheManager({ hotCacheMaxSize: options.cacheConfig?.hotCacheMaxSize || 10000, hotCacheEvictionThreshold: options.cacheConfig?.hotCacheEvictionThreshold || 0.9, warmCacheTTL: options.cacheConfig?.warmCacheTTL || 3600000 // 1 hour }) this.verbCacheManager = new CacheManager(options.cacheConfig) // Check for high-volume mode override if (typeof process !== 'undefined' && process.env?.BRAINY_FORCE_HIGH_VOLUME === 'true') { this.forceHighVolumeMode = true this.highVolumeMode = true prodLog.info('๐Ÿš€ R2: High-volume mode FORCED via environment variable') } } /** * Get R2-optimized batch configuration * * Cloudflare R2 has S3-compatible characteristics with some advantages: * - Zero egress fees (can cache more aggressively) * - Global edge network * - Similar throughput to S3 * * R2 benefits from the same configuration as S3: * - Larger batch sizes (100 items) * - Parallel processing * - Short delays (50ms) * * @returns R2-optimized batch configuration * @since v4.11.0 */ public getBatchConfig(): StorageBatchConfig { return { maxBatchSize: 100, batchDelayMs: 50, maxConcurrent: 100, supportsParallelWrites: true, // R2 handles parallel writes like S3 rateLimit: { operationsPerSecond: 3500, // Similar to S3 throughput burstCapacity: 1000 } } } /** * Initialize the storage adapter */ public async init(): Promise { if (this.isInitialized) { return } try { // Import AWS S3 SDK only when needed (R2 uses S3-compatible API) const { S3Client: S3ClientClass, HeadBucketCommand } = await import('@aws-sdk/client-s3') // Create S3 client configured for R2 this.s3Client = new S3ClientClass({ region: 'auto', // R2 uses 'auto' region endpoint: this.endpoint, credentials: { accessKeyId: this.accessKeyId, secretAccessKey: this.secretAccessKey } }) // Verify bucket exists and is accessible try { await this.s3Client.send(new HeadBucketCommand({ Bucket: this.bucketName })) } catch (error: any) { if (error.name === 'NotFound' || error.$metadata?.httpStatusCode === 404) { throw new Error(`R2 bucket ${this.bucketName} does not exist or is not accessible`) } throw error } prodLog.info(`โœ… Connected to R2 bucket: ${this.bucketName} (account: ${this.accountId})`) // Initialize write buffers for high-volume mode const storageId = `r2-${this.bucketName}` this.nounWriteBuffer = getWriteBuffer( `${storageId}-nouns`, 'noun', async (items) => { await this.flushNounBuffer(items) } ) this.verbWriteBuffer = getWriteBuffer( `${storageId}-verbs`, 'verb', async (items) => { await this.flushVerbBuffer(items) } ) // Initialize request coalescer for deduplication this.requestCoalescer = getCoalescer( storageId, async (batch) => { this.logger.trace(`Processing coalesced batch: ${batch.length} items`) } ) // Initialize counts from storage await this.initializeCounts() // Clear cache from previous runs prodLog.info('๐Ÿงน R2: Clearing cache from previous run') this.nounCacheManager.clear() this.verbCacheManager.clear() this.isInitialized = true } catch (error) { this.logger.error('Failed to initialize R2 storage:', error) throw new Error(`Failed to initialize R2 storage: ${error}`) } } /** * Get the R2 object key for a noun using UUID-based sharding */ private getNounKey(id: string): string { const shardId = getShardIdFromUuid(id) return `${this.nounPrefix}${shardId}/${id}.json` } /** * Get the R2 object key for a verb using UUID-based sharding */ private getVerbKey(id: string): string { const shardId = getShardIdFromUuid(id) return `${this.verbPrefix}${shardId}/${id}.json` } /** * Override base class method to detect R2-specific throttling errors */ protected isThrottlingError(error: any): boolean { // First check base class detection if (super.isThrottlingError(error)) { return true } // R2-specific throttling detection (uses S3 error codes) const errorName = error.name const statusCode = error.$metadata?.httpStatusCode return ( errorName === 'SlowDown' || errorName === 'ServiceUnavailable' || statusCode === 429 || statusCode === 503 ) } /** * Override base class to enable smart batching for cloud storage * R2 is cloud storage with network latency benefits from batching */ protected isCloudStorage(): boolean { return true } /** * Apply backpressure before starting an operation */ private async applyBackpressure(): Promise { const requestId = `${Date.now()}-${Math.random().toString(36).substr(2, 9)}` await this.backpressure.requestPermission(requestId, 1) this.pendingOperations++ return requestId } /** * Release backpressure after completing an operation */ private releaseBackpressure(success: boolean = true, requestId?: string): void { this.pendingOperations = Math.max(0, this.pendingOperations - 1) if (requestId) { this.backpressure.releasePermission(requestId, success) } } /** * Check if high-volume mode should be enabled */ private checkVolumeMode(): void { if (this.forceHighVolumeMode) { return } const now = Date.now() if (now - this.lastVolumeCheck < this.volumeCheckInterval) { return } this.lastVolumeCheck = now // R2 threshold: enable at 15 pending operations (lower than S3/GCS) const shouldEnable = this.pendingOperations > 15 if (shouldEnable && !this.highVolumeMode) { this.highVolumeMode = true prodLog.info('๐Ÿš€ R2: High-volume mode ENABLED (pending:', this.pendingOperations, ')') } else if (!shouldEnable && this.highVolumeMode && !this.forceHighVolumeMode) { this.highVolumeMode = false prodLog.info('๐ŸŒ R2: High-volume mode DISABLED (pending:', this.pendingOperations, ')') } } /** * Flush noun buffer to R2 */ private async flushNounBuffer(items: Map): Promise { const writes = Array.from(items.values()).map(async (noun) => { try { await this.saveNodeDirect(noun) } catch (error) { this.logger.error(`Failed to flush noun ${noun.id}:`, error) } }) await Promise.all(writes) } /** * Flush verb buffer to R2 */ private async flushVerbBuffer(items: Map): Promise { const writes = Array.from(items.values()).map(async (verb) => { try { await this.saveEdgeDirect(verb) } catch (error) { this.logger.error(`Failed to flush verb ${verb.id}:`, error) } }) await Promise.all(writes) } /** * Save a node to storage */ protected async saveNode(node: HNSWNode): Promise { await this.ensureInitialized() this.checkVolumeMode() // Use write buffer in high-volume mode if (this.highVolumeMode && this.nounWriteBuffer) { this.logger.trace(`๐Ÿ“ BUFFERING: Adding noun ${node.id} to write buffer`) await this.nounWriteBuffer.add(node.id, node) return } // Direct write in normal mode await this.saveNodeDirect(node) } /** * Save a node directly to R2 (bypass buffer) */ private async saveNodeDirect(node: HNSWNode): Promise { const requestId = await this.applyBackpressure() try { this.logger.trace(`Saving node ${node.id}`) // Convert connections Map to serializable format const serializableNode = { id: node.id, vector: node.vector, connections: Object.fromEntries( Array.from(node.connections.entries()).map(([level, nounIds]) => [ level, Array.from(nounIds) ]) ), level: node.level || 0 } // Get the R2 key with UUID-based sharding const key = this.getNounKey(node.id) // Save to R2 using S3 PutObject const { PutObjectCommand } = await import('@aws-sdk/client-s3') await this.s3Client!.send( new PutObjectCommand({ Bucket: this.bucketName, Key: key, Body: JSON.stringify(serializableNode, null, 2), ContentType: 'application/json' }) ) // Cache nodes with non-empty vectors (Phase 2 optimization) if (node.vector && Array.isArray(node.vector) && node.vector.length > 0) { this.nounCacheManager.set(node.id, node) } // Increment noun count const metadata = await this.getNounMetadata(node.id) if (metadata && metadata.type) { await this.incrementEntityCountSafe(metadata.type as string) } this.logger.trace(`Node ${node.id} saved successfully`) this.releaseBackpressure(true, requestId) } catch (error: any) { this.releaseBackpressure(false, requestId) if (this.isThrottlingError(error)) { await this.handleThrottling(error) throw error } this.logger.error(`Failed to save node ${node.id}:`, error) throw new Error(`Failed to save node ${node.id}: ${error}`) } } /** * Get a node from storage */ protected async getNode(id: string): Promise { await this.ensureInitialized() // Check cache first (Phase 2: aggressive caching for R2 zero-egress) const cached = await this.nounCacheManager.get(id) if (cached !== undefined && cached !== null) { if (!cached.id || !cached.vector || !Array.isArray(cached.vector) || cached.vector.length === 0) { this.logger.warn(`Invalid cached object for ${id.substring(0, 8)} - removing from cache`) this.nounCacheManager.delete(id) } else { this.logger.trace(`Cache hit for noun ${id}`) return cached } } const requestId = await this.applyBackpressure() try { this.logger.trace(`Getting node ${id}`) const key = this.getNounKey(id) // Get from R2 using S3 GetObject const { GetObjectCommand } = await import('@aws-sdk/client-s3') const response = await this.s3Client!.send( new GetObjectCommand({ Bucket: this.bucketName, Key: key }) ) const bodyContents = await response.Body!.transformToString() const data = JSON.parse(bodyContents) // Convert serialized connections back to Map const connections = new Map>() for (const [level, nounIds] of Object.entries(data.connections || {})) { connections.set(Number(level), new Set(nounIds as string[])) } const node: HNSWNode = { id: data.id, vector: data.vector, connections, level: data.level || 0 } // Cache valid nodes with non-empty vectors if (node && node.id && node.vector && Array.isArray(node.vector) && node.vector.length > 0) { this.nounCacheManager.set(id, node) } this.logger.trace(`Successfully retrieved node ${id}`) this.releaseBackpressure(true, requestId) return node } catch (error: any) { this.releaseBackpressure(false, requestId) // R2 returns NoSuchKey for 404 if (error.name === 'NoSuchKey' || error.$metadata?.httpStatusCode === 404) { return null } if (this.isThrottlingError(error)) { await this.handleThrottling(error) throw error } this.logger.error(`Failed to get node ${id}:`, error) throw BrainyError.fromError(error, `getNoun(${id})`) } } /** * Write an object to a specific path in R2 */ protected async writeObjectToPath(path: string, data: any): Promise { await this.ensureInitialized() try { this.logger.trace(`Writing object to path: ${path}`) const { PutObjectCommand } = await import('@aws-sdk/client-s3') await this.s3Client!.send( new PutObjectCommand({ Bucket: this.bucketName, Key: path, Body: JSON.stringify(data, null, 2), ContentType: 'application/json' }) ) this.logger.trace(`Object written successfully to ${path}`) } catch (error) { this.logger.error(`Failed to write object to ${path}:`, error) throw new Error(`Failed to write object to ${path}: ${error}`) } } /** * Read an object from a specific path in R2 */ protected async readObjectFromPath(path: string): Promise { await this.ensureInitialized() try { this.logger.trace(`Reading object from path: ${path}`) const { GetObjectCommand } = await import('@aws-sdk/client-s3') const response = await this.s3Client!.send( new GetObjectCommand({ Bucket: this.bucketName, Key: path }) ) const bodyContents = await response.Body!.transformToString() const data = JSON.parse(bodyContents) this.logger.trace(`Object read successfully from ${path}`) return data } catch (error: any) { if (error.name === 'NoSuchKey' || error.$metadata?.httpStatusCode === 404) { this.logger.trace(`Object not found at ${path}`) return null } this.logger.error(`Failed to read object from ${path}:`, error) throw BrainyError.fromError(error, `readObjectFromPath(${path})`) } } /** * Delete an object from a specific path in R2 */ protected async deleteObjectFromPath(path: string): Promise { await this.ensureInitialized() try { this.logger.trace(`Deleting object at path: ${path}`) const { DeleteObjectCommand } = await import('@aws-sdk/client-s3') await this.s3Client!.send( new DeleteObjectCommand({ Bucket: this.bucketName, Key: path }) ) this.logger.trace(`Object deleted successfully from ${path}`) } catch (error: any) { if (error.name === 'NoSuchKey' || error.$metadata?.httpStatusCode === 404) { this.logger.trace(`Object at ${path} not found (already deleted)`) return } this.logger.error(`Failed to delete object from ${path}:`, error) throw new Error(`Failed to delete object from ${path}: ${error}`) } } /** * List all objects under a specific prefix in R2 */ protected async listObjectsUnderPath(prefix: string): Promise { await this.ensureInitialized() try { this.logger.trace(`Listing objects under prefix: ${prefix}`) const { ListObjectsV2Command } = await import('@aws-sdk/client-s3') const response = await this.s3Client!.send( new ListObjectsV2Command({ Bucket: this.bucketName, Prefix: prefix, MaxKeys: MAX_R2_PAGE_SIZE }) ) const paths = (response.Contents || []) .map((obj: any) => obj.Key) .filter((key: string) => key && key.length > 0) this.logger.trace(`Found ${paths.length} objects under ${prefix}`) return paths } catch (error) { this.logger.error(`Failed to list objects under ${prefix}:`, error) throw new Error(`Failed to list objects under ${prefix}: ${error}`) } } // Verb storage methods (similar to noun methods - implementing key methods for space) protected async saveEdge(edge: Edge): Promise { await this.ensureInitialized() this.checkVolumeMode() if (this.highVolumeMode && this.verbWriteBuffer) { await this.verbWriteBuffer.add(edge.id, edge) return } await this.saveEdgeDirect(edge) } private async saveEdgeDirect(edge: Edge): Promise { const requestId = await this.applyBackpressure() try { // ARCHITECTURAL FIX (v3.50.1): Include core relational fields in verb vector file // These fields are essential for 90% of operations - no metadata lookup needed const serializableEdge = { id: edge.id, vector: edge.vector, connections: Object.fromEntries( Array.from(edge.connections.entries()).map(([level, verbIds]) => [ level, Array.from(verbIds) ]) ), // CORE RELATIONAL DATA (v3.50.1+) verb: edge.verb, sourceId: edge.sourceId, targetId: edge.targetId, // User metadata (if any) - saved separately for scalability // metadata field is saved separately via saveVerbMetadata() } const key = this.getVerbKey(edge.id) const { PutObjectCommand } = await import('@aws-sdk/client-s3') await this.s3Client!.send( new PutObjectCommand({ Bucket: this.bucketName, Key: key, Body: JSON.stringify(serializableEdge, null, 2), ContentType: 'application/json' }) ) this.verbCacheManager.set(edge.id, edge) // Count tracking happens in baseStorage.saveVerbMetadata_internal (v4.1.2) // This fixes the race condition where metadata didn't exist yet this.releaseBackpressure(true, requestId) } catch (error: any) { this.releaseBackpressure(false, requestId) if (this.isThrottlingError(error)) { await this.handleThrottling(error) throw error } throw new Error(`Failed to save edge ${edge.id}: ${error}`) } } protected async getEdge(id: string): Promise { await this.ensureInitialized() const cached = this.verbCacheManager.get(id) if (cached) { return cached } const requestId = await this.applyBackpressure() try { const key = this.getVerbKey(id) const { GetObjectCommand } = await import('@aws-sdk/client-s3') const response = await this.s3Client!.send( new GetObjectCommand({ Bucket: this.bucketName, Key: key }) ) const bodyContents = await response.Body!.transformToString() const data = JSON.parse(bodyContents) const connections = new Map>() for (const [level, verbIds] of Object.entries(data.connections || {})) { connections.set(Number(level), new Set(verbIds as string[])) } // v4.0.0: Return HNSWVerb with core relational fields (NO metadata field) const edge: Edge = { id: data.id, vector: data.vector, connections, // CORE RELATIONAL DATA (read from vector file) verb: data.verb, sourceId: data.sourceId, targetId: data.targetId // โœ… NO metadata field in v4.0.0 // User metadata retrieved separately via getVerbMetadata() } this.verbCacheManager.set(id, edge) this.releaseBackpressure(true, requestId) return edge } catch (error: any) { this.releaseBackpressure(false, requestId) if (error.name === 'NoSuchKey' || error.$metadata?.httpStatusCode === 404) { return null } if (this.isThrottlingError(error)) { await this.handleThrottling(error) throw error } throw BrainyError.fromError(error, `getVerb(${id})`) } } // Pagination and count management (simplified for space - full implementation similar to GCS) protected async initializeCounts(): Promise { const key = `${this.systemPrefix}counts.json` try { const counts = await this.readObjectFromPath(key) if (counts) { this.totalNounCount = counts.totalNounCount || 0 this.totalVerbCount = counts.totalVerbCount || 0 this.entityCounts = new Map(Object.entries(counts.entityCounts || {})) as Map this.verbCounts = new Map(Object.entries(counts.verbCounts || {})) as Map prodLog.info(`๐Ÿ“Š R2: Loaded counts: ${this.totalNounCount} nouns, ${this.totalVerbCount} verbs`) } else { prodLog.info('๐Ÿ“Š R2: No counts file found - initializing from scan') await this.initializeCountsFromScan() } } catch (error) { prodLog.error('โŒ R2: Failed to load counts:', error) await this.initializeCountsFromScan() } } private async initializeCountsFromScan(): Promise { try { prodLog.info('๐Ÿ“Š R2: Scanning bucket to initialize counts...') const { ListObjectsV2Command } = await import('@aws-sdk/client-s3') // Count nouns const nounResponse = await this.s3Client!.send( new ListObjectsV2Command({ Bucket: this.bucketName, Prefix: this.nounPrefix }) ) this.totalNounCount = (nounResponse.Contents || []).filter((obj: any) => obj.Key?.endsWith('.json') ).length // Count verbs const verbResponse = await this.s3Client!.send( new ListObjectsV2Command({ Bucket: this.bucketName, Prefix: this.verbPrefix }) ) this.totalVerbCount = (verbResponse.Contents || []).filter((obj: any) => obj.Key?.endsWith('.json') ).length if (this.totalNounCount > 0 || this.totalVerbCount > 0) { await this.persistCounts() prodLog.info(`โœ… R2: Initialized counts: ${this.totalNounCount} nouns, ${this.totalVerbCount} verbs`) } else { prodLog.warn('โš ๏ธ R2: No entities found during bucket scan') } } catch (error) { this.logger.error('โŒ R2: Failed to initialize counts from scan:', error) throw new Error(`Failed to initialize R2 storage counts: ${error}`) } } protected async persistCounts(): Promise { try { const key = `${this.systemPrefix}counts.json` const counts = { totalNounCount: this.totalNounCount, totalVerbCount: this.totalVerbCount, entityCounts: Object.fromEntries(this.entityCounts), verbCounts: Object.fromEntries(this.verbCounts), lastUpdated: new Date().toISOString() } await this.writeObjectToPath(key, counts) } catch (error) { this.logger.error('Error persisting counts:', error) } } // HNSW Index Persistence (Phase 2 support) public async getNounVector(id: string): Promise { const noun = await this.getNoun(id) return noun ? noun.vector : null } public async saveHNSWData(nounId: string, hnswData: { level: number connections: Record }): Promise { const lockKey = `hnsw/${nounId}` // Wait for pending operations 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 existingNoun = await this.getNoun(nounId) if (!existingNoun) { throw new Error(`Cannot save HNSW data: noun ${nounId} not found`) } const connectionsMap = new Map>() for (const [level, nodeIds] of Object.entries(hnswData.connections)) { connectionsMap.set(Number(level), new Set(nodeIds)) } const updatedNoun: HNSWNoun = { ...existingNoun, level: hnswData.level, connections: connectionsMap } await this.saveNoun(updatedNoun) } finally { this.hnswLocks.delete(lockKey) releaseLock() } } public async getHNSWData(nounId: string): Promise<{ level: number connections: Record } | null> { const noun = await this.getNoun(nounId) if (!noun) { return null } 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 } } public async saveHNSWSystem(systemData: { entryPointId: string | null maxLevel: number }): Promise { await this.ensureInitialized() const key = `${this.systemPrefix}hnsw-system.json` await this.writeObjectToPath(key, systemData) } public async getHNSWSystem(): Promise<{ entryPointId: string | null maxLevel: number } | null> { await this.ensureInitialized() const key = `${this.systemPrefix}hnsw-system.json` return await this.readObjectFromPath(key) } // Statistics support protected async saveStatisticsData(statistics: StatisticsData): Promise { await this.ensureInitialized() const key = `${this.systemPrefix}${STATISTICS_KEY}.json` await this.writeObjectToPath(key, statistics) } protected async getStatisticsData(): Promise { await this.ensureInitialized() const key = `${this.systemPrefix}${STATISTICS_KEY}.json` const stats = await this.readObjectFromPath(key) if (stats) { return { ...stats, totalNodes: this.totalNounCount, totalEdges: this.totalVerbCount, lastUpdated: new Date().toISOString() } } return { nounCount: {}, verbCount: {}, metadataCount: {}, hnswIndexSize: 0, totalNodes: this.totalNounCount, totalEdges: this.totalVerbCount, totalMetadata: 0, lastUpdated: new Date().toISOString() } } // Utility methods public async clear(): Promise { await this.ensureInitialized() prodLog.info('๐Ÿงน R2: Clearing all data from bucket...') // Clear all prefixes (v5.6.1: includes _cow/ for version control data) // _cow/ stores all git-like versioning data (commits, trees, blobs, refs) // Must be deleted to fully clear all data including version history for (const prefix of [this.nounPrefix, this.verbPrefix, this.metadataPrefix, this.verbMetadataPrefix, this.systemPrefix, '_cow/']) { const objects = await this.listObjectsUnderPath(prefix) for (const key of objects) { await this.deleteObjectFromPath(key) } } // CRITICAL: Reset COW state to prevent automatic reinitialization // When COW data is cleared, we must also clear the COW managers // Otherwise initializeCOW() will auto-recreate initial commit on next operation this.refManager = undefined this.blobStorage = undefined this.commitLog = undefined this.cowEnabled = false // v5.10.4: Create persistent marker object (CRITICAL FIX) // Bug: cowEnabled = false only affects current instance, not future instances // Fix: Create marker object that persists across instance restarts // When new instance calls initializeCOW(), it checks for this marker await this.createClearMarker() this.nounCacheManager.clear() this.verbCacheManager.clear() this.totalNounCount = 0 this.totalVerbCount = 0 this.entityCounts.clear() this.verbCounts.clear() prodLog.info('โœ… R2: All data cleared') } public async getStorageStatus(): Promise<{ type: string used: number quota: number | null details?: Record }> { return { type: 'r2', used: 0, quota: null, details: { bucket: this.bucketName, accountId: this.accountId, endpoint: this.endpoint, features: [ 'Zero egress fees', 'Global edge network', 'S3-compatible API', 'Type-aware HNSW support' ] } } } /** * Check if COW has been explicitly disabled via clear() * v5.10.4: Fixes bug where clear() doesn't persist across instance restarts * @returns true if marker object exists, false otherwise * @protected */ protected async checkClearMarker(): Promise { await this.ensureInitialized() try { const markerPath = `${this.systemPrefix}cow-disabled` const data = await this.readObjectFromPath(markerPath) return data !== null // Marker exists if we got any data } catch (error) { prodLog.warn('R2Storage.checkClearMarker: Error checking marker', error) return false } } /** * Create marker indicating COW has been explicitly disabled * v5.10.4: Called by clear() to prevent COW reinitialization on new instances * @protected */ protected async createClearMarker(): Promise { await this.ensureInitialized() try { const markerPath = `${this.systemPrefix}cow-disabled` // Create empty marker object await this.writeObjectToPath(markerPath, '') } catch (error) { prodLog.error('R2Storage.createClearMarker: Failed to create marker object', error) // Don't throw - marker creation failure shouldn't break clear() } } // v5.4.0: Removed getNounsWithPagination override - use BaseStorage's type-first implementation // v5.4.0: Removed 10 *_internal method overrides - now inherit from BaseStorage's type-first implementation }